Add a Confluence documentation lane to the Plane-2 pipeline, flag-gated behind AGENT_TEAM_CONFLUENCE_ENABLED (default off; daemon behavior unchanged when off). - confluence/client.py: OAuth 2LO + Basic REST client, dry-run-default writes - confluence/mermaid.py: vendored ADF-only Mermaid editor (macro-count + revert-diff guards, dry-run default) - nodes/confluence_writer.py(+_llm): conf_draft -> conf_gate -> conf_write, both direct (task_kind=confluence) and post-build documentation flows - task_model/graph/coordinator: new phases, state channels, route_after_intake, CONFLUENCE_APPROVAL_KIND gate delivery, task_kind forwarding - db schema v5: widen pending_questions kind CHECK (atomic rebuild) - tests for client, mermaid, writer node, ledger v5, coordinator gate, e2e
820 lines
34 KiB
Python
820 lines
34 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",
|
|
"TASK_TRANSITIONS_DDL",
|
|
"TASK_TRANSITIONS_INDEXES_DDL",
|
|
"SCHEMA_META_DDL",
|
|
"SCHEMA_VERSION",
|
|
"QUESTION_STATES",
|
|
"answer_question",
|
|
"assert_task_transitions_ready",
|
|
"connect",
|
|
"delete_issue_ingested",
|
|
"expire_question",
|
|
"find_open_question_by_channel_ref",
|
|
"find_open_question_kind_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 = 5
|
|
|
|
# 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,
|
|
kind TEXT NOT NULL DEFAULT 'clarify'
|
|
CHECK (kind IN ('clarify', 'plan_decision', 'confluence_approval'))
|
|
)
|
|
""".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()
|
|
|
|
# task_transitions: per-task pipeline history (schema v3). One row per node a
|
|
# task enters, recording the move from_phase -> to_phase with entry/exit
|
|
# timestamps and the task status at entry. The dashboard's /api/task drill-down
|
|
# reads this to render a task's journey through the pipeline (timestamps,
|
|
# per-node duration); per-stage cost is joined from budget_ledger at read time,
|
|
# not duplicated here. Written by the coordinator's instrumented graph nodes
|
|
# (db/transitions.py); read READ-ONLY by the status dashboard. ``exited_at`` is
|
|
# filled when the next transition lands OR when the task reaches a terminal
|
|
# status (the recorder's close_terminal). A single OPEN row per thread is the
|
|
# invariant the recorder maintains (close-open-before-insert), so a crash/resume
|
|
# replay cannot leave orphaned open rows.
|
|
TASK_TRANSITIONS_DDL: str = """
|
|
CREATE TABLE IF NOT EXISTS task_transitions (
|
|
transition_id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
thread_id TEXT NOT NULL,
|
|
from_phase TEXT,
|
|
to_phase TEXT NOT NULL,
|
|
entered_at TEXT NOT NULL,
|
|
exited_at TEXT,
|
|
status TEXT,
|
|
note TEXT
|
|
)
|
|
""".strip()
|
|
|
|
TASK_TRANSITIONS_INDEXES_DDL: str = """
|
|
CREATE INDEX IF NOT EXISTS idx_task_transitions_thread
|
|
ON task_transitions (thread_id, entered_at);
|
|
""".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 _pending_questions_has_kind(conn: sqlite3.Connection) -> bool:
|
|
"""Return True if ``pending_questions`` already has the ``kind`` column.
|
|
|
|
Inspects ``PRAGMA table_info(pending_questions)`` so the additive ``kind``
|
|
migration can be applied only when absent — making it idempotent on an
|
|
already-migrated (or freshly created) DB.
|
|
"""
|
|
rows = conn.execute("PRAGMA table_info(pending_questions)").fetchall()
|
|
return any(row["name"] == "kind" for row in rows)
|
|
|
|
|
|
def _ensure_pending_questions_kind(conn: sqlite3.Connection) -> None:
|
|
"""Idempotently add the ``kind`` discriminator column to ``pending_questions``.
|
|
|
|
Fresh DBs get ``kind`` from :data:`PENDING_QUESTIONS_DDL`; an existing (live
|
|
R720) ledger whose table predates the column gets it via an in-place additive
|
|
``ALTER TABLE``, guarded by :func:`_pending_questions_has_kind` so a second
|
|
run is a no-op. Existing rows take the ``'clarify'`` default. SQLite cannot
|
|
add a CHECK constraint via ALTER, so the added column carries only the
|
|
NOT NULL DEFAULT; the CHECK is enforced on fresh DBs via the CREATE DDL and
|
|
on writes via the typed insert helper.
|
|
"""
|
|
if not _pending_questions_has_kind(conn):
|
|
conn.execute(
|
|
"ALTER TABLE pending_questions "
|
|
"ADD COLUMN kind TEXT NOT NULL DEFAULT 'clarify'"
|
|
)
|
|
|
|
|
|
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)
|
|
# Additive in-place migration for an existing ledger whose
|
|
# pending_questions predates the ``kind`` column (CREATE IF NOT EXISTS
|
|
# above never alters an existing table). No-op on fresh/already-migrated.
|
|
_ensure_pending_questions_kind(conn)
|
|
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)
|
|
conn.execute(TASK_TRANSITIONS_DDL)
|
|
for stmt in _split_statements(TASK_TRANSITIONS_INDEXES_DDL):
|
|
conn.execute(stmt)
|
|
# Step the version stamp forward AND apply any version-gated migrations.
|
|
# init_db is the only schema entry point the daemon calls (Coordinator.
|
|
# setup -> init_db), so it MUST drive migrate() — otherwise an existing
|
|
# DB's schema_version is never advanced (migrate() upserts it; init_db's
|
|
# own writes do not) and version-gated steps in migrate() never run in
|
|
# production. migrate() is idempotent, so re-running the create/ensure
|
|
# statements above is harmless.
|
|
migrate(conn)
|
|
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
|
|
|
|
if current < 3:
|
|
# v3: per-task pipeline transition history (dashboard drill-down).
|
|
conn.execute(TASK_TRANSITIONS_DDL)
|
|
for stmt in _split_statements(TASK_TRANSITIONS_INDEXES_DDL):
|
|
conn.execute(stmt)
|
|
current = 3
|
|
|
|
if current < 4:
|
|
# v4: add the ``kind`` discriminator to pending_questions so the
|
|
# responder/resume layer can tell a clarifier question apart from a
|
|
# plan-review decision. Additive in-place ALTER (guarded), placed in its
|
|
# own version block ABOVE the unconditional tail per the ORDERING
|
|
# CONSTRAINT below — the tail's CREATE ... IF NOT EXISTS would NOT apply
|
|
# this alter. Existing rows take the 'clarify' default.
|
|
_ensure_pending_questions_kind(conn)
|
|
current = 4
|
|
|
|
if current < 5:
|
|
# v5: WIDEN the pending_questions ``kind`` CHECK from
|
|
# ('clarify','plan_decision') to add 'confluence_approval', so the new
|
|
# Confluence-approval human gate can persist its question rows. SQLite
|
|
# cannot ALTER a CHECK constraint in place, so this does the standard
|
|
# 12-step table rebuild: build a replacement table carrying the WIDENED
|
|
# CHECK, copy every existing row across, drop the old table, rename the
|
|
# new one into place, and recreate ALL of pending_questions' indexes
|
|
# (including the partial unique uq_pending_questions_open_channel_ref).
|
|
# Existing rows are preserved verbatim; re-running is safe because the
|
|
# version stamp gates this block (and the unconditional index tail below
|
|
# re-asserts the indexes idempotently). Placed in its own version block
|
|
# ABOVE the unconditional tail per the ORDERING CONSTRAINT below — the
|
|
# tail's CREATE ... IF NOT EXISTS would NOT rebuild an existing table.
|
|
#
|
|
# foreign_keys must be OFF during the table rebuild so the DROP/RENAME
|
|
# does not trip referential checks; connect() sets it ON. It is a no-op
|
|
# for this table (no FKs reference pending_questions) but follows the
|
|
# documented SQLite procedure. It is restored to ON afterwards.
|
|
#
|
|
# The whole DROP/RENAME swap is wrapped in an explicit BEGIN IMMEDIATE
|
|
# transaction so it is ATOMIC: connect() opens in autocommit mode
|
|
# (isolation_level=None), so without this each statement would commit on
|
|
# its own and a crash between DROP TABLE pending_questions and the RENAME
|
|
# would DESTROY the live ledger while schema_version stayed 4 — on restart
|
|
# the block re-runs and CREATE pending_questions_new fails ('already
|
|
# exists'), crash-looping with data loss. Wrapping makes a mid-rebuild
|
|
# failure roll back to the original table intact. PRAGMA foreign_keys
|
|
# cannot change inside a transaction, so it is toggled OUTSIDE the BEGIN.
|
|
# DROP ... _new IF EXISTS first clears any orphan left by a prior partial
|
|
# run (belt-and-suspenders; the transaction already prevents one).
|
|
conn.execute("PRAGMA foreign_keys=OFF")
|
|
try:
|
|
conn.execute("BEGIN IMMEDIATE")
|
|
try:
|
|
conn.execute("DROP TABLE IF EXISTS pending_questions_new")
|
|
conn.execute(
|
|
"""
|
|
CREATE TABLE pending_questions_new (
|
|
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,
|
|
kind TEXT NOT NULL DEFAULT 'clarify'
|
|
CHECK (kind IN ('clarify', 'plan_decision', 'confluence_approval'))
|
|
)
|
|
""".strip()
|
|
)
|
|
# EXPLICIT column list (not SELECT *) so the copy survives a
|
|
# future column reorder/add on either table.
|
|
conn.execute(
|
|
"INSERT INTO pending_questions_new "
|
|
"(question_id, thread_id, turn, status, transport, "
|
|
"channel_ref, posted_at, deadline_at, answer_json, "
|
|
"answered_at, answered_via, kind) "
|
|
"SELECT question_id, thread_id, turn, status, transport, "
|
|
"channel_ref, posted_at, deadline_at, answer_json, "
|
|
"answered_at, answered_via, kind FROM pending_questions"
|
|
)
|
|
conn.execute("DROP TABLE pending_questions")
|
|
conn.execute(
|
|
"ALTER TABLE pending_questions_new RENAME TO pending_questions"
|
|
)
|
|
for stmt in _split_statements(PENDING_QUESTIONS_INDEXES_DDL):
|
|
conn.execute(stmt)
|
|
conn.execute("COMMIT")
|
|
except Exception:
|
|
conn.execute("ROLLBACK")
|
|
raise
|
|
finally:
|
|
conn.execute("PRAGMA foreign_keys=ON")
|
|
current = 5
|
|
|
|
# Future steps go here: `if current < 6: ...; current = 6`.
|
|
|
|
# Applied UNCONDITIONALLY (idempotent IF NOT EXISTS) so an already-stamped DB
|
|
# — which skips the version blocks above — still gains these tables without a
|
|
# restamp. Safe on existing data: fresh empty tables / indexes only.
|
|
#
|
|
# ORDERING CONSTRAINT: a future migration that ALTERs one of these tables
|
|
# (e.g. `ALTER TABLE task_transitions ADD COLUMN ...`) MUST run in its own
|
|
# `if current < N` block placed ABOVE this tail — the unconditional CREATE
|
|
# ... IF NOT EXISTS here no-ops on an existing table and will NOT apply an
|
|
# alter. This tail is only for first-time creation on an already-stamped DB.
|
|
conn.execute(INGESTED_ISSUES_DDL)
|
|
conn.execute(TASK_TRANSITIONS_DDL)
|
|
for stmt in _split_statements(TASK_TRANSITIONS_INDEXES_DDL):
|
|
conn.execute(stmt)
|
|
|
|
# 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,),
|
|
)
|
|
|
|
|
|
# Columns task_transitions MUST expose for the dashboard drill-down to work.
|
|
# assert_task_transitions_ready checks these so a botched migration is caught at
|
|
# startup (loud) rather than surfacing as a half-broken /api/task at read time.
|
|
_TASK_TRANSITIONS_COLUMNS: frozenset[str] = frozenset(
|
|
{
|
|
"transition_id",
|
|
"thread_id",
|
|
"from_phase",
|
|
"to_phase",
|
|
"entered_at",
|
|
"exited_at",
|
|
"status",
|
|
"note",
|
|
}
|
|
)
|
|
|
|
|
|
def assert_task_transitions_ready(db_path: Path) -> None:
|
|
"""Assert the v3 ``task_transitions`` table exists with the expected columns.
|
|
|
|
The deploy convention (and BLOCK-2 in the WebUI-makeover plan) requires the
|
|
migration to be verified *before* declaring a deploy good — a missing or
|
|
malformed table must fail loudly at startup, not after a crash-loop or as a
|
|
silently-broken drill-down. Raises :class:`RuntimeError` if the table is
|
|
absent or any expected column is missing; returns ``None`` on success.
|
|
"""
|
|
conn = connect(db_path)
|
|
try:
|
|
rows = conn.execute("PRAGMA table_info(task_transitions)").fetchall()
|
|
finally:
|
|
conn.close()
|
|
if not rows:
|
|
raise RuntimeError(
|
|
f"task_transitions table missing in {db_path} after migration "
|
|
"(schema v3) — aborting; restore the ledger backup and re-deploy."
|
|
)
|
|
present = {row["name"] for row in rows}
|
|
missing = _TASK_TRANSITIONS_COLUMNS - present
|
|
if missing:
|
|
raise RuntimeError(
|
|
f"task_transitions in {db_path} is malformed — missing columns "
|
|
f"{sorted(missing)}; aborting, restore the ledger backup."
|
|
)
|
|
|
|
|
|
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 find_open_question_kind_by_channel_ref(
|
|
conn: sqlite3.Connection,
|
|
channel_ref: str,
|
|
) -> tuple[str, str] | None:
|
|
"""Map a transport ``channel_ref`` to its open ``(question_id, kind)``.
|
|
|
|
Like :func:`find_open_question_by_channel_ref` but also returns the ``kind``
|
|
column so the caller can route ``plan_decision`` rows through the
|
|
kind-aware decision normalizer without a second query.
|
|
|
|
Returns ``(question_id, kind)`` for the matching open row, or ``None`` if
|
|
``channel_ref`` is empty or matches no open row. The lookup is constrained
|
|
to ``status='open'`` (anti-replay) — same semantics as the parent function.
|
|
"""
|
|
if not channel_ref:
|
|
return None
|
|
row = conn.execute(
|
|
"SELECT question_id, kind FROM pending_questions "
|
|
"WHERE channel_ref=? AND status='open'",
|
|
(channel_ref,),
|
|
).fetchone()
|
|
if row is None:
|
|
return None
|
|
return str(row["question_id"]), str(row["kind"])
|
|
|
|
|
|
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()]
|