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/transport/github_intake.py
Adam Moussa ebc350aee5 fix(agent-team): harden github intake per /sh-security-review (CWE-918, idempotency)
Fixes from the high-recall detector fan-out on the durable-dedup change:

- INTAKE-LOGIC-01 (idempotency): switch the IngestStore seam from
  check-then-record (seen/mark) to claim-then-do (claim/release). The id is
  now reserved BEFORE the non-idempotent start_task side effect, so a crash in
  that window cannot re-spawn a duplicate task on the next run; a raising
  start_task releases the claim so transient failures stay retryable. Adds
  delete_issue_ingested to the schema layer for the release path.
- INTAKE-SSRF-001 / INTAKE-PATHSPLICE-002 (CWE-918) in build_default_issue_client:
  drop the caller-overridable api_root (hardcode GITHUB_API_ROOT) and validate
  owner/repo against an anchored charset before splicing them into the
  token-bearing API URL — mirrors the sibling ci_fetcher BLOCK-3/FIX-3 fixes.

Tests cover cross-process duplicate prevention, release-on-failure retry, and
the owner/repo + api_root rejection. 982 tests pass, ruff clean.

Follow-up (pre-existing, not introduced here): the label-only intake has no
author allowlist (cf. AGENT_TEAM_SLACK_OWNER_IDS on the Slack listener); the
Slack answer gate bounds the blast radius. Track as separate hardening.
2026-06-22 17:06:56 -04:00

457 lines
20 KiB
Python

"""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:
The poller checks an injected :class:`IngestStore` before starting a task
and records the issue id after ``start_task`` succeeds, so a re-poll over the
same open issue does not start a second task. Two stores ship:
* :class:`_InMemoryIngestStore` (the default) — best-effort, per-process; it
does NOT survive a restart. Fine for tests and one-off operator runs.
* the **durable** ledger store from :func:`build_ledger_ingest_store` — an
``ingested_issues`` SQLite table (schema v2) keyed by ``(source,
issue_id)``, mirroring the ``pending_questions`` durability discipline.
Production (a scheduled/cron intake — each run a fresh process) MUST use
this store, or every still-open labeled issue would be re-ingested on every
run and spawn duplicate tasks. The box is read-only (no write token to
remove the intake label), so durable de-dup is the only correct guard.
The poller uses **claim-then-do**: it reserves the id on the store BEFORE
the non-idempotent ``start_task`` call, and releases the claim if that call
raises. So a transient failure stays retryable, a re-run never double-spawns,
and only a hard crash in the (tiny) window between claim and call drops the
issue — the safe failure for a read-only source whose label persists.
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
import re
from typing import Any, Iterable, Protocol
__all__ = [
"GithubIntake",
"GithubIssueClient",
"IngestStore",
"build_default_issue_client",
"build_ledger_ingest_store",
"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"
# The GitHub REST host, HARDCODED (mirrors ci_fetcher FIX-3): the production
# client does NOT accept a caller-overridable api_root, so a token-bearing GET
# can never be redirected at an attacker host / file:// (CWE-918).
GITHUB_API_ROOT = "https://api.github.com"
# owner/repo are spliced into the API URL path; validate them (anchored) before
# URL construction so a crafted owner/repo cannot splice extra path/query
# segments or traversal into the authenticated request (mirrors ci_fetcher's
# _OWNER_REPO_RE / BLOCK-3). GitHub logins and repo names are a subset of this.
_OWNER_REPO_RE = re.compile(r"[A-Za-z0-9_.-]{1,100}")
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)."""
...
class IngestStore(Protocol):
"""De-dup memory for the poller: has this issue id already become a task?
A narrow seam so the poller's de-dup is swappable: a per-process in-memory
set for tests/one-off runs, or the durable SQLite ledger store for a
scheduled intake (see the module docstring).
The contract is **claim-then-do** (not check-then-record): :meth:`claim`
atomically reserves an id and reports whether THIS caller won the claim, so
the id is reserved BEFORE the (non-idempotent) ``start_task`` side effect
runs — closing the window where a crash between starting a task and recording
the id would re-spawn a duplicate. :meth:`release` undoes a claim when
``start_task`` raises, keeping a transiently-failed issue retryable.
"""
def claim(self, issue_id: str) -> bool:
"""Atomically reserve ``issue_id``. True if newly claimed (caller should
ingest); False if already claimed/ingested (caller skips)."""
...
def release(self, issue_id: str) -> None:
"""Undo a claim so the issue is retryable (used when ``start_task`` raises)."""
...
def snapshot(self) -> frozenset[str]:
"""Return a read-only snapshot of the claimed/ingested ids."""
...
class _InMemoryIngestStore:
"""Best-effort per-process de-dup; does NOT survive a restart (the default).
Reproduces the poller's original in-memory behaviour. Suitable for tests and
one-off operator runs; a scheduled/cron intake MUST use the durable store
(:func:`build_ledger_ingest_store`) instead.
"""
def __init__(self) -> None:
self._ids: set[str] = set()
def claim(self, issue_id: str) -> bool:
if issue_id in self._ids:
return False
self._ids.add(issue_id)
return True
def release(self, issue_id: str) -> None:
self._ids.discard(issue_id)
def snapshot(self) -> frozenset[str]:
return frozenset(self._ids)
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 delegated to the injected :class:`IngestStore` (see the module
docstring): the default in-memory store guards within one process; the
durable ledger store (:func:`build_ledger_ingest_store`) guards across
restarts and MUST be used for a scheduled/cron intake.
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,
store: IngestStore | None = None,
) -> None:
"""Bind the poller to one client, coordinator, intake label, and store.
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.
store: The de-dup memory. Defaults to a best-effort per-process
:class:`_InMemoryIngestStore`; a scheduled/cron intake MUST pass
the durable store from :func:`build_ledger_ingest_store` so a
restart does not re-ingest still-open labeled issues.
"""
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
# De-dup memory (default: in-memory, best-effort, NOT durable across a
# restart; production passes the durable ledger store). Tracks issue ids
# already turned into tasks so a re-poll does not double-ingest.
self._store: IngestStore = (
store if store is not None else _InMemoryIngestStore()
)
@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 ingested issue ids from the store (read-only)."""
return self._store.snapshot()
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, **claim** its id on the store first
(:meth:`IngestStore.claim`); if the claim is lost (already
claimed/ingested) the issue is skipped. Only the winning claim calls
``coordinator.start_task(task_text=<title+body>,
transport_name="github")``.
Claim-then-do (not record-after): the id is reserved BEFORE the
non-idempotent ``start_task`` side effect, so a crash between starting
the task and finishing the pass cannot re-spawn a duplicate on the next
run. If ``start_task`` *raises*, the claim is RELEASED so a transient
failure stays retryable (no silent drop). The only residual is a hard
crash (SIGKILL/OOM/reboot) between the claim and the call, which drops
the issue rather than duplicating it — the safer failure for a read-only
source whose label persists for an operator to re-trigger.
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 not self._store.claim(issue_id):
_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,
)
# The claim above reserved the id BEFORE this non-idempotent intake
# call (coordinator.start_task mints a fresh thread + an open question
# row). On a raise, RELEASE the claim so the transient failure is
# retryable on the next poll rather than a silent drop.
try:
self._coordinator.start_task(
task_text=task_text,
transport_name=GITHUB_TRANSPORT_NAME,
)
except Exception:
self._store.release(issue_id)
raise
ingested_now.append(issue_id)
return ingested_now
def build_default_issue_client(
*,
owner: str,
repo: str,
token_env: str = DEFAULT_TOKEN_ENV,
) -> 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.
Security (mirrors ci_fetcher BLOCK-3 / FIX-3, CWE-918): the GitHub host is
HARDCODED (:data:`GITHUB_API_ROOT`) — there is no caller-overridable
``api_root``, so the token-bearing GET can never be redirected at an
attacker host or ``file://``. ``owner`` and ``repo`` are validated against
an anchored charset (:data:`_OWNER_REPO_RE`) BEFORE being spliced into the
URL path, failing closed so a crafted value cannot inject extra path/query
segments or traversal.
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.
"""
for _field, _value in (("owner", owner), ("repo", repo)):
if not _OWNER_REPO_RE.fullmatch(_value):
raise ValueError(
f"invalid GitHub {_field} {_value!r}: must match "
f"{_OWNER_REPO_RE.pattern} (refusing to build an unsafe API URL)"
)
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
# Hardcoded host (no overridable api_root); owner/repo already validated.
self._api_root = GITHUB_API_ROOT
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")
# url is built from a fixed https GitHub API base; the dynamic part is
# the path, never the scheme, so there is no SSRF/file:// surface.
with _urlrequest.urlopen(request) as response: # noqa: S310 (trusted api host); nosemgrep
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()
def build_ledger_ingest_store(*, db_path: Any, source: str) -> IngestStore:
"""Build the DURABLE ledger-backed ingest store (schema v2 ``ingested_issues``).
Returns an :class:`IngestStore` whose de-dup memory is the
``ingested_issues`` SQLite table, so a scheduled/cron intake (each run a
fresh process) never re-ingests a still-open labeled issue across restarts.
``source`` namespaces by repo (e.g. ``"github:owner/repo"``) so per-repo
issue numbers from different repos cannot collide.
The DB connection and schema import are DEFERRED to call time (mirroring the
issue-client builder and the checkpointer discipline), so importing this
module never opens a connection. The store ensures its table exists on
construction, so it is robust even if ``init_db`` / ``migrate`` has not yet
run against the DB.
"""
from pathlib import Path
from agent_team.db import schema as _schema
if not source:
raise ValueError(
"build_ledger_ingest_store requires a non-empty source namespace"
)
class _LedgerIngestStore:
"""Durable ingest store backed by the ``ingested_issues`` table."""
def __init__(self) -> None:
self._source = source
self._conn = _schema.connect(Path(db_path))
# Ensure the table exists even if the DB predates schema v2 or has
# not been init_db'd yet (idempotent CREATE IF NOT EXISTS).
self._conn.execute(_schema.INGESTED_ISSUES_DDL)
def claim(self, issue_id: str) -> bool:
# INSERT OR IGNORE under the (source, issue_id) primary key: returns
# True only if THIS call inserted the row (won the claim), so two
# concurrent pollers or a re-run can never both ingest the same issue.
return _schema.record_issue_ingested(
self._conn, source=self._source, issue_id=issue_id
)
def release(self, issue_id: str) -> None:
_schema.delete_issue_ingested(
self._conn, source=self._source, issue_id=issue_id
)
def snapshot(self) -> frozenset[str]:
rows = self._conn.execute(
"SELECT issue_id FROM ingested_issues WHERE source = ?",
(self._source,),
).fetchall()
return frozenset(str(row["issue_id"]) for row in rows)
return _LedgerIngestStore()