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.
This commit is contained in:
parent
1b4d30e47f
commit
67b0f4c6ae
3 changed files with 497 additions and 2 deletions
|
|
@ -89,10 +89,16 @@ __all__ = [
|
||||||
"CLARIFY",
|
"CLARIFY",
|
||||||
"DEFAULT_CLARIFY_DEADLINE",
|
"DEFAULT_CLARIFY_DEADLINE",
|
||||||
"DISPATCH_NODE",
|
"DISPATCH_NODE",
|
||||||
|
"GATE_APPROVE_ROUTE",
|
||||||
|
"GATE_REVISE_ROUTE",
|
||||||
|
"GATE_TERMINAL_ROUTE",
|
||||||
"INTAKE",
|
"INTAKE",
|
||||||
|
"MAX_PLAN_GATE_VISITS",
|
||||||
"P1_PHASE_SEQUENCE",
|
"P1_PHASE_SEQUENCE",
|
||||||
"PARKED_ROUTE",
|
"PARKED_ROUTE",
|
||||||
"PLAN",
|
"PLAN",
|
||||||
|
"PLAN_DECISION_KIND",
|
||||||
|
"PLAN_GATE",
|
||||||
"REVIEW",
|
"REVIEW",
|
||||||
"VERIFY_NODE",
|
"VERIFY_NODE",
|
||||||
"build_checkpoint_serde",
|
"build_checkpoint_serde",
|
||||||
|
|
@ -102,9 +108,11 @@ __all__ = [
|
||||||
"get_pipeline_state",
|
"get_pipeline_state",
|
||||||
"intake_node",
|
"intake_node",
|
||||||
"pending_question",
|
"pending_question",
|
||||||
|
"plan_gate_node",
|
||||||
"plan_node",
|
"plan_node",
|
||||||
"plan_phase",
|
"plan_phase",
|
||||||
"resume_task",
|
"resume_task",
|
||||||
|
"route_after_plan_gate",
|
||||||
"start_task",
|
"start_task",
|
||||||
"thread_config",
|
"thread_config",
|
||||||
]
|
]
|
||||||
|
|
@ -125,6 +133,44 @@ REVIEW = "review"
|
||||||
BUILD_ROUTE = "build"
|
BUILD_ROUTE = "build"
|
||||||
PARKED_ROUTE = "parked"
|
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
|
# 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
|
# 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
|
# distinct from the route-id constants above (BUILD_ROUTE / PARKED_ROUTE) and
|
||||||
|
|
@ -287,6 +333,201 @@ def plan_node(state: PipelineState) -> PipelineState:
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
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]:
|
def _author_questions(state: PipelineState) -> list[str]:
|
||||||
"""Deterministic stand-in for the Claude clarifier's question authoring.
|
"""Deterministic stand-in for the Claude clarifier's question authoring.
|
||||||
|
|
||||||
|
|
@ -389,6 +630,7 @@ def build_graph(
|
||||||
]
|
]
|
||||||
| None = None,
|
| None = None,
|
||||||
dispatch_node: Callable[[PipelineState], Any] | None = None,
|
dispatch_node: Callable[[PipelineState], Any] | None = None,
|
||||||
|
plan_gate: bool = False,
|
||||||
) -> CompiledStateGraph:
|
) -> CompiledStateGraph:
|
||||||
"""Assemble + compile the P1 pipeline ``StateGraph`` (§3.3, §7.1).
|
"""Assemble + compile the P1 pipeline ``StateGraph`` (§3.3, §7.1).
|
||||||
|
|
||||||
|
|
@ -489,6 +731,13 @@ def build_graph(
|
||||||
"subgraph."
|
"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: StateGraph = StateGraph(PipelineState)
|
||||||
builder.add_node(INTAKE, _instrument(INTAKE, intake_node, transition_recorder))
|
builder.add_node(INTAKE, _instrument(INTAKE, intake_node, transition_recorder))
|
||||||
builder.add_node(CLARIFY, _instrument(CLARIFY, clarify, transition_recorder))
|
builder.add_node(CLARIFY, _instrument(CLARIFY, clarify, transition_recorder))
|
||||||
|
|
@ -506,12 +755,35 @@ def build_graph(
|
||||||
builder.add_node(REVIEW, _instrument(REVIEW, review_node, transition_recorder))
|
builder.add_node(REVIEW, _instrument(REVIEW, review_node, transition_recorder))
|
||||||
builder.add_edge(PLAN, REVIEW)
|
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:
|
if build_verify is None:
|
||||||
# P2: the review's "build" route is the approved-plan terminus.
|
# P2: the review's "build" route is the approved-plan terminus.
|
||||||
builder.add_conditional_edges(
|
builder.add_conditional_edges(
|
||||||
REVIEW,
|
REVIEW,
|
||||||
route_review,
|
route_review,
|
||||||
{BUILD_ROUTE: END, PLAN: PLAN, PARKED_ROUTE: END},
|
{BUILD_ROUTE: END, PLAN: PLAN, PARKED_ROUTE: parked_target},
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
# P3 (opt-in): repoint the review's "build" route at the BUILD node,
|
# P3 (opt-in): repoint the review's "build" route at the BUILD node,
|
||||||
|
|
@ -530,7 +802,7 @@ def build_graph(
|
||||||
builder.add_conditional_edges(
|
builder.add_conditional_edges(
|
||||||
REVIEW,
|
REVIEW,
|
||||||
route_review,
|
route_review,
|
||||||
{BUILD_ROUTE: BUILD_NODE, PLAN: PLAN, PARKED_ROUTE: END},
|
{BUILD_ROUTE: BUILD_NODE, PLAN: PLAN, PARKED_ROUTE: parked_target},
|
||||||
)
|
)
|
||||||
# Linear order is BUILD -> DISPATCH -> VERIFY so DISPATCH triggers CI
|
# Linear order is BUILD -> DISPATCH -> VERIFY so DISPATCH triggers CI
|
||||||
# and captures ``state["run_id"]`` BEFORE VERIFY reads it (design §4
|
# and captures ``state["run_id"]`` BEFORE VERIFY reads it (design §4
|
||||||
|
|
|
||||||
|
|
@ -114,6 +114,15 @@ class TaskRecord:
|
||||||
# parks once it reaches VerifierConfig.max_build_loops so a perpetually-
|
# parks once it reaches VerifierConfig.max_build_loops so a perpetually-
|
||||||
# failing task can never loop BUILD->DISPATCH->VERIFY forever (LOGIC-RACE-01).
|
# failing task can never loop BUILD->DISPATCH->VERIFY forever (LOGIC-RACE-01).
|
||||||
build_loops: int = 0
|
build_loops: int = 0
|
||||||
|
# Count of plan-review human-decision gate visits already consumed for THIS
|
||||||
|
# task (Phase B2a). The graph routes a review-cap dead-end to the plan gate,
|
||||||
|
# which interrupts for an owner approve / request_changes / abandon decision.
|
||||||
|
# A human ``request_changes`` re-enters plan<->review, which can hit the cap
|
||||||
|
# and gate AGAIN; this counter is the combined ceiling on human-driven gate
|
||||||
|
# loops (mirrors PipelineState.plan_gate_visits) so the loop always
|
||||||
|
# terminates: once it reaches ``MAX_PLAN_GATE_VISITS`` the gate stops
|
||||||
|
# offering request_changes and the task goes terminal PARKED.
|
||||||
|
plan_gate_visits: int = 0
|
||||||
transport: str = ""
|
transport: str = ""
|
||||||
created_at: str | None = None
|
created_at: str | None = None
|
||||||
updated_at: str | None = None
|
updated_at: str | None = None
|
||||||
|
|
@ -162,6 +171,12 @@ class PipelineState(TypedDict, total=False):
|
||||||
# budget is real (LOGIC-RACE-01: it was previously read from the shared
|
# budget is real (LOGIC-RACE-01: it was previously read from the shared
|
||||||
# wiring-time config and never advanced).
|
# wiring-time config and never advanced).
|
||||||
build_loops: int
|
build_loops: int
|
||||||
|
# Plan-review human-decision gate visits consumed for THIS task (mirrors
|
||||||
|
# TaskRecord.plan_gate_visits; Phase B2a). The combined ceiling on
|
||||||
|
# human-driven plan<->review loops: once it reaches MAX_PLAN_GATE_VISITS the
|
||||||
|
# gate stops offering request_changes and the task goes terminal PARKED, so
|
||||||
|
# the human-in-the-loop revision cycle can never spin forever.
|
||||||
|
plan_gate_visits: int
|
||||||
transport: str
|
transport: str
|
||||||
created_at: str | None
|
created_at: str | None
|
||||||
updated_at: str | None
|
updated_at: str | None
|
||||||
|
|
@ -198,6 +213,7 @@ def task_from_dict(data: dict[str, Any]) -> TaskRecord:
|
||||||
ci_correlation_tag=data.get("ci_correlation_tag"),
|
ci_correlation_tag=data.get("ci_correlation_tag"),
|
||||||
dispatched_at=data.get("dispatched_at"),
|
dispatched_at=data.get("dispatched_at"),
|
||||||
build_loops=data.get("build_loops", 0),
|
build_loops=data.get("build_loops", 0),
|
||||||
|
plan_gate_visits=data.get("plan_gate_visits", 0),
|
||||||
transport=data.get("transport", ""),
|
transport=data.get("transport", ""),
|
||||||
created_at=data.get("created_at"),
|
created_at=data.get("created_at"),
|
||||||
updated_at=data.get("updated_at"),
|
updated_at=data.get("updated_at"),
|
||||||
|
|
|
||||||
|
|
@ -501,6 +501,213 @@ def test_p2_graph_loops_then_escalates_on_persistent_changes(
|
||||||
assert len(final["review_verdicts"]) == 3 # looped to the cap, then escalated
|
assert len(final["review_verdicts"]) == 3 # looped to the cap, then escalated
|
||||||
|
|
||||||
|
|
||||||
|
# --- Plan-review human decision gate (Phase B2a). ---------------------------
|
||||||
|
|
||||||
|
|
||||||
|
def _plan_gate_graph(review_text: str):
|
||||||
|
"""Compile a P2 graph with the plan gate wired (plan_gate=True).
|
||||||
|
|
||||||
|
A reviewer that never approves drives the plan<->review loop to the review
|
||||||
|
cap, which — with the gate wired — suspends on the resumable PLAN_GATE
|
||||||
|
interrupt instead of terminally parking.
|
||||||
|
"""
|
||||||
|
from agent_team.nodes import review_loop
|
||||||
|
|
||||||
|
review_loop.set_review_invoker(lambda prompt, **kw: review_text)
|
||||||
|
return build_graph(
|
||||||
|
checkpointer=_Saver(),
|
||||||
|
live_plan_node=_p2_plan_stub,
|
||||||
|
review_node=review_loop.bind_review_node(),
|
||||||
|
route_review=review_loop.route_after_review,
|
||||||
|
plan_gate=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _drive_to_plan_gate(graph):
|
||||||
|
"""Start a task and resume the clarifier so it lands on the plan gate.
|
||||||
|
|
||||||
|
Returns ``(thread_id, gate_payload)`` where ``gate_payload`` is the pending
|
||||||
|
plan-decision interrupt payload.
|
||||||
|
"""
|
||||||
|
from agent_team.graph import pending_question
|
||||||
|
|
||||||
|
thread_id, _ = start_task(graph, transport="slack", slack_thread_ts="ROOT.1")
|
||||||
|
resume_task(graph, thread_id=thread_id, answer="scope is X")
|
||||||
|
payload = pending_question(graph, thread_id=thread_id)
|
||||||
|
return thread_id, payload
|
||||||
|
|
||||||
|
|
||||||
|
def test_plan_gate_requires_review_node() -> None:
|
||||||
|
with pytest.raises(ValueError, match="plan_gate requires review_node"):
|
||||||
|
build_graph(plan_gate=True)
|
||||||
|
|
||||||
|
|
||||||
|
def test_review_cap_interrupts_at_plan_gate(restore_review_invoker) -> None:
|
||||||
|
# Driving to the review-cap dead-end suspends on the plan gate (a pending
|
||||||
|
# plan_decision interrupt with the plan + findings) rather than parking.
|
||||||
|
from agent_team.graph import PLAN_DECISION_KIND
|
||||||
|
|
||||||
|
graph = _plan_gate_graph("VERDICT: REQUEST CHANGES\nstill not ready")
|
||||||
|
thread_id, payload = _drive_to_plan_gate(graph)
|
||||||
|
|
||||||
|
assert payload is not None
|
||||||
|
assert payload["kind"] == PLAN_DECISION_KIND
|
||||||
|
# Mirrors the clarifier contract: thread_id / question_id / turn / transport /
|
||||||
|
# deadline / slack_thread_ts all present (so pending_question + ResumeWorker
|
||||||
|
# drive it uniformly).
|
||||||
|
assert payload["thread_id"] == thread_id
|
||||||
|
assert payload["question_id"]
|
||||||
|
assert isinstance(payload["turn"], int)
|
||||||
|
assert payload["transport"] == "slack"
|
||||||
|
assert payload["deadline"]
|
||||||
|
assert payload["slack_thread_ts"] == "ROOT.1"
|
||||||
|
# Plus the review context the owner decides over.
|
||||||
|
assert payload["plan"] is not None
|
||||||
|
assert "still not ready" in payload["findings"]
|
||||||
|
|
||||||
|
# It did NOT terminally park: the task is suspended on the gate interrupt
|
||||||
|
# (resumable), not finished. (The status channel still reads the review
|
||||||
|
# node's carried-over "parked" until the gate's resume overwrites it; the
|
||||||
|
# load-bearing signal is the live pending interrupt.)
|
||||||
|
from agent_team.graph import pending_question
|
||||||
|
|
||||||
|
assert pending_question(graph, thread_id=thread_id) is not None
|
||||||
|
snapshot = graph.get_state(thread_config(thread_id))
|
||||||
|
assert snapshot.next # graph is suspended, not at a terminal END
|
||||||
|
|
||||||
|
|
||||||
|
def test_plan_gate_approve_settles_as_approved_plan(restore_review_invoker) -> None:
|
||||||
|
# Command(resume={"decision":"approve"}) -> the same terminal "approved plan"
|
||||||
|
# state an auto-approved plan reaches today (phase BUILD, status ACTIVE).
|
||||||
|
graph = _plan_gate_graph("VERDICT: REQUEST CHANGES\nnope")
|
||||||
|
thread_id, _ = _drive_to_plan_gate(graph)
|
||||||
|
final = resume_task(graph, thread_id=thread_id, answer={"decision": "approve"})
|
||||||
|
|
||||||
|
assert final["current_phase"] == Phase.BUILD.value
|
||||||
|
assert final["status"] == TaskStatus.ACTIVE.value
|
||||||
|
|
||||||
|
|
||||||
|
def test_plan_gate_request_changes_loops_back_with_notes(
|
||||||
|
restore_review_invoker,
|
||||||
|
) -> None:
|
||||||
|
# Command(resume={"decision":"request_changes","notes":"do X"}) loops back to
|
||||||
|
# the planner; the notes must reach _format_review_feedback / the planner.
|
||||||
|
from agent_team.nodes.planner import _format_review_feedback
|
||||||
|
|
||||||
|
captured: dict[str, object] = {}
|
||||||
|
|
||||||
|
def capturing_plan(state):
|
||||||
|
captured["feedback"] = _format_review_feedback(
|
||||||
|
list(state.get("review_verdicts") or [])
|
||||||
|
)
|
||||||
|
# After observing the folded-in notes, APPROVE on the re-plan so the
|
||||||
|
# graph settles (the review invoker is bound per-test below).
|
||||||
|
return _p2_plan_stub(state)
|
||||||
|
|
||||||
|
from agent_team.nodes import review_loop
|
||||||
|
|
||||||
|
# First review round REQUEST CHANGES (to reach the gate); after the human
|
||||||
|
# request_changes loops back, the next review APPROVES so the task settles.
|
||||||
|
texts = iter(
|
||||||
|
[
|
||||||
|
"VERDICT: REQUEST CHANGES\nnot ready",
|
||||||
|
"VERDICT: APPROVE\nnow good",
|
||||||
|
]
|
||||||
|
)
|
||||||
|
last = "VERDICT: APPROVE\nnow good"
|
||||||
|
|
||||||
|
def invoker(prompt, **kw):
|
||||||
|
nonlocal last
|
||||||
|
try:
|
||||||
|
last = next(texts)
|
||||||
|
except StopIteration:
|
||||||
|
pass
|
||||||
|
return last
|
||||||
|
|
||||||
|
review_loop.set_review_invoker(invoker)
|
||||||
|
graph = build_graph(
|
||||||
|
checkpointer=_Saver(),
|
||||||
|
live_plan_node=capturing_plan,
|
||||||
|
review_node=review_loop.bind_review_node({"max_review_rounds": 1}),
|
||||||
|
route_review=review_loop.route_after_review,
|
||||||
|
plan_gate=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
thread_id, _ = start_task(graph, transport="slack")
|
||||||
|
resume_task(graph, thread_id=thread_id, answer="scope is X")
|
||||||
|
final = resume_task(
|
||||||
|
graph,
|
||||||
|
thread_id=thread_id,
|
||||||
|
answer={"decision": "request_changes", "notes": "do X"},
|
||||||
|
)
|
||||||
|
|
||||||
|
# The human notes reached the planner's review-feedback formatter on re-plan.
|
||||||
|
assert "do X" in str(captured.get("feedback", ""))
|
||||||
|
# And the loop re-entered plan -> review and settled (not stuck at the gate).
|
||||||
|
assert final["current_phase"] == Phase.BUILD.value
|
||||||
|
|
||||||
|
|
||||||
|
def test_plan_gate_abandon_fails(restore_review_invoker) -> None:
|
||||||
|
# Command(resume={"decision":"abandon"}) -> terminal FAILED.
|
||||||
|
graph = _plan_gate_graph("VERDICT: REQUEST CHANGES\nnope")
|
||||||
|
thread_id, _ = _drive_to_plan_gate(graph)
|
||||||
|
final = resume_task(graph, thread_id=thread_id, answer={"decision": "abandon"})
|
||||||
|
|
||||||
|
assert final["status"] == TaskStatus.FAILED.value
|
||||||
|
|
||||||
|
|
||||||
|
def test_plan_gate_unrecognized_decision_fails_safe(restore_review_invoker) -> None:
|
||||||
|
# An unrecognized decision must NOT accidentally approve — it fails terminal.
|
||||||
|
graph = _plan_gate_graph("VERDICT: REQUEST CHANGES\nnope")
|
||||||
|
thread_id, _ = _drive_to_plan_gate(graph)
|
||||||
|
final = resume_task(graph, thread_id=thread_id, answer={"decision": "huh?"})
|
||||||
|
|
||||||
|
assert final["status"] == TaskStatus.FAILED.value
|
||||||
|
|
||||||
|
|
||||||
|
def test_plan_gate_request_changes_terminates_at_ceiling(
|
||||||
|
restore_review_invoker,
|
||||||
|
) -> None:
|
||||||
|
# Repeated request_changes resumes must eventually hit the gate ceiling and
|
||||||
|
# go terminal PARKED ("ceiling reached") — NOT loop unbounded. The reviewer
|
||||||
|
# NEVER approves, so every gate visit is request_changes until the ceiling.
|
||||||
|
from agent_team.graph import MAX_PLAN_GATE_VISITS, pending_question
|
||||||
|
|
||||||
|
graph = _plan_gate_graph("VERDICT: REQUEST CHANGES\nstill not ready")
|
||||||
|
thread_id, _ = _drive_to_plan_gate(graph)
|
||||||
|
|
||||||
|
# Each request_changes consumes exactly one gate visit; bound the loop well
|
||||||
|
# above the ceiling to prove it terminates on its own, not by our cap.
|
||||||
|
for _ in range(MAX_PLAN_GATE_VISITS + 5):
|
||||||
|
if pending_question(graph, thread_id=thread_id) is None:
|
||||||
|
break
|
||||||
|
resume_task(
|
||||||
|
graph,
|
||||||
|
thread_id=thread_id,
|
||||||
|
answer={"decision": "request_changes", "notes": "again"},
|
||||||
|
)
|
||||||
|
|
||||||
|
state = get_pipeline_state(graph, thread_id=thread_id)
|
||||||
|
assert pending_question(graph, thread_id=thread_id) is None
|
||||||
|
assert state["status"] == TaskStatus.PARKED.value
|
||||||
|
assert "ceiling" in (state.get("failure_reason") or "")
|
||||||
|
assert state.get("plan_gate_visits") == MAX_PLAN_GATE_VISITS
|
||||||
|
|
||||||
|
|
||||||
|
def test_auto_approved_plan_does_not_hit_gate(restore_review_invoker) -> None:
|
||||||
|
# No regression: an auto-APPROVED plan settles WITHOUT visiting the gate even
|
||||||
|
# when the gate is wired (gate only catches the review-cap dead-end).
|
||||||
|
from agent_team.graph import pending_question
|
||||||
|
|
||||||
|
graph = _plan_gate_graph("VERDICT: APPROVE\nlooks solid")
|
||||||
|
thread_id, _ = start_task(graph, transport="slack")
|
||||||
|
final = resume_task(graph, thread_id=thread_id, answer="scope is X")
|
||||||
|
|
||||||
|
assert final["current_phase"] == Phase.BUILD.value
|
||||||
|
assert pending_question(graph, thread_id=thread_id) is None
|
||||||
|
assert final.get("plan_gate_visits", 0) == 0
|
||||||
|
|
||||||
|
|
||||||
# --- P3 build -> verify subgraph wiring (opt-in). ---------------------------
|
# --- P3 build -> verify subgraph wiring (opt-in). ---------------------------
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
Reference in a new issue