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

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

530 lines
20 KiB
Python

"""Trusted apply-path dispatcher (§4.3): carry a box-emitted diff into org CI.
The read-only R720 box (D2) EMITS a candidate diff + the
``workflow_dispatch`` inputs but holds **no write token**. This dispatcher is
the ONLY component that writes, and it runs on a TRUSTED host (the operator's
Mac, where the GitHub App key / operator ``gh`` auth lives) — never on the box.
It does two things and nothing else:
1. Pushes the candidate diff as a short-lived **head branch** (the PR head).
2. Triggers the ``agent-team-apply-verify.yml`` ``workflow_dispatch``, passing
the diff as ``diff_b64`` plus the integrity inputs.
The CI workflow then RE-VERIFIES everything (materialize re-hashes; guard
re-hashes + denylist + scope; build-test builds/tests credential-less; the
pure-code gate decides pass/fail; the privileged job opens a DRAFT PR only on a
clean pass, gated by the ``agent-apply`` environment's required reviewer). This
dispatcher TRANSPORTS only — it makes no trust decision.
Design discipline (mirrors the rest of agent_team): the network/SDK/subprocess
side effects are behind INJECTED seams so the pure input-assembly + validation
logic is fully unit-testable with no git, no ``gh``, and no network. The real
default seams shell out to ``git`` / ``gh`` and are exercised only in
production.
"""
from __future__ import annotations
import base64
import re
from dataclasses import dataclass
from datetime import datetime, timezone
from typing import Any, Protocol
from agent_team.state_store import compute_content_hash
__all__ = [
"DispatchInputs",
"DispatchResult",
"DispatcherError",
"RunLocator",
"build_dispatch_inputs",
"dispatch_apply_verify",
"head_branch_for",
"select_run_id",
]
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
# steps look for ``candidate.diff`` (mirrors nodes.fixer.CANDIDATE_DIFF_FILENAME).
CANDIDATE_DIFF_FILENAME = "candidate.diff"
# Max candidate-diff size. NOT arbitrary: the diff is carried as the base64
# `diff_b64` workflow_dispatch INPUT, and GitHub caps total dispatch inputs at
# ~65,535 bytes. base64 inflates ~4/3, so a diff over ~45 KB cannot be
# dispatched at all. We cap at 40 KB (leaving headroom for the other inputs) and
# fail closed with a clear error rather than letting GitHub reject the dispatch
# opaquely — and to bound resource use on the operator host. Larger diffs need a
# branch-only transport (future), not an inline input.
MAX_DIFF_BYTES = 40_000
# task_id must match the workflow's gate-and-pr safety guard charset (it is
# embedded in the head ref and the PR text). A coordinator thread_id is a uuid,
# well within this set; we validate to fail closed on anything else.
_TASK_ID_RE = re.compile(r"\A[A-Za-z0-9._-]{1,200}\Z")
# owner/repo path-segment charset (mirrors ci_fetcher / github_intake).
_OWNER_REPO_RE = re.compile(r"\A[A-Za-z0-9_.-]{1,100}\Z")
class DispatcherError(Exception):
"""Raised on structurally invalid dispatch inputs (fail closed, never proceed)."""
def head_branch_for(task_id: str) -> str:
"""Return the head-branch ref the dispatcher pushes for ``task_id``.
A stable, namespaced ref so a re-dispatch for the same task reuses/replaces
one branch rather than littering refs. Validated to the safe charset so the
ref can never carry shell/ref metacharacters.
"""
if not _TASK_ID_RE.match(task_id or ""):
raise DispatcherError(
f"invalid task_id {task_id!r}: must match {_TASK_ID_RE.pattern}"
)
return f"agent-team/apply/{task_id}"
@dataclass(frozen=True)
class DispatchInputs:
"""The exact ``workflow_dispatch`` inputs for the apply/verify workflow.
Field names match ``on.workflow_dispatch.inputs`` 1:1 so :meth:`as_inputs`
can be handed straight to ``gh workflow run -f k=v``.
"""
task_id: str
diff_artifact_name: str
expected_diff_hash: str
declared_scope: str
diff_b64: str
head_branch: str
def as_inputs(self) -> dict[str, str]:
return {
"task_id": self.task_id,
"diff_artifact_name": self.diff_artifact_name,
"expected_diff_hash": self.expected_diff_hash,
"declared_scope": self.declared_scope,
"diff_b64": self.diff_b64,
"head_branch": self.head_branch,
}
def build_dispatch_inputs(
*,
task_id: str,
diff_text: str,
declared_scope: str,
artifact_prefix: str = "agent-team-diff",
) -> DispatchInputs:
"""Assemble the dispatch inputs from a task id + the candidate diff (PURE).
Computes the sha256 the workflow binds to (``expected_diff_hash``), base64s
the diff (``diff_b64``, carried as an input so the credential-less
materialize job can re-create the artifact), derives the head branch, and
names the artifact. No I/O, no network — fully testable.
Fails closed (:class:`DispatcherError`) on an empty diff / empty scope /
unsafe task id, so a malformed task can never be transported.
"""
if not isinstance(diff_text, str) or not diff_text.strip():
raise DispatcherError("diff_text must be a non-empty diff")
if not isinstance(declared_scope, str) or not declared_scope.strip():
raise DispatcherError(
"declared_scope must be non-empty (an empty scope allows any path)"
)
head_branch = head_branch_for(task_id) # validates task_id
raw = diff_text.encode("utf-8")
if len(raw) > MAX_DIFF_BYTES:
raise DispatcherError(
f"diff too large ({len(raw)} bytes > {MAX_DIFF_BYTES}); the diff is "
"carried as a base64 workflow_dispatch input (GitHub caps inputs at "
"~64 KB). Split the change or use a branch-only transport."
)
return DispatchInputs(
task_id=task_id,
diff_artifact_name=f"{artifact_prefix}-{task_id}",
expected_diff_hash=compute_content_hash(raw),
declared_scope=declared_scope,
diff_b64=base64.b64encode(raw).decode("ascii"),
head_branch=head_branch,
)
class BranchPusher(Protocol):
"""Pushes the candidate diff as the head branch (the only WRITE)."""
def __call__(
self, *, owner: str, repo: str, base: str, head_branch: str, diff_text: str
) -> None: ...
class WorkflowDispatcher(Protocol):
"""Triggers the apply/verify ``workflow_dispatch`` with the assembled inputs."""
def __call__(
self, *, owner: str, repo: str, inputs: dict[str, str], ref: str
) -> 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,
repo: str,
task_id: str,
diff_text: str,
declared_scope: str,
base: str = "main",
pusher: BranchPusher | None = None,
dispatcher: WorkflowDispatcher | None = None,
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``, 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. 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}")
inputs = build_dispatch_inputs(
task_id=task_id, diff_text=diff_text, declared_scope=declared_scope
)
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.
push(
owner=owner,
repo=repo,
base=base,
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)
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:
"""Real pusher: clone-free apply of the diff onto a fresh branch via git/gh.
Deferred to call time (no subprocess at import). Uses the operator's git
credentials / the GitHub App token present on the trusted host — NEVER a
box-held token. Implemented as a thin shell-out; the heavy lifting is the
pure :func:`build_dispatch_inputs` above, so this stays small.
"""
def _push(
*, owner: str, repo: str, base: str, head_branch: str, diff_text: str
) -> None:
import subprocess
import tempfile
from pathlib import Path
# Operator host only. Use a throwaway worktree, apply the diff, push the
# branch with the host's credentials. Kept intentionally minimal; the
# security comes from CI re-verifying the pushed content by hash.
with tempfile.TemporaryDirectory() as tmp:
tmpdir = Path(tmp)
diff_path = tmpdir / CANDIDATE_DIFF_FILENAME
diff_path.write_text(diff_text, encoding="utf-8")
clone = tmpdir / "repo"
subprocess.run(
[
"gh",
"repo",
"clone",
f"{owner}/{repo}",
str(clone),
"--",
"--depth",
"1",
"--branch",
base,
],
check=True,
capture_output=True,
)
subprocess.run(
["git", "-C", str(clone), "checkout", "-B", head_branch],
check=True,
capture_output=True,
)
# --index applies AND stages exactly the diff's changes (incl. new
# files) and NOTHING else — so the committed head tree is precisely
# base+diff, never stray untracked worktree content. This keeps the
# PR head bound to the same bytes CI hash-verified (LOGIC-1). No
# separate `git add -A` (which would stage unrelated content).
subprocess.run(
["git", "-C", str(clone), "apply", "--index", str(diff_path)],
check=True,
capture_output=True,
)
subprocess.run(
[
"git",
"-C",
str(clone),
"-c",
"user.name=agent-team",
"-c",
"user.email=agent-team@seahavenind.com",
"commit",
"-m",
f"agent-team apply: {head_branch}",
],
check=True,
capture_output=True,
)
subprocess.run(
[
"git",
"-C",
str(clone),
"push",
# --no-verify: skip the operator's LOCAL pre-push dev hook (the
# secrev scanners backstop, which flags pre-existing whole-repo
# findings like the .env.example FP). The apply path's security
# is enforced CI-side — the agent-team-apply-verify workflow
# (guard denylist/scope/hash + the credential-less build-test)
# and the draft PR's own required checks scan the actual
# content. The local human-commit hook is not the apply gate.
"--no-verify",
"--force-with-lease",
"origin",
head_branch,
],
check=True,
capture_output=True,
)
return _push
def _default_workflow_dispatcher() -> WorkflowDispatcher:
"""Real dispatcher: ``gh workflow run`` (operator auth has ``actions:write``)."""
def _fire(*, owner: str, repo: str, inputs: dict[str, str], ref: str) -> None:
import subprocess
args = [
"gh",
"workflow",
"run",
WORKFLOW_FILE,
"--repo",
f"{owner}/{repo}",
"--ref",
ref,
]
for key, value in inputs.items():
args += ["-f", f"{key}={value}"]
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
# 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 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:
import json
import subprocess
import time
from datetime import timedelta
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,status,conclusion",
"--limit",
"50",
],
check=True,
capture_output=True,
text=True,
)
runs = json.loads(proc.stdout or "[]")
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
return _locate