This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
orchestrator/agent-team/agent_team/graph.py
Adam Moussa 0eae5dbfc3 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.
2026-06-17 15:16:12 -04:00

405 lines
16 KiB
Python

"""LangGraph graph wiring for the Plane-2 SDLC pipeline (design §3.3, §7.1 P1).
This module is the **P1 skeleton + human gate** wiring (§7.1):
INTAKE ─► CLARIFY ─► PLAN (stop at an approved plan — no build yet)
It assembles the durable, resumable LangGraph graph whose state schema is the
foundation's :class:`~agent_team.task_model.PipelineState`. The clarifier raises
a LangGraph ``interrupt()`` carrying a :class:`~agent_team.transport.QuestionSet`
so the task suspends + checkpoints, a question-set is delivered over the chosen
transport (D10), and the task resumes via ``Command(resume=...)`` when Adam's
answer arrives (§3.3, §3.3.1). The planner is the P1 terminal stage: it lands an
approved plan and stops; builders/verifiers are later phases (P3).
What this module owns (Plane-2 P1 graph wiring only):
* :data:`P1_PHASE_SEQUENCE` — the ordered P1 stage list.
* the three pure node functions (:func:`intake_node`, :func:`clarify_node`,
:func:`plan_node`) operating on :class:`PipelineState`.
* :func:`build_graph` — assemble + compile the ``StateGraph`` over a *caller-
injected* checkpointer (tests inject an in-memory saver; production injects
the SQLite saver, which the foundation's :func:`agent_team.db.connect`
reserves the DB file for).
* :func:`build_sqlite_checkpointer` — the production checkpointer factory, with
the ``langgraph.checkpoint.sqlite`` import deferred so this module imports
cleanly even where that optional package is absent (pre-deploy scaffolding).
* :func:`thread_config` / :func:`start_task` / :func:`resume_task` /
:func:`get_pipeline_state` / :func:`pending_question` — the thin
``thread_id``-keyed driver seam the coordinator/responder call.
It imports the committed foundation contracts verbatim and does **not** redefine
them. No provisioning, no scheduling, no live SDK calls: the clarifier's
question authoring is a deterministic stub here (the real Claude clarifier binds
``billing.claude_invoke`` in a later phase), and the human-interaction *ledger*
(``pending_questions``) lives in :mod:`agent_team.db.schema` — this module only
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
from langgraph.graph import END, START, StateGraph
from langgraph.types import Command, interrupt
from agent_team.task_model import (
Phase,
PipelineState,
TaskStatus,
new_thread_id,
)
from agent_team.transport import QuestionSet
if TYPE_CHECKING: # pragma: no cover - typing only
from langgraph.checkpoint.base import BaseCheckpointSaver
from langgraph.graph.state import CompiledStateGraph
__all__ = [
"CLARIFY",
"DEFAULT_CLARIFY_DEADLINE",
"INTAKE",
"P1_PHASE_SEQUENCE",
"PLAN",
"build_graph",
"build_sqlite_checkpointer",
"clarify_node",
"get_pipeline_state",
"intake_node",
"pending_question",
"plan_node",
"plan_phase",
"resume_task",
"start_task",
"thread_config",
]
# --- Node names (graph vertices). ------------------------------------------
# Kept as constants so the driver/tests reference the wiring by name rather
# than by string literal.
INTAKE = "intake"
CLARIFY = "clarify"
PLAN = "plan"
# The P1 stage order (§7.1): intake -> clarify -> plan, then stop. Builders and
# verifiers (BUILD/VERIFY) are deliberately NOT wired here — P1 ends at an
# approved plan with no build (§7.1 "Stops at an approved plan, no build yet").
P1_PHASE_SEQUENCE: tuple[Phase, ...] = (Phase.INTAKE, Phase.CLARIFY, Phase.PLAN)
# How long a clarifier question-set stays open before the deadline policy runs
# (§3.3.1 ``deadline_at``). The driver records the concrete ``deadline_at`` on
# 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)."""
return datetime.now(timezone.utc).isoformat()
def _phase_value(phase: Phase) -> str:
"""Return the string value a phase is stored as in :class:`PipelineState`."""
return phase.value
# --- Nodes. -----------------------------------------------------------------
# Each node is a pure ``PipelineState -> partial PipelineState`` function. They
# write only the keys they change (PipelineState is ``total=False``), so a
# checkpoint transition stays minimal. None of them performs I/O.
def intake_node(state: PipelineState) -> PipelineState:
"""INTAKE stage: stamp the task ACTIVE and advance it into CLARIFY (§3.3).
A task enters as a new thread record (the driver mints ``thread_id`` and
seeds INTAKE). This node marks it ``ACTIVE`` and moves the current phase to
``CLARIFY`` so the next node runs the human gate. It never blocks.
"""
return PipelineState(
status=TaskStatus.ACTIVE.value,
current_phase=_phase_value(Phase.CLARIFY),
updated_at=_utc_now_iso(),
)
def clarify_node(state: PipelineState) -> PipelineState:
"""CLARIFY stage: the human gate (LangGraph ``interrupt()``) (§3.3, §3.3.1).
Authors a question-set, then suspends the graph with ``interrupt()`` so the
task checkpoints and waits for Adam's answer. The interrupt payload is the
:class:`~agent_team.transport.QuestionSet` plus the lifecycle metadata the
responder/ledger need (``transport``, ``deadline``); the responder posts it
over the chosen transport and resumes the task via ``Command(resume=...)``.
On resume, ``interrupt()`` returns Adam's answer; this node appends it to
``qa_history`` and advances to ``PLAN``. The P1 skeleton asks exactly one
question-set (``turn`` 0); the multi-turn "until 98% confident" loop is a
later phase, and this node's single-turn shape is forward-compatible with it
(the ``turn`` is read from existing history).
The actual Claude question authoring (``billing.claude_invoke``) is bound in
a later phase; here the question-set is a deterministic stub so the wiring
and the suspend/resume mechanic can be proven without a live model.
"""
history = list(state.get("qa_history", []))
turn = len(history)
thread_id = state.get("thread_id", "")
transport = state.get("transport", "")
# 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,
turn=turn,
questions=_author_questions(state),
context={"phase": _phase_value(Phase.CLARIFY)},
)
deadline = (datetime.now(timezone.utc) + DEFAULT_CLARIFY_DEADLINE).isoformat()
# Suspend here. The payload mirrors §3.3.1: {thread_id, question_id, turn,
# question_set, transport, deadline}. On resume, ``answer`` is whatever the
# responder passed to ``Command(resume=...)``.
answer = interrupt(
{
"thread_id": thread_id,
"question_id": question_id,
"turn": turn,
"question_set": question_set,
"transport": transport,
"deadline": deadline,
}
)
history.append({"turn": turn, "question_id": question_id, "answer": answer})
return PipelineState(
status=TaskStatus.ACTIVE.value,
current_phase=_phase_value(Phase.PLAN),
qa_history=history,
updated_at=_utc_now_iso(),
)
def plan_node(state: PipelineState) -> PipelineState:
"""PLAN stage: land an approved plan and stop — the P1 terminus (§7.1).
Produces the phased plan record and marks the task ``DONE`` for P1 purposes
(P1 "stops at an approved plan, no build yet"). The real planner is Claude
(§3.3); here the plan body is a deterministic stub derived from the gathered
Q&A so the terminal-state wiring is exercised. Builders/verifiers are wired
in P3.
"""
plan = plan_phase(state)
return PipelineState(
status=TaskStatus.DONE.value,
current_phase=_phase_value(Phase.DONE),
plan=plan,
updated_at=_utc_now_iso(),
)
def _author_questions(state: PipelineState) -> list[str]:
"""Deterministic stand-in for the Claude clarifier's question authoring.
The real clarifier gathers repo/memory/handbook context and asks until 98%
confident (§3.3); the P1 skeleton asks a single fixed question-set so the
suspend/resume mechanic is what's under test, not the model.
"""
return ["What problem should this task solve, and what is in scope?"]
def plan_phase(state: PipelineState) -> dict[str, Any]:
"""Build the deterministic P1 plan record from the clarifier Q&A.
Exposed (and unit-tested) separately from :func:`plan_node` so the plan
shape can be asserted without driving the whole graph. The real planner
replaces the body in a later phase.
"""
return {
"summary": "Approved P1 plan (skeleton).",
"phases": ["P1: skeleton + human gate"],
"qa_turns": len(state.get("qa_history", [])),
"approved": True,
}
# --- Graph assembly. --------------------------------------------------------
def build_graph(
checkpointer: BaseCheckpointSaver | None = None,
) -> CompiledStateGraph:
"""Assemble + compile the P1 pipeline ``StateGraph`` (§3.3, §7.1).
Wires ``START → intake → clarify → plan → END`` over
:class:`PipelineState`. The clarifier suspends on ``interrupt()`` for the
human gate; the planner is the P1 terminus (no build).
The ``checkpointer`` is **injected**, never constructed here: the design's
durable store is the SQLite checkpointer (D9), but pre-deploy scaffolding
must not provision it, and tests inject an in-memory saver. Production wires
:func:`build_sqlite_checkpointer`. A checkpointer is required for the
``interrupt()``/``resume`` mechanic to work, so callers that pass ``None``
get an uncheckpointed graph that can run straight-through but cannot
suspend; the driver functions therefore require a checkpointed graph.
"""
builder: StateGraph = StateGraph(PipelineState)
builder.add_node(INTAKE, intake_node)
builder.add_node(CLARIFY, clarify_node)
builder.add_node(PLAN, plan_node)
builder.add_edge(START, INTAKE)
builder.add_edge(INTAKE, CLARIFY)
builder.add_edge(CLARIFY, PLAN)
builder.add_edge(PLAN, END)
if checkpointer is None:
return builder.compile()
return builder.compile(checkpointer=checkpointer)
def build_sqlite_checkpointer(db_path: Path | str) -> BaseCheckpointSaver:
"""Construct the production SQLite checkpointer over ``db_path`` (D9, §3.3).
The import of ``langgraph.checkpoint.sqlite`` is deferred to call time so
this module imports cleanly in environments where that optional package is
not installed (pre-deploy scaffolding). The checkpointer creates its own
tables against the same DB file the foundation's
:func:`agent_team.db.init_db` reserves for it.
Raises a clear :class:`RuntimeError` if the optional package is missing, so
a misconfigured deploy fails loudly rather than silently running
uncheckpointed.
"""
try:
from langgraph.checkpoint.sqlite import SqliteSaver
except ImportError as exc: # pragma: no cover - depends on optional dep
raise RuntimeError(
"langgraph SQLite checkpointer is unavailable; install the "
"'langgraph-checkpoint-sqlite' package to use "
"build_sqlite_checkpointer (D9). Tests inject an in-memory saver."
) from exc
db_path = Path(db_path)
db_path.parent.mkdir(parents=True, exist_ok=True)
return SqliteSaver.from_conn_string(str(db_path))
# --- Driver seam (thread_id-keyed). -----------------------------------------
# Thin helpers the coordinator/responder call. They own the mapping between a
# task's ``thread_id`` and the LangGraph ``config``; the durable lifecycle
# ledger lives in agent_team.db, and the transport delivery in agent_team
# .transport. These keep that wiring in one tested place.
def thread_config(thread_id: str) -> dict[str, Any]:
"""Build the LangGraph ``config`` that scopes an invoke to ``thread_id``.
Every checkpointed invoke/resume for a task must carry the same
``{"configurable": {"thread_id": ...}}`` so it reads/writes that task's
checkpoint and no other (§3.3.1 per-thread isolation).
"""
return {"configurable": {"thread_id": thread_id}}
def start_task(
graph: CompiledStateGraph,
*,
thread_id: str | None = None,
transport: str = "",
) -> tuple[str, PipelineState]:
"""Start a new pipeline task and run it up to the first human gate (§3.3).
Mints a ``thread_id`` (unless one is supplied), seeds the INTAKE state, and
invokes the graph; it runs through INTAKE into CLARIFY and suspends on the
clarifier ``interrupt()``. Returns ``(thread_id, state)`` where ``state`` is
the checkpointed snapshot after the suspend (its ``__interrupt__`` carries
the pending question-set, surfaced by :func:`pending_question`).
The graph MUST be compiled with a checkpointer for the suspend to persist;
an uncheckpointed graph would run straight through without honouring the
interrupt.
"""
tid = thread_id or new_thread_id()
now = _utc_now_iso()
seed = PipelineState(
thread_id=tid,
status=TaskStatus.ACTIVE.value,
current_phase=_phase_value(Phase.INTAKE),
qa_history=[],
transport=transport,
created_at=now,
updated_at=now,
)
result = graph.invoke(seed, thread_config(tid))
return tid, result
def resume_task(
graph: CompiledStateGraph,
*,
thread_id: str,
answer: Any,
) -> PipelineState:
"""Resume a suspended task with Adam's ``answer`` (§3.3, §3.3.1).
Calls ``graph.invoke(Command(resume=answer), config)`` for ``thread_id``.
The clarifier's ``interrupt()`` returns ``answer``, the task appends it to
``qa_history`` and advances through PLAN to completion. The §3.3.1
first-answer-wins / turn-guard discipline lives in the responder + ledger;
this helper is the single-flight resume call the resume worker drives once
it has won the compare-and-set.
"""
return graph.invoke(Command(resume=answer), thread_config(thread_id))
def get_pipeline_state(
graph: CompiledStateGraph,
*,
thread_id: str,
) -> PipelineState:
"""Return the live checkpointed :class:`PipelineState` for ``thread_id``.
Reads the current checkpoint snapshot (post-suspend or post-completion).
Used by recovery + the manual CLI to inspect a task without resuming it.
"""
snapshot = graph.get_state(thread_config(thread_id))
return snapshot.values
def pending_question(
graph: CompiledStateGraph,
*,
thread_id: str,
) -> dict[str, Any] | None:
"""Return the pending interrupt payload for ``thread_id``, or ``None``.
When a task is suspended on the clarifier human gate, its checkpoint carries
an interrupt whose value is the §3.3.1 question payload (``thread_id``,
``question_id``, ``turn``, ``question_set``, ``transport``, ``deadline``).
The responder reads this to author the ledger row + transport post. Returns
``None`` when the task is not currently waiting on a human answer.
"""
snapshot = graph.get_state(thread_config(thread_id))
interrupts = getattr(snapshot, "interrupts", None) or ()
if not interrupts:
return None
return interrupts[0].value