#!/usr/bin/env python3
"""Audited backlog optimizer. Dry-run is the default; migration is append-only."""
from __future__ import annotations
import argparse,json,sqlite3
from datetime import datetime,timezone
from pathlib import Path
PIPE='/opt/struktur/obsidian-knowledge-pipeline/knowledge_pipeline.db'
REVIEW='/opt/struktur/knowledge-curator/review_pilot.db'
ANALYSIS='/opt/struktur/knowledge-curator/review-backlog-analysis.json'
VERSION='backlog-optimizer/1.0.0'
def now(): return datetime.now(timezone.utc).isoformat()
def main():
ap=argparse.ArgumentParser(); ap.add_argument('--migrate',action='store_true'); ap.add_argument('--limit',type=int,default=0); ap.add_argument('--report',default='/opt/struktur/knowledge-curator/review-backlog-dry-run.json'); ap.add_argument('--quiet',action='store_true'); args=ap.parse_args()
a=json.load(open(ANALYSIS)); items=a['items'];
action={}
for i in items:
cat=i['category']
if cat=='duplicate': action[i['candidate_id']]={'new_status':'duplicate','decision':'d','rule':'duplicate.normalized_key_and_evidence_hash','duplicate_of':i.get('canonical_candidate_id')}
elif cat=='excluded_non_knowledge': action[i['candidate_id']]={'new_status':'excluded_non_knowledge','decision':'r','rule':'nonknowledge.boilerplate_or_technical_metadata'}
elif cat=='invalid_fragment_or_insufficient_evidence': action[i['candidate_id']]={'new_status':'invalid_fragment','decision':'e','rule':'evidence.complete_sentence_and_subject_required'}
else: action[i['candidate_id']]={'new_status':'pending_review','decision':None,'rule':'retain_human_review'}
counts={}
for v in action.values(): counts[v['new_status']]=counts.get(v['new_status'],0)+1
out={'generated_at':now(),'optimizer_version':VERSION,'source_analysis':ANALYSIS,'before':len(items),'predicted_after_pending_review':counts.get('pending_review',0),'predicted_statuses':counts,'duplicate_groups':a['duplicate_groups'],'duplicate_items':a['duplicate_items'],'no_auto_truth':True,'migration_performed':False,'samples':{}}
for cat in ['duplicate','excluded_non_knowledge','invalid_fragment_or_insufficient_evidence','sensitive_or_high_impact','genuine_pending_review']:
pool=[i for i in items if i['category']==cat]; out['samples'][cat]={'available':len(pool),'sample_size':min(len(pool),100 if cat in ('sensitive_or_high_impact','genuine_pending_review') else 50),'candidate_ids':[i['candidate_id'] for i in pool[:100 if cat in ('sensitive_or_high_impact','genuine_pending_review') else 50]]}
if args.migrate:
pc=sqlite3.connect(PIPE); rc=sqlite3.connect(REVIEW); t=now();
pending_ids={r[0] for r in pc.execute("SELECT candidate_id FROM review_items WHERE status='pending_review'")}
selected=[i for i in items if i['candidate_id'] in pending_ids and action[i['candidate_id']]['new_status']!='pending_review'][:args.limit or len(items)]
for db in (pc,rc): db.execute('BEGIN IMMEDIATE')
try:
cols={r[1] for r in pc.execute('pragma table_info(review_items)')}
for name,typ in [('classification_status','TEXT'),('canonical_candidate_id','TEXT'),('classification_rule','TEXT'),('classification_version','TEXT'),('classified_at','TEXT'),('duplicate_of','TEXT')]:
if name not in cols: pc.execute(f'ALTER TABLE review_items ADD COLUMN {name} {typ}')
run_id='backlog-optimization-'+datetime.now(timezone.utc).strftime('%Y%m%dT%H%M%S%fZ')
rc.execute('INSERT INTO review_runs(run_id,curator_version,normalization_version,started_at,completed_at,candidate_count,reviewer,status) VALUES(?,?,?,?,?,?,?,?)',(run_id,VERSION,'kc-norm-v3',t,t,len(selected),'pipeline-optimizer','completed'))
protected=0; migrated=0
limit=args.limit or len(items)
for i in selected:
cid=i['candidate_id']; v=action[cid]
if v['new_status']=='pending_review': continue
old=pc.execute('SELECT status FROM review_items WHERE candidate_id=?',(cid,)).fetchone()
if not old or old[0] != 'pending_review': continue
pc.execute('UPDATE review_items SET status=?,classification_status=?,canonical_candidate_id=?,classification_rule=?,classification_version=?,classified_at=?,duplicate_of=?,updated_at=? WHERE candidate_id=?', (v['new_status'],v['new_status'],i.get('canonical_candidate_id'),v['rule'],VERSION,t,v.get('duplicate_of'),t,cid))
pc.execute('INSERT INTO audit(run_id,entity_type,entity_key,phase,old_status,new_status,details_json,created_at) VALUES(?,?,?,?,?,?,?,?)',(run_id,'review',cid,'backlog_classification','pending_review',v['new_status'],json.dumps({'rule':v['rule'],'optimizer_version':VERSION,'duplicate_of':v.get('duplicate_of')},ensure_ascii=False),t))
ex=rc.execute('SELECT 1 FROM review_decisions WHERE candidate_id=? LIMIT 1',(cid,)).fetchone()
if ex: protected+=1
else:
reason=f"Deterministische Backlog-Klassifikation {v['rule']}; keine Wahrheitsschreibung, keine Neo4j-Änderung."
did='auto-'+cid+'-'+VERSION.replace('/','-')
rc.execute('INSERT INTO review_decisions(decision_id,candidate_id,reviewer,decision,reason,corrected_candidate_json,reviewed_at,review_duration_seconds) VALUES(?,?,?,?,?,?,?,?)',(did,cid,'pipeline-optimizer',v['decision'],reason,None,t,0.0))
rc.execute('INSERT INTO review_audit(audit_id,candidate_id,previous_state,new_state,action,actor,reason,created_at) VALUES(?,?,?,?,?,?,?,?)',('auto-audit-'+cid+'-'+VERSION.replace('/','-'),cid,'pending_review',v['new_status'],'automatic_backlog_classification','pipeline-optimizer',reason,t))
migrated+=1
pc.commit(); rc.commit(); out['migration_performed']=True; out['migrated']=migrated; out['protected_existing_decisions']=protected; out['migration_run_id']=run_id
except Exception:
pc.rollback(); rc.rollback(); raise
finally: pc.close(); rc.close()
Path(args.report).write_text(json.dumps(out,ensure_ascii=False,indent=2)+'\n',encoding='utf-8')
if not args.quiet: print(json.dumps(out,ensure_ascii=False,indent=2))
if __name__=='__main__': main()