Prove P1 exit criteria against the real LangGraph graph; fix question_id stability

Reworks the P1 sim so the four §7.1 exit criteria are demonstrated against the
ACTUAL mechanic, not a model (resolves the verifier's "sim models the ledger,
not the LangGraph integration" finding).

- New tests/sim/test_p1_graph_integration.py drives the real agent_team.graph
  StateGraph (interrupt/Command(resume)) + the real langgraph SqliteSaver
  checkpointer + the committed pending_questions compare-and-set, proving:
  (a) suspend survives a simulated restart (drop saver/conn, rebuild over the
  same checkpoint DB) and resumes; (b) duplicate answer loses the CAS and the
  graph never double-advances; (c) a post-deadline answer loses to expire and
  the task is not resumed; (d) two concurrent tasks resume to the correct
  thread, with a turn-guarded no-double-apply check.
- graph.py: derive a STABLE question_id from uuid5(thread_id, turn). The
  clarifier node replays on resume, so the prior fresh-uuid id changed between
  the delivered/ledgered question and the qa_history entry — breaking the
  §3.3.1 identity contract. Now the delivered id == ledger key == history entry
  (unit-tested in test_graph.py).
- harness._connect() now uses the committed schema.connect() (WAL + busy_timeout)
  instead of a raw sqlite3.connect, so concurrent responders genuinely serialize;
  the criterion-(d) concurrency test no longer swallows OperationalError (it
  asserts zero errors + exactly one CAS winner).
- requirements.txt: pin langgraph-checkpoint-sqlite==3.1.0 (design D9 durable
  checkpointer), now exercised by the integration test.

Full suite: 564 passed; ruff + format clean.
This commit is contained in:
Adam Moussa 2026-06-17 14:47:54 -04:00
parent 8f28fbe6a5
commit 0eae5dbfc3
6 changed files with 305 additions and 9 deletions

View file

@ -38,6 +38,7 @@ shapes the interrupt payload that drives it.
from __future__ import annotations
import uuid
from datetime import datetime, timedelta, timezone
from pathlib import Path
from typing import TYPE_CHECKING, Any
@ -93,6 +94,20 @@ P1_PHASE_SEQUENCE: tuple[Phase, ...] = (Phase.INTAKE, Phase.CLARIFY, Phase.PLAN)
# the ledger row; this is only the default window the interrupt advertises.
DEFAULT_CLARIFY_DEADLINE = timedelta(hours=24)
# Fixed namespace for deriving a STABLE question_id from (thread_id, turn). The
# clarifier node re-executes from its start on resume (LangGraph replays the
# node, with interrupt() returning the answer the second time), so a fresh
# random id would change between the suspend that delivered/ledgered the
# question and the resume that records it — breaking the §3.3.1 identity
# contract. A uuid5 over (thread_id, turn) is uuid-shaped yet deterministic, so
# the delivered question_id, the ledger key, and the qa_history entry all agree.
_QUESTION_ID_NAMESPACE = uuid.UUID("a7b9c1d2-3e4f-5061-7283-94a5b6c7d8e9")
def _question_id_for(thread_id: str, turn: int) -> str:
"""Return the stable question_id for ``(thread_id, turn)`` (§3.3.1 identity)."""
return uuid.uuid5(_QUESTION_ID_NAMESPACE, f"{thread_id}:{turn}").hex
def _utc_now_iso() -> str:
"""Return the current UTC time as an ISO-8601 string (ledger-compatible)."""
@ -148,7 +163,9 @@ def clarify_node(state: PipelineState) -> PipelineState:
thread_id = state.get("thread_id", "")
transport = state.get("transport", "")
question_id = new_thread_id()
# Stable across the resume replay of this node (see _question_id_for): the
# id delivered at suspend == the ledger key == the qa_history entry.
question_id = _question_id_for(thread_id, turn)
question_set = QuestionSet(
thread_id=thread_id,
question_id=question_id,

View file

@ -48,6 +48,7 @@ from typing import Any
from agent_team.db.schema import (
answer_question,
connect,
expire_question,
init_db,
supersede_question,
@ -264,9 +265,11 @@ class SimPipeline:
# -- ledger helpers (real db.schema connection) -----------------------
def _connect(self) -> sqlite3.Connection:
return sqlite3.connect(
str(self._db_path), isolation_level=None, check_same_thread=False
)
# Use the committed foundation connection helper (WAL + busy_timeout +
# the stashed db path the compare-and-set relies on) rather than a raw
# sqlite3.connect — so concurrent responders genuinely serialize on the
# write lock (§3.3.1) instead of racing without a busy timeout.
return connect(self._db_path)
def ledger_row(self, question_id: str) -> sqlite3.Row | None:
"""Return the durable ``pending_questions`` row for ``question_id``."""

View file

@ -23,7 +23,6 @@ verify the real mechanic, not a harness convenience.
from __future__ import annotations
import json
import sqlite3
import threading
import pytest
@ -297,6 +296,7 @@ def test_d_concurrent_responders_only_one_wins_per_question(
barrier = threading.Barrier(2)
results: list[bool] = []
errors: list[BaseException] = []
lock = threading.Lock()
def race(via: str) -> None:
@ -305,10 +305,10 @@ def test_d_concurrent_responders_only_one_wins_per_question(
won = pipeline.submit_answer(
{"question_id": suspended.question_id, "answer": via, "via": via}
)
except (
sqlite3.OperationalError
): # pragma: no cover - lock contention is tolerated
won = False
except BaseException as exc: # noqa: BLE001 - record, must be empty
with lock:
errors.append(exc)
return
with lock:
results.append(won)
@ -318,7 +318,13 @@ def test_d_concurrent_responders_only_one_wins_per_question(
for thread in threads:
thread.join()
# No lock error is tolerated: the committed BEGIN IMMEDIATE + busy_timeout
# must *serialize* the responders, so the single winner is the compare-and
# -set, not a swallowed OperationalError loser (the bug the prior version
# masked).
assert errors == [], f"compare-and-set raised under contention: {errors!r}"
assert sum(1 for r in results if r) == 1 # exactly one winner
assert results.count(False) == 1 # the other genuinely lost the CAS (rowcount 0)
assert pipeline.ledger_row(suspended.question_id)["status"] == "answered"

View file

@ -0,0 +1,258 @@
"""P1 exit criteria proven against the REAL LangGraph graph (design §7.1, §3.3.1).
The sibling ``test_p1_exit_criteria.py`` proves the four §7.1 P1 criteria against
a faithful *model* of the ledger/state-store layer. This module proves the same
four criteria against the **actual** mechanic the design's P1 gate requires:
* the real :mod:`agent_team.graph` ``StateGraph`` with a real
``interrupt()`` / ``Command(resume=...)`` clarifier human gate, and
* the real ``langgraph.checkpoint.sqlite.SqliteSaver`` durable checkpointer
(design D9), so "kill the box" is modelled by dropping the saver/connection
and rebuilding the graph over the same checkpoint DB file, plus
* the real committed ``pending_questions`` ledger compare-and-set
(``answer_question`` / ``expire_question``) keyed by the graph's own
``question_id``.
The integration driver here mirrors what the responder/resume-worker do: write
the ledger row when the graph suspends, win the first-answer-wins compare-and-set
before resuming, and only resume a thread that still has a live interrupt (the
turn guard). Nothing is provisioned or networked; this is pre-deploy scaffolding.
"""
from __future__ import annotations
import sqlite3
from pathlib import Path
from typing import Any
import pytest
# The durable SQLite checkpointer is design-required (D9); skip cleanly if the
# optional package is absent so the rest of the suite still runs.
SqliteSaver = pytest.importorskip("langgraph.checkpoint.sqlite").SqliteSaver
from agent_team.db.schema import ( # noqa: E402 - after importorskip by design
answer_question,
connect,
expire_question,
init_db,
)
from agent_team.graph import ( # noqa: E402
build_graph,
get_pipeline_state,
pending_question,
resume_task,
start_task,
)
from agent_team.task_model import Phase, TaskStatus # noqa: E402
class _Pipeline:
"""Thin integration of the real graph + real SqliteSaver + real ledger.
Owns two on-disk SQLite files under ``root``: the LangGraph checkpoint DB
(driven by ``SqliteSaver``) and the committed ``pending_questions`` ledger.
The graph/checkpointer can be rebuilt over the same checkpoint file to model
a restart.
"""
def __init__(self, root: Path) -> None:
self._ckpt_db = root / "checkpoints.sqlite"
self._ledger_db = root / "agent_team.sqlite"
init_db(self._ledger_db)
self._conn: sqlite3.Connection | None = None
self.graph = self._boot()
def _boot(self) -> Any:
"""(Re)build the graph + checkpointer over the same checkpoint DB file."""
if self._conn is not None:
self._conn.close()
self._conn = sqlite3.connect(str(self._ckpt_db), check_same_thread=False)
saver = SqliteSaver(self._conn)
saver.setup()
return build_graph(checkpointer=saver)
def restart(self) -> None:
"""Model "kill the box": drop the saver/connection, rebuild from disk."""
self.graph = self._boot()
def start(self, *, deadline_at: str = "t9999", transport: str = "slack") -> str:
"""Start a task, suspend on the clarifier, and ledger the question row."""
thread_id, _ = start_task(self.graph, transport=transport)
payload = pending_question(self.graph, thread_id=thread_id)
assert payload is not None, "task did not suspend on the human gate"
conn = connect(self._ledger_db)
try:
conn.execute(
"INSERT INTO pending_questions "
"(question_id, thread_id, turn, status, transport, posted_at, "
"deadline_at) VALUES (?, ?, ?, 'open', ?, ?, ?)",
(
payload["question_id"],
thread_id,
payload["turn"],
payload["transport"],
"t0",
deadline_at,
),
)
finally:
conn.close()
return thread_id
def question_id(self, thread_id: str) -> str | None:
payload = pending_question(self.graph, thread_id=thread_id)
return None if payload is None else payload["question_id"]
def answer(self, question_id: str, value: Any, *, via: str = "slack") -> bool:
"""Run the first-answer-wins compare-and-set; True iff this call won."""
conn = connect(self._ledger_db)
try:
return answer_question(
conn,
question_id=question_id,
answer_json=f'{{"v": "{value}"}}',
answered_via=via,
)
finally:
conn.close()
def expire(self, question_id: str) -> bool:
conn = connect(self._ledger_db)
try:
return expire_question(conn, question_id=question_id)
finally:
conn.close()
def ledger_status(self, question_id: str) -> str | None:
conn = connect(self._ledger_db)
try:
row = conn.execute(
"SELECT status FROM pending_questions WHERE question_id = ?",
(question_id,),
).fetchone()
finally:
conn.close()
return None if row is None else row["status"]
def resume_if_won(self, thread_id: str, question_id: str, value: Any) -> bool:
"""Mirror the resume worker: win the CAS, then turn-guarded resume.
Resumes the graph only if (1) this call won the first-answer-wins
compare-and-set AND (2) the thread still has a live interrupt (the turn
guard — a thread already past the gate is never double-resumed).
Returns True iff the graph was actually advanced.
"""
if not self.answer(question_id, value):
return False
if self.question_id(thread_id) is None:
return False # already advanced; do not double-apply
resume_task(self.graph, thread_id=thread_id, answer={"v": value})
return True
def state(self, thread_id: str) -> Any:
return get_pipeline_state(self.graph, thread_id=thread_id)
@pytest.fixture()
def pipeline(tmp_path: Path) -> _Pipeline:
return _Pipeline(tmp_path)
# --- (a) kill the box mid-wait, resume after restart ------------------------
def test_a_suspend_survives_restart_and_resumes(pipeline: _Pipeline) -> None:
thread_id = pipeline.start()
qid = pipeline.question_id(thread_id)
assert qid is not None
pipeline.restart() # drop saver + connection; rebuild over the same DB file
# The interrupt persisted across the restart (durable checkpoint, D9).
assert pipeline.question_id(thread_id) == qid
assert pipeline.resume_if_won(thread_id, qid, "repo-x") is True
state = pipeline.state(thread_id)
assert state["current_phase"] == Phase.DONE.value
assert state["status"] == TaskStatus.DONE.value
assert state["qa_history"][0]["question_id"] == qid # identity held across resume
assert pipeline.question_id(thread_id) is None
# --- (b) duplicate answer no-ops --------------------------------------------
def test_b_duplicate_answer_noops(pipeline: _Pipeline) -> None:
thread_id = pipeline.start()
qid = pipeline.question_id(thread_id)
# First answer wins the CAS and drives the real resume to completion.
assert pipeline.resume_if_won(thread_id, qid, "first") is True
state_after_first = pipeline.state(thread_id)
assert state_after_first["current_phase"] == Phase.DONE.value
assert len(state_after_first["qa_history"]) == 1
# A duplicate answer loses the compare-and-set; no second resume, no change.
assert pipeline.resume_if_won(thread_id, qid, "second") is False
state_after_dup = pipeline.state(thread_id)
assert state_after_dup["qa_history"] == state_after_first["qa_history"]
assert state_after_dup["qa_history"][0]["answer"] == {"v": "first"}
# --- (c) answer after deadline rejected, task not resumed -------------------
def test_c_post_deadline_answer_rejected(pipeline: _Pipeline) -> None:
thread_id = pipeline.start(deadline_at="t1")
qid = pipeline.question_id(thread_id)
# The deadline timer wins the open->expired compare-and-set first.
assert pipeline.expire(qid) is True
assert pipeline.ledger_status(qid) == "expired"
# A late answer loses its compare-and-set, so no resume fires...
assert pipeline.resume_if_won(thread_id, qid, "too-late") is False
# ...and the task is still suspended on the human gate (never advanced).
assert pipeline.question_id(thread_id) == qid
state = pipeline.state(thread_id)
assert state["current_phase"] == Phase.CLARIFY.value
# --- (d) two concurrent tasks resume independently to the correct thread ----
def test_d_two_tasks_resume_to_correct_thread(pipeline: _Pipeline) -> None:
t1 = pipeline.start()
t2 = pipeline.start()
q1, q2 = pipeline.question_id(t1), pipeline.question_id(t2)
assert q1 != q2 # distinct identities per thread
# Resume each with a distinct answer; each must land on its own thread only.
assert pipeline.resume_if_won(t1, q1, "answer-1") is True
assert pipeline.resume_if_won(t2, q2, "answer-2") is True
s1, s2 = pipeline.state(t1), pipeline.state(t2)
assert s1["qa_history"][0]["answer"] == {"v": "answer-1"}
assert s2["qa_history"][0]["answer"] == {"v": "answer-2"}
assert s1["qa_history"][0]["question_id"] == q1
assert s2["qa_history"][0]["question_id"] == q2
assert s1["status"] == TaskStatus.DONE.value
assert s2["status"] == TaskStatus.DONE.value
def test_d_resume_after_completion_does_not_double_apply(pipeline: _Pipeline) -> None:
t1 = pipeline.start()
t2 = pipeline.start()
q1 = pipeline.question_id(t1)
assert pipeline.resume_if_won(t1, q1, "once") is True
before_t2 = pipeline.state(t2)
# A stale/redelivered resume for the already-completed t1 is a no-op (turn
# guard: no live interrupt), and never touches t2.
assert pipeline.resume_if_won(t1, q1, "again") is False
assert pipeline.state(t1)["qa_history"] == [
{"turn": 0, "question_id": q1, "answer": {"v": "once"}}
]
assert pipeline.state(t2) == before_t2

View file

@ -164,6 +164,16 @@ def test_answer_is_recorded_in_qa_history(compiled) -> None:
assert final["qa_history"][0]["turn"] == 0
def test_question_id_is_stable_across_resume(compiled) -> None:
# The clarifier node re-executes on resume; the question_id must NOT change
# between the id delivered at suspend (the ledger key) and the one recorded
# in qa_history, or the §3.3.1 identity contract breaks.
thread_id, _ = start_task(compiled, transport="slack")
delivered = pending_question(compiled, thread_id=thread_id)["question_id"]
final = resume_task(compiled, thread_id=thread_id, answer="ok")
assert final["qa_history"][0]["question_id"] == delivered
def test_no_pending_question_after_completion(compiled) -> None:
thread_id, _ = start_task(compiled, transport="slack")
resume_task(compiled, thread_id=thread_id, answer="ok")

View file

@ -1,4 +1,6 @@
langgraph==1.1.10
# Durable SQLite checkpointer for the R720 agent-team Plane-2 pipeline (design D9).
langgraph-checkpoint-sqlite==3.1.0
langchain-anthropic==1.4.3
langchain-openai==1.2.1
langchain-google-genai==4.2.2