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/tests/test_ci_watcher.py
Adam Moussa f0c5cfe57f feat(agent-team): P3 Phase-0 box-side build->dispatch->verify (WIP)
0c-binding: per-task expected_run_id bound from state (gate rejects substituted
  run_id; None -> BLOCK, never vacuous pass).
0e: fail-safe serve default (failsafe_production_p3_wiring) — inert on
  unprovisioned env (one WARNING + one #agent-team notice), never crash-loops.
0a: reorder P3 subgraph BUILD -> DISPATCH -> VERIFY (preserves _instrument).
0d: ci_watcher engine + VERIFY interrupt()-wait (async resume-on-CI-complete).

KNOWN-OPEN (adversarial review BLOCKs, to remediate next):
- CI-watcher not wired into run-team serve (ci_pending_provider/ci_poller None)
  -> a VERIFY-suspended task never resumes/parks.
- no durable ci_pending_provider enumerating threads suspended at VERIFY.
Branch only; not merged, not deployed.
2026-06-23 19:52:04 -04:00

323 lines
11 KiB
Python

"""Unit tests for agent_team.ci_watcher (design §3.3.2, P3 Decision 2).
Covers the async resume-on-CI-complete sweep: it RESUMES a task on a terminal
conclusion, PARKS on the dispatch timeout, and FAILS CLOSED (parks) on a fetch
error / unusable state. No network: every side effect (poll, resume, park) is an
injected callable, exactly as the module's contract promises.
"""
from __future__ import annotations
from datetime import datetime, timedelta, timezone
from agent_team.ci_watcher import (
CiPollResult,
CiWatchAction,
PendingCiTask,
default_ci_poller,
run_ci_watcher,
)
# A fixed "now" and a dispatched-at watermark; ISO strings mirror the ledger.
_NOW = datetime(2026, 6, 23, 12, 0, 0, tzinfo=timezone.utc)
_JUST_NOW = (_NOW - timedelta(minutes=1)).isoformat()
_LONG_AGO = (_NOW - timedelta(hours=2)).isoformat()
class _Recorder:
"""Records the tasks/results handed to an injected side-effect callback."""
def __init__(self) -> None:
self.resumed: list[tuple[PendingCiTask, dict]] = []
self.parked: list[PendingCiTask] = []
def on_resume(self, task: PendingCiTask, result: object) -> None:
self.resumed.append((task, dict(result))) # type: ignore[arg-type]
def on_park(self, task: PendingCiTask) -> None:
self.parked.append(task)
def _task(*, run_id: str | None = "12345", dispatched_at: str | None = _JUST_NOW):
return PendingCiTask(thread_id="t1", run_id=run_id, dispatched_at=dispatched_at)
# --------------------------------------------------------------------------- #
# TERMINAL -> resume
# --------------------------------------------------------------------------- #
def test_terminal_conclusion_resumes_the_task() -> None:
rec = _Recorder()
terminal = {"run_id": "12345", "conclusion": "success", "diff_hash": "abc"}
report = run_ci_watcher(
[_task()],
poll=lambda t: CiPollResult.terminal(terminal),
on_resume=rec.on_resume,
on_park=rec.on_park,
now=_NOW,
)
assert report.resumed == 1
assert report.parked == 0
assert rec.parked == []
assert len(rec.resumed) == 1
resumed_task, resumed_result = rec.resumed[0]
assert resumed_task.thread_id == "t1"
# The authenticated terminal result is handed to the resume side effect.
assert resumed_result == terminal
assert report.outcomes[0].action is CiWatchAction.RESUMED
def test_terminal_resume_even_past_timeout_prefers_resume() -> None:
"""A run that terminated should RESUME, not park, even if dispatched long ago."""
rec = _Recorder()
report = run_ci_watcher(
[_task(dispatched_at=_LONG_AGO)],
poll=lambda t: CiPollResult.terminal({"conclusion": "failure"}),
on_resume=rec.on_resume,
on_park=rec.on_park,
timeout=timedelta(minutes=30),
now=_NOW,
)
assert report.resumed == 1
assert rec.parked == []
# --------------------------------------------------------------------------- #
# PENDING + timeout -> park ; PENDING within window -> wait
# --------------------------------------------------------------------------- #
def test_pending_within_window_waits() -> None:
rec = _Recorder()
report = run_ci_watcher(
[_task(dispatched_at=_JUST_NOW)],
poll=lambda t: CiPollResult.pending(),
on_resume=rec.on_resume,
on_park=rec.on_park,
timeout=timedelta(minutes=30),
now=_NOW,
)
assert report.waiting == 1
assert report.parked == 0
assert rec.parked == []
assert rec.resumed == []
assert report.outcomes[0].action is CiWatchAction.WAITING
def test_pending_past_timeout_parks() -> None:
rec = _Recorder()
report = run_ci_watcher(
[_task(dispatched_at=_LONG_AGO)],
poll=lambda t: CiPollResult.pending(),
on_resume=rec.on_resume,
on_park=rec.on_park,
timeout=timedelta(minutes=30),
now=_NOW,
)
assert report.parked_timeout == 1
assert report.parked == 1
assert len(rec.parked) == 1
assert rec.resumed == []
assert report.outcomes[0].action is CiWatchAction.PARKED_TIMEOUT
# --------------------------------------------------------------------------- #
# ERROR / None result / unusable state -> fail closed (park)
# --------------------------------------------------------------------------- #
def test_poll_error_parks_fail_closed() -> None:
rec = _Recorder()
report = run_ci_watcher(
[_task(dispatched_at=_JUST_NOW)], # within window: error still parks
poll=lambda t: CiPollResult.error(),
on_resume=rec.on_resume,
on_park=rec.on_park,
now=_NOW,
)
assert report.parked_error == 1
assert len(rec.parked) == 1
assert rec.resumed == []
assert report.outcomes[0].action is CiWatchAction.PARKED_ERROR
def test_raising_poll_is_isolated_and_parks() -> None:
rec = _Recorder()
def boom(task: PendingCiTask) -> CiPollResult:
raise RuntimeError("github exploded")
report = run_ci_watcher(
[_task()],
poll=boom,
on_resume=rec.on_resume,
on_park=rec.on_park,
now=_NOW,
)
assert report.parked_error == 1
assert len(rec.parked) == 1
assert "RuntimeError" in (report.outcomes[0].error or "")
def test_missing_run_id_parks_without_polling() -> None:
rec = _Recorder()
polled: list[PendingCiTask] = []
def tracking_poll(task: PendingCiTask) -> CiPollResult:
polled.append(task)
return CiPollResult.pending()
report = run_ci_watcher(
[_task(run_id=None)],
poll=tracking_poll,
on_resume=rec.on_resume,
on_park=rec.on_park,
now=_NOW,
)
assert report.parked_error == 1
assert polled == [] # unusable state parks BEFORE any poll
assert len(rec.parked) == 1
def test_unparseable_dispatched_at_parks_without_polling() -> None:
rec = _Recorder()
polled: list[PendingCiTask] = []
report = run_ci_watcher(
[_task(dispatched_at="not-a-timestamp")],
poll=lambda t: (polled.append(t), CiPollResult.pending())[1],
on_resume=rec.on_resume,
on_park=rec.on_park,
now=_NOW,
)
assert report.parked_error == 1
assert polled == []
assert len(rec.parked) == 1
# --------------------------------------------------------------------------- #
# Side-effect isolation across the batch
# --------------------------------------------------------------------------- #
def test_one_task_failure_does_not_abort_the_sweep() -> None:
"""A resume callback that raises for one task does not stop the others."""
rec = _Recorder()
raised_for: list[str] = []
def flaky_resume(task: PendingCiTask, result: object) -> None:
if task.thread_id == "bad":
raised_for.append(task.thread_id)
raise RuntimeError("resume failed")
rec.resumed.append((task, dict(result))) # type: ignore[arg-type]
tasks = [
PendingCiTask(thread_id="bad", run_id="1", dispatched_at=_JUST_NOW),
PendingCiTask(thread_id="good", run_id="2", dispatched_at=_JUST_NOW),
]
report = run_ci_watcher(
tasks,
poll=lambda t: CiPollResult.terminal({"conclusion": "success"}),
on_resume=flaky_resume,
on_park=rec.on_park,
now=_NOW,
)
assert report.examined == 2
# The bad task is recorded as a PARKED_ERROR; the good one resumed.
actions = {o.thread_id: o.action for o in report.outcomes}
assert actions["bad"] is CiWatchAction.PARKED_ERROR
assert actions["good"] is CiWatchAction.RESUMED
assert raised_for == ["bad"]
def test_park_callback_failure_downgrades_to_error_outcome() -> None:
def park_boom(task: PendingCiTask) -> None:
raise RuntimeError("park write failed")
report = run_ci_watcher(
[_task(dispatched_at=_LONG_AGO)],
poll=lambda t: CiPollResult.pending(),
on_resume=lambda t, r: None,
on_park=park_boom,
timeout=timedelta(minutes=30),
now=_NOW,
)
# The intended action was a timeout-park, but the park callback raised, so the
# outcome is recorded as PARKED_ERROR (the sweep continues either way).
assert report.outcomes[0].action is CiWatchAction.PARKED_ERROR
assert "RuntimeError" in (report.outcomes[0].error or "")
# --------------------------------------------------------------------------- #
# default_ci_poller — reuses ci_fetcher read-only, classifies the result
# --------------------------------------------------------------------------- #
class _FakeResponse:
def __init__(self, status: int, body: dict) -> None:
self.status_code = status
self._body = body
def json(self) -> dict:
return self._body
class _FakeClient:
"""A read-only ``requests``-like client: records GETs, never writes."""
def __init__(self, response: _FakeResponse) -> None:
self._response = response
self.gets: list[str] = []
def get(self, url: str, *, timeout: float) -> _FakeResponse:
self.gets.append(url)
return self._response
def test_default_poller_terminal_on_success_conclusion() -> None:
client = _FakeClient(_FakeResponse(200, {"id": 12345, "conclusion": "success"}))
poll = default_ci_poller(owner="o", repo="r", client=client)
result = poll(_task(run_id="12345"))
assert result.outcome.value == "terminal"
assert result.result is not None
assert result.result["conclusion"] == "success"
# Exactly one read-only GET issued.
assert len(client.gets) == 1
def test_default_poller_pending_on_in_progress_run() -> None:
# An in-progress run has conclusion=None -> fetch_ci_result returns None.
client = _FakeClient(_FakeResponse(200, {"id": 12345, "conclusion": None}))
poll = default_ci_poller(owner="o", repo="r", client=client)
result = poll(_task(run_id="12345"))
assert result.outcome.value == "pending"
assert result.result is None
def test_default_poller_pending_on_http_error_status() -> None:
# A 404/5xx makes fetch_ci_result return None; the poller classifies it as
# pending (the timeout branch in the sweep is the fail-closed backstop, and a
# never-resolving run parks on timeout).
client = _FakeClient(_FakeResponse(404, {}))
poll = default_ci_poller(owner="o", repo="r", client=client)
result = poll(_task(run_id="12345"))
assert result.outcome.value == "pending"
def test_default_poller_error_when_fetch_raises() -> None:
class _BoomClient:
def get(self, url: str, *, timeout: float):
raise RuntimeError("network down")
# fetch_ci_result itself swallows GET errors to None, so the poller sees
# pending; but a defensive wrapper still classifies a raised fetch as error.
# Here we drive the error path by passing a client whose .get raises AND
# bypass fetch_ci_result's own swallow via a poller that re-raises is not
# possible; instead assert the read-only GET error degrades to pending (the
# timeout backstop parks it). This documents the boundary.
poll = default_ci_poller(owner="o", repo="r", client=_BoomClient())
result = poll(_task(run_id="12345"))
assert result.outcome.value == "pending"