feat(agent-team): durable GitHub-issue intake de-dup + intake hardening (#32)
* 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). * fix(agent-team): harden github intake per /sh-security-review (CWE-918, idempotency) Fixes from the high-recall detector fan-out on the durable-dedup change: - INTAKE-LOGIC-01 (idempotency): switch the IngestStore seam from check-then-record (seen/mark) to claim-then-do (claim/release). The id is now reserved BEFORE the non-idempotent start_task side effect, so a crash in that window cannot re-spawn a duplicate task on the next run; a raising start_task releases the claim so transient failures stay retryable. Adds delete_issue_ingested to the schema layer for the release path. - INTAKE-SSRF-001 / INTAKE-PATHSPLICE-002 (CWE-918) in build_default_issue_client: drop the caller-overridable api_root (hardcode GITHUB_API_ROOT) and validate owner/repo against an anchored charset before splicing them into the token-bearing API URL — mirrors the sibling ci_fetcher BLOCK-3/FIX-3 fixes. Tests cover cross-process duplicate prevention, release-on-failure retry, and the owner/repo + api_root rejection. 982 tests pass, ruff clean. Follow-up (pre-existing, not introduced here): the label-only intake has no author allowlist (cf. AGENT_TEAM_SLACK_OWNER_IDS on the Slack listener); the Slack answer gate bounds the blast radius. Track as separate hardening.
This commit is contained in:
parent
2f7d12a431
commit
f916c03818
6 changed files with 559 additions and 47 deletions
|
|
@ -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,
|
||||
*,
|
||||
|
|
|
|||
|
|
@ -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 (
|
||||
|
|
|
|||
|
|
@ -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=<title+body>,
|
||||
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=<label>`` and returns the
|
||||
|
|
@ -254,6 +352,12 @@ def build_default_issue_client(
|
|||
already carries ``id`` / ``number`` / ``title`` / ``body``). It performs no
|
||||
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:
|
||||
"""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._repo = repo
|
||||
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]]:
|
||||
import json
|
||||
|
|
@ -293,3 +398,60 @@ 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 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()
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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,8 @@ import pytest
|
|||
from agent_team.transport.github_intake import (
|
||||
GITHUB_TRANSPORT_NAME,
|
||||
GithubIntake,
|
||||
build_default_issue_client,
|
||||
build_ledger_ingest_store,
|
||||
issue_task_text,
|
||||
)
|
||||
|
||||
|
|
@ -231,3 +234,159 @@ 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="")
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# 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")
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
Reference in a new issue