#!/usr/bin/env python3
"""Restore the tail of p2i_e2e.py that fix_upsert truncated (it cut from the
manual upsert to the first kc.close(), deleting commit/assertions/sections
4-9). Append the missing parts."""
src = open('/opt/struktur/obsidian-graphiti-import-staging-p2i/p2i_e2e.py').read()
assert src.rstrip().endswith('post_request_id})'), 'unexpected file end'
tail = '''
kc.commit()
# ADD_KNOWLEDGE assertions (P2Q contract):
all_p2i = kc.execute("SELECT * FROM graphiti_import_queue WHERE obsidian_path='P2I/e2e-note.md'").fetchall()
ak = [r for r in all_p2i if r["action"] == "ADD_KNOWLEDGE"]
assert len(ak) == 1, [(r["reconciliation_record_id"], r["action"]) for r in all_p2i]
row = ak[0]
assert row["reconciliation_record_id"].startswith("aa043:"), row["reconciliation_record_id"]
assert row["payload_json"] and row["payload_hash"] and row["post_request_id"]
import hashlib as _hh
assert _hh.sha256(row["payload_json"].encode()).hexdigest() == row["payload_hash"]
kc.close()
print(" rid:", row["reconciliation_record_id"])
ok("ADD_KNOWLEDGE [queued/ADD_KNOWLEDGE]")
# --- 4. CHANGED -> UPDATE_VERSION -------------------------------------------
res, rows = run_case("CHANGED->UPDATE_VERSION", fid=9003, cvhash=HEX_B,
expect_status="queued", expect_action="UPDATE_VERSION")
# --- 5. AMBIGUOUS -> Block ---------------------------------------------------
from identity.graphiti_identity import GraphitiIdentityAdapter
desc = ('source_identity:obsidian:/ambig.md content_version_hash:' + HEX_A +
' import_identity:iidA canonical_obsidian_path:/ambig.md')
inv = lambda: {"episodes": [
{"uuid": "x1", "name": "a", "source_description": desc},
{"uuid": "x2", "name": "b", "source_description": desc}]}
adapter = GraphitiIdentityAdapter(inv)
m = adapter.lookup(import_identity="iidA", source_identity="obsidian:/ambig.md",
content_version_hash=HEX_A, canonical_obsidian_path="/ambig.md")
assert m.status == "AMBIGUOUS"
from decision_contract import action_for, BLOCKED_CLASSIFICATIONS
assert action_for("AMBIGUOUS") == "BLOCKED"
assert "AMBIGUOUS" in BLOCKED_CLASSIFICATIONS
ok("AMBIGUOUS [block]")
# --- 6. Idempotency ----------------------------------------------------------
kc = kconn(); before = counts(kc); kc.close()
res = bridge_post(9001, HEX_A)
assert res["status"] == "no_write", res
res2 = bridge_post(9003, HEX_B)
kc = kconn(); after = counts(kc); kc.close()
delta = {k: after[k] - before[k] for k in before}
assert delta["graphiti_import_queue"] == 0, delta
ok("Idempotency [no new logical work item]")
# --- 7. Single Writer DRY-RUN -------------------------------------------------
spec = importlib.util.spec_from_file_location(
"p2i_writer", "/tmp/aa043-p2i/aa043_single_writer_staging.py")
wmod = importlib.util.module_from_spec(spec)
sys.modules["p2i_writer"] = wmod
spec.loader.exec_module(wmod)
from pathlib import Path
assert str(wmod.DB) == KNOWLEDGE_DB, str(wmod.DB)
wmod.inventory = lambda: []
os.makedirs("/tmp/P2I", exist_ok=True)
with open("/tmp/P2I/e2e-note.md", "w") as _f:
_f.write("# P2I E2E Note\\n")
os.chdir("/tmp")
kc = kconn()
qid = kc.execute("SELECT id FROM graphiti_import_queue WHERE obsidian_path='P2I/e2e-note.md' AND action='CREATE'").fetchone()["id"]
kc.close()
import io, contextlib
buf = io.StringIO()
with contextlib.redirect_stdout(buf):
wmod.run(qid, dry_run=True)
out = json.loads(buf.getvalue().strip().splitlines()[-1])
assert out["dry_run"] is True and out["action"] == "WRITE_REQUIRED", out
ok("SingleWriter DRY-RUN [WRITE_REQUIRED, injected test DB]")
print(" writer:", json.dumps(out))
# --- 8. No network -----------------------------------------------------------
assert not _calls, _calls
ok("No network POST attempted")
print("=== P2I E2E:", len(PASS), "checks passed ===")
'''
open('/opt/struktur/obsidian-graphiti-import-staging-p2i/p2i_e2e.py', 'w').write(src + tail)
import py_compile
py_compile.compile('/opt/struktur/obsidian-graphiti-import-staging-p2i/p2i_e2e.py', doraise=True)
print('tail restored + SYNTAX_OK')