#!/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'; SEL='/opt/struktur/reports/aa043-p35/20260828T074105Z/p35-selection.json'; EX='/opt/struktur/AA-043-P33-extractor.py'; OUTDIR='/opt/struktur/reports/aa043-p35/20260828T074105Z'
def load(p,n): s=importlib.util.spec_from_file_location(n,p); m=importlib.util.module_from_spec(s); s.loader.exec_module(m); return m
def main():
 ap=argparse.ArgumentParser(); ap.add_argument('--batch',type=int,required=True); a=ap.parse_args(); m=load(EX,'p33'); D=json.load(open(SEL,encoding='utf8')); batch=next(x for x in D['batches'] if x['batch']==a.batch); ids=batch['source_ids']; assert len(ids)==10
 meta={x['source_id']:x for x in D['sources']}; c=sqlite3.connect('file:'+DB+'?mode=ro&immutable=1',uri=True); c.row_factory=sqlite3.Row; existing={r[0] for r in c.execute('select distinct source_identity from knowledge_units where legacy_placeholder=0 and source_identity is not null')}; c.close()
 sources=[]
 for sid in ids:
  if sid.startswith('youtube:'):
   yid=sid.split(':',1)[1]; rows=[r for r in m.youtube_rows() if r['youtube_id']==yid];
   if len(rows)!=1: raise RuntimeError('source_resolution:'+sid+':'+str(len(rows)))
   s=m.canonical_youtube(rows[0],'p35_batch_'+str(a.batch))
  else: s=m.canonical_obs(Path('/opt/obsidian-vault')/sid.removeprefix('obsidian:'))
  frozen=meta[sid]
  if s['content_version']!=frozen['content_version']: raise RuntimeError('content_version_drift:'+sid)
  if len(m.words(s['text']))<8: raise RuntimeError('unusable_content:'+sid)
  sources.append(s)
 allc=[]; 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); (allc if ok else rejects).append((s,k,rr))
 accepted=[]; dup=[]; cross_potential=[]
 for s,k,_ in allc:
  same=[(ps,pk) for ps,pk in accepted if pk['statement'].strip().casefold()==k['statement'].strip().casefold() and ps['source_id']!=s['source_id']]
  if same: cross_potential.append({'source_id':s['source_id'],'statement':k['statement'],'other_source_id':same[0][0]['source_id']})
  found=None
  for ps,pk in accepted:
   if pk['knowledge_unit_identity']==k['knowledge_unit_identity'] or m.jaccard(pk['statement'],k['statement'])>=.88: found=pk; break
  if found: dup.append((s,k,'exact_or_semantic_duplicate'))
  else: accepted.append((s,k))
 if set(ids)&existing: raise RuntimeError('allowlist_contains_processed:'+json.dumps(sorted(set(ids)&existing)))
 result={'batch':a.batch,'started_at':datetime.now(timezone.utc).isoformat(),'allowlist':ids,'source_count':10,'candidate_count':len(allc)+len(rejects),'valid_count':len(allc),'reject_count':len(rejects),'duplicate_count':len(dup),'semantic_duplicate_count':sum(1 for _,_,r in dup if r=='exact_or_semantic_duplicate'),'cross_source_potential_count':len(cross_potential),'cross_source_potentials':cross_potential,'unique_count':len(accepted),'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)},'provider_calls':0,'graphiti_posts':0,'neo4j_writes':0,'queue_writes':0,'sources':[{'source_id':s['source_id'],'source_type':s['source_type'],'content_version':s['content_version'],'youtube_id':s.get('youtube_id'),'obsidian_path':s.get('obsidian_path'),'word_count':len(m.words(s['text'])),'text_bytes':len(s['text'].encode())} for s in sources]}
 c=sqlite3.connect(DB,timeout=20); c.row_factory=sqlite3.Row; c.execute('PRAGMA foreign_keys=ON'); c.execute('BEGIN IMMEDIATE')
 try:
  inserted=[]
  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:
    sv=(s.get('source_versions') or []); sv_id=sv[0].get('version_id') if sv else None
    prov={'prompt_version':m.PROMPT_VERSION,'source_only':True,'batch':a.batch,'source_locator':k['source_locator'],'content_version':k['content_version']}
    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,ensure_ascii=False,sort_keys=True),'P35 controlled historical batch '+str(a.batch)+'; no automatic Gold promotion',datetime.now(timezone.utc).isoformat(),k['source_locator'],k['evidence_class'],k['confidence_reason'],'PENDING','NONE')); ku_id=cur.lastrowid; result['new_kus']+=1; inserted.append((ku_id,s,k))
   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())); result['new_provenance']+=max(0,cur.rowcount)
  c.commit()
 except: c.rollback(); c.close(); raise
 # Review in separate short transaction, all new rows only
 c.execute('BEGIN IMMEDIATE')
 try:
  for ku_id,s,k in inserted:
   lines=s['text'].splitlines(); n=int(k['source_locator'].split(':',1)[1]); line_ok=1<=n<=len(lines) and k['statement'] in m.norm(lines[n-1]).strip('` '); verdict='PASS' if line_ok and m.valid(k)[0] else 'FAIL'; reason='P35 source-only review: valid structure and exact frozen source line match' if verdict=='PASS' else 'P35 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,'AA-043-P35-source-only-review-v1',datetime.now(timezone.utc).isoformat())); result['new_review_audit']+=max(0,cur.rowcount)
   if cur.rowcount: result['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(); c.close(); raise
 c.close(); result['finished_at']=datetime.now(timezone.utc).isoformat(); Path(OUTDIR,'batch-'+str(a.batch)+'-result.json').write_text(json.dumps(result,ensure_ascii=False,indent=2,sort_keys=True)); print(json.dumps(result,ensure_ascii=False))
if __name__=='__main__': main()
