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.
305 lines
11 KiB
Python
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
|