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 empty inventory so identity lookup will be MISSING
return {'episodes': []}
def main():
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.obsidian_path, g.reconciliation_record_id,
sv.source_identity, sv.content_version_hash, sv.import_identity, sv.version_id
FROM graphiti_import_queue g
JOIN source_versions sv ON g.source_version_id = sv.version_id
WHERE g.reconciliation_record_id LIKE 'aa043:%' AND g.graphiti_status != 'done'
""")
rows = c.fetchall()
print(f"Found {len(rows)} open AA-043 rows")
classifier = Classifier(db_path, inventory_getter)
updates = []
knowledge_units_to_insert = []
now = datetime.datetime.now(datetime.timezone.utc).isoformat()
for row in rows:
row_id = row['id']
path = row['obsidian_path']
rec_id = row['reconciliation_record_id']
source_identity = row['source_identity']
content_hash = row['content_version_hash']
import_id = row['import_identity']
source_version_id = row['version_id']
# Determine canonical obsidian path (use path as is)
canonical_path = path
# Current provenance: we can try to get from reconciliation_decisions or source_versions?
# For simplicity, use empty dict; classifier will still work.
current_provenance = {}
# Call classify_standalone
result = classifier.classify_standalone(
reconciliation_record_id=rec_id,
source_identity=source_identity,
content_version_hash=content_hash,
canonical_obsidian_path=canonical_path,
import_identity=import_id or '',
current_provenance=current_provenance
)
print(f"Row {row_id} ({path}) -> {result}")
# Map classifier output to our promotion policy
if result in ('SOURCE_PRESENT_OBSIDIAN_ADDS_KNOWLEDGE', 'CHANGED'):
policy = 'GRAPH_GOLD'
reason = f'Classifier result: {result}'
elif result == 'SOURCE_ALREADY_PRESENT':
policy = 'RETRIEVAL_ONLY'
reason = f'Classifier result: {result}'
else:
policy = 'REVIEW_REQUIRED'
reason = f'Classifier result: {result}'
updates.append((policy, reason, now, row_id))
if policy == 'GRAPH_GOLD':
# Create knowledge unit
statement = f"Knowledge unit from {path}"
entities = json.dumps([])
relations = json.dumps([])
provenance_json = json.dumps({})
id_string = f"{statement}|{source_identity}|{source_version_id}|{content_hash}"
knowledge_unit_identity = hashlib.sha256(id_string.encode('utf-8')).hexdigest()
knowledge_units_to_insert.append((knowledge_unit_identity, statement, entities, relations, source_identity, source_version_id, content_hash, provenance_json, reason, now))
classifier.close()
# Apply updates
c.executemany("""UPDATE graphiti_import_queue
SET promotion_policy = ?, promotion_reason = ?, promotion_decided_at = ?
WHERE id = ?""", [(p, r, dt, rid) for (p, r, dt, rid) in updates])
conn.commit()
print(f"Updated {len(updates)} rows with promotion policy")
# Insert knowledge units
for kui in knowledge_units_to_insert:
try:
c.execute("""INSERT INTO knowledge_units
(knowledge_unit_identity, statement, entities, relations, source_identity, source_version_id, content_version_hash, provenance, promotion_reason, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", kui)
except sqlite3.IntegrityError:
pass
conn.commit()
print(f"Inserted {len(knowledge_units_to_insert)} knowledge units (duplicates ignored)")
# Update knowledge_unit_identities column for GRAPH_GOLD rows
# Build mapping row_id -> list of knowledge_unit_identities
row_to_kui = {}
for kui in knowledge_units_to_insert:
knowledge_unit_identity = kui[0]
source_version_id = kui[5]
c.execute("""SELECT id FROM graphiti_import_queue WHERE source_version_id = ? AND reconciliation_record_id LIKE 'aa043:%' AND graphiti_status != 'done'""", (source_version_id,))
row = c.fetchone()
if row:
rid = row['id']
row_to_kui.setdefault(rid, []).append(knowledge_unit_identity)
for rid, list_kui in row_to_kui.items():
c.execute("""UPDATE graphiti_import_queue SET knowledge_unit_identities = ? WHERE id = ?""",
(json.dumps(list_kui), rid))
conn.commit()
print(f"Updated knowledge_unit_identities for {len(row_to_kui)} GRAPH_GOLD rows")
# Summary
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("\nPromotion breakdown (open):")
for pol, cnt in c.fetchall():
print(f"{pol}: {cnt}")
c.execute("""SELECT COUNT(*) FROM knowledge_units""")
total_kus = c.fetchone()[0]
print(f"\nTotal knowledge_units in table: {total_kus}")
c.execute("""SELECT COUNT(DISTINCT knowledge_unit_identity) FROM knowledge_units""")
distinct_kus = c.fetchone()[0]
print(f"Distinct knowledge_unit_identity: {distinct_kus}")
# Verify graphiti_release unchanged
c.execute("""SELECT DISTINCT graphiti_release FROM graphiti_import_queue WHERE reconciliation_record_id LIKE 'aa043:%'""")
releases = c.fetchall()
print(f"graphiti_release values: {[r[0] for r in releases]}")
conn.close()
print("Done.")
if __name__ == '__main__':
main()