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/ledger.py
Adam Moussa 1b4d30e47f feat(agent-team): add pending_questions.kind discriminator + migration (Phase B1)
The plan-review gate (coming next) needs to tell its decision questions apart
from clarifier questions in the durable ledger. Add a `kind` column to
pending_questions (values 'clarify' | 'plan_decision').

- Fresh DBs: `kind TEXT NOT NULL DEFAULT 'clarify'` (+ CHECK) in the DDL.
- Live ledger: idempotent additive migration (SCHEMA_VERSION 3→4) — a guarded
  ALTER (PRAGMA table_info) run from both migrate() and init_db; legacy rows
  take the 'clarify' default, never null. (SQLite can't add a CHECK via ALTER,
  so the migrated column is NOT NULL DEFAULT only; value constraint is enforced
  on fresh DBs by the CHECK and on all writes by the typed helper.)
- ledger.post_question gains a keyword-only `kind="clarify"` (backward
  compatible — existing callers unchanged); PendingQuestion.from_row reads it.

Tests: fresh-DB column+default, idempotent init_db, legacy-DB backfill to
'clarify', plan_decision round-trip. 1396 passed.
2026-06-23 21:03:53 -04:00

318 lines
12 KiB
Python

"""Pending-questions ledger ops — the §3.3.1 durable source of truth.
The LangGraph SQLite checkpointer suspends/resumes the graph, but the
checkpoint alone does not track the human-interaction lifecycle (delivery,
duplicate/late answers, expiry). This module is the thin, transport-agnostic
operations layer over the ``pending_questions`` table declared in
:mod:`agent_team.db.schema`. It is what delivery, the responder, the deadline
timer, the resume worker, restart recovery, and the manual CLI all read and
write.
Design (§3.3.1) mapping:
* **Delivery (and lost-post).** :func:`post_question` writes the row ``open``
*first* (no ``channel_ref``); the caller then posts to the transport and
records the ref via :func:`set_channel_ref`. If the post fails the row stays
``open`` with no ref and the reconcile loop (:func:`open_questions_needing_ref`)
retries idempotently.
* **Answer / expiry / supersede.** The atomic compare-and-set helpers live in
:mod:`agent_team.db.schema` and run under ``BEGIN IMMEDIATE``. They are
imported here verbatim and re-exported as the ledger's public mutation API
(:func:`answer_question`, :func:`expire_question`, :func:`supersede_question`)
so callers depend on one module — the SQL is **not** redefined here.
* **Deadline / no-answer.** :func:`overdue_open_questions` returns the ``open``
rows past their ``deadline_at`` for the timer loop to flip with
:func:`expire_question`.
* **Resume (turn-guarded).** :func:`answered_questions` feeds the resume worker
the ``answered`` rows whose graph may still be interrupted on that turn.
* **Restart recovery.** The startup sweep reads
:func:`open_questions_needing_ref` (retry delivery), :func:`answered_questions`
(re-enqueue resume), and :func:`overdue_open_questions` (run deadline policy).
All state is durable, so recovery is purely a function of the table.
* **Manual path (CLI).** :func:`list_questions` (optionally filtered by status)
and :func:`get_question` back the operator CLI that lists ``open`` questions,
re-delivers, force-expires, or answers on a task's behalf.
This module is pure stdlib + the foundation modules; it performs no transport
I/O of its own.
"""
from __future__ import annotations
import sqlite3
from dataclasses import dataclass
from datetime import datetime, timezone
# Compare-and-set primitives + lifecycle constants are imported VERBATIM from
# the committed foundation. They are NOT redefined here; the ledger re-exports
# them so callers depend on a single operations module.
from agent_team.db.schema import (
QUESTION_STATES,
answer_question,
expire_question,
supersede_question,
)
__all__ = [
"PendingQuestion",
"QUESTION_STATES",
"answer_question",
"answered_questions",
"count_by_status",
"expire_question",
"get_question",
"list_questions",
"open_questions_needing_ref",
"overdue_open_questions",
"post_question",
"set_channel_ref",
"supersede_question",
]
@dataclass(frozen=True)
class PendingQuestion:
"""A typed view of one ``pending_questions`` row (§3.3.1 ledger).
Mirrors the table columns declared in :data:`agent_team.db.schema.
PENDING_QUESTIONS_DDL`. ``frozen`` because rows are read snapshots; mutation
goes through the compare-and-set helpers, never by editing an instance.
"""
question_id: str
thread_id: str
turn: int
status: str
transport: str
channel_ref: str | None = None
posted_at: str | None = None
deadline_at: str | None = None
answer_json: str | None = None
answered_at: str | None = None
answered_via: str | None = None
kind: str = "clarify"
@classmethod
def from_row(cls, row: sqlite3.Row) -> PendingQuestion:
"""Build a :class:`PendingQuestion` from a ``sqlite3.Row``."""
keys = row.keys()
return cls(
question_id=row["question_id"],
thread_id=row["thread_id"],
turn=int(row["turn"]),
status=row["status"],
transport=row["transport"],
channel_ref=row["channel_ref"],
posted_at=row["posted_at"],
deadline_at=row["deadline_at"],
answer_json=row["answer_json"],
answered_at=row["answered_at"],
answered_via=row["answered_via"],
# Tolerate a row read before the kind column exists (legacy/partial
# SELECT): fall back to the 'clarify' default rather than KeyError.
kind=row["kind"] if "kind" in keys else "clarify",
)
def _utc_now_iso() -> str:
"""Return the current UTC time as an ISO-8601 string."""
return datetime.now(timezone.utc).isoformat()
def post_question(
conn: sqlite3.Connection,
*,
question_id: str,
thread_id: str,
turn: int,
transport: str,
deadline_at: str | None = None,
posted_at: str | None = None,
kind: str = "clarify",
) -> None:
"""Insert a new ``open`` question row (delivery step 1 of §3.3.1).
Writes the row ``open`` with **no** ``channel_ref`` *before* the caller
posts to the transport. If the subsequent post fails, the row stays ``open``
with no ref and the reconcile loop (:func:`open_questions_needing_ref`)
retries delivery idempotently. The caller records the ref via
:func:`set_channel_ref` once the post succeeds.
``kind`` discriminates the human gate this question belongs to —
``'clarify'`` (the clarifier, the default so existing callers are unchanged)
or ``'plan_decision'`` (the plan-review gate). The keyword-only default keeps
every existing call site writing clarifier rows with no signature change.
Raises :class:`sqlite3.IntegrityError` if ``question_id`` already exists
(PK) — re-posting the same question is the caller's reconcile concern, not a
silent overwrite. ``posted_at`` defaults to now (UTC ISO-8601).
"""
conn.execute(
"INSERT INTO pending_questions "
"(question_id, thread_id, turn, status, transport, posted_at, "
"deadline_at, kind) "
"VALUES (?, ?, ?, 'open', ?, ?, ?, ?)",
(
question_id,
thread_id,
int(turn),
transport,
posted_at or _utc_now_iso(),
deadline_at,
kind,
),
)
def set_channel_ref(
conn: sqlite3.Connection,
*,
question_id: str,
channel_ref: str,
) -> bool:
"""Record the transport ``channel_ref`` after a successful post (§3.3.1).
Sets the Slack message ts / GitHub issue-comment id / Claude session id on
a still-``open`` row. Guarded on ``status='open'`` so a late post-confirm
cannot resurrect a ref on an already-answered/expired/superseded question.
Returns ``True`` if exactly one open row was updated, ``False`` otherwise
(unknown id, or no longer open) — letting the reconcile loop decide whether
to retry.
"""
cur = conn.execute(
"UPDATE pending_questions SET channel_ref=? "
"WHERE question_id=? AND status='open'",
(channel_ref, question_id),
)
return cur.rowcount == 1
def get_question(
conn: sqlite3.Connection,
question_id: str,
) -> PendingQuestion | None:
"""Return one question by id, or ``None`` if absent."""
row = conn.execute(
"SELECT * FROM pending_questions WHERE question_id=?",
(question_id,),
).fetchone()
return PendingQuestion.from_row(row) if row is not None else None
def list_questions(
conn: sqlite3.Connection,
*,
status: str | None = None,
thread_id: str | None = None,
) -> list[PendingQuestion]:
"""List questions, optionally filtered by ``status`` and/or ``thread_id``.
Backs the manual CLI's ``list`` view (§3.3.1). ``status`` must be one of
:data:`QUESTION_STATES` when given; an unknown status raises
:class:`ValueError` rather than silently returning nothing. Results are
ordered oldest-first by ``posted_at`` so the operator sees the longest-open
questions first.
"""
if status is not None and status not in QUESTION_STATES:
raise ValueError(
f"unknown status {status!r}; expected one of {QUESTION_STATES}"
)
clauses: list[str] = []
params: list[object] = []
if status is not None:
clauses.append("status=?")
params.append(status)
if thread_id is not None:
clauses.append("thread_id=?")
params.append(thread_id)
where = f" WHERE {' AND '.join(clauses)}" if clauses else ""
rows = conn.execute(
"SELECT * FROM pending_questions"
f"{where} ORDER BY posted_at ASC, question_id ASC",
params,
).fetchall()
return [PendingQuestion.from_row(r) for r in rows]
def open_questions_needing_ref(
conn: sqlite3.Connection,
) -> list[PendingQuestion]:
"""Return ``open`` rows that have no ``channel_ref`` (lost-post reconcile).
The reconcile/recovery sweep (§3.3.1) retries delivery for these
idempotently: an ``open`` row with no ref means the row was written but the
transport post never confirmed.
"""
rows = conn.execute(
"SELECT * FROM pending_questions "
"WHERE status='open' AND channel_ref IS NULL "
"ORDER BY posted_at ASC, question_id ASC"
).fetchall()
return [PendingQuestion.from_row(r) for r in rows]
def overdue_open_questions(
conn: sqlite3.Connection,
*,
now: str | None = None,
) -> list[PendingQuestion]:
"""Return ``open`` rows whose ``deadline_at`` has passed (deadline loop).
Feeds the §3.3.1 timer loop, which flips each returned row with
:func:`expire_question` (the same compare-and-set, so an answer racing the
timer is resolved deterministically). Rows with a NULL ``deadline_at`` never
expire and are excluded. ``now`` defaults to the current UTC ISO-8601 time;
deadlines are compared as ISO-8601 strings, which sort lexicographically iff
written in a consistent UTC offset (the ledger always writes UTC).
"""
cutoff = now or _utc_now_iso()
rows = conn.execute(
"SELECT * FROM pending_questions "
"WHERE status='open' AND deadline_at IS NOT NULL AND deadline_at <= ? "
"ORDER BY deadline_at ASC, question_id ASC",
(cutoff,),
).fetchall()
return [PendingQuestion.from_row(r) for r in rows]
def answered_questions(
conn: sqlite3.Connection,
*,
thread_id: str | None = None,
) -> list[PendingQuestion]:
"""Return ``answered`` rows (resume-worker / restart-recovery feed).
The turn-guarded resume worker (§3.3.1) re-enqueues a resume for each
``answered`` row whose graph is still interrupted on that turn; the guard
makes a redelivered job idempotent. Optionally scoped to one ``thread_id``
(the resume worker serializes per thread). Ordered by ``answered_at`` so the
oldest pending resume is handled first.
"""
clauses = ["status='answered'"]
params: list[object] = []
if thread_id is not None:
clauses.append("thread_id=?")
params.append(thread_id)
rows = conn.execute(
"SELECT * FROM pending_questions "
f"WHERE {' AND '.join(clauses)} "
"ORDER BY answered_at ASC, question_id ASC",
params,
).fetchall()
return [PendingQuestion.from_row(r) for r in rows]
def count_by_status(conn: sqlite3.Connection) -> dict[str, int]:
"""Return a ``{status: count}`` map over all :data:`QUESTION_STATES`.
Backs CLI/telemetry summaries. Every state in :data:`QUESTION_STATES` is
present in the result (zero when absent) so callers get a stable shape.
"""
counts = dict.fromkeys(QUESTION_STATES, 0)
for row in conn.execute(
"SELECT status, COUNT(*) AS n FROM pending_questions GROUP BY status"
).fetchall():
counts[row["status"]] = int(row["n"])
return counts