#!/usr/bin/env python3 """``run-team.py`` — R720 agent-team operator entry CLI (design §3.3.1, §7.1 P1). This is the **entry CLI** named in the design (§2 "entry CLI ``run-team.py``"; §9 "operator CLI"). It is the small manual path over the durable ``pending_questions`` ledger that §3.3.1 ("Manual path") requires:: A small CLI over the ledger lets an operator list ``open``/``parked`` questions, re-deliver, force-expire, or answer on a task's behalf; a stuck task parks rather than spins. Destructive CLI actions (force-expire, answer-on-behalf, force-resume) are audit-logged and require an explicit confirmation flag. It imports the committed FOUNDATION contracts verbatim — it does not redefine them: * :mod:`agent_team.db.schema` — :func:`connect`, :func:`init_db`, :func:`answer_question`, :func:`expire_question`, :func:`supersede_question`, :data:`QUESTION_STATES`. * :mod:`agent_team.state_store` — :func:`atomic_write` for the append-only, crash-safe audit log of destructive actions (§6.7 discipline). Per the build constraints this is **pre-deployment scaffolding**: it provisions nothing, enables no live CI, and performs no network or rsync. It only reads and mutates the local SQLite ledger and writes a local audit log. Subcommands (P1 surface): * ``init-db`` — create/upgrade the agent-team tables in the ledger DB (idempotent; wraps :func:`init_db`). * ``list`` — list ``open`` (default) or any-status pending questions; with ``--parked`` it lists questions whose status is read as parked context. Pure read; no confirmation needed. * ``show`` — print one question row by ``question_id``. Pure read. * ``expire`` — force-expire an ``open`` question (DESTRUCTIVE: requires ``--confirm``; audit-logged). Maps to :func:`expire_question`. * ``answer`` — answer a question on a task's behalf (DESTRUCTIVE: requires ``--confirm``; audit-logged). Maps to :func:`answer_question`. * ``supersede`` — mark a stale question ``superseded`` (DESTRUCTIVE: requires ``--confirm``; audit-logged). Maps to :func:`supersede_question`. Exit codes: ``0`` success, ``1`` operational failure (e.g. row not found, the compare-and-set lost the race), ``2`` usage error (argparse). """ from __future__ import annotations import argparse import getpass import json import os import sqlite3 import sys from datetime import datetime, timezone from pathlib import Path from typing import Any, Sequence # ``run-team.py`` lives in ``agent-team/`` next to the importable ``agent_team`` # package. The hyphenated filename cannot itself be imported, so when run as a # script we make the sibling package importable without an editable install # (mirrors tests/conftest.py). _CLI_DIR = Path(__file__).resolve().parent if str(_CLI_DIR) not in sys.path: sys.path.insert(0, str(_CLI_DIR)) from agent_team.db.schema import ( # noqa: E402 (path bootstrap must precede) QUESTION_STATES, answer_question, connect, expire_question, init_db, reopen_question, supersede_question, ) __all__ = [ "build_parser", "main", ] # Default ledger DB location. Kept out of the repo (the package .gitignore # excludes ``state/`` and ``*.sqlite``) so durable state is never committed. _DEFAULT_DB = _CLI_DIR / "state" / "agent_team.sqlite" # Default audit log for destructive actions, alongside the ledger DB. _DEFAULT_AUDIT_LOG = _CLI_DIR / "state" / "audit.log.jsonl" # Columns selected for list/show rendering, in display order. _QUESTION_COLUMNS: tuple[str, ...] = ( "question_id", "thread_id", "turn", "status", "transport", "channel_ref", "posted_at", "deadline_at", "answered_at", "answered_via", ) # Destructive subcommands that require ``--confirm`` and are audit-logged. # ``force-resume`` is the design-named operator verb (§3.3.1/§6.6); ``supersede`` # is kept as its lower-level alias. ``redeliver`` is NOT here — it is idempotent # and non-destructive (it only clears a delivery ref), though it is still # audit-logged for provenance. _DESTRUCTIVE_ACTIONS: frozenset[str] = frozenset( {"expire", "answer", "supersede", "force-resume"} ) def _utc_now_iso() -> str: """Return the current UTC time as an ISO-8601 string (audit timestamps).""" return datetime.now(timezone.utc).isoformat() def _default_operator() -> str: """Best-effort OS login for audit attribution (never an empty string). A previous empty default left destructive actions non-attributable (the audit record named no one). Defaulting to the OS login keeps the §3.3.1 "audit-logged AND attributable" guarantee even when --operator is omitted; falls back to "unknown" only if the login cannot be resolved. """ try: user = getpass.getuser() except Exception: # noqa: BLE001 - getuser can raise on odd environments return "unknown" return user or "unknown" def _row_to_dict(row: sqlite3.Row) -> dict[str, Any]: """Project a ``pending_questions`` row to a plain dict for display.""" return {col: row[col] for col in _QUESTION_COLUMNS if col in row.keys()} def _append_audit(audit_log: Path, entry: dict[str, Any]) -> None: """Append one JSON audit record with an atomic ``O_APPEND`` single write. Destructive actions (§3.3.1) must leave an attributable trail. A previous read-modify-rewrite design lost records under concurrent operators (two processes each read the same bytes and the last rewrite wins). Instead each record is one line written with ``O_APPEND``: the kernel serializes the append and a write below ``PIPE_BUF`` is atomic on POSIX, so concurrent appends never clobber each other. The file is created mode ``0600`` (operator identity / action content is sensitive) and re-chmod'd in case it pre-existed wider. A failure here raises ``OSError`` BEFORE any ledger mutation, preserving the audit-before-mutate guarantee. """ audit_log = Path(audit_log) audit_log.parent.mkdir(parents=True, exist_ok=True) line = (json.dumps(entry, sort_keys=True) + "\n").encode("utf-8") fd = os.open(str(audit_log), os.O_WRONLY | os.O_CREAT | os.O_APPEND, 0o600) try: os.write(fd, line) finally: os.close(fd) os.chmod(audit_log, 0o600) def _audit_attempt( audit_log: Path, action: str, *, question_id: str, operator: str, detail: dict[str, Any] | None = None, ) -> None: """Record the *intent* to perform a destructive action BEFORE it mutates. §3.3.1 requires every destructive action to be audit-logged. Writing the attempt before the ledger mutation closes the "mutation applied with no audit record" gap: if this append fails (e.g. an unwritable audit path) it raises before any ledger row is touched, so the action aborts cleanly with nothing changed. The matching :func:`_audit_outcome` records what happened. """ _append_audit( audit_log, { "ts": _utc_now_iso(), "action": action, "phase": "attempt", "question_id": question_id, "operator": operator, **(detail or {}), }, ) def _audit_outcome( audit_log: Path, action: str, *, question_id: str, operator: str, applied: bool, detail: dict[str, Any] | None = None, ) -> None: """Record the *result* of a destructive action AFTER it ran. Carries ``applied`` (did the compare-and-set change a row). Pairs with the :func:`_audit_attempt` record written before the mutation, so even if this outcome append fails the attempt already proves the action was made. """ _append_audit( audit_log, { "ts": _utc_now_iso(), "action": action, "phase": "outcome", "question_id": question_id, "operator": operator, "applied": applied, **(detail or {}), }, ) def _require_confirm(action: str, *, confirm: bool) -> None: """Raise unless a destructive ``action`` was explicitly confirmed. Mirrors the §3.3.1 rule: force-expire, answer-on-behalf, and force-resume are audit-logged AND require an explicit confirmation flag. Failing closed here means a typo can never silently mutate a live task's ledger row. """ if action in _DESTRUCTIVE_ACTIONS and not confirm: raise PermissionError( f"refusing destructive action '{action}' without --confirm " f"(force-expire / answer-on-behalf / supersede are gated, §3.3.1)" ) def _fetch_question(conn: sqlite3.Connection, question_id: str) -> sqlite3.Row | None: """Return the ledger row for ``question_id`` or ``None`` if absent.""" return conn.execute( "SELECT * FROM pending_questions WHERE question_id = ?", (question_id,), ).fetchone() # --------------------------------------------------------------------------- # # Subcommand handlers. Each returns a process exit code (0 ok, 1 op failure). # --------------------------------------------------------------------------- # def _cmd_init_db(args: argparse.Namespace, *, out: Any) -> int: """Create/upgrade the agent-team tables (idempotent).""" init_db(args.db) print(f"initialized ledger DB at {args.db}", file=out) return 0 def _cmd_list(args: argparse.Namespace, *, out: Any) -> int: """List pending questions, optionally filtered by status. Default lists ``open`` questions (the operator's "what is waiting" view). ``--status STATE`` narrows to one lifecycle state; ``--all`` lists every state. ``--parked`` is a convenience alias that surfaces the parked-task context an operator chases: questions that are no longer ``open`` (answered but never resumed, expired, or superseded) and so may back a parked task. """ conn = connect(args.db) try: if args.all: rows = conn.execute( "SELECT * FROM pending_questions ORDER BY thread_id, turn" ).fetchall() elif args.parked: placeholders = ",".join("?" for _ in _PARKED_STATES) rows = conn.execute( f"SELECT * FROM pending_questions WHERE status IN ({placeholders}) " "ORDER BY thread_id, turn", tuple(_PARKED_STATES), ).fetchall() else: rows = conn.execute( "SELECT * FROM pending_questions WHERE status = ? " "ORDER BY thread_id, turn", (args.status,), ).fetchall() finally: conn.close() payload = [_row_to_dict(row) for row in rows] print(json.dumps(payload, indent=2, sort_keys=True), file=out) return 0 def _cmd_show(args: argparse.Namespace, *, out: Any) -> int: """Print one question row by ``question_id`` (pure read).""" conn = connect(args.db) try: row = _fetch_question(conn, args.question_id) finally: conn.close() if row is None: print(f"no such question: {args.question_id}", file=sys.stderr) return 1 print(json.dumps(_row_to_dict(row), indent=2, sort_keys=True), file=out) return 0 def _cmd_expire(args: argparse.Namespace, *, out: Any) -> int: """Force-expire an ``open`` question (destructive; audit-logged). Audits the attempt BEFORE mutating so a mutation can never land without a trail (§3.3.1); records the outcome after. """ _require_confirm("expire", confirm=args.confirm) _audit_attempt( args.audit_log, "expire", question_id=args.question_id, operator=args.operator ) conn = connect(args.db) try: changed = expire_question(conn, question_id=args.question_id) finally: conn.close() _audit_outcome( args.audit_log, "expire", question_id=args.question_id, operator=args.operator, applied=changed, ) if not changed: print( f"expire no-op: question {args.question_id} was not 'open' " "(already answered/expired/superseded or absent)", file=sys.stderr, ) return 1 print(f"expired question {args.question_id}", file=out) return 0 def _cmd_redeliver(args: argparse.Namespace, *, out: Any) -> int: """Clear an ``open`` question's ``channel_ref`` so it is re-posted (§3.3.1). The design's "re-deliver" operator action. Re-delivery itself is performed by the transport reconcile loop; clearing ``channel_ref`` makes that loop re-post and record a fresh ref. Idempotent and non-destructive (the question stays ``open``), so it needs no ``--confirm`` — but it is audit-logged for provenance. Returns ``1`` if the question is absent or not ``open``. """ _audit_attempt( args.audit_log, "redeliver", question_id=args.question_id, operator=args.operator, ) conn = connect(args.db) try: row = _fetch_question(conn, args.question_id) if row is None: applied = False prior_ref = None status = None elif row["status"] != "open": applied = False prior_ref = row["channel_ref"] status = row["status"] else: prior_ref = row["channel_ref"] status = "open" conn.execute( "UPDATE pending_questions SET channel_ref=NULL WHERE question_id=?", (args.question_id,), ) applied = True finally: conn.close() _audit_outcome( args.audit_log, "redeliver", question_id=args.question_id, operator=args.operator, applied=applied, detail={"prior_channel_ref": prior_ref}, ) if not applied: reason = "absent" if status is None else f"status={status}, not open" print( f"redeliver no-op: question {args.question_id} ({reason}); " "nothing to re-post", file=sys.stderr, ) return 1 print( f"cleared channel_ref for {args.question_id}; reconcile loop will re-post", file=out, ) return 0 def _cmd_answer(args: argparse.Namespace, *, out: Any) -> int: """Answer a question on a task's behalf (destructive; audit-logged). Uses the foundation first-answer-wins compare-and-set: succeeds only if the question is still ``open``. ``--answer`` is stored verbatim as the answer payload string; ``--via`` records the answering identity for the audit trail. The audit log records the operator regardless of outcome. """ _require_confirm("answer", confirm=args.confirm) via = args.via or f"cli:{args.operator}" _audit_attempt( args.audit_log, "answer", question_id=args.question_id, operator=args.operator, detail={"answered_via": via}, ) conn = connect(args.db) try: changed = answer_question( conn, question_id=args.question_id, answer_json=args.answer, answered_via=via, ) finally: conn.close() _audit_outcome( args.audit_log, "answer", question_id=args.question_id, operator=args.operator, applied=changed, detail={"answered_via": via}, ) if not changed: print( f"answer no-op: question {args.question_id} was not 'open' " "(already answered/expired/superseded or absent)", file=sys.stderr, ) return 1 print(f"answered question {args.question_id} (via {via})", file=out) return 0 def _cmd_supersede(args: argparse.Namespace, *, out: Any) -> int: """Mark a stale ``open``/``answered`` question ``superseded`` (destructive).""" _require_confirm("supersede", confirm=args.confirm) _audit_attempt( args.audit_log, "supersede", question_id=args.question_id, operator=args.operator, ) conn = connect(args.db) try: changed = supersede_question(conn, question_id=args.question_id) finally: conn.close() _audit_outcome( args.audit_log, "supersede", question_id=args.question_id, operator=args.operator, applied=changed, ) if not changed: print( f"supersede no-op: question {args.question_id} was not " "'open'/'answered' (already expired/superseded or absent)", file=sys.stderr, ) return 1 print(f"superseded question {args.question_id}", file=out) return 0 def _cmd_force_resume(args: argparse.Namespace, *, out: Any) -> int: """Force-resume a parked task's question (destructive; audit-logged). The design-named operator verb (§3.3.1 / §6.6 "an operator can force-resume ... a parked task via the CLI"). A task parks when its clarifier question EXPIRES with no answer, so the un-park action is to RE-OPEN that expired question (:func:`agent_team.db.schema.reopen_question`) so the normal delivery → answer → resume flow can proceed. Crucially this does NOT ``supersede`` the row: superseding an ``answered`` row would flip it out of the state the recovery sweep resumes from, making a stuck-but-answered task permanently un-resumable — the opposite of force-resume. So: * ``expired`` (the parked case) → reopened; returns 0. * ``answered`` (answered but not yet resumed) → already eligible for the recovery resume sweep; intent is recorded and we report that, no mutation. * ``open`` / ``superseded`` / absent → nothing to force; reported as a no-op. """ _require_confirm("force-resume", confirm=args.confirm) _audit_attempt( args.audit_log, "force-resume", question_id=args.question_id, operator=args.operator, detail={"resume_requested": True}, ) conn = connect(args.db) try: row = _fetch_question(conn, args.question_id) status = None if row is None else row["status"] reopened = False if status == "expired": reopened = reopen_question(conn, question_id=args.question_id) finally: conn.close() _audit_outcome( args.audit_log, "force-resume", question_id=args.question_id, operator=args.operator, applied=reopened, detail={"resume_requested": True, "prior_status": status}, ) if reopened: print( f"force-resume: reopened expired question {args.question_id}; " "it will be re-delivered for an answer", file=out, ) return 0 if status == "answered": print( f"force-resume: question {args.question_id} is answered and pending " "resume; the recovery sweep will resume it (intent recorded)", file=out, ) return 0 print( f"force-resume no-op: question {args.question_id} " f"({'absent' if status is None else f'status={status}'}) is not parked", file=sys.stderr, ) return 1 # Statuses an operator treats as "parked context": a task whose only pending # question is no longer open may be parked (answered-but-unresumed, expired, or # superseded). ``open`` is excluded — that is the live-waiting view (default # ``list``). Derived from the foundation QUESTION_STATES so it stays in sync. _PARKED_STATES: tuple[str, ...] = tuple(s for s in QUESTION_STATES if s != "open") def build_parser() -> argparse.ArgumentParser: """Construct the argparse parser for ``run-team.py`` (no side effects).""" parser = argparse.ArgumentParser( prog="run-team.py", description=( "R720 agent-team operator CLI — manual path over the durable " "pending_questions ledger (design §3.3.1)." ), ) parser.add_argument( "--db", type=Path, default=_DEFAULT_DB, help=f"path to the agent-team SQLite ledger (default: {_DEFAULT_DB})", ) parser.add_argument( "--audit-log", type=Path, default=_DEFAULT_AUDIT_LOG, dest="audit_log", help=( "append-only JSONL audit log for destructive actions " f"(default: {_DEFAULT_AUDIT_LOG})" ), ) parser.add_argument( "--operator", default=_default_operator(), help="operator identity recorded in the audit log for destructive actions " "(defaults to the OS login so the trail is always attributable)", ) sub = parser.add_subparsers(dest="command", required=True) p_init = sub.add_parser("init-db", help="create/upgrade the ledger tables") p_init.set_defaults(func=_cmd_init_db) p_list = sub.add_parser("list", help="list pending questions (read-only)") list_filter = p_list.add_mutually_exclusive_group() list_filter.add_argument( "--status", choices=QUESTION_STATES, default="open", help="lifecycle status to list (default: open)", ) list_filter.add_argument( "--all", action="store_true", help="list questions in every lifecycle status", ) list_filter.add_argument( "--parked", action="store_true", help="list non-open questions (parked-task context)", ) p_list.set_defaults(func=_cmd_list) p_show = sub.add_parser("show", help="print one question row (read-only)") p_show.add_argument("question_id", help="the question_id to show") p_show.set_defaults(func=_cmd_show) p_redeliver = sub.add_parser( "redeliver", help="clear an open question's channel_ref so it is re-posted", ) p_redeliver.add_argument("question_id", help="the question_id to re-deliver") p_redeliver.set_defaults(func=_cmd_redeliver) p_expire = sub.add_parser( "expire", help="force-expire an open question (destructive)" ) p_expire.add_argument("question_id", help="the question_id to expire") p_expire.add_argument( "--confirm", action="store_true", help="required: confirm this destructive, audit-logged action", ) p_expire.set_defaults(func=_cmd_expire) p_answer = sub.add_parser( "answer", help="answer a question on a task's behalf (destructive)" ) p_answer.add_argument("question_id", help="the question_id to answer") p_answer.add_argument( "--answer", required=True, help="the answer payload (stored verbatim as answer_json)", ) p_answer.add_argument( "--via", default="", help="answering identity for answered_via (default: cli:)", ) p_answer.add_argument( "--confirm", action="store_true", help="required: confirm this destructive, audit-logged action", ) p_answer.set_defaults(func=_cmd_answer) p_supersede = sub.add_parser( "supersede", help="mark a stale question superseded (destructive)" ) p_supersede.add_argument("question_id", help="the question_id to supersede") p_supersede.add_argument( "--confirm", action="store_true", help="required: confirm this destructive, audit-logged action", ) p_supersede.set_defaults(func=_cmd_supersede) p_resume = sub.add_parser( "force-resume", help="force-resume a parked task's question (destructive)", ) p_resume.add_argument("question_id", help="the question_id to force-resume") p_resume.add_argument( "--confirm", action="store_true", help="required: confirm this destructive, audit-logged action", ) p_resume.set_defaults(func=_cmd_force_resume) return parser def main(argv: Sequence[str] | None = None, *, out: Any = None) -> int: """CLI entry point. Returns a process exit code. ``argv`` defaults to ``sys.argv[1:]``; ``out`` defaults to ``sys.stdout`` (injectable for tests). Operational failures return ``1``; a missing ``--confirm`` on a destructive action raises :class:`PermissionError`, surfaced as exit code ``1`` with a stderr message. """ out = out if out is not None else sys.stdout parser = build_parser() args = parser.parse_args(argv) try: return int(args.func(args, out=out)) except PermissionError as exc: # A refused destructive action (no --confirm) or an unwritable audit # path. The attempt-before-mutate ordering means nothing was mutated. print(f"error: {exc}", file=sys.stderr) return 1 except OSError as exc: # Any other audit-log / filesystem failure (e.g. the audit append could # not be written). Surfaced cleanly instead of as an uncaught traceback; # if the attempt record was written, the action is on the trail. print(f"error: audit/IO failure: {exc}", file=sys.stderr) return 1 if __name__ == "__main__": # pragma: no cover raise SystemExit(main())