"""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", "INTAKE", "P1_PHASE_SEQUENCE", "PARKED_ROUTE", "PLAN", "REVIEW", "VERIFY_NODE", "build_checkpoint_serde", "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" # 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" # 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 carries an approved diff into org CI. # Wired by build_graph only when the caller injects a dispatch_node callable; the # default (None) leaves APPROVED_ROUTE → END unchanged so the graph is inert. 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 _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): new_status = result.get("status") if new_status in _TERMINAL_STATUS_VALUES: recorder.close_terminal(thread_id=thread_id, status=new_status) 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, ) -> 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, ``BUILD -> VERIFY`` 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). """ 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 repoints the " "verifier's APPROVED_ROUTE, so there is nothing to repoint without a " "build->verify subgraph." ) 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) 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: END}, ) 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: END}, ) builder.add_edge(BUILD_NODE, VERIFY_NODE) if dispatch_node is None: # P3 default: APPROVED_ROUTE is the terminus (no dispatch). builder.add_conditional_edges( VERIFY_NODE, route_after_verify, {APPROVED_ROUTE: END, BUILD_ROUTE: BUILD_NODE, PARKED_ROUTE: END}, ) else: # P3+: repoint APPROVED_ROUTE at the dispatch node, then END. builder.add_node( DISPATCH_NODE, _instrument(DISPATCH_NODE, dispatch_node, transition_recorder), ) builder.add_conditional_edges( VERIFY_NODE, route_after_verify, { APPROVED_ROUTE: DISPATCH_NODE, BUILD_ROUTE: BUILD_NODE, PARKED_ROUTE: END, }, ) builder.add_edge(DISPATCH_NODE, 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