#!/usr/bin/env python3
import argparse, importlib.util, json, sqlite3, sys
from datetime import datetime, timezone
from pathlib import Path
DB='/opt/struktur/youtube-research/knowledge.db'
EX='/opt/struktur/AA-043-P33-extractor.py'
ALLOW=['youtube:FVuqEATzMWI','youtube:qmPAFHeiU1A','youtube:z-C_q4LzPlA','youtube:3yqQ664R1hg','obsidian:Systemstruktur/SYSTEM-REFERENZ.md','obsidian:Systemstruktur/Review-Routing-Decision-Units.md','obsidian:Systemstruktur/Knowledge-Pipeline-Review-Unter-100-2026-07-22.md','obsidian:05-Ressourcen/Hermes-Sessions-2026-07/TranscriptBackfillWorker-untersuchen.md','youtube:fxc4yds-DUs','youtube:ZyNdWhKqlEs']
REVIEWER='AA-043-P34-R2-source-only-review-v1'
def load_extractor():
spec=importlib.util.spec_from_file_location('p33_extractor',EX); m=importlib.util.module_from_spec(spec); spec.loader.exec_module(m); return m
def main():
ap=argparse.ArgumentParser(); ap.add_argument('--dry-run',action='store_true'); a=ap.parse_args(); m=load_extractor(); sources=m.select_sources(); ids=[s['source_id'] for s in sources]
if ids!=ALLOW: raise RuntimeError('allowlist_mismatch:'+json.dumps({'expected':ALLOW,'actual':ids},ensure_ascii=False))
candidates=[]; rejects=[]
for s in sources:
for st,loc in m.sentence_candidates(s):
ev,reason=m.classify(st); k={'statement':st,'source_id':s['source_id'],'content_version':s['content_version'],'source_locator':loc,'evidence_class':ev,'confidence_reason':reason,'knowledge_unit_identity':m.identity(s['source_id'],s['content_version'],st)}
ok,rr=m.valid(k)
(candidates if ok else rejects).append((s,k,rr))
# same conservative dedup as P33, cross-source semantic duplicates retained only if exact identity differs and similarity is not >= .88
accepted=[]; duplicates=[]
for s,k,_ in candidates:
dup=None
for ps,pk in accepted:
if pk['knowledge_unit_identity']==k['knowledge_unit_identity'] or m.jaccard(pk['statement'],k['statement'])>=.88: dup=(ps,pk); break
if dup: duplicates.append((s,k,'exact_or_semantic_duplicate'))
else: accepted.append((s,k))
summary={'allowlist':ALLOW,'source_count':len(sources),'candidate_count':len(candidates)+len(rejects),'valid_count':len(candidates),'reject_count':len(rejects),'duplicate_count':len(duplicates),'unique_count':len(accepted),'dry_run':a.dry_run,'provider_calls':0,'graphiti_posts':0,'neo4j_writes':0,'queue_writes':0,'new_kus':0,'new_provenance':0,'new_review_audit':0,'review_verdicts':{'PASS':0,'TEILWEISE':0,'FAIL':0},'evidence':{e:sum(1 for _,k in accepted if k['evidence_class']==e) for e in sorted(m.VALID_EVIDENCE)}}
if a.dry_run:
print(json.dumps(summary,ensure_ascii=False)); return
c=sqlite3.connect(DB,timeout=20); c.row_factory=sqlite3.Row; c.execute('PRAGMA foreign_keys=ON'); c.execute('BEGIN IMMEDIATE')
try:
for s,k in accepted:
row=c.execute('SELECT id FROM knowledge_units WHERE knowledge_unit_identity=?',(k['knowledge_unit_identity'],)).fetchone()
if row:
ku_id=row['id']
else:
pv=(s.get('source_versions') or [])
sv_id=pv[0].get('version_id') if pv else None
prov_obj={'prompt_version':m.PROMPT_VERSION,'source_only':True,'source_type':s['source_type'],'source_locator':k['source_locator'],'content_version':k['content_version'],'review_basis':'statement is source-derived and locator resolves to the frozen source record'}
cur=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,source_locator,evidence_class,confidence_reason,review_state,gold_state,legacy_placeholder) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,0)''',(k['knowledge_unit_identity'],k['statement'],'[]','[]',s['source_id'],sv_id,k['content_version'],json.dumps(prov_obj,ensure_ascii=False,sort_keys=True),'P34-R2 controlled 10-source pilot; no automatic Gold promotion',datetime.now(timezone.utc).isoformat(),k['source_locator'],k['evidence_class'],k['confidence_reason'],'PENDING','NONE'))
ku_id=cur.lastrowid; summary['new_kus']+=1
cur=c.execute('''INSERT OR IGNORE INTO knowledge_unit_provenance(knowledge_unit_id,source_identity,source_version,source_locator,representation_type,created_at) VALUES(?,?,?,?,?,?)''',(ku_id,s['source_id'],k['content_version'],k['source_locator'],s['selection_kind'],datetime.now(timezone.utc).isoformat()))
if cur.rowcount: summary['new_provenance']+=1
# Source-only review: every accepted candidate has all required fields, exact locator, and source text containing the statement.
source_text=s['text']; source_lines=source_text.splitlines(); line_no=int(k['source_locator'].split(':',1)[1]); line_ok=1<=line_no<=len(source_lines) and k['statement'] in m.norm(source_lines[line_no-1]).strip('` ')
verdict='PASS' if line_ok and m.valid(k)[0] else 'FAIL'; reason='source-only review: valid structure and exact frozen source line match' if verdict=='PASS' else 'source-only review failed exact frozen source line match'
cur=c.execute('''INSERT OR IGNORE INTO knowledge_unit_review_audit(knowledge_unit_id,verdict,reason,reviewer,reviewed_at) VALUES(?,?,?,?,?)''',(ku_id,verdict,reason,REVIEWER,datetime.now(timezone.utc).isoformat()))
if cur.rowcount:
summary['new_review_audit']+=1; summary['review_verdicts'][verdict]+=1
c.execute('UPDATE knowledge_units SET review_state=? WHERE id=?',('APPROVED' if verdict=='PASS' else 'REJECTED',ku_id))
c.commit()
except: c.rollback(); raise
finally: c.close()
print(json.dumps(summary,ensure_ascii=False))
if __name__=='__main__': main()