#!/usr/bin/env python3
"""Explicit, append-only group review and safe backlog migration CLI.

Default operation is read-only.  Status changes require --apply and are
limited to deterministic non-knowledge or insufficient-evidence classes.
Group approvals are recorded separately and never write truth automatically.
"""
from __future__ import annotations

import argparse
import json
import sqlite3
from datetime import datetime, timezone
from pathlib import Path

PIPE = Path('/opt/struktur/obsidian-knowledge-pipeline/knowledge_pipeline.db')
REVIEW = Path('/opt/struktur/knowledge-curator/review_pilot.db')
MHR = Path('/opt/struktur/knowledge-curator/minimal-human-review.db')
ANALYSIS = Path('/opt/struktur/knowledge-curator/minimal-human-review-dry-run.json')
VERSION = 'minimal-human-review/1.0.0'


def now() -> str:
    return datetime.now(timezone.utc).isoformat()


def db() -> sqlite3.Connection:
    c = sqlite3.connect(MHR); c.row_factory = sqlite3.Row; return c


def ensure_schema(c: sqlite3.Connection) -> None:
    c.executescript('''
      CREATE TABLE IF NOT EXISTS group_decisions(
        group_decision_id TEXT PRIMARY KEY, cluster_id TEXT NOT NULL, decision TEXT NOT NULL,
        reviewer TEXT NOT NULL, reason TEXT NOT NULL, scope TEXT NOT NULL, applied INTEGER NOT NULL DEFAULT 0,
        reversible INTEGER NOT NULL DEFAULT 1, created_at TEXT NOT NULL
      );
      CREATE TABLE IF NOT EXISTS group_decision_members(
        group_decision_id TEXT NOT NULL, candidate_id TEXT NOT NULL, PRIMARY KEY(group_decision_id,candidate_id)
      );
      CREATE TABLE IF NOT EXISTS rule_proposals(
        rule_id TEXT PRIMARY KEY, version TEXT NOT NULL, pattern_json TEXT NOT NULL, status TEXT NOT NULL,
        creator TEXT NOT NULL, approval_basis TEXT, approved_by TEXT, approved_at TEXT,
        test_population INTEGER, sample_results_json TEXT, exclusions_json TEXT, sensitivity_limit TEXT,
        rollback_reference TEXT, created_at TEXT NOT NULL
      );
      CREATE TABLE IF NOT EXISTS rule_audit(
        audit_id TEXT PRIMARY KEY, rule_id TEXT NOT NULL, action TEXT NOT NULL, actor TEXT NOT NULL,
        details_json TEXT, created_at TEXT NOT NULL
      );
    '''); c.commit()


def latest_run(c: sqlite3.Connection) -> str:
    row = c.execute('SELECT run_id FROM analysis_runs ORDER BY completed_at DESC LIMIT 1').fetchone()
    if not row: raise SystemExit('Keine Analyse vorhanden; zuerst minimal_review_analysis.py ausführen.')
    return row['run_id']


def show_cluster(c: sqlite3.Connection, cluster_id: str) -> None:
    row = c.execute('SELECT * FROM clusters WHERE cluster_id=?',(cluster_id,)).fetchone()
    if not row: raise SystemExit(f'Unbekannter Cluster: {cluster_id}')
    members=c.execute('SELECT ca.* FROM candidate_analysis ca JOIN cluster_members cm ON cm.candidate_id=ca.candidate_id WHERE cm.cluster_id=? ORDER BY ca.path,ca.candidate_id',(cluster_id,)).fetchall()
    print(json.dumps({'cluster':dict(row),'members':[dict(m) for m in members[:100]],'member_count':len(members),'requires_explicit_human_decision':True},ensure_ascii=False,indent=2))


def list_clusters(c: sqlite3.Connection, limit: int) -> None:
    rows=c.execute('SELECT cluster_id,candidate_count,document_count,status FROM clusters ORDER BY candidate_count DESC,cluster_id LIMIT ?',(limit,)).fetchall()
    print(json.dumps([dict(x) for x in rows],ensure_ascii=False,indent=2))


def safe_migrate(args: argparse.Namespace) -> None:
    analysis=json.loads(ANALYSIS.read_text(encoding='utf-8'))
    selected=[x for x in analysis['items'] if x['category'] in ('excluded_non_knowledge','invalid_fragment')]
    if args.limit: selected=selected[:args.limit]
    if not args.apply:
        print(json.dumps({'dry_run':True,'selected':len(selected),'categories':{k:sum(x['category']==k for x in selected) for k in ('excluded_non_knowledge','invalid_fragment')},'no_status_changes':True},ensure_ascii=False,indent=2)); return
    pc=sqlite3.connect(PIPE); rc=sqlite3.connect(REVIEW); mc=db(); ensure_schema(mc); run_id='mhr-safe-migration-'+datetime.now(timezone.utc).strftime('%Y%m%dT%H%M%S%fZ')
    changed=0; protected=0; batches=0
    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}')
        pc.commit()
        for start in range(0,len(selected),500):
            batch=selected[start:start+500]; t=now(); batches+=1
            for conn in (pc,rc,mc): conn.execute('BEGIN IMMEDIATE')
            try:
                batch_run=run_id+'-'+str(batches)
                rc.execute('INSERT OR IGNORE INTO review_runs(run_id,curator_version,normalization_version,started_at,completed_at,candidate_count,reviewer,status) VALUES(?,?,?,?,?,?,?,?)',(batch_run,VERSION,'kc-norm-v4',t,t,len(batch),'pipeline-optimizer','completed'))
                for item in batch:
                    old=pc.execute('SELECT status FROM review_items WHERE candidate_id=?',(item['candidate_id'],)).fetchone()
                    if not old or old[0]!='pending_review': continue
                    new='excluded_non_knowledge' if item['category']=='excluded_non_knowledge' else 'invalid_fragment'; decision='r' if new=='excluded_non_knowledge' else 'e'; rule='minimal-review.'+item['classification_reason']
                    pc.execute('UPDATE review_items SET status=?,classification_status=?,classification_rule=?,classification_version=?,classified_at=?,updated_at=? WHERE candidate_id=?',(new,new,rule,VERSION,t,t,item['candidate_id']))
                    pc.execute('INSERT INTO audit(run_id,entity_type,entity_key,phase,old_status,new_status,details_json,created_at) VALUES(?,?,?,?,?,?,?,?)',(batch_run,'review',item['candidate_id'],'minimal_review_safe_classification','pending_review',new,json.dumps({'rule':rule,'no_truth_write':True},ensure_ascii=False),t))
                    prior=rc.execute('SELECT 1 FROM review_decisions WHERE candidate_id=? LIMIT 1',(item['candidate_id'],)).fetchone()
                    if prior: protected+=1
                    else:
                        reason=f'Deterministic minimal-review classification {rule}; no truth decision.'
                        rc.execute('INSERT INTO review_decisions(decision_id,candidate_id,reviewer,decision,reason,corrected_candidate_json,reviewed_at,review_duration_seconds) VALUES(?,?,?,?,?,?,?,?)',('auto-mhr-'+item['candidate_id'],item['candidate_id'],'pipeline-optimizer',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-mhr-audit-'+item['candidate_id'],item['candidate_id'],'pending_review',new,'minimal_safe_classification','pipeline-optimizer',reason,t))
                    changed+=1
                pc.commit(); rc.commit(); mc.commit()
            except Exception:
                pc.rollback(); rc.rollback(); mc.rollback(); raise
        print(json.dumps({'migration_run_id':run_id,'dry_run':False,'changed':changed,'protected_existing_decisions':protected,'selected':len(selected),'batches':batches,'batch_limit':500,'no_truth_write':True,'no_neo4j_write':True},ensure_ascii=False,indent=2))
    finally: pc.close(); rc.close(); mc.close()


def apply_group_members(candidates: list[str], decision: str, reason: str, reviewer: str, group_id: str) -> int:
    if decision not in {'approve_pattern','reject_pattern','mark_non_knowledge_pattern','mark_insufficient_evidence_pattern'}:
        return 0
    mapped={'approve_pattern':('approved_for_pilot', 'approved_for_pilot'), 'reject_pattern':('rejected','r'), 'mark_non_knowledge_pattern':('excluded_non_knowledge','r'), 'mark_insufficient_evidence_pattern':('invalid_fragment','e')}
    pipeline_status, review_code = mapped[decision]
    pc=sqlite3.connect(PIPE); rc=sqlite3.connect(REVIEW); t=now(); changed=0
    for conn in (pc,rc): conn.execute('BEGIN IMMEDIATE')
    try:
        for cid in candidates:
            old=pc.execute('SELECT status FROM review_items WHERE candidate_id=?',(cid,)).fetchone()
            if not old or old[0] != 'pending_review': continue
            if decision != 'approve_pattern':
                pc.execute('UPDATE review_items SET status=?,classification_status=?,classification_rule=?,classification_version=?,classified_at=?,updated_at=? WHERE candidate_id=?',(pipeline_status,pipeline_status,'group_decision:'+group_id,VERSION,t,t,cid))
                pc.execute('INSERT INTO audit(run_id,entity_type,entity_key,phase,old_status,new_status,details_json,created_at) VALUES(?,?,?,?,?,?,?,?)',(group_id,'review',cid,'group_decision','pending_review',pipeline_status,json.dumps({'decision':decision,'group_id':group_id,'no_truth_write':True},ensure_ascii=False),t))
            prior=rc.execute('SELECT 1 FROM review_decisions WHERE candidate_id=? LIMIT 1',(cid,)).fetchone()
            if not prior:
                rc.execute('INSERT INTO review_decisions(decision_id,candidate_id,reviewer,decision,reason,corrected_candidate_json,reviewed_at,review_duration_seconds) VALUES(?,?,?,?,?,?,?,?)',('group-'+group_id+'-'+cid,cid,reviewer,review_code,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(?,?,?,?,?,?,?,?)',('group-audit-'+group_id+'-'+cid,cid,'pending_review',pipeline_status,'human_group_decision',reviewer,reason,t))
            changed+=1
        pc.commit(); rc.commit(); return changed
    except Exception:
        pc.rollback(); rc.rollback(); raise
    finally: pc.close(); rc.close()


def record_group(args: argparse.Namespace) -> None:
    c=db(); ensure_schema(c); row=c.execute('SELECT * FROM clusters WHERE cluster_id=?',(args.cluster_id,)).fetchone()
    if not row: raise SystemExit(f'Unbekannter Cluster: {args.cluster_id}')
    if args.decision not in {'approve_pattern','reject_pattern','mark_non_knowledge_pattern','mark_insufficient_evidence_pattern','require_individual_review','split_cluster','approve_only_selected'}: raise SystemExit('Ungültige Gruppenentscheidung')
    members=[r['candidate_id'] for r in c.execute('SELECT candidate_id FROM cluster_members WHERE cluster_id=?',(args.cluster_id,))]
    gid='group-decision-'+datetime.now(timezone.utc).strftime('%Y%m%dT%H%M%S%fZ')
    c.execute('INSERT INTO group_decisions VALUES(?,?,?,?,?,?,?,?,?)',(gid,args.cluster_id,args.decision,args.reviewer,args.reason,args.scope,0,1,now()))
    c.executemany('INSERT INTO group_decision_members VALUES(?,?)',[(gid,m) for m in members]); c.commit()
    c.close()
    changed=apply_group_members(members,args.decision,args.reason,args.reviewer,gid) if args.apply else 0
    if args.apply:
        c=db(); c.execute('UPDATE group_decisions SET applied=1 WHERE group_decision_id=?',(gid,)); c.commit(); c.close()
    print(json.dumps({'group_decision_id':gid,'cluster_id':args.cluster_id,'decision':args.decision,'candidate_count':len(members),'applied':bool(args.apply),'changed':changed,'note':'Explicit group decision recorded append-only. No automatic future scope; current/superseded/verified is never assigned by this command.'},ensure_ascii=False,indent=2))


def suggest_rule(args: argparse.Namespace) -> None:
    c=db(); ensure_schema(c); row=c.execute('SELECT * FROM clusters WHERE cluster_id=?',(args.cluster_id,)).fetchone()
    if not row: raise SystemExit(f'Unbekannter Cluster: {args.cluster_id}')
    members=[r['candidate_id'] for r in c.execute('SELECT candidate_id FROM cluster_members WHERE cluster_id=?',(args.cluster_id,))]
    rule_id='mhr-rule-'+datetime.now(timezone.utc).strftime('%Y%m%dT%H%M%S%fZ')
    pattern={'cluster_id':args.cluster_id,'cluster_key':row['cluster_key'],'candidate_count':row['candidate_count'],'candidate_ids':members,'scope':'existing_only','automatic_truth':False}
    c.execute('INSERT INTO rule_proposals VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?)',(rule_id,'mhr-rule-v1',json.dumps(pattern,ensure_ascii=False),'suggested',args.creator,args.basis,None,None,len(members),None,json.dumps({'sensitive_or_conflict':True},ensure_ascii=False),'never_without_explicit_review',None,now()))
    c.execute('INSERT INTO rule_audit VALUES(?,?,?,?,?,?)',('rule-audit-'+rule_id,rule_id,'suggest_rule',args.creator,json.dumps(pattern,ensure_ascii=False),now())); c.commit(); c.close()
    print(json.dumps({'rule_id':rule_id,'status':'suggested','candidate_count':len(members),'requires_explicit_approval':True,'automatic_truth':False},ensure_ascii=False,indent=2))


def rule_action(args: argparse.Namespace) -> None:
    c=db(); ensure_schema(c); row=c.execute('SELECT * FROM rule_proposals WHERE rule_id=?',(args.rule_id,)).fetchone()
    if not row: raise SystemExit(f'Unbekannte Regel: {args.rule_id}')
    mapping={'approve_rule':'active','reject_rule':'rejected','disable_rule':'disabled','rollback_rule':'rolled_back'}
    if args.action not in mapping: raise SystemExit('Ungültige Regelaktion')
    status=mapping[args.action]; t=now(); c.execute('UPDATE rule_proposals SET status=?,approved_by=?,approved_at=?,approval_basis=? WHERE rule_id=?',(status,args.actor,t,args.reason,args.rule_id))
    c.execute('INSERT INTO rule_audit VALUES(?,?,?,?,?,?)',('rule-audit-'+args.action+'-'+args.rule_id,args.rule_id,args.action,args.actor,json.dumps({'reason':args.reason,'automatic_truth':False},ensure_ascii=False),t)); c.commit(); c.close()
    print(json.dumps({'rule_id':args.rule_id,'action':args.action,'status':status,'automatic_future_application':False,'note':'Rule status is audited; pipeline application remains disabled until separately integrated and approved.'},ensure_ascii=False,indent=2))


def test_rule(args: argparse.Namespace) -> None:
    c=db(); row=c.execute('SELECT * FROM rule_proposals WHERE rule_id=?',(args.rule_id,)).fetchone()
    if not row: raise SystemExit(f'Unbekannte Regel: {args.rule_id}')
    pattern=json.loads(row['pattern_json']); ids=pattern.get('candidate_ids',[]); c.close()
    print(json.dumps({'rule_id':args.rule_id,'status':row['status'],'test_population':len(ids),'sample_size':min(100,len(ids)),'false_positive_review_required':True,'no_status_changes':True},ensure_ascii=False,indent=2))


def main() -> None:
    ap=argparse.ArgumentParser(); sub=ap.add_subparsers(dest='cmd',required=True)
    s=sub.add_parser('list-clusters'); s.add_argument('--limit',type=int,default=50)
    s=sub.add_parser('show-cluster'); s.add_argument('cluster_id')
    s=sub.add_parser('migrate-safe'); s.add_argument('--apply',action='store_true'); s.add_argument('--limit',type=int,default=0)
    s=sub.add_parser('record-group'); s.add_argument('cluster_id'); s.add_argument('--decision',required=True); s.add_argument('--reviewer',required=True); s.add_argument('--reason',required=True); s.add_argument('--scope',choices=['apply_to_existing_only','apply_to_existing_and_future'],default='apply_to_existing_only'); s.add_argument('--apply',action='store_true')
    s=sub.add_parser('suggest-rule'); s.add_argument('cluster_id'); s.add_argument('--creator',required=True); s.add_argument('--basis',required=True)
    s=sub.add_parser('rule-action'); s.add_argument('rule_id'); s.add_argument('--action',choices=['approve_rule','reject_rule','disable_rule','rollback_rule'],required=True); s.add_argument('--actor',required=True); s.add_argument('--reason',required=True)
    s=sub.add_parser('test-rule'); s.add_argument('rule_id')
    args=ap.parse_args(); c=db(); ensure_schema(c)
    if args.cmd=='list-clusters': list_clusters(c,args.limit); c.close()
    elif args.cmd=='show-cluster': show_cluster(c,args.cluster_id); c.close()
    elif args.cmd=='migrate-safe': c.close(); safe_migrate(args)
    elif args.cmd=='record-group': c.close(); record_group(args)
    elif args.cmd=='suggest-rule': c.close(); suggest_rule(args)
    elif args.cmd=='rule-action': c.close(); rule_action(args)
    elif args.cmd=='test-rule': c.close(); test_rule(args)


if __name__=='__main__': main()
