feat(agent-team): durable GitHub-issue intake de-dup + intake hardening #32

Merged
amoussa1229 merged 2 commits from feature/agent-team-github-intake-durable-dedup into main 2026-06-22 21:20:44 +00:00
6 changed files with 559 additions and 47 deletions

View file

@ -31,21 +31,25 @@ __all__ = [
"PENDING_QUESTIONS_DDL", "PENDING_QUESTIONS_DDL",
"PENDING_QUESTIONS_INDEXES_DDL", "PENDING_QUESTIONS_INDEXES_DDL",
"BUDGET_LEDGER_INDEXES_DDL", "BUDGET_LEDGER_INDEXES_DDL",
"INGESTED_ISSUES_DDL",
"SCHEMA_META_DDL", "SCHEMA_META_DDL",
"SCHEMA_VERSION", "SCHEMA_VERSION",
"QUESTION_STATES", "QUESTION_STATES",
"answer_question", "answer_question",
"connect", "connect",
"delete_issue_ingested",
"expire_question", "expire_question",
"find_open_question_by_channel_ref", "find_open_question_by_channel_ref",
"init_db", "init_db",
"issue_already_ingested",
"migrate", "migrate",
"record_issue_ingested",
"reopen_question", "reopen_question",
"supersede_question", "supersede_question",
] ]
# Bump when the DDL below changes; migrate() steps a connection forward. # 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 # Default SQLite busy timeout (ms) so concurrent writers wait for the write
# lock rather than failing immediately. # lock rather than failing immediately.
@ -105,6 +109,24 @@ CREATE INDEX IF NOT EXISTS idx_budget_ledger_thread
ON budget_ledger (thread_id); ON budget_ledger (thread_id);
""".strip() """.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 = """ SCHEMA_META_DDL: str = """
CREATE TABLE IF NOT EXISTS schema_meta ( CREATE TABLE IF NOT EXISTS schema_meta (
id INTEGER PRIMARY KEY CHECK (id = 1), id INTEGER PRIMARY KEY CHECK (id = 1),
@ -180,7 +202,11 @@ def init_db(db_path: Path) -> None:
conn.execute(BUDGET_LEDGER_DDL) conn.execute(BUDGET_LEDGER_DDL)
for stmt in _split_statements(BUDGET_LEDGER_INDEXES_DDL): for stmt in _split_statements(BUDGET_LEDGER_INDEXES_DDL):
conn.execute(stmt) 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( conn.execute(
"INSERT INTO schema_meta (id, schema_version) VALUES (1, ?) " "INSERT INTO schema_meta (id, schema_version) VALUES (1, ?) "
"ON CONFLICT(id) DO NOTHING", "ON CONFLICT(id) DO NOTHING",
@ -213,7 +239,17 @@ def migrate(conn: sqlite3.Connection) -> None:
conn.execute(stmt) conn.execute(stmt)
current = 1 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) # Defense-in-depth index, applied UNCONDITIONALLY (idempotent IF NOT EXISTS)
# so an already-stamped v1 DB — which skips the `current < 1` block above — # 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( def answer_question(
conn: sqlite3.Connection, conn: sqlite3.Connection,
*, *,

View file

@ -64,6 +64,21 @@ CREATE INDEX IF NOT EXISTS idx_budget_ledger_day
CREATE INDEX IF NOT EXISTS idx_budget_ledger_thread CREATE INDEX IF NOT EXISTS idx_budget_ledger_thread
ON budget_ledger (thread_id); 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 -- schema_meta: single-row table recording the applied schema version so
-- migrate() can detect and step forward. -- migrate() can detect and step forward.
CREATE TABLE IF NOT EXISTS schema_meta ( CREATE TABLE IF NOT EXISTS schema_meta (

View file

@ -21,15 +21,26 @@ socket:
transport_name=...)``: the live :class:`~agent_team.coordinator.Coordinator` transport_name=...)``: the live :class:`~agent_team.coordinator.Coordinator`
in production, a stub in tests. No model or transport is touched here. in production, a stub in tests. No model or transport is touched here.
De-duplication (P3 scope note): De-duplication:
The poller tracks already-ingested issue ids in an **in-memory** set, so a The poller checks an injected :class:`IngestStore` before starting a task
re-poll over the same open issue does not start a second task. This is and records the issue id after ``start_task`` succeeds, so a re-poll over the
deliberately simple for now: it does NOT survive a process restart. Durable same open issue does not start a second task. Two stores ship:
de-dup (a ledger table of ingested issue ids, mirroring the
``pending_questions`` discipline) is a FOLLOW-UP and is intentionally not * :class:`_InMemoryIngestStore` (the default) — best-effort, per-process; it
shipped here. After a restart an already-ingested-but-still-open issue would does NOT survive a restart. Fine for tests and one-off operator runs.
be re-ingested; document that and treat the in-memory set as a best-effort * the **durable** ledger store from :func:`build_ledger_ingest_store` — an
guard, not a durable contract. ``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): Design constraints honoured here (pre-deployment scaffolding):
* **No live infrastructure.** Nothing is provisioned or called at import. * **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 from __future__ import annotations
import logging import logging
import re
from typing import Any, Iterable, Protocol from typing import Any, Iterable, Protocol
__all__ = [ __all__ = [
"GithubIntake", "GithubIntake",
"GithubIssueClient", "GithubIssueClient",
"IngestStore",
"build_default_issue_client", "build_default_issue_client",
"build_ledger_ingest_store",
"issue_task_text", "issue_task_text",
] ]
@ -67,6 +81,17 @@ GITHUB_TRANSPORT_NAME = "github"
# time (never stored in source/state). Mirrors the github_adapter default. # time (never stored in source/state). Mirrors the github_adapter default.
DEFAULT_TOKEN_ENV = "GITHUB_TOKEN" 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): class GithubIssueClient(Protocol):
"""Injected GitHub issue source: list open issues carrying a label. """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: def issue_task_text(issue: dict[str, Any]) -> str:
"""Render one issue's intake ``task_text`` from its title + body. """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 operator loop or a cron); each call lists the open labeled issues and starts
a task for every one not yet ingested. a task for every one not yet ingested.
De-dup is in-memory only (see the module docstring): the set of ingested De-dup is delegated to the injected :class:`IngestStore` (see the module
issue ids lives on the instance, so a re-poll within one process never docstring): the default in-memory store guards within one process; the
double-ingests, but a restart loses the set. Durable de-dup is a follow-up. 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 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 it is injected), so the whole poller is unit-testable with a fake client and
@ -142,8 +221,9 @@ class GithubIntake:
client: GithubIssueClient, client: GithubIssueClient,
coordinator: Any, coordinator: Any,
label: str, label: str,
store: IngestStore | None = 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: Args:
client: The injected issue source. Its ``list_open_issues`` is the 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 issues the client returns for this label are considered; an
empty label is rejected so a misconfiguration cannot ingest empty label is rejected so a misconfiguration cannot ingest
every open issue. 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: if not label:
raise ValueError( raise ValueError(
@ -164,10 +248,12 @@ class GithubIntake:
self._client = client self._client = client
self._coordinator = coordinator self._coordinator = coordinator
self._label = label self._label = label
# In-memory de-dup set (P3 scope: best-effort, NOT durable across a # De-dup memory (default: in-memory, best-effort, NOT durable across a
# restart; see the module docstring). Tracks issue ids already turned # restart; production passes the durable ledger store). Tracks issue ids
# into tasks so a re-poll does not double-ingest. # already turned into tasks so a re-poll does not double-ingest.
self._ingested: set[str] = set() self._store: IngestStore = (
store if store is not None else _InMemoryIngestStore()
)
@property @property
def label(self) -> str: def label(self) -> str:
@ -176,8 +262,8 @@ class GithubIntake:
@property @property
def ingested_ids(self) -> frozenset[str]: def ingested_ids(self) -> frozenset[str]:
"""A snapshot of the issue ids already ingested this process (read-only).""" """A snapshot of the ingested issue ids from the store (read-only)."""
return frozenset(self._ingested) return self._store.snapshot()
def poll_once(self) -> list[str]: def poll_once(self) -> list[str]:
"""List the labeled open issues and start a task for each new one. """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 1. Ask the injected client for the open issues carrying the configured
label (:meth:`GithubIssueClient.list_open_issues`). 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=<title+body>, ``coordinator.start_task(task_text=<title+body>,
transport_name="github")`` and record its id so a subsequent poll transport_name="github")``.
does not re-ingest it.
Issues already ingested this process are skipped (the in-memory de-dup), Claim-then-do (not record-after): the id is reserved BEFORE the
and any issue the client returns without the label is *not* expected non-idempotent ``start_task`` side effect, so a crash between starting
(the client filters by label) but is ignored defensively if present. the task and finishing the pass cannot re-spawn a duplicate on the next
An issue id is recorded as ingested ONLY after ``start_task`` returns, run. If ``start_task`` *raises*, the claim is RELEASED so a transient
so a failing intake leaves the issue eligible for retry on the next poll failure stays retryable (no silent drop). The only residual is a hard
rather than silently dropping it. 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 Returns the list of issue ids ingested on THIS pass (empty when nothing
new), so an operator loop can log/meter intake volume. new), so an operator loop can log/meter intake volume.
@ -204,7 +293,7 @@ class GithubIntake:
ingested_now: list[str] = [] ingested_now: list[str] = []
for issue in self._client.list_open_issues(label=self._label): for issue in self._client.list_open_issues(label=self._label):
issue_id = _issue_id(issue) 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) _LOG.debug("github-intake: issue %s already ingested; skip", issue_id)
continue continue
@ -214,15 +303,18 @@ class GithubIntake:
issue_id, issue_id,
self._label, self._label,
) )
# start_task is the committed coordinator intake entry; the resulting # The claim above reserved the id BEFORE this non-idempotent intake
# task's clarifier question-sets route over the GitHub adapter. Record # call (coordinator.start_task mints a fresh thread + an open question
# the id only after the call returns so a raise leaves the issue # row). On a raise, RELEASE the claim so the transient failure is
# eligible for retry on the next poll (no silent drop). # retryable on the next poll rather than a silent drop.
self._coordinator.start_task( try:
task_text=task_text, self._coordinator.start_task(
transport_name=GITHUB_TRANSPORT_NAME, task_text=task_text,
) transport_name=GITHUB_TRANSPORT_NAME,
self._ingested.add(issue_id) )
except Exception:
self._store.release(issue_id)
raise
ingested_now.append(issue_id) ingested_now.append(issue_id)
return ingested_now return ingested_now
@ -233,7 +325,6 @@ def build_default_issue_client(
owner: str, owner: str,
repo: str, repo: str,
token_env: str = DEFAULT_TOKEN_ENV, token_env: str = DEFAULT_TOKEN_ENV,
api_root: str = "https://api.github.com",
) -> GithubIssueClient: ) -> GithubIssueClient:
"""Build the production read-only issue client (deferred SDK/HTTP import). """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. 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 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 credential; it is never stored in source or logged.
overridable for GitHub Enterprise.
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 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 the issues endpoint with ``state=open&labels=<label>`` and returns the
@ -254,6 +352,12 @@ def build_default_issue_client(
already carries ``id`` / ``number`` / ``title`` / ``body``). It performs no already carries ``id`` / ``number`` / ``title`` / ``body``). It performs no
CI, OIDC, write, or GitHub-Actions call; it only reads issues. 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: class _RestIssueClient:
"""Stdlib-only GitHub REST issue lister (built lazily, no import-time HTTP).""" """Stdlib-only GitHub REST issue lister (built lazily, no import-time HTTP)."""
@ -262,7 +366,8 @@ def build_default_issue_client(
self._owner = owner self._owner = owner
self._repo = repo self._repo = repo
self._token_env = token_env self._token_env = token_env
self._api_root = api_root.rstrip("/") # 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]]: def list_open_issues(self, *, label: str) -> Iterable[dict[str, Any]]:
import json import json
@ -293,3 +398,60 @@ def build_default_issue_client(
return [item for item in data if "pull_request" not in item] return [item for item in data if "pull_request" not in item]
return _RestIssueClient() 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()

View file

@ -700,8 +700,9 @@ def _cmd_intake_github(args: argparse.Namespace, *, out: Any) -> int:
(:func:`~agent_team.transport.github_intake.build_default_issue_client`, (:func:`~agent_team.transport.github_intake.build_default_issue_client`,
deferred-import, reads ``GITHUB_TOKEN`` at call time) and runs ONE deferred-import, reads ``GITHUB_TOKEN`` at call time) and runs ONE
:meth:`~agent_team.transport.github_intake.GithubIntake.poll_once`. An :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 operator (or a cron) re-runs the command on a cadence; de-dup is DURABLE
per process, so each run is a single pass. (the ``ingested_issues`` ledger table), so re-running over a still-open
labeled issue never starts a second task — safe for a scheduled timer.
OPT-IN and INERT: this only reads labeled issues and calls the committed 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 coordinator intake entry — no CI, OIDC, git/patch apply, or GitHub-Actions
@ -711,16 +712,25 @@ def _cmd_intake_github(args: argparse.Namespace, *, out: Any) -> int:
from agent_team.transport.github_intake import ( from agent_team.transport.github_intake import (
GithubIntake, GithubIntake,
build_default_issue_client, build_default_issue_client,
build_ledger_ingest_store,
) )
coordinator = _build_coordinator(args) coordinator = _build_coordinator(args)
coordinator.setup() coordinator.setup()
client = build_default_issue_client(owner=args.owner, repo=args.repo) client = build_default_issue_client(owner=args.owner, repo=args.repo)
# Durable de-dup: a scheduled/cron intake runs as a fresh process each time,
# so the in-memory default would re-ingest every still-open labeled issue and
# spawn duplicate tasks. The ledger store keys on (source, issue_id) in the
# same SQLite DB the coordinator uses, making intake idempotent across runs.
store = build_ledger_ingest_store(
db_path=args.db, source=f"github:{args.owner}/{args.repo}"
)
intake = GithubIntake( intake = GithubIntake(
client=client, client=client,
coordinator=coordinator, coordinator=coordinator,
label=args.label, label=args.label,
store=store,
) )
ingested = intake.poll_once() ingested = intake.poll_once()
for issue_id in ingested: for issue_id in ingested:

View file

@ -9,6 +9,7 @@ filters by label), and the intake text is the issue title + body.
from __future__ import annotations from __future__ import annotations
from pathlib import Path
from typing import Any from typing import Any
import pytest import pytest
@ -16,6 +17,8 @@ import pytest
from agent_team.transport.github_intake import ( from agent_team.transport.github_intake import (
GITHUB_TRANSPORT_NAME, GITHUB_TRANSPORT_NAME,
GithubIntake, GithubIntake,
build_default_issue_client,
build_ledger_ingest_store,
issue_task_text, issue_task_text,
) )
@ -231,3 +234,159 @@ def test_failed_start_task_leaves_issue_eligible_for_retry() -> None:
ingested = intake.poll_once() ingested = intake.poll_once()
assert ingested == ["1"] assert ingested == ["1"]
assert coordinator.attempts == 2 assert coordinator.attempts == 2
# --------------------------------------------------------------------------- #
# Durable de-dup (the ledger-backed IngestStore)
# --------------------------------------------------------------------------- #
def test_durable_store_dedup_survives_a_fresh_intake_process(tmp_path: Path) -> None:
"""A new GithubIntake (= a new cron process) over the same ledger DB does not
re-ingest a still-open labeled issue — the bug a scheduled timer would hit
with the in-memory default."""
db = tmp_path / "ledger.sqlite"
source = "github:o/r"
# Process 1: ingests issue 1.
store1 = build_ledger_ingest_store(db_path=db, source=source)
coord1 = FakeCoordinator()
intake1 = GithubIntake(
client=FakeIssueClient([_issue(1)]),
coordinator=coord1,
label=INTAKE_LABEL,
store=store1,
)
assert intake1.poll_once() == ["1"]
assert len(coord1.calls) == 1
# Process 2 (restart): brand-new intake + store over the SAME DB; the issue
# is still open, but durable de-dup means it is NOT ingested again.
store2 = build_ledger_ingest_store(db_path=db, source=source)
coord2 = FakeCoordinator()
intake2 = GithubIntake(
client=FakeIssueClient([_issue(1)]),
coordinator=coord2,
label=INTAKE_LABEL,
store=store2,
)
assert intake2.poll_once() == []
assert coord2.calls == []
assert intake2.ingested_ids == frozenset({"1"})
def test_durable_store_namespaces_by_source(tmp_path: Path) -> None:
"""The same issue id under a different repo source is ingested independently."""
db = tmp_path / "ledger.sqlite"
store_a = build_ledger_ingest_store(db_path=db, source="github:o/repo-a")
coord_a = FakeCoordinator()
GithubIntake(
client=FakeIssueClient([_issue(1)]),
coordinator=coord_a,
label=INTAKE_LABEL,
store=store_a,
).poll_once()
assert len(coord_a.calls) == 1
# Different repo, same issue number → not de-duped against repo-a.
store_b = build_ledger_ingest_store(db_path=db, source="github:o/repo-b")
coord_b = FakeCoordinator()
ingested = GithubIntake(
client=FakeIssueClient([_issue(1)]),
coordinator=coord_b,
label=INTAKE_LABEL,
store=store_b,
).poll_once()
assert ingested == ["1"]
assert len(coord_b.calls) == 1
def test_build_ledger_ingest_store_rejects_empty_source(tmp_path: Path) -> None:
with pytest.raises(ValueError):
build_ledger_ingest_store(db_path=tmp_path / "x.sqlite", source="")
# --------------------------------------------------------------------------- #
# Claim-then-do durability (INTAKE-LOGIC-01 fix) + client hardening
# --------------------------------------------------------------------------- #
def test_durable_claim_without_release_prevents_duplicate(tmp_path: Path) -> None:
"""Claim-then-do: a claim recorded but never released (a crash around
start_task) must NOT re-ingest on a fresh process — no duplicate task."""
db = tmp_path / "ledger.sqlite"
source = "github:o/r"
# Simulate a poller that won the claim then 'crashed' before releasing.
store1 = build_ledger_ingest_store(db_path=db, source=source)
assert store1.claim("42") is True
assert store1.claim("42") is False # same store: re-claim loses
# Fresh process over the same DB: the claim persists, so poll skips it.
store2 = build_ledger_ingest_store(db_path=db, source=source)
coord = FakeCoordinator()
intake = GithubIntake(
client=FakeIssueClient([_issue(42)]),
coordinator=coord,
label=INTAKE_LABEL,
store=store2,
)
assert intake.poll_once() == []
assert coord.calls == []
def test_failed_start_task_releases_durable_claim_for_retry(tmp_path: Path) -> None:
"""A raising start_task releases the durable claim so a later poll retries."""
db = tmp_path / "ledger.sqlite"
source = "github:o/r"
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"
coord = FlakyCoordinator()
intake = GithubIntake(
client=FakeIssueClient([_issue(1)]),
coordinator=coord,
label=INTAKE_LABEL,
store=build_ledger_ingest_store(db_path=db, source=source),
)
with pytest.raises(RuntimeError):
intake.poll_once()
# The claim was released, so a fresh process sees the issue as unclaimed.
intake2 = GithubIntake(
client=FakeIssueClient([_issue(1)]),
coordinator=coord,
label=INTAKE_LABEL,
store=build_ledger_ingest_store(db_path=db, source=source),
)
assert intake2.poll_once() == ["1"]
assert coord.attempts == 2
def test_build_default_issue_client_rejects_unsafe_owner_repo() -> None:
"""owner/repo are validated before URL construction (path-splice / CWE-918)."""
with pytest.raises(ValueError):
build_default_issue_client(owner="../../..", repo="r")
with pytest.raises(ValueError):
build_default_issue_client(owner="o", repo="r/issues?labels=admin")
with pytest.raises(ValueError):
build_default_issue_client(owner="", repo="r")
# A normal owner/repo builds fine.
client = build_default_issue_client(
owner="Sea-Haven-Industries", repo="orchestrator"
)
assert hasattr(client, "list_open_issues")
def test_build_default_issue_client_rejects_api_root_kwarg() -> None:
"""No caller-overridable api_root (SSRF/host-redirect, mirrors ci_fetcher FIX-3)."""
with pytest.raises(TypeError):
build_default_issue_client(owner="o", repo="r", api_root="https://evil.example")

View file

@ -17,12 +17,21 @@ from agent_team.db.schema import (
connect, connect,
expire_question, expire_question,
init_db, init_db,
issue_already_ingested,
migrate, migrate,
record_issue_ingested,
reopen_question, reopen_question,
supersede_question, supersede_question,
) )
def _tables(conn: sqlite3.Connection) -> set[str]:
return {
row[0]
for row in conn.execute("SELECT name FROM sqlite_master WHERE type='table'")
}
def _insert_open_question(conn: sqlite3.Connection, qid: str, turn: int = 0) -> None: def _insert_open_question(conn: sqlite3.Connection, qid: str, turn: int = 0) -> None:
conn.execute( conn.execute(
"INSERT INTO pending_questions " "INSERT INTO pending_questions "
@ -463,3 +472,67 @@ def test_migrate_adds_open_channel_ref_partial_unique(tmp_path: Path) -> None:
_insert_open_question_with_ref(conn, "q2", "ts-300") _insert_open_question_with_ref(conn, "q2", "ts-300")
finally: finally:
conn.close() conn.close()
# --------------------------------------------------------------------------- #
# ingested_issues durable de-dup (schema v2)
# --------------------------------------------------------------------------- #
def test_schema_version_is_at_least_2() -> None:
assert SCHEMA_VERSION >= 2
def test_init_db_creates_ingested_issues(tmp_path: Path) -> None:
db = tmp_path / "db.sqlite"
init_db(db)
conn = connect(db)
try:
assert "ingested_issues" in _tables(conn)
finally:
conn.close()
def test_issue_ingest_helpers_record_and_detect(tmp_path: Path) -> None:
db = tmp_path / "db.sqlite"
init_db(db)
conn = connect(db)
try:
src = "github:o/r"
assert not issue_already_ingested(conn, source=src, issue_id="1")
# first record inserts a new row
assert record_issue_ingested(conn, source=src, issue_id="1") is True
assert issue_already_ingested(conn, source=src, issue_id="1")
# idempotent: a repeat record is a no-op (False) but still "seen"
assert record_issue_ingested(conn, source=src, issue_id="1") is False
assert issue_already_ingested(conn, source=src, issue_id="1")
# source namespacing: same id under a different repo is independent
assert not issue_already_ingested(conn, source="github:o/other", issue_id="1")
finally:
conn.close()
def test_migrate_adds_ingested_issues_to_a_legacy_v1_db(tmp_path: Path) -> None:
"""A DB stamped at v1 (no ingested_issues) gains the table + a v2 stamp."""
db = tmp_path / "legacy.sqlite"
conn = connect(db)
try:
# Simulate a legacy v1 DB: pending_questions + a schema_meta stamped at 1,
# WITHOUT the v2 ingested_issues table.
conn.execute(PENDING_QUESTIONS_DDL)
conn.execute(
"CREATE TABLE IF NOT EXISTS schema_meta "
"(id INTEGER PRIMARY KEY CHECK (id = 1), schema_version INTEGER NOT NULL)"
)
conn.execute("INSERT INTO schema_meta (id, schema_version) VALUES (1, 1)")
assert "ingested_issues" not in _tables(conn)
migrate(conn)
assert "ingested_issues" in _tables(conn)
ver = conn.execute(
"SELECT schema_version FROM schema_meta WHERE id = 1"
).fetchone()[0]
assert ver == SCHEMA_VERSION
finally:
conn.close()