295 lines
9.8 KiB
Python
295 lines
9.8 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
|
|
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",
|
|
"SCHEMA_META_DDL",
|
|
"SCHEMA_VERSION",
|
|
"QUESTION_STATES",
|
|
"answer_question",
|
|
"connect",
|
|
"expire_question",
|
|
"init_db",
|
|
"migrate",
|
|
"supersede_question",
|
|
]
|
|
|
|
# Bump when the DDL below changes; migrate() steps a connection forward.
|
|
SCHEMA_VERSION: int = 1
|
|
|
|
# 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);
|
|
""".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()
|
|
|
|
SCHEMA_META_DDL: str = """
|
|
CREATE TABLE IF NOT EXISTS schema_meta (
|
|
id INTEGER PRIMARY KEY CHECK (id = 1),
|
|
schema_version INTEGER NOT NULL
|
|
)
|
|
""".strip()
|
|
|
|
|
|
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,
|
|
)
|
|
conn.row_factory = sqlite3.Row
|
|
conn.execute("PRAGMA journal_mode=WAL")
|
|
conn.execute("PRAGMA foreign_keys=ON")
|
|
conn.execute(f"PRAGMA busy_timeout={_BUSY_TIMEOUT_MS}")
|
|
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)
|
|
# Record the schema version (single-row table).
|
|
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
|
|
|
|
# Future steps go here: `if current < 2: ...; current = 2`.
|
|
|
|
conn.execute(
|
|
"INSERT INTO schema_meta (id, schema_version) VALUES (1, ?) "
|
|
"ON CONFLICT(id) DO UPDATE SET schema_version = excluded.schema_version",
|
|
(current,),
|
|
)
|
|
|
|
|
|
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 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,),
|
|
)
|
|
|
|
|
|
def _compare_and_set(
|
|
conn: sqlite3.Connection,
|
|
*,
|
|
sql: str,
|
|
params: tuple[Any, ...],
|
|
) -> bool:
|
|
"""Run a single compare-and-set UPDATE under ``BEGIN IMMEDIATE``.
|
|
|
|
Returns ``True`` iff exactly one row changed. Takes the write lock up front
|
|
so concurrent responders cannot both observe ``status='open'`` (SQLite's
|
|
default deferred isolation would not serialize them).
|
|
"""
|
|
conn.execute("BEGIN IMMEDIATE")
|
|
try:
|
|
cur = conn.execute(sql, params)
|
|
changed = cur.rowcount == 1
|
|
conn.execute("COMMIT")
|
|
except BaseException:
|
|
conn.execute("ROLLBACK")
|
|
raise
|
|
return changed
|
|
|
|
|
|
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()]
|