This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
orchestrator/agent-team/agent_team/db/schema.py
Adam Moussa f916c03818
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.
2026-06-22 17:20:43 -04:00

572 lines
22 KiB
Python

"""SQLite schema, DDL constants, and connection helpers (design §3.3.1, §6.7).
This module is the single source of truth for the R720 agent-team durable
SQL. It declares:
* the ``pending_questions`` human-interaction lifecycle ledger (§3.3.1),
* the ``budget_ledger`` shared Claude budget ledger (§6.1, §6.6),
* a ``schema_meta`` version row driving :func:`migrate`.
The companion ``schema.sql`` mirrors these statements verbatim for tooling.
SQL DDL lives ONLY here. The LangGraph ``SqliteSaver`` checkpointer creates
its OWN tables against this same connection / database file; the design
reserves this DB for it but does not declare its tables.
The atomic compare-and-set helpers used by the responder and the deadline
timer (§3.3.1) take the write lock up front via ``BEGIN IMMEDIATE`` so
concurrent responders are serialized — SQLite's default deferred isolation
does not serialize a check-and-set.
"""
from __future__ import annotations
import sqlite3
import time
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
__all__ = [
"BUDGET_LEDGER_DDL",
"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 = 2
# Default SQLite busy timeout (ms) so concurrent writers wait for the write
# lock rather than failing immediately.
_BUSY_TIMEOUT_MS: int = 5000
# Allowed lifecycle states for a pending question (§3.3.1). Mirrors the DDL
# CHECK constraint; exported so leaves can validate without re-listing them.
QUESTION_STATES: tuple[str, ...] = ("open", "answered", "expired", "superseded")
PENDING_QUESTIONS_DDL: str = """
CREATE TABLE IF NOT EXISTS pending_questions (
question_id TEXT PRIMARY KEY,
thread_id TEXT NOT NULL,
turn INTEGER NOT NULL,
status TEXT NOT NULL
CHECK (status IN ('open', 'answered', 'expired', 'superseded')),
transport TEXT NOT NULL,
channel_ref TEXT,
posted_at TEXT,
deadline_at TEXT,
answer_json TEXT,
answered_at TEXT,
answered_via TEXT
)
""".strip()
PENDING_QUESTIONS_INDEXES_DDL: str = """
CREATE INDEX IF NOT EXISTS idx_pending_questions_thread
ON pending_questions (thread_id, turn);
CREATE INDEX IF NOT EXISTS idx_pending_questions_status
ON pending_questions (status);
CREATE UNIQUE INDEX IF NOT EXISTS uq_pending_questions_open_channel_ref
ON pending_questions (channel_ref)
WHERE channel_ref IS NOT NULL AND status = 'open';
""".strip()
BUDGET_LEDGER_DDL: str = """
CREATE TABLE IF NOT EXISTS budget_ledger (
entry_id INTEGER PRIMARY KEY AUTOINCREMENT,
thread_id TEXT,
stage TEXT,
model TEXT NOT NULL,
billing_mode TEXT NOT NULL,
input_tokens INTEGER NOT NULL DEFAULT 0,
output_tokens INTEGER NOT NULL DEFAULT 0,
usd_cost REAL NOT NULL DEFAULT 0.0,
recorded_at TEXT NOT NULL,
day_bucket TEXT NOT NULL
)
""".strip()
BUDGET_LEDGER_INDEXES_DDL: str = """
CREATE INDEX IF NOT EXISTS idx_budget_ledger_day
ON budget_ledger (day_bucket);
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),
schema_version INTEGER NOT NULL
)
""".strip()
class _Connection(sqlite3.Connection):
"""``sqlite3.Connection`` subclass that can carry its backing file path.
The base ``Connection`` has no ``__dict__``, so a path cannot be stashed on
it. This thin subclass (passed as ``factory=`` to :func:`sqlite3.connect`)
lets :func:`connect` record the db file for a thread-safe attribute lookup by
the compare-and-set, avoiding a ``PRAGMA`` on a connection shared across
threads.
"""
agent_team_db_path: str = ""
def connect(db_path: Path) -> sqlite3.Connection:
"""Open ``db_path`` with WAL, foreign keys, and a busy timeout.
WAL (``journal_mode=WAL``) lets the resume worker read while a responder
writes; ``foreign_keys=ON`` enforces referential integrity; the busy
timeout makes concurrent writers wait for the write lock instead of
failing. ``isolation_level=None`` puts the connection in autocommit mode so
the compare-and-set helpers can drive transactions explicitly with
``BEGIN IMMEDIATE`` (§3.3.1).
"""
db_path = Path(db_path)
db_path.parent.mkdir(parents=True, exist_ok=True)
conn = sqlite3.connect(
str(db_path),
isolation_level=None,
check_same_thread=False,
factory=_Connection,
)
conn.row_factory = sqlite3.Row
# Set the busy timeout FIRST so every subsequent statement — including the
# journal-mode pragma below, which briefly needs the write lock — waits for
# the lock instead of failing immediately when another connection is mid
# -write. (Without this, opening a connection under concurrent writers could
# raise "database is locked" before the timeout was ever applied.)
conn.execute(f"PRAGMA busy_timeout={_BUSY_TIMEOUT_MS}")
conn.execute("PRAGMA journal_mode=WAL")
conn.execute("PRAGMA foreign_keys=ON")
# Record the backing file path so the compare-and-set can derive it via a
# thread-safe attribute read instead of running a PRAGMA on a connection that
# callers share across threads (a sqlite3.Connection is not safe for
# concurrent use — even a read would corrupt its transaction state). Empty
# for an in-memory DB (no file to reopen on a second connection).
conn.agent_team_db_path = "" if str(db_path) == ":memory:" else str(db_path)
return conn
def init_db(db_path: Path) -> None:
"""Create the agent-team tables in ``db_path`` if absent.
Creates ``pending_questions`` (+ indexes), the budget ledger (+ indexes),
and the ``schema_meta`` version row, and reserves the same DB file for the
LangGraph ``SqliteSaver`` checkpointer (which creates its own tables on
first use against this connection). Idempotent: safe to call on every
startup.
"""
conn = connect(db_path)
try:
conn.execute(SCHEMA_META_DDL)
conn.execute(PENDING_QUESTIONS_DDL)
for stmt in _split_statements(PENDING_QUESTIONS_INDEXES_DDL):
conn.execute(stmt)
conn.execute(BUDGET_LEDGER_DDL)
for stmt in _split_statements(BUDGET_LEDGER_INDEXES_DDL):
conn.execute(stmt)
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",
(SCHEMA_VERSION,),
)
finally:
conn.close()
def migrate(conn: sqlite3.Connection) -> None:
"""Step ``conn``'s schema forward to :data:`SCHEMA_VERSION`.
Reads the recorded version from ``schema_meta`` (treating an empty/absent
row as version 0), applies any forward steps, and records the new version.
At ``SCHEMA_VERSION == 1`` there are no prior versions to migrate from, so
this ensures the base tables exist and stamps the version. Future versions
add ordered ``if current < N`` blocks here.
"""
conn.execute(SCHEMA_META_DDL)
row = conn.execute("SELECT schema_version FROM schema_meta WHERE id = 1").fetchone()
current = int(row["schema_version"]) if row is not None else 0
if current < 1:
# Base schema (v1): ensure all tables/indexes exist.
conn.execute(PENDING_QUESTIONS_DDL)
for stmt in _split_statements(PENDING_QUESTIONS_INDEXES_DDL):
conn.execute(stmt)
conn.execute(BUDGET_LEDGER_DDL)
for stmt in _split_statements(BUDGET_LEDGER_INDEXES_DDL):
conn.execute(stmt)
current = 1
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 —
# still gains the partial unique index on (channel_ref) WHERE open. This is a
# pure add-on guard (no version bump): two OPEN rows can never share a
# non-null channel_ref. Safe on existing data — NULL channel_refs are fine
# (the WHERE clause excludes them and SQLite treats NULLs as distinct), and
# there should be no existing duplicate non-null OPEN channel_refs (the
# responder records one channel_ref per posted open question, and recovery
# clears unposted ones). If a legacy DB *did* hold a duplicate, this CREATE
# would raise IntegrityError loudly rather than silently — flag, don't force.
for stmt in _split_statements(PENDING_QUESTIONS_INDEXES_DDL):
conn.execute(stmt)
conn.execute(
"INSERT INTO schema_meta (id, schema_version) VALUES (1, ?) "
"ON CONFLICT(id) DO UPDATE SET schema_version = excluded.schema_version",
(current,),
)
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,
*,
question_id: str,
answer_json: str,
answered_via: str,
answered_at: str | None = None,
) -> bool:
"""First-answer-wins compare-and-set: flip an ``open`` question to answered.
Runs the §3.3.1 atomic statement inside a ``BEGIN IMMEDIATE`` transaction
so concurrent responders are serialized (the check-and-set takes the write
lock up front). Returns ``True`` when rowcount == 1 (this caller recorded
the first valid answer; enqueue a resume job), ``False`` when rowcount == 0
(the question was not ``open`` — already answered/expired/superseded — so
the answer is a duplicate or late and must be ignored).
"""
stamp = answered_at or _utc_now_iso()
return _compare_and_set(
conn,
sql=(
"UPDATE pending_questions "
"SET status='answered', answer_json=?, answered_via=?, answered_at=? "
"WHERE question_id=? AND status='open'"
),
params=(answer_json, answered_via, stamp, question_id),
)
def expire_question(
conn: sqlite3.Connection,
*,
question_id: str,
) -> bool:
"""Deadline race: flip an overdue ``open`` question to ``expired``.
Same compare-and-set discipline as :func:`answer_question` (§3.3.1): an
answer that arrives for an already-expired question loses the race and is
ignored. Returns ``True`` if this call expired the question.
"""
return _compare_and_set(
conn,
sql=(
"UPDATE pending_questions SET status='expired' "
"WHERE question_id=? AND status='open'"
),
params=(question_id,),
)
def reopen_question(
conn: sqlite3.Connection,
*,
question_id: str,
deadline_at: str | None = None,
) -> bool:
"""Un-park: flip an ``expired`` question back to ``open`` (operator action).
The §6.6 operator force-resume path for a parked task whose clarifier
question expired with no answer: re-open it so the normal delivery → answer →
resume flow can proceed, instead of destructively superseding it (which would
remove it from the recovery sweep's reach). Same compare-and-set discipline —
only an ``expired`` row is reopened; an already-answered/open/superseded row
loses the CAS and is untouched. ``deadline_at`` sets a fresh window (``NULL``
means no deadline until one is set, so it will not immediately re-expire).
Returns ``True`` if this call reopened the question.
"""
return _compare_and_set(
conn,
sql=(
"UPDATE pending_questions "
"SET status='open', deadline_at=?, channel_ref=NULL, "
"answer_json=NULL, answered_via=NULL, answered_at=NULL "
"WHERE question_id=? AND status='expired'"
),
params=(deadline_at, question_id),
)
def find_open_question_by_channel_ref(
conn: sqlite3.Connection,
channel_ref: str,
) -> str | None:
"""Map a transport ``channel_ref`` to its still-``open`` ``question_id``.
The clarifier posts a question message and stores that message's ``ts`` as
the ledger row's ``channel_ref``; an inbound thread reply carries that same
value as its ``thread_ts``. When a reply carries no explicit
``callback_id`` / ``question_id`` / metadata (the real free-text-reply
shape), this resolves WHICH question the reply answers by its thread anchor.
The lookup is CONSTRAINED to ``status='open'`` (anti-replay): a stale or
replayed ``thread_ts`` pointing at a closed / expired / answered /
superseded row resolves to ``None`` and is a no-op for the caller. This
only resolves which question a reply targets; it is NEVER authorization —
the caller authorizes the sender first and fails closed.
Returns the ``question_id`` of the matching open row, or ``None`` if
``channel_ref`` is empty or matches no open row.
"""
if not channel_ref:
return None
row = conn.execute(
"SELECT question_id FROM pending_questions "
"WHERE channel_ref=? AND status='open'",
(channel_ref,),
).fetchone()
return None if row is None else str(row["question_id"])
def supersede_question(
conn: sqlite3.Connection,
*,
question_id: str,
) -> bool:
"""Mark a stale ``open``/``answered`` question ``superseded``.
Used by the turn-guarded resume worker: if the graph already advanced past
this turn, the question is superseded and the resume is skipped (§3.3.1).
Returns ``True`` if this call superseded the question.
"""
return _compare_and_set(
conn,
sql=(
"UPDATE pending_questions SET status='superseded' "
"WHERE question_id=? AND status IN ('open', 'answered')"
),
params=(question_id,),
)
# Bounded retry if the write lock is still contended after ``busy_timeout``
# elapses, so transient over-timeout contention does not surface as an error to
# the responder / deadline-timer callers.
_CAS_RETRY_ATTEMPTS: int = 3
_CAS_RETRY_BACKOFF_S: float = 0.05
def _main_db_file(conn: sqlite3.Connection) -> str | None:
"""Return the file backing ``conn``'s ``main`` database, or ``None``.
``None`` signals an in-memory database (no file to reopen on a second
connection). Prefers the path stashed by :func:`connect` — a thread-safe
attribute read, so it is safe even when callers share ``conn`` across
threads. Falls back to ``PRAGMA database_list`` (rows of ``(seq, name,
file)``, indexed positionally to be ``row_factory``-agnostic) only for a
connection not opened via :func:`connect`; such a connection must not be
shared across threads.
"""
stashed = getattr(conn, "agent_team_db_path", None)
if stashed is not None:
return stashed or None
for row in conn.execute("PRAGMA database_list"):
if row[1] == "main":
return row[2] or None
return None
def _compare_and_set(
conn: sqlite3.Connection,
*,
sql: str,
params: tuple[Any, ...],
) -> bool:
"""Run a single compare-and-set UPDATE under ``BEGIN IMMEDIATE`` (§3.3.1).
Returns ``True`` iff exactly one row changed. The check-and-set takes the
write lock up front so concurrent responders cannot both observe
``status='open'`` (SQLite's default deferred isolation would not serialize
them).
**Concurrency safety.** The write runs on a private, short-lived connection
to the same database file — never on the passed ``conn``. A single SQLite
connection cannot hold two explicit transactions at once, so if a caller
shares one ``conn`` across threads (the responder and resume worker do, and
``connect()`` sets ``check_same_thread=False``), two concurrent
``BEGIN IMMEDIATE`` statements on it would raise "cannot start a transaction
within a transaction". Giving each call its own connection makes the
compare-and-set safe under that sharing; WAL serializes the writers via the
busy handler. A lock that outlasts ``busy_timeout`` is retried a bounded
number of times before propagating. ``BEGIN IMMEDIATE`` runs inside the
guarded path so its lock error is caught and retried, not raised uncaught.
For an in-memory database (no file to reopen) the call falls back to the
passed ``conn``; in-memory DBs are single-connection and not the concurrent
production path.
"""
db_file = _main_db_file(conn)
if db_file is None:
return _cas_once(conn, sql, params)
last_err: sqlite3.OperationalError | None = None
for attempt in range(_CAS_RETRY_ATTEMPTS):
write = connect(Path(db_file))
try:
return _cas_once(write, sql, params)
except sqlite3.OperationalError as err:
if "locked" not in str(err).lower():
raise
last_err = err
finally:
write.close()
time.sleep(_CAS_RETRY_BACKOFF_S * (attempt + 1))
assert last_err is not None # loop only exits early via return or raise
raise last_err
def _cas_once(
conn: sqlite3.Connection,
sql: str,
params: tuple[Any, ...],
) -> bool:
"""Execute one ``BEGIN IMMEDIATE`` compare-and-set on ``conn``.
``BEGIN IMMEDIATE`` is issued before the try so a lock-acquisition error
propagates to the caller's retry loop with no transaction to unwind; once
the transaction is open, any failure rolls it back (best-effort) and
re-raises.
"""
conn.execute("BEGIN IMMEDIATE")
try:
cur = conn.execute(sql, params)
changed = cur.rowcount == 1
conn.execute("COMMIT")
return changed
except BaseException:
try:
conn.execute("ROLLBACK")
except sqlite3.OperationalError:
pass
raise
def _utc_now_iso() -> str:
"""Return the current UTC time as an ISO-8601 string."""
return datetime.now(timezone.utc).isoformat()
def _split_statements(ddl: str) -> list[str]:
"""Split a multi-statement DDL blob into individual statements."""
return [stmt.strip() for stmt in ddl.split(";") if stmt.strip()]