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).
This commit is contained in:
Adam Moussa 2026-06-22 16:58:45 -04:00
parent 2f7d12a431
commit 3c8317b06e
6 changed files with 389 additions and 26 deletions

View file

@ -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,
*,

View file

@ -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 (

View file

@ -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()

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`,
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:

View file

@ -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="")

View file

@ -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()