Explorer
/opt/struktur/knowledge-curator/backlog_optimizer.py
← Zurück ↓ Download
#!/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()