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=