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/nodes/planner.py
Adam Moussa 15a416d31a Add Plane-2 leaf scaffold (pipeline graph, nodes, HITL, transports, CI)
Consolidates the 18 leaf modules from the r720-plane2-scaffold workflow onto
the foundation commit. Full suite: 535 passed, 1 skipped; ruff + format clean.

Built (pre-deployment scaffold only — nothing provisioned/enabled):
- LangGraph pipeline graph.py (INTAKE->CLARIFY->PLAN, interrupt()/resume, checkpointer-injectable)
- nodes: clarifier (98% gate), planner, review_loop (GPT-4.1), builders->candidate diff, verifier
- §3.3.1 HITL: ledger ops, resume_worker, deadline_timer, recovery sweep, responder
- transports: slack / github / claude_code adapters
- ci_gate (pure-code pass/fail), operator_cli, run-team.py entry, P1 sim harness
- ci/agent-team-apply-verify.yml (split untrusted/privileged jobs) — authored, disabled

KNOWN OPEN FINDINGS (verifier/cross-review, not yet fixed — see follow-up):
- builders denylist: 4 execution-proven bypasses (delete, mode-change, copy-to, out-of-scope delete)
- §3.3.1 CAS: BEGIN IMMEDIATE outside try/except; shared-connection txn nesting unsafe under concurrency
- operator_cli: missing re-deliver/force-resume; audit-after-mutate ordering gap
- ci yaml: GPT-4.1 cross-review PASS w/ 4 FIX items (symlink path escape, etc.)
- P1 sim harness models the ledger layer, not real LangGraph interrupt/resume; P1 exit criteria not yet truly proven

Deploy-gated (NOT done): IAM/step-ca/Roles Anywhere/confluence-bot provisioning,
/sh-security-review sign-off, live Slack/CI, rsync, live dry-runs, Adam approval.
2026-06-17 15:16:12 -04:00

293 lines
11 KiB
Python

"""Planner node — Plane-2 pipeline PLAN stage (design §3.3, §7.1 P2).
The planner turns the clarified task context into a **phased plan** (the format
these design docs use) by calling Claude through the billing seam, then hands
the plan to the adversarial review loop.
Pipeline position (§3.3)::
INTAKE -> CLARIFIER -> [PLANNER] -> REVIEW LOOP -> BUILDERS -> VERIFIERS
Responsibilities of this leaf (Phase P2, §7.1):
* Read the clarified context out of :class:`~agent_team.task_model.PipelineState`
(the task description + the full clarifier Q&A history).
* On a **review loop-back** (§3.3 "Loops back to the planner on REQUEST
CHANGES"), fold the prior ``review_verdicts`` into the re-plan prompt so the
next plan answers the reviewer's objections.
* Call :func:`agent_team.billing.claude_invoke` to draft the plan, then parse
the model's reply into a structured ``plan`` dict.
* Return a **partial** ``PipelineState`` update advancing the phase to
:attr:`~agent_team.task_model.Phase.REVIEW`.
* Enforce the convergence bound (§3.3 "escalates to Adam if it cannot
converge"): after :data:`MAX_PLAN_REVISIONS` REQUEST-CHANGES loop-backs the
task is parked + ALARM-ed rather than spun (status ``PARKED``, phase
``PARKED``).
This module imports the committed foundation contracts verbatim
(:mod:`agent_team.task_model`, :mod:`agent_team.billing`) and does no I/O of its
own beyond the injected Claude seam — keeping the node a pure LangGraph
state-transition function that is trivially unit-testable.
"""
from __future__ import annotations
import json
import re
from typing import Any
from agent_team.billing import ClaudeResult, claude_invoke
from agent_team.task_model import Phase, PipelineState, TaskStatus
__all__ = [
"MAX_PLAN_REVISIONS",
"PlannerError",
"build_plan_prompt",
"parse_plan",
"plan_node",
]
# §3.3 / §6.6 — the planner must converge or park. After this many
# REQUEST-CHANGES loop-backs from the review stage the task is escalated to
# Adam (parked + ALARM) instead of looping forever.
MAX_PLAN_REVISIONS = 3
# The reviewer verdict string that sends a plan back to the planner. Kept here
# (rather than imported from a review node that does not exist yet) so this leaf
# stays self-contained; the review leaf will emit this same token.
_REQUEST_CHANGES = "REQUEST_CHANGES"
class PlannerError(Exception):
"""Raised when the planner cannot produce a usable plan.
Distinct from a *converged-but-rejected* plan (which loops back through
review) — this signals the planner itself failed (empty/garbled model
output), so the coordinator can fail the task rather than advance it.
"""
def _task_description(state: PipelineState) -> str:
"""Pull the task description out of the graph state.
Intake writes the originating ask; we look in the conventional places and
fall back to an empty string so a malformed state surfaces as an empty
prompt section rather than a ``KeyError`` inside the node.
"""
plan = state.get("plan") or {}
if isinstance(plan, dict):
desc = plan.get("task") or plan.get("description")
if isinstance(desc, str) and desc.strip():
return desc.strip()
desc = state.get("task") # type: ignore[call-overload]
if isinstance(desc, str) and desc.strip():
return desc.strip()
return ""
def _format_qa_history(qa_history: list[Any]) -> str:
"""Render the clarifier Q&A history into prompt text.
Each entry may be a ``{"question": ..., "answer": ...}`` mapping (the
clarifier's shape) or a plain string; both are handled so the planner does
not couple tightly to the clarifier's internal record format.
"""
lines: list[str] = []
for idx, entry in enumerate(qa_history, start=1):
if isinstance(entry, dict):
question = str(entry.get("question", "")).strip()
answer = str(entry.get("answer", "")).strip()
if question or answer:
lines.append(f"{idx}. Q: {question}\n A: {answer}")
else:
text = str(entry).strip()
if text:
lines.append(f"{idx}. {text}")
return "\n".join(lines)
def _format_review_feedback(review_verdicts: list[Any]) -> str:
"""Render prior review verdicts into re-plan guidance (§3.3 loop-back).
Only the most recent verdict drives the re-plan, but earlier ones are
summarised so the planner does not re-introduce already-rejected ideas.
"""
if not review_verdicts:
return ""
lines: list[str] = []
for idx, verdict in enumerate(review_verdicts, start=1):
if isinstance(verdict, dict):
decision = str(verdict.get("decision", "")).strip()
notes = str(verdict.get("notes") or verdict.get("comment") or "").strip()
lines.append(f"Review {idx} [{decision}]: {notes}".rstrip())
else:
lines.append(f"Review {idx}: {str(verdict).strip()}")
return "\n".join(lines)
def build_plan_prompt(state: PipelineState) -> str:
"""Build the Claude prompt that drafts (or re-drafts) the phased plan.
Pure string assembly over the graph state — no I/O — so the prompt shape is
directly unit-testable. On a loop-back (``review_verdicts`` present) the
prompt instructs Claude to revise the prior plan against the feedback rather
than start from scratch.
"""
description = _task_description(state)
qa = _format_qa_history(list(state.get("qa_history", [])))
feedback = _format_review_feedback(list(state.get("review_verdicts", [])))
prior_plan = state.get("plan")
sections = [
"You are the PLANNER stage of an agentic SDLC pipeline. Produce a "
"phased implementation plan for the task below. The clarifier has "
"already reached confidence with the human, so do not ask questions — "
"plan.",
"",
"## Task",
description or "(no task description provided)",
]
if qa:
sections += ["", "## Clarified context (Q&A)", qa]
if feedback:
# Loop-back: the review stage sent the prior plan back for changes.
sections += [
"",
"## Reviewer feedback on the previous plan (address every point)",
feedback,
]
if isinstance(prior_plan, dict) and prior_plan.get("phases"):
sections += [
"",
"## Previous plan (revise; do not restart from scratch)",
json.dumps(prior_plan, sort_keys=True, indent=2),
]
sections += [
"",
"## Output format",
"Return ONLY a JSON object with keys: "
'"summary" (string), "phases" (a non-empty list of objects each with '
'"name" and "steps", where "steps" is a non-empty list of strings). '
"Do not include prose outside the JSON.",
]
return "\n".join(sections)
def _strip_code_fence(text: str) -> str:
"""Strip a leading/trailing Markdown code fence if Claude wrapped the JSON."""
fenced = re.match(
r"^\s*```(?:json)?\s*\n(?P<body>.*?)\n?\s*```\s*$",
text,
flags=re.DOTALL | re.IGNORECASE,
)
if fenced:
return fenced.group("body")
return text
def parse_plan(text: str) -> dict[str, Any]:
"""Parse the model reply into a validated phased-plan dict.
Raises :class:`PlannerError` if the reply is not a JSON object with a
non-empty ``phases`` list of well-formed phases — a garbled plan must fail
loudly so the coordinator does not advance an empty plan into review.
"""
candidate = _strip_code_fence(text or "").strip()
if not candidate:
raise PlannerError("planner returned an empty response")
try:
data = json.loads(candidate)
except json.JSONDecodeError as exc:
raise PlannerError(f"planner reply was not valid JSON: {exc}") from exc
if not isinstance(data, dict):
raise PlannerError("planner reply JSON must be an object")
phases = data.get("phases")
if not isinstance(phases, list) or not phases:
raise PlannerError("plan must contain a non-empty 'phases' list")
normalized_phases: list[dict[str, Any]] = []
for idx, phase in enumerate(phases, start=1):
if not isinstance(phase, dict):
raise PlannerError(f"phase {idx} must be an object")
name = phase.get("name")
steps = phase.get("steps")
if not isinstance(name, str) or not name.strip():
raise PlannerError(f"phase {idx} is missing a non-empty 'name'")
if not isinstance(steps, list) or not steps:
raise PlannerError(f"phase {idx} ('{name}') has no steps")
normalized_steps = [str(step).strip() for step in steps if str(step).strip()]
if not normalized_steps:
raise PlannerError(f"phase {idx} ('{name}') has no non-empty steps")
normalized_phases.append({"name": name.strip(), "steps": normalized_steps})
summary = data.get("summary")
return {
"summary": str(summary).strip() if isinstance(summary, str) else "",
"phases": normalized_phases,
}
def _revision_count(state: PipelineState) -> int:
"""How many REQUEST-CHANGES loop-backs have happened so far (§3.3).
Counts the request-changes verdicts already recorded in the state; the
planner uses this to decide whether the next attempt is still within the
convergence bound.
"""
count = 0
for verdict in state.get("review_verdicts", []):
decision = (
verdict.get("decision") if isinstance(verdict, dict) else str(verdict)
)
if isinstance(decision, str) and decision.strip().upper() == _REQUEST_CHANGES:
count += 1
return count
def plan_node(
state: PipelineState,
config: dict[str, Any] | None = None,
) -> PipelineState:
"""LangGraph node: draft/refine the phased plan, then advance to REVIEW.
Returns a **partial** :class:`~agent_team.task_model.PipelineState` (the
keys this node owns) — LangGraph merges it into the checkpointed state.
Behaviour:
* **Convergence bound (§3.3, §6.6).** If the review stage has already sent
the plan back :data:`MAX_PLAN_REVISIONS` times, the planner does **not**
burn more Claude budget: it parks the task (status ``PARKED``, phase
``PARKED``) so the coordinator ALARMs Adam. This is the "escalate to Adam
if it cannot converge" path.
* **Plan / re-plan.** Otherwise it builds the prompt (folding in any review
feedback on a loop-back), calls :func:`claude_invoke`, parses the reply,
and returns the new ``plan`` with the phase advanced to ``REVIEW``.
``config`` is forwarded to the billing seam so the caller can pin the
billing mode (it is threaded through to :func:`claude_invoke` as ``config``).
"""
revisions = _revision_count(state)
if revisions >= MAX_PLAN_REVISIONS:
# Escalation ladder (§5/§3.3): stop looping, hand to the human.
return PipelineState(
status=TaskStatus.PARKED.value,
current_phase=Phase.PARKED.value,
)
prompt = build_plan_prompt(state)
result: ClaudeResult = claude_invoke(prompt, config=config)
plan = parse_plan(result.text)
# Record how many times we have planned so review/observability can see it.
plan["revision"] = revisions
return PipelineState(
plan=plan,
current_phase=Phase.REVIEW.value,
status=TaskStatus.ACTIVE.value,
)