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.
405 lines
16 KiB
Python
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
|