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 61ab1f0cd5 fix(agent-team): dispatch via GitHub App so P3 reaches CI (run_id resolves)
The P3 dispatcher's default seams shell out to gh/git, but the R720 box has
no gh and a read-only PAT with no Actions scope — so dispatch_apply_verify
returned no run_id and every task parked at verify ("dispatch unresolved").

Add a GitHub-App auth path: the box mints short-lived (~1h) installation
access tokens from the App private key and uses them for the three dispatch
seams, removing the gh dependency.

- agent_team/github_app.py (new): mint_installation_token (RS256 App JWT,
  iss=app_id, iat backdated 60s, exp 9 min; POST /access_tokens) + a lazy
  TokenProvider that caches and re-mints near expiry. Secret-safe: the JWT
  and token are never logged, never in an exception message, never persisted.
- dispatcher.py: app_branch_pusher / app_workflow_dispatcher / app_run_locator
  (additive; gh/git _default_* left untouched). Push auth rides a host-scoped
  http.extraHeader via GIT_CONFIG_* env (token never in argv/ps); the REST
  run locator maps id->databaseId / created_at->createdAt into select_run_id
  and surfaces 4xx promptly instead of silently exhausting the poll window.
- coordinator.py: default_dispatch_node_factory binds the App seams when
  AGENT_TEAM_GH_APP_ID / _INSTALLATION_ID / _PRIVATE_KEY are all set; partial
  or unreadable config logs one warning and falls back to gh-default (never
  raises at serve-start).
- requirements.txt: pin PyJWT, cryptography, requests (App seams + CI fetcher).
- DEPLOY-R720.md / README.md: App dispatch config, permission/scope audit,
  env-precedence check, key rotation/revocation + incident response.

Tests: +18 (test_github_app.py new; dispatcher/coordinator additions) covering
JWT claims, cache/re-mint, token-scrub-on-error, REST field mapping + run-name
correlation, and the partial-env inert fallback. Full suite 1523 passing.
2026-06-24 17:07:01 -04:00

826 lines
34 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",
"app_branch_pusher",
"app_run_locator",
"app_workflow_dispatcher",
"build_dispatch_inputs",
"dispatch_apply_verify",
"head_branch_for",
"select_run_id",
]
# GitHub REST API root. The App-token seams (:func:`app_branch_pusher` /
# :func:`app_workflow_dispatcher` / :func:`app_run_locator`) talk to the REST API
# directly with an installation token instead of shelling out to ``gh`` — so the
# box can dispatch with a freshly-minted, short-lived App token and no operator
# ``gh`` auth. Mirrors :data:`agent_team.ci_fetcher.GITHUB_API_ROOT`.
GITHUB_API_ROOT = "https://api.github.com"
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
# ───────────────────────────────────────────────────────────────────────────
# GitHub App-token seams (box-side). Same three behaviours as the ``gh``/``git``
# defaults above, but authenticated with a short-lived installation token from an
# injected ``token_provider`` (``.token() -> str``) instead of the operator's
# ``gh`` auth. This lets the read-only box mint a write token on demand and
# dispatch directly via the REST API + a token-embedding clone URL.
#
# SECRET HYGIENE (BLOCKING): the installation token — and any clone URL that
# embeds it — MUST NEVER be logged, placed in an exception message / ``str()``,
# or written to ledger / graph / task state. The git seam never lets a raw
# ``CalledProcessError`` propagate (its ``.cmd`` / ``.output`` carry the token):
# it scrubs the token + remote URL to ``***`` and re-raises a
# :class:`DispatcherError` with ``from None``. The HTTP seams never include the
# token in any message (status only).
# ───────────────────────────────────────────────────────────────────────────
def app_branch_pusher(token_provider: Any, *, _run: Any = None) -> BranchPusher:
"""App-token :class:`BranchPusher`: clone + apply + push with an installation token.
Mirrors the apply/commit/push semantics of :func:`_default_branch_pusher`
exactly — only the auth + clone differ: this clones the plain HTTPS remote
(``https://github.com/owner/repo.git``) and injects the installation token
minted on demand from ``token_provider.token()`` via an in-memory
``-c http.extraHeader='Authorization: Bearer <token>'`` on the network-facing
git steps (clone + push), rather than embedding the token in the remote URL or
relying on the operator's ``gh``/``git`` credentials. Carrying auth in a
config header keeps the token OUT of the remote URL and therefore out of the
visible process argv / stored ``origin`` — it is never in a position to leak
via ``ps``/``/proc/<pid>/cmdline``. The diff is applied with ``--index``
(stages exactly the diff, nothing stray) so the committed head tree is
precisely base+diff — the same bytes CI hash-verifies.
``_run`` injects the subprocess runner for tests; the default shells out to
``git`` with ``check=True`` + ``capture_output=True``.
SECRET HYGIENE: the token still appears in a git argv element (the
``http.extraHeader`` config value), so every git step is wrapped so a raw
``CalledProcessError`` (whose ``.cmd`` / ``.output`` may carry that header)
NEVER propagates. Any failure is re-raised as a :class:`DispatcherError` whose
message has the token scrubbed to ``***`` (``from None`` so the original —
token-bearing — exception is not chained).
"""
def _default_run(
cmd: list[str], *, cwd: str | None = None, env: dict[str, str] | None = None
) -> None:
import os
import subprocess
full_env = {**os.environ, **env} if env else None
subprocess.run(cmd, cwd=cwd, check=True, capture_output=True, env=full_env)
run = _run if _run is not None else _default_run
def _push(
*, owner: str, repo: str, base: str, head_branch: str, diff_text: str
) -> None:
import tempfile
from pathlib import Path
token = token_provider.token()
# Plain remote — the token is NEVER in the URL/argv/stored origin.
remote = f"https://github.com/{owner}/{repo}.git"
# Auth is injected via http.extraHeader through git's GIT_CONFIG_* env
# vars (NOT argv) on the network steps (clone + push). Env is visible
# only to the same uid via /proc/<pid>/environ, never via ps argv. The
# header key is SCOPED to the github.com host
# (``http.https://github.com/.extraHeader``, the GitHub-Actions checkout
# pattern) so git never sends the Authorization header to any other host
# it might be redirected to.
auth_env = {
"GIT_CONFIG_COUNT": "1",
"GIT_CONFIG_KEY_0": "http.https://github.com/.extraHeader",
"GIT_CONFIG_VALUE_0": f"Authorization: Bearer {token}",
}
def _scrub(s: str) -> str:
return s.replace(token, "***")
with tempfile.TemporaryDirectory() as tmp:
tmpdir = Path(tmp)
diff_path = tmpdir / CANDIDATE_DIFF_FILENAME
diff_path.write_text(diff_text, encoding="utf-8")
clonedir = tmpdir / "repo"
try:
run(
[
"git",
"clone",
"--depth",
"1",
"--branch",
base,
remote,
str(clonedir),
],
env=auth_env,
)
except Exception as exc: # noqa: BLE001 - scrub token before surfacing
raise DispatcherError(f"git clone failed: {_scrub(str(exc))}") from None
try:
run(["git", "-C", str(clonedir), "checkout", "-B", head_branch])
except Exception as exc: # noqa: BLE001
raise DispatcherError(
f"git checkout failed: {_scrub(str(exc))}"
) from None
try:
# --index applies AND stages exactly the diff (incl. new files) and
# NOTHING else, so the head tree is precisely base+diff (LOGIC-1).
run(["git", "-C", str(clonedir), "apply", "--index", str(diff_path)])
except Exception as exc: # noqa: BLE001
raise DispatcherError(f"git apply failed: {_scrub(str(exc))}") from None
try:
run(
[
"git",
"-C",
str(clonedir),
"-c",
"user.name=agent-team",
"-c",
"user.email=agent-team@seahavenind.com",
"commit",
"-m",
f"agent-team apply: {head_branch}",
]
)
except Exception as exc: # noqa: BLE001
raise DispatcherError(
f"git commit failed: {_scrub(str(exc))}"
) from None
try:
run(
[
"git",
"-C",
str(clonedir),
"push",
"--no-verify",
"--force-with-lease",
"origin",
head_branch,
],
env=auth_env,
)
except Exception as exc: # noqa: BLE001
raise DispatcherError(f"git push failed: {_scrub(str(exc))}") from None
return _push
def app_workflow_dispatcher(
token_provider: Any, *, _http: Any = None
) -> WorkflowDispatcher:
"""App-token :class:`WorkflowDispatcher`: trigger the workflow via the REST API.
Same trigger as :func:`_default_workflow_dispatcher` (the apply/verify
``workflow_dispatch``), but issues
``POST /repos/{owner}/{repo}/actions/workflows/{WORKFLOW_FILE}/dispatches``
directly with an installation token from ``token_provider.token()`` instead of
shelling out to ``gh``. ``_http`` injects a ``requests``-like client for tests
(``.post(url, *, json=..., headers=..., timeout=...)`` -> response exposing
``.status_code``); the default lazily imports ``requests``.
GitHub returns ``204`` on success; any other status fails closed with a
:class:`DispatcherError` carrying the STATUS only (NEVER the token).
"""
def _fire(*, owner: str, repo: str, inputs: dict[str, str], ref: str) -> None:
url = (
f"{GITHUB_API_ROOT}/repos/{owner}/{repo}"
f"/actions/workflows/{WORKFLOW_FILE}/dispatches"
)
headers = {
"Accept": "application/vnd.github+json",
"X-GitHub-Api-Version": "2022-11-28",
"Authorization": f"Bearer {token_provider.token()}",
}
body = {"ref": ref, "inputs": inputs}
if _http is not None:
resp = _http.post(url, json=body, headers=headers, timeout=15.0)
else:
import requests # deferred: optional dependency
resp = requests.post(url, json=body, headers=headers, timeout=15.0)
status = getattr(resp, "status_code", None)
if status != 204:
raise DispatcherError(f"workflow dispatch failed: status={status}")
return _fire
def app_run_locator(
token_provider: Any,
*,
_http: Any = None,
_sleep: Any = None,
_now: Any = None,
) -> RunLocator:
"""App-token :class:`RunLocator`: match the triggered run via the REST API.
Same anti-stale SELECTION as :func:`_default_run_locator` — it floors on the
dispatched-at watermark (minus :data:`_LOCATE_SKEW_S` for clock skew),
delegates to the pure :func:`select_run_id`, and polls a bounded number of
times — but lists runs via
``GET /repos/{owner}/{repo}/actions/runs`` with an installation token instead
of ``gh run list``. The REST rows are mapped to the field names
:func:`select_run_id` reads (``id`` -> ``databaseId``, ``created_at`` ->
``createdAt``). Returns ``None`` if no matching run registers within the poll
window (fails closed).
A non-200 list response (e.g. a 401/403/404 auth/scope edge) raises a
:class:`DispatcherError` carrying the STATUS only (NEVER the token) so the
dispatch node parks immediately with an actionable signal, rather than
treating the error body as "no runs" and silently exhausting the poll window.
``_http`` injects a ``requests``-like client (``.get(url, *, params=...,
headers=..., timeout=...)`` -> response exposing ``.json()``); ``_sleep`` /
``_now`` are injected for tests. No token ever appears in a log or message.
"""
sleep = _sleep if _sleep is not None else None
def _locate(*, owner: str, repo: str, task_id: str, since_iso: str) -> str | None:
import time
from datetime import timedelta
do_sleep = sleep if sleep is not None else time.sleep
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
headers = {
"Accept": "application/vnd.github+json",
"X-GitHub-Api-Version": "2022-11-28",
"Authorization": f"Bearer {token_provider.token()}",
}
url = f"{GITHUB_API_ROOT}/repos/{owner}/{repo}/actions/runs"
params = {
"event": "workflow_dispatch",
"created": f">={floor_iso}",
"per_page": 50,
}
for attempt in range(_LOCATE_ATTEMPTS):
if _http is not None:
resp = _http.get(url, params=params, headers=headers, timeout=15.0)
else:
import requests # deferred: optional dependency
resp = requests.get(url, params=params, headers=headers, timeout=15.0)
# Surface auth/4xx promptly with the STATUS only (NEVER the token)
# instead of treating a 401/403/404 error body as "no runs" and
# silently exhausting the ~60s poll window. A genuine 200 with no
# matching run still falls through to the None-on-no-match path below.
status = getattr(resp, "status_code", 200)
if status != 200:
raise DispatcherError(f"run list failed: status={status}")
body = resp.json()
runs_raw = body.get("workflow_runs") or []
mapped = [
{
"databaseId": r.get("id"),
"name": r.get("name"),
"createdAt": r.get("created_at"),
"status": r.get("status"),
"conclusion": r.get("conclusion"),
}
for r in runs_raw
]
run_id = select_run_id(mapped, task_id=task_id, floor_iso=floor_iso)
if run_id is not None:
return run_id
if attempt < _LOCATE_ATTEMPTS - 1:
do_sleep(_LOCATE_DELAY_S)
return None
return _locate