diff --git a/agent-team/agent_team/ci_gate.py b/agent-team/agent_team/ci_gate.py index d789be7..7774807 100644 --- a/agent-team/agent_team/ci_gate.py +++ b/agent-team/agent_team/ci_gate.py @@ -473,7 +473,7 @@ def evaluate_ci_gate( candidate_diff: str, ledger_hash: str | None, ci_result: Mapping[str, Any] | None, - expected_run_id: str, + expected_run_id: str | None, allowed_scope: Sequence[str] | None = None, ) -> GateResult: """Make the deterministic pass/fail/block decision (§3.3.2 boundary #4). @@ -486,26 +486,47 @@ def evaluate_ci_gate( read-only PAT (GitHub Checks/Actions API). The gate reads only ``run_id``, ``conclusion``, and (optionally) ``diff_hash`` from it; it **never** reads a patch-written success file/artifact. - * ``expected_run_id`` — the run id the verifier dispatched for this exact - diff; the conclusion must be keyed to it (a stale/substituted run id is a - BLOCK). + * ``expected_run_id`` — the run id THIS task dispatched the apply/verify + workflow under (per-task, sourced from ``state["run_id"]`` at node-run + time, NOT a static wiring-time constant). The conclusion must be keyed to + it (a stale/substituted run id is a BLOCK). A ``None``/empty value means + the dispatcher captured no run id to bind to (e.g. the run-id poll fell + through) — there is nothing to anchor the verdict against, so the gate + BLOCKs (refuse-to-proceed, never a vacuous pass). * ``allowed_scope`` — optional declared-scope prefixes for the task. Decision order (a trust violation always wins over a CI verdict): - 1. **Denylist / scope** — any violation -> :data:`GateDecision.BLOCK`. - 2. **Diff-hash integrity** — hash mismatch (ledger or CI-verified) -> + 1. **Missing binding** — a ``None``/empty ``expected_run_id`` -> + :data:`GateDecision.BLOCK` (no per-task run to bind the verdict to). + 2. **Denylist / scope** — any violation -> :data:`GateDecision.BLOCK`. + 3. **Diff-hash integrity** — hash mismatch (ledger or CI-verified) -> ``BLOCK``. - 3. **Authenticated conclusion** — missing result, a ``run_id`` that does not + 4. **Authenticated conclusion** — missing result, a ``run_id`` that does not match ``expected_run_id``, or an unrecognised/ambiguous conclusion -> ``BLOCK``; a recognised failure -> :data:`GateDecision.FAIL`; ``success`` -> :data:`GateDecision.PASS`. Returns a :class:`GateResult` with the decision and the reasons behind it. - Raises :class:`CiGateError` on structurally invalid inputs. + Raises :class:`CiGateError` on structurally invalid inputs (a non-string + ``expected_run_id`` is structurally invalid; ``None``/empty is a BLOCK, not + an exception, because it is the legitimate "no run captured" runtime state). """ - if not isinstance(expected_run_id, str) or not expected_run_id: - raise CiGateError("expected_run_id must be a non-empty string") + if expected_run_id is not None and not isinstance(expected_run_id, str): + raise CiGateError("expected_run_id must be a string or None") + + # (0) Missing per-task binding: the dispatcher captured no run id for this + # task (None/empty). There is nothing to anchor the verdict to, so refuse to + # proceed rather than gating against an empty string (which a substituted + # ``ci_result`` with no/empty run_id could otherwise vacuously satisfy). + if not expected_run_id: + return GateResult( + decision=GateDecision.BLOCK, + reasons=["no per-task run_id to bind the verdict to (dispatch unresolved)"], + run_id=None, + diff_hash=ledger_hash, + ci_conclusion=None, + ) reasons: list[str] = [] diff --git a/agent-team/agent_team/ci_watcher.py b/agent-team/agent_team/ci_watcher.py new file mode 100644 index 0000000..9777370 --- /dev/null +++ b/agent-team/agent_team/ci_watcher.py @@ -0,0 +1,496 @@ +"""CI-watcher sweep — async resume-on-CI-complete for suspended VERIFY (design §3.3.2, P3 Decision 2). + +A CI apply/verify run takes ~7 minutes. The coordinator is a single durable +daemon: a multi-minute *blocking* VERIFY node would stall the ``tick()`` loop and +every other task. So the §4 Decision-2 shape is **async resume-on-CI-complete**, +NOT a blocking poll: + + BUILD → DISPATCH (push branch, trigger CI, capture run_id, suspend at VERIFY) + → [CI-watcher resumes on terminal conclusion] → VERIFY (read result, gate) + +This module is that CI-watcher: a ``tick()``-driven sweep that **mirrors +:mod:`agent_team.deadline_timer`'s shape** (a pure, restart-safe maintenance pass +whose side effects are injected callables, so it is unit-testable with no network +and no live graph). For each task suspended awaiting CI it: + +1. reads the task's ``run_id`` / ``dispatched_at`` (written ONLY by the trusted + dispatch node — never by an LLM/builder/verifier node, mirroring the + ``ci_fetcher`` TRUST SOURCE note); +2. polls that run **read-only** via an injected poller (the production default + reuses :func:`agent_team.ci_fetcher.fetch_ci_result` — a single read-only GET; + it NEVER writes, dispatches, or applies anything); and +3. acts on the poll outcome: + + * **TERMINAL** (the run concluded with a recognised conclusion) → RESUME the + suspended task. VERIFY then re-reads the now-terminal authenticated result + and the pure-code gate (:mod:`agent_team.ci_gate`) decides PASS/FAIL. The + watcher itself NEVER decides the verdict — it only un-suspends the task. + * **PENDING** (the run has no conclusion yet) → if ``dispatched_at + timeout`` + has elapsed, PARK the task (a CI run that never terminates must not wait + forever); otherwise leave it suspended for a later pass. + * **ERROR** (the read-only poll raised / could not produce a result) → PARK + the task. **Fail-closed:** a broken poll parks for a human rather than + spinning or fabricating progress. + +Fail-closed discipline (mirroring ``ci_fetcher`` / ``deadline_timer``): + +* A task whose state carries no usable ``run_id`` or ``dispatched_at`` is + **parked immediately** (a suspended-awaiting-CI task with no run to watch is + unrecoverable by polling — it can never resume on a conclusion that has no run + id to read). This is the "None result" fail-closed case. +* Any poll exception is isolated per task (caught → PARK that one task) so one + task's poll failure never aborts the whole sweep or crashes the daemon. +* All inputs are read fresh from the durable task state each pass; the watcher + holds no in-memory-only state, so a reboot mid-sweep simply re-runs the + remaining suspended tasks on the next pass. RESUME is idempotent via the + resume worker's turn guard; PARK is idempotent (a parked task is no longer + surfaced as awaiting-CI). + +Side effects (RESUME, PARK, ALARM) are injected as callables — exactly what makes +the terminal/timeout/error branches testable without a live graph, a network, or +a real coordinator. +""" + +from __future__ import annotations + +import logging +from collections.abc import Callable, Mapping +from dataclasses import dataclass, field +from datetime import datetime, timedelta, timezone +from enum import Enum +from typing import Any + +__all__ = [ + "AlarmFn", + "CiPollOutcome", + "CiPollResult", + "CiWatchAction", + "CiWatchOutcome", + "CiWatchReport", + "ParkFn", + "PendingCiTask", + "PollFn", + "ResumeFn", + "default_ci_poller", + "run_ci_watcher", +] + +logger = logging.getLogger(__name__) + +# Conservative default: a CI apply/verify run that has not reached a terminal +# conclusion within this window is treated as stuck and the task parks (§4 +# Decision 2 "dispatched_at + timeout elapses -> park"). A normal run is ~7 min; +# 30 min leaves generous headroom for queueing / retries before we ALARM. +DEFAULT_CI_TIMEOUT = timedelta(minutes=30) + + +class CiPollOutcome(Enum): + """The classification of one read-only CI poll (the poller's verdict). + + ``TERMINAL`` — the run reached a recognised, authenticated conclusion (the + poll ``result`` mapping is populated); the task should RESUME so VERIFY can + gate it. ``PENDING`` — the run is still in progress (no conclusion yet); keep + waiting unless the timeout has elapsed. ``ERROR`` — the poll could not produce + a result (a fetch error / fail-closed ``None``); the task PARKS. + """ + + TERMINAL = "terminal" + PENDING = "pending" + ERROR = "error" + + +@dataclass(frozen=True) +class CiPollResult: + """The outcome of one read-only poll of a task's CI run. + + ``outcome`` is the :class:`CiPollOutcome`; ``result`` carries the + authenticated CI conclusion mapping (``run_id`` / ``conclusion`` / + ``diff_hash``) when (and only when) ``outcome`` is ``TERMINAL`` — the same + DATA shape :func:`agent_team.ci_gate.evaluate_ci_gate` consumes. The watcher + never derives a verdict from it; it only decides resume-vs-wait-vs-park. + """ + + outcome: CiPollOutcome + result: Mapping[str, Any] | None = None + + @classmethod + def terminal(cls, result: Mapping[str, Any]) -> CiPollResult: + """A run that reached a recognised terminal conclusion.""" + return cls(outcome=CiPollOutcome.TERMINAL, result=result) + + @classmethod + def pending(cls) -> CiPollResult: + """A run still in progress (no authenticated conclusion yet).""" + return cls(outcome=CiPollOutcome.PENDING) + + @classmethod + def error(cls) -> CiPollResult: + """A poll that could not produce a result (fail-closed).""" + return cls(outcome=CiPollOutcome.ERROR) + + +class CiWatchAction(Enum): + """The action the sweep actually took for one suspended-awaiting-CI task. + + ``RESUMED`` — the run terminated, the task was resumed so VERIFY can gate it. + ``PARKED_TIMEOUT`` — the run never terminated within the timeout; the task + parked. ``PARKED_ERROR`` — the poll failed / the state was unusable; the task + parked (fail-closed). ``WAITING`` — the run is still in progress and within + the timeout; the task stays suspended for a later pass. + """ + + RESUMED = "resumed" + PARKED_TIMEOUT = "parked_timeout" + PARKED_ERROR = "parked_error" + WAITING = "waiting" + + +@dataclass(frozen=True) +class PendingCiTask: + """A minimal read-snapshot of one task suspended awaiting CI (§3.3.2). + + Only the fields the watcher needs: identity (``thread_id``) and the trusted + dispatch watermarks (``run_id`` / ``dispatched_at``) the poll + timeout key + off. Frozen because it is a snapshot — the watcher never mutates a task + in-place; it acts through the injected resume/park callables. + """ + + thread_id: str + run_id: str | None = None + dispatched_at: str | None = None + + @classmethod + def from_state(cls, thread_id: str, state: Mapping[str, Any]) -> PendingCiTask: + """Build a :class:`PendingCiTask` from a durable task ``state`` mapping.""" + run_id = state.get("run_id") + dispatched_at = state.get("dispatched_at") + return cls( + thread_id=thread_id, + run_id=str(run_id) if run_id is not None and str(run_id) != "" else None, + dispatched_at=( + str(dispatched_at) + if dispatched_at is not None and str(dispatched_at) != "" + else None + ), + ) + + +@dataclass(frozen=True) +class CiWatchOutcome: + """The result of processing one suspended-awaiting-CI task this pass.""" + + thread_id: str + action: CiWatchAction + run_id: str | None = None + error: str | None = None + + +@dataclass +class CiWatchReport: + """Aggregate result of one CI-watcher pass (mirrors :class:`~agent_team.deadline_timer.TimerLoopReport`). + + ``outcomes`` is one entry per suspended task examined. The summary counters + let the coordinator decide whether to ALARM (any ``parked_error``) without + re-walking the list. + """ + + outcomes: list[CiWatchOutcome] = field(default_factory=list) + + @property + def examined(self) -> int: + """Number of suspended-awaiting-CI tasks examined this pass.""" + return len(self.outcomes) + + @property + def resumed(self) -> int: + """Tasks whose run terminated and were resumed.""" + return sum(1 for o in self.outcomes if o.action is CiWatchAction.RESUMED) + + @property + def parked_timeout(self) -> int: + """Tasks parked because their run never terminated within the timeout.""" + return sum(1 for o in self.outcomes if o.action is CiWatchAction.PARKED_TIMEOUT) + + @property + def parked_error(self) -> int: + """Tasks parked because the poll failed / the state was unusable.""" + return sum(1 for o in self.outcomes if o.action is CiWatchAction.PARKED_ERROR) + + @property + def waiting(self) -> int: + """Tasks still in progress and left suspended for a later pass.""" + return sum(1 for o in self.outcomes if o.action is CiWatchAction.WAITING) + + @property + def parked(self) -> int: + """Total tasks parked this pass (timeout + error).""" + return self.parked_timeout + self.parked_error + + +# Injected seams. Keeping these as callables means the watcher performs no +# transport / graph / network I/O of its own (testable, and faithful to the +# deadline_timer shape). +PollFn = Callable[[PendingCiTask], CiPollResult] +"""Poll one task's CI run read-only and classify it. The production default +(:func:`default_ci_poller`) reuses :func:`agent_team.ci_fetcher.fetch_ci_result` +(a single read-only GET; never a write).""" + +ResumeFn = Callable[[PendingCiTask, Mapping[str, Any]], None] +"""Called once per *terminated* task to RESUME it (drive the suspended VERIFY +node forward). Receives the task snapshot + the authenticated terminal result.""" + +ParkFn = Callable[[PendingCiTask], None] +"""Called once per task that must PARK (timeout, poll error, or unusable +state). Parks the task and (typically) raises an ALARM.""" + +AlarmFn = Callable[[str], None] +"""Optional ALARM hook for the coordinator (one message per parked task).""" + + +def _utc_now() -> datetime: + """Return the current UTC time (injectable via ``now`` in the loop).""" + return datetime.now(timezone.utc) + + +def _parse_iso(value: str | None) -> datetime | None: + """Parse an ISO-8601 timestamp to an aware UTC datetime, or ``None``. + + A missing / unparseable ``dispatched_at`` yields ``None`` so the caller fails + closed (treats the task as unusable → park) rather than crashing the sweep. + """ + if not value: + return None + try: + parsed = datetime.fromisoformat(value) + except (TypeError, ValueError): + return None + if parsed.tzinfo is None: + # The ledger always writes UTC; treat a naive stamp as UTC rather than + # raising on the aware/naive compare below. + parsed = parsed.replace(tzinfo=timezone.utc) + return parsed + + +def default_ci_poller( + *, + owner: str, + repo: str, + client: Any = None, +) -> PollFn: + """Build the production read-only CI poller (reuses :mod:`agent_team.ci_fetcher`). + + Returns a :data:`PollFn` that, given a :class:`PendingCiTask`, performs ONE + read-only GET of the task's run via + :func:`agent_team.ci_fetcher.fetch_ci_result` (closing over ``owner`` / ``repo`` + and an optional injected ``client`` test double) and classifies it: + + * a populated mapping (a recognised, authenticated terminal conclusion) → + :meth:`CiPollResult.terminal`; + * ``None`` (the run is still in progress — ``fetch_ci_result`` returns ``None`` + for a null/unrecognised conclusion) → :meth:`CiPollResult.pending`; the + watcher keeps waiting until the timeout, so an in-progress run never blocks + the daemon and never parks prematurely; + * any exception raised by the fetch → :meth:`CiPollResult.error` (fail-closed: + a broken read-only poll parks the task). + + The poller NEVER writes: ``fetch_ci_result`` issues exactly one read-only GET + and fails closed to ``None`` on every error path. It NEVER decides the + verdict — that stays with the pure-code gate after VERIFY resumes. + """ + from agent_team.ci_fetcher import fetch_ci_result + + def poll(task: PendingCiTask) -> CiPollResult: + # The fetcher reads ``state["run_id"]`` / ``state["diff_hash"]``; rebuild + # the minimal state it needs from the trusted dispatch watermarks. (We do + # not have the full task state here — only what the watcher snapshotted — + # which is exactly the read-only run identity the fetcher requires.) + state = {"run_id": task.run_id} + try: + result = fetch_ci_result(state, owner=owner, repo=repo, client=client) + except Exception: # noqa: BLE001 - any fetch failure fails closed -> park + logger.warning( + "ci_watcher: read-only poll for task %s raised; failing closed", + task.thread_id, + exc_info=True, + ) + return CiPollResult.error() + if isinstance(result, Mapping): + return CiPollResult.terminal(result) + # None: the run has no recognised conclusion yet (in progress). Keep + # waiting; the timeout branch parks a run that never terminates. + return CiPollResult.pending() + + return poll + + +def run_ci_watcher( + pending_tasks: list[PendingCiTask], + *, + poll: PollFn, + on_resume: ResumeFn, + on_park: ParkFn, + timeout: timedelta = DEFAULT_CI_TIMEOUT, + now: datetime | None = None, +) -> CiWatchReport: + """Run one CI-watcher pass over the suspended-awaiting-CI tasks (§3.3.2 Decision 2). + + ``pending_tasks`` is the set of tasks currently suspended at VERIFY awaiting + CI (the coordinator supplies them from the durable task store). For each: + + 1. **Unusable state → PARK (fail-closed).** A task with no ``run_id`` or no + parseable ``dispatched_at`` cannot be watched (there is no run to poll, or + no watermark to time out against), so it parks immediately rather than + waiting forever on a run it can never read. + 2. **Poll the run read-only.** Call ``poll`` (the production default reuses + :func:`agent_team.ci_fetcher.fetch_ci_result`; never a write). A poll that + raises is isolated to this one task (→ PARK) so it cannot abort the sweep. + 3. **Act on the outcome:** + * ``TERMINAL`` → ``on_resume`` (resume the task; VERIFY re-reads the + now-terminal result and the pure-code gate decides). The watcher never + decides the verdict. + * ``PENDING`` → if ``now >= dispatched_at + timeout`` → ``on_park`` + (timeout); else leave suspended (``WAITING``) for a later pass. + * ``ERROR`` → ``on_park`` (fail-closed). + + Side effects are isolated per task: if ``on_resume`` / ``on_park`` raises, the + failure is recorded for that task and the sweep continues with the rest of the + batch rather than aborting the whole pass. + + Restart-safety: inputs are the durable task snapshots, resume is idempotent + (the resume worker's turn guard), and park is idempotent, so re-running the + pass after a crash safely processes only the still-suspended tasks. + + Returns a :class:`CiWatchReport` describing what happened to each task. + """ + current = now or _utc_now() + report = CiWatchReport() + + for task in pending_tasks: + outcome = _process_task( + task, + poll=poll, + on_resume=on_resume, + on_park=on_park, + timeout=timeout, + now=current, + ) + report.outcomes.append(outcome) + + return report + + +def _process_task( + task: PendingCiTask, + *, + poll: PollFn, + on_resume: ResumeFn, + on_park: ParkFn, + timeout: timedelta, + now: datetime, +) -> CiWatchOutcome: + """Process one suspended task: poll its run, then resume / wait / park. + + Isolates each task's side effect: a poll OR a resume/park callback that raises + yields a :attr:`CiWatchAction.PARKED_ERROR` outcome (best-effort: the task is + parked if it can be) instead of crashing the whole sweep. + """ + # (1) Unusable state -> fail closed (park). No run to poll / no watermark to + # time out against means polling can never recover this task. + if task.run_id is None or _parse_iso(task.dispatched_at) is None: + logger.warning( + "ci_watcher: task %s has no usable run_id/dispatched_at; parking " + "(fail-closed)", + task.thread_id, + ) + return _park(task, on_park, CiWatchAction.PARKED_ERROR) + + # (2) Read-only poll. A raising poll is isolated to this task (-> park). + try: + poll_result = poll(task) + except Exception as exc: # noqa: BLE001 - one task's poll failure -> park it + logger.warning( + "ci_watcher: poll for task %s raised (%s); parking (fail-closed)", + task.thread_id, + type(exc).__name__, + ) + return _park( + task, + on_park, + CiWatchAction.PARKED_ERROR, + error=f"{type(exc).__name__}: {exc}", + ) + + # (3) Act on the classified outcome. + if poll_result.outcome is CiPollOutcome.TERMINAL: + result = poll_result.result or {} + try: + on_resume(task, result) + except Exception as exc: # noqa: BLE001 - isolate one task's resume failure + logger.exception( + "ci_watcher: resume of task %s (terminal run) raised", task.thread_id + ) + return CiWatchOutcome( + thread_id=task.thread_id, + action=CiWatchAction.PARKED_ERROR, + run_id=task.run_id, + error=f"{type(exc).__name__}: {exc}", + ) + return CiWatchOutcome( + thread_id=task.thread_id, + action=CiWatchAction.RESUMED, + run_id=task.run_id, + ) + + if poll_result.outcome is CiPollOutcome.ERROR: + # Fail-closed: a poll that could not produce a result parks the task. + return _park(task, on_park, CiWatchAction.PARKED_ERROR) + + # PENDING: still in progress. Park only if the dispatch timeout has elapsed; + # otherwise leave it suspended for a later pass (the async-wait, not a block). + dispatched = _parse_iso(task.dispatched_at) + assert dispatched is not None # narrowed by the unusable-state guard above + if now >= dispatched + timeout: + logger.warning( + "ci_watcher: task %s run %s did not terminate within %s; parking (timeout)", + task.thread_id, + task.run_id, + timeout, + ) + return _park(task, on_park, CiWatchAction.PARKED_TIMEOUT) + + return CiWatchOutcome( + thread_id=task.thread_id, + action=CiWatchAction.WAITING, + run_id=task.run_id, + ) + + +def _park( + task: PendingCiTask, + on_park: ParkFn, + action: CiWatchAction, + *, + error: str | None = None, +) -> CiWatchOutcome: + """Invoke the injected park side effect, isolating a callback failure. + + A ``on_park`` that raises is downgraded to a ``PARKED_ERROR`` outcome (the + sweep continues) rather than crashing the whole pass — the same per-row + side-effect isolation :mod:`agent_team.deadline_timer` uses. + """ + try: + on_park(task) + except Exception as exc: # noqa: BLE001 - isolate one task's park failure + logger.exception("ci_watcher: park of task %s raised", task.thread_id) + return CiWatchOutcome( + thread_id=task.thread_id, + action=CiWatchAction.PARKED_ERROR, + run_id=task.run_id, + error=f"{type(exc).__name__}: {exc}", + ) + return CiWatchOutcome( + thread_id=task.thread_id, + action=action, + run_id=task.run_id, + error=error, + ) diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 474e112..e746d5f 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -75,6 +75,7 @@ __all__ = [ "default_clarify_node_factory", "default_dispatch_node_factory", "default_slack_listener_factory", + "failsafe_production_p3_wiring", "gated_build_verify_wiring", ] @@ -355,7 +356,7 @@ def gated_build_verify_wiring( *, owner: str, repo: str, - expected_run_id: str, + expected_run_id: str | None = None, allowed_scope: list[str] | None = None, diff_builder: Any = None, ci_client: Any = None, @@ -386,6 +387,14 @@ def gated_build_verify_wiring( task parks). ``ci_client`` injects a test double; the real path builds a read-only ``requests`` session at call time from the read-only token env var. + Per-task run-id binding (design §4 Decision 4): the gate binds each task's + verdict to the run id THAT TASK dispatched (persisted as ``state["run_id"]`` + by the dispatch node and read by the verifier at node-run time), NOT a + static wiring-time constant. ``expected_run_id`` is therefore optional and + defaults to ``None``; when supplied it is only a static fallback for a + harness that drives the verifier without a per-task ``state["run_id"]``. A + task whose dispatch left no run id BLOCKs (never a vacuous pass). + Lazy-imported (ci_fetcher pulls the subgraph + verifier leaves) for the same import-hygiene reason as the other factories. """ @@ -439,6 +448,97 @@ def default_dispatch_node_factory() -> "Callable[[Any], Any]": return make_dispatch_node(owner=owner, repo=repo, base=base) +# The one-line operator notice posted to #agent-team when the live P3 wiring +# cannot bind (missing owner/repo/CI-read token) and the daemon degrades to the +# inert P3 path. Goes through the lifecycle NOTIFY sink (not the park ALARM): an +# unconfigured box is an operational state, not a parked task. +_P3_INERT_NOTICE = ( + "ℹ️ P3 build→verify/dispatch is INERT this run: " + "AGENT_TEAM_REPO_OWNER / AGENT_TEAM_REPO_NAME (and a CI-read token: " + "AGENT_TEAM_CI_READ_TOKEN or GITHUB_TOKEN) are not all set. The daemon is " + "up and tasks run through PLAN/REVIEW; any task reaching BUILD/VERIFY will " + "PARK (fail-closed) until the P3 env is provisioned." +) + + +def _p3_env_is_configured() -> bool: + """Return True iff the live P3 wiring can bind from the environment. + + Live P3 (gated build→verify + auto-dispatch) needs the dispatch target + (``AGENT_TEAM_REPO_OWNER`` / ``AGENT_TEAM_REPO_NAME`` — the same vars + :func:`default_dispatch_node_factory` requires) AND a read-only CI token for + the verifier's authenticated conclusion read (``AGENT_TEAM_CI_READ_TOKEN``, + falling back to ``GITHUB_TOKEN`` — mirrors + :func:`agent_team.ci_fetcher._resolve_read_token`). Any missing piece means + the gate could never read an authenticated pass, so we keep the whole P3 + subgraph OFF rather than wire a half-configured, always-BLOCKing path. + """ + owner = os.environ.get("AGENT_TEAM_REPO_OWNER", "").strip() + repo = os.environ.get("AGENT_TEAM_REPO_NAME", "").strip() + ci_token = ( + os.environ.get("AGENT_TEAM_CI_READ_TOKEN", "").strip() + or os.environ.get("GITHUB_TOKEN", "").strip() + ) + return bool(owner and repo and ci_token) + + +def failsafe_production_p3_wiring( + *, + notify: "Callable[..., None] | None" = None, +) -> "tuple[BuildVerifyWiring | None, DispatchNodeFactory | None]": + """Resolve the production P3 wiring fail-safe (Decision 5; serve default). + + Per the Phase-0 design the bound P3 wiring is the new ``serve`` default, but + its factories are called EAGERLY at graph-build (``setup`` calls + ``self._build_verify_wiring()`` / ``self._dispatch_node_wiring()``), and the + live dispatch factory RAISES when ``AGENT_TEAM_REPO_OWNER`` / + ``AGENT_TEAM_REPO_NAME`` are unset. A raise there would crash-loop the + daemon at serve-start — exactly the failure mode this wrapper exists to + prevent. + + So this resolver decides ONCE, up front, from the environment: + + * **Configured** (:func:`_p3_env_is_configured` — owner + repo + a CI-read + token all present) → returns the LIVE pair: a + :func:`gated_build_verify_wiring` bound to the env owner/repo (so the + verifier reads the authenticated conclusion via the read-only fetcher) and + :func:`default_dispatch_node_factory` (which re-reads the same env at + build time). Tasks reaching P3 run BUILD → DISPATCH → VERIFY. + * **Unconfigured** → returns ``(None, None)`` — the INERT P3 path: no + build→verify subgraph, no dispatch, so a task reaching P3 simply PARKS + (fail-closed; never a fabricated pass). Logs exactly ONE WARNING and emits + ONE ``#agent-team`` inert-mode notice via the lifecycle ``notify`` sink + (NOT the park-ALARM path: an unprovisioned box is an operational state, + not a parked task). The notify sink is best-effort and fully guarded so a + Slack failure here never blocks serve-start. + + NEVER raises: serve-start must come up either fully wired or inert, but it + must always come up. + """ + if _p3_env_is_configured(): + owner = os.environ.get("AGENT_TEAM_REPO_OWNER", "").strip() + repo = os.environ.get("AGENT_TEAM_REPO_NAME", "").strip() + return ( + lambda: gated_build_verify_wiring(owner=owner, repo=repo), + default_dispatch_node_factory, + ) + + _LOG.warning( + "P3 build→verify/dispatch wiring is INERT: AGENT_TEAM_REPO_OWNER / " + "AGENT_TEAM_REPO_NAME (and a CI-read token) are not all set. The " + "coordinator starts and runs PLAN/REVIEW; any task reaching BUILD/VERIFY " + "will park (fail-closed) until the P3 env is provisioned." + ) + if notify is not None: + try: + notify(_P3_INERT_NOTICE) + except Exception: # noqa: BLE001 - an inert-notice failure must not block serve + _LOG.warning( + "inert-mode notify failed; serve still starting inert", exc_info=True + ) + return None, None + + class Coordinator: """Owns the live Plane-2 runtime: graph + resume worker + transport (§3.3). @@ -470,6 +570,9 @@ class Coordinator: build_listener: ListenerFactory | None = None, new_task_callback: "Callable[[str, str, str], str] | None" = None, notify: "Callable[..., None] | None" = None, + ci_pending_provider: "Callable[[], list[Any]] | None" = None, + ci_poller: "Callable[[Any], Any] | None" = None, + ci_timeout: timedelta | None = None, ) -> None: self._db_path = Path(db_path) self._transport = transport @@ -507,6 +610,17 @@ class Coordinator: # behavior). The serve path injects a Slack poster so a task is never a # black box: the human sees parked / needs-more-input / plan-ready. self._notify = notify + # CI-watcher seams (P3 async resume-on-CI-complete; §4 Decision 2). Both + # OPT-IN and default None, so the CI sweep in tick() is a NO-OP unless the + # live P3 path provides them: ``ci_pending_provider`` enumerates tasks + # suspended at VERIFY awaiting CI (as + # :class:`agent_team.ci_watcher.PendingCiTask`), and ``ci_poller`` is the + # read-only poll seam (the default reuses + # :func:`agent_team.ci_fetcher.fetch_ci_result`). Left None, no CI sweep + # runs — exactly the INERT default and the unit-test path. + self._ci_pending_provider = ci_pending_provider + self._ci_poller = ci_poller + self._ci_timeout = ci_timeout # Built by setup(). self._graph: Any = None @@ -1076,14 +1190,19 @@ class Coordinator: ) def tick(self) -> list[ResumeResult]: - """One maintenance pass: deadline sweep + park policy, then drain (§3.3.1). + """One maintenance pass: deadline sweep + CI-watch sweep, then drain (§3.3.1, §3.3.2). Runs :func:`agent_team.responder.deadline_sweep` to flip overdue ``open`` questions to ``expired`` (the deterministic answer-vs-expiry race), then applies the park policy to each newly-expired id (raise the ALARM hook — §6.6 "ALARM rather than spin"; the durable ledger row is already - ``expired``, which is the task's parked state for P1). Finally drains any - resume jobs that landed. Returns the drain results. + ``expired``, which is the task's parked state for P1). It then runs the + CI-watcher sweep (:meth:`_ci_watch`) alongside the deadline sweep — the P3 + async resume-on-CI-complete pass that resumes tasks whose CI run + terminated and parks tasks whose run timed out / failed to poll (§4 + Decision 2). Both sweeps are fail-soft. Finally drains any resume jobs + that landed (including resumes the CI-watch sweep enqueued). Returns the + drain results. """ conn = connect(self._db_path) try: @@ -1094,10 +1213,117 @@ class Coordinator: for question_id in expired: self._park(question_id) + # CI-watcher sweep alongside the deadline sweep (§3.3.2 Decision 2). NO-OP + # unless the live P3 seams are wired; fail-soft so a CI-watch error never + # breaks the maintenance loop. + self._ci_watch() + results = self.drain_resumes() self._post_resume_followups(results) return results + def _ci_watch(self) -> Any: + """Run one CI-watcher sweep over tasks suspended awaiting CI (§3.3.2 Decision 2). + + NO-OP unless BOTH CI-watcher seams are wired (``ci_pending_provider`` + + ``ci_poller``) — the INERT default and the unit-test path skip it + entirely. When wired, it: + + 1. enumerates the tasks currently suspended at VERIFY awaiting CI (via + ``ci_pending_provider``, as + :class:`agent_team.ci_watcher.PendingCiTask`); and + 2. runs :func:`agent_team.ci_watcher.run_ci_watcher` with the read-only + ``ci_poller`` and injected RESUME / PARK side effects: RESUME enqueues + a turn-guarded resume onto the shared queue (drained in the same tick), + so the suspended VERIFY node re-reads the now-terminal CI result and + the pure-code gate decides; PARK marks the task parked + ALARMs. + + Fail-soft: any error in enumeration or the sweep is logged and swallowed + so a CI-watch failure never breaks the tick loop. Returns the + :class:`~agent_team.ci_watcher.CiWatchReport` (or ``None`` when skipped / + on error) for logging/tests. + """ + if self._ci_pending_provider is None or self._ci_poller is None: + return None + + from agent_team.ci_watcher import DEFAULT_CI_TIMEOUT, run_ci_watcher + + try: + pending = self._ci_pending_provider() + except Exception: # noqa: BLE001 - an enumeration failure must not break tick + _LOG.warning("ci-watch: pending-task enumeration raised", exc_info=True) + return None + + if not pending: + return None + + try: + report = run_ci_watcher( + list(pending), + poll=self._ci_poller, + on_resume=self._ci_resume, + on_park=self._ci_park, + timeout=self._ci_timeout or DEFAULT_CI_TIMEOUT, + ) + except Exception: # noqa: BLE001 - a sweep failure must not break tick + _LOG.warning("ci-watch: sweep raised", exc_info=True) + return None + + if report.parked_error: + _LOG.error( + "ci-watch: %d task(s) parked on a poll/state error (fail-closed)", + report.parked_error, + ) + return report + + def _ci_resume(self, task: Any, result: Any) -> None: + """RESUME a task whose CI run terminated: enqueue a turn-guarded resume. + + The suspended VERIFY node interrupted with a payload but no pending + ledger question (it is a machine gate, not a human gate), so there is no + ``question_id`` / ``answer`` to thread through the responder's + first-answer-wins flip. We drive the graph directly via the resume worker, + which re-runs VERIFY; that node RE-FETCHES the now-terminal authenticated + CI result (it never trusts the resume payload) and the pure-code gate + decides. Best-effort and isolated: a resume failure for one task parks it + rather than crashing the sweep (the watcher records PARKED_ERROR). + """ + from agent_team.resume_worker import build_resume_command # noqa: PLC0415 + + if self._graph is None: + raise RuntimeError("Coordinator._ci_resume called before setup()") + # The resume value is intentionally ignored by VERIFY (it re-fetches the + # authenticated conclusion), so any payload works; pass the terminal + # result for operator-log provenance. + self._graph.invoke( + build_resume_command(result), + graph_mod.thread_config(task.thread_id), + ) + + def _ci_park(self, task: Any) -> None: + """PARK a task whose CI run timed out / could not be polled (fail-closed). + + Flips the durable task status to PARKED via ``update_state`` (so the + terminal state is durable across a reboot) and raises the ALARM hook so + the stall is surfaced rather than silently spun on (§6.6). Best-effort: + a write failure is logged; the watcher still records the park outcome. + """ + from agent_team.task_model import Phase, TaskStatus # noqa: PLC0415 + + try: + self._graph.update_state( + graph_mod.thread_config(task.thread_id), + { + "status": TaskStatus.PARKED.value, + "current_phase": Phase.PARKED.value, + }, + ) + except Exception: # noqa: BLE001 - best-effort durable park + _LOG.warning( + "ci-watch: could not mark task %s parked", task.thread_id, exc_info=True + ) + self._alarm_hook(task.thread_id) + def _park(self, question_id: str) -> None: """Apply the park policy to one expired question (§6.6 ALARM, not spin). diff --git a/agent-team/agent_team/graph.py b/agent-team/agent_team/graph.py index 9983f3a..efb8d19 100644 --- a/agent-team/agent_team/graph.py +++ b/agent-team/agent_team/graph.py @@ -135,9 +135,11 @@ PARKED_ROUTE = "parked" # 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. +# 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 @@ -445,14 +447,23 @@ def build_graph( 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"`` 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 @@ -472,9 +483,10 @@ def build_graph( 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." + "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." ) builder: StateGraph = StateGraph(PipelineState) @@ -520,30 +532,29 @@ def build_graph( route_review, {BUILD_ROUTE: BUILD_NODE, PLAN: PLAN, PARKED_ROUTE: END}, ) - builder.add_edge(BUILD_NODE, VERIFY_NODE) + # 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: - # 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}, - ) + builder.add_edge(BUILD_NODE, VERIFY_NODE) 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) + 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() diff --git a/agent-team/agent_team/nodes/build_verify_subgraph.py b/agent-team/agent_team/nodes/build_verify_subgraph.py index fba59ec..24fe356 100644 --- a/agent-team/agent_team/nodes/build_verify_subgraph.py +++ b/agent-team/agent_team/nodes/build_verify_subgraph.py @@ -2,7 +2,13 @@ This module is the **wiring topology** for the Plane-2 build -> verify stage:: - ... -> REVIEW (route "build") -> BUILD -> VERIFY -> {approved | build | parked} + ... -> REVIEW (route "build") -> BUILD -> [DISPATCH] -> VERIFY + -> {approved | build | parked} + +DISPATCH is the optional CI-trigger vertex :func:`agent_team.graph.build_graph` +splices between BUILD and VERIFY when a dispatch node is wired (design §4 +Decision 1): it triggers the org CI run and captures ``state["run_id"]`` BEFORE +VERIFY reads it. With no dispatch node the order is simply ``BUILD -> VERIFY``. It produces the BUILD node, the VERIFY node, and the :func:`route_after_verify` conditional-edge function so the Integrate phase can @@ -189,18 +195,58 @@ def make_verify_node( than smuggling something past the gate. The real read-only-PAT fetcher is bound via :func:`bind_ci_result_fetcher` only after the §3.3.2 trust boundary clears its security gate. + + PER-TASK run-id binding (design §4 Decision 4): the expected run id the gate + binds the verdict to is NOT baked into ``config`` at factory time. The + dispatch node persists the run id THIS task dispatched as ``state["run_id"]``, + and :func:`agent_team.nodes.verifier.verifier_node` reads it from the state + threaded through below (``config.expected_run_id`` is only a static fallback + for harnesses with no per-task run id). So a single ``config`` shared across + tasks still gates each task against its OWN dispatched run, and a task whose + dispatch left no run id BLOCKs — never a vacuous pass. + + ASYNC CI-WAIT (design §4 Decision 2): a CI apply/verify run takes ~7 minutes, + and a multi-minute *blocking* fetch here would stall the coordinator daemon's + tick loop and every other task. So when there IS a dispatched run to wait for + (``state["run_id"]`` is set) but the fetch yields no terminal result yet (the + run is still in progress → ``None``), the node SUSPENDS via + :func:`~langgraph.types.interrupt` — exactly the durable suspend/resume shape + the clarify human-gate node uses. The CI-watcher sweep + (:func:`agent_team.ci_watcher.run_ci_watcher`) polls the run read-only and + RESUMES this node once the run reaches a terminal conclusion; on resume the + node RE-FETCHES the now-terminal result and the pure-code gate decides. If the + re-fetched result is still not terminal (e.g. a spurious resume), the node + falls through to the gate, which BLOCKs/parks — fail-closed, never a vacuous + pass. With NO dispatched run (``state["run_id"]`` absent — the INERT path or a + harness), the node does NOT suspend: a ``None`` fetch flows straight to the + gate, which BLOCKs and parks exactly as before. """ fetcher: CiResultFetcher = ( ci_result_fetcher if ci_result_fetcher is not None else _no_ci_result ) def node(state: PipelineState) -> PipelineState: - fetched = fetcher(state) - ci_result = fetched if isinstance(fetched, Mapping) else None + ci_result = _fetch_ci_result(fetcher, state) - # Merge the (possibly None) fetched CI result into the state the node - # reads from, WITHOUT mutating the caller's state object. The node reads - # ``ci_results``; a None result leaves the gate with nothing to pass on. + # Async CI-wait: only when a run was actually dispatched (state["run_id"] + # is set) AND it has no terminal result yet do we suspend, so the daemon + # never blocks on an in-progress run. The CI-watcher resumes us on a + # terminal conclusion; we re-fetch once after resume. The INERT/no-run + # path (no run_id) skips this and lets the gate BLOCK/park as before. + run_id = state.get("run_id") + if ci_result is None and isinstance(run_id, str) and run_id: + # Suspend + checkpoint; the CI-watcher's resume payload is the signal + # that the run terminated. We do not trust the payload's contents — + # we RE-FETCH the authenticated conclusion below so the gate reads a + # patch-independent, freshly-fetched result, never a resume-supplied + # one. + _await_ci(run_id) + ci_result = _fetch_ci_result(fetcher, state) + + # Merge the (possibly still-None) fetched CI result into the state the + # node reads from, WITHOUT mutating the caller's state object. The node + # reads ``ci_results``; a None result leaves the gate with nothing to pass + # on (it BLOCKs → park), so a never-terminal run fails closed. scoped_state: dict[str, Any] = dict(state) scoped_state["ci_results"] = ci_result return verifier_node(scoped_state, config) @@ -208,6 +254,40 @@ def make_verify_node( return node +def _fetch_ci_result( + fetcher: CiResultFetcher, state: PipelineState +) -> Mapping[str, Any] | None: + """Call the injected CI fetcher and normalise its result. + + Any value other than a mapping is treated as "no terminal result" (``None``), + so a malformed fetcher fails SAFE (the gate BLOCKs) rather than smuggling a + non-mapping past the gate. + """ + fetched = fetcher(state) + return fetched if isinstance(fetched, Mapping) else None + + +def _await_ci(run_id: str) -> None: + """Suspend the VERIFY node until the CI-watcher resumes it (§4 Decision 2). + + Mirrors the clarify human-gate node's durable suspend: calls + :func:`langgraph.types.interrupt` so the graph checkpoints and the daemon's + tick loop is freed while a multi-minute CI run is in flight. The + :func:`agent_team.ci_watcher.run_ci_watcher` sweep polls the run read-only and + drives the resume once it terminates. The interrupt payload carries only the + ``run_id`` being awaited (provenance for the watcher / operator logs); the + resume VALUE is intentionally ignored — the node re-fetches the authenticated + conclusion so the gate never reads a resume-supplied verdict. + + ``langgraph`` is imported lazily here to preserve this module's "no SDK at + module top" discipline (the topology stays importable where ``langgraph`` is + absent; the interrupt is only reached on the live, dispatched path). + """ + from langgraph.types import interrupt + + interrupt({"awaiting_ci": True, "run_id": run_id}) + + def route_after_verify(state: PipelineState) -> str: """LangGraph conditional-edge: the next route id after the VERIFY node. @@ -275,12 +355,16 @@ def bind_ci_result_fetcher( # A module-level note for the Integrate phase (no execution): the build->verify -# subgraph is hung off the review loop's "build" route. The conditional-edge map +# subgraph is hung off the review loop's "build" route. The linear stage order is +# BUILD -> [DISPATCH] -> VERIFY — DISPATCH (the optional CI-trigger vertex +# build_graph splices in when a dispatch node is wired) fires the CI run and +# captures ``state["run_id"]`` BEFORE VERIFY reads it (design §4 Decision 1); with +# no dispatch node the order is just BUILD -> VERIFY. The conditional-edge map # from VERIFY should send APPROVED_ROUTE to the PR/draft terminus, BUILD_ROUTE # back to the BUILD node (the bounded build<->verify loop, capped by # VerifierConfig.max_build_loops), and PARKED_ROUTE to the escalation terminus. # build_graph wires this in opt-in; this module never assembles it itself. _INTEGRATE_NOTE = ( - "review('build') -> BUILD -> VERIFY -> route_after_verify -> " + "review('build') -> BUILD -> [DISPATCH] -> VERIFY -> route_after_verify -> " "{approved: PR terminus, build: BUILD (loop), parked: escalation}" ) diff --git a/agent-team/agent_team/nodes/verifier.py b/agent-team/agent_team/nodes/verifier.py index 6e59a75..069f5c2 100644 --- a/agent-team/agent_team/nodes/verifier.py +++ b/agent-team/agent_team/nodes/verifier.py @@ -65,15 +65,19 @@ DEFAULT_MAX_BUILD_LOOPS: int = 3 class VerifierConfig: """Per-invocation knobs for the verifier node (§3.3, §3.3.2). - ``expected_run_id`` keys the gate to the exact CI run the verifier - dispatched for this diff (a stale/substituted run id is a BLOCK). + ``expected_run_id`` is a STATIC fallback only. The gate binds to the run id + THIS task dispatched, read from ``state["run_id"]`` at node-run time (the + dispatcher persists it there); the config value is consulted only when state + carries no ``run_id`` (e.g. a unit harness that drives the node directly). A + per-task ``state["run_id"]`` therefore always wins over this constant, and a + ``None`` effective run id is a BLOCK (never a vacuous pass). ``allowed_scope`` is the task's declared-scope path prefixes for the denylist boundary. ``max_build_loops`` caps build<->verify retries before the task parks. ``build_loops`` is the loops already consumed for this task (the coordinator threads it through state). """ - expected_run_id: str + expected_run_id: str | None = None allowed_scope: list[str] | None = None max_build_loops: int = DEFAULT_MAX_BUILD_LOOPS build_loops: int = 0 @@ -148,9 +152,14 @@ def verifier_node( Decision flow (§3.3.2 boundary #4 + §3.3 autonomy bounds): - 1. Call :func:`agent_team.ci_gate.evaluate_ci_gate` with the candidate diff, - the ledger hash (``diff_hash``), the authenticated ``ci_results``, the - ``expected_run_id``, and the task's ``allowed_scope``. + 1. Resolve the PER-TASK expected run id from ``state["run_id"]`` (the id the + dispatch node captured for THIS task), falling back to + ``config.expected_run_id`` only when state carries none. Call + :func:`agent_team.ci_gate.evaluate_ci_gate` with the candidate diff, the + ledger hash (``diff_hash``), the authenticated ``ci_results``, that + per-task expected run id, and the task's ``allowed_scope``. A ``None`` + effective run id BLOCKs (anti-substitution: the verdict has nothing to + bind to), never a vacuous pass. 2. PASS -> advance to DONE (draft PR). The advisor is NOT consulted. 3. FAIL -> consult the fix-advisor for a hint, then loop back to BUILD — unless ``build_loops`` has reached ``max_build_loops``, in which case @@ -162,6 +171,18 @@ def verifier_node( ledger_hash = state.get("diff_hash") ci_results = state.get("ci_results") + # Per-task binding: the gate must compare CI's run_id against the id THIS + # task dispatched (persisted by the dispatch node as ``state["run_id"]``), + # not a static wiring-time constant. State wins; the config value is only a + # fallback for harnesses that drive the node without a per-task run_id. A + # blank/None effective run id is left as None so the gate BLOCKs. + state_run_id = state.get("run_id") + expected_run_id = ( + state_run_id + if isinstance(state_run_id, str) and state_run_id + else config.expected_run_id + ) + if not isinstance(candidate_diff, str): # No diff to verify is itself a refuse-to-proceed: park for a human # rather than declaring anything. (A builder must have produced a diff @@ -169,7 +190,7 @@ def verifier_node( gate_result = GateResult( decision=GateDecision.BLOCK, reasons=["no candidate_diff present in state to verify"], - run_id=config.expected_run_id, + run_id=expected_run_id, diff_hash=ledger_hash, ci_conclusion=None, ) @@ -178,7 +199,7 @@ def verifier_node( candidate_diff=candidate_diff, ledger_hash=ledger_hash, ci_result=ci_results, - expected_run_id=config.expected_run_id, + expected_run_id=expected_run_id, allowed_scope=config.allowed_scope, ) diff --git a/agent-team/run-team.py b/agent-team/run-team.py index 58326f3..01d20c7 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -587,6 +587,7 @@ def _build_coordinator(args: argparse.Namespace) -> Any: default_clarify_node_factory, default_plan_node_factory, default_review_wiring, + failsafe_production_p3_wiring, ) transport = _build_transport(args) @@ -604,6 +605,21 @@ def _build_coordinator(args: argparse.Namespace) -> Any: # silent (notify None) so import + ledger commands need no token. notify, alarm_hook = _build_notifiers(args) + # P3 fail-safe serve default (design Decision 5): the bound build→verify + + # dispatch wiring is now the production ``serve`` default, but it must NEVER + # crash-loop serve-start. ``failsafe_production_p3_wiring`` resolves the env + # ONCE: configured -> the live P3 pair; unconfigured -> (None, None) inert (a + # task reaching P3 parks), with one WARNING + one #agent-team inert notice + # via the lifecycle ``notify`` sink. Scoped to ``serve`` (the daemon): the + # one-shot ``start`` / ``intake-*`` paths never auto-bind P3 — they run to the + # first human gate and exit, well short of BUILD/VERIFY. + build_verify_wiring = None + dispatch_node_wiring = None + if getattr(args, "command", None) == "serve": + build_verify_wiring, dispatch_node_wiring = failsafe_production_p3_wiring( + notify=notify + ) + # Production runs the full P2 graph: the wrapped real planner + the bound # GPT-4.1 review loop (Plane-2 depth-first). These factories are lazy and # only build/bind the model seams when a task actually runs. @@ -617,6 +633,8 @@ def _build_coordinator(args: argparse.Namespace) -> Any: context_provider=context_provider ), review_wiring=default_review_wiring, + build_verify_wiring=build_verify_wiring, + dispatch_node_wiring=dispatch_node_wiring, notify=notify, alarm_hook=alarm_hook, ) diff --git a/agent-team/tests/test_build_verify_subgraph.py b/agent-team/tests/test_build_verify_subgraph.py index 30a5976..815dee6 100644 --- a/agent-team/tests/test_build_verify_subgraph.py +++ b/agent-team/tests/test_build_verify_subgraph.py @@ -102,6 +102,19 @@ def test_route_ids_mirror_graph_by_value() -> None: assert APPROVED_ROUTE == "approved" +def test_integrate_note_documents_build_dispatch_verify_order() -> None: + """The topology note records BUILD -> [DISPATCH] -> VERIFY (design §4 D1). + + DISPATCH must precede VERIFY so it captures ``state["run_id"]`` before VERIFY + reads it. Pin the documented order so a future reorder back to the old + BUILD -> VERIFY -> DISPATCH topology is caught. + """ + note = bvs._INTEGRATE_NOTE + assert "BUILD -> [DISPATCH] -> VERIFY" in note + # DISPATCH appears before VERIFY in the recorded stage order. + assert note.index("DISPATCH") < note.index("VERIFY") + + # --------------------------------------------------------------------------- # # BUILD node: proposes a diff via an injected fake builder # --------------------------------------------------------------------------- # @@ -373,3 +386,146 @@ def test_verify_node_does_not_mutate_caller_state() -> None: node(state) # Caller's state is untouched (the node wrote into a dict copy). assert state["ci_results"] is sentinel + + +# --------------------------------------------------------------------------- # +# Per-task run-id binding through the subgraph wrapper (design §4 Decision 4) +# --------------------------------------------------------------------------- # + + +def test_verify_node_binds_to_per_task_state_run_id() -> None: + """make_verify_node gates against state["run_id"], not a static config id. + + A single VerifierConfig with no expected_run_id is shared; the per-task + state["run_id"] supplies the binding the fetched CI conclusion must match. + """ + diff = _diff_for("src/foo.py") + + def fetcher_keyed_to_task(s): + # The (real) fetcher keys its conclusion to the task's dispatched run id. + return { + "run_id": s["run_id"], + "conclusion": "success", + "diff_hash": _hash(diff), + } + + node = make_verify_node( + VerifierConfig(expected_run_id=None, allowed_scope=["src"]), + ci_result_fetcher=fetcher_keyed_to_task, + ) + + state = _verify_state(diff) + state["run_id"] = "dispatched-123" + out = node(state) + assert out["status"] == TaskStatus.DONE.value + assert out["review_verdicts"][0]["run_id"] == "dispatched-123" + + +def test_verify_node_blocks_substituted_run_id_through_wrapper() -> None: + """A fetched CI result keyed to a DIFFERENT run id than state["run_id"] blocks.""" + diff = _diff_for("src/foo.py") + + def substituting_fetcher(s): + # CI result grafted from another task (run id mismatch) with success. + return { + "run_id": "someone-elses-run", + "conclusion": "success", + "diff_hash": _hash(diff), + } + + node = make_verify_node( + VerifierConfig(expected_run_id=None, allowed_scope=["src"]), + ci_result_fetcher=substituting_fetcher, + ) + + state = _verify_state(diff) + state["run_id"] = "my-run" + out = node(state) + assert out["status"] == TaskStatus.PARKED.value + assert out["ci_results"]["gate_decision"] == "block" + + +def test_verify_node_blocks_when_no_run_id_anywhere() -> None: + """No state run_id and no config fallback -> BLOCK/park, never a vacuous pass.""" + diff = _diff_for("src/foo.py") + + def pass_fetcher(s): + return {"run_id": "r1", "conclusion": "success", "diff_hash": _hash(diff)} + + node = make_verify_node( + VerifierConfig(expected_run_id=None, allowed_scope=["src"]), + ci_result_fetcher=pass_fetcher, + ) + out = node(_verify_state(diff)) # no state["run_id"] + assert out["status"] == TaskStatus.PARKED.value + assert out["ci_results"]["gate_decision"] == "block" + + +# --------------------------------------------------------------------------- # +# Async CI-wait: VERIFY suspends via interrupt() on an in-progress run (§4 Dec 2) +# --------------------------------------------------------------------------- # + + +def test_verify_suspends_on_in_progress_run_then_resumes_to_gate() -> None: + """With a dispatched run_id and a not-yet-terminal fetch, VERIFY interrupts; + once the CI-watcher resumes it, the re-fetched terminal result is gated.""" + from langgraph.checkpoint.memory import MemorySaver + from langgraph.graph import END, START, StateGraph + from langgraph.types import Command + + diff = _diff_for("src/foo.py") + + # The fetcher is "in progress" on the first call, then terminal on the second + # (mirroring a real run that concluded while the task was suspended). + calls = {"n": 0} + + def progressing_fetcher(state): + calls["n"] += 1 + if calls["n"] == 1: + return None # still running -> VERIFY must suspend + return { + "run_id": state["run_id"], + "conclusion": "success", + "diff_hash": _hash(diff), + } + + node = make_verify_node( + VerifierConfig(expected_run_id=None, allowed_scope=["src"]), + ci_result_fetcher=progressing_fetcher, + ) + + builder = StateGraph(PipelineState) + builder.add_node("verify", node) + builder.add_edge(START, "verify") + builder.add_edge("verify", END) + graph = builder.compile(checkpointer=MemorySaver()) + + cfg = {"configurable": {"thread_id": "tw1"}} + state = _verify_state(diff) + state["run_id"] = "dispatched-77" + + first = graph.invoke(state, cfg) + # The task suspended at VERIFY's interrupt (awaiting CI) rather than gating. + assert "__interrupt__" in first + assert calls["n"] == 1 + + # The CI-watcher resumes the task once the run terminated. + out = graph.invoke(Command(resume={"awaiting_ci": "done"}), cfg) + # On resume the node re-fetched the now-terminal result and the gate PASSed. + assert calls["n"] == 2 + assert out["status"] == TaskStatus.DONE.value + assert out["current_phase"] == Phase.DONE.value + + +def test_verify_no_run_id_does_not_suspend_and_parks() -> None: + """The INERT/no-run path (no state run_id) never suspends: a None fetch flows + straight to the gate, which BLOCKs and parks (existing behavior preserved).""" + diff = _diff_for("src/foo.py") + + node = make_verify_node( + VerifierConfig(expected_run_id="r1", allowed_scope=["src"]), + ci_result_fetcher=lambda s: None, + ) + out = node(_verify_state(diff)) # no state["run_id"] -> no interrupt + assert out["current_phase"] == Phase.PARKED.value + assert route_after_verify(out) == PARKED_ROUTE diff --git a/agent-team/tests/test_ci_gate.py b/agent-team/tests/test_ci_gate.py index f6f0ec9..47c9c77 100644 --- a/agent-team/tests/test_ci_gate.py +++ b/agent-team/tests/test_ci_gate.py @@ -389,14 +389,46 @@ def test_gate_scope_violation_blocks() -> None: assert any("scope" in r for r in result.reasons) -def test_gate_empty_run_id_raises() -> None: +def test_gate_empty_run_id_blocks_never_passes() -> None: + # Per design §4 (per-task binding): an empty expected_run_id means the + # dispatcher captured no run id for this task. There is nothing to bind the + # verdict to, so the gate BLOCKs (fail-closed) rather than vacuously gating + # against "" — even when a substituted ci_result reports a matching success. + diff = _diff_for("src/foo.py") + result = evaluate_ci_gate( + candidate_diff=diff, + ledger_hash=_ledger_hash(diff), + ci_result=_good_ci("", diff), + expected_run_id="", + ) + assert result.decision is GateDecision.BLOCK + assert result.run_id is None + assert any("bind" in r for r in result.reasons) + + +def test_gate_none_run_id_blocks_never_passes() -> None: + # A None expected_run_id (the legitimate "dispatch unresolved" runtime state) + # is a BLOCK, not an exception — and never a pass even on a success ci_result. + diff = _diff_for("src/foo.py") + result = evaluate_ci_gate( + candidate_diff=diff, + ledger_hash=_ledger_hash(diff), + ci_result=_good_ci("run-1", diff, "success"), + expected_run_id=None, + ) + assert result.decision is GateDecision.BLOCK + assert any("bind" in r for r in result.reasons) + + +def test_gate_non_string_run_id_raises() -> None: + # A non-string, non-None expected_run_id is structurally invalid -> raise. diff = _diff_for("src/foo.py") with pytest.raises(CiGateError): evaluate_ci_gate( candidate_diff=diff, ledger_hash=_ledger_hash(diff), ci_result=_good_ci("run-1", diff), - expected_run_id="", + expected_run_id=123, # type: ignore[arg-type] ) diff --git a/agent-team/tests/test_ci_watcher.py b/agent-team/tests/test_ci_watcher.py new file mode 100644 index 0000000..cec2c0d --- /dev/null +++ b/agent-team/tests/test_ci_watcher.py @@ -0,0 +1,323 @@ +"""Unit tests for agent_team.ci_watcher (design §3.3.2, P3 Decision 2). + +Covers the async resume-on-CI-complete sweep: it RESUMES a task on a terminal +conclusion, PARKS on the dispatch timeout, and FAILS CLOSED (parks) on a fetch +error / unusable state. No network: every side effect (poll, resume, park) is an +injected callable, exactly as the module's contract promises. +""" + +from __future__ import annotations + +from datetime import datetime, timedelta, timezone + +from agent_team.ci_watcher import ( + CiPollResult, + CiWatchAction, + PendingCiTask, + default_ci_poller, + run_ci_watcher, +) + +# A fixed "now" and a dispatched-at watermark; ISO strings mirror the ledger. +_NOW = datetime(2026, 6, 23, 12, 0, 0, tzinfo=timezone.utc) +_JUST_NOW = (_NOW - timedelta(minutes=1)).isoformat() +_LONG_AGO = (_NOW - timedelta(hours=2)).isoformat() + + +class _Recorder: + """Records the tasks/results handed to an injected side-effect callback.""" + + def __init__(self) -> None: + self.resumed: list[tuple[PendingCiTask, dict]] = [] + self.parked: list[PendingCiTask] = [] + + def on_resume(self, task: PendingCiTask, result: object) -> None: + self.resumed.append((task, dict(result))) # type: ignore[arg-type] + + def on_park(self, task: PendingCiTask) -> None: + self.parked.append(task) + + +def _task(*, run_id: str | None = "12345", dispatched_at: str | None = _JUST_NOW): + return PendingCiTask(thread_id="t1", run_id=run_id, dispatched_at=dispatched_at) + + +# --------------------------------------------------------------------------- # +# TERMINAL -> resume +# --------------------------------------------------------------------------- # + + +def test_terminal_conclusion_resumes_the_task() -> None: + rec = _Recorder() + terminal = {"run_id": "12345", "conclusion": "success", "diff_hash": "abc"} + + report = run_ci_watcher( + [_task()], + poll=lambda t: CiPollResult.terminal(terminal), + on_resume=rec.on_resume, + on_park=rec.on_park, + now=_NOW, + ) + + assert report.resumed == 1 + assert report.parked == 0 + assert rec.parked == [] + assert len(rec.resumed) == 1 + resumed_task, resumed_result = rec.resumed[0] + assert resumed_task.thread_id == "t1" + # The authenticated terminal result is handed to the resume side effect. + assert resumed_result == terminal + assert report.outcomes[0].action is CiWatchAction.RESUMED + + +def test_terminal_resume_even_past_timeout_prefers_resume() -> None: + """A run that terminated should RESUME, not park, even if dispatched long ago.""" + rec = _Recorder() + report = run_ci_watcher( + [_task(dispatched_at=_LONG_AGO)], + poll=lambda t: CiPollResult.terminal({"conclusion": "failure"}), + on_resume=rec.on_resume, + on_park=rec.on_park, + timeout=timedelta(minutes=30), + now=_NOW, + ) + assert report.resumed == 1 + assert rec.parked == [] + + +# --------------------------------------------------------------------------- # +# PENDING + timeout -> park ; PENDING within window -> wait +# --------------------------------------------------------------------------- # + + +def test_pending_within_window_waits() -> None: + rec = _Recorder() + report = run_ci_watcher( + [_task(dispatched_at=_JUST_NOW)], + poll=lambda t: CiPollResult.pending(), + on_resume=rec.on_resume, + on_park=rec.on_park, + timeout=timedelta(minutes=30), + now=_NOW, + ) + assert report.waiting == 1 + assert report.parked == 0 + assert rec.parked == [] + assert rec.resumed == [] + assert report.outcomes[0].action is CiWatchAction.WAITING + + +def test_pending_past_timeout_parks() -> None: + rec = _Recorder() + report = run_ci_watcher( + [_task(dispatched_at=_LONG_AGO)], + poll=lambda t: CiPollResult.pending(), + on_resume=rec.on_resume, + on_park=rec.on_park, + timeout=timedelta(minutes=30), + now=_NOW, + ) + assert report.parked_timeout == 1 + assert report.parked == 1 + assert len(rec.parked) == 1 + assert rec.resumed == [] + assert report.outcomes[0].action is CiWatchAction.PARKED_TIMEOUT + + +# --------------------------------------------------------------------------- # +# ERROR / None result / unusable state -> fail closed (park) +# --------------------------------------------------------------------------- # + + +def test_poll_error_parks_fail_closed() -> None: + rec = _Recorder() + report = run_ci_watcher( + [_task(dispatched_at=_JUST_NOW)], # within window: error still parks + poll=lambda t: CiPollResult.error(), + on_resume=rec.on_resume, + on_park=rec.on_park, + now=_NOW, + ) + assert report.parked_error == 1 + assert len(rec.parked) == 1 + assert rec.resumed == [] + assert report.outcomes[0].action is CiWatchAction.PARKED_ERROR + + +def test_raising_poll_is_isolated_and_parks() -> None: + rec = _Recorder() + + def boom(task: PendingCiTask) -> CiPollResult: + raise RuntimeError("github exploded") + + report = run_ci_watcher( + [_task()], + poll=boom, + on_resume=rec.on_resume, + on_park=rec.on_park, + now=_NOW, + ) + assert report.parked_error == 1 + assert len(rec.parked) == 1 + assert "RuntimeError" in (report.outcomes[0].error or "") + + +def test_missing_run_id_parks_without_polling() -> None: + rec = _Recorder() + polled: list[PendingCiTask] = [] + + def tracking_poll(task: PendingCiTask) -> CiPollResult: + polled.append(task) + return CiPollResult.pending() + + report = run_ci_watcher( + [_task(run_id=None)], + poll=tracking_poll, + on_resume=rec.on_resume, + on_park=rec.on_park, + now=_NOW, + ) + assert report.parked_error == 1 + assert polled == [] # unusable state parks BEFORE any poll + assert len(rec.parked) == 1 + + +def test_unparseable_dispatched_at_parks_without_polling() -> None: + rec = _Recorder() + polled: list[PendingCiTask] = [] + + report = run_ci_watcher( + [_task(dispatched_at="not-a-timestamp")], + poll=lambda t: (polled.append(t), CiPollResult.pending())[1], + on_resume=rec.on_resume, + on_park=rec.on_park, + now=_NOW, + ) + assert report.parked_error == 1 + assert polled == [] + assert len(rec.parked) == 1 + + +# --------------------------------------------------------------------------- # +# Side-effect isolation across the batch +# --------------------------------------------------------------------------- # + + +def test_one_task_failure_does_not_abort_the_sweep() -> None: + """A resume callback that raises for one task does not stop the others.""" + rec = _Recorder() + raised_for: list[str] = [] + + def flaky_resume(task: PendingCiTask, result: object) -> None: + if task.thread_id == "bad": + raised_for.append(task.thread_id) + raise RuntimeError("resume failed") + rec.resumed.append((task, dict(result))) # type: ignore[arg-type] + + tasks = [ + PendingCiTask(thread_id="bad", run_id="1", dispatched_at=_JUST_NOW), + PendingCiTask(thread_id="good", run_id="2", dispatched_at=_JUST_NOW), + ] + report = run_ci_watcher( + tasks, + poll=lambda t: CiPollResult.terminal({"conclusion": "success"}), + on_resume=flaky_resume, + on_park=rec.on_park, + now=_NOW, + ) + assert report.examined == 2 + # The bad task is recorded as a PARKED_ERROR; the good one resumed. + actions = {o.thread_id: o.action for o in report.outcomes} + assert actions["bad"] is CiWatchAction.PARKED_ERROR + assert actions["good"] is CiWatchAction.RESUMED + assert raised_for == ["bad"] + + +def test_park_callback_failure_downgrades_to_error_outcome() -> None: + def park_boom(task: PendingCiTask) -> None: + raise RuntimeError("park write failed") + + report = run_ci_watcher( + [_task(dispatched_at=_LONG_AGO)], + poll=lambda t: CiPollResult.pending(), + on_resume=lambda t, r: None, + on_park=park_boom, + timeout=timedelta(minutes=30), + now=_NOW, + ) + # The intended action was a timeout-park, but the park callback raised, so the + # outcome is recorded as PARKED_ERROR (the sweep continues either way). + assert report.outcomes[0].action is CiWatchAction.PARKED_ERROR + assert "RuntimeError" in (report.outcomes[0].error or "") + + +# --------------------------------------------------------------------------- # +# default_ci_poller — reuses ci_fetcher read-only, classifies the result +# --------------------------------------------------------------------------- # + + +class _FakeResponse: + def __init__(self, status: int, body: dict) -> None: + self.status_code = status + self._body = body + + def json(self) -> dict: + return self._body + + +class _FakeClient: + """A read-only ``requests``-like client: records GETs, never writes.""" + + def __init__(self, response: _FakeResponse) -> None: + self._response = response + self.gets: list[str] = [] + + def get(self, url: str, *, timeout: float) -> _FakeResponse: + self.gets.append(url) + return self._response + + +def test_default_poller_terminal_on_success_conclusion() -> None: + client = _FakeClient(_FakeResponse(200, {"id": 12345, "conclusion": "success"})) + poll = default_ci_poller(owner="o", repo="r", client=client) + result = poll(_task(run_id="12345")) + assert result.outcome.value == "terminal" + assert result.result is not None + assert result.result["conclusion"] == "success" + # Exactly one read-only GET issued. + assert len(client.gets) == 1 + + +def test_default_poller_pending_on_in_progress_run() -> None: + # An in-progress run has conclusion=None -> fetch_ci_result returns None. + client = _FakeClient(_FakeResponse(200, {"id": 12345, "conclusion": None})) + poll = default_ci_poller(owner="o", repo="r", client=client) + result = poll(_task(run_id="12345")) + assert result.outcome.value == "pending" + assert result.result is None + + +def test_default_poller_pending_on_http_error_status() -> None: + # A 404/5xx makes fetch_ci_result return None; the poller classifies it as + # pending (the timeout branch in the sweep is the fail-closed backstop, and a + # never-resolving run parks on timeout). + client = _FakeClient(_FakeResponse(404, {})) + poll = default_ci_poller(owner="o", repo="r", client=client) + result = poll(_task(run_id="12345")) + assert result.outcome.value == "pending" + + +def test_default_poller_error_when_fetch_raises() -> None: + class _BoomClient: + def get(self, url: str, *, timeout: float): + raise RuntimeError("network down") + + # fetch_ci_result itself swallows GET errors to None, so the poller sees + # pending; but a defensive wrapper still classifies a raised fetch as error. + # Here we drive the error path by passing a client whose .get raises AND + # bypass fetch_ci_result's own swallow via a poller that re-raises is not + # possible; instead assert the read-only GET error degrades to pending (the + # timeout backstop parks it). This documents the boundary. + poll = default_ci_poller(owner="o", repo="r", client=_BoomClient()) + result = poll(_task(run_id="12345")) + assert result.outcome.value == "pending" diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index 15afb11..476f6e6 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -1318,3 +1318,117 @@ def test_parked_message_infers_phase_when_current_phase_is_parked( coord._post_resume_followups([_resume_result("abcd1234ef00")]) assert "Reached phase: review" in msgs[0] assert "Reached phase: parked" not in msgs[0] + + +# --------------------------------------------------------------------------- # +# CI-watcher sweep in tick (§3.3.2 Decision 2) — async resume-on-CI-complete +# --------------------------------------------------------------------------- # + + +def test_tick_ci_watch_noop_without_seams(db_path: Path) -> None: + """With no CI-watcher seams wired, tick() runs no CI sweep (the default).""" + coord = _make_coordinator(db_path) + coord.setup() + # _ci_watch returns None (skipped) when the seams are absent. + assert coord._ci_watch() is None + # tick still works as before. + assert coord.tick() == [] + + +def _ci_coordinator(db_path: Path, *, pending, poller): + """A Coordinator wired with CI-watcher seams + an in-memory saver.""" + saver = _Saver() + return Coordinator( + db_path=db_path, + transport=FakeTransport(), + build_clarify_node=lambda: graph_mod.clarify_node, + build_checkpointer=lambda _path: saver, + ci_pending_provider=lambda: pending, + ci_poller=poller, + ) + + +def test_tick_ci_watch_resumes_on_terminal_conclusion(db_path: Path) -> None: + """A terminated run drives a graph resume of the suspended VERIFY task.""" + from agent_team.ci_watcher import CiPollResult, PendingCiTask + + task = PendingCiTask( + thread_id="t-ci", run_id="999", dispatched_at="2026-06-23T11:00:00+00:00" + ) + coord = _ci_coordinator( + db_path, + pending=[task], + poller=lambda t: CiPollResult.terminal( + {"run_id": "999", "conclusion": "success", "diff_hash": "h"} + ), + ) + coord.setup() + + invoked: list[tuple[Any, Any]] = [] + coord._graph.invoke = lambda inp, cfg: invoked.append((inp, cfg)) # type: ignore[assignment] + + report = coord._ci_watch() + assert report is not None + assert report.resumed == 1 + # The graph was driven forward for the suspended task's thread. + assert len(invoked) == 1 + _inp, cfg = invoked[0] + assert cfg["configurable"]["thread_id"] == "t-ci" + + +def test_tick_ci_watch_parks_on_timeout(db_path: Path) -> None: + """A run that never terminates within the timeout parks the task + ALARMs.""" + from agent_team.ci_watcher import CiPollResult, PendingCiTask + + alarms: list[str] = [] + task = PendingCiTask( + thread_id="t-slow", run_id="42", dispatched_at="2000-01-01T00:00:00+00:00" + ) + saver = _Saver() + coord = Coordinator( + db_path=db_path, + transport=FakeTransport(), + build_clarify_node=lambda: graph_mod.clarify_node, + build_checkpointer=lambda _path: saver, + ci_pending_provider=lambda: [task], + ci_poller=lambda t: CiPollResult.pending(), + alarm_hook=alarms.append, + ) + coord.setup() + + updated: list[tuple[Any, Any]] = [] + coord._graph.update_state = lambda cfg, delta: updated.append((cfg, delta)) # type: ignore[assignment] + + report = coord._ci_watch() + assert report is not None + assert report.parked_timeout == 1 + # The task was durably parked and an ALARM was raised. + assert updated and updated[0][1]["status"] == "parked" + assert alarms == ["t-slow"] + + +def test_tick_ci_watch_parks_fail_closed_on_error(db_path: Path) -> None: + """A poll error parks the task (fail-closed), even within the timeout window.""" + from agent_team.ci_watcher import CiPollResult, PendingCiTask + + alarms: list[str] = [] + task = PendingCiTask( + thread_id="t-err", run_id="7", dispatched_at="2026-06-23T11:59:00+00:00" + ) + saver = _Saver() + coord = Coordinator( + db_path=db_path, + transport=FakeTransport(), + build_clarify_node=lambda: graph_mod.clarify_node, + build_checkpointer=lambda _path: saver, + ci_pending_provider=lambda: [task], + ci_poller=lambda t: CiPollResult.error(), + alarm_hook=alarms.append, + ) + coord.setup() + coord._graph.update_state = lambda cfg, delta: None # type: ignore[assignment] + + report = coord._ci_watch() + assert report is not None + assert report.parked_error == 1 + assert alarms == ["t-err"] diff --git a/agent-team/tests/test_graph.py b/agent-team/tests/test_graph.py index b4b58be..afb8a15 100644 --- a/agent-team/tests/test_graph.py +++ b/agent-team/tests/test_graph.py @@ -616,6 +616,128 @@ def test_p3_graph_inert_default_parks_at_verify(restore_review_invoker) -> None: assert final["status"] == TaskStatus.PARKED.value +def _p3_graph_with_dispatch(review_text: str, *, ci_result_fetcher, dispatch_node): + """Compile a P3+ graph: review -> build -> DISPATCH -> verify. + + Same as ``_p3_graph`` but splices a ``dispatch_node`` between BUILD and + VERIFY so the reordered topology (design §4 Decision 1) can be driven end to + end — DISPATCH writes ``state["run_id"]`` before VERIFY reads it. + """ + from agent_team.nodes import review_loop + from agent_team.nodes.build_verify_subgraph import ( + make_build_node, + make_verify_node, + route_after_verify, + ) + from agent_team.nodes.verifier import VerifierConfig + + review_loop.set_review_invoker(lambda prompt, **kw: review_text) + + def fake_builder(*, plan, config): + return _p3_diff() + + build_node = make_build_node(diff_builder=fake_builder) + # No static expected_run_id: the gate must bind to the run id DISPATCH wrote. + verify_node = make_verify_node( + VerifierConfig(expected_run_id=None, allowed_scope=["src"]), + ci_result_fetcher=ci_result_fetcher, + ) + return build_graph( + checkpointer=_Saver(), + live_plan_node=_p3_plan_stub, + review_node=review_loop.bind_review_node(), + route_review=review_loop.route_after_review, + build_verify=(build_node, verify_node, route_after_verify), + dispatch_node=dispatch_node, + ) + + +def test_build_graph_dispatch_node_requires_build_verify() -> None: + """dispatch_node without build_verify is a wiring error (nothing to splice).""" + with pytest.raises(ValueError, match="dispatch_node requires build_verify"): + build_graph(dispatch_node=lambda state: {}) + + +def test_p3_dispatch_runs_before_verify_and_supplies_run_id( + restore_review_invoker, +) -> None: + """BUILD -> DISPATCH -> VERIFY: DISPATCH captures run_id BEFORE VERIFY reads it. + + The verify node is wired with NO static expected_run_id, so the only way the + authenticated-pass gate can bind a verdict is if DISPATCH wrote + ``state["run_id"]`` first. The fetcher keys its conclusion to the dispatched + run id and asserts it observes that id — proving DISPATCH ran before VERIFY. + """ + from agent_team.state_store import compute_content_hash + + observed: dict[str, object] = {} + dispatched_run_id = "r-dispatched-007" + + def dispatch_node(state): + # Mirror dispatch_invoker's contract: persist the located run identity so + # the downstream verifier binds the gate to THIS task's dispatched run. + return {"run_id": dispatched_run_id, "dispatched_at": "2026-06-23T00:00:00Z"} + + def pass_fetcher(state): + # The fetcher only sees state["run_id"] if DISPATCH already ran. + observed["run_id"] = state.get("run_id") + diff_hash = compute_content_hash(_p3_diff().encode("utf-8")) + return { + "run_id": state.get("run_id"), + "conclusion": "success", + "diff_hash": diff_hash, + } + + graph = _p3_graph_with_dispatch( + "VERDICT: APPROVE\nlooks solid", + ci_result_fetcher=pass_fetcher, + dispatch_node=dispatch_node, + ) + thread_id, _ = start_task(graph, transport="slack") + final = resume_task(graph, thread_id=thread_id, answer="scope is X") + + # VERIFY observed the run id DISPATCH wrote -> DISPATCH ran first. + assert observed["run_id"] == dispatched_run_id + # And the per-task-bound authenticated pass cleared the gate -> DONE terminus. + assert final["current_phase"] == Phase.DONE.value + assert final["status"] == TaskStatus.DONE.value + assert final["run_id"] == dispatched_run_id + + +def test_p3_dispatch_node_present_in_graph_topology(restore_review_invoker) -> None: + """The DISPATCH vertex is wired between BUILD and VERIFY when injected.""" + from agent_team.graph import BUILD_NODE, DISPATCH_NODE, VERIFY_NODE + + graph = _p3_graph_with_dispatch( + "VERDICT: APPROVE\nlooks solid", + ci_result_fetcher=lambda state: None, + dispatch_node=lambda state: {"run_id": "r"}, + ) + g = graph.get_graph() + nodes = set(g.nodes) + assert {BUILD_NODE, DISPATCH_NODE, VERIFY_NODE} <= nodes + + # The linear order is BUILD -> DISPATCH -> VERIFY (no direct BUILD -> VERIFY). + edges = {(e.source, e.target) for e in g.edges} + assert (BUILD_NODE, DISPATCH_NODE) in edges + assert (DISPATCH_NODE, VERIFY_NODE) in edges + assert (BUILD_NODE, VERIFY_NODE) not in edges + + +def test_p3_no_dispatch_falls_back_to_build_then_verify(restore_review_invoker) -> None: + """Without a dispatch node the order falls back to BUILD -> VERIFY directly.""" + from agent_team.graph import BUILD_NODE, DISPATCH_NODE, VERIFY_NODE + + graph = _p3_graph("VERDICT: APPROVE\nlooks solid", ci_result_fetcher=lambda s: None) + g = graph.get_graph() + nodes = set(g.nodes) + assert {BUILD_NODE, VERIFY_NODE} <= nodes + assert DISPATCH_NODE not in nodes + + edges = {(e.source, e.target) for e in g.edges} + assert (BUILD_NODE, VERIFY_NODE) in edges + + def test_p3_graph_route_constants_mirror_subgraph_by_value() -> None: """graph.py's P3 route ids match the subgraph module by value (no cycle).""" from agent_team.nodes import build_verify_subgraph as bvs diff --git a/agent-team/tests/test_p3_failsafe_serve_default.py b/agent-team/tests/test_p3_failsafe_serve_default.py new file mode 100644 index 0000000..8664f8f --- /dev/null +++ b/agent-team/tests/test_p3_failsafe_serve_default.py @@ -0,0 +1,252 @@ +"""P3 fail-safe serve default (design Decision 5; UNIT 0e). + +The bound P3 build→verify + dispatch wiring is the new production ``serve`` +default, but its factories are called EAGERLY at graph-build and the live +dispatch factory RAISES when ``AGENT_TEAM_REPO_OWNER`` / ``AGENT_TEAM_REPO_NAME`` +are unset. These tests pin the fail-safe contract: + +* :func:`agent_team.coordinator._p3_env_is_configured` truth table. +* :func:`agent_team.coordinator.failsafe_production_p3_wiring` — + env-unset degrades to the INERT ``(None, None)`` pair with exactly ONE WARNING + and ONE ``#agent-team`` inert notice via the lifecycle ``notify`` sink (never + the ALARM path), and NEVER raises; env-set returns the live wiring pair. +* a Coordinator built with the env-unset (inert) result sets up cleanly and an + approved task settles at the P2 BUILD terminus (no P3 nodes, no dispatch, no + exception) — i.e. a task that would reach P3 parks short of build/dispatch. +* ``run-team.py serve`` binds the fail-safe pair; ``start`` / ``intake`` do not. +""" + +from __future__ import annotations + +import logging +from pathlib import Path +from typing import Any + +import pytest + +from agent_team import graph as graph_mod +from agent_team.coordinator import ( + Coordinator, + _p3_env_is_configured, + default_dispatch_node_factory, + failsafe_production_p3_wiring, +) + +try: # InMemorySaver is the modern name; fall back on older langgraph. + from langgraph.checkpoint.memory import InMemorySaver as _Saver +except ImportError: # pragma: no cover - older langgraph + from langgraph.checkpoint.memory import MemorySaver as _Saver + + +_P3_ENV = ("AGENT_TEAM_REPO_OWNER", "AGENT_TEAM_REPO_NAME") +_TOKEN_ENV = ("AGENT_TEAM_CI_READ_TOKEN", "GITHUB_TOKEN") + + +@pytest.fixture +def clean_p3_env(monkeypatch: pytest.MonkeyPatch) -> None: + """Start each test from a fully unset P3 environment.""" + for name in (*_P3_ENV, *_TOKEN_ENV): + monkeypatch.delenv(name, raising=False) + + +# --------------------------------------------------------------------------- # +# _p3_env_is_configured truth table +# --------------------------------------------------------------------------- # + + +def test_env_unset_is_not_configured(clean_p3_env: None) -> None: + assert _p3_env_is_configured() is False + + +def test_owner_repo_without_token_is_not_configured( + clean_p3_env: None, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "org") + monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "repo") + # No CI-read token -> the verifier could never read an authenticated pass. + assert _p3_env_is_configured() is False + + +def test_token_without_owner_repo_is_not_configured( + clean_p3_env: None, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test") + assert _p3_env_is_configured() is False + + +def test_owner_repo_and_ci_token_is_configured( + clean_p3_env: None, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "org") + monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "repo") + monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test") + assert _p3_env_is_configured() is True + + +def test_github_token_fallback_satisfies_ci_token( + clean_p3_env: None, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "org") + monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "repo") + monkeypatch.setenv("GITHUB_TOKEN", "ghp_fallback") + assert _p3_env_is_configured() is True + + +def test_blank_env_values_are_not_configured( + clean_p3_env: None, monkeypatch: pytest.MonkeyPatch +) -> None: + # Whitespace-only values must not count as configured (fail-closed). + monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", " ") + monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "repo") + monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test") + assert _p3_env_is_configured() is False + + +# --------------------------------------------------------------------------- # +# failsafe_production_p3_wiring — env-unset degrade +# --------------------------------------------------------------------------- # + + +def test_env_unset_returns_inert_pair_warns_and_notifies( + clean_p3_env: None, caplog: pytest.LogCaptureFixture +) -> None: + notices: list[str] = [] + + with caplog.at_level(logging.WARNING, logger="agent_team.coordinator"): + build_verify, dispatch = failsafe_production_p3_wiring( + notify=lambda msg, **_: notices.append(msg) + ) + + # Inert P3: no build->verify subgraph, no dispatch. + assert build_verify is None + assert dispatch is None + + # Exactly one WARNING about the inert P3 wiring. + inert_warnings = [ + r + for r in caplog.records + if r.levelno == logging.WARNING and "INERT" in r.getMessage() + ] + assert len(inert_warnings) == 1 + + # Exactly one #agent-team inert notice via the (non-ALARM) notify sink. + assert len(notices) == 1 + assert "INERT" in notices[0] + + +def test_env_unset_without_notify_does_not_raise(clean_p3_env: None) -> None: + # A token-less / channel-less serve still comes up inert with no notify sink. + build_verify, dispatch = failsafe_production_p3_wiring(notify=None) + assert build_verify is None + assert dispatch is None + + +def test_inert_notify_failure_is_swallowed(clean_p3_env: None) -> None: + def _boom(_msg: str, **_kw: Any) -> None: + raise RuntimeError("slack down") + + # The notify sink raising must not propagate out of serve-start. + build_verify, dispatch = failsafe_production_p3_wiring(notify=_boom) + assert build_verify is None + assert dispatch is None + + +# --------------------------------------------------------------------------- # +# failsafe_production_p3_wiring — env-set binds live wiring +# --------------------------------------------------------------------------- # + + +def test_env_set_returns_live_wiring_pair( + clean_p3_env: None, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "test-org") + monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "test-repo") + monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test") + + notices: list[str] = [] + build_verify, dispatch = failsafe_production_p3_wiring( + notify=lambda msg, **_: notices.append(msg) + ) + + # Live pair: both factories present, no inert notice emitted. + assert callable(build_verify) + assert callable(dispatch) + assert notices == [] + + # The dispatch factory is the env-reading production factory. + assert dispatch is default_dispatch_node_factory + + # Each live factory builds without raising and yields the expected shapes. + bv_tuple = build_verify() + assert isinstance(bv_tuple, tuple) and len(bv_tuple) == 3 + assert all(callable(part) for part in bv_tuple) + assert callable(dispatch()) + + +# --------------------------------------------------------------------------- # +# Coordinator built with the inert result sets up + a P3-bound task does not +# dispatch / crash (it settles at the P2 BUILD terminus). +# --------------------------------------------------------------------------- # + + +def test_coordinator_with_inert_wiring_builds_and_task_parks_short_of_p3( + clean_p3_env: None, tmp_path: Path +) -> None: + """env-unset -> graph builds, coordinator constructs, an approved task settles + at the P2 BUILD terminus with no P3 nodes, no dispatch, and no exception.""" + from agent_team.graph import resume_task, start_task + from agent_team.nodes import review_loop + from agent_team.task_model import Phase, PipelineState, TaskStatus + + build_verify, dispatch = failsafe_production_p3_wiring(notify=None) + assert (build_verify, dispatch) == (None, None) + + saver = _Saver() + saved_invoker = review_loop._review_invoker + try: + review_loop.set_review_invoker(lambda prompt, **kw: "VERDICT: APPROVE\nok") + + def _plan_stub(state: PipelineState) -> PipelineState: + return PipelineState( + plan={"title": "t", "scope": ["src"], "phases": ["P1"]}, + current_phase=Phase.REVIEW.value, + status=TaskStatus.ACTIVE.value, + ) + + coord = Coordinator( + db_path=tmp_path / "inert.db", + transport=_FakeTransport(), + build_clarify_node=lambda: graph_mod.clarify_node, + build_plan_node=lambda: _plan_stub, + review_wiring=lambda: ( + review_loop.bind_review_node(), + review_loop.route_after_review, + ), + build_verify_wiring=build_verify, + dispatch_node_wiring=dispatch, + build_checkpointer=lambda _path: saver, + ) + coord.setup() + + # No P3 nodes were wired (inert): the graph stops at the P2 terminus. + nodes = coord.graph.get_graph().nodes + assert graph_mod.BUILD_NODE not in nodes + assert graph_mod.VERIFY_NODE not in nodes + + # An approved task runs to the P2 BUILD terminus without dispatching or + # raising (it never reaches a live build/dispatch). + thread_id, _ = start_task(coord.graph, transport="slack") + final = resume_task(coord.graph, thread_id=thread_id, answer="scope is X") + assert final["current_phase"] == Phase.BUILD.value + finally: + review_loop._review_invoker = saved_invoker + + +class _FakeTransport: + """Minimal non-posting transport for the inert-wiring coordinator test.""" + + def post_question(self, **_kwargs: Any) -> str: + return "ref" + + def parse_answer(self, raw: Any) -> tuple[str, Any, str]: # pragma: no cover + raise NotImplementedError diff --git a/agent-team/tests/test_run_team.py b/agent-team/tests/test_run_team.py index ad6de38..6fa46cb 100644 --- a/agent-team/tests/test_run_team.py +++ b/agent-team/tests/test_run_team.py @@ -647,6 +647,8 @@ class _FakeCoordinator: build_clarify_node: Any = None, build_plan_node: Any = None, review_wiring: Any = None, + build_verify_wiring: Any = None, + dispatch_node_wiring: Any = None, notify: Any = None, alarm_hook: Any = None, ) -> None: @@ -659,6 +661,10 @@ class _FakeCoordinator: self.build_clarify_node = build_clarify_node self.build_plan_node = build_plan_node self.review_wiring = review_wiring + # P3 fail-safe serve default (Decision 5): only the ``serve`` command + # auto-binds these; start/intake leave them None. + self.build_verify_wiring = build_verify_wiring + self.dispatch_node_wiring = dispatch_node_wiring self.setup_called = False self.start_kwargs: dict[str, Any] | None = None self.new_task_callback: Any = None @@ -863,6 +869,87 @@ def test_intake_github_label_required(cli: ModuleType) -> None: parser.parse_args(["intake-github", "--owner", "o", "--repo", "r"]) +# --------------------------------------------------------------------------- # +# serve fail-safe P3 wiring default (design Decision 5; UNIT 0e) +# --------------------------------------------------------------------------- # + + +def _serve_args(db_path: Path, command: str) -> argparse.Namespace: + """A minimal args namespace for ``_build_coordinator`` (dry-run, no token).""" + return argparse.Namespace( + command=command, + db=db_path, + transport="slack", + dry_run=True, + ) + + +def test_serve_binds_inert_p3_wiring_when_env_unset( + cli: ModuleType, + db_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + """``serve`` is the new P3 default but degrades to inert (None) when the P3 + env is unset — never crashing serve-start.""" + for name in ( + "AGENT_TEAM_REPO_OWNER", + "AGENT_TEAM_REPO_NAME", + "AGENT_TEAM_CI_READ_TOKEN", + "GITHUB_TOKEN", + ): + monkeypatch.delenv(name, raising=False) + _FakeCoordinator.instances.clear() + monkeypatch.setattr( + "agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True + ) + + coord = cli._build_coordinator(_serve_args(db_path, "serve")) + + assert coord.build_verify_wiring is None + assert coord.dispatch_node_wiring is None + + +def test_serve_binds_live_p3_wiring_when_env_set( + cli: ModuleType, + db_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + """With the P3 env provisioned, ``serve`` binds the live wiring pair.""" + monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "Sea-Haven-Industries") + monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "orchestrator") + monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test") + _FakeCoordinator.instances.clear() + monkeypatch.setattr( + "agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True + ) + + coord = cli._build_coordinator(_serve_args(db_path, "serve")) + + assert callable(coord.build_verify_wiring) + assert callable(coord.dispatch_node_wiring) + + +def test_start_does_not_bind_p3_wiring_even_when_env_set( + cli: ModuleType, + db_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Only ``serve`` (the daemon) auto-binds P3; the one-shot ``start`` path runs + to the first human gate and never wires build/verify/dispatch.""" + monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "Sea-Haven-Industries") + monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "orchestrator") + monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test") + _FakeCoordinator.instances.clear() + monkeypatch.setattr( + "agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True + ) + + coord = cli._build_coordinator(_serve_args(db_path, "start")) + + assert coord.build_verify_wiring is None + assert coord.dispatch_node_wiring is None + + def test_build_transport_dry_run_returns_dry_run_transport(cli: ModuleType) -> None: """``dry_run=True`` yields a _DryRunTransport whose post returns a synthetic ref.""" args = argparse.Namespace(dry_run=True, transport="slack") diff --git a/agent-team/tests/test_verifier.py b/agent-team/tests/test_verifier.py index dc362b5..0fcf3ee 100644 --- a/agent-team/tests/test_verifier.py +++ b/agent-team/tests/test_verifier.py @@ -202,3 +202,80 @@ def test_allowed_scope_threaded_to_gate() -> None: cfg = VerifierConfig(expected_run_id="r1", allowed_scope=["src/"]) out = verifier_node(_state(diff, ci), cfg) assert out["status"] == TaskStatus.PARKED.value # out-of-scope -> BLOCK -> park + + +# --------------------------------------------------------------------------- # +# Per-task run-id binding (design §4 Decision 4) +# --------------------------------------------------------------------------- # + + +def test_state_run_id_overrides_static_config_run_id() -> None: + # The gate must bind to the run id THIS task dispatched (state["run_id"]), + # not the static config constant. A CI conclusion keyed to the per-task + # run id passes even though config carries a different (stale) run id. + diff = _diff_for("src/foo.py") + state = _state( + diff, {"run_id": "task-run", "conclusion": "success", "diff_hash": _hash(diff)} + ) + state["run_id"] = "task-run" + # config.expected_run_id is a DIFFERENT, stale value — state must win. + out = verifier_node(state, VerifierConfig(expected_run_id="stale-wiring-run")) + assert out["status"] == TaskStatus.DONE.value + assert out["review_verdicts"][0]["run_id"] == "task-run" + + +def test_substituted_run_id_is_rejected() -> None: + # Anti-substitution: a CI result whose run_id != state["run_id"] is a BLOCK, + # even with a success conclusion (someone tried to graft a passing run from + # another task onto this one). + diff = _diff_for("src/foo.py") + state = _state( + diff, + {"run_id": "other-task-run", "conclusion": "success", "diff_hash": _hash(diff)}, + ) + state["run_id"] = "my-task-run" + out = verifier_node(state, VerifierConfig(expected_run_id=None)) + assert out["status"] == TaskStatus.PARKED.value + assert out["ci_results"]["gate_decision"] == "block" + assert any("run-id mismatch" in r for r in out["review_verdicts"][0]["reasons"]) + + +def test_two_concurrent_tasks_each_gate_against_own_run_id() -> None: + # Two tasks share ONE VerifierConfig but each gates against its OWN + # state["run_id"]. Task A's CI matches A's run id (PASS); task B's CI is + # keyed to A's run id (substitution) so B BLOCKs. + shared_cfg = VerifierConfig(expected_run_id=None) + + diff_a = _diff_for("src/a.py") + state_a = _state( + diff_a, {"run_id": "run-A", "conclusion": "success", "diff_hash": _hash(diff_a)} + ) + state_a["run_id"] = "run-A" + + diff_b = _diff_for("src/b.py") + # B's fetched CI is wrongly keyed to run-A (a leaked/substituted run). + state_b = _state( + diff_b, {"run_id": "run-A", "conclusion": "success", "diff_hash": _hash(diff_b)} + ) + state_b["run_id"] = "run-B" + + out_a = verifier_node(state_a, shared_cfg) + out_b = verifier_node(state_b, shared_cfg) + + assert out_a["status"] == TaskStatus.DONE.value + assert out_a["review_verdicts"][0]["run_id"] == "run-A" + assert out_b["status"] == TaskStatus.PARKED.value + assert out_b["review_verdicts"][0]["run_id"] == "run-B" + + +def test_none_run_id_at_gate_time_blocks_never_passes() -> None: + # No per-task run_id in state AND no config fallback -> BLOCK (park), never a + # vacuous pass, even when CI reports success. + diff = _diff_for("src/foo.py") + state = _state( + diff, {"run_id": "", "conclusion": "success", "diff_hash": _hash(diff)} + ) + # state has no "run_id" key; config fallback is None. + out = verifier_node(state, VerifierConfig(expected_run_id=None)) + assert out["status"] == TaskStatus.PARKED.value + assert out["ci_results"]["gate_decision"] == "block"