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 15a416d31a Add Plane-2 leaf scaffold (pipeline graph, nodes, HITL, transports, CI)
Consolidates the 18 leaf modules from the r720-plane2-scaffold workflow onto
the foundation commit. Full suite: 535 passed, 1 skipped; ruff + format clean.

Built (pre-deployment scaffold only — nothing provisioned/enabled):
- LangGraph pipeline graph.py (INTAKE->CLARIFY->PLAN, interrupt()/resume, checkpointer-injectable)
- nodes: clarifier (98% gate), planner, review_loop (GPT-4.1), builders->candidate diff, verifier
- §3.3.1 HITL: ledger ops, resume_worker, deadline_timer, recovery sweep, responder
- transports: slack / github / claude_code adapters
- ci_gate (pure-code pass/fail), operator_cli, run-team.py entry, P1 sim harness
- ci/agent-team-apply-verify.yml (split untrusted/privileged jobs) — authored, disabled

KNOWN OPEN FINDINGS (verifier/cross-review, not yet fixed — see follow-up):
- builders denylist: 4 execution-proven bypasses (delete, mode-change, copy-to, out-of-scope delete)
- §3.3.1 CAS: BEGIN IMMEDIATE outside try/except; shared-connection txn nesting unsafe under concurrency
- operator_cli: missing re-deliver/force-resume; audit-after-mutate ordering gap
- ci yaml: GPT-4.1 cross-review PASS w/ 4 FIX items (symlink path escape, etc.)
- P1 sim harness models the ledger layer, not real LangGraph interrupt/resume; P1 exit criteria not yet truly proven

Deploy-gated (NOT done): IAM/step-ca/Roles Anywhere/confluence-bot provisioning,
/sh-security-review sign-off, live Slack/CI, rsync, live dry-runs, Adam approval.
2026-06-17 15:16:12 -04:00

305 lines
11 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
@classmethod
def from_row(cls, row: sqlite3.Row) -> PendingQuestion:
"""Build a :class:`PendingQuestion` from a ``sqlite3.Row``."""
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"],
)
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,
) -> 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.
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) "
"VALUES (?, ?, ?, 'open', ?, ?, ?)",
(
question_id,
thread_id,
int(turn),
transport,
posted_at or _utc_now_iso(),
deadline_at,
),
)
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