"""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": }``). 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": }`` (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, )