diff --git a/agent-team/agent_team/db/schema.py b/agent-team/agent_team/db/schema.py index d2d0067..106def1 100644 --- a/agent-team/agent_team/db/schema.py +++ b/agent-team/agent_team/db/schema.py @@ -31,6 +31,7 @@ __all__ = [ "PENDING_QUESTIONS_DDL", "PENDING_QUESTIONS_INDEXES_DDL", "BUDGET_LEDGER_INDEXES_DDL", + "INGESTED_ISSUES_DDL", "SCHEMA_META_DDL", "SCHEMA_VERSION", "QUESTION_STATES", @@ -39,13 +40,15 @@ __all__ = [ "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 +108,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 +201,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 +238,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 +270,47 @@ 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 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..468a825 100644 --- a/agent-team/agent_team/transport/github_intake.py +++ b/agent-team/agent_team/transport/github_intake.py @@ -21,15 +21,23 @@ 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 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. Design constraints honoured here (pre-deployment scaffolding): * **No live infrastructure.** Nothing is provisioned or called at import. @@ -53,7 +61,9 @@ from typing import Any, Iterable, Protocol __all__ = [ "GithubIntake", "GithubIssueClient", + "IngestStore", "build_default_issue_client", + "build_ledger_ingest_store", "issue_task_text", ] @@ -83,6 +93,50 @@ 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). ``seen``/``mark`` mirror the + check-then-record discipline; ``snapshot`` exposes the current id set for + introspection/logging. + """ + + def seen(self, issue_id: str) -> bool: + """Return True if ``issue_id`` has already been ingested.""" + ... + + def mark(self, issue_id: str) -> None: + """Record ``issue_id`` as ingested (idempotent).""" + ... + + def snapshot(self) -> frozenset[str]: + """Return a read-only snapshot of the 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 seen(self, issue_id: str) -> bool: + return issue_id in self._ids + + def mark(self, issue_id: str) -> None: + self._ids.add(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 +181,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 +197,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 +211,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 +224,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 +238,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. @@ -204,7 +266,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 self._store.seen(issue_id): _LOG.debug("github-intake: issue %s already ingested; skip", issue_id) continue @@ -222,7 +284,7 @@ class GithubIntake: task_text=task_text, transport_name=GITHUB_TRANSPORT_NAME, ) - self._ingested.add(issue_id) + self._store.mark(issue_id) ingested_now.append(issue_id) return ingested_now @@ -293,3 +355,57 @@ def build_default_issue_client( 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 seen(self, issue_id: str) -> bool: + return _schema.issue_already_ingested( + self._conn, source=self._source, issue_id=issue_id + ) + + def mark(self, issue_id: str) -> None: + _schema.record_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() diff --git a/agent-team/run-team.py b/agent-team/run-team.py index 785949c..2986149 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -700,8 +700,9 @@ def _cmd_intake_github(args: argparse.Namespace, *, out: Any) -> int: (:func:`~agent_team.transport.github_intake.build_default_issue_client`, deferred-import, reads ``GITHUB_TOKEN`` at call time) and runs ONE :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 - per process, so each run is a single pass. + operator (or a cron) re-runs the command on a cadence; de-dup is DURABLE + (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 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 ( GithubIntake, build_default_issue_client, + build_ledger_ingest_store, ) coordinator = _build_coordinator(args) coordinator.setup() 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( client=client, coordinator=coordinator, label=args.label, + store=store, ) ingested = intake.poll_once() for issue_id in ingested: diff --git a/agent-team/tests/test_github_intake.py b/agent-team/tests/test_github_intake.py index dfe3fea..b041adc 100644 --- a/agent-team/tests/test_github_intake.py +++ b/agent-team/tests/test_github_intake.py @@ -9,6 +9,7 @@ filters by label), and the intake text is the issue title + body. from __future__ import annotations +from pathlib import Path from typing import Any import pytest @@ -16,6 +17,7 @@ import pytest from agent_team.transport.github_intake import ( GITHUB_TRANSPORT_NAME, GithubIntake, + build_ledger_ingest_store, issue_task_text, ) @@ -231,3 +233,74 @@ def test_failed_start_task_leaves_issue_eligible_for_retry() -> None: ingested = intake.poll_once() assert ingested == ["1"] 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="") diff --git a/agent-team/tests/test_schema.py b/agent-team/tests/test_schema.py index 52c6e61..1dca826 100644 --- a/agent-team/tests/test_schema.py +++ b/agent-team/tests/test_schema.py @@ -17,12 +17,21 @@ from agent_team.db.schema import ( connect, expire_question, init_db, + issue_already_ingested, migrate, + record_issue_ingested, reopen_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: conn.execute( "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") finally: 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()