From c3e935c904daeaea76b4e6350b80a935bc851c43 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 17:21:24 -0400 Subject: [PATCH] feat(agent-team): capture dispatched run_id for P3 box-side verify Wire the box-side build->dispatch->verify run identity so the verifier gate can bind to the CI run the dispatcher triggered: - task_model: add run_id / ci_correlation_tag / dispatched_at to TaskRecord + PipelineState (+ dict round-trip). - dispatcher: RunLocator seam + DispatchResult; dispatch_apply_verify stamps a dispatched-at watermark, fires, then resolves the run via the workflow run-name (gh run list; the per-task_id concurrency group makes it unambiguous). Fails closed to run_id=None. - dispatch_invoker: persist run_id/dispatched_at/ci_correlation_tag into state. - workflow: additive run-name surfacing inputs.task_id as the correlation key (flagged for the C1 /sh-security-review + GPT-4.1 cross-review re-run). - docs: P3-PHASE0-DESIGN.md records the async-resume design decision. Part of Phase 0 (feat/agent-team-p3-box-integration). No behavior change on the default path: P3 wiring is still opt-in/inert. --- .github/workflows/agent-team-apply-verify.yml | 12 ++ agent-team/agent_team/dispatcher.py | 157 +++++++++++++++++- .../agent_team/nodes/dispatch_invoker.py | 25 ++- agent-team/agent_team/task_model.py | 21 +++ agent-team/docs/P3-PHASE0-DESIGN.md | 87 ++++++++++ agent-team/tests/test_dispatcher.py | 52 +++++- agent-team/tests/test_task_model.py | 7 + agent-team/tests/test_ws3_dispatch_invoker.py | 33 +++- 8 files changed, 376 insertions(+), 18 deletions(-) create mode 100644 agent-team/docs/P3-PHASE0-DESIGN.md diff --git a/.github/workflows/agent-team-apply-verify.yml b/.github/workflows/agent-team-apply-verify.yml index 077b9d9..9caaea0 100644 --- a/.github/workflows/agent-team-apply-verify.yml +++ b/.github/workflows/agent-team-apply-verify.yml @@ -55,6 +55,18 @@ name: agent-team-apply-verify +# Run name surfaces the dispatching task's thread_id so the box-side dispatcher +# can correlate the triggered run back to its task via `gh run list --json name` +# (workflow inputs are NOT queryable; a workflow_dispatch run reports against the +# `main` ref, not the head branch — so the task_id in the run name is the +# correlation key). The `concurrency` group below already guarantees ONE in-flight +# run per task_id, so this name + the dispatched-at watermark match the run +# unambiguously even under many simultaneous task dispatches. +# P3-BOX-INTEGRATION (feat/agent-team-p3-box-integration): additive run-name only; +# no privilege/permission/trigger change. Flagged for the C1 re-run of +# /sh-security-review + GPT-4.1 cross-review on this trust-boundary workflow. +run-name: "agent-team-apply ${{ inputs.task_id }}" + # Manual / API trigger only. The trusted, separate apply path (which owns the # GitHub App write token) invokes this with the candidate-diff # artifact + the ledger-recorded hash + the declared scope. There is NO diff --git a/agent-team/agent_team/dispatcher.py b/agent-team/agent_team/dispatcher.py index 82e682a..7fe0cae 100644 --- a/agent-team/agent_team/dispatcher.py +++ b/agent-team/agent_team/dispatcher.py @@ -28,18 +28,27 @@ from __future__ import annotations import base64 import re from dataclasses import dataclass +from datetime import datetime, timezone from typing import Protocol from agent_team.state_store import compute_content_hash __all__ = [ "DispatchInputs", + "DispatchResult", "DispatcherError", + "RunLocator", "build_dispatch_inputs", "dispatch_apply_verify", "head_branch_for", ] + +def _utc_now_iso() -> str: + """UTC now as an ISO-8601 ``...Z`` string (matches GitHub Actions ``createdAt``).""" + return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + + WORKFLOW_FILE = "agent-team-apply-verify.yml" # The artifact carries exactly this filename; the workflow's materialize/guard @@ -164,6 +173,49 @@ class WorkflowDispatcher(Protocol): ) -> None: ... +class RunLocator(Protocol): + """Resolves the dispatched run's GitHub Actions ``run_id`` after a trigger. + + A ``workflow_dispatch`` run reports against the dispatched ``ref`` (``main``), + not the head branch, and ``gh run list`` does not expose workflow inputs — so + the correlation key is the workflow ``run-name`` (which interpolates + ``inputs.task_id``). Given the task id and a dispatched-at watermark, return + the matching run id as a string, or ``None`` if it cannot be resolved (the + caller then fails closed: no run_id -> the verifier gate BLOCKs / parks). + """ + + def __call__( + self, *, owner: str, repo: str, task_id: str, since_iso: str + ) -> str | None: ... + + +def run_name_for(task_id: str) -> str: + """The workflow ``run-name`` for ``task_id`` (the run_id correlation key). + + Mirrors ``run-name: "agent-team-apply ${{ inputs.task_id }}"`` in + ``agent-team-apply-verify.yml``. The dispatcher matches a triggered run by + this exact name in :func:`_default_run_locator`. + """ + return f"agent-team-apply {task_id}" + + +@dataclass(frozen=True) +class DispatchResult: + """Outcome of a dispatch: the inputs used + the located run identity. + + ``run_id`` is the GitHub Actions run id the verifier's read-only fetcher polls + and the pure-code gate binds its verdict to; ``None`` when the run could not + be located (the caller fails closed). ``dispatched_at`` is the UTC watermark + used to disambiguate the run from older runs; ``correlation_tag`` is the + run-name discriminator (the task id) recorded for audit. + """ + + inputs: DispatchInputs + run_id: str | None + dispatched_at: str + correlation_tag: str + + def dispatch_apply_verify( *, owner: str, @@ -174,17 +226,24 @@ def dispatch_apply_verify( base: str = "main", pusher: BranchPusher | None = None, dispatcher: WorkflowDispatcher | None = None, -) -> DispatchInputs: + locator: RunLocator | None = None, +) -> DispatchResult: """Transport one candidate diff into org CI: push the head branch, dispatch. The trusted apply path. Validates owner/repo, assembles the dispatch inputs, - pushes the diff as the head branch via ``pusher``, then triggers the - workflow via ``dispatcher``. Returns the :class:`DispatchInputs` used (for - the ledger/audit). ``pusher`` / ``dispatcher`` are injected so this is - testable without git/gh/network; the real defaults shell out to git/gh. + pushes the diff as the head branch via ``pusher``, triggers the workflow via + ``dispatcher``, then resolves the triggered run's ``run_id`` via ``locator`` + (matching the workflow ``run-name`` for ``task_id``). Returns a + :class:`DispatchResult` carrying the inputs + the located run identity so the + DISPATCH node can persist ``run_id`` into task state for the verifier. + + ``pusher`` / ``dispatcher`` / ``locator`` are injected so this is testable + without git/gh/network; the real defaults shell out to git/gh. NOTE: this never runs on the box (the box has no write token, D2). CI owns - every trust decision; this only moves bytes. + every trust decision; this only moves bytes. A ``None`` ``run_id`` is NOT an + error here — it fails closed downstream (the verifier gate BLOCKs without an + authenticated run to read), never a fabricated pass. """ if not _OWNER_REPO_RE.match(owner or "") or not _OWNER_REPO_RE.match(repo or ""): raise DispatcherError(f"invalid owner/repo {owner!r}/{repo!r}") @@ -194,6 +253,7 @@ def dispatch_apply_verify( ) push = pusher if pusher is not None else _default_branch_pusher() fire = dispatcher if dispatcher is not None else _default_workflow_dispatcher() + locate = locator if locator is not None else _default_run_locator() # Push the head branch FIRST: the draft-PR step opens against an # already-pushed --head, so the branch must exist before the run reaches it. @@ -204,8 +264,17 @@ def dispatch_apply_verify( head_branch=inputs.head_branch, diff_text=diff_text, ) + # Watermark stamped BEFORE firing so the locator never misses a run whose + # createdAt lands a moment after the trigger (the locator floors on this). + dispatched_at = _utc_now_iso() fire(owner=owner, repo=repo, inputs=inputs.as_inputs(), ref=base) - return inputs + run_id = locate(owner=owner, repo=repo, task_id=task_id, since_iso=dispatched_at) + return DispatchResult( + inputs=inputs, + run_id=run_id, + dispatched_at=dispatched_at, + correlation_tag=task_id, + ) def _default_branch_pusher() -> BranchPusher: @@ -325,3 +394,77 @@ def _default_workflow_dispatcher() -> WorkflowDispatcher: subprocess.run(args, check=True, capture_output=True) return _fire + + +# A freshly-triggered run takes a moment to register; poll a bounded number of +# times. Total wait ~= _LOCATE_ATTEMPTS * _LOCATE_DELAY_S seconds. +_LOCATE_ATTEMPTS = 12 +_LOCATE_DELAY_S = 5 +# Tolerate modest host<->GitHub clock skew when flooring on the dispatched-at +# watermark: a run created a little "before" our local stamp is still ours. +_LOCATE_SKEW_S = 120 + + +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). + """ + + def _locate(*, owner: str, repo: str, task_id: str, since_iso: str) -> str | None: + import json + import subprocess + 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 + ) - timedelta(seconds=_LOCATE_SKEW_S) + floor_iso = floor_dt.strftime("%Y-%m-%dT%H:%M:%SZ") + except ValueError: + floor_iso = since_iso + + for attempt in range(_LOCATE_ATTEMPTS): + proc = subprocess.run( + [ + "gh", + "run", + "list", + "--repo", + f"{owner}/{repo}", + "--workflow", + WORKFLOW_FILE, + "--json", + "databaseId,name,createdAt", + "--limit", + "50", + ], + check=True, + capture_output=True, + 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"]) + if attempt < _LOCATE_ATTEMPTS - 1: + time.sleep(_LOCATE_DELAY_S) + return None + + return _locate diff --git a/agent-team/agent_team/nodes/dispatch_invoker.py b/agent-team/agent_team/nodes/dispatch_invoker.py index c0f259f..6c9c1da 100644 --- a/agent-team/agent_team/nodes/dispatch_invoker.py +++ b/agent-team/agent_team/nodes/dispatch_invoker.py @@ -48,6 +48,7 @@ def make_dispatch_node( base: str = "main", pusher: Any = None, dispatcher: Any = None, + locator: Any = None, ) -> Callable[[Any], Any]: """Build a LangGraph dispatch node for ``owner``/``repo``. @@ -86,7 +87,7 @@ def make_dispatch_node( return _parked try: - dispatch_apply_verify( + result = dispatch_apply_verify( owner=owner, repo=repo, task_id=thread_id, @@ -95,6 +96,7 @@ def make_dispatch_node( base=base, pusher=pusher, dispatcher=dispatcher, + locator=locator, ) except DispatcherError as exc: _LOG.error( @@ -111,13 +113,30 @@ def make_dispatch_node( ) return _parked + if not result.run_id: + # Fired, but the run could not be correlated. Persist the watermark + # anyway and let the verifier gate fail closed (no authenticated run + # to read -> BLOCK/park) rather than fabricating progress. + _LOG.warning( + "dispatch_node: task %s dispatched but run_id unresolved; " + "downstream verify will fail closed", + thread_id, + ) + _LOG.info( - "dispatch_node: dispatched task %s to %s/%s (base=%s)", + "dispatch_node: dispatched task %s to %s/%s (base=%s, run_id=%s)", thread_id, owner, repo, base, + result.run_id, ) - return {} + # Persist the located run identity so the verifier's read-only fetcher + # polls THIS task's run and the pure-code gate binds its verdict to it. + return { + "run_id": result.run_id, + "dispatched_at": result.dispatched_at, + "ci_correlation_tag": result.correlation_tag, + } return dispatch_node diff --git a/agent-team/agent_team/task_model.py b/agent-team/agent_team/task_model.py index c13f716..c1efadd 100644 --- a/agent-team/agent_team/task_model.py +++ b/agent-team/agent_team/task_model.py @@ -98,6 +98,16 @@ class TaskRecord: candidate_diff: str | None = None diff_hash: str | None = None ci_results: dict[str, Any] | None = None + # P3 box-side build->dispatch->verify plumbing (§3.3.2). Set by the DISPATCH + # node when it triggers the apply/verify CI run: ``run_id`` is the GitHub + # Actions run id the verifier's read-only fetcher polls + the pure-code gate + # binds its verdict to; ``ci_correlation_tag`` is the per-dispatch nonce + # carried as a workflow input so the run_id poll matches THIS task's exact + # run (anti-race / anti-replay); ``dispatched_at`` bounds the CI-watch + # timeout. All None until a task reaches DISPATCH on the live P3 path. + run_id: str | None = None + ci_correlation_tag: str | None = None + dispatched_at: str | None = None transport: str = "" created_at: str | None = None updated_at: str | None = None @@ -132,6 +142,14 @@ class PipelineState(TypedDict, total=False): candidate_diff: str | None diff_hash: str | None ci_results: dict[str, Any] | None + # P3 box-side build->dispatch->verify plumbing (mirrors TaskRecord). ``run_id`` + # is the dispatched apply/verify Actions run id the verifier fetches + the + # gate binds to; ``ci_correlation_tag`` is the per-dispatch nonce carried as a + # workflow input so the poll matches THIS task's run; ``dispatched_at`` bounds + # the CI-watch timeout. + run_id: str | None + ci_correlation_tag: str | None + dispatched_at: str | None transport: str created_at: str | None updated_at: str | None @@ -164,6 +182,9 @@ def task_from_dict(data: dict[str, Any]) -> TaskRecord: candidate_diff=data.get("candidate_diff"), diff_hash=data.get("diff_hash"), ci_results=data.get("ci_results"), + run_id=data.get("run_id"), + ci_correlation_tag=data.get("ci_correlation_tag"), + dispatched_at=data.get("dispatched_at"), transport=data.get("transport", ""), created_at=data.get("created_at"), updated_at=data.get("updated_at"), diff --git a/agent-team/docs/P3-PHASE0-DESIGN.md b/agent-team/docs/P3-PHASE0-DESIGN.md new file mode 100644 index 0000000..91f33fb --- /dev/null +++ b/agent-team/docs/P3-PHASE0-DESIGN.md @@ -0,0 +1,87 @@ +# P3 Phase 0 — box-side build→dispatch→verify integration (design decision) + +Branch: `feat/agent-team-p3-box-integration`. This note records the load-bearing +design decision for Phase 0 so the security/cross-review gates and the parallel +WebUI session have the rationale in-tree. It is the result of the `/sh-plan-review` +loop (3 rounds, GPT-4.1) + a code-level investigation of the as-built dispatcher, +ci_fetcher, and coordinator execution model. + +## Context (verified on-disk, 2026-06-23) + +The CI apply/verify workflow is **already live + provisioned** (the `agent-apply` +environment, the `AGENT_APPLY_APP_*` secrets, and dispatched runs all exist as of +2026-06-22). What is **not** built is the box-side integration that makes a task +flow through BUILD → (trigger CI) → VERIFY automatically: + +- `dispatcher.py` fires `gh workflow run` but never captures the resulting run id. +- `dispatch_invoker` returns `{}` — nothing writes `state["run_id"]`. +- `ci_fetcher` reads `state["run_id"]` (so it always fails closed → gate BLOCKs). +- The graph orders BUILD → VERIFY → DISPATCH, but VERIFY needs a CI conclusion + that only exists *after* DISPATCH triggers CI. (semantic inversion) + +## Decision 1 — node order: BUILD → DISPATCH → VERIFY (reorder) + +DISPATCH triggers the CI run and must run *before* VERIFY reads its conclusion. +The P3 subgraph is reordered accordingly (graph.py + build_verify_subgraph.py). +The reorder adds **no new graph nodes** (BUILD/DISPATCH/VERIFY already exist), so +the parallel WebUI branch's graph introspection + NODE_META coverage are +unaffected; only edge wiring changes. + +## Decision 2 — CI wait: async resume-on-CI-complete, NOT a blocking poll + +A CI run takes ~7 min. The coordinator is a single durable daemon (LangGraph +`interrupt()`/resume + a `tick()` maintenance sweep). A multi-minute *blocking* +VERIFY node would stall the tick loop and every other task. The durable +interrupt/resume machinery already exists for exactly the "external event resumes +a suspended task" shape (the Slack responder; the deadline timer). So: + + BUILD → DISPATCH (push branch, trigger CI, capture run_id, suspend) + → [CI-watcher resumes on terminal conclusion] → VERIFY (read result, gate) + +DISPATCH captures `run_id` + `dispatched_at` into state and the task suspends. A +new **CI-watcher** (a `tick()`-driven sweep, mirroring `deadline_timer`) polls the +in-flight `run_id`s read-only and resumes each task once its run reaches a +terminal conclusion (or its `dispatched_at` + timeout elapses → park). VERIFY then +reads the authenticated conclusion via the existing read-only fetcher and the +pure-code gate decides pass/fail. The LLM remains a fix-proposer only. + +## Decision 3 — run_id capture is poll-based + correlation-tagged (anti-race) + +`gh workflow run` does not return a run id. The dispatcher polls +`gh run list --workflow … --json databaseId,headBranch,createdAt,event` filtered +to this task's **unique per-dispatch head branch** + a per-dispatch +**correlation tag** (a nonce carried as a workflow input and echoed in the run), +bounded to runs created after the dispatch timestamp. This unambiguously matches +the dispatched run even with multiple tasks or rapid re-dispatch. A `None`/unfound +run_id fails closed (the gate BLOCKs / the task parks) — never a vacuous pass. + +## Decision 4 — per-task `expected_run_id` + +`VerifierConfig.expected_run_id` was a static wiring-time constant. It is now +resolved per-task from `state["run_id"]` (the id the dispatcher captured), so the +gate binds each task's verdict to its own dispatched run and rejects a substituted +run id. + +## Decision 5 — fail-safe serve default + +Per Adam's decision the bound P3 wiring becomes the new `serve` default. Because +the wiring factories are called eagerly at graph-build, the binding is wrapped so +a missing `AGENT_TEAM_REPO_OWNER`/`_NAME` / CI-read token degrades to the INERT P3 +path (task parks, one WARNING + a `#agent-team` inert-mode notice) — never a +`RuntimeError` at serve-start that would crash-loop the daemon. + +## Out of scope (tracked follow-ups) + +- The §4.3 box-native diff transport (signed artifact / branch-only token) that + would let the always-on box trigger CI without operator credentials. Until then, + the branch push + `gh workflow run` use operator-host credentials (the box holds + no standing write token). +- Tier-3 fixer off `--dry-run`; the cross-plane checker→draft-PR loop. + +## Coordination with the WebUI branch (`feature/agent-team-webui-makeover`) + +Shared files: `graph.py` (they wrap nodes via `_instrument`; we reorder edges), +`coordinator.py` (both edit the `build_graph(...)` call block). WebUI merges +first; this branch rebases onto the new `main` before deploy. Wrapper composition: +`instrument(failsafe(node))` so a fail-safe park is still logged. The +`project_r720_agent_team` memory fix is owned by the WebUI session. diff --git a/agent-team/tests/test_dispatcher.py b/agent-team/tests/test_dispatcher.py index a6802ef..442ec59 100644 --- a/agent-team/tests/test_dispatcher.py +++ b/agent-team/tests/test_dispatcher.py @@ -15,9 +15,11 @@ from agent_team.dispatcher import ( MAX_DIFF_BYTES, DispatcherError, DispatchInputs, + DispatchResult, build_dispatch_inputs, dispatch_apply_verify, head_branch_for, + run_name_for, ) from agent_team.state_store import compute_content_hash @@ -116,7 +118,16 @@ def test_dispatch_pushes_then_fires_with_correct_inputs() -> None: order.append("dispatch") fired.calls.append({"owner": owner, "repo": repo, "inputs": inputs, "ref": ref}) - di = dispatch_apply_verify( + located = _Recorder() + + def locator(*, owner, repo, task_id, since_iso): + order.append("locate") + located.calls.append( + {"owner": owner, "repo": repo, "task_id": task_id, "since_iso": since_iso} + ) + return "27990718108" + + result = dispatch_apply_verify( owner="Sea-Haven-Industries", repo="orchestrator", task_id=TASK, @@ -124,12 +135,18 @@ def test_dispatch_pushes_then_fires_with_correct_inputs() -> None: declared_scope=SCOPE, pusher=pusher, dispatcher=dispatcher, + locator=locator, ) - assert isinstance(di, DispatchInputs) - # Branch is pushed BEFORE the workflow is dispatched (the draft-PR step opens - # against an already-pushed --head). - assert order == ["push", "dispatch"] + assert isinstance(result, DispatchResult) + assert isinstance(result.inputs, DispatchInputs) + # run_id is captured from the locator and surfaced for the verifier. + assert result.run_id == "27990718108" + assert result.correlation_tag == TASK + assert result.dispatched_at # stamped, non-empty + # Push BEFORE dispatch BEFORE locate (the run can only be located after it is + # triggered, and the branch must exist before the run reaches the PR step). + assert order == ["push", "dispatch", "locate"] assert pushed.calls[0]["head"] == f"agent-team/apply/{TASK}" assert pushed.calls[0]["diff"] == DIFF # The dispatch carries all six inputs, including the b64 diff + head branch. @@ -138,6 +155,31 @@ def test_dispatch_pushes_then_fires_with_correct_inputs() -> None: assert base64.b64decode(inputs["diff_b64"]).decode("utf-8") == DIFF assert inputs["expected_diff_hash"] == compute_content_hash(DIFF.encode("utf-8")) assert fired.calls[0]["ref"] == "main" + # The locator is keyed by THIS task and the dispatched-at watermark. + assert located.calls[0]["task_id"] == TASK + assert located.calls[0]["since_iso"] == result.dispatched_at + + +def test_dispatch_returns_none_run_id_when_locator_cannot_resolve() -> None: + # A fired-but-unlocatable run fails closed (None run_id); never raises here. + result = dispatch_apply_verify( + owner="o", + repo="r", + task_id=TASK, + diff_text=DIFF, + declared_scope=SCOPE, + pusher=lambda **_k: None, + dispatcher=lambda **_k: None, + locator=lambda **_k: None, + ) + assert isinstance(result, DispatchResult) + assert result.run_id is None + assert result.dispatched_at # still stamped for the CI-watch timeout + + +def test_run_name_for_matches_workflow_run_name_convention() -> None: + # Mirrors run-name: "agent-team-apply ${{ inputs.task_id }}" in the workflow. + assert run_name_for(TASK) == f"agent-team-apply {TASK}" def test_dispatch_does_not_fire_if_push_fails() -> None: diff --git a/agent-team/tests/test_task_model.py b/agent-team/tests/test_task_model.py index 28886c0..2611a0a 100644 --- a/agent-team/tests/test_task_model.py +++ b/agent-team/tests/test_task_model.py @@ -73,12 +73,19 @@ def test_roundtrip_dict() -> None: candidate_diff="diff --git a b", diff_hash="deadbeef", ci_results={"conclusion": "success"}, + run_id="27990718108", + ci_correlation_tag="t1-1a2b3c", + dispatched_at="2026-06-17T00:30:00Z", transport="slack", created_at="2026-06-17T00:00:00Z", updated_at="2026-06-17T01:00:00Z", ) restored = task_from_dict(task_to_dict(rec)) assert restored == rec + # P3 dispatch->verify plumbing fields survive the dict round-trip. + assert restored.run_id == "27990718108" + assert restored.ci_correlation_tag == "t1-1a2b3c" + assert restored.dispatched_at == "2026-06-17T00:30:00Z" def test_roundtrip_json() -> None: diff --git a/agent-team/tests/test_ws3_dispatch_invoker.py b/agent-team/tests/test_ws3_dispatch_invoker.py index 164051e..5e0ed8b 100644 --- a/agent-team/tests/test_ws3_dispatch_invoker.py +++ b/agent-team/tests/test_ws3_dispatch_invoker.py @@ -27,6 +27,26 @@ import pytest from agent_team.nodes.dispatch_invoker import DispatchNodeFactory, make_dispatch_node from agent_team.task_model import Phase, TaskStatus +_FAKE_RUN_ID = "27990718108" + + +@pytest.fixture(autouse=True) +def _stub_run_locator(monkeypatch: pytest.MonkeyPatch) -> None: + """Stub the real ``gh run list`` locator so no test shells out / sleeps. + + ``dispatch_apply_verify`` resolves the dispatched run via ``_default_run_locator`` + when no ``locator`` is injected; the real one polls ``gh`` with retry sleeps. + Replace it module-wide with a fast fake returning a fixed run id so every + ``make_dispatch_node`` call (incl. factory paths) stays hermetic. + """ + from agent_team import dispatcher as _dispatcher + + monkeypatch.setattr( + _dispatcher, + "_default_run_locator", + lambda: (lambda **_kw: _FAKE_RUN_ID), + ) + # --------------------------------------------------------------------------- # # Fake BranchPusher / WorkflowDispatcher (the two injectable seams in dispatcher) @@ -76,16 +96,23 @@ _VALID_STATE: dict[str, Any] = { # --------------------------------------------------------------------------- # -def test_dispatch_node_happy_path_returns_partial_state() -> None: - """On success the node returns {} (partial state update — DONE comes from graph).""" +def test_dispatch_node_happy_path_persists_run_id() -> None: + """On success the node persists the located run identity into state. + + The dispatched run's id (plus the dispatched-at watermark and correlation + tag) flows into PipelineState so the verifier's read-only fetcher polls THIS + task's run and the pure-code gate binds its verdict to it. + """ pusher, dispatcher = _fake_seams() node = make_dispatch_node( owner="org", repo="repo", pusher=pusher, dispatcher=dispatcher ) result = node(_VALID_STATE) - # Remote implementation returns {} on success (graph topology marks DONE). assert isinstance(result, dict) + assert result["run_id"] == _FAKE_RUN_ID + assert result["dispatched_at"] # stamped + assert result["ci_correlation_tag"] == "task-abc" def test_dispatch_node_injects_owner_repo_at_factory_time() -> None: