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.
258 lines
10 KiB
Python
258 lines
10 KiB
Python
"""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
|