From f916c03818d4d7b7f269cef64608c72f2513b069 Mon Sep 17 00:00:00 2001 From: Adam Moussa <166072409+amoussa1229@users.noreply.github.com> Date: Mon, 22 Jun 2026 17:20:43 -0400 Subject: [PATCH] feat(agent-team): durable GitHub-issue intake de-dup + intake hardening (#32) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat(agent-team): durable GitHub-issue intake de-dup (schema v2) The intake poller de-duped ingested issues in an in-memory set that does not survive a process restart. A scheduled/cron intake (each run a fresh process) would therefore re-ingest every still-open labeled issue on every run and spawn duplicate pipeline tasks. Since the box is read-only (no write token to remove the intake label), durable de-dup is the only correct guard. - schema v2: new ingested_issues(source, issue_id, ingested_at) table + issue_already_ingested / record_issue_ingested helpers; migrate() adds the table to a legacy v1 DB and restamps; init_db creates it. - github_intake: pluggable IngestStore seam (in-memory default preserved for tests/one-off; durable build_ledger_ingest_store for production). Record is after start_task succeeds, so a failed intake stays retryable. - run-team.py intake-github wires the ledger store keyed by github:owner/repo, making a scheduled timer idempotent across runs. 978 tests pass (ruff clean). * 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. --- agent-team/agent_team/db/schema.py | 99 ++++++- agent-team/agent_team/db/schema.sql | 15 ++ .../agent_team/transport/github_intake.py | 246 +++++++++++++++--- agent-team/run-team.py | 14 +- agent-team/tests/test_github_intake.py | 159 +++++++++++ agent-team/tests/test_schema.py | 73 ++++++ 6 files changed, 559 insertions(+), 47 deletions(-) diff --git a/agent-team/agent_team/db/schema.py b/agent-team/agent_team/db/schema.py index d2d0067..8b7cf47 100644 --- a/agent-team/agent_team/db/schema.py +++ b/agent-team/agent_team/db/schema.py @@ -31,21 +31,25 @@ __all__ = [ "PENDING_QUESTIONS_DDL", "PENDING_QUESTIONS_INDEXES_DDL", "BUDGET_LEDGER_INDEXES_DDL", + "INGESTED_ISSUES_DDL", "SCHEMA_META_DDL", "SCHEMA_VERSION", "QUESTION_STATES", "answer_question", "connect", + "delete_issue_ingested", "expire_question", "find_open_question_by_channel_ref", "init_db", + "issue_already_ingested", "migrate", + "record_issue_ingested", "reopen_question", "supersede_question", ] # Bump when the DDL below changes; migrate() steps a connection forward. -SCHEMA_VERSION: int = 1 +SCHEMA_VERSION: int = 2 # Default SQLite busy timeout (ms) so concurrent writers wait for the write # lock rather than failing immediately. @@ -105,6 +109,24 @@ CREATE INDEX IF NOT EXISTS idx_budget_ledger_thread ON budget_ledger (thread_id); """.strip() +# ingested_issues: durable de-dup ledger for the GitHub-issue intake poller +# (schema v2). The in-memory poller set does not survive a process restart, so a +# scheduled/cron intake (each run a fresh process) would re-ingest every still +# -open labeled issue and spawn duplicate pipeline tasks. This table records, per +# (source, issue_id), the issues already turned into tasks so intake is idempotent +# across restarts — mirroring the pending_questions durability discipline. The box +# is read-only (no write token to remove the intake label), so durable de-dup is +# the only correct guard. `source` namespaces by repo (e.g. "github:owner/repo") +# so per-repo issue numbers from different repos can never collide. +INGESTED_ISSUES_DDL: str = """ +CREATE TABLE IF NOT EXISTS ingested_issues ( + source TEXT NOT NULL, + issue_id TEXT NOT NULL, + ingested_at TEXT NOT NULL, + PRIMARY KEY (source, issue_id) +) +""".strip() + SCHEMA_META_DDL: str = """ CREATE TABLE IF NOT EXISTS schema_meta ( id INTEGER PRIMARY KEY CHECK (id = 1), @@ -180,7 +202,11 @@ def init_db(db_path: Path) -> None: conn.execute(BUDGET_LEDGER_DDL) for stmt in _split_statements(BUDGET_LEDGER_INDEXES_DDL): conn.execute(stmt) - # Record the schema version (single-row table). + conn.execute(INGESTED_ISSUES_DDL) + # Record the schema version (single-row table). DO NOTHING leaves an + # existing row's version untouched (an already-stamped DB just gains any + # IF-NOT-EXISTS tables above); migrate() is what steps the version stamp + # forward on an existing DB. conn.execute( "INSERT INTO schema_meta (id, schema_version) VALUES (1, ?) " "ON CONFLICT(id) DO NOTHING", @@ -213,7 +239,17 @@ def migrate(conn: sqlite3.Connection) -> None: conn.execute(stmt) current = 1 - # Future steps go here: `if current < 2: ...; current = 2`. + if current < 2: + # v2: durable GitHub-issue intake de-dup ledger. + conn.execute(INGESTED_ISSUES_DDL) + current = 2 + + # Future steps go here: `if current < 3: ...; current = 3`. + + # Applied UNCONDITIONALLY (idempotent IF NOT EXISTS) so an already-stamped DB + # — which skips the version blocks above — still gains the v2 table without a + # restamp. Safe on existing data: a fresh empty table only. + conn.execute(INGESTED_ISSUES_DDL) # Defense-in-depth index, applied UNCONDITIONALLY (idempotent IF NOT EXISTS) # so an already-stamped v1 DB — which skips the `current < 1` block above — @@ -235,6 +271,63 @@ def migrate(conn: sqlite3.Connection) -> None: ) +def issue_already_ingested( + conn: sqlite3.Connection, *, source: str, issue_id: str +) -> bool: + """Return True if ``(source, issue_id)`` was already turned into a task. + + The durable counterpart of the intake poller's in-memory set: survives a + process restart so a scheduled/cron intake never re-ingests a still-open + labeled issue. ``source`` namespaces by repo (e.g. ``"github:owner/repo"``) + so per-repo issue numbers from different repos do not collide. + """ + row = conn.execute( + "SELECT 1 FROM ingested_issues WHERE source = ? AND issue_id = ?", + (source, issue_id), + ).fetchone() + return row is not None + + +def record_issue_ingested( + conn: sqlite3.Connection, + *, + source: str, + issue_id: str, + ingested_at: str | None = None, +) -> bool: + """Durably record that ``(source, issue_id)`` has been ingested. + + ``INSERT OR IGNORE`` against the composite primary key, so a concurrent or + repeated record is a harmless no-op (first-record-wins, mirroring the + pending_questions compare-and-set discipline). Returns True if this call + inserted a new row, False if the pair was already present. Callers record + AFTER ``start_task`` succeeds so a failed intake leaves the issue eligible + for retry (no silent drop). + """ + cur = conn.execute( + "INSERT OR IGNORE INTO ingested_issues (source, issue_id, ingested_at) " + "VALUES (?, ?, ?)", + (source, issue_id, ingested_at or _utc_now_iso()), + ) + return cur.rowcount > 0 + + +def delete_issue_ingested( + conn: sqlite3.Connection, *, source: str, issue_id: str +) -> bool: + """Undo a recorded ingest of ``(source, issue_id)``; True if a row was removed. + + Used to RELEASE a claim made via :func:`record_issue_ingested` when the + downstream intake (``start_task``) raises, so a transiently-failed issue + stays eligible for retry on the next poll rather than being silently dropped. + """ + cur = conn.execute( + "DELETE FROM ingested_issues WHERE source = ? AND issue_id = ?", + (source, issue_id), + ) + return cur.rowcount > 0 + + def answer_question( conn: sqlite3.Connection, *, diff --git a/agent-team/agent_team/db/schema.sql b/agent-team/agent_team/db/schema.sql index db4f056..6690104 100644 --- a/agent-team/agent_team/db/schema.sql +++ b/agent-team/agent_team/db/schema.sql @@ -64,6 +64,21 @@ CREATE INDEX IF NOT EXISTS idx_budget_ledger_day CREATE INDEX IF NOT EXISTS idx_budget_ledger_thread ON budget_ledger (thread_id); +-- ingested_issues (schema v2): durable de-dup ledger for the GitHub-issue intake +-- poller. The in-memory poller set does not survive a restart, so a scheduled +-- intake (each run a fresh process) would re-ingest every still-open labeled +-- issue and spawn duplicate tasks. This table records, per (source, issue_id), +-- the issues already turned into tasks so intake is idempotent across restarts. +-- The box is read-only (no write token to remove the intake label), so durable +-- de-dup is the only correct guard. `source` namespaces by repo (e.g. +-- "github:owner/repo") so per-repo issue numbers from different repos cannot collide. +CREATE TABLE IF NOT EXISTS ingested_issues ( + source TEXT NOT NULL, + issue_id TEXT NOT NULL, + ingested_at TEXT NOT NULL, + PRIMARY KEY (source, issue_id) +); + -- schema_meta: single-row table recording the applied schema version so -- migrate() can detect and step forward. CREATE TABLE IF NOT EXISTS schema_meta ( diff --git a/agent-team/agent_team/transport/github_intake.py b/agent-team/agent_team/transport/github_intake.py index 14e0598..263d586 100644 --- a/agent-team/agent_team/transport/github_intake.py +++ b/agent-team/agent_team/transport/github_intake.py @@ -21,15 +21,26 @@ socket: 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. +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. @@ -48,12 +59,15 @@ Design constraints honoured here (pre-deployment scaffolding): 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", ] @@ -67,6 +81,17 @@ GITHUB_TRANSPORT_NAME = "github" # 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. @@ -83,6 +108,59 @@ class GithubIssueClient(Protocol): ... +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. @@ -127,9 +205,10 @@ class GithubIntake: 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. + 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 @@ -142,8 +221,9 @@ class GithubIntake: client: GithubIssueClient, coordinator: Any, label: str, + store: IngestStore | None = None, ) -> None: - """Bind the poller to one client, coordinator, and intake label. + """Bind the poller to one client, coordinator, intake label, and store. Args: client: The injected issue source. Its ``list_open_issues`` is the @@ -155,6 +235,10 @@ class GithubIntake: 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( @@ -164,10 +248,12 @@ class GithubIntake: 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() + # 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: @@ -176,8 +262,8 @@ class GithubIntake: @property def ingested_ids(self) -> frozenset[str]: - """A snapshot of the issue ids already ingested this process (read-only).""" - return frozenset(self._ingested) + """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. @@ -186,17 +272,20 @@ class GithubIntake: 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 + 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=, - transport_name="github")`` and record its id so a subsequent poll - does not re-ingest it. + transport_name="github")``. - 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. + 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. @@ -204,7 +293,7 @@ class GithubIntake: 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: + if not self._store.claim(issue_id): _LOG.debug("github-intake: issue %s already ingested; skip", issue_id) continue @@ -214,15 +303,18 @@ class GithubIntake: 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) + # 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 @@ -233,7 +325,6 @@ 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). @@ -245,8 +336,15 @@ def build_default_issue_client( 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. + 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=