diff --git a/agent-team/agent_team/db/schema.py b/agent-team/agent_team/db/schema.py index 106def1..8b7cf47 100644 --- a/agent-team/agent_team/db/schema.py +++ b/agent-team/agent_team/db/schema.py @@ -37,6 +37,7 @@ __all__ = [ "QUESTION_STATES", "answer_question", "connect", + "delete_issue_ingested", "expire_question", "find_open_question_by_channel_ref", "init_db", @@ -311,6 +312,22 @@ def record_issue_ingested( 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/transport/github_intake.py b/agent-team/agent_team/transport/github_intake.py index 468a825..263d586 100644 --- a/agent-team/agent_team/transport/github_intake.py +++ b/agent-team/agent_team/transport/github_intake.py @@ -36,8 +36,11 @@ De-duplication: 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 id is recorded only AFTER ``start_task`` returns, so a failing intake - leaves the issue eligible for retry on the next poll rather than a silent drop. + 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. @@ -56,6 +59,7 @@ Design constraints honoured here (pre-deployment scaffolding): from __future__ import annotations import logging +import re from typing import Any, Iterable, Protocol __all__ = [ @@ -77,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. @@ -98,21 +113,27 @@ class IngestStore(Protocol): 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). ``seen``/``mark`` mirror the - check-then-record discipline; ``snapshot`` exposes the current id set for - introspection/logging. + 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 seen(self, issue_id: str) -> bool: - """Return True if ``issue_id`` has already been ingested.""" + 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 mark(self, issue_id: str) -> None: - """Record ``issue_id`` as ingested (idempotent).""" + 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 ingested ids.""" + """Return a read-only snapshot of the claimed/ingested ids.""" ... @@ -127,11 +148,14 @@ class _InMemoryIngestStore: def __init__(self) -> None: self._ids: set[str] = set() - def seen(self, issue_id: str) -> bool: - return issue_id in self._ids - - def mark(self, issue_id: str) -> None: + 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) @@ -248,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. @@ -266,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 self._store.seen(issue_id): + if not self._store.claim(issue_id): _LOG.debug("github-intake: issue %s already ingested; skip", issue_id) continue @@ -276,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._store.mark(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 @@ -295,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). @@ -307,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=