Explorer
/proc/5461/root/tmp/retry_recovery_orchestrator.py
← Zurück ↓ Download
import sqlite3, json, subprocess, urllib.request, time, os, signal, sys, datetime, hashlib
from pathlib import Path
BASE=Path('/opt/struktur/obsidian-graphiti-import'); DB=BASE/'obsidian_graphiti_import.db'; BLOCK=BASE/'import_run_block.json'; REQ=Path('/opt/struktur/graphiti/request-data/requests.db'); EVENTS=Path('/opt/struktur/graphiti/request-data/processing-events.jsonl'); REPORT=Path('/opt/struktur/knowledge-pipeline-aggregator-staging/AA-RETRY-RECOVERY-FINAL.md')
start=time.time(); original=json.loads(BLOCK.read_text()); results=[]; stopped=None

def now(): return datetime.datetime.now(datetime.timezone.utc).isoformat()
def db(): return sqlite3.connect(f'file:{DB}?mode=ro',uri=True)
def qrow(qid):
 c=db(); c.row_factory=sqlite3.Row; r=c.execute('select * from files where id=?',(qid,)).fetchone(); c.close(); return dict(r) if r else None
def reqrows(qid):
 c=sqlite3.connect(f'file:{REQ}?mode=ro',uri=True); c.row_factory=sqlite3.Row; rs=[dict(x) for x in c.execute('select * from graphiti_requests where queue_id=? order by started_at',(str(qid),))]; c.close(); return rs
def health(): return json.loads(urllib.request.urlopen('http://127.0.0.1:8644/health',timeout=30).read())
def check(row):
 data=json.dumps({'name':'obsidian_'+row['import_key'],'import_key':row['import_key'],'obsidian_path':row['relative_path'],'content_hash':row['content_hash']}).encode(); r=urllib.request.Request('http://127.0.0.1:8644/episodes/check',data=data,headers={'Content-Type':'application/json'}); return json.loads(urllib.request.urlopen(r,timeout=30).read()).get('episodes',[])
def setblock(payload):
 tmp=Path(str(BLOCK)+'.tmp'); tmp.write_text(json.dumps(payload,ensure_ascii=False,indent=2)+'\n'); os.replace(tmp,BLOCK)
def clearblock(qid): setblock({'blocked':False,'reason':'authorized_single_retry_recovery','queue_id':str(qid),'authorized_at':now()})
def active_events(rid):
 out=[]
 if not EVENTS.exists(): return out
 for line in EVENTS.read_text(errors='replace').splitlines():
  try:
   x=json.loads(line)
   if x.get('request_id')==rid: out.append(x)
  except: pass
 return out
def md():
 lines=['# AA-RETRY-RECOVERY-FINAL','',f'Prüf-/Startzeit: `{datetime.datetime.now(datetime.timezone.utc).isoformat()}` UTC','', '## Abschlussstatus','']
 if stopped: lines += ['**RETRY-RECOVERY ABGEBROCHEN –**  ','**FAIL-STOP AKTIV –**  ','**WEITERE QUELLEN NICHT GESTARTET**','']
 else: lines += ['**RETRY-RECOVERY 44/44 ERFOLGREICH –**  ','**RETRY-BESTAND VOLLSTÄNDIG ABGEARBEITET –**  ','**FAILED-BESTAND NOCH GESPERRT**','']
 lines += ['## Ausgangsbestand','', 'retry=44, done=16, duplicate=1, failed=123, processing=0, Requests=745, active_request_count=0','', '## Reihenfolge und Ergebnis','', '| Queue | Pfad | Request | Run | Start | Ende | Dauer s | HTTP | Status | LLM | Edge total/completed | Rate wait s | Episode | Queue-Endstatus |','|---:|---|---|---|---|---|---:|---:|---|---:|---|---:|---|---|']
 for x in results: lines.append('| {queue_id} | {relative_path} | {request_id} | {run_id} | {started_at} | {completed_at} | {duration} | {http_status} | {status} | {llm_request_count} | {edge_total}/{edge_completed} | {rate_gate_wait_seconds} | {episode_uuid} | {queue_status} |'.format(**{k:str(x.get(k,'' )).replace('|','/') for k in ['queue_id','relative_path','request_id','run_id','started_at','completed_at','http_status','status','llm_request_count','edge_total','edge_completed','rate_gate_wait_seconds','episode_uuid','queue_status']},duration=x.get('duration','')))
 if stopped: lines += ['', '## Exakter Abbruchpunkt', f'Queue `{stopped["queue_id"]}`: {stopped["reason"]}', 'Keine weitere Quelle wurde gestartet.']
 lines += ['', '## Schutz- und Endstatus', 'Alle Selektionsbefehle verwendeten ausschließlich `--once --ids X --force-id X --limit 1 --hold-after-success`.', 'Keine failed-Quelle wurde als Ziel ausgewählt; die Allowlist stammte aus `files.status="retry" ORDER BY id ASC`.', 'Provider-/Rate-Gate-Status wurde vor jeder Quelle geprüft; bei Fehler wäre sofort Fail-Stop erfolgt.', 'Run-Block-Endzustand:', '', '```json', json.dumps(json.loads(BLOCK.read_text()),ensure_ascii=False,indent=2), '```', '', '## Endbestand']
 try:
  c=db(); counts=dict(c.execute('select status,count(*) from files group by status').fetchall()); requests=c.execute('select count(*) from graphiti_requests').fetchone()[0]; c.close(); lines.append(f'`{counts}`; Requests=`{requests}`')
 except Exception as e: lines.append(f'Endbestand nicht lesbar: {e}')
 REPORT.write_text('\n'.join(lines)+'\n')

def main():
 global stopped
 c=db(); ids=[r[0] for r in c.execute("select id from files where status='retry' order by id")]; baseline_req=c.execute('select count(*) from graphiti_requests').fetchone() if False else None; c.close()
 if len(ids)!=44: stopped={'queue_id':'PRECHECK','reason':f'expected 44 retry IDs, found {len(ids)}'}; setblock({'blocked':True,'reason':'preflight_failed','blocked_at':now()}); md(); return 1
 for qid in ids:
  row=qrow(qid); h=health(); eps=check(row)
  if row['status']!='retry' or not row['content_hash'] or not row['import_key'] or eps or not h.get('provider_ready') or h.get('provider_rate_limited') or h.get('active_request_count')!=0:
   stopped={'queue_id':qid,'reason':json.dumps({'status':row['status'],'hash':bool(row['content_hash']),'import_key':bool(row['import_key']),'episodes':len(eps),'provider_ready':h.get('provider_ready'),'provider_rate_limited':h.get('provider_rate_limited'),'active_request_count':h.get('active_request_count')},ensure_ascii=False)}; setblock({'blocked':True,'reason':'preflight_failed','queue_id':str(qid),'blocked_at':now()}); md(); return 1
  before=reqrows(qid); clearblock(qid); cmd=['python3',str(BASE/'importer.py'),'--once','--ids',str(qid),'--force-id',str(qid),'--limit','1','--hold-after-success']; t0=time.time(); proc=subprocess.run(cmd,text=True,capture_output=True,timeout=600); t1=time.time(); after=reqrows(qid); new=[x for x in after if x.get('request_id') not in {y.get('request_id') for y in before}]
  rr=new[-1] if len(new)==1 else (new[-1] if new else {})
  verified=[]
  try: verified=check(row)
  except Exception: verified=[]
  events=active_events(rr.get('request_id','')) if rr else []
  phases={x.get('phase') for x in events}; edge_completed=sum(1 for x in events if x.get('phase')=='edge_completed' and x.get('outcome')=='completed'); edge_total=max([int(x.get('edge_total')) for x in events if str(x.get('edge_total','')).isdigit()] or [0]); waits=sum(float(x.get('duration_ms',0))/1000 for x in events if x.get('phase')=='rate_gate_wait' and x.get('outcome')=='completed')
  endrow=qrow(qid); reason=[]
  if proc.returncode!=0: reason.append(f'importer_rc={proc.returncode}')
  if len(new)!=1: reason.append(f'new_requests_for_queue={len(new)}')
  if rr.get('queue_id')!=str(qid): reason.append(f'queue_id_mismatch={rr.get("queue_id")}')
  if rr.get('status')!='completed' or rr.get('http_status')!=200: reason.append(f'request={rr.get("status")}/{rr.get("http_status")}')
  if len(verified)!=1: reason.append(f'episodes={len(verified)}')
  if endrow.get('status')!='done': reason.append(f'queue_endstatus={endrow.get("status")}')
  if any(x in phases for x in ['edge_resolution_failed','edge_operation_failed']): reason.append('terminal_edge_failure_event')
  result={'queue_id':qid,'relative_path':row['relative_path'],'request_id':rr.get('request_id',''),'run_id':rr.get('run_id',''),'started_at':rr.get('started_at',''),'completed_at':rr.get('completed_at',''),'duration':round(t1-t0,3),'http_status':rr.get('http_status',''),'status':rr.get('status',''),'llm_request_count':rr.get('llm_request_count',''),'edge_total':edge_total,'edge_completed':edge_completed,'rate_gate_wait_seconds':round(waits,3),'episode_uuid':(verified[0].get('uuid') if verified else rr.get('episode_uuid','')),'queue_status':endrow.get('status',''),'reason':'; '.join(reason)}
  if reason:
   results.append(result); stopped={'queue_id':qid,'reason':'; '.join(reason)}; setblock({'blocked':True,'reason':'source_failure','queue_id':str(qid),'run_id':rr.get('run_id',''),'blocked_at':now()}); md(); print(json.dumps({'STOP':stopped,'result':result},ensure_ascii=False)); return 1
  results.append(result); setblock({'blocked':True,'reason':'source_verified_wait_next','queue_id':str(qid),'run_id':rr.get('run_id',''),'blocked_at':now()}); print(json.dumps({'SUCCESS':result},ensure_ascii=False),flush=True)
 setblock({'blocked':True,'reason':'retry_recovery_completed','blocked_at':now(),'completed_count':len(results),'completed_queue_ids':[str(x['queue_id']) for x in results]}); md(); print(json.dumps({'COMPLETE':len(results)},ensure_ascii=False)); return 0
try: sys.exit(main())
except BaseException as e:
 stopped={'queue_id':'UNHANDLED','reason':repr(e)}; setblock({'blocked':True,'reason':'unhandled_fail_stop','blocked_at':now()}); md(); raise