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.
This commit is contained in:
parent
f0c5cfe57f
commit
1f8c7e1ee3
12 changed files with 953 additions and 42 deletions
|
|
@ -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": <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;
|
||||
|
|
@ -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": <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,
|
||||
|
|
|
|||
|
|
@ -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).
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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": <run>}`` 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.
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)."""
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
)
|
||||
|
|
|
|||
245
agent-team/tests/test_p3_async_resume.py
Normal file
245
agent-team/tests/test_p3_async_resume.py
Normal file
|
|
@ -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()
|
||||
|
|
@ -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
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
|
|
|||
Reference in a new issue