#!/usr/bin/env python3
"""P2W fix v2 — clean rewrite of the patch script."""
p = '/tmp/aa043-p2i/aa043_single_writer_staging.py'
src = open(p).read()
helper = '''
def _success_state(c, qid, vid, eid):
"""AA-043-P2W: atomarer Success-State nach bestaetigtem Graphiti-Write
oder ADOPT_EXISTING. 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 and '_success_state' not in src
src = src.replace(anchor, helper + "\n" + anchor)
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'))
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, 'adopt not found'
src = src.replace(old_adopt, new_adopt)
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, 'confirmed not found'
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')