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 67b0f4c6ae feat(agent-team): resumable plan-review gate in the graph (Phase B2a)
Replace the terminal review-cap PARK with a resumable human decision gate,
opt-in via build_graph(plan_gate=True) (default False → all existing P1/P2/P3
wiring unchanged).

- plan_gate_node interrupt()s mirroring the clarifier contract (same payload
  keys → existing pending_question() extractor + turn-guarded ResumeWorker drive
  it with zero special-casing) plus a kind="plan_decision" discriminator and the
  plan + latest review findings as context.
- Decision contract {"decision": approve|request_changes|abandon, "notes": ...}:
  approve → the same terminal state an auto-approved plan reaches (ACTIVE/BUILD);
  request_changes → append a synthetic human verdict to review_verdicts (so the
  planner's _format_review_feedback surfaces the notes) and loop back to PLAN;
  abandon / unrecognized → terminal FAILED (safe default, never accidental
  approve).
- Bounded termination: MAX_PLAN_GATE_VISITS=3 combined ceiling on plan_gate_visits
  (new channel on PipelineState + TaskRecord); on exhaustion the gate goes
  terminal PARKED ("revision ceiling reached") WITHOUT interrupting. Proven by a
  loop-past-ceiling test.

Notes for the coordinator wiring (B2b): while suspended at the gate the status
channel still reads 'parked' (carried over from review_node's escalate branch) —
the load-bearing "awaiting decision, not terminal" signal is the live pending
interrupt + kind="plan_decision", NOT the status channel.

Tests: interrupt-at-cap, approve/request_changes(notes)/abandon routing,
ceiling-terminates, auto-approve still bypasses the gate. 1404 passed.
2026-06-23 21:03:53 -04:00

1046 lines
47 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 functools
import inspect
import uuid
from collections.abc import Callable
from contextlib import contextmanager
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 collections.abc import Iterator
from contextlib import AbstractContextManager
from langgraph.checkpoint.base import BaseCheckpointSaver
from langgraph.graph.state import CompiledStateGraph
# The custom (non-builtin) types that travel inside a LangGraph checkpoint and
# must therefore be on the msgpack deserialization allowlist (D9, §3.3.1). Today
# the only such type is the clarifier's QuestionSet, which rides in the
# ``interrupt()`` payload (graph.clarify_node) and so is msgpack-encoded by the
# checkpoint serializer as ``(module, name, kwargs)`` and reconstructed on
# resume. Without an explicit allowlist the serializer runs in permissive mode
# and logs "Deserializing unregistered type agent_team.transport.base.QuestionSet
# ... will be blocked in a future version" on every checkpoint load; passing the
# allowlist registers the type so it deserializes silently AND keeps working when
# LangGraph flips the future default to block-unregistered. Every PipelineState
# field is a primitive/dict/list and TaskStatus/Phase are stored as ``.value``
# strings, so QuestionSet is the complete set; add any new custom checkpoint type
# here when one is introduced (otherwise it would be blocked under strict mode).
_CHECKPOINT_ALLOWED_MSGPACK_TYPES: tuple[tuple[str, str], ...] = (
("agent_team.transport.base", "QuestionSet"),
)
__all__ = [
"APPROVED_ROUTE",
"BUILD_NODE",
"BUILD_ROUTE",
"CLARIFY",
"DEFAULT_CLARIFY_DEADLINE",
"DISPATCH_NODE",
"GATE_APPROVE_ROUTE",
"GATE_REVISE_ROUTE",
"GATE_TERMINAL_ROUTE",
"INTAKE",
"MAX_PLAN_GATE_VISITS",
"P1_PHASE_SEQUENCE",
"PARKED_ROUTE",
"PLAN",
"PLAN_DECISION_KIND",
"PLAN_GATE",
"REVIEW",
"VERIFY_NODE",
"build_checkpoint_serde",
"build_graph",
"build_sqlite_checkpointer",
"clarify_node",
"get_pipeline_state",
"intake_node",
"pending_question",
"plan_gate_node",
"plan_node",
"plan_phase",
"resume_task",
"route_after_plan_gate",
"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"
# P2 (review loop) node ids. REVIEW is the adversarial-review vertex; BUILD_ROUTE
# and PARKED_ROUTE are the *route ids* the injected route function returns (they
# mirror agent_team.nodes.review_loop.BUILD_NODE / PARKED_NODE by value, so the
# conditional-edge map matches without graph.py importing review_loop). In P2
# both terminate the graph (no builders yet); P3 replaces BUILD_ROUTE's target
# with the real builders subgraph.
REVIEW = "review"
BUILD_ROUTE = "build"
PARKED_ROUTE = "parked"
# Plan-review human-decision gate (Phase B2a). PLAN_GATE is the graph vertex the
# review-cap dead-end (the route the review loop returns as PARKED_ROUTE/ESCALATE)
# is repointed at when the gate is wired: instead of terminally parking, the task
# suspends on a resumable ``interrupt()`` so the owner can decide. The gate
# mirrors the clarifier's interrupt/resume contract exactly (same payload shape +
# the same stable question_id/turn helpers), so the existing pending_question(...)
# extractor and the turn-guarded ResumeWorker drive it uniformly. The route ids
# the gate's conditional-edge function returns are kept distinct from the vertex
# id so a route key never collides with a node name.
PLAN_GATE = "plan_gate"
GATE_APPROVE_ROUTE = "gate_approve"
GATE_REVISE_ROUTE = "gate_revise"
GATE_TERMINAL_ROUTE = "gate_terminal"
# The interrupt-payload discriminator that tells the responder/ledger this is a
# plan-decision gate (vs the clarifier's question-set). The clarifier payload has
# no ``kind`` (legacy = "clarify"); the gate stamps this so a mixed-state ledger
# can tell the two apart.
PLAN_DECISION_KIND = "plan_decision"
# Combined ceiling on plan-gate visits (Phase B2a, bounded termination). A human
# ``request_changes`` re-enters plan<->review, which can hit the review cap and
# gate AGAIN. Each gate visit increments ``plan_gate_visits``; once it reaches
# this constant the gate stops offering request_changes and the task goes
# terminal PARKED with a "revision ceiling reached" marker. This is what proves
# the human-in-the-loop revision cycle terminates: the only non-terminal gate
# decision (request_changes) strictly consumes one of a finite number of visits.
MAX_PLAN_GATE_VISITS = 3
# Disjoint namespace base for the gate's stable ``turn`` derivation. The
# clarifier numbers its turns 0..N from ``qa_history`` length; the gate numbers
# its turns from a high base offset by ``plan_gate_visits`` so a gate turn can
# never collide with a clarifier turn for the same thread (the turn guard in
# ResumeWorker matches a resume to the open interrupt by ``turn``). Only one
# interrupt is ever open per thread, but keeping the spaces disjoint makes the
# stable-question_id derivation unambiguous across the task's whole life.
_PLAN_GATE_TURN_BASE = 1_000_000
# P3 (build -> verify subgraph) vertex ids. These are the GRAPH VERTEX names the
# opt-in P3 subgraph hangs off the review loop's "build" route; they are kept
# distinct from the route-id constants above (BUILD_ROUTE / PARKED_ROUTE) and
# from PLAN/REVIEW so the conditional-edge maps never collide a route key with a
# vertex id. The subgraph itself is supplied wholesale by the Integrate-phase
# caller (the ``build_verify`` tuple), so graph.py does not import
# agent_team.nodes.build_verify_subgraph (no wiring import cycle); it only owns
# the topology that connects the injected nodes.
BUILD_NODE = "build_node"
VERIFY_NODE = "verify_node"
# P3+ dispatch vertex id: the node that triggers org CI and captures the run id.
# Wired by build_graph only when the caller injects a dispatch_node callable, in
# which case it is spliced between BUILD and VERIFY (BUILD -> DISPATCH -> VERIFY)
# so DISPATCH fires the CI run + writes ``state["run_id"]`` BEFORE VERIFY reads it
# (design §4 Decision 1). The default (None) falls back to BUILD -> VERIFY.
DISPATCH_NODE = "dispatch_node"
# P3 route ids returned by the injected ``route_after_verify`` function. They
# mirror agent_team.nodes.build_verify_subgraph.APPROVED_ROUTE / BUILD_ROUTE /
# PARKED_ROUTE by VALUE so this module wires the VERIFY conditional-edge map
# without importing that module. APPROVED_ROUTE is the build->verify-specific
# PASS terminus (the draft-PR endpoint); BUILD_ROUTE loops back to the builders
# under the build-loop budget; PARKED_ROUTE is the fail-safe escalation.
APPROVED_ROUTE = "approved"
# 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", "")
# The Slack root-message ts (one-thread-per-task). Empty for non-/new-task
# origins; surfaced in the interrupt payload so the responder can thread the
# question post under the root message (and use it as the channel_ref).
slack_thread_ts = state.get("slack_thread_ts", "")
# 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,
"slack_thread_ts": slack_thread_ts,
}
)
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 _plan_gate_turn(visits: int) -> int:
"""Return the stable gate ``turn`` for the ``visits``-th gate visit (B2a).
Offset into a high, disjoint namespace (:data:`_PLAN_GATE_TURN_BASE`) so a
gate turn can never collide with a clarifier turn (``qa_history`` length) for
the same thread. Monotonic in ``visits`` so each successive gate suspend has
its own stable ``(thread_id, turn)`` identity (and thus its own question_id).
"""
return _PLAN_GATE_TURN_BASE + visits
def plan_gate_node(state: PipelineState) -> PipelineState:
"""PLAN-GATE stage: the resumable human decision at the review-cap dead-end.
Replaces the terminal PARKED escalation. When the plan<->review loop cannot
converge (the review loop returned the cap/ESCALATE route) — or a salvaged
partial plan is presented — this node suspends with ``interrupt()`` exactly
like the clarifier so the owner can decide. The interrupt payload mirrors the
clarifier's §3.3.1 shape (``thread_id``, ``question_id``, ``turn``,
``transport``, ``deadline``, ``slack_thread_ts``) so the existing
:func:`pending_question` extractor and the turn-guarded
:class:`~agent_team.resume_worker.ResumeWorker` drive it with no special
casing, PLUS:
* ``kind`` = :data:`PLAN_DECISION_KIND` — the discriminator that tells the
responder/ledger this is a plan-decision (not a clarifier question-set);
* ``plan`` / ``findings`` — the review context the owner decides over (the
current plan and the latest review verdict's findings).
**Bounded termination.** Before suspending, the node bumps
``plan_gate_visits``. If the ceiling (:data:`MAX_PLAN_GATE_VISITS`) is already
exhausted it does **not** interrupt: it returns terminal PARKED with a
"revision ceiling reached" marker. So every gate suspend strictly consumes
one of a finite number of visits and the human loop always terminates.
On resume, ``interrupt()`` returns whatever was passed to
``Command(resume=...)`` — the decision contract
``{"decision": "approve"|"request_changes"|"abandon", "notes": str}``. The
node consumes it and writes the routing state (see
:func:`route_after_plan_gate`):
* ``approve`` -> the same terminal "approved plan" state an auto-approved
plan reaches (phase ``BUILD``, status ``ACTIVE``);
* ``request_changes`` -> phase ``PLAN``, status ``ACTIVE`` (loop back), with
the human ``notes`` folded into ``review_verdicts`` as a synthetic
REQUEST_CHANGES verdict so the planner's ``_format_review_feedback`` reads
it on the re-plan;
* ``abandon`` (or any unrecognized decision) -> terminal ``FAILED``.
"""
prior_visits = int(state.get("plan_gate_visits", 0) or 0)
# Ceiling guard FIRST: an exhausted budget means the gate must not offer
# another request_changes loop. Park terminally rather than suspend.
if prior_visits >= MAX_PLAN_GATE_VISITS:
return PipelineState(
status=TaskStatus.PARKED.value,
current_phase=_phase_value(Phase.PARKED),
failure_reason=(
"plan-gate revision ceiling reached "
f"({prior_visits}/{MAX_PLAN_GATE_VISITS} gate visits)"
),
updated_at=_utc_now_iso(),
)
visits = prior_visits + 1
thread_id = state.get("thread_id", "")
transport = state.get("transport", "")
slack_thread_ts = state.get("slack_thread_ts", "")
turn = _plan_gate_turn(visits)
question_id = _question_id_for(thread_id, turn)
deadline = (datetime.now(timezone.utc) + DEFAULT_CLARIFY_DEADLINE).isoformat()
decision = interrupt(
{
"thread_id": thread_id,
"question_id": question_id,
"turn": turn,
"kind": PLAN_DECISION_KIND,
"transport": transport,
"deadline": deadline,
"slack_thread_ts": slack_thread_ts,
"plan": state.get("plan"),
"findings": _latest_review_findings(state),
}
)
return _apply_plan_decision(state, decision, visits=visits)
def _latest_review_findings(state: PipelineState) -> str:
"""Return the latest review verdict's findings text (gate presentation)."""
verdicts = state.get("review_verdicts") or []
if not verdicts:
return ""
last = verdicts[-1]
if isinstance(last, dict):
return str(last.get("findings") or last.get("notes") or "")
return str(last)
def _apply_plan_decision(
state: PipelineState, decision: Any, *, visits: int
) -> PipelineState:
"""Consume the owner's resume decision and return the routing state (B2a).
``decision`` is the value passed to ``Command(resume=...)`` — the contract
``{"decision": ..., "notes": ...}``. A bare string is tolerated as the
decision verb. The default for an unrecognized/empty decision is the SAFE
direction: terminal ``FAILED`` (never an accidental approve).
"""
verb, notes = _parse_decision(decision)
now = _utc_now_iso()
if verb == "approve":
# Settle exactly as an auto-APPROVED plan does today (review_node's
# APPROVE branch: phase BUILD, status ACTIVE) so the terminus matches.
return PipelineState(
status=TaskStatus.ACTIVE.value,
current_phase=_phase_value(Phase.BUILD),
plan_gate_visits=visits,
updated_at=now,
)
if verb == "request_changes":
# Loop back to the planner. Fold the human notes into review_verdicts as
# a synthetic REQUEST_CHANGES verdict whose ``findings`` carry the notes,
# so the planner's _format_review_feedback surfaces them on the re-plan.
verdicts = list(state.get("review_verdicts") or [])
verdicts.append(
{
"verdict": "request_changes",
"outcome": "loop_back",
"findings": notes,
"reviewer": "human_plan_gate",
"created_at": now,
}
)
return PipelineState(
status=TaskStatus.ACTIVE.value,
current_phase=_phase_value(Phase.PLAN),
review_verdicts=verdicts,
plan_gate_visits=visits,
updated_at=now,
)
# abandon (or any unrecognized decision) -> terminal FAILED.
return PipelineState(
status=TaskStatus.FAILED.value,
current_phase=_phase_value(Phase.PARKED),
plan_gate_visits=visits,
failure_reason=f"plan abandoned at human gate: {notes}".rstrip(": "),
updated_at=now,
)
def _parse_decision(decision: Any) -> tuple[str, str]:
"""Normalize a resume decision into ``(verb, notes)`` (B2a).
Accepts the ``{"decision": ..., "notes": ...}`` mapping or a bare string. The
verb is lowercased/stripped and mapped to one of ``approve`` /
``request_changes`` / ``abandon``; anything unrecognized maps to ``abandon``
(the SAFE terminal direction — never an accidental approve).
"""
raw_decision: Any = ""
notes = ""
if isinstance(decision, dict):
raw_decision = decision.get("decision", "")
notes = str(decision.get("notes", "") or "")
else:
raw_decision = decision
verb = str(raw_decision or "").strip().lower()
if verb in {"approve", "request_changes", "abandon"}:
return verb, notes
return "abandon", notes
def route_after_plan_gate(state: PipelineState) -> str:
"""LangGraph conditional-edge after the plan gate (B2a).
Reads the routing state :func:`plan_gate_node` wrote on resume (or on the
ceiling-reached terminal park) and maps it to a route id:
* status ACTIVE + phase BUILD -> :data:`GATE_APPROVE_ROUTE` (END/approved);
* status ACTIVE + phase PLAN -> :data:`GATE_REVISE_ROUTE` (loop to planner);
* anything else (FAILED, or PARKED ceiling) -> :data:`GATE_TERMINAL_ROUTE`.
"""
status = state.get("status")
phase = state.get("current_phase")
if status == TaskStatus.ACTIVE.value and phase == _phase_value(Phase.PLAN):
return GATE_REVISE_ROUTE
if status == TaskStatus.ACTIVE.value and phase == _phase_value(Phase.BUILD):
return GATE_APPROVE_ROUTE
return GATE_TERMINAL_ROUTE
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,
}
# --- Node instrumentation (per-task transition history). --------------------
# Terminal status VALUES (TaskStatus.value) that close a task's open transition
# row — the last node never gets an exited_at from a *next* transition (plan §N2).
_TERMINAL_STATUS_VALUES: frozenset[str] = frozenset(
{TaskStatus.DONE.value, TaskStatus.PARKED.value, TaskStatus.FAILED.value}
)
def _instrument(
name: str,
fn: Callable[..., Any],
recorder: Any | None,
) -> Callable[..., Any]:
"""Wrap a graph node so entering it records a ``task_transitions`` row.
Returns ``fn`` unchanged when ``recorder`` is ``None`` (the default — current
tests and the uninstrumented graph are untouched). Otherwise returns a
signature-preserving wrapper that, on entry, calls ``recorder.record_entry``
(idempotent under LangGraph's resume replay) and, when the node returns a
terminal status, calls ``recorder.close_terminal`` to stamp ``exited_at``.
**Signature preservation (plan §B4).** LangGraph's ``add_node`` inspects the
callable's signature to decide whether to inject a ``RunnableConfig`` second
argument (there is a prior fixed bug of this exact class). ``functools.wraps``
sets ``__wrapped__`` (which ``inspect.signature`` follows) and we ALSO set
``__signature__`` explicitly to ``fn``'s, so LangGraph sees ``fn``'s real
arity and passes exactly the arguments ``fn`` expects; the wrapper forwards
them verbatim via ``*args, **kwargs``. Recording is best-effort: the recorder
is itself fail-soft, and the node call is never gated on it.
"""
if recorder is None:
return fn
@functools.wraps(fn)
def wrapped(state: Any, *args: Any, **kwargs: Any) -> Any:
thread_id = ""
status: Any = None
if isinstance(state, dict):
thread_id = state.get("thread_id", "") or ""
status = state.get("status")
recorder.record_entry(thread_id=thread_id, to_phase=name, status=status)
result = fn(state, *args, **kwargs)
if isinstance(result, dict):
# Nodes store status as TaskStatus.value strings; coerce defensively
# so an enum member (should one slip through) still closes the row.
new_status = result.get("status")
status_value = getattr(new_status, "value", new_status)
if status_value in _TERMINAL_STATUS_VALUES:
recorder.close_terminal(thread_id=thread_id, status=status_value)
return result
# Belt-and-suspenders for B4: present fn's exact signature to LangGraph.
try:
wrapped.__signature__ = inspect.signature(fn) # type: ignore[attr-defined]
except (TypeError, ValueError): # pragma: no cover - exotic callables
pass
return wrapped
# --- Graph assembly. --------------------------------------------------------
def build_graph(
checkpointer: BaseCheckpointSaver | None = None,
*,
transition_recorder: Any | None = None,
live_clarify_node: Callable[[PipelineState], PipelineState] | None = None,
live_plan_node: Callable[[PipelineState], PipelineState] | None = None,
review_node: Callable[[PipelineState], PipelineState] | None = None,
route_review: Callable[[PipelineState], str] | None = None,
build_verify: tuple[
Callable[[PipelineState], dict[str, Any]],
Callable[[PipelineState], PipelineState],
Callable[[PipelineState], str],
]
| None = None,
dispatch_node: Callable[[PipelineState], Any] | None = None,
plan_gate: bool = False,
) -> 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
the **entered** saver yielded by :func:`build_sqlite_checkpointer` (which
returns a context manager the caller must enter and hold, not a bare saver).
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.
``live_clarify_node`` is the **injected real clarifier** (P1a): the live
coordinator passes the Claude-backed multi-turn node built from
:func:`agent_team.nodes.clarifier.make_clarifier_node`, while tests and the
pre-deploy scaffold fall back to the deterministic single-turn
:func:`clarify_node` stub. Either node honours the same ``interrupt()``
suspend/resume contract, so the durable human gate is identical; only the
question authoring differs. Defaulting to the stub keeps the graph wiring
model-free and the existing tests unchanged.
``live_plan_node`` / ``review_node`` / ``route_review`` wire **P2** (planner
+ adversarial review loop). All are injected so this module stays decoupled
from the model + review code (the coordinator passes the wrapped
:func:`agent_team.nodes.planner.plan_node`, the
:func:`agent_team.nodes.review_loop.review_node`, and its
:func:`~agent_team.nodes.review_loop.route_after_review`):
* **P1 (default):** ``review_node`` is ``None`` -> ``plan -> END``. The plan
stage is the terminus (no review, no build), exactly as before.
* **P2:** ``review_node`` is given -> ``plan -> review -> {build|plan|parked}``.
The injected ``route_review`` reads the latest verdict and returns a route
id; the conditional-edge map sends ``"plan"`` back to the planner
(loop-back), and ``"build"`` / ``"parked"`` to ``END`` (P2 stops at an
approved-or-escalated plan; P3 will repoint ``"build"`` at the real
builders subgraph). The plan<->review cycle is bounded by the planner's
revision cap and the review round cap, so the loop always terminates.
``review_node`` requires ``route_review`` (and a real ``live_plan_node`` that
advances to REVIEW); passing one without the other is a wiring error.
``build_verify`` wires the **OPT-IN P3** build -> verify subgraph and is
supplied wholesale as the tuple
``(build_node, verify_node, route_after_verify)`` the coordinator composes
from :mod:`agent_team.nodes.build_verify_subgraph` (passed in so this module
never imports that module — no wiring import cycle). It only takes effect
when the review loop is also wired (it hangs off the review's ``"build"``
route):
* **P2 (default):** ``build_verify`` is ``None`` -> the review's ``"build"``
route terminates at ``END`` (the approved-plan terminus), exactly as
before. Production stays P2 (clarify -> plan -> review).
* **P3:** ``build_verify`` is given (with ``review_node``) -> the review's
``"build"`` route is REPOINTED at the BUILD node, the linear stage order
is wired, and ``route_after_verify`` maps ``{approved -> END (PR terminus),
build -> BUILD (bounded build<->verify loop), parked -> END (escalation)}``.
The subgraph stays INERT unless the caller binds real diff-builder / CI
seams (held for the §3.3.2 security gate); with the default INERT seams the
verifier gate has no authenticated pass and parks. Passing
``build_verify`` without ``review_node`` is a wiring error (there is no
``"build"`` route to repoint).
``dispatch_node`` (P3+, design §4 Decision 1) inserts the CI-dispatch vertex
BETWEEN BUILD and VERIFY: ``BUILD -> DISPATCH -> VERIFY``. DISPATCH triggers
the org CI run and captures the run id into ``state["run_id"]`` BEFORE VERIFY
reads it (VERIFY gates on a CI conclusion that only exists once DISPATCH has
fired the run). When ``dispatch_node`` is ``None`` (the default) the order
falls back to ``BUILD -> VERIFY`` directly. ``dispatch_node`` requires
``build_verify`` (there is no BUILD/VERIFY pair to splice it between
otherwise).
"""
clarify = live_clarify_node if live_clarify_node is not None else clarify_node
plan = live_plan_node if live_plan_node is not None else plan_node
if review_node is not None and route_review is None:
raise ValueError(
"build_graph: review_node requires route_review (the conditional-edge "
"function, e.g. review_loop.route_after_review)."
)
if build_verify is not None and review_node is None:
raise ValueError(
"build_graph: build_verify (the P3 build->verify subgraph) requires "
"review_node — it hangs off the review loop's 'build' route, so there "
"is nothing to repoint without a review loop."
)
if dispatch_node is not None and build_verify is None:
raise ValueError(
"build_graph: dispatch_node requires build_verify — it is spliced "
"between BUILD and VERIFY (BUILD -> DISPATCH -> VERIFY), so there is "
"no BUILD/VERIFY pair to wire it between without a build->verify "
"subgraph."
)
if plan_gate and review_node is None:
raise ValueError(
"build_graph: plan_gate requires review_node — the gate is the "
"resumable replacement for the review loop's terminal PARKED route, "
"so there is no review-cap dead-end to repoint without a review loop."
)
builder: StateGraph = StateGraph(PipelineState)
builder.add_node(INTAKE, _instrument(INTAKE, intake_node, transition_recorder))
builder.add_node(CLARIFY, _instrument(CLARIFY, clarify, transition_recorder))
builder.add_node(PLAN, _instrument(PLAN, plan, transition_recorder))
builder.add_edge(START, INTAKE)
builder.add_edge(INTAKE, CLARIFY)
builder.add_edge(CLARIFY, PLAN)
if review_node is None:
# P1: the plan stage is the terminus.
builder.add_edge(PLAN, END)
else:
# P2/P3: plan -> review -> {loop-back to plan | build | END}.
builder.add_node(REVIEW, _instrument(REVIEW, review_node, transition_recorder))
builder.add_edge(PLAN, REVIEW)
# The review loop's PARKED/ESCALATE route normally terminates at END. When
# the plan gate is wired it is REPOINTED at the PLAN_GATE vertex instead:
# the cap dead-end suspends on a resumable human decision rather than
# terminally parking. The gate's own conditional edges then route
# approve -> END, request_changes -> PLAN (loop back), terminal -> END.
if plan_gate:
parked_target = PLAN_GATE
builder.add_node(
PLAN_GATE,
_instrument(PLAN_GATE, plan_gate_node, transition_recorder),
)
builder.add_conditional_edges(
PLAN_GATE,
route_after_plan_gate,
{
GATE_APPROVE_ROUTE: END,
GATE_REVISE_ROUTE: PLAN,
GATE_TERMINAL_ROUTE: END,
},
)
else:
parked_target = END
if build_verify is None:
# P2: the review's "build" route is the approved-plan terminus.
builder.add_conditional_edges(
REVIEW,
route_review,
{BUILD_ROUTE: END, PLAN: PLAN, PARKED_ROUTE: parked_target},
)
else:
# P3 (opt-in): repoint the review's "build" route at the BUILD node,
# wire BUILD -> VERIFY, and route the verifier verdict to
# {approved -> END (PR terminus), build -> BUILD (loop), parked ->
# END (escalation)}. The subgraph nodes + router are injected (the
# ``build_verify`` tuple) so this module imports no P3 code.
build_node, verify_node, route_after_verify = build_verify
builder.add_node(
BUILD_NODE, _instrument(BUILD_NODE, build_node, transition_recorder)
)
builder.add_node(
VERIFY_NODE, _instrument(VERIFY_NODE, verify_node, transition_recorder)
)
builder.add_conditional_edges(
REVIEW,
route_review,
{BUILD_ROUTE: BUILD_NODE, PLAN: PLAN, PARKED_ROUTE: parked_target},
)
# Linear order is BUILD -> DISPATCH -> VERIFY so DISPATCH triggers CI
# and captures ``state["run_id"]`` BEFORE VERIFY reads it (design §4
# Decision 1: reorder; the CI conclusion VERIFY gates on only exists
# after DISPATCH fires the run). When no dispatch node is wired, fall
# back to BUILD -> VERIFY directly (the INERT default: the verifier's
# CI fetcher yields no authenticated pass and the task parks).
if dispatch_node is None:
builder.add_edge(BUILD_NODE, VERIFY_NODE)
else:
builder.add_node(
DISPATCH_NODE,
_instrument(DISPATCH_NODE, dispatch_node, transition_recorder),
)
builder.add_edge(BUILD_NODE, DISPATCH_NODE)
builder.add_edge(DISPATCH_NODE, VERIFY_NODE)
# VERIFY's verdict routes to {approved -> END (PR terminus),
# build -> BUILD (bounded build<->verify loop), parked -> END
# (escalation)} regardless of whether DISPATCH is wired.
builder.add_conditional_edges(
VERIFY_NODE,
route_after_verify,
{APPROVED_ROUTE: END, BUILD_ROUTE: BUILD_NODE, PARKED_ROUTE: END},
)
if checkpointer is None:
return builder.compile()
return builder.compile(checkpointer=checkpointer)
def build_checkpoint_serde() -> Any:
"""Build the checkpoint serializer with the QuestionSet msgpack allowlist.
Returns a :class:`~langgraph.checkpoint.serde.jsonplus.JsonPlusSerializer`
constructed with an explicit ``allowed_msgpack_modules`` covering every
custom type that rides inside a checkpoint (today only
:class:`~agent_team.transport.QuestionSet`; see
:data:`_CHECKPOINT_ALLOWED_MSGPACK_TYPES`).
Why this exists: the default serializer runs in *permissive* msgpack mode
(``allowed_msgpack_modules=True``), which deserializes any type but logs
``"Deserializing unregistered type agent_team.transport.base.QuestionSet ...
This will be blocked in a future version"`` on every checkpoint load — and a
future LangGraph release will turn that into a hard block, breaking durable
resume. Passing the explicit allowlist is the registration path that warning
recommends: the listed type deserializes silently, and the config is
already block-clean for when the default flips. (Per
``langgraph/checkpoint/serde/jsonplus.py``: ``_create_msgpack_ext_hook``
only emits the warning while the allowlist is the ``True`` sentinel; once an
explicit collection is supplied, an allowlisted ``(module, name)`` returns
True with no warning.)
The ``langgraph.checkpoint.serde.jsonplus`` import is deferred to call time
so this module still imports cleanly where the optional checkpoint package
is absent (pre-deploy scaffolding).
"""
from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer
return JsonPlusSerializer(
allowed_msgpack_modules=list(_CHECKPOINT_ALLOWED_MSGPACK_TYPES)
)
def build_sqlite_checkpointer(
db_path: Path | str,
) -> AbstractContextManager[BaseCheckpointSaver]:
"""Construct the production SQLite checkpointer over ``db_path`` (D9, §3.3).
Returns a **context manager**, not an entered saver: the caller MUST enter
it (``with`` it, or ``__enter__`` and retain it for the graph's lifetime)
before passing the yielded saver to :func:`build_graph`. The live
coordinator owns that lifecycle (it enters the CM at setup and holds it for
the daemon's life); passing the raw return value straight into
``build_graph`` would compile a graph whose checkpointer is an un-entered CM
and break ``get_state`` / ``invoke`` at runtime.
We open the connection and construct ``SqliteSaver(conn, serde=...)``
ourselves rather than using ``SqliteSaver.from_conn_string`` because the
latter has no seam to inject a serializer (it always builds the default,
warning-emitting one). The injected serde is :func:`build_checkpoint_serde`,
whose msgpack allowlist registers :class:`QuestionSet` so resume no longer
logs the "unregistered type" warning and stays forward-compatible with
LangGraph's coming block-by-default. The connection is opened with
``check_same_thread=False`` (matching ``from_conn_string``) and closed when
the context manager exits.
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)
serde = build_checkpoint_serde()
@contextmanager
def _saver_cm() -> Iterator[BaseCheckpointSaver]:
import sqlite3
from contextlib import closing
# Mirror SqliteSaver.from_conn_string's connection settings, but build
# the saver with our allowlisted serde (from_conn_string offers no serde
# seam). closing() guarantees the connection is released on exit.
with closing(sqlite3.connect(str(db_path), check_same_thread=False)) as conn:
yield SqliteSaver(conn, serde=serde)
return _saver_cm()
# --- 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 = "",
task: str = "",
slack_thread_ts: 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`).
``task`` is the intake description (e.g. the Slack ``/new-task`` text or a
GitHub issue body). It is written into the seed ``PipelineState`` so the
clarifier can reason about it; ``intake_node`` returns only a partial state
(status/phase), so the seeded ``task`` channel persists into CLARIFY. An
empty ``task`` (the default) seeds no description — the clarifier then asks
for one.
``slack_thread_ts`` is the Slack root-message ``ts`` for a ``/new-task`` task
(the "📥 Task received" ack post) — when set, every clarifier question and
lifecycle notification for this task threads under it (one-thread-per-task).
Like ``task`` it is seeded into the initial invoke and persists through
INTAKE into CLARIFY (``intake_node`` returns only a partial state). Empty
(the default) for a task with no root post — posts are top-level as before.
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),
task=task,
slack_thread_ts=slack_thread_ts,
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