diff --git a/agent-team/agent_team/graph.py b/agent-team/agent_team/graph.py index 2bf8dbf..eeda2f9 100644 --- a/agent-team/agent_team/graph.py +++ b/agent-team/agent_team/graph.py @@ -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, diff --git a/agent-team/tests/sim/harness.py b/agent-team/tests/sim/harness.py index 5ac0411..66f4a59 100644 --- a/agent-team/tests/sim/harness.py +++ b/agent-team/tests/sim/harness.py @@ -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``.""" diff --git a/agent-team/tests/sim/test_p1_exit_criteria.py b/agent-team/tests/sim/test_p1_exit_criteria.py index 03212f6..d67b533 100644 --- a/agent-team/tests/sim/test_p1_exit_criteria.py +++ b/agent-team/tests/sim/test_p1_exit_criteria.py @@ -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" diff --git a/agent-team/tests/sim/test_p1_graph_integration.py b/agent-team/tests/sim/test_p1_graph_integration.py new file mode 100644 index 0000000..b547588 --- /dev/null +++ b/agent-team/tests/sim/test_p1_graph_integration.py @@ -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 diff --git a/agent-team/tests/test_graph.py b/agent-team/tests/test_graph.py index b94b3eb..fa6ed27 100644 --- a/agent-team/tests/test_graph.py +++ b/agent-team/tests/test_graph.py @@ -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") diff --git a/requirements.txt b/requirements.txt index b4c01cf..0f02aa4 100644 --- a/requirements.txt +++ b/requirements.txt @@ -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