#!/usr/bin/env python3
"""
AA-043-P2W: Writer-Success-State korrigieren (Staging-Kopie).
Reale Vertragswerte (aus produktiver DB + AA-043-Code belegt, nichts erfunden):
source_versions.graphiti_state : Erfolgswert 'done' (675 Zeilen belegt)
source_versions.reconciliation_state: bleibt 'DECIDED' (unverändert)
source_processing_registry.current_state: 'QUEUED' ist der höchste bisher
verwendete Wert; es existiert KEIN weiterer vertraglicher Erfolgsstatus.
-> current_state wird NICHT geändert (keine neuen Statuswerte erfinden);
der Graph-Erfolg wird ausschließlich über current_graph_version_id
ausgedrückt, was der Classifier-Baseline entspricht.
source_versions hat episode_uuid + verified_at Spalten (vorhanden, genutzt).
"""
p = '/tmp/aa043-p2i/aa043_single_writer_staging.py'
src = open(p).read()
# --- Success helper (atomar) ---
helper = '''
def _success_state(c, qid, vid, eid, adopted=False):
"""AA-043-P2W: atomarer Success-State nach bestätigtem Graphiti-Write
oder ADOPT_EXISTING. Setzt current_graph_version_id + graphiti_state='done'
in derselben lokalen Transaktion wie das queue-done Update."""
reg = c.execute("SELECT registry_id FROM source_versions WHERE version_id=?", (vid,)).fetchone()
if reg is None:
return
rid = reg["registry_id"]
c.execute("UPDATE graphiti_import_queue SET graphiti_status='done',graphiti_episode_id=?,completed_at=?,updated_at=? where id=?", (eid, now(), now(), qid))
c.execute("UPDATE source_versions SET graphiti_state='done', episode_uuid=?, verified_at=COALESCE(verified_at,?) where version_id=?", (eid, now(), vid))
c.execute("UPDATE source_processing_registry SET current_graph_version_id=CASE WHEN current_graph_version_id IS NULL OR current_graph_version_id < ? THEN ? ELSE current_graph_version_id END, updated_at=? where id=?", (vid, vid, now(), rid))
'''
anchor = "def inventory():"
assert anchor in src
src = src.replace(anchor, helper + "\n" + anchor)
# --- ADOPT_EXISTING path: use helper ---
old_adopt = """ if not dry_run:
c.execute("update graphiti_import_queue set graphiti_status='done',graphiti_episode_id=?,completed_at=?,updated_at=? where id=?",(match.get('uuid'),now(),now(),qid)); c.commit()"""
new_adopt = """ if not dry_run:
vid=row['source_version_id'] if 'source_version_id' in keys else None
if vid:
_success_state(c,qid,vid,match.get('uuid'),adopted=True)
else:
c.execute("update graphiti_import_queue set graphiti_status='done',graphiti_episode_id=?,completed_at=?,updated_at=? where id=?",(match.get('uuid'),now(),now(),qid))
c.commit()"""
assert old_adopt in src
src = src.replace(old_adopt, new_adopt)
" keys=set(row.keys())\n if status=='AMBIGUOUS': raise RuntimeError('identity_review_multiple_matches')")
# --- WRITE_CONFIRMED path: use helper ---
old_ok = """ c.execute("update graphiti_import_queue set graphiti_status='done',graphiti_episode_id=?,completed_at=?,remote_outcome='WRITE_CONFIRMED',lock_owner=NULL,lock_token=NULL,lease_expires_at=NULL,updated_at=? where id=?",(eid,now(),now(),qid)); c.commit(); print(json.dumps({'action':'WRITE_CONFIRMED','queue_id':qid,'episode_uuid':eid}));"""
new_ok = """ vid=row['source_version_id'] if 'source_version_id' in keys else None
if vid:
_success_state(c,qid,vid,eid)
else:
c.execute("update graphiti_import_queue set graphiti_status='done',graphiti_episode_id=?,completed_at=?,updated_at=? where id=?",(eid,now(),now(),qid))
c.execute("update graphiti_import_queue set remote_outcome='WRITE_CONFIRMED',lock_owner=NULL,lock_token=NULL,lease_expires_at=NULL where id=?",(qid,)); c.commit(); print(json.dumps({'action':'WRITE_CONFIRMED','queue_id':qid,'episode_uuid':eid}));"""
assert old_ok in src
src = src.replace(old_ok, new_ok)
open(p, 'w').write(src)
import py_compile
py_compile.compile(p, doraise=True)
print('writer fixed + SYNTAX_OK')