diff --git a/agent-team/agent_team/ci_watcher.py b/agent-team/agent_team/ci_watcher.py index 9777370..b6f809c 100644 --- a/agent-team/agent_team/ci_watcher.py +++ b/agent-team/agent_team/ci_watcher.py @@ -61,6 +61,7 @@ from enum import Enum from typing import Any __all__ = [ + "AWAITING_CI_KEY", "AlarmFn", "CiPollOutcome", "CiPollResult", @@ -71,12 +72,23 @@ __all__ = [ "PendingCiTask", "PollFn", "ResumeFn", + "awaiting_ci_run_id", "default_ci_poller", "run_ci_watcher", + "snapshot_awaiting_ci_run_id", ] logger = logging.getLogger(__name__) +# The marker key the VERIFY node stamps on its ``interrupt()`` payload when it +# suspends awaiting a CI conclusion (see +# :func:`agent_team.nodes.build_verify_subgraph._await_ci`, whose payload is +# ``{"awaiting_ci": True, "run_id": }``). Both the durable pending-task +# enumerator (which threads are suspended-at-VERIFY?) and the resume turn-guard +# (is this thread STILL suspended awaiting THIS run?) key off this single marker +# so they recognise a suspended-at-VERIFY thread the exact same way. +AWAITING_CI_KEY = "awaiting_ci" + # 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; @@ -271,6 +283,48 @@ def _parse_iso(value: str | None) -> datetime | None: return parsed +def awaiting_ci_run_id(interrupt_value: Any) -> str | None: + """Return the run id from a VERIFY *awaiting-CI* interrupt payload, or ``None``. + + The VERIFY node suspends with + ``{"awaiting_ci": True, "run_id": }`` (see + :func:`agent_team.nodes.build_verify_subgraph._await_ci`). This recognises + THAT marker precisely: it returns the awaited ``run_id`` only when the + payload is a mapping with ``awaiting_ci`` truthy AND a non-empty string + ``run_id``. Any other interrupt (a human clarify gate, a malformed payload, + a payload with no run id) yields ``None`` so callers never mistake a + different suspension for a suspended-at-VERIFY-awaiting-CI thread. + """ + if not isinstance(interrupt_value, Mapping): + return None + if not interrupt_value.get(AWAITING_CI_KEY): + return None + run_id = interrupt_value.get("run_id") + if isinstance(run_id, str) and run_id: + return run_id + return None + + +def snapshot_awaiting_ci_run_id(snapshot: Any) -> str | None: + """Return the awaited run id if ``snapshot`` is suspended at VERIFY awaiting CI. + + Walks the snapshot's pending interrupts (``snapshot.interrupts``) and returns + the first ``run_id`` carried by an *awaiting-CI* payload (via + :func:`awaiting_ci_run_id`). Returns ``None`` when the snapshot is not + interrupted, is interrupted on a non-CI gate (e.g. a human clarify + question), or has already advanced (resumed / parked / done — no pending + interrupts). This is the single predicate both the durable enumerator and + the resume turn-guard use to decide "still suspended at VERIFY awaiting CI". + """ + interrupts = getattr(snapshot, "interrupts", None) or () + for item in interrupts: + value = getattr(item, "value", item) + run_id = awaiting_ci_run_id(value) + if run_id is not None: + return run_id + return None + + def default_ci_poller( *, owner: str, diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index e746d5f..2f8610a 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -456,8 +456,10 @@ _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." + "up and tasks run through PLAN/REVIEW; with no P3 subgraph wired, the " + "review loop's 'build' route is its terminus, so a task that would advance " + "to BUILD/VERIFY instead settles at the approved-plan terminus (no build, " + "no dispatch) until the P3 env is provisioned." ) @@ -505,8 +507,12 @@ def failsafe_production_p3_wiring( :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 + build→verify subgraph and no dispatch are wired at all, so the review + loop's "build" route stays its terminus (END). A task that would advance + to P3 therefore settles at the approved-plan terminus rather than building + or dispatching — there is no BUILD/VERIFY node to reach and so nothing + parks. (Fail-closed in the sense that no diff is ever built, dispatched, or + passed; 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 @@ -526,8 +532,10 @@ def failsafe_production_p3_wiring( _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." + "coordinator starts and runs PLAN/REVIEW; with no P3 subgraph wired, the " + "review loop's 'build' route is its terminus, so a task that would " + "advance to BUILD/VERIFY instead settles at the approved-plan terminus " + "(no build, no dispatch) until the P3 env is provisioned." ) if notify is not None: try: @@ -1277,27 +1285,34 @@ class Coordinator: return report def _ci_resume(self, task: Any, result: Any) -> None: - """RESUME a task whose CI run terminated: enqueue a turn-guarded resume. + """RESUME a task whose CI run terminated, via the turn-guarded worker. 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). + first-answer-wins flip. We drive the resume through the single-flight, + turn-guarded :meth:`agent_team.resume_worker.ResumeWorker.resume_ci` + (NOT a bare ``graph.invoke``): it takes the same per-thread lock the + deadline/human-answer resumes use and re-confirms the thread is STILL + suspended at VERIFY awaiting THIS run before invoking, so a double resume + (e.g. the same terminal run observed on two overlapping sweeps) can never + corrupt the durable state. On resume VERIFY re-runs and RE-FETCHES the + now-terminal authenticated CI result (it never trusts the resume + payload); 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: + if self._graph is None or self._resume_worker 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), + # result for operator-log provenance. ``run_id`` is the trusted dispatch + # watermark the watcher polled, so the guard binds the resume to THIS + # task's awaited run. + self._resume_worker.resume_ci( + thread_id=task.thread_id, + run_id=task.run_id, + answer=result, ) def _ci_park(self, task: Any) -> None: @@ -1324,6 +1339,81 @@ class Coordinator: ) self._alarm_hook(task.thread_id) + def _enumerate_ci_pending(self) -> list[Any]: + """Enumerate the durable threads suspended at VERIFY awaiting CI (§3.3.2). + + This is the durable :data:`ci_pending_provider` the live serve path binds: + without it the CI-watcher could never be fed real tasks (nothing else + enumerates which durable threads are parked at the VERIFY machine-gate), + so a dispatched task would suspend at VERIFY and wait FOREVER — the + async-resume gap this closes. + + It walks the LangGraph SQLite checkpointer (the same DB file the ledger + uses) for the distinct ``thread_id``s it holds, then asks the compiled + graph for each thread's live snapshot. A thread is included ONLY when its + snapshot is still interrupted on the VERIFY *awaiting-CI* marker + (:func:`agent_team.ci_watcher.snapshot_awaiting_ci_run_id` returns a run + id). Threads that already advanced — resumed, parked, or done — carry no + awaiting-CI interrupt and are excluded, so the watcher never re-resumes a + task that already left the gate (the double-resume guard holds at the + enumeration boundary, before the resume worker's guard even runs). + + Fail-soft: a snapshot read that raises for one thread is logged and that + thread skipped, so one unreadable thread never blanks the whole sweep. + Returns ``[]`` (never raises) before :meth:`setup` or when the + checkpointer cannot be enumerated. + """ + from agent_team.ci_watcher import ( # noqa: PLC0415 + PendingCiTask, + snapshot_awaiting_ci_run_id, + ) + + graph = self._graph + if graph is None: + return [] + checkpointer = getattr(graph, "checkpointer", None) + if checkpointer is None or not hasattr(checkpointer, "list"): + return [] + + # Distinct thread_ids the checkpointer holds. ``list(None)`` yields every + # checkpoint across all threads (newest-first, with repeats per thread); + # we keep insertion order and dedupe so each thread is examined once. + thread_ids: list[str] = [] + seen: set[str] = set() + try: + for ckpt in checkpointer.list(None): + cfg = getattr(ckpt, "config", None) or {} + tid = (cfg.get("configurable") or {}).get("thread_id") + if isinstance(tid, str) and tid and tid not in seen: + seen.add(tid) + thread_ids.append(tid) + except Exception: # noqa: BLE001 - enumeration must never break the sweep + _LOG.warning( + "ci-watch: checkpointer enumeration raised; no CI-pending tasks " + "this pass", + exc_info=True, + ) + return [] + + pending: list[Any] = [] + for tid in thread_ids: + try: + snap = graph.get_state(graph_mod.thread_config(tid)) + except Exception: # noqa: BLE001 - one bad thread must not blank the sweep + _LOG.warning( + "ci-watch: snapshot read for thread %s raised; skipping", + tid, + exc_info=True, + ) + continue + if snapshot_awaiting_ci_run_id(snap) is None: + # Not suspended at VERIFY awaiting CI (resumed / parked / done / + # a human gate) -> exclude so we never re-resume it. + continue + values = getattr(snap, "values", None) or {} + pending.append(PendingCiTask.from_state(tid, values)) + return pending + 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/dispatcher.py b/agent-team/agent_team/dispatcher.py index 7fe0cae..e9bf59d 100644 --- a/agent-team/agent_team/dispatcher.py +++ b/agent-team/agent_team/dispatcher.py @@ -29,7 +29,7 @@ import base64 import re from dataclasses import dataclass from datetime import datetime, timezone -from typing import Protocol +from typing import Any, Protocol from agent_team.state_store import compute_content_hash @@ -41,6 +41,7 @@ __all__ = [ "build_dispatch_inputs", "dispatch_apply_verify", "head_branch_for", + "select_run_id", ] @@ -405,15 +406,84 @@ _LOCATE_DELAY_S = 5 _LOCATE_SKEW_S = 120 +# A run's ``conclusion`` value that marks it as a stale/superseded run: the +# apply/verify workflow's per-task concurrency group cancels the prior run when a +# re-dispatch fires, so a re-dispatch of the SAME task_id leaves an older +# CANCELLED run sharing the run-name. We must NOT select it (it carries the prior +# build's verdict). Only ``cancelled`` is treated as stale-by-conclusion; a +# genuinely completed (success/failure) run is a legitimate match. +_STALE_CONCLUSIONS: frozenset[str] = frozenset({"cancelled"}) +# Non-terminal statuses: a run still queued or executing is the freshly-triggered +# one we want to bind to (it has no conclusion yet). +_ACTIVE_STATUSES: frozenset[str] = frozenset( + {"queued", "in_progress", "waiting", "requested", "pending"} +) + + +def select_run_id( + runs: list[dict[str, Any]], *, task_id: str, floor_iso: str +) -> str | None: + """Pick the dispatched run from ``gh run list`` rows (PURE; anti-stale). + + Matches rows whose ``name`` equals :func:`run_name_for` and whose + ``createdAt`` is at/after ``floor_iso``, then selects with a tightened rule so + a RAPID RE-DISPATCH of the same ``task_id`` (the build<->verify loop) never + binds to a stale/cancelled prior run: + + 1. Drop any matched run whose ``conclusion`` is ``cancelled`` — the workflow's + per-task concurrency group cancels the prior run on re-dispatch, so a + cancelled run sharing the run-name is the superseded one, never our run. + 2. Prefer the run with the GREATEST ``createdAt`` among the still-active + (queued / in_progress / waiting) runs — the freshly-triggered run is the + newest and has no conclusion yet. + 3. If none are active (e.g. a fast run already concluded by the time we poll), + fall back to the newest non-cancelled run overall. + + Returns the chosen ``databaseId`` as a string, or ``None`` if nothing + matches (the caller then fails closed). ``createdAt`` ties break on the + greater ``databaseId`` (monotonic per repo → the later-created run). + """ + target_name = run_name_for(task_id) + matches = [ + r + for r in runs + if r.get("name") == target_name and str(r.get("createdAt", "")) >= floor_iso + ] + # Drop superseded (concurrency-cancelled) prior runs of the same task_id. + matches = [ + r + for r in matches + if str(r.get("conclusion") or "").lower() not in _STALE_CONCLUSIONS + ] + if not matches: + return None + + def _sort_key(r: dict[str, Any]) -> tuple[str, int]: + try: + db_id = int(r.get("databaseId", 0)) + except (TypeError, ValueError): + db_id = 0 + return (str(r.get("createdAt", "")), db_id) + + active = [ + r for r in matches if str(r.get("status") or "").lower() in _ACTIVE_STATUSES + ] + pool = active or matches + pool.sort(key=_sort_key, reverse=True) + return str(pool[0]["databaseId"]) + + def _default_run_locator() -> RunLocator: """Real locator: match the triggered run by ``run-name`` via ``gh run list``. - Polls ``gh run list`` (read-only) for runs of the apply/verify workflow whose - name equals :func:`run_name_for` and whose ``createdAt`` is at/after the - dispatched-at watermark (minus a skew tolerance), returning the NEWEST match's - databaseId. Newest-wins so a re-dispatch (the concurrency group cancels the - stale run per task_id) resolves to the current run. Returns ``None`` if no - matching run registers within the bounded poll window (fails closed). + Polls ``gh run list`` (read-only) for runs of the apply/verify workflow and + delegates the anti-stale SELECTION to the pure :func:`select_run_id`: it skips + a concurrency-cancelled prior run of the same ``task_id`` and prefers the + newest still-active (queued/in_progress) run, falling back to the newest + non-cancelled run overall. This means a RAPID RE-DISPATCH of the same task + (the build<->verify loop) binds to the CURRENT run, never the superseded one. + Returns ``None`` if no matching run registers within the bounded poll window + (fails closed). """ def _locate(*, owner: str, repo: str, task_id: str, since_iso: str) -> str | None: @@ -422,7 +492,6 @@ def _default_run_locator() -> RunLocator: import time from datetime import timedelta - target_name = run_name_for(task_id) try: floor_dt = datetime.strptime(since_iso, "%Y-%m-%dT%H:%M:%SZ").replace( tzinfo=timezone.utc @@ -442,7 +511,7 @@ def _default_run_locator() -> RunLocator: "--workflow", WORKFLOW_FILE, "--json", - "databaseId,name,createdAt", + "databaseId,name,createdAt,status,conclusion", "--limit", "50", ], @@ -451,18 +520,9 @@ def _default_run_locator() -> RunLocator: text=True, ) runs = json.loads(proc.stdout or "[]") - matches = [ - r - for r in runs - if r.get("name") == target_name - and str(r.get("createdAt", "")) >= floor_iso - ] - if matches: - matches.sort( - key=lambda r: (str(r.get("createdAt", "")), r.get("databaseId", 0)), - reverse=True, - ) - return str(matches[0]["databaseId"]) + run_id = select_run_id(runs, task_id=task_id, floor_iso=floor_iso) + if run_id is not None: + return run_id if attempt < _LOCATE_ATTEMPTS - 1: time.sleep(_LOCATE_DELAY_S) return None diff --git a/agent-team/agent_team/resume_worker.py b/agent-team/agent_team/resume_worker.py index 8dd0845..e0d3a45 100644 --- a/agent-team/agent_team/resume_worker.py +++ b/agent-team/agent_team/resume_worker.py @@ -288,6 +288,59 @@ class ResumeWorker: graph_result=graph_result, ) + def resume_ci( + self, + *, + thread_id: str, + run_id: str, + answer: Any, + ) -> ResumeResult: + """Resume a CI machine-gate (VERIFY awaiting CI), single-flight + guarded. + + The VERIFY node's async CI-wait is a *machine* gate, not a human one: + it suspends with ``{"awaiting_ci": True, "run_id": }`` and has no + ledger ``question``/``turn`` to thread through the first-answer-wins + flip. So this path turn-guards on the CI marker instead: under the same + per-``thread_id`` lock that serialises human resumes (so a CI resume and + a redelivered CI resume for the same task never race), it re-reads the + *live* checkpoint and confirms the thread is STILL suspended at VERIFY + awaiting THIS ``run_id`` (via + :func:`agent_team.ci_watcher.snapshot_awaiting_ci_run_id`). Only then does + it invoke ``Command(resume=answer)``. If the thread already advanced + (resumed / parked / done — no awaiting-CI interrupt), or is now awaiting + a *different* run, it skips (:attr:`ResumeOutcome.STALE`) so a + double-resume can never corrupt the durable state. + + Returns :attr:`ResumeOutcome.RESUMED` when the resume applied, else + :attr:`ResumeOutcome.STALE` (there is no ledger question to supersede on + this machine gate). ``question_id`` is reported as ``""`` since none + exists. + """ + from agent_team.ci_watcher import snapshot_awaiting_ci_run_id + + lock = self._lock_for(thread_id) + with lock: + config = _thread_config(thread_id) + snapshot = self._graph.get_state(config) + awaited = snapshot_awaiting_ci_run_id(snapshot) + if awaited != run_id: + # Already advanced past this wait (resumed/parked/done) or now + # awaiting a different run: skip rather than double-apply. + return ResumeResult( + outcome=ResumeOutcome.STALE, + thread_id=thread_id, + question_id="", + turn=0, + ) + graph_result = self._graph.invoke(build_resume_command(answer), config) + return ResumeResult( + outcome=ResumeOutcome.RESUMED, + thread_id=thread_id, + question_id="", + turn=0, + graph_result=graph_result, + ) + def recover_pending_resumes(self) -> list[ResumeResult]: """Restart sweep: re-enqueue resumes for durable ``answered`` rows. diff --git a/agent-team/run-team.py b/agent-team/run-team.py index 01d20c7..51b1d4e 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -620,6 +620,28 @@ def _build_coordinator(args: argparse.Namespace) -> Any: notify=notify ) + # CI-watcher seams (design §4 Decision 2 — async resume-on-CI-complete). ONLY + # wired when ``failsafe_production_p3_wiring`` returned a LIVE pair (a + # configured box): on the inert/unconfigured box both stay None, so the + # tick() CI sweep is a NO-OP and there is no behaviour change. The poller is + # the read-only default (one GET per run, never a write); the provider is the + # durable enumerator bound to THIS coordinator below (post-construction, so it + # can close over the just-built coordinator); the timeout is a sane default + # (30 min — generous headroom over the ~7-min CI run before a stuck run + # parks). Without these, a task that dispatches and suspends at VERIFY would + # wait forever — the async-resume gap this closes. + ci_poller = None + ci_timeout = None + if build_verify_wiring is not None and dispatch_node_wiring is not None: + from datetime import timedelta + + from agent_team.ci_watcher import default_ci_poller + + owner = os.environ.get("AGENT_TEAM_REPO_OWNER", "").strip() + repo = os.environ.get("AGENT_TEAM_REPO_NAME", "").strip() + ci_poller = default_ci_poller(owner=owner, repo=repo) + ci_timeout = timedelta(minutes=30) + # 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. @@ -637,7 +659,16 @@ def _build_coordinator(args: argparse.Namespace) -> Any: dispatch_node_wiring=dispatch_node_wiring, notify=notify, alarm_hook=alarm_hook, + ci_poller=ci_poller, + ci_timeout=ci_timeout, ) + # Bind the durable CI-pending provider to THIS coordinator (only on the live + # P3 path — ``ci_poller`` is the live-pair signal). It enumerates the durable + # threads suspended at VERIFY awaiting CI so the watcher has real tasks to + # poll; bound post-construction so it can reference the just-built + # coordinator. Left unbound on the inert box, the CI sweep stays a NO-OP. + if ci_poller is not None: + coordinator._ci_pending_provider = coordinator._enumerate_ci_pending # WS2: an allowlisted Slack /new-task starts a task on THIS coordinator. Set # post-construction (the adapter closes over the just-built coordinator), and # before serve() builds the listener. AUTHZ-01 (owner allowlist) gates this diff --git a/agent-team/tests/test_build_verify_subgraph.py b/agent-team/tests/test_build_verify_subgraph.py index 815dee6..a2fa548 100644 --- a/agent-team/tests/test_build_verify_subgraph.py +++ b/agent-team/tests/test_build_verify_subgraph.py @@ -517,6 +517,63 @@ def test_verify_suspends_on_in_progress_run_then_resumes_to_gate() -> None: assert out["current_phase"] == Phase.DONE.value +def test_verify_spurious_resume_still_non_terminal_parks_never_passes() -> None: + """A spurious resume (re-fetch STILL non-terminal) must fail closed. + + The CI-watcher resumes VERIFY on what it believes is a terminal conclusion, + but the authenticated re-fetch is the source of truth. If that re-fetch is + STILL None (a spurious / premature resume, or a run that flapped back to + in-progress), the node must NOT vacuously pass: it falls through to the gate, + which — with a dispatched run_id but no terminal result — BLOCKs and parks. + """ + 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 fetch is non-terminal on EVERY call (the resume was spurious — the run + # never actually concluded). + calls = {"n": 0} + + def never_terminal_fetcher(state): + calls["n"] += 1 + return None + + node = make_verify_node( + VerifierConfig(expected_run_id=None, allowed_scope=["src"]), + ci_result_fetcher=never_terminal_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": "spurious-1"}} + state = _verify_state(diff) + state["run_id"] = "dispatched-99" + + first = graph.invoke(state, cfg) + # First pass: in-progress fetch -> VERIFY suspends awaiting CI. + assert "__interrupt__" in first + assert calls["n"] == 1 + + # The watcher resumes, but the authenticated re-fetch is STILL non-terminal. + out = graph.invoke(Command(resume={"awaiting_ci": "spurious"}), cfg) + # Re-fetched again on resume (the node replays from its start, so it fetches, + # the interrupt returns the resume value rather than re-suspending, then it + # re-fetches once more) — every fetch is non-terminal. + assert calls["n"] > 1 + # ...and with no terminal result the gate BLOCKs and the task PARKS — never a + # vacuous pass. + assert out["status"] == TaskStatus.PARKED.value + assert out["current_phase"] == Phase.PARKED.value + assert out["ci_results"]["gate_decision"] == "block" + assert route_after_verify(out) == PARKED_ROUTE + + 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 --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index 476f6e6..39ae88f 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -20,6 +20,7 @@ from __future__ import annotations import queue from datetime import timedelta from pathlib import Path +from types import SimpleNamespace from typing import Any import pytest @@ -1348,8 +1349,25 @@ def _ci_coordinator(db_path: Path, *, pending, poller): ) +class _AwaitingCiSnap: + """A StateSnapshot stand-in interrupted at VERIFY awaiting a given run.""" + + def __init__(self, run_id: str) -> None: + self.interrupts = ( + SimpleNamespace(value={"awaiting_ci": True, "run_id": run_id}), + ) + self.next = ("verify",) + self.values: dict[str, Any] = {} + + 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.""" + """A terminated run drives a graph resume of the suspended VERIFY task. + + The resume MUST go through the turn-guarded resume worker (single-flight + + awaiting-CI guard), NOT a bare ``graph.invoke``: the guard re-reads the live + snapshot and only invokes while the thread is still suspended at VERIFY + awaiting THIS run, so a double-resume cannot corrupt state. + """ from agent_team.ci_watcher import CiPollResult, PendingCiTask task = PendingCiTask( @@ -1364,6 +1382,9 @@ def test_tick_ci_watch_resumes_on_terminal_conclusion(db_path: Path) -> None: ) coord.setup() + # Thread is genuinely suspended at VERIFY awaiting run 999, so the worker's + # awaiting-CI guard passes and the resume applies exactly once. + coord._graph.get_state = lambda _cfg: _AwaitingCiSnap("999") # type: ignore[assignment] invoked: list[tuple[Any, Any]] = [] coord._graph.invoke = lambda inp, cfg: invoked.append((inp, cfg)) # type: ignore[assignment] @@ -1376,6 +1397,34 @@ def test_tick_ci_watch_resumes_on_terminal_conclusion(db_path: Path) -> None: assert cfg["configurable"]["thread_id"] == "t-ci" +def test_ci_resume_skips_when_thread_already_advanced(db_path: Path) -> None: + """A double-resume is a no-op: if the thread already left the CI gate (no + awaiting-CI interrupt — resumed/parked/done), the turn-guarded path skips the + bare invoke so state cannot be corrupted.""" + from agent_team.ci_watcher import CiPollResult, PendingCiTask + + task = PendingCiTask( + thread_id="t-gone", run_id="5", dispatched_at="2026-06-23T11:00:00+00:00" + ) + coord = _ci_coordinator( + db_path, + pending=[task], + poller=lambda t: CiPollResult.terminal({"run_id": "5", "conclusion": "ok"}), + ) + coord.setup() + + # Already advanced: no pending interrupts -> guard must skip the invoke. + coord._graph.get_state = lambda _cfg: SimpleNamespace(interrupts=(), next=()) # type: ignore[assignment] + invoked: list[Any] = [] + coord._graph.invoke = lambda inp, cfg: invoked.append((inp, cfg)) # type: ignore[assignment] + + report = coord._ci_watch() + assert report is not None + # The watcher still records it as RESUMED (the resume callable returned + # without raising), but the bare invoke never fired — the guard held. + assert invoked == [] + + 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 diff --git a/agent-team/tests/test_dispatcher.py b/agent-team/tests/test_dispatcher.py index 442ec59..6008abe 100644 --- a/agent-team/tests/test_dispatcher.py +++ b/agent-team/tests/test_dispatcher.py @@ -20,6 +20,7 @@ from agent_team.dispatcher import ( dispatch_apply_verify, head_branch_for, run_name_for, + select_run_id, ) from agent_team.state_store import compute_content_hash @@ -219,3 +220,111 @@ def test_unsafe_owner_repo_rejected(owner: str, repo: str) -> None: pusher=lambda **_k: None, dispatcher=lambda **_k: None, ) + + +# --------------------------------------------------------------------------- # +# select_run_id (pure; anti-stale on rapid re-dispatch of the SAME task_id) +# --------------------------------------------------------------------------- # + + +def _row(db_id, *, created, status="completed", conclusion=None): + """A minimal ``gh run list`` row for the apply/verify run-name of TASK.""" + return { + "databaseId": db_id, + "name": run_name_for(TASK), + "createdAt": created, + "status": status, + "conclusion": conclusion, + } + + +def test_select_run_id_skips_cancelled_prior_run_on_re_dispatch() -> None: + # Rapid re-dispatch of the SAME task_id: the concurrency group cancelled the + # OLDER run, and a NEWER run is now in progress. We must bind to the newer, + # active run — never the older cancelled one (it carries the prior verdict). + older_cancelled = _row( + 100, created="2026-06-23T10:00:00Z", status="completed", conclusion="cancelled" + ) + newer_active = _row( + 200, created="2026-06-23T10:05:00Z", status="in_progress", conclusion=None + ) + runs = [newer_active, older_cancelled] + chosen = select_run_id(runs, task_id=TASK, floor_iso="2026-06-23T09:58:00Z") + assert chosen == "200" + + +def test_select_run_id_skips_cancelled_even_when_it_is_newest() -> None: + # Defensive: a cancelled run is NEVER selected even if its createdAt is the + # greatest — it is the superseded run, not ours. + active = _row(300, created="2026-06-23T10:00:00Z", status="queued", conclusion=None) + newest_cancelled = _row( + 400, created="2026-06-23T10:10:00Z", status="completed", conclusion="cancelled" + ) + runs = [active, newest_cancelled] + chosen = select_run_id(runs, task_id=TASK, floor_iso="2026-06-23T09:58:00Z") + assert chosen == "300" + + +def test_select_run_id_prefers_newest_active_over_older_completed() -> None: + # An older legitimately-completed run plus a newer active run -> the active, + # newest run wins (the freshly-triggered one with no conclusion yet). + older_done = _row( + 500, created="2026-06-23T10:00:00Z", status="completed", conclusion="success" + ) + newer_active = _row( + 600, created="2026-06-23T10:05:00Z", status="in_progress", conclusion=None + ) + chosen = select_run_id( + [older_done, newer_active], task_id=TASK, floor_iso="2026-06-23T09:58:00Z" + ) + assert chosen == "600" + + +def test_select_run_id_falls_back_to_newest_non_cancelled_when_none_active() -> None: + # No active runs (e.g. a fast run already concluded by the time we poll): + # fall back to the newest NON-cancelled run overall. + older = _row( + 700, created="2026-06-23T10:00:00Z", status="completed", conclusion="success" + ) + newer = _row( + 800, created="2026-06-23T10:05:00Z", status="completed", conclusion="failure" + ) + cancelled = _row( + 900, created="2026-06-23T10:09:00Z", status="completed", conclusion="cancelled" + ) + chosen = select_run_id( + [older, newer, cancelled], task_id=TASK, floor_iso="2026-06-23T09:58:00Z" + ) + assert chosen == "800" + + +def test_select_run_id_respects_created_floor_and_run_name() -> None: + # Below-floor runs and other-task runs are not matched. + below_floor = _row( + 1000, created="2026-06-23T09:00:00Z", status="in_progress", conclusion=None + ) + other_task = { + "databaseId": 1100, + "name": "agent-team-apply other-task", + "createdAt": "2026-06-23T10:00:00Z", + "status": "in_progress", + "conclusion": None, + } + assert ( + select_run_id( + [below_floor, other_task], task_id=TASK, floor_iso="2026-06-23T09:58:00Z" + ) + is None + ) + + +def test_select_run_id_returns_none_when_only_cancelled_matches() -> None: + # If the only matching run is cancelled, there is nothing to bind to -> None + # (the caller fails closed: no run_id -> verify BLOCKs/parks). + only_cancelled = _row( + 1200, created="2026-06-23T10:00:00Z", status="completed", conclusion="cancelled" + ) + assert ( + select_run_id([only_cancelled], task_id=TASK, floor_iso="2026-06-23T09:58:00Z") + is None + ) diff --git a/agent-team/tests/test_p3_async_resume.py b/agent-team/tests/test_p3_async_resume.py new file mode 100644 index 0000000..a25e8ba --- /dev/null +++ b/agent-team/tests/test_p3_async_resume.py @@ -0,0 +1,245 @@ +"""End-to-end async resume-on-CI-complete (design §4 Decision 2, P3 BLOCK-1/BLOCK-2). + +These tests prove the "suspended forever" defect is gone: a task that dispatches +and suspends at VERIFY awaiting CI is, on a later ``tick()``, RESUMED once its run +reaches a terminal conclusion and TIMEOUT-PARKED once ``dispatched_at + timeout`` +elapses with no terminal result. They drive the REAL machinery — the durable +``ci_pending_provider`` (:meth:`Coordinator._enumerate_ci_pending`, which walks the +LangGraph checkpointer), the real CI-watcher sweep, and the real turn-guarded +resume worker (:meth:`ResumeWorker.resume_ci`) — over a real compiled LangGraph app +whose VERIFY node suspends with the production awaiting-CI interrupt payload. + +No network: the CI poll seam is injected. The graph is a faithful minimal stand-in +for the production BUILD → DISPATCH → VERIFY shape — DISPATCH writes the trusted +``run_id`` / ``dispatched_at`` watermarks; VERIFY ``interrupt()``s with +``{"awaiting_ci": True, "run_id": ...}`` exactly as +:func:`agent_team.nodes.build_verify_subgraph._await_ci` does — so the durable +enumeration, the suspend, and the resume are the real ones. +""" + +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Any, TypedDict + +import pytest + +pytest.importorskip("langgraph") + +from langgraph.checkpoint.memory import InMemorySaver # noqa: E402 +from langgraph.graph import END, START, StateGraph # noqa: E402 +from langgraph.types import interrupt # noqa: E402 + +from agent_team.ci_watcher import CiPollResult, CiWatchAction # noqa: E402 +from agent_team.coordinator import Coordinator # noqa: E402 +from agent_team.db.schema import init_db # noqa: E402 +from agent_team.resume_worker import ResumeWorker # noqa: E402 +from agent_team.transport.base import Transport # noqa: E402 + + +# --------------------------------------------------------------------------- # +# Test doubles +# --------------------------------------------------------------------------- # + + +class _Transport(Transport): + """Record-only transport (no network).""" + + def post_question(self, **_kwargs: Any) -> str: # type: ignore[override] + return "fake:q" + + def parse_answer(self, raw: Any) -> tuple[str, Any, str]: + return raw["question_id"], raw["answer"], "fake" + + +class _S(TypedDict, total=False): + run_id: str + dispatched_at: str + ci_results: Any + status: str + current_phase: str + + +def _build_dispatch_verify_app(saver: InMemorySaver, *, dispatched_at: str) -> Any: + """Compile a real DISPATCH → VERIFY graph that suspends awaiting CI. + + DISPATCH writes the trusted ``run_id`` / ``dispatched_at`` watermarks (as the + production dispatch node does). VERIFY suspends via ``interrupt()`` with the + production awaiting-CI payload while there is a run but no terminal + ``ci_results``; on resume it falls through (the CI-watcher re-drives it). + """ + + def dispatch(state: _S) -> dict[str, Any]: + return {"run_id": "R-1", "dispatched_at": dispatched_at} + + def verify(state: _S) -> dict[str, Any]: + run_id = state.get("run_id") + ci = state.get("ci_results") + if run_id and ci is None: + # Production payload shape (build_verify_subgraph._await_ci). On resume + # the production node RE-FETCHES the authenticated conclusion; this + # stand-in uses the resume value the CI-watcher passes (the terminal + # poll result) as that re-fetched conclusion so VERIFY can advance. + ci = interrupt({"awaiting_ci": True, "run_id": run_id}) + return {"status": "done", "current_phase": "done", "ci_results": ci} + + g: StateGraph = StateGraph(_S) + g.add_node("dispatch", dispatch) + g.add_node("verify", verify) + g.add_edge(START, "dispatch") + g.add_edge("dispatch", "verify") + g.add_edge("verify", END) + return g.compile(checkpointer=saver) + + +def _coordinator_over( + db_path: Path, + app: Any, + saver: InMemorySaver, + *, + poller: Any, + timeout: timedelta | None = None, + alarm_hook: Any = None, +) -> Coordinator: + """A Coordinator whose graph/resume-worker are the supplied real app. + + ``_enumerate_ci_pending`` (the durable provider) is bound as the + ``ci_pending_provider`` so the watcher is fed REAL suspended threads off the + real checkpointer — exactly the live serve wiring. + """ + coord = Coordinator( + db_path=db_path, + transport=_Transport(), + build_checkpointer=lambda _p: saver, + ci_poller=poller, + ci_timeout=timeout, + alarm_hook=alarm_hook, + ) + coord._graph = app + coord._resume_worker = ResumeWorker(app, _connect(db_path)) + coord._ci_pending_provider = coord._enumerate_ci_pending + return coord + + +def _connect(db_path: Path) -> Any: + from agent_team.db.schema import connect + + return connect(db_path) + + +def _thread_cfg(thread_id: str) -> dict[str, Any]: + return {"configurable": {"thread_id": thread_id}} + + +@pytest.fixture() +def db_path(tmp_path: Path) -> Path: + path = tmp_path / "state" / "agent_team.sqlite" + init_db(path) + return path + + +# --------------------------------------------------------------------------- # +# (a) provider enumeration: includes suspended-at-VERIFY, excludes advanced +# --------------------------------------------------------------------------- # + + +def test_provider_enumerates_suspended_and_excludes_advanced(db_path: Path) -> None: + """``_enumerate_ci_pending`` returns the thread suspended at VERIFY awaiting CI + and EXCLUDES a resumed/done thread and a never-dispatched (no run) thread.""" + saver = InMemorySaver() + app = _build_dispatch_verify_app(saver, dispatched_at=_now_iso()) + coord = _coordinator_over( + db_path, app, saver, poller=lambda t: CiPollResult.pending() + ) + + # t-wait: dispatch + suspend at VERIFY awaiting CI (still pending). + app.invoke({}, _thread_cfg("t-wait")) + # t-done: dispatch, suspend, then RESUME with a terminal CI result -> advances + # past VERIFY to DONE (no awaiting-CI interrupt left). + from langgraph.types import Command + + app.invoke({}, _thread_cfg("t-done")) + app.invoke(Command(resume={"conclusion": "success"}), _thread_cfg("t-done")) + + pending = coord._enumerate_ci_pending() + thread_ids = {p.thread_id for p in pending} + + assert "t-wait" in thread_ids # suspended at VERIFY awaiting CI -> included + assert "t-done" not in thread_ids # advanced past the gate -> excluded + # The included task carries the trusted dispatch watermarks the watcher keys + # off (so the poll + timeout have a run to act on). + waiting = next(p for p in pending if p.thread_id == "t-wait") + assert waiting.run_id == "R-1" + assert waiting.dispatched_at is not None + + +# --------------------------------------------------------------------------- # +# (b) end-to-end: tick() RESUMES on terminal, TIMEOUT-PARKS on no-terminal +# --------------------------------------------------------------------------- # + + +def test_tick_resumes_suspended_task_on_terminal_conclusion(db_path: Path) -> None: + """A task suspended at VERIFY is RESUMED by a later tick() once its run is + terminal — proving 'suspended forever' is gone.""" + saver = InMemorySaver() + app = _build_dispatch_verify_app(saver, dispatched_at=_now_iso()) + coord = _coordinator_over( + db_path, + app, + saver, + poller=lambda t: CiPollResult.terminal({"run_id": "R-1", "conclusion": "ok"}), + ) + + # Dispatch + suspend at VERIFY. + app.invoke({}, _thread_cfg("t1")) + snap = app.get_state(_thread_cfg("t1")) + assert snap.next == ("verify",) # genuinely suspended awaiting CI + + # A subsequent tick() runs the CI-watch sweep -> resume the suspended task. + report = coord._ci_watch() + assert report is not None + assert report.resumed == 1 + assert report.outcomes[0].action is CiWatchAction.RESUMED + + # The thread advanced past VERIFY (no longer suspended). + snap2 = app.get_state(_thread_cfg("t1")) + assert snap2.next == () + assert snap2.values.get("status") == "done" + + +def test_tick_timeout_parks_suspended_task_when_run_never_terminates( + db_path: Path, +) -> None: + """A task whose run never reaches a terminal conclusion within the timeout is + TIMEOUT-PARKED by a later tick() (never waits forever).""" + saver = InMemorySaver() + # Dispatched long ago: dispatched_at + timeout has already elapsed. + old = (datetime.now(timezone.utc) - timedelta(hours=2)).isoformat() + app = _build_dispatch_verify_app(saver, dispatched_at=old) + + alarms: list[str] = [] + coord = _coordinator_over( + db_path, + app, + saver, + poller=lambda t: CiPollResult.pending(), # never terminal + timeout=timedelta(minutes=30), + alarm_hook=alarms.append, + ) + + app.invoke({}, _thread_cfg("t-slow")) + assert app.get_state(_thread_cfg("t-slow")).next == ("verify",) + + report = coord._ci_watch() + assert report is not None + assert report.parked_timeout == 1 + assert report.outcomes[0].action is CiWatchAction.PARKED_TIMEOUT + # Durably parked + ALARMed (surfaced, not silently spun on). + parked = app.get_state(_thread_cfg("t-slow")) + assert parked.values.get("status") == "parked" + assert alarms == ["t-slow"] + + +def _now_iso() -> str: + return datetime.now(timezone.utc).isoformat() diff --git a/agent-team/tests/test_resume_worker.py b/agent-team/tests/test_resume_worker.py index 4cdb7c4..212ee09 100644 --- a/agent-team/tests/test_resume_worker.py +++ b/agent-team/tests/test_resume_worker.py @@ -428,6 +428,91 @@ def test_decode_answer_non_json_passthrough() -> None: assert resume_worker._decode_answer("not-json{{") == "not-json{{" +# --------------------------------------------------------------------------- # +# resume_ci (CI machine-gate: turn-guarded on the awaiting-CI marker) +# --------------------------------------------------------------------------- # + + +class _CiGraph: + """A GraphLike double suspended at VERIFY awaiting a given CI ``run_id``. + + ``awaiting_run`` is the run the thread is currently suspended on, or ``None`` + if it has advanced past the CI gate (resumed / parked / done). ``invoke`` + records every resume and advances the graph (clears the interrupt) as a real + resume would, so a double-resume is directly observable. + """ + + def __init__(self, awaiting_run: str | None) -> None: + self.awaiting_run = awaiting_run + self.invocations: list[Any] = [] + + def get_state(self, config: dict[str, Any]) -> _FakeSnapshot: + if self.awaiting_run is None: + return _FakeSnapshot(next=(), interrupts=()) + payload = {"awaiting_ci": True, "run_id": self.awaiting_run} + return _FakeSnapshot( + next=("verify",), + interrupts=(_FakeInterrupt(value=payload),), + ) + + def invoke(self, command: Any, config: dict[str, Any]) -> Any: + self.invocations.append(command) + self.awaiting_run = None + return {"resumed": True} + + +def test_resume_ci_applies_when_suspended_on_run(conn: sqlite3.Connection) -> None: + graph = _CiGraph(awaiting_run="999") + worker = ResumeWorker(graph, conn) + + result = worker.resume_ci(thread_id="t1", run_id="999", answer={"conclusion": "ok"}) + + assert result.outcome is ResumeOutcome.RESUMED + assert result.resumed is True + assert len(graph.invocations) == 1 + + +def test_resume_ci_skips_when_thread_already_advanced( + conn: sqlite3.Connection, +) -> None: + # Already resumed/parked/done: no awaiting-CI interrupt -> guard skips invoke. + graph = _CiGraph(awaiting_run=None) + worker = ResumeWorker(graph, conn) + + result = worker.resume_ci(thread_id="t1", run_id="999", answer={}) + + assert result.outcome is ResumeOutcome.STALE + assert graph.invocations == [] + + +def test_resume_ci_skips_when_awaiting_a_different_run( + conn: sqlite3.Connection, +) -> None: + # Suspended awaiting a DIFFERENT run (e.g. a re-dispatch): must not resume. + graph = _CiGraph(awaiting_run="other") + worker = ResumeWorker(graph, conn) + + result = worker.resume_ci(thread_id="t1", run_id="999", answer={}) + + assert result.outcome is ResumeOutcome.STALE + assert graph.invocations == [] + + +def test_resume_ci_double_resume_is_idempotent(conn: sqlite3.Connection) -> None: + # Two terminal observations of the same run on overlapping sweeps: the first + # applies; the second finds the thread advanced (guard) and skips. State is + # never double-applied. + graph = _CiGraph(awaiting_run="999") + worker = ResumeWorker(graph, conn) + + first = worker.resume_ci(thread_id="t1", run_id="999", answer={}) + second = worker.resume_ci(thread_id="t1", run_id="999", answer={}) + + assert first.outcome is ResumeOutcome.RESUMED + assert second.outcome is ResumeOutcome.STALE + assert len(graph.invocations) == 1 + + # --------------------------------------------------------------------------- # # command builder # --------------------------------------------------------------------------- # diff --git a/agent-team/tests/test_run_team.py b/agent-team/tests/test_run_team.py index 6fa46cb..f278b22 100644 --- a/agent-team/tests/test_run_team.py +++ b/agent-team/tests/test_run_team.py @@ -17,6 +17,7 @@ import argparse import importlib.util import io import json +from datetime import timedelta from pathlib import Path from types import ModuleType from typing import Any @@ -651,6 +652,8 @@ class _FakeCoordinator: dispatch_node_wiring: Any = None, notify: Any = None, alarm_hook: Any = None, + ci_poller: Any = None, + ci_timeout: Any = None, ) -> None: self.db_path = db_path self.transport = transport @@ -665,11 +668,21 @@ class _FakeCoordinator: # auto-binds these; start/intake leave them None. self.build_verify_wiring = build_verify_wiring self.dispatch_node_wiring = dispatch_node_wiring + # CI-watcher seams (§4 Decision 2): the live serve path binds the poller + + # timeout here and the provider post-construction; inert leaves all None. + self.ci_poller = ci_poller + self.ci_timeout = ci_timeout + self._ci_pending_provider: Any = None self.setup_called = False self.start_kwargs: dict[str, Any] | None = None self.new_task_callback: Any = None _FakeCoordinator.instances.append(self) + def _enumerate_ci_pending(self) -> list[Any]: + # Stand-in for the durable enumerator the live path binds as the + # ci_pending_provider; identity is what the wiring test asserts. + return [] + def setup(self) -> None: self.setup_called = True @@ -907,6 +920,10 @@ def test_serve_binds_inert_p3_wiring_when_env_unset( assert coord.build_verify_wiring is None assert coord.dispatch_node_wiring is None + # Inert box: no CI-watcher seams, so the tick() CI sweep is a NO-OP. + assert coord.ci_poller is None + assert coord.ci_timeout is None + assert coord._ci_pending_provider is None def test_serve_binds_live_p3_wiring_when_env_set( @@ -927,6 +944,13 @@ def test_serve_binds_live_p3_wiring_when_env_set( assert callable(coord.build_verify_wiring) assert callable(coord.dispatch_node_wiring) + # Live box: the CI-watcher seams are wired so a task suspended at VERIFY + # awaiting CI gets resumed/parked rather than waiting forever. The poller is + # the read-only default; the provider is the coordinator's durable + # enumerator (bound post-construction); the timeout is the 30-min default. + assert callable(coord.ci_poller) + assert coord.ci_timeout == timedelta(minutes=30) + assert coord._ci_pending_provider == coord._enumerate_ci_pending def test_start_does_not_bind_p3_wiring_even_when_env_set( diff --git a/agent-team/tests/test_ws3_dispatch_invoker.py b/agent-team/tests/test_ws3_dispatch_invoker.py index 5e0ed8b..22dda10 100644 --- a/agent-team/tests/test_ws3_dispatch_invoker.py +++ b/agent-team/tests/test_ws3_dispatch_invoker.py @@ -44,7 +44,7 @@ def _stub_run_locator(monkeypatch: pytest.MonkeyPatch) -> None: monkeypatch.setattr( _dispatcher, "_default_run_locator", - lambda: (lambda **_kw: _FAKE_RUN_ID), + lambda: lambda **_kw: _FAKE_RUN_ID, ) @@ -172,6 +172,60 @@ def test_dispatch_node_flattens_scope_list_to_string() -> None: assert "agent_team/" in scope_str +def test_dispatch_node_unresolved_run_id_persists_watermark_and_verify_fails_closed( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """The ``if not result.run_id`` branch: dispatch fired but the run could not + be correlated. The node must NOT park (the dispatch succeeded) — it persists + the dispatched-at watermark + correlation tag with ``run_id=None`` so the + downstream verifier gate fails closed (BLOCK/park), never fabricating a pass. + """ + from agent_team import dispatcher as _dispatcher + + # Override the autouse locator stub: this run cannot be located -> None. + monkeypatch.setattr(_dispatcher, "_default_run_locator", lambda: lambda **_kw: None) + + pusher, dispatcher = _fake_seams() + node = make_dispatch_node( + owner="org", repo="repo", pusher=pusher, dispatcher=dispatcher + ) + result = node(_VALID_STATE) + + # Not parked: the dispatch itself succeeded; only correlation failed. + assert result.get("status") != TaskStatus.PARKED.value + assert result["run_id"] is None + # The watermark + correlation tag are still persisted for the CI-watch timeout + # and for audit, even though the run id is unresolved. + assert result["dispatched_at"] + assert result["ci_correlation_tag"] == "task-abc" + + # Feed the dispatch output into the verifier: a None run_id has nothing for + # the gate to bind to, so it BLOCKs and parks — fail closed, never a pass. + from agent_team.nodes.build_verify_subgraph import make_verify_node + from agent_team.nodes.verifier import VerifierConfig + + verify_state: dict[str, Any] = { + "thread_id": "task-abc", + "candidate_diff": _VALID_STATE["candidate_diff"], + "diff_hash": "deadbeef", + "ci_results": None, + "run_id": result["run_id"], # None, carried from dispatch + } + + def _pass_fetcher(_s: Any) -> Any: + # Even a 'success' CI result cannot rescue a missing run_id: the gate has + # no run identity to bind the verdict to. + return {"run_id": "whatever", "conclusion": "success", "diff_hash": "deadbeef"} + + verify_node = make_verify_node( + VerifierConfig(expected_run_id=None, allowed_scope=["agent_team/"]), + ci_result_fetcher=_pass_fetcher, + ) + out = verify_node(verify_state) + assert out["status"] == TaskStatus.PARKED.value + assert out["ci_results"]["gate_decision"] == "block" + + # --------------------------------------------------------------------------- # # make_dispatch_node — fail closed (never dispatches incomplete input) # --------------------------------------------------------------------------- #