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/operator_cli.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

664 lines
24 KiB
Python

"""Audit-logged operator CLI over the ``pending_questions`` ledger (§3.3.1, §6.6).
A stuck task *parks* rather than spins, and an operator needs a manual path to
unstick it without poking SQLite by hand. §3.3.1 specifies that path: "a small
CLI over the ledger lets an operator list ``open``/``parked`` questions,
re-deliver, force-expire, or answer on a task's behalf; ... Destructive CLI
actions (force-expire, answer-on-behalf, force-resume) are audit-logged and
require an explicit confirmation flag." §6.6 adds that an operator can
force-resume or re-prioritize a parked task via the CLI.
This module is the leaf that implements that CLI. Two hard rules from the design
are encoded as code, not convention:
1. **Every destructive action is audit-logged** — the *attempt* is recorded
before the mutation runs, and the *outcome* after, so a refused, failed, or
no-op action still leaves a trail. The audit log is an append-only JSONL file
written through the foundation's atomic ``state_store.atomic_write`` primitive
(write-temp → fsync → rename), created mode ``600`` because it records
operator identity and answer content.
2. **Destructive actions require an explicit confirmation flag** — ``--confirm``
on the CLI, ``confirm=True`` in the API. Without it the action raises
:class:`ConfirmationRequired` and performs no mutation (the refused attempt
is still audit-logged). Read-only actions (``list``) are never gated.
The ledger lifecycle transitions reuse the committed compare-and-set helpers
from :mod:`agent_team.db.schema` verbatim (``expire_question``,
``answer_question``, ``supersede_question``); this module does NOT re-implement
the lifecycle SQL. ``force-resume`` records the operator's intent durably and
marks the question ``superseded`` per the turn-guard discipline (§3.3.1) so the
resume worker (a later phase) picks it up idempotently; it does not fabricate a
resume-queue table the foundation schema does not declare.
"""
from __future__ import annotations
import argparse
import getpass
import json
import sqlite3
import sys
from collections.abc import Sequence
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from agent_team.db.schema import (
answer_question,
connect,
expire_question,
supersede_question,
)
from agent_team.state_store import atomic_write
__all__ = [
"DESTRUCTIVE_ACTIONS",
"AuditEntry",
"AuditLog",
"ConfirmationRequired",
"OperatorCli",
"OperatorError",
"QuestionNotFound",
"build_parser",
"main",
]
# The destructive verbs the design singles out (§3.3.1): each requires an
# explicit confirmation flag and is audit-logged. ``redeliver`` and ``list`` are
# deliberately NOT here — re-delivery is idempotent and read-only listing is
# harmless, so neither is gated, though ``redeliver`` is still audit-logged for
# provenance.
DESTRUCTIVE_ACTIONS: frozenset[str] = frozenset(
{"force-expire", "answer-on-behalf", "force-resume"}
)
# The audit log records operator identity and answer payloads, so it is created
# with owner-only permissions (mode 600), matching the design's mode-600 report
# / state convention (§6.7).
_AUDIT_FILE_MODE: int = 0o600
class OperatorError(Exception):
"""Base class for operator-CLI errors."""
class ConfirmationRequired(OperatorError):
"""Raised when a destructive action is invoked without explicit confirmation.
The §3.3.1 rule: force-expire, answer-on-behalf, and force-resume "require an
explicit confirmation flag". The refused attempt is still audit-logged before
this is raised, so an operator who forgets ``--confirm`` leaves a trail.
"""
class QuestionNotFound(OperatorError):
"""Raised when an action targets a ``question_id`` absent from the ledger."""
def _utc_now_iso() -> str:
"""Return the current UTC time as an ISO-8601 string (audit timestamps)."""
return datetime.now(timezone.utc).isoformat()
@dataclass(frozen=True)
class AuditEntry:
"""One immutable audit-log record (§3.3.1 "audit-logged").
``phase`` is ``"attempt"`` for the pre-mutation record and ``"outcome"`` for
the post-mutation record, so a refused or failed action still leaves the
attempt on disk. ``confirmed`` records whether the explicit confirmation flag
was present; ``detail`` carries action-specific context (target ids, answer
payload, rowcount result). Frozen so a constructed record cannot be mutated
before it is written.
"""
timestamp: str
actor: str
action: str
phase: str
confirmed: bool
question_id: str | None = None
thread_id: str | None = None
detail: dict[str, Any] | None = None
def to_dict(self) -> dict[str, Any]:
"""Return a JSON-safe dict for serialization into the JSONL log."""
return {
"timestamp": self.timestamp,
"actor": self.actor,
"action": self.action,
"phase": self.phase,
"confirmed": self.confirmed,
"question_id": self.question_id,
"thread_id": self.thread_id,
"detail": self.detail or {},
}
def to_json(self) -> str:
"""Serialize this entry to a single-line JSON string (one JSONL row)."""
return json.dumps(self.to_dict(), sort_keys=True)
class AuditLog:
"""Append-only, atomically-written JSONL audit trail for operator actions.
Backed by the foundation's :func:`agent_team.state_store.atomic_write`
(write-temp → fsync → rename), so a crash mid-append leaves either the prior
log or the appended-to log, never a torn file. The file is created mode
``600`` (operator identity + answer content are sensitive). The log is
append-only by contract: :meth:`append` reads the current bytes, adds one
line, and atomically rewrites — it never truncates or rewrites history.
"""
def __init__(self, path: Path) -> None:
self._path = Path(path)
@property
def path(self) -> Path:
"""The on-disk path of the JSONL audit log."""
return self._path
def append(self, entry: AuditEntry) -> None:
"""Atomically append ``entry`` as one JSONL line.
Reads the existing log (empty if absent), appends the serialized entry
plus a trailing newline, and atomically rewrites the whole file via
:func:`atomic_write`. Then tightens the file mode to ``600`` so the
recreated file never widens its permissions.
"""
existing = b""
if self._path.exists():
existing = self._path.read_bytes()
line = (entry.to_json() + "\n").encode("utf-8")
atomic_write(self._path, existing + line)
# atomic_write recreates the file via a fresh temp; enforce 600 each time
# so the log never ends up group/world-readable.
self._path.chmod(_AUDIT_FILE_MODE)
def read_all(self) -> list[dict[str, Any]]:
"""Return all audit records as a list of dicts (oldest first)."""
if not self._path.exists():
return []
text = self._path.read_text(encoding="utf-8")
return [json.loads(line) for line in text.splitlines() if line.strip()]
@dataclass(frozen=True)
class _ActionResult:
"""Internal carrier for an action's outcome detail + human-readable message."""
ok: bool
message: str
detail: dict[str, Any]
class OperatorCli:
"""Programmatic operator interface over the ledger + the audit log.
Each public method maps to one CLI verb. Destructive methods take a
``confirm`` flag and raise :class:`ConfirmationRequired` when it is false —
after audit-logging the refused attempt. The class owns the SQLite
connection lifetime so a single process can issue several actions; a context
manager closes it.
"""
def __init__(
self,
db_path: Path,
audit_log_path: Path,
*,
actor: str | None = None,
) -> None:
self._conn = connect(Path(db_path))
self._audit = AuditLog(Path(audit_log_path))
# Default the actor to the OS login so an operator who omits --actor is
# still attributed in the trail.
self._actor = actor or _default_actor()
# -- lifecycle --------------------------------------------------------
def close(self) -> None:
"""Close the underlying SQLite connection."""
self._conn.close()
def __enter__(self) -> OperatorCli:
return self
def __exit__(self, *exc: object) -> None:
self.close()
@property
def audit_log(self) -> AuditLog:
"""The :class:`AuditLog` this CLI writes to (exposed for inspection)."""
return self._audit
# -- helpers ----------------------------------------------------------
def _audit_attempt(
self,
action: str,
*,
confirmed: bool,
question_id: str | None = None,
thread_id: str | None = None,
detail: dict[str, Any] | None = None,
) -> None:
self._audit.append(
AuditEntry(
timestamp=_utc_now_iso(),
actor=self._actor,
action=action,
phase="attempt",
confirmed=confirmed,
question_id=question_id,
thread_id=thread_id,
detail=detail,
)
)
def _audit_outcome(
self,
action: str,
*,
confirmed: bool,
question_id: str | None = None,
thread_id: str | None = None,
detail: dict[str, Any] | None = None,
) -> None:
self._audit.append(
AuditEntry(
timestamp=_utc_now_iso(),
actor=self._actor,
action=action,
phase="outcome",
confirmed=confirmed,
question_id=question_id,
thread_id=thread_id,
detail=detail,
)
)
def _require_confirm(
self,
action: str,
confirm: bool,
*,
question_id: str | None = None,
thread_id: str | None = None,
detail: dict[str, Any] | None = None,
) -> None:
"""Audit-log the attempt; raise if a destructive action is unconfirmed.
Always records the attempt first (so a refused action is on disk), then
enforces the explicit-confirmation rule for destructive verbs.
"""
self._audit_attempt(
action,
confirmed=confirm,
question_id=question_id,
thread_id=thread_id,
detail=detail,
)
if action in DESTRUCTIVE_ACTIONS and not confirm:
self._audit_outcome(
action,
confirmed=confirm,
question_id=question_id,
thread_id=thread_id,
detail={"refused": "confirmation flag required"},
)
raise ConfirmationRequired(
f"action {action!r} is destructive and requires explicit "
f"confirmation (pass --confirm / confirm=True)"
)
def _fetch_row(self, question_id: str) -> sqlite3.Row:
row = self._conn.execute(
"SELECT * FROM pending_questions WHERE question_id=?",
(question_id,),
).fetchone()
if row is None:
raise QuestionNotFound(f"no ledger row for question_id={question_id!r}")
return row
# -- read-only --------------------------------------------------------
def list_questions(
self,
*,
statuses: Sequence[str] | None = None,
thread_id: str | None = None,
) -> list[dict[str, Any]]:
"""List ledger rows, optionally filtered by status and/or thread.
Read-only and ungated; defaults to the operator-relevant ``open`` set
(the questions awaiting action). Pass ``statuses`` to widen or narrow.
Returns plain dicts (oldest-posted first) so the result is JSON-safe.
"""
clauses: list[str] = []
params: list[Any] = []
if statuses:
placeholders = ",".join("?" for _ in statuses)
clauses.append(f"status IN ({placeholders})")
params.extend(statuses)
if thread_id is not None:
clauses.append("thread_id = ?")
params.append(thread_id)
where = f" WHERE {' AND '.join(clauses)}" if clauses else ""
sql = (
"SELECT * FROM pending_questions"
+ where
+ " ORDER BY posted_at IS NULL, posted_at, question_id"
)
rows = self._conn.execute(sql, tuple(params)).fetchall()
return [dict(row) for row in rows]
# -- destructive ------------------------------------------------------
def redeliver(self, question_id: str) -> _ActionResult:
"""Mark an ``open`` question for re-delivery (idempotent, audit-logged).
Re-delivery itself is performed by the transport reconcile loop (§3.3.1);
this clears the stored ``channel_ref`` so the reconcile loop re-posts and
records a fresh ref. Not destructive (the question stays ``open``), so it
does not require confirmation, but it is audit-logged for provenance.
"""
self._require_confirm("redeliver", True, question_id=question_id)
row = self._fetch_row(question_id)
if row["status"] != "open":
result = _ActionResult(
ok=False,
message=f"question {question_id} is {row['status']}, not open; "
"nothing to re-deliver",
detail={"status": row["status"]},
)
else:
self._conn.execute(
"UPDATE pending_questions SET channel_ref=NULL WHERE question_id=?",
(question_id,),
)
result = _ActionResult(
ok=True,
message=f"cleared channel_ref for {question_id}; reconcile loop "
"will re-post",
detail={"prior_channel_ref": row["channel_ref"]},
)
self._audit_outcome(
"redeliver",
confirmed=True,
question_id=question_id,
thread_id=row["thread_id"],
detail=result.detail,
)
return result
def force_expire(self, question_id: str, *, confirm: bool = False) -> _ActionResult:
"""Force an ``open`` question to ``expired`` (destructive, confirmed).
Delegates the lifecycle transition to the committed compare-and-set
:func:`agent_team.db.schema.expire_question`, so an answer racing in
still loses deterministically (§3.3.1). Returns ``ok`` reflecting the
compare-and-set rowcount.
"""
self._require_confirm("force-expire", confirm, question_id=question_id)
row = self._fetch_row(question_id)
changed = expire_question(self._conn, question_id=question_id)
result = _ActionResult(
ok=changed,
message=(
f"expired {question_id}"
if changed
else f"{question_id} was not open (status={row['status']}); no change"
),
detail={"changed": changed, "prior_status": row["status"]},
)
self._audit_outcome(
"force-expire",
confirmed=confirm,
question_id=question_id,
thread_id=row["thread_id"],
detail=result.detail,
)
return result
def answer_on_behalf(
self,
question_id: str,
answer: Any,
*,
confirm: bool = False,
) -> _ActionResult:
"""Answer an ``open`` question on the task's behalf (destructive, confirmed).
Routes through the committed first-answer-wins
:func:`agent_team.db.schema.answer_question`; ``answered_via`` is stamped
``operator:<actor>`` so the audit trail and the ledger agree on who
answered. The answer is JSON-encoded into ``answer_json``. A late answer
(question already closed) loses the compare-and-set and returns
``ok=False``.
"""
answer_json = json.dumps(answer, sort_keys=True)
via = f"operator:{self._actor}"
self._require_confirm(
"answer-on-behalf",
confirm,
question_id=question_id,
detail={"answer": answer, "via": via},
)
row = self._fetch_row(question_id)
changed = answer_question(
self._conn,
question_id=question_id,
answer_json=answer_json,
answered_via=via,
)
result = _ActionResult(
ok=changed,
message=(
f"answered {question_id} on behalf of the task"
if changed
else f"{question_id} was not open (status={row['status']}); "
"answer ignored"
),
detail={"changed": changed, "prior_status": row["status"], "via": via},
)
self._audit_outcome(
"answer-on-behalf",
confirmed=confirm,
question_id=question_id,
thread_id=row["thread_id"],
detail=result.detail,
)
return result
def force_resume(self, question_id: str, *, confirm: bool = False) -> _ActionResult:
"""Force-resume a parked task's question (destructive, confirmed).
§6.6: an operator can force-resume a parked task via the CLI. The resume
worker proper is a later phase; here we record the operator's intent
durably in the audit log and mark the stale question ``superseded`` via
the committed turn-guard helper
:func:`agent_team.db.schema.supersede_question` so a redelivered/stale
resume job is skipped idempotently (§3.3.1). The resume worker consumes
the audit intent on its next sweep; this method never invokes the
LangGraph runtime directly (out of scope for the foundation).
"""
self._require_confirm("force-resume", confirm, question_id=question_id)
row = self._fetch_row(question_id)
superseded = supersede_question(self._conn, question_id=question_id)
result = _ActionResult(
ok=True,
message=(
f"recorded force-resume intent for {question_id} "
f"(thread {row['thread_id']}); "
+ (
"superseded stale question"
if superseded
else "no open/answered question to supersede"
)
),
detail={
"thread_id": row["thread_id"],
"turn": row["turn"],
"superseded": superseded,
"prior_status": row["status"],
"resume_requested": True,
},
)
self._audit_outcome(
"force-resume",
confirmed=confirm,
question_id=question_id,
thread_id=row["thread_id"],
detail=result.detail,
)
return result
def _default_actor() -> str:
"""Best-effort OS login name for audit attribution.
Falls back to ``"unknown"`` rather than raising in the rare environment with
no resolvable login, so an action is never blocked purely on attribution
lookup (the attempt is still recorded, just as ``unknown``).
"""
try:
return getpass.getuser()
except Exception: # noqa: BLE001 - getuser can raise on odd environments
return "unknown"
def build_parser() -> argparse.ArgumentParser:
"""Build the ``argparse`` parser for the operator CLI.
Exposed separately so tests can parse argv without dispatching. The entry
CLI (``run-team.py operator ...``) and direct ``python -m`` both route here.
"""
parser = argparse.ArgumentParser(
prog="operator",
description="Audit-logged operator CLI over the pending_questions ledger "
"(R720 agent-team §3.3.1, §6.6).",
)
parser.add_argument(
"--db",
required=True,
type=Path,
help="path to the agent-team SQLite database",
)
parser.add_argument(
"--audit-log",
required=True,
type=Path,
help="path to the append-only JSONL audit log (created mode 600)",
)
parser.add_argument(
"--actor",
default=None,
help="operator identity for the audit trail (defaults to OS login)",
)
sub = parser.add_subparsers(dest="command", required=True)
p_list = sub.add_parser("list", help="list ledger questions (read-only)")
p_list.add_argument(
"--status",
action="append",
dest="statuses",
choices=("open", "answered", "expired", "superseded"),
help="filter by status (repeatable); default is open",
)
p_list.add_argument("--thread", default=None, help="filter by thread_id")
p_redeliver = sub.add_parser("redeliver", help="clear channel_ref to re-post")
p_redeliver.add_argument("question_id")
p_expire = sub.add_parser(
"force-expire", help="force an open question to expired (destructive)"
)
p_expire.add_argument("question_id")
p_expire.add_argument(
"--confirm", action="store_true", help="explicit confirmation (required)"
)
p_answer = sub.add_parser(
"answer-on-behalf",
help="answer an open question on the task's behalf (destructive)",
)
p_answer.add_argument("question_id")
p_answer.add_argument(
"answer", help="answer payload; parsed as JSON if valid, else a raw string"
)
p_answer.add_argument(
"--confirm", action="store_true", help="explicit confirmation (required)"
)
p_resume = sub.add_parser(
"force-resume", help="force-resume a parked task's question (destructive)"
)
p_resume.add_argument("question_id")
p_resume.add_argument(
"--confirm", action="store_true", help="explicit confirmation (required)"
)
return parser
def _coerce_answer(raw: str) -> Any:
"""Parse an answer arg as JSON, falling back to the raw string.
Lets an operator pass ``'{"approve": true}'`` or just ``approve`` from the
shell; structured JSON round-trips, a bare token is kept verbatim.
"""
try:
return json.loads(raw)
except (ValueError, TypeError):
return raw
def main(argv: Sequence[str] | None = None) -> int:
"""CLI entrypoint. Returns a process exit code (0 ok, non-zero on error).
Dispatches the parsed subcommand against an :class:`OperatorCli`. Confirmation
refusals and missing questions are reported on stderr with a non-zero exit;
the attempt is already audit-logged before the error surfaces.
"""
parser = build_parser()
args = parser.parse_args(argv)
try:
with OperatorCli(args.db, args.audit_log, actor=args.actor) as cli:
if args.command == "list":
rows = cli.list_questions(
statuses=args.statuses or ["open"],
thread_id=args.thread,
)
print(json.dumps(rows, indent=2, sort_keys=True))
return 0
if args.command == "redeliver":
result = cli.redeliver(args.question_id)
elif args.command == "force-expire":
result = cli.force_expire(args.question_id, confirm=args.confirm)
elif args.command == "answer-on-behalf":
result = cli.answer_on_behalf(
args.question_id,
_coerce_answer(args.answer),
confirm=args.confirm,
)
elif args.command == "force-resume":
result = cli.force_resume(args.question_id, confirm=args.confirm)
else: # pragma: no cover - argparse enforces the choices
parser.error(f"unknown command {args.command!r}")
print(result.message)
return 0 if result.ok else 1
except ConfirmationRequired as exc:
print(f"refused: {exc}", file=sys.stderr)
return 2
except QuestionNotFound as exc:
print(f"error: {exc}", file=sys.stderr)
return 3
if __name__ == "__main__": # pragma: no cover - module-level entry shim
raise SystemExit(main())