import sys
sys.path.insert(0, '/opt/struktur/obsidian-graphiti-import/aa043')
import sqlite3
import hashlib
import json
import datetime
from classifier import Classifier
def inventory_getter():
return {}
db_path = '/opt/struktur/youtube-research/knowledge.db'
conn = sqlite3.connect(db_path)
conn.row_factory = sqlite3.Row
c = conn.cursor()
# Get open AA-043 rows with needed fields
c.execute("""
SELECT
g.id,
g.reconciliation_record_id,
g.obsidian_path,
g.source_version_id,
g.graphiti_status,
sv.content_version_hash,
sv.import_identity,
rd.source_identity,
rd.canonical_obsidian_path,
rd.promotion_policy as existing_promotion_policy,
rd.provenance_json
FROM graphiti_import_queue g
JOIN source_versions sv ON g.source_version_id = sv.version_id
JOIN reconciliation_decisions rd ON
rd.source_version_id = sv.version_id
AND rd.content_version_hash = sv.content_version_hash
AND rd.import_identity = sv.import_identity
WHERE g.reconciliation_record_id LIKE 'aa043:%'
AND g.graphiti_status != 'done'
""")
rows = c.fetchall()
print(f'Open AA-043 rows: {len(rows)}')
# Initial promotion distribution
c.execute("""
SELECT promotion_policy, COUNT(*)
FROM graphiti_import_queue
WHERE reconciliation_record_id LIKE 'aa043:%' AND graphiti_status != 'done'
GROUP BY promotion_policy
""")
print('Initial promotion distribution:', c.fetchall())
# Initial knowledge units count
c.execute("SELECT COUNT(*) FROM knowledge_units")
print(f'Initial knowledge units: {c.fetchone()[0]}')
# Initialize classifier
classifier = Classifier(db_path=db_path, inventory_getter=inventory_getter)
updates = []
knowledge_units = []
for r in rows:
rid = r['reconciliation_record_id']
source_identity = r['source_identity']
content_version_hash = r['content_version_hash']
obsidian_path = r['obsidian_path']
canonical_obsidian_path = r['canonical_obsidian_path']
import_identity = r['import_identity']
provenance_json = r['provenance_json']
# Parse provenance_json to dict if not null
import json
try:
provenance = json.loads(provenance_json) if provenance_json else {}
except:
provenance = {}
# Call classify_standalone with required args
result = classifier.classify_standalone(
reconciliation_record_id=rid,
source_identity=source_identity,
content_version_hash=content_version_hash,
canonical_obsidian_path=canonical_obsidian_path,
import_identity=import_identity,
current_provenance=provenance,
conn=conn
)
policy = result.get('policy', 'REVIEW_REQUIRED')
reason = result.get('reason', 'Default')
kus = result.get('knowledge_units', [])
updates.append((rid, policy, reason, datetime.datetime.now().isoformat()))
for ku in kus:
stmt = ku.get('statement', '')
ent = json.dumps(ku.get('entities', []), sort_keys=True)
rel = json.dumps(ku.get('relations', []), sort_keys=True)
prov = json.dumps(ku.get('provenance', {}), sort_keys=True)
# Build deterministic identity per spec: statement + source_identity + source_version_id + content_version_hash
hash_input = f"{stmt}|{source_identity}|{r['source_version_id']}|{content_version_hash}"
kuid = hashlib.sha256(hash_input.encode('utf-8')).hexdigest()
knowledge_units.append({
'knowledge_unit_identity': kuid,
'statement': stmt,
'entities': ent,
'relations': rel,
'provenance': prov,
'promotion_reason': reason,
'source_identity': source_identity,
'source_version_id': r['source_version_id'],
'content_version_hash': content_version_hash,
'created_at': datetime.datetime.now().isoformat()
})
# Apply updates
for rid, policy, reason, decided_at in updates:
c.execute("""
UPDATE graphiti_import_queue
SET promotion_policy = ?, promotion_reason = ?, promotion_decided_at = ?
WHERE reconciliation_record_id = ?
""", (policy, reason, decided_at, rid))
# Insert knowledge units (ignore duplicates)
for ku in knowledge_units:
c.execute("""
INSERT OR IGNORE INTO knowledge_units
(knowledge_unit_identity, statement, entities, relations, provenance, promotion_reason,
source_identity, source_version_id, content_version_hash, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""", (ku['knowledge_unit_identity'], ku['statement'], ku['entities'], ku['relations'],
ku['provenance'], ku['promotion_reason'], ku['source_identity'], ku['source_version_id'],
ku['content_version_hash'], ku['created_at']))
conn.commit()
print(f'Applied {len(updates)} promotion updates')
print(f'Inserted {len(knowledge_units)} knowledge units (duplicates ignored)')
# Final promotion distribution
c.execute("""
SELECT promotion_policy, COUNT(*)
FROM graphiti_import_queue
WHERE reconciliation_record_id LIKE 'aa043:%' AND graphiti_status != 'done'
GROUP BY promotion_policy
""")
print('Final promotion distribution:', c.fetchall())
# Final knowledge units stats
c.execute("SELECT COUNT(*) FROM knowledge_units")
print(f'Final knowledge units total: {c.fetchone()[0]}')
c.execute("SELECT COUNT(DISTINCT knowledge_unit_identity) FROM knowledge_units")
print(f'Final distinct knowledge unit identities: {c.fetchone()[0]}')
# List some knowledge units for verification (limit 5)
c.execute("""
SELECT ku.knowledge_unit_identity, ku.statement, ku.entities, ku.relations, ku.provenance, ku.promotion_reason,
ku.source_identity, ku.source_version_id, ku.content_version_hash
FROM knowledge_units ku
ORDER BY ku.id DESC
LIMIT 5
""")
print('Sample knowledge units:')
for row in c.fetchall():
print(row)
conn.close()
classifier.close()