feat(agent-team): P4 live github/claude_code transports + GitHub-issue intake

Live github (issue-comment poster) and claude_code (file-drop) transports, plus GithubIntake (labeled issue -> coordinator.start_task, de-duped). run-team _build_transport now wires github/claude_code live (was SystemExit) + adds the intake-github subcommand. claude_code drop-path also neutralizes backslash (defense-in-depth).
This commit is contained in:
Adam Moussa 2026-06-18 13:23:02 -04:00
parent 8d9babe52e
commit 0842ff778b
8 changed files with 1935 additions and 14 deletions

View file

@ -0,0 +1,265 @@
"""Live Claude-Code-on-the-Mac delivery wiring (design §3.3.1, §7.1 P4, D10).
The :mod:`agent_team.transport.claude_code_adapter` module ships the §3.3.1
transport contract with a dependency-injected ``delivery`` seam: the adapter
renders the question-set into a prompt body and hands it to a sink called as
``delivery(session_hint=..., prompt=...) -> str`` whose job is to surface the
prompt inside a Claude-Code session and return the Claude **session id** used as
the ``channel_ref`` locator. The foundation's default sink refuses to act so
nothing ships provisioned; this module supplies the **production** sink, backed
by a local **file drop** the Mac harness reads, that the P4 (Claude-Code) live
wiring injects.
Channel model (file drop)
-------------------------
The R720 box has no standing write path into Adam's interactive Claude-Code
session, so delivery is **SSH-invoked from the Mac side** (D10). The live sink
writes the rendered prompt to a drop directory on the Mac filesystem; the
Claude-Code harness polls that directory, surfaces the prompt inline, and Adam
answers it. The sink returns the Claude **session id** (the file stem, derived
from the embedded ``question_id``) which §3.3.1 names as this transport's
``channel_ref`` locator. The answer travels back the same way: the harness drops
an answer file the box reads and feeds to
:meth:`ClaudeCodeAdapter.parse_answer`.
Deferred / injected I/O (mirrors :func:`agent_team.graph.build_sqlite_checkpointer`)
-----------------------------------------------------------------------------------
All filesystem access is injected or deferred so this module imports cleanly in
pre-deploy / test environments and is unit-testable with fakes:
* the prompt-writer and answer-reader are injected callables (``writer`` /
``reader``) defaulting to thin wrappers over the local filesystem; and
* the standard-library ``pathlib`` import is the only hard dependency — there is
no SDK, network, SSH, or secret touched here. A missing / unwritable drop
directory raises a clear :class:`ClaudeCodeDeliveryError` so a misconfigured
deploy fails loudly rather than silently reporting a delivery that never
reached Adam.
Scope (P2 default; P3 inert)
----------------------------
This is the human-in-the-loop *question delivery* path used by the production
P2 (clarify -> plan -> review) pipeline. It performs no CI, OIDC, git/patch
apply, or network calls; the P3 builder/verifier apply-verify workflow is held
for a separate review gate and is neither shipped nor enabled here.
"""
from __future__ import annotations
import os
from pathlib import Path
from typing import Any, Callable
from agent_team.transport.claude_code_adapter import (
ClaudeCodeAdapter,
ClaudeCodeDeliveryError,
parse_channel_ref,
)
__all__ = [
"PROMPT_SUFFIX",
"build_claude_code_delivery",
"build_file_drop_reader",
"build_file_drop_writer",
"build_live_claude_code_transport",
"read_answer_payload",
]
# File extension for a dropped prompt file. The Mac harness globs this suffix in
# the drop directory to discover pending questions.
PROMPT_SUFFIX = ".prompt.txt"
# File extension for a dropped answer file written back by the Mac harness. The
# box globs this suffix to discover answered questions.
ANSWER_SUFFIX = ".answer.txt"
# Type of the injected prompt-writer: surfaces ``prompt`` at ``path``.
PromptWriter = Callable[[Path, str], None]
# Type of the injected answer-reader: returns the answer text at ``path`` (or
# ``None`` if the harness has not dropped an answer yet).
AnswerReader = Callable[[Path], "str | None"]
def _drop_path(drop_dir: Path, session_id: str, suffix: str) -> Path:
"""Resolve the on-disk path for a drop file.
The ``session_id`` is the file stem, so a prompt and its answer share a stem
and differ only by suffix. ``session_id`` is sanitized to a single path
component (no separators) so a hostile/garbled id cannot escape ``drop_dir``.
"""
# Replace os.sep + os.altsep AND an explicit backslash: on POSIX os.altsep
# is None, so a literal backslash would otherwise survive (defense-in-depth,
# even though stripping "/" already prevents traversal on the deploy targets).
safe = session_id.strip().replace(os.sep, "_").replace("\\", "_")
if os.altsep:
safe = safe.replace(os.altsep, "_")
safe = safe.lstrip(".") or "session"
return drop_dir / f"{safe}{suffix}"
def build_file_drop_writer(drop_dir: Path | str) -> PromptWriter:
"""Build the default filesystem prompt-writer for the live sink.
The returned callable writes ``prompt`` to ``path`` (UTF-8), creating the
drop directory if needed. Filesystem access is deferred to call time so this
factory is side-effect-free at import. A write failure (unwritable / missing
parent) propagates so the caller can treat the post as failed.
"""
base = Path(drop_dir)
def _writer(path: Path, prompt: str) -> None:
base.mkdir(parents=True, exist_ok=True)
path.write_text(prompt, encoding="utf-8")
return _writer
def build_file_drop_reader(drop_dir: Path | str) -> AnswerReader:
"""Build the default filesystem answer-reader for the inbound path.
The returned callable reads the answer text at ``path`` (UTF-8), returning
``None`` when the harness has not yet dropped an answer file. Filesystem
access is deferred to call time.
"""
Path(drop_dir) # validate / normalize eagerly; read happens at call time.
def _reader(path: Path) -> str | None:
try:
return path.read_text(encoding="utf-8")
except FileNotFoundError:
return None
return _reader
def build_claude_code_delivery(
drop_dir: Path | str,
*,
writer: PromptWriter | None = None,
) -> Callable[..., str]:
"""Build a live file-drop ``delivery`` sink (§3.3.1, P4, D10).
The returned callable matches the adapter's injected-sink contract,
``delivery(session_hint=..., prompt=...) -> str``: it writes the rendered
``prompt`` to a file in ``drop_dir`` for the Mac harness to read and returns
the Claude **session id** the adapter folds into the ``channel_ref``.
The session id is derived from the ``question_id`` embedded in ``prompt`` via
the shared ``<!-- shq:<question_id> -->`` marker (so the prompt and its
answer share a deterministic file stem), falling back to ``session_hint``
when present. A prompt with neither an embedded marker nor a ``session_hint``
cannot be addressed, so the sink raises :class:`ClaudeCodeDeliveryError`
rather than dropping an unaddressable file.
``writer`` (optional) injects the prompt-writer for testability; any
``Callable[[Path, str], None]`` works. When omitted, a filesystem writer over
``drop_dir`` is built via :func:`build_file_drop_writer`. A write failure is
wrapped in :class:`ClaudeCodeDeliveryError` so a failed drop leaves the
ledger row ``open`` (no ``channel_ref``) for idempotent reconcile retry,
mirroring a failed Slack post.
"""
base = Path(drop_dir)
write = writer if writer is not None else build_file_drop_writer(base)
def _delivery(*, session_hint: str, prompt: str) -> str:
session_id = _session_id_for(prompt=prompt, session_hint=session_hint)
path = _drop_path(base, session_id, PROMPT_SUFFIX)
try:
write(path, prompt)
except OSError as exc:
raise ClaudeCodeDeliveryError(
f"failed to write Claude-Code prompt drop at {path!s}: {exc}"
) from exc
return session_id
return _delivery
def build_live_claude_code_transport(
drop_dir: Path | str,
*,
session_hint: str = "",
writer: PromptWriter | None = None,
) -> ClaudeCodeAdapter:
"""Build a :class:`ClaudeCodeAdapter` wired to a live file-drop sink.
Convenience constructor for the P4 live coordinator: equivalent to
``ClaudeCodeAdapter(build_claude_code_delivery(drop_dir, writer=writer),
session_hint=session_hint)``. See :func:`build_claude_code_delivery` for the
drop-directory / writer / session-id semantics.
"""
return ClaudeCodeAdapter(
build_claude_code_delivery(drop_dir, writer=writer),
session_hint=session_hint,
)
def read_answer_payload(
channel_ref: str,
drop_dir: Path | str,
*,
reader: AnswerReader | None = None,
) -> dict[str, Any] | None:
"""Read a dropped answer for ``channel_ref`` into a ``parse_answer`` payload.
Resolves the answer file (the prompt's stem + :data:`ANSWER_SUFFIX`) for the
Claude session in ``channel_ref``, reads it via the injected ``reader``
(defaulting to a filesystem reader over ``drop_dir``), and returns a payload
dict ready for :meth:`ClaudeCodeAdapter.parse_answer`. The payload echoes the
original ``channel_ref`` so the ``question_id`` round-trips from the ref
alone even if the answer text carries no marker.
Returns ``None`` when no answer has been dropped yet (the harness has not
answered), so the reconcile loop can poll idempotently. Raises
:class:`ValueError` for a ``channel_ref`` that is not a Claude-Code ref, so a
caller cannot silently read the wrong transport's drop.
"""
parsed = parse_channel_ref(channel_ref)
if parsed is None:
raise ValueError(
f"not a Claude-Code channel_ref; cannot read answer drop: {channel_ref!r}"
)
session_id, _question_id = parsed
base = Path(drop_dir)
read = reader if reader is not None else build_file_drop_reader(base)
path = _drop_path(base, session_id, ANSWER_SUFFIX)
answer_text = read(path)
if answer_text is None:
return None
return {"channel_ref": channel_ref, "answer": answer_text}
def _session_id_for(*, prompt: str, session_hint: str) -> str:
"""Derive the file-stem session id for a prompt drop.
Prefers the ``question_id`` embedded in ``prompt`` via the shared marker (so
the prompt and its answer share a deterministic stem), then a non-blank
``session_hint``. Raises :class:`ClaudeCodeDeliveryError` when neither is
available, since an unaddressable drop would be unrecoverable.
"""
embedded = _embedded_question_id(prompt)
if embedded:
return embedded
hint = session_hint.strip()
if hint:
return hint
raise ClaudeCodeDeliveryError(
"cannot address a Claude-Code prompt drop: no embedded question_id marker "
"and no session_hint"
)
def _embedded_question_id(prompt: str) -> str | None:
"""Recover the ``question_id`` embedded in a rendered prompt body, or ``None``.
Delegates to the adapter's own marker parsing through a throwaway payload so
the marker grammar stays single-sourced in
:mod:`agent_team.transport.claude_code_adapter`.
"""
adapter = ClaudeCodeAdapter()
try:
question_id, _answer, _via = adapter.parse_answer({"prompt": prompt})
except ValueError:
return None
return question_id

View file

@ -0,0 +1,293 @@
"""GitHub-issue INTAKE poller: a labeled issue becomes a pipeline task (§3.3.1).
This is the *inbound front door* for the GitHub transport. Where
:mod:`agent_team.transport.github_adapter` delivers clarifier question-sets
*outbound* (and parses answers back), this leaf runs the other direction: it
polls a repository for open issues carrying a configured label and turns each
not-yet-ingested issue into one pipeline task by calling the coordinator's
intake entry,
:meth:`agent_team.coordinator.Coordinator.start_task` (``task_text=<issue
title+body>``, ``transport_name="github"``).
The shape mirrors the slack_listener seam: everything network/SDK is
**injected** so the poller is fully unit-testable with no GitHub SDK and no
socket:
* ``client`` is a small :class:`GithubIssueClient` protocol:
``list_open_issues(label) -> iterable of issue mappings``. Production wires a
thin client over the GitHub REST API (deferred import, see
:func:`build_default_issue_client`); tests pass an in-memory fake.
* ``coordinator`` is anything exposing ``start_task(task_text=...,
transport_name=...)``: the live :class:`~agent_team.coordinator.Coordinator`
in production, a stub in tests. No model or transport is touched here.
De-duplication (P3 scope note):
The poller tracks already-ingested issue ids in an **in-memory** set, so a
re-poll over the same open issue does not start a second task. This is
deliberately simple for now: it does NOT survive a process restart. Durable
de-dup (a ledger table of ingested issue ids, mirroring the
``pending_questions`` discipline) is a FOLLOW-UP and is intentionally not
shipped here. After a restart an already-ingested-but-still-open issue would
be re-ingested; document that and treat the in-memory set as a best-effort
guard, not a durable contract.
Design constraints honoured here (pre-deployment scaffolding):
* **No live infrastructure.** Nothing is provisioned or called at import.
The GitHub client is dependency-injected; the default client's SDK/HTTP
import is DEFERRED (mirrors
:func:`agent_team.graph.build_sqlite_checkpointer` and the slack_listener
SDK discipline), so this module imports cleanly with no optional SDK
present and the unit tests stay fully hermetic.
* **P2 stays the production default; this is OPT-IN and INERT.** This module
does no CI, no OIDC, no git/patch apply, and no network to GitHub Actions.
It only reads issues and calls the coordinator's existing intake entry.
* **Secrets never committed.** The default client reads the GitHub token
from the environment at call time, never from source.
"""
from __future__ import annotations
import logging
from typing import Any, Iterable, Protocol
__all__ = [
"GithubIntake",
"GithubIssueClient",
"build_default_issue_client",
"issue_task_text",
]
_LOG = logging.getLogger(__name__)
# The transport name handed to the coordinator's intake entry so the resulting
# task's clarifier question-sets route over the GitHub adapter (§3.3.1 D10).
GITHUB_TRANSPORT_NAME = "github"
# Environment variable the default client reads the GitHub token from at call
# time (never stored in source/state). Mirrors the github_adapter default.
DEFAULT_TOKEN_ENV = "GITHUB_TOKEN"
class GithubIssueClient(Protocol):
"""Injected GitHub issue source: list open issues carrying a label.
A narrow read-only seam so the poller has no hard dependency on any GitHub
SDK and the tests pass a pure in-memory fake. Each returned issue is a
mapping with at least an ``id`` (or ``number``) and a ``title``; ``body`` is
optional. The mapping shape mirrors the GitHub REST issue object so the
production client can return the API JSON unchanged.
"""
def list_open_issues(self, *, label: str) -> Iterable[dict[str, Any]]:
"""Return the open issues carrying ``label`` (most-recent-first is fine)."""
...
def issue_task_text(issue: dict[str, Any]) -> str:
"""Render one issue's intake ``task_text`` from its title + body.
The pipeline's task description is the issue title followed by its body (a
blank line between them when both are present). A missing/empty body yields
just the title; a missing/empty title falls back to ``issue #<id>`` so the
task is never an empty string. Whitespace is stripped at the edges so a
trailing-newline body does not produce trailing blank lines.
"""
title = str(issue.get("title") or "").strip()
body = str(issue.get("body") or "").strip()
if not title:
title = f"issue #{_issue_id(issue)}"
if body:
return f"{title}\n\n{body}"
return title
def _issue_id(issue: dict[str, Any]) -> str:
"""Return the de-dup identity for ``issue`` as a string.
Prefers the GitHub global ``id`` (stable across renames); falls back to the
per-repo ``number`` when ``id`` is absent (some payload shapes / fakes carry
only ``number``). Stringified so heterogeneous int/str ids compare cleanly
in the ingested set.
"""
raw = issue.get("id")
if raw is None:
raw = issue.get("number")
return str(raw)
class GithubIntake:
"""Poll a repo for labeled issues and start one pipeline task per new issue.
Construct with an injected ``client`` (a :class:`GithubIssueClient`), an
injected ``coordinator`` (anything exposing
``start_task(task_text=..., transport_name=...)``), and the ``label`` that
flags an issue as pipeline intake. Call :meth:`poll_once` on a cadence (an
operator loop or a cron); each call lists the open labeled issues and starts
a task for every one not yet ingested.
De-dup is in-memory only (see the module docstring): the set of ingested
issue ids lives on the instance, so a re-poll within one process never
double-ingests, but a restart loses the set. Durable de-dup is a follow-up.
Nothing here touches the network or any SDK directly (the client does, and
it is injected), so the whole poller is unit-testable with a fake client and
a stub coordinator.
"""
def __init__(
self,
*,
client: GithubIssueClient,
coordinator: Any,
label: str,
) -> None:
"""Bind the poller to one client, coordinator, and intake label.
Args:
client: The injected issue source. Its ``list_open_issues`` is the
only GitHub call the poller makes.
coordinator: The intake target. Must expose
``start_task(task_text=..., transport_name=...)``: the live
:class:`~agent_team.coordinator.Coordinator` in production.
label: The issue label that marks an issue as pipeline intake. Only
issues the client returns for this label are considered; an
empty label is rejected so a misconfiguration cannot ingest
every open issue.
"""
if not label:
raise ValueError(
"GithubIntake requires a non-empty intake label; an empty label "
"would ingest every open issue"
)
self._client = client
self._coordinator = coordinator
self._label = label
# In-memory de-dup set (P3 scope: best-effort, NOT durable across a
# restart; see the module docstring). Tracks issue ids already turned
# into tasks so a re-poll does not double-ingest.
self._ingested: set[str] = set()
@property
def label(self) -> str:
"""The configured intake label (read-only)."""
return self._label
@property
def ingested_ids(self) -> frozenset[str]:
"""A snapshot of the issue ids already ingested this process (read-only)."""
return frozenset(self._ingested)
def poll_once(self) -> list[str]:
"""List the labeled open issues and start a task for each new one.
One maintenance pass:
1. Ask the injected client for the open issues carrying the configured
label (:meth:`GithubIssueClient.list_open_issues`).
2. For each issue NOT already in the in-memory ingested set, call
``coordinator.start_task(task_text=<title+body>,
transport_name="github")`` and record its id so a subsequent poll
does not re-ingest it.
Issues already ingested this process are skipped (the in-memory de-dup),
and any issue the client returns without the label is *not* expected
(the client filters by label) but is ignored defensively if present.
An issue id is recorded as ingested ONLY after ``start_task`` returns,
so a failing intake leaves the issue eligible for retry on the next poll
rather than silently dropping it.
Returns the list of issue ids ingested on THIS pass (empty when nothing
new), so an operator loop can log/meter intake volume.
"""
ingested_now: list[str] = []
for issue in self._client.list_open_issues(label=self._label):
issue_id = _issue_id(issue)
if issue_id in self._ingested:
_LOG.debug("github-intake: issue %s already ingested; skip", issue_id)
continue
task_text = issue_task_text(issue)
_LOG.info(
"github-intake: starting task for issue %s (label=%s)",
issue_id,
self._label,
)
# start_task is the committed coordinator intake entry; the resulting
# task's clarifier question-sets route over the GitHub adapter. Record
# the id only after the call returns so a raise leaves the issue
# eligible for retry on the next poll (no silent drop).
self._coordinator.start_task(
task_text=task_text,
transport_name=GITHUB_TRANSPORT_NAME,
)
self._ingested.add(issue_id)
ingested_now.append(issue_id)
return ingested_now
def build_default_issue_client(
*,
owner: str,
repo: str,
token_env: str = DEFAULT_TOKEN_ENV,
api_root: str = "https://api.github.com",
) -> GithubIssueClient:
"""Build the production read-only issue client (deferred SDK/HTTP import).
Returns a :class:`GithubIssueClient` that lists a repo's open issues by
label over the GitHub REST API. The HTTP machinery (``urllib``) and the
token read are deferred to call time (mirroring
:func:`agent_team.graph.build_sqlite_checkpointer` and the slack_listener
SDK discipline), so importing this module never touches the network and the
unit tests (which inject a fake client) never reach this path.
The token is read from ``token_env`` at call time and sent as a bearer
credential; it is never stored in source or logged. ``api_root`` is
overridable for GitHub Enterprise.
This is intentionally a thin, read-only lister: it issues a single GET to
the issues endpoint with ``state=open&labels=<label>`` and returns the
parsed JSON array unchanged (each element is a GitHub issue object, which
already carries ``id`` / ``number`` / ``title`` / ``body``). It performs no
CI, OIDC, write, or GitHub-Actions call; it only reads issues.
"""
class _RestIssueClient:
"""Stdlib-only GitHub REST issue lister (built lazily, no import-time HTTP)."""
def __init__(self) -> None:
self._owner = owner
self._repo = repo
self._token_env = token_env
self._api_root = api_root.rstrip("/")
def list_open_issues(self, *, label: str) -> Iterable[dict[str, Any]]:
import json
import os
from urllib import parse as _urlparse
from urllib import request as _urlrequest
token = os.environ.get(self._token_env)
if not token:
raise RuntimeError(
f"no GitHub token available (env {self._token_env!r} unset); "
"cannot list issues for intake"
)
query = _urlparse.urlencode({"state": "open", "labels": label})
url = f"{self._api_root}/repos/{self._owner}/{self._repo}/issues?{query}"
request = _urlrequest.Request(url, method="GET")
request.add_header("Authorization", f"Bearer {token}")
request.add_header("Accept", "application/vnd.github+json")
request.add_header("X-GitHub-Api-Version", "2022-11-28")
with _urlrequest.urlopen(request) as response: # noqa: S310 (trusted api host)
raw = response.read().decode("utf-8")
data = json.loads(raw) if raw else []
# The issues endpoint can include pull requests (they share the
# endpoint); filter them out so a PR is never intaken as an issue.
return [item for item in data if "pull_request" not in item]
return _RestIssueClient()

View file

@ -0,0 +1,229 @@
"""Live ``requests``-backed GitHub poster (design §3.3.1, §7.1 P4 — GitHub).
The :mod:`agent_team.transport.github_adapter` module ships the §3.3.1 transport
contract with a dependency-injected ``http_post`` seam: the adapter renders the
question-set into a Markdown comment body (embedding the
``<!-- shq:<question_id> -->`` marker so an inbound answer maps back) and hands
the REST POST to an
``HttpPost = (url, *, headers, json_body) -> (status, data)`` whose job is to
perform the real ``POST /repos/{owner}/{repo}/issues/{n}/comments`` and return
the GitHub comment payload carrying the new comment ``id``. The foundation's
default ``http_post`` is a stdlib-only (``urllib``) poster invoked only on an
actual delivery, so nothing ships provisioned; this module supplies the
**production** poster, backed by a thin ``requests`` session, that the P4 (GitHub)
live wiring injects.
Why a thin ``requests`` shim (and not PyGithub):
The adapter already speaks the GitHub REST API directly — it builds the
comments URL, the ``Authorization: Bearer`` / ``X-GitHub-Api-Version``
headers, and the ``{"body": ...}`` JSON itself, then hands a plain
``(url, headers, json_body)`` POST to the seam. The live poster therefore
only needs a minimal HTTP client, not the full PyGithub object model. A
``requests`` session keeps the shim small and fully mirrors the seam the
adapter already accepts.
Deferred import (mirrors :func:`agent_team.graph.build_sqlite_checkpointer` and
:func:`agent_team.transport.slack_live.build_slack_poster`):
``requests`` is an optional dependency that may be absent in pre-deploy / test
environments, so this module imports cleanly without it. The import is deferred
to the moment a live client is actually constructed, and a missing package
raises a clear :class:`RuntimeError` so a misconfigured deploy fails loudly
rather than silently. The token is likewise resolved at build time (falling back
to ``GITHUB_TOKEN``); a missing token raises a clear :class:`RuntimeError`.
The marker round-trip is owned by the adapter, not the poster: the poster is the
pure network seam. :meth:`GitHubTransport.post_question` embeds the
``question_id`` marker in the comment body it hands to this poster, and
:meth:`GitHubTransport.parse_answer` recovers it from an inbound reply, so the
``question_id`` survives end-to-end without the poster needing to know about it.
Production-default safety (P2 stays clarify->plan->review): this module is live
**transport I/O only** for the human gate. It performs exactly one operation —
post an issue/PR comment — and contains NO CI calls, NO OIDC, NO git/patch
apply, and NO GitHub Actions network. The P3 builders/verifier live apply/verify
workflow is held for a separate review gate and is deliberately absent here.
"""
from __future__ import annotations
import os
from typing import Any
from agent_team.transport.github_adapter import (
GITHUB_API_ROOT,
GitHubApiError,
GitHubTransport,
HttpPost,
)
__all__ = [
"build_github_poster",
"build_live_github_transport",
]
def build_github_poster(token: str | None = None, *, client: Any = None) -> HttpPost:
"""Build a live ``requests``-backed :data:`HttpPost` seam (§3.3.1, P4).
The returned callable implements the adapter's
``(url, *, headers, json_body) -> (status, data)`` seam: it performs the
real ``POST`` against the GitHub REST API and returns the response status
plus parsed JSON body so :meth:`GitHubTransport.post_question` can read the
new comment ``id`` as the ``channel_ref`` the ledger stores. A non-2xx
response is surfaced as :class:`GitHubApiError` (the adapter also guards the
status, but the poster raises early so a failed POST never looks like a
success with an empty body).
``client`` (optional) injects a pre-built HTTP client for testability; any
object exposing ``post(url, *, headers, json)`` and returning a response
with ``status_code`` plus a ``json()`` method works (the ``requests``
``Session`` shape). When omitted, a ``requests.Session`` is constructed
lazily from ``token`` (falling back to the ``GITHUB_TOKEN`` environment
variable). The ``requests`` import is deferred so this module imports
cleanly without the optional package; a missing package or a missing token
raises a clear :class:`RuntimeError`.
The poster is the pure network seam: the ``question_id`` marker is embedded
by the adapter in the comment body it passes through ``json_body``, so the
posted comment carries the marker without the poster handling it. See the
module docstring for the full rationale.
"""
if client is None:
client = _build_session(token)
def _poster(
url: str,
*,
headers: dict[str, str],
json_body: dict[str, Any],
) -> tuple[int, dict[str, Any]]:
response = client.post(url, headers=headers, json=json_body)
status = _status_of(response)
data = _json_of(response)
if not (200 <= status < 300):
raise GitHubApiError(status, _body_repr(data))
return status, data
return _poster
def build_live_github_transport(
*,
owner: str,
repo: str,
issue_number: int,
token: str | None = None,
client: Any = None,
api_root: str = GITHUB_API_ROOT,
) -> GitHubTransport:
"""Build a :class:`GitHubTransport` wired to a live ``requests`` poster.
Convenience constructor for the P4 live coordinator: builds the live poster
and binds it to the issue (or PR) thread where the human gate posts its
question-set comment and reads the reply. See :func:`build_github_poster`
for the token / client / deferred-import semantics.
The transport builds the ``Authorization`` header itself (from a
``token_provider``), so the resolved token is threaded in once and shared
with the poster's session: a single token source backs the whole live path.
When neither ``token`` nor ``client`` is given, the token still resolves from
``GITHUB_TOKEN`` at call time so the deploy fails loudly if it is unset.
``api_root`` is forwarded to the transport so a GitHub Enterprise host can be
targeted; the poster itself is endpoint-agnostic (the adapter builds the URL).
"""
resolved = _resolve_token(token)
return GitHubTransport(
owner=owner,
repo=repo,
issue_number=issue_number,
http_post=build_github_poster(resolved, client=client),
api_root=api_root,
token_provider=(lambda: resolved) if resolved is not None else None,
)
def _build_session(token: str | None) -> Any:
"""Lazily construct a ``requests.Session`` (deferred optional import).
Raises a clear :class:`RuntimeError` if ``requests`` is not installed or no
token is resolvable (neither ``token`` nor ``GITHUB_TOKEN``), so a
misconfigured deploy fails loudly rather than silently. The token is held on
the session purely so the same authenticated client can be reused; the
adapter still sets the ``Authorization`` header per call.
"""
try:
import requests
except ImportError as exc: # pragma: no cover - depends on optional dep
raise RuntimeError(
"requests is unavailable; install the 'requests' package to build "
"a live GitHub poster (P4), or inject a 'client' for testing."
) from exc
resolved = _resolve_token(token)
if not resolved:
raise RuntimeError(
"No GitHub token available; pass 'token' or set the GITHUB_TOKEN "
"environment variable to build a live GitHub poster."
)
session = requests.Session()
session.headers.update({"Authorization": f"Bearer {resolved}"})
return session
def _resolve_token(token: str | None) -> str | None:
"""Resolve the GitHub token, falling back to ``GITHUB_TOKEN`` (read at call time).
Returns ``None`` when no token is available so callers can decide whether a
missing token is fatal (the poster path) or deferrable (an injected-client
test path that never needs one).
"""
return token or os.environ.get("GITHUB_TOKEN")
def _status_of(response: Any) -> int:
"""Read the HTTP status code from a ``requests``-like response.
``requests.Response`` exposes ``status_code``; a test double may instead
expose ``status``. Either is accepted so the poster stays usable with a
minimal fake.
"""
for attr in ("status_code", "status"):
value = getattr(response, attr, None)
if value is not None:
return int(value)
raise TypeError(
"GitHub response exposes no 'status_code'/'status'; expected a "
f"requests-like response, got {type(response)!r}"
)
def _json_of(response: Any) -> dict[str, Any]:
"""Parse the JSON body of a ``requests``-like response to a dict.
A successful create returns the new comment object (carrying ``id``); an
empty body coerces to ``{}`` so the adapter's missing-id guard fires with a
clear message rather than an attribute error.
"""
parser = getattr(response, "json", None)
if not callable(parser):
raise TypeError(
"GitHub response exposes no callable 'json()'; expected a "
f"requests-like response, got {type(response)!r}"
)
data = parser()
return data if isinstance(data, dict) else {}
def _body_repr(data: dict[str, Any]) -> str:
"""Render a response body for a :class:`GitHubApiError` message.
Best-effort JSON; falls back to ``repr`` so the error is always constructible
even for an exotic body.
"""
try:
import json
return json.dumps(data)
except (TypeError, ValueError): # pragma: no cover - exotic body
return repr(data)

View file

@ -38,6 +38,9 @@ Subcommands (P1 surface):
``--confirm``; audit-logged). Maps to :func:`answer_question`.
* ``supersede`` — mark a stale question ``superseded`` (DESTRUCTIVE: requires
``--confirm``; audit-logged). Maps to :func:`supersede_question`.
* ``intake-github`` — poll a repo for labeled open issues and start one pipeline
task per not-yet-ingested issue (one pass). Reads issues + calls the committed
coordinator intake entry only; no CI, OIDC, or git/patch apply. Opt-in/inert.
Exit codes: ``0`` success, ``1`` operational failure (e.g. row not found, the
compare-and-set lost the race), ``2`` usage error (argparse).
@ -74,9 +77,12 @@ from agent_team.db.schema import ( # noqa: E402 (path bootstrap must precede)
)
from agent_team.transport.base import Transport # noqa: E402 (path bootstrap)
# Transport choices the start/serve commands accept (§3.3.1 D10). Only ``slack``
# has a live adapter wired for the P1 CLI; the others are accepted for forward
# compatibility and gated in _build_transport.
# Transport choices the start/serve/intake commands accept (§3.3.1 D10). All
# three now have a live human-gate adapter wired in _build_transport: ``slack``
# (slack_sdk), ``github`` (requests issue/PR comment), and ``claude_code`` (local
# file drop). Each builder defers its optional SDK / token resolution to call
# time, so a missing dep/credential fails loudly only when that transport is
# actually selected.
_TRANSPORT_CHOICES: tuple[str, ...] = ("slack", "github", "claude_code")
__all__ = [
@ -511,10 +517,23 @@ def _build_transport(args: argparse.Namespace) -> Any:
"""Build the transport for a coordinator command (lazy; token-tolerant).
``--dry-run`` (or any transport in dry-run) yields a non-posting transport so
intake works without credentials. Otherwise the live Slack transport is
constructed lazily from ``SLACK_BOT_TOKEN`` / ``SLACK_CHANNEL``; GitHub and
Claude-Code live transports are not wired for the P1 CLI surface and raise a
clear error rather than pretending to post.
intake works without credentials. Otherwise the live transport is built
lazily from the environment for the chosen ``--transport`` (so import,
``--help``, and ledger commands never need a token):
* ``slack`` -> :func:`build_live_slack_transport` over ``SLACK_BOT_TOKEN`` /
``SLACK_CHANNEL``.
* ``github`` -> :func:`build_live_github_transport` over ``GITHUB_TOKEN`` and
the issue thread ``GITHUB_OWNER`` / ``GITHUB_REPO`` /
``GITHUB_ISSUE_NUMBER``. This is the §3.3.1 human-gate I/O only (post an
issue/PR comment); it carries NO CI, OIDC, git/patch apply, or GitHub
Actions network (the P3 apply/verify workflow is held for the security
gate).
* ``claude_code`` -> :func:`build_live_claude_code_transport` over the local
file-drop directory ``CLAUDE_CODE_DROP_DIR`` the Mac harness polls (D10).
Each live builder defers its optional SDK / token resolution to call time, so
a missing dependency or credential fails loudly here rather than at import.
"""
if getattr(args, "dry_run", False):
return _DryRunTransport()
@ -523,12 +542,74 @@ def _build_transport(args: argparse.Namespace) -> Any:
channel = os.environ.get("SLACK_CHANNEL", "")
return build_live_slack_transport(channel)
if args.transport == "github":
from agent_team.transport.github_live import build_live_github_transport
owner, repo, issue_number = _github_thread_from_env()
# Token resolves from GITHUB_TOKEN inside the builder (call-time read);
# a missing token fails loudly there rather than being captured here.
return build_live_github_transport(
owner=owner,
repo=repo,
issue_number=issue_number,
token=os.environ.get("GITHUB_TOKEN") or None,
)
if args.transport == "claude_code":
from agent_team.transport.claude_code_live import (
build_live_claude_code_transport,
)
drop_dir = os.environ.get("CLAUDE_CODE_DROP_DIR", "")
if not drop_dir:
raise SystemExit(
"live transport 'claude_code' requires CLAUDE_CODE_DROP_DIR (the "
"Mac file-drop directory the Claude-Code harness polls); set it, "
"or use --dry-run for a no-token dry run"
)
return build_live_claude_code_transport(drop_dir)
raise SystemExit(
f"live transport '{args.transport}' is not wired for the run-team CLI; "
"use --transport slack, or --dry-run for a no-token dry run"
"use --transport slack/github/claude_code, or --dry-run for a no-token "
"dry run"
)
def _github_thread_from_env() -> tuple[str, str, int]:
"""Resolve the GitHub issue thread (owner/repo/issue) from the environment.
The live GitHub transport posts the clarifier question-set as a comment on a
fixed ``owner/repo#issue_number`` thread, so the thread is configured via
``GITHUB_OWNER`` / ``GITHUB_REPO`` / ``GITHUB_ISSUE_NUMBER`` (read lazily so
the value is never captured at import). A missing or non-integer value raises
a clear :class:`SystemExit` rather than building a half-configured transport.
"""
owner = os.environ.get("GITHUB_OWNER", "")
repo = os.environ.get("GITHUB_REPO", "")
raw_issue = os.environ.get("GITHUB_ISSUE_NUMBER", "")
missing = [
name
for name, value in (
("GITHUB_OWNER", owner),
("GITHUB_REPO", repo),
("GITHUB_ISSUE_NUMBER", raw_issue),
)
if not value
]
if missing:
raise SystemExit(
"live transport 'github' requires "
f"{', '.join(missing)}; set the issue thread (owner/repo/issue), or "
"use --dry-run for a no-token dry run"
)
try:
issue_number = int(raw_issue)
except ValueError:
raise SystemExit(
f"GITHUB_ISSUE_NUMBER must be an integer, got {raw_issue!r}"
) from None
return owner, repo, issue_number
class _DryRunTransport(Transport):
"""A non-posting transport for ``--dry-run`` intake (no token, no Slack).
@ -591,6 +672,46 @@ def _cmd_serve(args: argparse.Namespace, *, out: Any) -> int:
return 0 # pragma: no cover - serve() loops until interrupted
def _cmd_intake_github(args: argparse.Namespace, *, out: Any) -> int:
"""Poll a repo for labeled issues and start one task per new issue (§3.3.1).
The GitHub-issue intake front door: builds a :class:`Coordinator` (transport
from the lazy factory; ``--dry-run`` posts nowhere), runs ``setup``, then
constructs a :class:`agent_team.transport.github_intake.GithubIntake` over a
read-only REST issue client
(:func:`~agent_team.transport.github_intake.build_default_issue_client`,
deferred-import, reads ``GITHUB_TOKEN`` at call time) and runs ONE
:meth:`~agent_team.transport.github_intake.GithubIntake.poll_once`. An
operator (or a cron) re-runs the command on a cadence; de-dup is in-memory
per process, so each run is a single pass.
OPT-IN and INERT: this only reads labeled issues and calls the committed
coordinator intake entry — no CI, OIDC, git/patch apply, or GitHub-Actions
network. Owner/repo/label come from CLI flags; the token comes from
``GITHUB_TOKEN``. Prints the issue ids ingested on this pass (one per line).
"""
from agent_team.transport.github_intake import (
GithubIntake,
build_default_issue_client,
)
coordinator = _build_coordinator(args)
coordinator.setup()
client = build_default_issue_client(owner=args.owner, repo=args.repo)
intake = GithubIntake(
client=client,
coordinator=coordinator,
label=args.label,
)
ingested = intake.poll_once()
for issue_id in ingested:
print(issue_id, file=out)
if not ingested:
print("github-intake: no new labeled issues to ingest", file=sys.stderr)
return 0
def _cmd_force_resume(args: argparse.Namespace, *, out: Any) -> int:
"""Force-resume a parked task's question (destructive; audit-logged).
@ -828,6 +949,39 @@ def build_parser() -> argparse.ArgumentParser:
)
p_serve.set_defaults(func=_cmd_serve)
p_intake = sub.add_parser(
"intake-github",
help="poll a repo for labeled issues and start one task per new issue",
)
p_intake.add_argument(
"--owner",
required=True,
help="GitHub repository owner / org login to poll for intake issues",
)
p_intake.add_argument(
"--repo",
required=True,
help="GitHub repository name to poll for intake issues",
)
p_intake.add_argument(
"--label",
required=True,
help="issue label that flags an issue as pipeline intake (non-empty)",
)
p_intake.add_argument(
"--transport",
choices=_TRANSPORT_CHOICES,
default="github",
help="channel for delivering clarifier questions (default: github)",
)
p_intake.add_argument(
"--dry-run",
action="store_true",
dest="dry_run",
help="use a non-posting transport (no token needed; ingest still runs)",
)
p_intake.set_defaults(func=_cmd_intake_github)
return parser

View file

@ -0,0 +1,280 @@
"""Unit tests for agent_team.transport.claude_code_live (§3.3.1, §7.1 P4, D10).
The live file-drop wiring is the production backing for the §3.3.1 injected
Claude-Code ``delivery`` seam. These tests prove the contract entirely with
mocks (no network, no SDK, and the real filesystem is exercised only through
``tmp_path`` or injected fakes): the sink writes a prompt drop and returns the
session id, the ``channel_ref`` round-trips the ``question_id`` through a real
``ClaudeCodeAdapter``, a dropped answer is read back into a ``parse_answer``
payload that maps to the original question, and an unaddressable / unwritable
drop fails loudly.
"""
from __future__ import annotations
import importlib
from pathlib import Path
from typing import Any
import pytest
from agent_team.transport.base import QuestionSet
from agent_team.transport.claude_code_adapter import (
VIA,
ClaudeCodeAdapter,
ClaudeCodeDeliveryError,
build_channel_ref,
render_prompt,
)
from agent_team.transport.claude_code_live import (
ANSWER_SUFFIX,
PROMPT_SUFFIX,
build_claude_code_delivery,
build_file_drop_reader,
build_file_drop_writer,
build_live_claude_code_transport,
read_answer_payload,
)
# --------------------------------------------------------------------------- #
# Test doubles / helpers #
# --------------------------------------------------------------------------- #
class _RecordingWriter:
"""A fake prompt-writer recording the (path, prompt) it was handed."""
def __init__(self) -> None:
self.calls: list[tuple[Path, str]] = []
def __call__(self, path: Path, prompt: str) -> None:
self.calls.append((path, prompt))
def _question_set(**overrides: Any) -> QuestionSet:
defaults: dict[str, Any] = {
"thread_id": "t1",
"question_id": "q1",
"turn": 0,
"questions": ["Proceed with the dependency bump?"],
}
defaults.update(overrides)
return QuestionSet(**defaults)
def _prompt_for(question_id: str) -> str:
return render_prompt(
question_id=question_id,
turn=0,
question_set=_question_set(question_id=question_id),
deadline="2026-06-18T00:00:00Z",
)
# --------------------------------------------------------------------------- #
# Clean import (no optional SDK, no network) #
# --------------------------------------------------------------------------- #
def test_module_imports_cleanly() -> None:
"""The module reloads without any optional dependency or network."""
module = importlib.reload(
importlib.import_module("agent_team.transport.claude_code_live")
)
assert hasattr(module, "build_claude_code_delivery")
assert hasattr(module, "build_live_claude_code_transport")
assert hasattr(module, "read_answer_payload")
# --------------------------------------------------------------------------- #
# delivery sink: addressing + session id #
# --------------------------------------------------------------------------- #
def test_delivery_derives_session_id_from_embedded_marker() -> None:
"""The session id is the question_id embedded in the prompt marker."""
writer = _RecordingWriter()
delivery = build_claude_code_delivery("/drop", writer=writer)
session_id = delivery(session_hint="", prompt=_prompt_for("qEmbed"))
assert session_id == "qEmbed"
def test_delivery_falls_back_to_session_hint_without_marker() -> None:
"""A markerless prompt is addressed by the session_hint."""
writer = _RecordingWriter()
delivery = build_claude_code_delivery("/drop", writer=writer)
session_id = delivery(session_hint="mac-sess-9", prompt="bare prompt, no marker")
assert session_id == "mac-sess-9"
def test_delivery_unaddressable_prompt_raises() -> None:
"""No marker and no session_hint cannot be addressed: fail loudly."""
delivery = build_claude_code_delivery("/drop", writer=_RecordingWriter())
with pytest.raises(ClaudeCodeDeliveryError):
delivery(session_hint="", prompt="bare prompt, no marker")
def test_delivery_writes_prompt_to_drop_path() -> None:
"""The sink hands the writer a path under the drop dir with PROMPT_SUFFIX."""
writer = _RecordingWriter()
delivery = build_claude_code_delivery("/drop", writer=writer)
delivery(session_hint="", prompt=_prompt_for("qWrite"))
assert len(writer.calls) == 1
path, prompt = writer.calls[0]
assert path == Path("/drop") / f"qWrite{PROMPT_SUFFIX}"
assert "qWrite" in prompt
def test_delivery_sanitizes_separators_in_session_id() -> None:
"""A session_hint with path separators cannot escape the drop directory."""
writer = _RecordingWriter()
delivery = build_claude_code_delivery("/drop", writer=writer)
delivery(session_hint="../../etc/passwd", prompt="no marker here")
path, _prompt = writer.calls[0]
# The drop must stay inside the drop directory: separators are flattened so
# the file is a single component under /drop, not a traversal out of it.
assert path.parent == Path("/drop")
assert path.name == f"_.._etc_passwd{PROMPT_SUFFIX}"
assert path == Path("/drop") / path.name
def test_delivery_wraps_writer_oserror() -> None:
"""A writer OSError surfaces as ClaudeCodeDeliveryError (failed post)."""
def _boom(path: Path, prompt: str) -> None:
raise OSError("disk full")
delivery = build_claude_code_delivery("/drop", writer=_boom)
with pytest.raises(ClaudeCodeDeliveryError, match="failed to write"):
delivery(session_hint="", prompt=_prompt_for("qBoom"))
# --------------------------------------------------------------------------- #
# Wired through a real adapter: post -> channel_ref #
# --------------------------------------------------------------------------- #
def test_post_question_round_trips_question_id_as_channel_ref() -> None:
"""Through a real adapter, the drop's session id becomes the channel_ref."""
writer = _RecordingWriter()
transport = ClaudeCodeAdapter(
build_claude_code_delivery("/drop", writer=writer),
)
channel_ref = transport.post_question(
thread_id="t1",
question_id="q1",
turn=0,
question_set=_question_set(),
deadline="2026-06-18T00:00:00Z",
)
assert channel_ref == build_channel_ref("q1", "q1")
assert writer.calls, "a prompt drop must have been written"
def test_convenience_factory_returns_adapter() -> None:
"""``build_live_claude_code_transport`` yields a ClaudeCodeAdapter."""
transport = build_live_claude_code_transport("/drop", session_hint="mac-1")
assert isinstance(transport, ClaudeCodeAdapter)
# --------------------------------------------------------------------------- #
# Inbound: read dropped answer -> parse_answer payload #
# --------------------------------------------------------------------------- #
def test_read_answer_payload_maps_back_to_question_id() -> None:
"""A dropped answer reads into a payload that parse_answer maps correctly."""
channel_ref = build_channel_ref("q1", "q1")
def _reader(path: Path) -> str | None:
assert path == Path("/drop") / f"q1{ANSWER_SUFFIX}"
return "ship it"
payload = read_answer_payload(channel_ref, "/drop", reader=_reader)
assert payload == {"channel_ref": channel_ref, "answer": "ship it"}
# The payload must feed parse_answer and recover the original question_id.
adapter = ClaudeCodeAdapter()
assert adapter.parse_answer(payload) == ("q1", "ship it", VIA)
def test_read_answer_payload_none_when_no_answer_dropped() -> None:
"""No dropped answer yet returns None so reconcile can poll idempotently."""
def _reader(path: Path) -> str | None:
return None
payload = read_answer_payload(
build_channel_ref("q1", "q1"), "/drop", reader=_reader
)
assert payload is None
def test_read_answer_payload_rejects_foreign_channel_ref() -> None:
"""A non-Claude-Code channel_ref is rejected rather than mis-read."""
with pytest.raises(ValueError):
read_answer_payload("1718000000.001100", "/drop", reader=lambda p: None)
# --------------------------------------------------------------------------- #
# Full round-trip over the real filesystem (tmp_path, no network) #
# --------------------------------------------------------------------------- #
def test_filesystem_writer_and_reader_round_trip(tmp_path: Path) -> None:
"""The default filesystem writer/reader round-trip a prompt and answer."""
drop = tmp_path / "claude-drop"
# Post a question with the live filesystem-backed sink.
transport = build_live_claude_code_transport(drop)
channel_ref = transport.post_question(
thread_id="t1",
question_id="qFS",
turn=0,
question_set=_question_set(question_id="qFS"),
deadline="2026-06-18T00:00:00Z",
)
assert channel_ref == build_channel_ref("qFS", "qFS")
prompt_file = drop / f"qFS{PROMPT_SUFFIX}"
assert prompt_file.exists()
assert "qFS" in prompt_file.read_text(encoding="utf-8")
# Before the harness answers, the reader yields None.
assert read_answer_payload(channel_ref, drop) is None
# The harness drops an answer file; the reader picks it up.
(drop / f"qFS{ANSWER_SUFFIX}").write_text("done", encoding="utf-8")
payload = read_answer_payload(channel_ref, drop)
assert payload == {"channel_ref": channel_ref, "answer": "done"}
assert transport.parse_answer(payload) == ("qFS", "done", VIA)
def test_build_file_drop_writer_creates_dir(tmp_path: Path) -> None:
"""The filesystem writer creates a missing drop directory on first write."""
drop = tmp_path / "nested" / "drop"
writer = build_file_drop_writer(drop)
writer(drop / f"qX{PROMPT_SUFFIX}", "hello")
assert (drop / f"qX{PROMPT_SUFFIX}").read_text(encoding="utf-8") == "hello"
def test_build_file_drop_reader_missing_file_returns_none(tmp_path: Path) -> None:
"""The filesystem reader returns None for an absent answer file."""
reader = build_file_drop_reader(tmp_path)
assert reader(tmp_path / "absent.answer.txt") is None

View file

@ -0,0 +1,233 @@
"""Unit tests for agent_team.transport.github_intake (§3.3.1, INTAKE poller).
Fully hermetic: both the GitHub issue client and the coordinator are injected
in-memory fakes, so no network call, token, GitHub SDK, or model is exercised.
The tests pin the poller's contract: a labeled issue creates exactly one task,
a re-poll does not double-ingest, unlabeled issues are never seen (the client
filters by label), and the intake text is the issue title + body.
"""
from __future__ import annotations
from typing import Any
import pytest
from agent_team.transport.github_intake import (
GITHUB_TRANSPORT_NAME,
GithubIntake,
issue_task_text,
)
INTAKE_LABEL = "agent-team"
# --------------------------------------------------------------------------- #
# Fakes
# --------------------------------------------------------------------------- #
class FakeIssueClient:
"""In-memory ``GithubIssueClient`` returning only issues carrying ``label``.
Mirrors the production client's contract: ``list_open_issues(label=...)``
returns the subset of the configured issues whose ``labels`` include the
requested label. Records each requested label so a test can assert the
poller queries with the configured label.
"""
def __init__(self, issues: list[dict[str, Any]]) -> None:
self.issues = issues
self.requested_labels: list[str] = []
def list_open_issues(self, *, label: str) -> list[dict[str, Any]]:
self.requested_labels.append(label)
return [issue for issue in self.issues if label in (issue.get("labels") or [])]
class FakeCoordinator:
"""In-memory coordinator double recording every ``start_task`` call.
Captures the keyword arguments of each call so a test can assert exactly one
task was started, with the expected ``task_text`` / ``transport_name``.
Returns a synthetic ``thread_id`` like the real coordinator.
"""
def __init__(self) -> None:
self.calls: list[dict[str, Any]] = []
def start_task(self, *, task_text: str, transport_name: str) -> str:
self.calls.append({"task_text": task_text, "transport_name": transport_name})
return f"thread-{len(self.calls)}"
def _issue(
issue_id: int,
*,
title: str = "Do the thing",
body: str = "with details",
labels: list[str] | None = None,
) -> dict[str, Any]:
"""Build a minimal GitHub-issue-shaped mapping for the fakes."""
return {
"id": issue_id,
"number": issue_id,
"title": title,
"body": body,
"labels": [INTAKE_LABEL] if labels is None else labels,
}
# --------------------------------------------------------------------------- #
# issue_task_text
# --------------------------------------------------------------------------- #
def test_issue_task_text_joins_title_and_body() -> None:
text = issue_task_text(_issue(1, title="Add poller", body="for GitHub intake"))
assert text == "Add poller\n\nfor GitHub intake"
def test_issue_task_text_title_only_when_body_empty() -> None:
assert issue_task_text(_issue(1, title="Title only", body="")) == "Title only"
assert issue_task_text(_issue(1, title="Title only", body=" ")) == "Title only"
def test_issue_task_text_falls_back_to_id_when_title_empty() -> None:
text = issue_task_text(_issue(42, title="", body=""))
assert text == "issue #42"
def test_issue_task_text_strips_surrounding_whitespace() -> None:
text = issue_task_text(_issue(1, title=" Trim me ", body="\n body \n"))
assert text == "Trim me\n\nbody"
# --------------------------------------------------------------------------- #
# GithubIntake construction
# --------------------------------------------------------------------------- #
def test_empty_label_is_rejected() -> None:
with pytest.raises(ValueError):
GithubIntake(
client=FakeIssueClient([]), coordinator=FakeCoordinator(), label=""
)
# --------------------------------------------------------------------------- #
# poll_once: the core contract
# --------------------------------------------------------------------------- #
def test_labeled_issue_creates_exactly_one_task() -> None:
client = FakeIssueClient([_issue(1, title="Build it", body="now")])
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
ingested = intake.poll_once()
assert ingested == ["1"]
assert len(coordinator.calls) == 1
call = coordinator.calls[0]
assert call["task_text"] == "Build it\n\nnow"
assert call["transport_name"] == GITHUB_TRANSPORT_NAME
assert client.requested_labels == [INTAKE_LABEL]
def test_repoll_does_not_double_ingest() -> None:
client = FakeIssueClient([_issue(1)])
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
first = intake.poll_once()
second = intake.poll_once()
assert first == ["1"]
assert second == [] # already ingested -> no new task
assert len(coordinator.calls) == 1
assert intake.ingested_ids == frozenset({"1"})
def test_unlabeled_issues_are_ignored() -> None:
client = FakeIssueClient(
[
_issue(1, labels=[INTAKE_LABEL]),
_issue(2, labels=["bug"]),
_issue(3, labels=[]),
]
)
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
ingested = intake.poll_once()
assert ingested == ["1"]
assert len(coordinator.calls) == 1
assert coordinator.calls[0]["task_text"].startswith("Do the thing")
def test_new_issue_on_second_poll_is_ingested() -> None:
issues = [_issue(1)]
client = FakeIssueClient(issues)
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
first = intake.poll_once()
issues.append(_issue(2, title="Second", body="task"))
second = intake.poll_once()
assert first == ["1"]
assert second == ["2"]
assert len(coordinator.calls) == 2
assert coordinator.calls[1]["task_text"] == "Second\n\ntask"
def test_multiple_labeled_issues_each_create_one_task() -> None:
client = FakeIssueClient([_issue(1), _issue(2), _issue(3)])
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
ingested = intake.poll_once()
assert ingested == ["1", "2", "3"]
assert len(coordinator.calls) == 3
def test_id_falls_back_to_number_when_id_absent() -> None:
issue = {"number": 7, "title": "No id", "body": "", "labels": [INTAKE_LABEL]}
client = FakeIssueClient([issue])
coordinator = FakeCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
ingested = intake.poll_once()
assert ingested == ["7"]
assert intake.poll_once() == [] # de-dup on number-derived id
def test_failed_start_task_leaves_issue_eligible_for_retry() -> None:
"""A raising start_task must NOT mark the issue ingested (no silent drop)."""
class FlakyCoordinator:
def __init__(self) -> None:
self.attempts = 0
def start_task(self, *, task_text: str, transport_name: str) -> str:
self.attempts += 1
if self.attempts == 1:
raise RuntimeError("transient intake failure")
return "thread-ok"
client = FakeIssueClient([_issue(1)])
coordinator = FlakyCoordinator()
intake = GithubIntake(client=client, coordinator=coordinator, label=INTAKE_LABEL)
with pytest.raises(RuntimeError):
intake.poll_once()
assert intake.ingested_ids == frozenset() # not recorded -> retryable
# The retry succeeds and ingests the issue exactly once.
ingested = intake.poll_once()
assert ingested == ["1"]
assert coordinator.attempts == 2

View file

@ -0,0 +1,333 @@
"""Unit tests for agent_team.transport.github_live (§3.3.1, §7.1 P4).
The live poster is the production ``requests`` backing for the §3.3.1 injected
``HttpPost`` seam. These tests prove the contract entirely with mocks (no
network, and ``requests`` itself is never required): the poster performs the
REST POST and returns ``(status, data)``; the posted comment carries the
``<!-- shq:<question_id> -->`` marker; the new comment ``id`` round-trips as the
``channel_ref`` through a real ``GitHubTransport``; a missing package / token
fails loudly; and a non-2xx response surfaces as ``GitHubApiError``.
No CI, OIDC, git-apply, or GitHub Actions surface is touched — this is human-gate
transport I/O only (production default stays P2: clarify->plan->review).
"""
from __future__ import annotations
import importlib
from typing import Any
import pytest
from agent_team.transport.base import GITHUB_MARKER_TEMPLATE, QuestionSet, Transport
from agent_team.transport.github_adapter import GitHubApiError, GitHubTransport
from agent_team.transport.github_live import (
build_github_poster,
build_live_github_transport,
)
# --------------------------------------------------------------------------- #
# Test doubles #
# --------------------------------------------------------------------------- #
class _FakeResponse:
"""A ``requests.Response``-like object: ``status_code`` + ``json()``."""
def __init__(self, status_code: int, data: dict[str, Any]) -> None:
self.status_code = status_code
self._data = data
def json(self) -> dict[str, Any]:
return self._data
class _FakeSession:
"""A fake ``requests.Session`` recording ``post`` kwargs and scripting a reply."""
def __init__(
self, status_code: int = 201, data: dict[str, Any] | None = None
) -> None:
self.response = _FakeResponse(
status_code, {"id": 987654321} if data is None else data
)
self.calls: list[dict[str, Any]] = []
def post(
self, url: str, *, headers: dict[str, str], json: dict[str, Any]
) -> _FakeResponse:
self.calls.append({"url": url, "headers": headers, "json": json})
return self.response
def _question_set() -> QuestionSet:
return QuestionSet(
thread_id="task-7",
question_id="q-42",
turn=1,
questions=["Ship it?"],
context={"repo": "agent-team"},
)
# --------------------------------------------------------------------------- #
# Clean import without requests #
# --------------------------------------------------------------------------- #
def test_module_imports_without_requests(monkeypatch: pytest.MonkeyPatch) -> None:
"""The module imports cleanly even when ``requests`` cannot be imported."""
import builtins
real_import = builtins.__import__
def _blocked_import(name: str, *args: Any, **kwargs: Any) -> Any:
if name == "requests" or name.startswith("requests."):
raise ImportError("requests is blocked for this test")
return real_import(name, *args, **kwargs)
monkeypatch.setattr(builtins, "__import__", _blocked_import)
module = importlib.reload(
importlib.import_module("agent_team.transport.github_live")
)
assert hasattr(module, "build_github_poster")
assert hasattr(module, "build_live_github_transport")
# --------------------------------------------------------------------------- #
# Happy path: injected fake client #
# --------------------------------------------------------------------------- #
def test_poster_returns_status_and_data() -> None:
"""The poster forwards to the client and returns ``(status, data)``."""
session = _FakeSession()
poster = build_github_poster(client=session)
status, data = poster(
"https://api.github.com/repos/o/r/issues/1/comments",
headers={"Authorization": "Bearer x"},
json_body={"body": "hi"},
)
assert status == 201
assert data["id"] == 987654321
assert len(session.calls) == 1
assert session.calls[0]["json"] == {"body": "hi"}
def test_post_question_round_trips_comment_id_as_channel_ref() -> None:
"""Wired through a real ``GitHubTransport``, the comment id is the channel_ref."""
session = _FakeSession(data={"id": 555})
transport = GitHubTransport(
owner="Sea-Haven-Industries",
repo="agent-team",
issue_number=1,
http_post=build_github_poster(client=session),
token_provider=lambda: "ghp_fake",
)
channel_ref = transport.post_question(
thread_id="task-7",
question_id="q-42",
turn=1,
question_set=_question_set(),
deadline="2026-06-18T00:00:00Z",
)
assert channel_ref == "555"
def test_posted_comment_carries_question_id_marker() -> None:
"""The posted comment body embeds ``<!-- shq:<question_id> -->``."""
session = _FakeSession()
transport = GitHubTransport(
owner="Sea-Haven-Industries",
repo="agent-team",
issue_number=1,
http_post=build_github_poster(client=session),
token_provider=lambda: "ghp_fake",
)
transport.post_question(
thread_id="task-7",
question_id="q-42",
turn=1,
question_set=_question_set(),
deadline="2026-06-18T00:00:00Z",
)
assert len(session.calls) == 1
body = session.calls[0]["json"]["body"]
assert GITHUB_MARKER_TEMPLATE.format(question_id="q-42") in body
def test_convenience_transport_factory_round_trips_and_parses() -> None:
"""``build_live_github_transport`` wires the poster and the marker round-trips."""
session = _FakeSession(data={"id": 777})
transport = build_live_github_transport(
owner="Sea-Haven-Industries",
repo="agent-team",
issue_number=1,
token="ghp_fake",
client=session,
)
assert isinstance(transport, Transport)
channel_ref = transport.post_question(
thread_id="task-7",
question_id="q-42",
turn=1,
question_set=_question_set(),
deadline="2026-06-18T00:00:00Z",
)
assert channel_ref == "777"
# The factory threads one token into the transport's auth header.
assert session.calls[0]["headers"]["Authorization"] == "Bearer ghp_fake"
# The marker the poster shipped round-trips back through parse_answer: a
# human reply quoting the question comment recovers the same question_id.
posted_body = session.calls[0]["json"]["body"]
reply = "> " + posted_body.replace("\n", "\n> ") + "\nLooks good, ship it."
question_id, answer, via = transport.parse_answer(
{"comment": {"body": reply, "user": {"login": "adam"}}}
)
assert question_id == "q-42"
assert answer == "Looks good, ship it."
assert via == "github:adam"
def test_response_with_status_attr_accepted() -> None:
"""A response exposing ``status`` (not ``status_code``) is also accepted."""
class _StatusOnly:
status = 200
def json(self) -> dict[str, Any]:
return {"id": 1}
class _Session:
def post(self, url: str, **kwargs: Any) -> _StatusOnly:
return _StatusOnly()
poster = build_github_poster(client=_Session())
status, data = poster("u", headers={}, json_body={})
assert status == 200
assert data["id"] == 1
def test_empty_body_coerces_to_dict() -> None:
"""A non-dict JSON body coerces to ``{}`` (adapter's missing-id guard fires)."""
class _NullJson:
status_code = 201
def json(self) -> Any:
return None
class _Session:
def post(self, url: str, **kwargs: Any) -> _NullJson:
return _NullJson()
poster = build_github_poster(client=_Session())
status, data = poster("u", headers={}, json_body={})
assert status == 201
assert data == {}
# --------------------------------------------------------------------------- #
# Failure modes #
# --------------------------------------------------------------------------- #
def test_non_2xx_response_raises_github_api_error() -> None:
"""A non-2xx response surfaces as ``GitHubApiError`` carrying the status."""
session = _FakeSession(status_code=403, data={"message": "Forbidden"})
poster = build_github_poster(client=session)
with pytest.raises(GitHubApiError) as excinfo:
poster("u", headers={}, json_body={"body": "x"})
assert excinfo.value.status == 403
def test_unsupported_response_raises_type_error() -> None:
"""A response with neither status nor json() is fatal (not a silent success)."""
class _Session:
def post(self, url: str, **kwargs: Any) -> object:
return object()
poster = build_github_poster(client=_Session())
with pytest.raises(TypeError, match="status_code"):
poster("u", headers={}, json_body={})
def test_missing_token_raises_runtime_error(monkeypatch: pytest.MonkeyPatch) -> None:
"""No token and no GITHUB_TOKEN raises a clear RuntimeError.
Stub ``requests`` into ``sys.modules`` so the deferred import SUCCEEDS and
the no-token branch is what's under test (avoids local-vs-CI drift where a
missing package would otherwise mask the token check).
"""
import sys
from types import ModuleType
fake = ModuleType("requests")
fake.Session = lambda: type(
"S", (), {"headers": {}, "post": lambda self, *a, **k: None}
)() # type: ignore[attr-defined]
monkeypatch.setitem(sys.modules, "requests", fake)
monkeypatch.delenv("GITHUB_TOKEN", raising=False)
with pytest.raises(RuntimeError, match="GitHub token"):
build_github_poster()
def test_missing_package_raises_runtime_error(monkeypatch: pytest.MonkeyPatch) -> None:
"""A missing ``requests`` package raises a clear RuntimeError."""
import builtins
real_import = builtins.__import__
def _blocked_import(name: str, *args: Any, **kwargs: Any) -> Any:
if name == "requests" or name.startswith("requests."):
raise ImportError("requests is blocked for this test")
return real_import(name, *args, **kwargs)
monkeypatch.setattr(builtins, "__import__", _blocked_import)
monkeypatch.setenv("GITHUB_TOKEN", "ghp_present")
with pytest.raises(RuntimeError, match="requests is unavailable"):
build_github_poster()
def test_token_falls_back_to_env(monkeypatch: pytest.MonkeyPatch) -> None:
"""With ``requests`` stubbed, a GITHUB_TOKEN env var builds a session cleanly."""
import sys
from types import ModuleType
captured: dict[str, Any] = {}
class _Session:
def __init__(self) -> None:
self.headers: dict[str, str] = {}
def post(self, *a: Any, **k: Any) -> None: # pragma: no cover - unused
return None
def _make_session() -> _Session:
session = _Session()
captured["session"] = session
return session
fake = ModuleType("requests")
fake.Session = _make_session # type: ignore[attr-defined]
monkeypatch.setitem(sys.modules, "requests", fake)
monkeypatch.setenv("GITHUB_TOKEN", "ghp_from_env")
poster = build_github_poster()
assert callable(poster)
assert captured["session"].headers["Authorization"] == "Bearer ghp_from_env"

View file

@ -705,20 +705,154 @@ def test_start_runs_setup_and_start_task_and_prints_thread_id(
assert callable(coord.review_wiring)
def test_build_transport_live_github_raises_system_exit(cli: ModuleType) -> None:
"""A non-slack live transport is not wired and raises a clear SystemExit."""
def test_build_transport_live_github_builds_github_transport(
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
) -> None:
"""``--transport github`` builds the live GitHub transport from the env thread.
Token resolution is deferred to call time, so with GITHUB_TOKEN set and the
issue thread configured the builder constructs a GitHubTransport bound to
that ``owner/repo#issue`` (no network — only the issue thread is wired).
"""
monkeypatch.setenv("GITHUB_TOKEN", "ghp_test")
monkeypatch.setenv("GITHUB_OWNER", "Sea-Haven-Industries")
monkeypatch.setenv("GITHUB_REPO", "orchestrator")
monkeypatch.setenv("GITHUB_ISSUE_NUMBER", "42")
from agent_team.transport.github_adapter import GitHubTransport
args = argparse.Namespace(dry_run=False, transport="github")
with pytest.raises(SystemExit, match="is not wired for the run-team CLI"):
transport = cli._build_transport(args)
assert isinstance(transport, GitHubTransport)
assert transport.owner == "Sea-Haven-Industries"
assert transport.repo == "orchestrator"
assert transport.issue_number == 42
def test_build_transport_live_github_missing_thread_raises_system_exit(
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A github transport with no issue thread configured fails loudly."""
monkeypatch.delenv("GITHUB_OWNER", raising=False)
monkeypatch.delenv("GITHUB_REPO", raising=False)
monkeypatch.delenv("GITHUB_ISSUE_NUMBER", raising=False)
args = argparse.Namespace(dry_run=False, transport="github")
with pytest.raises(SystemExit, match="GITHUB_OWNER"):
cli._build_transport(args)
def test_build_transport_live_claude_code_raises_system_exit(cli: ModuleType) -> None:
"""claude_code is likewise un-wired for the P1 CLI surface."""
def test_build_transport_live_github_non_integer_issue_raises_system_exit(
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
) -> None:
"""A non-integer GITHUB_ISSUE_NUMBER fails loudly rather than half-building."""
monkeypatch.setenv("GITHUB_OWNER", "o")
monkeypatch.setenv("GITHUB_REPO", "r")
monkeypatch.setenv("GITHUB_ISSUE_NUMBER", "not-a-number")
args = argparse.Namespace(dry_run=False, transport="github")
with pytest.raises(SystemExit, match="must be an integer"):
cli._build_transport(args)
def test_build_transport_live_claude_code_builds_file_drop_transport(
cli: ModuleType, monkeypatch: pytest.MonkeyPatch, tmp_path: Path
) -> None:
"""``--transport claude_code`` builds the live file-drop transport from env."""
monkeypatch.setenv("CLAUDE_CODE_DROP_DIR", str(tmp_path / "drops"))
from agent_team.transport.claude_code_adapter import ClaudeCodeAdapter
args = argparse.Namespace(dry_run=False, transport="claude_code")
with pytest.raises(SystemExit, match="is not wired for the run-team CLI"):
transport = cli._build_transport(args)
assert isinstance(transport, ClaudeCodeAdapter)
def test_build_transport_live_claude_code_missing_drop_dir_raises_system_exit(
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
) -> None:
"""claude_code with no CLAUDE_CODE_DROP_DIR fails loudly."""
monkeypatch.delenv("CLAUDE_CODE_DROP_DIR", raising=False)
args = argparse.Namespace(dry_run=False, transport="claude_code")
with pytest.raises(SystemExit, match="CLAUDE_CODE_DROP_DIR"):
cli._build_transport(args)
def test_intake_github_polls_and_starts_tasks_dry_run(
cli: ModuleType,
db_path: Path,
audit_log: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""``intake-github --dry-run`` builds a coordinator, polls, and ingests issues.
The lazily-imported ``Coordinator`` is replaced with the fake (so no graph /
SDK is built), and ``build_default_issue_client`` is patched to return an
in-memory fake lister (so no GitHub network). The fake coordinator records
the ``start_task`` calls the intake leaf makes — one per labeled issue.
"""
_FakeCoordinator.instances.clear()
monkeypatch.setattr(
"agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True
)
class _FakeIssueClient:
def list_open_issues(self, *, label: str) -> list[dict[str, Any]]:
assert label == "agent-team"
return [
{"id": 1001, "title": "Do thing A", "body": "details A"},
{"id": 1002, "title": "Do thing B", "body": ""},
]
def _fake_build_client(*, owner: str, repo: str, **_kw: Any) -> Any:
assert owner == "Sea-Haven-Industries"
assert repo == "orchestrator"
return _FakeIssueClient()
monkeypatch.setattr(
"agent_team.transport.github_intake.build_default_issue_client",
_fake_build_client,
raising=True,
)
code, out = _run(
cli,
db_path,
audit_log,
"intake-github",
"--dry-run",
"--owner",
"Sea-Haven-Industries",
"--repo",
"orchestrator",
"--label",
"agent-team",
)
assert code == 0
assert len(_FakeCoordinator.instances) == 1
coord = _FakeCoordinator.instances[0]
assert coord.setup_called is True
# Both labeled issues were ingested; the leaf passes title+body and the
# GitHub transport name. (The fake records only the LAST call's kwargs.)
assert coord.start_kwargs == {
"task_text": "Do thing B",
"transport_name": "github",
}
# The ingested issue ids are printed (one per line).
assert out.split() == ["1001", "1002"]
def test_intake_github_label_required(cli: ModuleType) -> None:
"""``intake-github`` requires --owner/--repo/--label (argparse usage error)."""
parser = cli.build_parser()
with pytest.raises(SystemExit):
parser.parse_args(["intake-github", "--owner", "o", "--repo", "r"])
def test_build_transport_dry_run_returns_dry_run_transport(cli: ModuleType) -> None:
"""``dry_run=True`` yields a _DryRunTransport whose post returns a synthetic ref."""
args = argparse.Namespace(dry_run=True, transport="slack")