Explorer
/tmp/p2w_fix2.py
← Zurück ↓ Download
#!/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')