This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
orchestrator/agent-team/agent_team/ci_watcher.py
Adam Moussa 1f8c7e1ee3 fix(agent-team): close P3 async-resume BLOCKs (durable CI-watcher wiring)
Remediates the Phase-0 adversarial BLOCKs:
- Durable ci_pending_provider (_enumerate_ci_pending) walks the LangGraph
  SQLite checkpointer to enumerate threads suspended at VERIFY awaiting CI;
  re-derives across restart. Excludes human-clarify gates + advanced threads.
- run-team serve wires ci_pending_provider + ci_poller + ci_timeout ONLY on a
  configured box; inert path unchanged. Closes the 'VERIFY suspended forever'
  defect: tick()->_ci_watch resumes on terminal CI or timeout-parks.
- CI resume routes through the single-flight, turn-guarded ResumeWorker.
- FIXes: run-locator skips cancelled/stale runs on rapid re-dispatch; inert-mode
  wording matches behavior; added node-level fail-closed + spurious-resume tests.
- end-to-end async-resume proof (test_p3_async_resume.py, real checkpointer).

Suite: 1270 passed, ruff clean. Branch only; not merged/deployed.
2026-06-23 19:52:04 -04:00

550 lines
22 KiB
Python

"""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__ = [
"AWAITING_CI_KEY",
"AlarmFn",
"CiPollOutcome",
"CiPollResult",
"CiWatchAction",
"CiWatchOutcome",
"CiWatchReport",
"ParkFn",
"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": <run>}``). 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;
# 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 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": <run>}`` (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,
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,
)