diff --git a/agent-team/.gitignore b/agent-team/.gitignore new file mode 100644 index 0000000..1a8546b --- /dev/null +++ b/agent-team/.gitignore @@ -0,0 +1,18 @@ +# Local durable state — never commit (SQLite stores + integrity sidecars) +*.sqlite +*.sqlite-wal +*.sqlite-shm +*.db +*.meta.json +state/ +.state/ + +# Secrets — never commit +*.env +secrev.env + +# Python artifacts +__pycache__/ +*.pyc +.pytest_cache/ +.ruff_cache/ diff --git a/agent-team/README.md b/agent-team/README.md new file mode 100644 index 0000000..12bb29d --- /dev/null +++ b/agent-team/README.md @@ -0,0 +1,65 @@ +# agent-team — R720 Plane-2 FOUNDATION + +Pre-deployment scaffolding for the R720 agent-team SDLC pipeline (design: +`../docs/r720-agent-team-design.md`). This commit ships the **Plane-2 +FOUNDATION** layer only — the durable, transport-agnostic **contracts** the leaf +builders import verbatim. Nothing here is provisioned, scheduled, or wired to +live infrastructure. + +> Status: FOUNDATION modules only. No coordinator, no transports' concrete +> adapters, no CI workflow, no provisioning. Those are later phases (§7). + +## Layout + +``` +agent-team/ + agent_team/ # importable package (snake_case) + state_store.py # §6.7 atomic write + integrity-checked read + billing.py # §3.1 claude_invoke billing-mode seam + task_model.py # §3.3 TaskRecord / Phase / PipelineState + db/ + schema.py # §3.3.1/§6.7 SQLite DDL + connect/init/migrate + schema.sql # raw DDL, mirrors schema.py verbatim + transport/ + base.py # §3.3.1 Transport ABC + QuestionSet/NormalizedAnswer + tests/ # pytest unit tests, one module per source module +``` + +The top directory is kebab-case (`agent-team/`); the importable package is +snake_case (`agent_team/`), per the engineering handbook. + +## Modules (contracts) + +| Module | Design ref | What it provides | +|---|---|---| +| `state_store` | §6.7 | `atomic_write(path, data)` (write-temp → fsync → rename), `read_checked(path, *, schema_version)` (schema-version + content-hash integrity check, raises `IntegrityError`), `compute_content_hash(data)`. Pure stdlib; no other `agent_team` deps. | +| `db.schema` | §3.3.1, §6.7 | `SCHEMA_VERSION`, `PENDING_QUESTIONS_DDL`, `BUDGET_LEDGER_DDL`, `connect()` (WAL + foreign_keys + busy_timeout), `init_db()`, `migrate()`, and the `BEGIN IMMEDIATE` compare-and-set helpers (`answer_question`/`expire_question`/`supersede_question`). SQL DDL lives **only** here. | +| `billing` | §3.1 | `BillingMode{SUBSCRIPTION,API,BEDROCK}`, `claude_invoke(prompt, *, mode=None, **kw) -> ClaudeResult`, `resolve_mode(config)`. Single seam; subscription mode pops any stray `ANTHROPIC_API_KEY` so OAuth can't be overridden. | +| `transport.base` | §3.3.1 | `Transport` ABC (`post_question` → `channel_ref`; `parse_answer` → `(question_id, answer, via)`), `QuestionSet`, `NormalizedAnswer`. Transport-independent; Slack/GitHub/Claude-Code adapters subclass in the leaves. | +| `task_model` | §3.3 | `TaskRecord`, `TaskStatus`, `Phase{INTAKE…DONE}`, `new_thread_id()`, `PipelineState` TypedDict (LangGraph state schema), JSON serialization helpers. Pure model, no I/O. | + +### Durable human-in-the-loop (§3.3.1) + +The `pending_questions` ledger is the single durable source of truth for the +question lifecycle. Every race (duplicate answers, transport redelivery, +answer-vs-timeout) resolves via one atomic compare-and-set against the `status` +column, run inside a `BEGIN IMMEDIATE` transaction so concurrent responders are +serialized — first-answer-wins (`rowcount == 1`), late/duplicate ignored +(`rowcount == 0`). The LangGraph `SqliteSaver` checkpointer creates its own +tables against the **same** DB file. + +## Running the tests + +``` +cd agent-team +python3 -m pytest tests/ -q +``` + +`tests/conftest.py` puts the package on `sys.path`, so no install is required. + +## Not in this commit (later phases) + +Coordinator/brain, concrete Slack/GitHub/Claude-Code transport adapters, the +CI apply/verify workflow (§3.3.2), step-ca / Roles Anywhere, scheduling, and +the operator CLI (`run-team.py`). See `../docs/r720-agent-team-design.md` §7 +for the phased rollout. Secrets are never committed. diff --git a/agent-team/agent_team/__init__.py b/agent-team/agent_team/__init__.py new file mode 100644 index 0000000..b89a20d --- /dev/null +++ b/agent-team/agent_team/__init__.py @@ -0,0 +1,11 @@ +"""R720 agent-team — Plane-2 FOUNDATION modules (design v2). + +This package holds the durable, transport-agnostic contracts the leaf builders +import verbatim: the atomic state-store (:mod:`agent_team.state_store`), the +SQLite schema (:mod:`agent_team.db`), the Claude billing seam +(:mod:`agent_team.billing`), the transport interface +(:mod:`agent_team.transport`), and the task/thread model +(:mod:`agent_team.task_model`). +""" + +__all__: list[str] = [] diff --git a/agent-team/agent_team/billing.py b/agent-team/agent_team/billing.py new file mode 100644 index 0000000..ca84b2e --- /dev/null +++ b/agent-team/agent_team/billing.py @@ -0,0 +1,148 @@ +"""Billing-mode abstraction — the single ``claude_invoke`` seam (design §3.1). + +Every Claude-calling node imports :func:`claude_invoke` from here. The seam +selects the Claude auth/billing path from config: + +* ``SUBSCRIPTION`` — OAuth token (the R720 default; headless Agent SDK), +* ``API`` — metered ``ANTHROPIC_API_KEY``, +* ``BEDROCK`` — cross-account Bedrock (the rare cross-family tiebreak). + +Switching modes is a config flip, not a code change. In ``SUBSCRIPTION`` mode +the seam pops/unsets any stray ``ANTHROPIC_API_KEY`` from the environment +before invoking, so an inherited key cannot silently override OAuth (§3.1). + +This module is the contract leaf builders import verbatim; the actual SDK call +is delegated to an injectable ``_invoker`` so the seam stays testable and the +transport/SDK wiring lives in the leaves. +""" + +from __future__ import annotations + +import os +from dataclasses import dataclass, field +from enum import Enum +from typing import Any, Callable, Mapping + +__all__ = [ + "BillingMode", + "ClaudeResult", + "claude_invoke", + "resolve_mode", + "set_invoker", +] + +# Environment variable that carries the metered API key. Popped in +# subscription mode so OAuth cannot be silently overridden. +_API_KEY_ENV = "ANTHROPIC_API_KEY" + +# Config key (env or mapping) naming the desired billing mode. +_MODE_ENV = "AGENT_TEAM_BILLING_MODE" + + +class BillingMode(Enum): + """Claude auth/billing path selector (§3.1).""" + + SUBSCRIPTION = "subscription" + API = "api" + BEDROCK = "bedrock" + + +@dataclass +class ClaudeResult: + """Result of a :func:`claude_invoke` call. + + ``text`` is the model's response text. ``mode`` records which billing path + served the call. ``usage`` carries token/cost accounting for the budget + ledger (§6.6); ``raw`` is the untouched provider response for callers that + need more. + """ + + text: str + mode: BillingMode + usage: dict[str, Any] = field(default_factory=dict) + raw: Any = None + + +# Pluggable invoker: signature (prompt, mode, **kw) -> ClaudeResult. The +# default raises so an un-wired environment fails loudly rather than silently +# returning nothing; leaves call set_invoker() to bind the real SDK path. +Invoker = Callable[..., ClaudeResult] + + +def _unconfigured_invoker(prompt: str, *, mode: BillingMode, **kw: Any) -> ClaudeResult: + raise RuntimeError( + "claude_invoke has no invoker bound; call billing.set_invoker(fn) to " + "wire the Claude SDK path (subscription OAuth / API / Bedrock)." + ) + + +_invoker: Invoker = _unconfigured_invoker + + +def set_invoker(invoker: Invoker) -> None: + """Bind the function that performs the actual Claude SDK call. + + Leaves call this once at startup with an implementation that honours the + resolved :class:`BillingMode`. Keeping the SDK call injectable keeps this + seam dependency-free and unit-testable. + """ + global _invoker + _invoker = invoker + + +def resolve_mode(config: Mapping[str, Any] | None) -> BillingMode: + """Resolve the billing mode from ``config`` (falling back to env). + + Precedence: an explicit ``billing_mode`` in ``config`` (a + :class:`BillingMode` or its string value), then the + ``AGENT_TEAM_BILLING_MODE`` env var, then the ``SUBSCRIPTION`` default. + """ + raw: Any = None + if config is not None: + raw = config.get("billing_mode") + if raw is None: + raw = os.environ.get(_MODE_ENV) + if raw is None: + return BillingMode.SUBSCRIPTION + if isinstance(raw, BillingMode): + return raw + try: + return BillingMode(str(raw).strip().lower()) + except ValueError as exc: + valid = ", ".join(m.value for m in BillingMode) + raise ValueError( + f"unknown billing mode {raw!r}; expected one of: {valid}" + ) from exc + + +def claude_invoke( + prompt: str, + *, + mode: BillingMode | None = None, + config: Mapping[str, Any] | None = None, + **kw: Any, +) -> ClaudeResult: + """Invoke Claude through the configured billing path (§3.1). + + ``mode`` overrides config when given; otherwise it is resolved via + :func:`resolve_mode`. In ``SUBSCRIPTION`` mode any stray + ``ANTHROPIC_API_KEY`` is popped from ``os.environ`` for the duration of the + call so OAuth cannot be silently overridden, then restored afterward. + + The actual SDK call is delegated to the bound invoker (see + :func:`set_invoker`); this function owns only mode selection and the + subscription-mode env hygiene that the design mandates. + """ + effective = mode if mode is not None else resolve_mode(config) + + if effective is BillingMode.SUBSCRIPTION: + # Pop the stray key for the duration of the call; restore on exit so we + # don't mutate the caller's environment permanently. + stashed = os.environ.pop(_API_KEY_ENV, None) + try: + return _invoker(prompt, mode=effective, **kw) + finally: + if stashed is not None: + os.environ[_API_KEY_ENV] = stashed + + return _invoker(prompt, mode=effective, **kw) diff --git a/agent-team/agent_team/db/__init__.py b/agent-team/agent_team/db/__init__.py new file mode 100644 index 0000000..c62f898 --- /dev/null +++ b/agent-team/agent_team/db/__init__.py @@ -0,0 +1,24 @@ +"""SQLite schema and connection helpers for the R720 agent-team pipeline. + +SQL DDL lives ONLY in this subpackage (``schema.py`` constants mirrored by the +companion ``schema.sql``). The LangGraph ``SqliteSaver`` checkpointer creates +its own tables against the same database file. +""" + +from agent_team.db.schema import ( + BUDGET_LEDGER_DDL, + PENDING_QUESTIONS_DDL, + SCHEMA_VERSION, + connect, + init_db, + migrate, +) + +__all__ = [ + "BUDGET_LEDGER_DDL", + "PENDING_QUESTIONS_DDL", + "SCHEMA_VERSION", + "connect", + "init_db", + "migrate", +] diff --git a/agent-team/agent_team/db/schema.py b/agent-team/agent_team/db/schema.py new file mode 100644 index 0000000..691ed4c --- /dev/null +++ b/agent-team/agent_team/db/schema.py @@ -0,0 +1,295 @@ +"""SQLite schema, DDL constants, and connection helpers (design §3.3.1, §6.7). + +This module is the single source of truth for the R720 agent-team durable +SQL. It declares: + +* the ``pending_questions`` human-interaction lifecycle ledger (§3.3.1), +* the ``budget_ledger`` shared Claude budget ledger (§6.1, §6.6), +* a ``schema_meta`` version row driving :func:`migrate`. + +The companion ``schema.sql`` mirrors these statements verbatim for tooling. +SQL DDL lives ONLY here. The LangGraph ``SqliteSaver`` checkpointer creates +its OWN tables against this same connection / database file; the design +reserves this DB for it but does not declare its tables. + +The atomic compare-and-set helpers used by the responder and the deadline +timer (§3.3.1) take the write lock up front via ``BEGIN IMMEDIATE`` so +concurrent responders are serialized — SQLite's default deferred isolation +does not serialize a check-and-set. +""" + +from __future__ import annotations + +import sqlite3 +from datetime import datetime, timezone +from pathlib import Path +from typing import Any + +__all__ = [ + "BUDGET_LEDGER_DDL", + "PENDING_QUESTIONS_DDL", + "PENDING_QUESTIONS_INDEXES_DDL", + "BUDGET_LEDGER_INDEXES_DDL", + "SCHEMA_META_DDL", + "SCHEMA_VERSION", + "QUESTION_STATES", + "answer_question", + "connect", + "expire_question", + "init_db", + "migrate", + "supersede_question", +] + +# Bump when the DDL below changes; migrate() steps a connection forward. +SCHEMA_VERSION: int = 1 + +# Default SQLite busy timeout (ms) so concurrent writers wait for the write +# lock rather than failing immediately. +_BUSY_TIMEOUT_MS: int = 5000 + +# Allowed lifecycle states for a pending question (§3.3.1). Mirrors the DDL +# CHECK constraint; exported so leaves can validate without re-listing them. +QUESTION_STATES: tuple[str, ...] = ("open", "answered", "expired", "superseded") + + +PENDING_QUESTIONS_DDL: str = """ +CREATE TABLE IF NOT EXISTS pending_questions ( + question_id TEXT PRIMARY KEY, + thread_id TEXT NOT NULL, + turn INTEGER NOT NULL, + status TEXT NOT NULL + CHECK (status IN ('open', 'answered', 'expired', 'superseded')), + transport TEXT NOT NULL, + channel_ref TEXT, + posted_at TEXT, + deadline_at TEXT, + answer_json TEXT, + answered_at TEXT, + answered_via TEXT +) +""".strip() + +PENDING_QUESTIONS_INDEXES_DDL: str = """ +CREATE INDEX IF NOT EXISTS idx_pending_questions_thread + ON pending_questions (thread_id, turn); +CREATE INDEX IF NOT EXISTS idx_pending_questions_status + ON pending_questions (status); +""".strip() + +BUDGET_LEDGER_DDL: str = """ +CREATE TABLE IF NOT EXISTS budget_ledger ( + entry_id INTEGER PRIMARY KEY AUTOINCREMENT, + thread_id TEXT, + stage TEXT, + model TEXT NOT NULL, + billing_mode TEXT NOT NULL, + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + usd_cost REAL NOT NULL DEFAULT 0.0, + recorded_at TEXT NOT NULL, + day_bucket TEXT NOT NULL +) +""".strip() + +BUDGET_LEDGER_INDEXES_DDL: str = """ +CREATE INDEX IF NOT EXISTS idx_budget_ledger_day + ON budget_ledger (day_bucket); +CREATE INDEX IF NOT EXISTS idx_budget_ledger_thread + ON budget_ledger (thread_id); +""".strip() + +SCHEMA_META_DDL: str = """ +CREATE TABLE IF NOT EXISTS schema_meta ( + id INTEGER PRIMARY KEY CHECK (id = 1), + schema_version INTEGER NOT NULL +) +""".strip() + + +def connect(db_path: Path) -> sqlite3.Connection: + """Open ``db_path`` with WAL, foreign keys, and a busy timeout. + + WAL (``journal_mode=WAL``) lets the resume worker read while a responder + writes; ``foreign_keys=ON`` enforces referential integrity; the busy + timeout makes concurrent writers wait for the write lock instead of + failing. ``isolation_level=None`` puts the connection in autocommit mode so + the compare-and-set helpers can drive transactions explicitly with + ``BEGIN IMMEDIATE`` (§3.3.1). + """ + db_path = Path(db_path) + db_path.parent.mkdir(parents=True, exist_ok=True) + conn = sqlite3.connect( + str(db_path), + isolation_level=None, + check_same_thread=False, + ) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA journal_mode=WAL") + conn.execute("PRAGMA foreign_keys=ON") + conn.execute(f"PRAGMA busy_timeout={_BUSY_TIMEOUT_MS}") + return conn + + +def init_db(db_path: Path) -> None: + """Create the agent-team tables in ``db_path`` if absent. + + Creates ``pending_questions`` (+ indexes), the budget ledger (+ indexes), + and the ``schema_meta`` version row, and reserves the same DB file for the + LangGraph ``SqliteSaver`` checkpointer (which creates its own tables on + first use against this connection). Idempotent: safe to call on every + startup. + """ + conn = connect(db_path) + try: + conn.execute(SCHEMA_META_DDL) + conn.execute(PENDING_QUESTIONS_DDL) + for stmt in _split_statements(PENDING_QUESTIONS_INDEXES_DDL): + conn.execute(stmt) + conn.execute(BUDGET_LEDGER_DDL) + for stmt in _split_statements(BUDGET_LEDGER_INDEXES_DDL): + conn.execute(stmt) + # Record the schema version (single-row table). + conn.execute( + "INSERT INTO schema_meta (id, schema_version) VALUES (1, ?) " + "ON CONFLICT(id) DO NOTHING", + (SCHEMA_VERSION,), + ) + finally: + conn.close() + + +def migrate(conn: sqlite3.Connection) -> None: + """Step ``conn``'s schema forward to :data:`SCHEMA_VERSION`. + + Reads the recorded version from ``schema_meta`` (treating an empty/absent + row as version 0), applies any forward steps, and records the new version. + At ``SCHEMA_VERSION == 1`` there are no prior versions to migrate from, so + this ensures the base tables exist and stamps the version. Future versions + add ordered ``if current < N`` blocks here. + """ + conn.execute(SCHEMA_META_DDL) + row = conn.execute("SELECT schema_version FROM schema_meta WHERE id = 1").fetchone() + current = int(row["schema_version"]) if row is not None else 0 + + if current < 1: + # Base schema (v1): ensure all tables/indexes exist. + conn.execute(PENDING_QUESTIONS_DDL) + for stmt in _split_statements(PENDING_QUESTIONS_INDEXES_DDL): + conn.execute(stmt) + conn.execute(BUDGET_LEDGER_DDL) + for stmt in _split_statements(BUDGET_LEDGER_INDEXES_DDL): + conn.execute(stmt) + current = 1 + + # Future steps go here: `if current < 2: ...; current = 2`. + + conn.execute( + "INSERT INTO schema_meta (id, schema_version) VALUES (1, ?) " + "ON CONFLICT(id) DO UPDATE SET schema_version = excluded.schema_version", + (current,), + ) + + +def answer_question( + conn: sqlite3.Connection, + *, + question_id: str, + answer_json: str, + answered_via: str, + answered_at: str | None = None, +) -> bool: + """First-answer-wins compare-and-set: flip an ``open`` question to answered. + + Runs the §3.3.1 atomic statement inside a ``BEGIN IMMEDIATE`` transaction + so concurrent responders are serialized (the check-and-set takes the write + lock up front). Returns ``True`` when rowcount == 1 (this caller recorded + the first valid answer; enqueue a resume job), ``False`` when rowcount == 0 + (the question was not ``open`` — already answered/expired/superseded — so + the answer is a duplicate or late and must be ignored). + """ + stamp = answered_at or _utc_now_iso() + return _compare_and_set( + conn, + sql=( + "UPDATE pending_questions " + "SET status='answered', answer_json=?, answered_via=?, answered_at=? " + "WHERE question_id=? AND status='open'" + ), + params=(answer_json, answered_via, stamp, question_id), + ) + + +def expire_question( + conn: sqlite3.Connection, + *, + question_id: str, +) -> bool: + """Deadline race: flip an overdue ``open`` question to ``expired``. + + Same compare-and-set discipline as :func:`answer_question` (§3.3.1): an + answer that arrives for an already-expired question loses the race and is + ignored. Returns ``True`` if this call expired the question. + """ + return _compare_and_set( + conn, + sql=( + "UPDATE pending_questions SET status='expired' " + "WHERE question_id=? AND status='open'" + ), + params=(question_id,), + ) + + +def supersede_question( + conn: sqlite3.Connection, + *, + question_id: str, +) -> bool: + """Mark a stale ``open``/``answered`` question ``superseded``. + + Used by the turn-guarded resume worker: if the graph already advanced past + this turn, the question is superseded and the resume is skipped (§3.3.1). + Returns ``True`` if this call superseded the question. + """ + return _compare_and_set( + conn, + sql=( + "UPDATE pending_questions SET status='superseded' " + "WHERE question_id=? AND status IN ('open', 'answered')" + ), + params=(question_id,), + ) + + +def _compare_and_set( + conn: sqlite3.Connection, + *, + sql: str, + params: tuple[Any, ...], +) -> bool: + """Run a single compare-and-set UPDATE under ``BEGIN IMMEDIATE``. + + Returns ``True`` iff exactly one row changed. Takes the write lock up front + so concurrent responders cannot both observe ``status='open'`` (SQLite's + default deferred isolation would not serialize them). + """ + conn.execute("BEGIN IMMEDIATE") + try: + cur = conn.execute(sql, params) + changed = cur.rowcount == 1 + conn.execute("COMMIT") + except BaseException: + conn.execute("ROLLBACK") + raise + return changed + + +def _utc_now_iso() -> str: + """Return the current UTC time as an ISO-8601 string.""" + return datetime.now(timezone.utc).isoformat() + + +def _split_statements(ddl: str) -> list[str]: + """Split a multi-statement DDL blob into individual statements.""" + return [stmt.strip() for stmt in ddl.split(";") if stmt.strip()] diff --git a/agent-team/agent_team/db/schema.sql b/agent-team/agent_team/db/schema.sql new file mode 100644 index 0000000..5dc4479 --- /dev/null +++ b/agent-team/agent_team/db/schema.sql @@ -0,0 +1,65 @@ +-- R720 agent-team durable SQLite schema (design §3.3.1, §6.7). +-- +-- This file holds the raw DDL statements ONLY. The authoritative copies live +-- as string constants in agent_team/db/schema.py; this companion file mirrors +-- them verbatim for tooling / direct inspection. SQL DDL lives only in these +-- two places. +-- +-- The LangGraph SqliteSaver checkpointer creates its OWN tables against this +-- same database file/connection; those are intentionally NOT declared here. + +-- pending_questions: the durable human-interaction lifecycle ledger. +-- The LangGraph checkpoint holds graph state; this table holds the question +-- lifecycle (delivery, duplicate/late answers, expiry) and is what delivery, +-- the responder, and restart recovery read. Every race resolves via an atomic +-- compare-and-set against the `status` column under BEGIN IMMEDIATE. +CREATE TABLE IF NOT EXISTS pending_questions ( + question_id TEXT PRIMARY KEY, + thread_id TEXT NOT NULL, + turn INTEGER NOT NULL, + status TEXT NOT NULL + CHECK (status IN ('open', 'answered', 'expired', 'superseded')), + transport TEXT NOT NULL, + channel_ref TEXT, + posted_at TEXT, + deadline_at TEXT, + answer_json TEXT, + answered_at TEXT, + answered_via TEXT +); + +CREATE INDEX IF NOT EXISTS idx_pending_questions_thread + ON pending_questions (thread_id, turn); + +CREATE INDEX IF NOT EXISTS idx_pending_questions_status + ON pending_questions (status); + +-- budget_ledger: the persistent shared Claude budget ledger (§6.1, §6.6). +-- One row per accounted spend event; the shared daily cap and the +-- interactive-first reserve are computed by summing over a UTC day. Spend is +-- recorded across ALL R720 Claude work (pipeline + Plane-1 sweeps). +CREATE TABLE IF NOT EXISTS budget_ledger ( + entry_id INTEGER PRIMARY KEY AUTOINCREMENT, + thread_id TEXT, + stage TEXT, + model TEXT NOT NULL, + billing_mode TEXT NOT NULL, + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + usd_cost REAL NOT NULL DEFAULT 0.0, + recorded_at TEXT NOT NULL, + day_bucket TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_budget_ledger_day + ON budget_ledger (day_bucket); + +CREATE INDEX IF NOT EXISTS idx_budget_ledger_thread + ON budget_ledger (thread_id); + +-- schema_meta: single-row table recording the applied schema version so +-- migrate() can detect and step forward. +CREATE TABLE IF NOT EXISTS schema_meta ( + id INTEGER PRIMARY KEY CHECK (id = 1), + schema_version INTEGER NOT NULL +); diff --git a/agent-team/agent_team/state_store.py b/agent-team/agent_team/state_store.py new file mode 100644 index 0000000..55fb849 --- /dev/null +++ b/agent-team/agent_team/state_store.py @@ -0,0 +1,176 @@ +"""Atomic state-store utilities and integrity checking (design §6.7). + +All durable state on the R720 (LangGraph SQLite checkpoint, the +``pending_questions`` ledger, the budget ledger, the Plane-1 rotation/coverage +pointer) is written atomically (write-temp-then-fsync-then-rename) and +integrity-checked on load. "Integrity-checked" is concrete here: a +schema-version match plus a stored content hash. On any mismatch the loader +refuses to proceed silently and raises :class:`IntegrityError` so the +coordinator can park the affected task with an ALARM rather than acting on +corrupt state. + +This module is pure stdlib (``os``, ``tempfile``, ``hashlib``, ``pathlib``) +and depends on no other ``agent_team`` module. The leaf builders import these +signatures verbatim, so they are intentionally explicit and final. +""" + +from __future__ import annotations + +import hashlib +import json +import os +import tempfile +from pathlib import Path + +__all__ = [ + "IntegrityError", + "atomic_write", + "compute_content_hash", + "read_checked", +] + +# Sidecar files sit next to the protected payload and carry the integrity +# metadata (schema version + content hash). Keeping them separate from the +# payload means the payload bytes round-trip unchanged. +_META_SUFFIX = ".meta.json" + +# Algorithm used for the stored content hash. Recorded in the sidecar so a +# future algorithm change stays backward-readable. +_HASH_ALGO = "sha256" + + +class IntegrityError(Exception): + """Raised when durable state fails its integrity check on load. + + Signals a schema-version mismatch, a missing/garbled integrity sidecar, or + a stored-content-hash mismatch (corruption or tampering). Callers treat + this as "refuse to proceed silently": park the task and ALARM rather than + restart blindly (§6.7). + """ + + +def compute_content_hash(data: bytes) -> str: + """Return the hex content hash for ``data`` (sha256). + + The same routine is used when writing the sidecar and when verifying on + load, so the two are guaranteed consistent. + """ + return hashlib.new(_HASH_ALGO, data).hexdigest() + + +def _meta_path(path: Path) -> Path: + """Return the sidecar metadata path for a payload ``path``.""" + return path.with_name(path.name + _META_SUFFIX) + + +def atomic_write(path: Path, data: bytes) -> None: + """Atomically write ``data`` to ``path`` (write-temp -> fsync -> rename). + + The bytes are written to a temporary file in the same directory, flushed + and ``fsync``-ed to durable storage, then ``os.replace``-d onto the final + path. ``os.replace`` is atomic on POSIX within a filesystem, so a reader + never observes a half-written file and a crash mid-write leaves either the + old payload or the new one, never a torn one. The containing directory is + ``fsync``-ed afterward so the rename itself is durable. + + This writes only the payload; integrity metadata is written by callers via + :func:`write_checked` / read back by :func:`read_checked`. (The sidecar is + written through this same primitive, so it is equally crash-safe.) + """ + path = Path(path) + directory = path.parent + directory.mkdir(parents=True, exist_ok=True) + + # delete=False so we control the rename; same dir guarantees same fs. + fd, tmp_name = tempfile.mkstemp( + prefix=path.name + ".", suffix=".tmp", dir=directory + ) + tmp_path = Path(tmp_name) + try: + with os.fdopen(fd, "wb") as handle: + handle.write(data) + handle.flush() + os.fsync(handle.fileno()) + os.replace(tmp_path, path) + except BaseException: + # Best-effort cleanup of the temp file on any failure. + try: + os.unlink(tmp_path) + except FileNotFoundError: + pass + raise + + # Make the rename itself durable by fsync-ing the directory. + dir_fd = os.open(directory, os.O_RDONLY) + try: + os.fsync(dir_fd) + except OSError: + # Some filesystems disallow directory fsync; the rename is still + # atomic, only its durability across power-loss is weakened. + pass + finally: + os.close(dir_fd) + + +def write_checked(path: Path, data: bytes, *, schema_version: int) -> None: + """Atomically write ``data`` plus its integrity sidecar. + + Writes the payload first, then the sidecar carrying ``schema_version`` and + the content hash. :func:`read_checked` verifies both. Both writes go + through :func:`atomic_write`, so each is crash-safe; if a crash lands + between them the sidecar is simply stale/absent and :func:`read_checked` + fails closed with :class:`IntegrityError`, which is the intended + refuse-to-proceed behaviour. + """ + path = Path(path) + atomic_write(path, data) + meta = { + "schema_version": int(schema_version), + "hash_algo": _HASH_ALGO, + "content_hash": compute_content_hash(data), + } + atomic_write(_meta_path(path), json.dumps(meta, sort_keys=True).encode("utf-8")) + + +def read_checked(path: Path, *, schema_version: int) -> bytes: + """Read and integrity-check ``path``, returning its bytes. + + Verifies the integrity sidecar exists, that its recorded + ``schema_version`` matches the expected ``schema_version``, and that the + stored content hash matches a freshly computed hash of the payload bytes. + Any mismatch (missing/garbled sidecar, schema drift, corruption/tampering) + raises :class:`IntegrityError`. + """ + path = Path(path) + try: + data = path.read_bytes() + except FileNotFoundError as exc: + raise IntegrityError(f"state payload missing: {path}") from exc + + meta_path = _meta_path(path) + try: + raw_meta = meta_path.read_bytes() + except FileNotFoundError as exc: + raise IntegrityError(f"integrity sidecar missing: {meta_path}") from exc + + try: + meta = json.loads(raw_meta) + except (ValueError, UnicodeDecodeError) as exc: + raise IntegrityError(f"integrity sidecar unreadable: {meta_path}") from exc + + stored_version = meta.get("schema_version") + if stored_version != schema_version: + raise IntegrityError( + f"schema-version mismatch for {path}: " + f"stored={stored_version!r} expected={schema_version!r}" + ) + + stored_hash = meta.get("content_hash") + actual_hash = compute_content_hash(data) + if stored_hash != actual_hash: + raise IntegrityError( + f"content-hash mismatch for {path}: " + f"stored={stored_hash!r} actual={actual_hash!r}" + ) + + return data diff --git a/agent-team/agent_team/task_model.py b/agent-team/agent_team/task_model.py new file mode 100644 index 0000000..2e09a91 --- /dev/null +++ b/agent-team/agent_team/task_model.py @@ -0,0 +1,155 @@ +"""Task-record / thread model + LangGraph graph-state schema (design §3.3). + +A task is a long-lived, resumable record (a LangGraph thread). This module is +the pure model layer — no I/O — defining: + +* :class:`TaskStatus` / :class:`Phase` — task lifecycle enums. +* :class:`TaskRecord` — the durable task record (§3.3 "the task record holds: + status, current phase, the full Q&A history, the plan, review verdicts, the + candidate diff, and CI results"). +* :class:`PipelineState` — a ``TypedDict`` used as the LangGraph graph state + schema; its keys mirror :class:`TaskRecord` fields. +* :func:`new_thread_id` — uuid thread-id minting. +* JSON serialization helpers (:func:`task_to_dict` / :func:`task_from_dict` / + :func:`task_to_json` / :func:`task_from_json`). + +The signatures here are CONTRACTS leaf builders import verbatim. +""" + +from __future__ import annotations + +import json +import uuid +from dataclasses import asdict, dataclass, field +from enum import Enum +from typing import Any, TypedDict + +__all__ = [ + "Phase", + "PipelineState", + "TaskRecord", + "TaskStatus", + "new_thread_id", + "task_from_dict", + "task_from_json", + "task_to_dict", + "task_to_json", +] + + +class TaskStatus(Enum): + """Top-level task lifecycle status. + + ``ACTIVE`` — progressing through stages. ``WAITING_HUMAN`` — suspended on a + LangGraph ``interrupt()`` awaiting Adam's answer. ``PARKED`` — stalled + (no answer in window, N failed build loops, or budget contention) and + ALARM-ed rather than spinning (§3.3, §6.6). ``DONE`` — draft PR + report + produced. ``FAILED`` — terminal failure. + """ + + ACTIVE = "active" + WAITING_HUMAN = "waiting_human" + PARKED = "parked" + DONE = "done" + FAILED = "failed" + + +class Phase(Enum): + """Pipeline phase the task is currently in (§3.3).""" + + INTAKE = "intake" + CLARIFY = "clarify" + PLAN = "plan" + REVIEW = "review" + BUILD = "build" + VERIFY = "verify" + PARKED = "parked" + DONE = "done" + + +def new_thread_id() -> str: + """Mint a fresh unique ``thread_id`` (uuid4 hex).""" + return uuid.uuid4().hex + + +@dataclass +class TaskRecord: + """The durable per-task record (§3.3, §3.3.1). + + Mirrors the LangGraph thread state; the SQLite checkpointer persists the + graph state while this record is the logical view the coordinator reasons + over. ``qa_history`` is the full clarifier Q&A; ``review_verdicts`` the + adversarial review outcomes; ``candidate_diff`` + ``diff_hash`` the builder + output and its ledger-recorded hash (§3.3.2); ``ci_results`` the + authenticated CI conclusion the verifier reads. + """ + + thread_id: str + status: TaskStatus + current_phase: Phase + qa_history: list[Any] = field(default_factory=list) + plan: dict[str, Any] | None = None + review_verdicts: list[Any] = field(default_factory=list) + candidate_diff: str | None = None + diff_hash: str | None = None + ci_results: dict[str, Any] | None = None + transport: str = "" + created_at: str | None = None + updated_at: str | None = None + + +class PipelineState(TypedDict, total=False): + """LangGraph graph-state schema; keys mirror :class:`TaskRecord` (§3.3). + + Used as the graph's state type. ``total=False`` so a node may write a + subset of keys per checkpoint transition. + """ + + thread_id: str + status: str + current_phase: str + qa_history: list[Any] + plan: dict[str, Any] | None + review_verdicts: list[Any] + candidate_diff: str | None + diff_hash: str | None + ci_results: dict[str, Any] | None + transport: str + created_at: str | None + updated_at: str | None + + +def task_to_dict(record: TaskRecord) -> dict[str, Any]: + """Serialize a :class:`TaskRecord` to a JSON-safe dict (enums -> values).""" + data = asdict(record) + data["status"] = record.status.value + data["current_phase"] = record.current_phase.value + return data + + +def task_from_dict(data: dict[str, Any]) -> TaskRecord: + """Rebuild a :class:`TaskRecord` from a :func:`task_to_dict` dict.""" + return TaskRecord( + thread_id=data["thread_id"], + status=TaskStatus(data["status"]), + current_phase=Phase(data["current_phase"]), + qa_history=list(data.get("qa_history", [])), + plan=data.get("plan"), + review_verdicts=list(data.get("review_verdicts", [])), + candidate_diff=data.get("candidate_diff"), + diff_hash=data.get("diff_hash"), + ci_results=data.get("ci_results"), + transport=data.get("transport", ""), + created_at=data.get("created_at"), + updated_at=data.get("updated_at"), + ) + + +def task_to_json(record: TaskRecord) -> str: + """Serialize a :class:`TaskRecord` to a JSON string.""" + return json.dumps(task_to_dict(record), sort_keys=True) + + +def task_from_json(payload: str | bytes) -> TaskRecord: + """Deserialize a :class:`TaskRecord` from a JSON string/bytes.""" + return task_from_dict(json.loads(payload)) diff --git a/agent-team/agent_team/transport/__init__.py b/agent-team/agent_team/transport/__init__.py new file mode 100644 index 0000000..b3d3acd --- /dev/null +++ b/agent-team/agent_team/transport/__init__.py @@ -0,0 +1,17 @@ +"""Transport seam for the durable human-in-the-loop responder (design §3.3.1). + +The ledger + resume logic are transport-independent; concrete Slack / GitHub / +Claude-Code adapters subclass :class:`Transport` in the leaves. +""" + +from agent_team.transport.base import ( + NormalizedAnswer, + QuestionSet, + Transport, +) + +__all__ = [ + "NormalizedAnswer", + "QuestionSet", + "Transport", +] diff --git a/agent-team/agent_team/transport/base.py b/agent-team/agent_team/transport/base.py new file mode 100644 index 0000000..015cfda --- /dev/null +++ b/agent-team/agent_team/transport/base.py @@ -0,0 +1,110 @@ +"""Transport interface ABC + payload dataclasses (design §3.3.1). + +The durable human-in-the-loop responder owns a notify+resume seam that is +transport-agnostic. This module defines the contract every adapter implements: + +* :class:`Transport` — abstract base with ``post_question`` (deliver a + question-set, return a ``channel_ref`` that embeds the ``question_id``) and + ``parse_answer`` (normalize an inbound raw payload to + ``(question_id, answer, via)``). +* :class:`QuestionSet` — the question-set payload carried by a LangGraph + ``interrupt()``. +* :class:`NormalizedAnswer` — the normalized inbound answer the responder + feeds into the §3.3.1 first-answer-wins compare-and-set. + +Concrete Slack / GitHub / Claude-Code adapters subclass :class:`Transport` in +the leaves. The signatures here are CONTRACTS the leaf builders import +verbatim, so they are explicit and final. +""" + +from __future__ import annotations + +from abc import ABC, abstractmethod +from dataclasses import dataclass, field +from typing import Any + +__all__ = [ + "NormalizedAnswer", + "QuestionSet", + "Transport", +] + +# Marker template embedded in transports without native callback metadata +# (e.g. a GitHub issue comment), so an inbound answer can be mapped back to its +# question. Slack embeds the question_id in ``callback_id`` instead. +GITHUB_MARKER_TEMPLATE = "" + + +@dataclass +class QuestionSet: + """A set of questions delivered to Adam for one ``turn`` of a task (§3.3.1). + + Carried in the ``interrupt()`` payload alongside ``thread_id``, + ``question_id``, ``turn``, ``transport``, and ``deadline``. ``questions`` is + the ordered list of prompts; ``context`` is optional rendering metadata + (repo, summary) the adapter may surface. + """ + + thread_id: str + question_id: str + turn: int + questions: list[str] + context: dict[str, Any] = field(default_factory=dict) + + +@dataclass +class NormalizedAnswer: + """A transport-normalized inbound answer (§3.3.1). + + The responder maps this into the first-answer-wins compare-and-set: + ``UPDATE ... SET status='answered' ... WHERE question_id=? AND + status='open'``. ``via`` records the answering channel/identity for the + audit trail (``answered_via``). + """ + + question_id: str + answer: Any + via: str + + +class Transport(ABC): + """Abstract transport adapter (§3.3.1). + + Subclasses implement delivery and answer parsing for one channel. The + ledger and resume worker depend only on this interface, so an adapter can + ship first (Slack) and others follow without touching the durable core. + """ + + @abstractmethod + def post_question( + self, + *, + thread_id: str, + question_id: str, + turn: int, + question_set: QuestionSet, + deadline: str, + ) -> str: + """Deliver ``question_set`` and return its ``channel_ref``. + + The posted message MUST embed ``question_id`` so an inbound answer can + be mapped back (Slack ``callback_id``; a + ```` marker in a GitHub comment). The returned + ``channel_ref`` is the transport's locator for the post (Slack message + ``ts`` / issue-comment id / Claude session id) and is stored on the + ledger row so reconcile/recovery can act on it (§3.3.1). + """ + raise NotImplementedError + + @abstractmethod + def parse_answer(self, raw: Any) -> tuple[str, Any, str]: + """Normalize an inbound ``raw`` payload to ``(question_id, answer, via)``. + + Extracts the embedded ``question_id`` (from the Slack ``callback_id`` / + the GitHub marker / the Claude session), the answer value, and the + ``via`` channel identity. The responder feeds the result into the + atomic compare-and-set. Implementations may build a + :class:`NormalizedAnswer` internally and return its fields as the tuple + the contract specifies. + """ + raise NotImplementedError diff --git a/agent-team/tests/__init__.py b/agent-team/tests/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/agent-team/tests/conftest.py b/agent-team/tests/conftest.py new file mode 100644 index 0000000..900c650 --- /dev/null +++ b/agent-team/tests/conftest.py @@ -0,0 +1,13 @@ +"""Pytest configuration: make the ``agent_team`` package importable. + +Adds the ``agent-team/`` project root (the directory containing the +``agent_team`` package) to ``sys.path`` so tests can run without an editable +install. +""" + +import sys +from pathlib import Path + +_PROJECT_ROOT = Path(__file__).resolve().parents[1] +if str(_PROJECT_ROOT) not in sys.path: + sys.path.insert(0, str(_PROJECT_ROOT)) diff --git a/agent-team/tests/test_billing.py b/agent-team/tests/test_billing.py new file mode 100644 index 0000000..39097b8 --- /dev/null +++ b/agent-team/tests/test_billing.py @@ -0,0 +1,119 @@ +"""Unit tests for agent_team.billing (§3.1).""" + +from __future__ import annotations + +import pytest + +from agent_team import billing +from agent_team.billing import ( + BillingMode, + ClaudeResult, + claude_invoke, + resolve_mode, + set_invoker, +) + + +@pytest.fixture(autouse=True) +def _restore_invoker(): + """Restore the module invoker after each test.""" + original = billing._invoker + yield + billing._invoker = original + + +def test_resolve_mode_default_is_subscription(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv("AGENT_TEAM_BILLING_MODE", raising=False) + assert resolve_mode(None) is BillingMode.SUBSCRIPTION + assert resolve_mode({}) is BillingMode.SUBSCRIPTION + + +def test_resolve_mode_from_config_string() -> None: + assert resolve_mode({"billing_mode": "api"}) is BillingMode.API + assert resolve_mode({"billing_mode": "BEDROCK"}) is BillingMode.BEDROCK + + +def test_resolve_mode_from_config_enum() -> None: + assert resolve_mode({"billing_mode": BillingMode.API}) is BillingMode.API + + +def test_resolve_mode_from_env(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("AGENT_TEAM_BILLING_MODE", "bedrock") + assert resolve_mode(None) is BillingMode.BEDROCK + + +def test_resolve_mode_config_overrides_env(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("AGENT_TEAM_BILLING_MODE", "bedrock") + assert resolve_mode({"billing_mode": "api"}) is BillingMode.API + + +def test_resolve_mode_invalid_raises() -> None: + with pytest.raises(ValueError): + resolve_mode({"billing_mode": "carrier-pigeon"}) + + +def test_claude_invoke_delegates_with_resolved_mode() -> None: + captured: dict = {} + + def fake(prompt: str, *, mode: BillingMode, **kw): + captured["prompt"] = prompt + captured["mode"] = mode + captured["kw"] = kw + return ClaudeResult(text="ok", mode=mode) + + set_invoker(fake) + result = claude_invoke("hi", mode=BillingMode.API, temperature=0.2) + assert result.text == "ok" + assert captured["mode"] is BillingMode.API + assert captured["prompt"] == "hi" + assert captured["kw"] == {"temperature": 0.2} + + +def test_subscription_mode_pops_stray_api_key(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("ANTHROPIC_API_KEY", "sk-should-be-hidden") + seen: dict = {} + + def fake(prompt: str, *, mode: BillingMode, **kw): + import os + + seen["key_present"] = "ANTHROPIC_API_KEY" in os.environ + return ClaudeResult(text="ok", mode=mode) + + set_invoker(fake) + claude_invoke("hi", mode=BillingMode.SUBSCRIPTION) + assert seen["key_present"] is False + # Restored after the call. + import os + + assert os.environ.get("ANTHROPIC_API_KEY") == "sk-should-be-hidden" + + +def test_api_mode_does_not_pop_api_key(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("ANTHROPIC_API_KEY", "sk-metered") + seen: dict = {} + + def fake(prompt: str, *, mode: BillingMode, **kw): + import os + + seen["key_present"] = "ANTHROPIC_API_KEY" in os.environ + return ClaudeResult(text="ok", mode=mode) + + set_invoker(fake) + claude_invoke("hi", mode=BillingMode.API) + assert seen["key_present"] is True + + +def test_subscription_with_no_key_is_safe(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv("ANTHROPIC_API_KEY", raising=False) + set_invoker(lambda prompt, *, mode, **kw: ClaudeResult(text="ok", mode=mode)) + assert claude_invoke("hi", mode=BillingMode.SUBSCRIPTION).text == "ok" + + +def test_unconfigured_invoker_raises() -> None: + billing._invoker = billing._unconfigured_invoker + with pytest.raises(RuntimeError): + claude_invoke("hi", mode=BillingMode.API) + + +def test_billing_mode_enum_members() -> None: + assert {m.name for m in BillingMode} == {"SUBSCRIPTION", "API", "BEDROCK"} diff --git a/agent-team/tests/test_schema.py b/agent-team/tests/test_schema.py new file mode 100644 index 0000000..35e2c01 --- /dev/null +++ b/agent-team/tests/test_schema.py @@ -0,0 +1,228 @@ +"""Unit tests for agent_team.db.schema (§3.3.1, §6.7).""" + +from __future__ import annotations + +import sqlite3 +import threading +from pathlib import Path + +import pytest + +from agent_team.db.schema import ( + BUDGET_LEDGER_DDL, + PENDING_QUESTIONS_DDL, + QUESTION_STATES, + SCHEMA_VERSION, + answer_question, + connect, + expire_question, + init_db, + migrate, + supersede_question, +) + + +def _insert_open_question(conn: sqlite3.Connection, qid: str, turn: int = 0) -> None: + conn.execute( + "INSERT INTO pending_questions " + "(question_id, thread_id, turn, status, transport) " + "VALUES (?, 'thread-1', ?, 'open', 'slack')", + (qid, turn), + ) + + +def test_ddl_constants_are_nonempty_strings() -> None: + assert isinstance(PENDING_QUESTIONS_DDL, str) and PENDING_QUESTIONS_DDL + assert isinstance(BUDGET_LEDGER_DDL, str) and BUDGET_LEDGER_DDL + assert "pending_questions" in PENDING_QUESTIONS_DDL + assert "budget_ledger" in BUDGET_LEDGER_DDL + + +def test_schema_version_is_int() -> None: + assert isinstance(SCHEMA_VERSION, int) + + +def test_question_states_match_ddl_check() -> None: + assert QUESTION_STATES == ("open", "answered", "expired", "superseded") + for state in QUESTION_STATES: + assert f"'{state}'" in PENDING_QUESTIONS_DDL + + +def test_connect_sets_pragmas(tmp_path: Path) -> None: + conn = connect(tmp_path / "db.sqlite") + try: + assert conn.execute("PRAGMA journal_mode").fetchone()[0].lower() == "wal" + assert conn.execute("PRAGMA foreign_keys").fetchone()[0] == 1 + assert conn.execute("PRAGMA busy_timeout").fetchone()[0] >= 1 + finally: + conn.close() + + +def test_init_db_creates_tables(tmp_path: Path) -> None: + db = tmp_path / "db.sqlite" + init_db(db) + conn = connect(db) + try: + names = { + r[0] + for r in conn.execute( + "SELECT name FROM sqlite_master WHERE type='table'" + ).fetchall() + } + finally: + conn.close() + assert {"pending_questions", "budget_ledger", "schema_meta"} <= names + + +def test_init_db_is_idempotent(tmp_path: Path) -> None: + db = tmp_path / "db.sqlite" + init_db(db) + init_db(db) # must not raise + conn = connect(db) + try: + version = conn.execute( + "SELECT schema_version FROM schema_meta WHERE id=1" + ).fetchone()[0] + finally: + conn.close() + assert version == SCHEMA_VERSION + + +def test_pending_questions_status_check_constraint(tmp_path: Path) -> None: + db = tmp_path / "db.sqlite" + init_db(db) + conn = connect(db) + try: + with pytest.raises(sqlite3.IntegrityError): + conn.execute( + "INSERT INTO pending_questions " + "(question_id, thread_id, turn, status, transport) " + "VALUES ('q', 't', 0, 'bogus', 'slack')" + ) + finally: + conn.close() + + +def test_migrate_stamps_version(tmp_path: Path) -> None: + db = tmp_path / "db.sqlite" + conn = connect(db) + try: + migrate(conn) + version = conn.execute( + "SELECT schema_version FROM schema_meta WHERE id=1" + ).fetchone()[0] + # Tables exist after migrate. + conn.execute("SELECT 1 FROM pending_questions LIMIT 1") + conn.execute("SELECT 1 FROM budget_ledger LIMIT 1") + finally: + conn.close() + assert version == SCHEMA_VERSION + + +def test_answer_question_first_wins(tmp_path: Path) -> None: + db = tmp_path / "db.sqlite" + init_db(db) + conn = connect(db) + try: + _insert_open_question(conn, "q1") + first = answer_question( + conn, question_id="q1", answer_json='{"a":1}', answered_via="slack" + ) + second = answer_question( + conn, question_id="q1", answer_json='{"a":2}', answered_via="github" + ) + assert first is True + assert second is False # duplicate/late loses the compare-and-set + row = conn.execute( + "SELECT status, answer_json, answered_via, answered_at " + "FROM pending_questions WHERE question_id='q1'" + ).fetchone() + finally: + conn.close() + assert row["status"] == "answered" + assert row["answer_json"] == '{"a":1}' # first answer retained + assert row["answered_via"] == "slack" + assert row["answered_at"] + + +def test_answer_after_expire_loses(tmp_path: Path) -> None: + db = tmp_path / "db.sqlite" + init_db(db) + conn = connect(db) + try: + _insert_open_question(conn, "q2") + assert expire_question(conn, question_id="q2") is True + assert ( + answer_question( + conn, question_id="q2", answer_json="{}", answered_via="slack" + ) + is False + ) + status = conn.execute( + "SELECT status FROM pending_questions WHERE question_id='q2'" + ).fetchone()["status"] + finally: + conn.close() + assert status == "expired" + + +def test_expire_only_open(tmp_path: Path) -> None: + db = tmp_path / "db.sqlite" + init_db(db) + conn = connect(db) + try: + _insert_open_question(conn, "q3") + answer_question(conn, question_id="q3", answer_json="{}", answered_via="slack") + # already answered -> cannot expire + assert expire_question(conn, question_id="q3") is False + finally: + conn.close() + + +def test_supersede_open_or_answered(tmp_path: Path) -> None: + db = tmp_path / "db.sqlite" + init_db(db) + conn = connect(db) + try: + _insert_open_question(conn, "q4") + answer_question(conn, question_id="q4", answer_json="{}", answered_via="slack") + assert supersede_question(conn, question_id="q4") is True + # already superseded -> no-op + assert supersede_question(conn, question_id="q4") is False + finally: + conn.close() + + +def test_concurrent_answers_single_winner(tmp_path: Path) -> None: + """Two threads racing to answer the same open question: exactly one wins.""" + db = tmp_path / "db.sqlite" + init_db(db) + seed = connect(db) + try: + _insert_open_question(seed, "race") + finally: + seed.close() + + results: list[bool] = [] + barrier = threading.Barrier(2) + lock = threading.Lock() + + def worker(via: str) -> None: + conn = connect(db) + try: + barrier.wait() + won = answer_question( + conn, question_id="race", answer_json='{"v":1}', answered_via=via + ) + with lock: + results.append(won) + finally: + conn.close() + + threads = [threading.Thread(target=worker, args=(f"c{i}",)) for i in range(2)] + for t in threads: + t.start() + for t in threads: + t.join() + + assert sorted(results) == [False, True] diff --git a/agent-team/tests/test_state_store.py b/agent-team/tests/test_state_store.py new file mode 100644 index 0000000..4a6d2ae --- /dev/null +++ b/agent-team/tests/test_state_store.py @@ -0,0 +1,111 @@ +"""Unit tests for agent_team.state_store (§6.7).""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from agent_team import state_store +from agent_team.state_store import ( + IntegrityError, + atomic_write, + compute_content_hash, + read_checked, + write_checked, +) + + +def test_compute_content_hash_is_sha256_hex() -> None: + import hashlib + + data = b"hello world" + assert compute_content_hash(data) == hashlib.sha256(data).hexdigest() + + +def test_compute_content_hash_distinguishes_inputs() -> None: + assert compute_content_hash(b"a") != compute_content_hash(b"b") + + +def test_atomic_write_creates_file_with_exact_bytes(tmp_path: Path) -> None: + target = tmp_path / "state.bin" + payload = b"\x00\x01binary\xff" + atomic_write(target, payload) + assert target.read_bytes() == payload + + +def test_atomic_write_creates_missing_parent_dirs(tmp_path: Path) -> None: + target = tmp_path / "nested" / "deep" / "state.bin" + atomic_write(target, b"x") + assert target.read_bytes() == b"x" + + +def test_atomic_write_overwrites_existing(tmp_path: Path) -> None: + target = tmp_path / "state.bin" + atomic_write(target, b"old") + atomic_write(target, b"new-and-longer") + assert target.read_bytes() == b"new-and-longer" + + +def test_atomic_write_leaves_no_temp_files(tmp_path: Path) -> None: + target = tmp_path / "state.bin" + atomic_write(target, b"data") + leftovers = [p for p in tmp_path.iterdir() if p.name != "state.bin"] + assert leftovers == [] + + +def test_write_then_read_checked_roundtrip(tmp_path: Path) -> None: + target = tmp_path / "state.bin" + payload = json.dumps({"k": "v"}).encode() + write_checked(target, payload, schema_version=3) + assert read_checked(target, schema_version=3) == payload + + +def test_read_checked_schema_version_mismatch_raises(tmp_path: Path) -> None: + target = tmp_path / "state.bin" + write_checked(target, b"data", schema_version=1) + with pytest.raises(IntegrityError): + read_checked(target, schema_version=2) + + +def test_read_checked_content_corruption_raises(tmp_path: Path) -> None: + target = tmp_path / "state.bin" + write_checked(target, b"original", schema_version=1) + # Corrupt the payload without touching the sidecar -> hash mismatch. + target.write_bytes(b"tampered") + with pytest.raises(IntegrityError): + read_checked(target, schema_version=1) + + +def test_read_checked_missing_payload_raises(tmp_path: Path) -> None: + with pytest.raises(IntegrityError): + read_checked(tmp_path / "nope.bin", schema_version=1) + + +def test_read_checked_missing_sidecar_raises(tmp_path: Path) -> None: + target = tmp_path / "state.bin" + # Plain atomic_write writes payload but NOT the integrity sidecar. + atomic_write(target, b"data") + with pytest.raises(IntegrityError): + read_checked(target, schema_version=1) + + +def test_read_checked_garbled_sidecar_raises(tmp_path: Path) -> None: + target = tmp_path / "state.bin" + write_checked(target, b"data", schema_version=1) + meta_path = target.with_name(target.name + ".meta.json") + meta_path.write_bytes(b"not-json{{{") + with pytest.raises(IntegrityError): + read_checked(target, schema_version=1) + + +def test_module_exports_public_contract() -> None: + for name in ( + "IntegrityError", + "atomic_write", + "compute_content_hash", + "read_checked", + ): + assert name in state_store.__all__ + assert hasattr(state_store, name) diff --git a/agent-team/tests/test_task_model.py b/agent-team/tests/test_task_model.py new file mode 100644 index 0000000..28886c0 --- /dev/null +++ b/agent-team/tests/test_task_model.py @@ -0,0 +1,110 @@ +"""Unit tests for agent_team.task_model (§3.3, §3.3.1).""" + +from __future__ import annotations + +from agent_team.task_model import ( + Phase, + PipelineState, + TaskRecord, + TaskStatus, + new_thread_id, + task_from_dict, + task_from_json, + task_to_dict, + task_to_json, +) + + +def test_new_thread_id_unique_hex() -> None: + a = new_thread_id() + b = new_thread_id() + assert a != b + assert len(a) == 32 + int(a, 16) # must be valid hex + + +def test_phase_members() -> None: + assert {p.name for p in Phase} == { + "INTAKE", + "CLARIFY", + "PLAN", + "REVIEW", + "BUILD", + "VERIFY", + "PARKED", + "DONE", + } + + +def test_task_record_defaults() -> None: + rec = TaskRecord( + thread_id="t1", + status=TaskStatus.ACTIVE, + current_phase=Phase.INTAKE, + ) + assert rec.qa_history == [] + assert rec.plan is None + assert rec.review_verdicts == [] + assert rec.candidate_diff is None + assert rec.diff_hash is None + assert rec.ci_results is None + assert rec.transport == "" + + +def test_to_dict_serializes_enums_to_values() -> None: + rec = TaskRecord( + thread_id="t1", + status=TaskStatus.WAITING_HUMAN, + current_phase=Phase.CLARIFY, + ) + data = task_to_dict(rec) + assert data["status"] == "waiting_human" + assert data["current_phase"] == "clarify" + + +def test_roundtrip_dict() -> None: + rec = TaskRecord( + thread_id="t1", + status=TaskStatus.PARKED, + current_phase=Phase.PLAN, + qa_history=[{"q": "x", "a": "y"}], + plan={"phases": [1, 2]}, + review_verdicts=["REQUEST_CHANGES"], + candidate_diff="diff --git a b", + diff_hash="deadbeef", + ci_results={"conclusion": "success"}, + transport="slack", + created_at="2026-06-17T00:00:00Z", + updated_at="2026-06-17T01:00:00Z", + ) + restored = task_from_dict(task_to_dict(rec)) + assert restored == rec + + +def test_roundtrip_json() -> None: + rec = TaskRecord( + thread_id="t2", + status=TaskStatus.DONE, + current_phase=Phase.DONE, + diff_hash="abc", + ) + restored = task_from_json(task_to_json(rec)) + assert restored == rec + assert restored.status is TaskStatus.DONE + assert restored.current_phase is Phase.DONE + + +def test_pipeline_state_keys_mirror_task_record() -> None: + # Every PipelineState key should be a TaskRecord field. + state_keys = set(PipelineState.__annotations__) + record_fields = set(TaskRecord.__dataclass_fields__) + assert state_keys == record_fields + + +def test_pipeline_state_usable_as_dict() -> None: + state: PipelineState = { + "thread_id": "t1", + "status": "active", + "current_phase": "intake", + } + assert state["thread_id"] == "t1" diff --git a/agent-team/tests/test_transport_base.py b/agent-team/tests/test_transport_base.py new file mode 100644 index 0000000..a46364e --- /dev/null +++ b/agent-team/tests/test_transport_base.py @@ -0,0 +1,92 @@ +"""Unit tests for agent_team.transport.base (§3.3.1).""" + +from __future__ import annotations + +from typing import Any + +import pytest + +from agent_team.transport.base import ( + GITHUB_MARKER_TEMPLATE, + NormalizedAnswer, + QuestionSet, + Transport, +) + + +def test_transport_is_abstract() -> None: + with pytest.raises(TypeError): + Transport() # type: ignore[abstract] + + +def test_question_set_fields() -> None: + qs = QuestionSet( + thread_id="t1", + question_id="q1", + turn=2, + questions=["a?", "b?"], + context={"repo": "x"}, + ) + assert qs.thread_id == "t1" + assert qs.question_id == "q1" + assert qs.turn == 2 + assert qs.questions == ["a?", "b?"] + assert qs.context == {"repo": "x"} + + +def test_question_set_context_defaults_empty() -> None: + qs = QuestionSet(thread_id="t", question_id="q", turn=0, questions=[]) + assert qs.context == {} + + +def test_normalized_answer_fields() -> None: + ans = NormalizedAnswer(question_id="q1", answer={"choice": 1}, via="slack") + assert ans.question_id == "q1" + assert ans.answer == {"choice": 1} + assert ans.via == "slack" + + +def test_concrete_subclass_implements_contract() -> None: + class FakeTransport(Transport): + def __init__(self) -> None: + self.posted: dict[str, Any] = {} + + def post_question( + self, *, thread_id, question_id, turn, question_set, deadline + ) -> str: + ref = f"slack-ts-{question_id}" + self.posted = { + "thread_id": thread_id, + "question_id": question_id, + "turn": turn, + "deadline": deadline, + "ref": ref, + } + return ref + + def parse_answer(self, raw) -> tuple[str, Any, str]: + na = NormalizedAnswer( + question_id=raw["callback_id"], answer=raw["value"], via="slack" + ) + return na.question_id, na.answer, na.via + + t = FakeTransport() + qs = QuestionSet(thread_id="t1", question_id="q1", turn=0, questions=["?"]) + ref = t.post_question( + thread_id="t1", + question_id="q1", + turn=0, + question_set=qs, + deadline="2026-06-18T00:00:00Z", + ) + assert ref == "slack-ts-q1" + assert t.posted["question_id"] == "q1" + + parsed = t.parse_answer({"callback_id": "q1", "value": "yes"}) + assert parsed == ("q1", "yes", "slack") + + +def test_github_marker_embeds_question_id() -> None: + marker = GITHUB_MARKER_TEMPLATE.format(question_id="abc123") + assert marker == "" + assert "abc123" in marker diff --git a/docs/r720-agent-team-design.md b/docs/r720-agent-team-design.md new file mode 100644 index 0000000..a4e4530 --- /dev/null +++ b/docs/r720-agent-team-design.md @@ -0,0 +1,490 @@ +# R720 Agent Team — Design (v2) + +Status: **DESIGN LOCKED — ready to build. Not built yet.** **v5** folded the final plan-review refinements +(SQLite `BEGIN IMMEDIATE` for the compare-and-set §3.3.1; CI denylist defense-in-depth + authenticated-only +pass/fail gate + compromised-box honesty §3.3.2; contention reserve + parked-task aging §6.6; backup integrity +definition + post-restore reconciliation §6.7; tested rollbacks §7; per-role canaries §6.4; CLI audit/confirm +§3.3.1). The design went through 3 GPT-4.1 `sh-plan-review` cycles (v1, v3, v4); the architecture was stable +throughout and remaining grain is now build-time implementation detail captured in the P1/P3 exit gates. **No +further gate runs by decision; build may begin with Phase 0 / P1.** Earlier status: Drafted 2026-06-17. **v3** reframes around Adam's clarified north star: the R720 +is not just scheduled checkers, it is a self-hosted, human-gated **agentic SDLC pipeline** +(intake -> clarify -> plan -> review -> build -> verify), fed from the Mac harness and tickets, with the +scheduled checkers as one task source. v2's roster work becomes "Plane 1"; the pipeline is "Plane 2" and the +centerpiece. **v4** folds in every v3 plan-review finding: B2 (§3.3.1) and B4 (§3.3.2) resolved with P1/P3 exit +gates; B1 billing realism + contention (§6.6); B3/B5/F5 provisioning + rollback + cross-review gate (§7); F1 +state durability + backup (§6.7); F2 canary update process (§6.4); F3 service-account lifecycle + F1 backups +(§9); F4 runbook incident handling (Phase 6); Q3 escalation ladder (§5); Q1/Q2 in §3.3.1. Ready for a re-run of +`sh-plan-review`. No build until that gate passes and Adam approves. + +Extends the `sh-secrev` pattern (see `security-review/DEPLOY-R720.md`) from a single security sweep into a small +roster of scheduled, unattended agents that check, plan, and (carefully) build, coordinated by a Claude `claude +-p` brain and using GPT / Gemini / DeepSeek where each is the better fit. + +This doc is the plan, not a runbook. It deliberately reuses the security agent's proven substrate. + +## 0. Locked decisions (this revision) + +| # | Decision | Choice | +|---|---|---| +| D1 | Headless `claude -p` under Max | **Permitted for now** (Anthropic pushed the disallowing ToS change to a later, unannounced date). Build the auth path **swappable** so the cutover is a config flip. See `reference_claude_subscription_billing`. | +| D2 | Fixer write path | **Option B**: the always-on box stays read-only; it emits a patch + opens an issue, and a trusted org CI workflow (OIDC) applies the patch on a branch and opens the **draft** PR. No standing write token on the box. | +| D3 | Checker/planner output mode | **Report + ALARM-only to start.** Clean nights post nothing; confirmed criticals alarm Slack; everything else lands in a mode-600 report. No auto-Jira/Notion writes until signal quality is trusted. | +| D4 | Billing/auth resilience | **Build a billing-mode abstraction now** (subscription OAuth default, API-key / Bedrock fallback ready). | +| D5 | aws-posture | **Resident on the box via IAM Roles Anywhere**, with a new **step-ca** internal CA for automated short-lived leaf-cert rotation (no long-lived AWS key on the box). | +| D6 | Confluence write identity | **Dedicated `confluence-bot` Atlassian service account, edit scoped to the IT space only** (Confluence API tokens inherit the whole user's permissions, so a scoped service account is how we bound blast radius). Token in `~/secrev.env` (mode 600). Costs one Confluence seat. | +| D7 | Confluence agent modes + Mermaid | **Scheduled = read + recommend only** (gaps/staleness into the report, per D3, never auto-writes). **On-demand = SSH-invoked from Adam's Mac** for an actual write. Mermaid map edits go through `~/.claude/scripts/confluence_mermaid.py` (ADF-only, dry-run-default, macro-count + revert-diff guarded). | +| D8 | Two-plane architecture | **Plane 1** = scheduled checkers (v2 roster), which also act as a task source. **Plane 2** = the agentic SDLC task pipeline (the centerpiece). Shared substrate. | +| D9 | Pipeline foundation | **LangGraph (the open-source library, runs in-process, NOT SaaS) + a local SQLite checkpointer.** Durable, resumable graph; `interrupt()` for the human gate. Reuses the existing stack. | +| D10 | Human-in-the-loop transport | **Pluggable, all three adapters** (Slack Block Kit, GitHub/Jira ticket comments, Claude Code on the Mac); Adam picks the channel per task. A transport-agnostic notify/resume layer maps answers back to the task thread. | +| D11 | Build/verify execution | **Org CI is the primary sandbox.** Builders emit candidate diffs; the Option-B OIDC workflow builds/tests/security-reviews; the verifier agent reads CI results. Keeps the 4GB box light + read-only. Dedicated builder VM only if fast local loops prove necessary. | +| D12 | Observability (LangSmith) | **Deprecate LangSmith (the SaaS tracer).** Keep LangGraph (framework, local). Local JSONL (`telemetry.py`, already exists) is the default; self-hosted Phoenix optional later. No SaaS dependency. | +| D13 | Task intake | **Mac harness first** (SSH-invoke enqueues onto the box), **GitHub issues next** (phased). | + +## 1. Goal and scope + +Stand up a coordinated team of agents on the existing `sh-secrev` R720 VM that runs unattended on a schedule, +operates across every Sea-Haven-Industries org repo automatically, stays cost-bounded, and reports through one +alarm channel. Claude (subscription OAuth) coordinates and does deep reasoning; Gemini does broad scans; GPT does +adversarial cross-checks; DeepSeek does mechanical code edits. + +In scope: read-mostly checkers, a low-blast-radius planner, a resident AWS posture check (D5), and a CI-applied +draft-PR fixer (D2). Out of scope: anything interactive or needing back-and-forth, and (per the secrev ethos) any +standing **write** credential on the always-on box. + +This does **not** replace the security review agent; it sits beside it and reuses its plumbing. + +## 2. What we reuse vs. what is new + +Reused from `sh-secrev` as-is: the VM, the systemd-timer model, the OAuth billing path, clean-clone +auto-discovery into `~/repo-mirrors` (read-only PAT, scrubbed post-fetch; the team scans the **same mirrors**, it +does not re-clone), the budget primitives (per-call + total caps, fail-toward-over-reporting), ALARM-only Slack, +mode-600 reports, the anti-complacency canary + coverage-rotation idea, the `orchestrator/` rsync deploy, and the +non-Claude provider keys in `~/orchestrator/.env`. + +New: a **coordinator** runner (`agent-team/` dir; entry CLI `run-team.py`, importable modules stay snake_case per +Python rules; the resource/dir name is kebab-case per handbook), per-role prompt/checklist modules each with its +own canary, a **billing-mode abstraction** (D4), the **fixer patch -> CI -> draft-PR** path (D2), and the +**step-ca + Roles Anywhere** setup for aws-posture (D5). + +## 3. Architecture + +``` +systemd timer (shared with secrev — see §8) + │ + ▼ + mirror step (reuse sh-secrev discovery) ──► ~/repo-mirrors (read-only) + │ + ▼ + COORDINATOR (claude -p via billing-mode abstraction; read-only tools + Bash to call orchestrator/gh) + │ shared budget ledger + versioned rotation/coverage state + ├──────────────┬───────────────┬────────────────┬──────────────┬───────────┐ + ▼ ▼ ▼ ▼ ▼ ▼ + drift checker dep/CVE checker doc-drift checker aws-posture planner (each role + (Gemini scan (Claude + (Gemini large- (Sonnet via (Claude) has its own + + Claude judge) GPT tiebreak) context) Roles Anywhere) canary) + │ │ │ │ │ + └──────────────┴───────────────┴────────────────┴──────────────┘ + │ structured JSON per agent + ▼ + coordinator: dedup + prioritize + route + │ + ┌──────────────────────────┼───────────────────────────┐ + ▼ ▼ ▼ + Slack ALARM mode-600 report fix-spec queue (D2) + (confirmed crit) (everything else; │ + no auto-Jira/Notion yet, D3) ▼ + FIXER: DeepSeek edit + Claude + spec + GPT review ──► patch + issue + │ + ▼ org CI (OIDC) applies patch, + opens DRAFT PR, runs pre-push + hooks + CI gates + Claude Code App +``` + +Model assignment (matches how `orchestrator` already splits them): **Claude (subscription)** = coordinator, deep +checks, fix-spec authoring; **Gemini 2.5 Pro** = broad whole-repo scans; **GPT-4.1** = adversarial cross-check / +tiebreak / PR review; **DeepSeek** = mechanical patch writing. + +### 3.1 Billing-mode abstraction (D4) +A single `claude_invoke(...)` seam selects the Claude auth/billing path from config: `subscription` (OAuth token, +default today), `api` (metered `ANTHROPIC_API_KEY`), or `bedrock` (cross-account Bedrock, already used by secrev +for the rare cross-family tiebreak). Switching modes is a config flip, not a code change. The box still pops any +stray `ANTHROPIC_API_KEY` in `subscription` mode so OAuth cannot be silently overridden. + +### 3.2 Relationship to the LangSmith orchestrator (what stays vs what the R720 hosts) + +Two orchestrators coexist after this plan; neither replaces the other. The split is **trigger + Claude billing**, +not capability. + +- **LangSmith orchestrator (Mac-hosted, unchanged).** The existing `orchestrator/` (LangGraph router + memory + retriever + Composio connector + LangSmith tracing on the `orchestration` project) stays the **on-demand, + interactive** delegation path: one task -> retrieve memory -> route -> one agent -> result, API-billed. This is + the CLAUDE.md hybrid-delegation path Claude Code uses for cross-family review, large scans, fast coding, and + connector actions. It stays one-shot and stateless (Q1: not refactored for persistence). +- **R720 orchestrator (the agent-team coordinator, new).** The scheduled, unattended, multi-agent layer: + cadence, the coordinator brain, shared budget ledger, versioned rotation/coverage state, the clean-clone mirror + corpus, and per-role canaries. Claude work here runs **headless under subscription OAuth** (Agent SDK), + billing-mode-swappable (§3.1). + +| Component | Today (LangSmith orchestrator, Mac) | After this plan | +|---|---|---| +| Trigger | On-demand from Claude Code / CLI | + scheduled (systemd timer) and SSH-invoked on-demand, on the R720 | +| Execution shape | One task -> one agent (stateless) | + multi-agent coordination with shared budget + versioned state (R720) | +| Claude billing | Metered API key | **Subscription OAuth on the R720** (swappable to api/bedrock per §3.1) | +| Non-Claude (GPT-4.1 / Gemini / DeepSeek) | `run.py` router, API-billed, LangSmith-traced | **unchanged in shape** — the R720 coordinator calls the **local** `~/orchestrator/run.py` (already rsync'd to the box) for these single-shot sub-tasks, so they keep API billing + LangSmith tracing | +| Memory retriever + embeddings cache | Mac | reused read-only by both (the box's rsync'd copy embeds the same memory store) | +| Composio connector (Slack/Notion/GitHub) | Mac | reused; the R720 routes Slack/Jira/Notion through it. **Confluence stays OUT of the connector** — native Atlassian MCP on the Mac for interactive edits, `confluence_mermaid.py` + REST for the box | +| Observability | LangSmith SaaS tracing (`orchestration` project) | **LangSmith deprecated (D12)** — local JSONL (`telemetry.py`) default, self-hosted Phoenix optional. LangGraph framework stays (it is not SaaS) | +| `models.py` factories + model-ID constants | Mac | shared code (rsync'd); single source of truth for both | + +**What does NOT migrate (stays Mac / interactive):** the daily-driver Claude Code sessions, the hybrid on-demand +delegation, and interactive Confluence edits via the native Atlassian MCP. + +**What is genuinely NEW on the R720 (not a migration — these never existed in the LangSmith orchestrator):** +scheduling, the coordinator + shared state/budget, canary/coverage, and subscription-OAuth Claude. + +Net: the LangSmith orchestrator keeps its job (on-demand routing, non-Claude execution, tracing, connector); the +R720 becomes the host for everything **scheduled, stateful, and subscription-billed**, and it **reuses the +LangSmith orchestrator in place** (the local rsync'd copy) for the non-Claude single-shots rather than +re-implementing them. + +### 3.3 Plane 2 — the agentic SDLC pipeline (the centerpiece) + +A durable, human-gated task pipeline hosted on the R720. A task is a long-lived, resumable record; the +coordinator drives it through stages, asking Adam for input when it is not confident and handing off to the org +CI to actually build and verify. + +``` +INTAKE ─► CLARIFIER ─► PLANNER ─► REVIEW LOOP ─► BUILDERS ─► VERIFIERS ─► draft PR + report + │ │ │ │ │ │ +Mac harness asks Adam phased plan GPT-4.1 + Claude spec org CI builds/tests/ +(SSH-invoke) question- (Claude) multi-model + DeepSeek security-review; +GitHub issue sets until adversarial; edits ─► verifier reads results; +checker 98%+, then loops back candidate loops back to builders +finding HUMAN GATE to planner diff on failure +``` + +**Stages and model per stage:** +- **Intake** — a task enters from the Mac harness (D13, first), a GitHub issue (next), or a Plane-1 checker + finding. It is written as a new task record (LangGraph thread) with a unique `thread_id`. +- **Clarifier (Claude)** — gathers context (repo, memory, handbook), then asks Adam **question-sets until 98%+ + confident**. This is a LangGraph `interrupt()`: the task suspends and checkpoints, a question-set is delivered + over the chosen transport (D10), and the task resumes via `Command(resume=...)` when the answer arrives. The + **human gate**: no progression to build without the clarifier clearing the bar and Adam approving the plan. +- **Planner (Claude)** — produces a phased plan (the format these design docs use). +- **Review loop (GPT-4.1 + optional multi-model)** — adversarial plan review (the `sh-plan-review` / + `cross_reviewer` discipline). Loops back to the planner on REQUEST CHANGES; escalates to Adam if it cannot + converge. +- **Builders (Claude spec + DeepSeek edits)** — turn the approved plan into a **candidate diff**. They do not + write to repos; per D2/D11 they emit the diff for CI. +- **Verifiers (org CI + a Claude/GPT reader)** — CI (Option-B OIDC) applies the diff on a branch, builds, runs + tests + the security review + lint; the verifier agent reads the CI results and either loops back to builders + or advances. Confirmed pass produces a **draft PR** plus a report to Adam. + +**Durable state (D9).** LangGraph (local) + a SQLite checkpointer. Each stage transition is checkpointed, so a +crash, a budget pause, or an overnight wait on a human answer all resume cleanly instead of restarting. The task +record holds: status, current phase, the full Q&A history, the plan, review verdicts, the candidate diff, and CI +results. + +**Human-in-the-loop (D10).** A small transport-agnostic responder service owns the notify+resume seam: it posts +the interrupt's question-set to the channel Adam chose for that task (Slack Block Kit / ticket comment / a Claude +Code session) and maps his reply back to the right `thread_id` to resume it. Adapters are independent so one can +ship first (Slack) and the others follow. + +**Stability + autonomy bounds (the "stable" requirement).** Hard gates, not vibes: (1) no build before the +clarifier hits 98% AND Adam approves the plan; (2) draft PRs only, never auto-merge; (3) the verifier must pass +or the task loops/holds, never ships; (4) per-task budget cap inside the shared nightly cap (§6.1); (5) every +stage checkpointed so failures resume, not restart; (6) a task that stalls (no human answer within a window, or +N failed build loops) parks and ALARMs rather than spinning. Plane-1's canary/coverage discipline applies to the +pipeline's agents too. + +### 3.3.1 Durable human-in-the-loop suspend/resume (resolves B2) + +LangGraph `interrupt()` + the SQLite checkpointer suspend and resume the graph, but the checkpoint alone does not +track the human-interaction lifecycle (delivery, duplicate/late answers, expiry). So the pipeline adds one +durable source of truth, a SQLite `pending_questions` table, and resolves every race with an atomic +compare-and-set against it. This is the riskiest mechanic, so it is specified here and P1 must prove it. + +- **Identity.** Each task is a graph `thread_id`. Each question-set gets a `question_id` (uuid) and a monotonic + `turn` within the task. The interrupt payload carries `{thread_id, question_id, turn, question_set, transport, + deadline}`. +- **Ledger.** `pending_questions(question_id PK, thread_id, turn, status[open|answered|expired|superseded], + transport, channel_ref, posted_at, deadline_at, answer_json, answered_at, answered_via)`. The LangGraph + checkpoint holds graph state; this table holds the question lifecycle and is what delivery, the responder, and + recovery read. +- **Delivery (and lost-post).** On interrupt, write the row `open` first, then post to the chosen transport and + store its `channel_ref` (Slack message ts / issue-comment id / Claude session id). The posted question embeds + the `question_id` (Slack `callback_id`; a `` marker in a GitHub comment). If the post + fails, the row stays `open` with no ref and a reconcile loop retries idempotently. +- **Answer mapping + idempotency (first-answer-wins).** Each transport's inbound adapter normalizes an answer to + `(question_id, answer, via)`. The responder then runs one atomic statement: + `UPDATE pending_questions SET status='answered', answer_json=?, answered_via=? WHERE question_id=? AND + status='open'`. rowcount 1 = first valid answer, enqueue a resume job; rowcount 0 = the question was not open + (already answered/expired/superseded), so the answer is a duplicate or late and is ignored with a "already + closed" reply. This single compare-and-set makes duplicate clicks, transport redelivery, answers via two + channels, and answer-after-timeout all safe. The statement runs inside a `BEGIN IMMEDIATE` transaction (SQLite's + default deferred isolation does not serialize concurrent responders, so the check-and-set must take the write + lock up front). +- **Resume (single-flight, turn-guarded).** A resume worker serializes per `thread_id` and calls + `graph.invoke(Command(resume=answer), {configurable:{thread_id}})`. Before resuming it checks the live + checkpoint is still interrupted on this `turn`; if the graph already advanced (stale/redelivered job) it marks + the question `superseded` and skips. A resume can never double-apply. +- **Deadline / no-answer.** Each open question has `deadline_at`. A timer loop flips overdue `open` rows to + `expired` (same compare-and-set) and applies the task policy: park + ALARM Adam, or apply a defined default + answer. An answer arriving for an already-`expired` question loses the compare-and-set and is ignored. Timeout + vs answer is a deterministic race on flipping `open`. +- **Restart recovery.** All state is durable (both SQLite stores), so a reboot converges via a startup sweep: + retry delivery for `open` rows lacking a ref; re-enqueue resume for `answered` rows whose graph is still + interrupted on that turn (idempotent via the turn guard); run the deadline policy for overdue `open` rows. No + in-memory-only state. +- **Parallel tasks + isolation (Q1).** Each task is its own `thread_id` with its own checkpoint and ledger rows; + the resume worker serializes per thread but runs different threads concurrently within the budget cap. +- **Manual 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; 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. +- **Transport seam (Q2 fallback).** A `Transport` interface (`post_question(...) -> channel_ref`, + `parse_answer(raw) -> (question_id, answer, via)`) with Slack / GitHub / Claude Code adapters; the ledger + + resume logic are transport-independent. Adam picks the channel per task at intake. If the chosen transport is + unreachable, reconcile retries and, after N failures, falls back to a Slack ALARM pointing at the task. + +### 3.3.2 CI-as-verifier trust boundary (resolves B4) + +The builders are semi-trusted at best: an LLM that read repo content can be wrong or prompt-injected, so **the +candidate diff is treated as untrusted code.** The threat is that executing it in CI with org credentials lets a +bad diff exfiltrate secrets, assume the deploy role, or tamper with other repos. Five boundaries bound it: + +1. **Split CI: untrusted execution is credential-less; privileged steps never see the patch.** The job that + checks out and runs the diff (install/build/test) runs with `permissions: contents: read`, **no secrets, no + OIDC, no write token**, and restricted network egress. The patch executes only here, where there is nothing + to steal and nothing to assume. Any privileged action (the OIDC role, authoritative status, opening the PR) + runs in a **separate job that does not check out or execute patch-controlled code**; it consumes the + build/test report as data only. This is the standard untrusted-code-in-CI ("pwn request") mitigation, so the + workflow must NOT use `pull_request_target` with a checkout of the head ref. +2. **The patch may not touch the trust-control surface.** A box-side check and a CI guard both reject any + candidate diff that modifies `.github/workflows/**`, IAM/policy/permission IaC (CDK/SAM), branch-protection / + `CODEOWNERS` / Dependabot config, or files outside the task's declared scope. Such a diff is escalated to + mandatory human review + GPT cross-review, never auto-built (those files are the mandatory-cross-review + surface regardless). The match is not naive: enforcement is a **CI-side hard fail** (not only the box check), + it resolves symlinks and canonicalizes paths, and it rejects renames into denied paths and build steps that + generate files into them, so a path match cannot be bypassed by indirection. +3. **Diff integrity, box to CI.** The builder hashes the candidate diff and records the hash in the task ledger. + The box has no write token (D2), so the diff reaches CI as a signed artifact / via a short-lived + single-purpose token that can only create the agent branch; CI verifies the diff hash matches the ledger + before applying. A tampered or substituted diff fails the hash check. +4. **Pass/fail is a pure-code gate over authenticated CI results, not the LLM verifier.** Mirroring secrev's "one + pure-code script owns the block decision," a deterministic gate reads the CI run conclusion via the GitHub + Checks/Actions API (authenticated as the box read-only PAT), keyed to the specific run id + diff hash. It + consumes **only** that authenticated, patch-independent conclusion; it never trusts a success/failure file or + artifact the patch could have written. The verifier **agent** only reads failures to propose the next fix; it + cannot declare success. +5. **Branch protection is the final enforcement.** The draft PR targets a protected branch; the locked-down + build/test/security-review are required status checks; merge needs them green plus the security-review, the + Claude Code App review, and human approval. The agent never auto-merges (D2). + +The CI apply/verify workflow is itself IaC + untrusted-input handling, so before it ships it goes through +**`/sh-security-review` and the mandatory GPT-4.1 cross-review** (in P3), and its actions are SHA-pinned per the +handbook Pinning Principle. Provenance (diff hash, CI run id, checks consumed, gate decision) is recorded per +task; a diff that tried to touch the denylist, a build job that attempted egress, or a checks/hash mismatch is an +ALARM. Egress restriction in the untrusted job is enforced at the runner level (a harden-runner-style egress +policy or a self-hosted-runner network deny), but the primary mitigation is that the job holds no secrets, so a +successful exfil yields nothing of value. + +**Threat-model honesty (the diff-hash does not cover a compromised box).** The hash protects the box->CI transit +(tamper/substitution), not a box that is itself compromised: the box is in the trusted computing base for what it +*proposes*. The real backstops against a bad box are the credential-less build job, the trust-control-surface +denylist, and above all the **human review of the draft PR** plus required checks (nothing merges without Adam). +A signed-artifact/attestation step hardens transit further but is not claimed to defend a compromised box. + +## 4. Roster + +### Tier 1 — checkers (read-only) +| Agent | Models | Cadence | Output / gate | +|---|---|---|---| +| **compliance-drift** | Gemini scan + Claude judge | nightly | Drift vs engineering-handbook (naming, secrets placement, CI/CD present, Dependabot, branch protection). Report + Slack ALARM on violations (no auto-Jira yet, D3) | +| **dependency-cve** | Claude + GPT tiebreak | nightly | Cross-ref lockfiles vs advisories org-wide; report + feed fixer. Complements Dependabot | +| **doc-drift** | Gemini (large context) | weekly | Flags repos whose architecture moved but Confluence/README did not | + +### Tier 2 — aws-posture (resident, D5) + planner +| Agent | Models | Cadence | Output | +|---|---|---|---| +| **aws-posture** | Sonnet collectors + judge | weekly | Idle/anomalous spend (≈$330/mo flagged) + reasoning layer over baseline findings. Auths via **Roles Anywhere** (short-lived leaf certs, auto-rotated by step-ca). Complements existing GuardDuty/Security Hub/Config, does not replace them | +| **plan-groomer** | Claude | weekly | Drafts a groomed weekly plan **into the mode-600 report** for now (D3); auto-write to Notion/Jira is a later toggle once trusted | +| **confluence-doc** | Gemini scan + Claude judge | weekly (scheduled) + on-demand | **Scheduled:** diffs repos + AWS inventory + the page-ID map (`project_confluence_migration`) against Confluence, reports doc gaps / stale pages / missing runbooks (recommend-only, D3). **On-demand (SSH-invoked):** performs an actual update, including Mermaid map edits via `confluence_mermaid.py`. Writes as the IT-space-scoped `confluence-bot` (D6). Overlaps the existing `sh-confluence-audit`/`sh-confluence` skills; the box adds unattended cross-repo scope + the tested Mermaid script | + +### Tier 3 — fixer (D2) +| Agent | Models | Trigger | Output | +|---|---|---|---| +| **fixer** | DeepSeek edit + Claude spec + GPT review | on a confirmed, low-risk finding | Emits a patch + opens an issue; org CI applies it and opens a **draft** PR. Never auto-merges; pre-push hooks + CI + Claude Code App gate it (F3) | + +## 5. Coordination model + +Nightly, after the mirror refresh, the coordinator: (1) loads the shared budget ledger and the **versioned** +rotation/coverage state (F1); (2) runs the **canary suite first**, one planted-fault corpus per role, a miss is a +COMPLACENCY ALARM and that role is skipped; (3) fans out the scheduled agents over the mirrors, each with a +per-call cap, all drawing from **one shared total cap** (critical for the Claude subscription draw, §6.1), with +budget-exhausted roles deferred via the rotation pointer (never dropped) and a COVERAGE ALARM if a role slips +past `MAX_CYCLE_NIGHTS`; (4) collects each agent's structured JSON, dedups across agents, prioritizes; (5) routes +per D3: Slack ALARM for confirmed criticals, everything else to the mode-600 report, and fix-specs to the fixer +queue if Tier 3 is enabled; (6) a fully clean night posts nothing. + +**Escalation ladder (resolves Q3).** A confirmed critical that stays unaddressed escalates beyond a one-shot +Slack ALARM: it re-alarms on a backoff each night it persists, and after `N` nights (default 3) the coordinator +opens a tracking Jira ticket (INFRA) so it cannot quietly linger. The same ladder applies to a COMPLACENCY or +COVERAGE alarm that does not clear. Escalation stays ALARM-only in spirit (nothing posts on a clean state). + +The coordinator holds its own state; it does not rely on the orchestrator (one-shot, stateless). It may *call* +`orchestrator/run.py` for GPT/Gemini/DeepSeek single-shot sub-tasks, or call those providers directly (Q1: the +team implements its own coordination; the orchestrator is not refactored for persistence). + +## 6. Constraints and how this revision answers them + +### 6.1 Subscription billing draw — shared cap (B1 resolved by D1/D4) +Headless under Max is permitted for now (D1). Every Claude SDK call still draws from the **same Max pool as +interactive Claude Code**, so the team runs under **one shared nightly cap across all agents**, and pushes volume +to Gemini/GPT/DeepSeek (own-account billing) where quality allows. The billing-mode abstraction (D4) lets us flip +to `api`/`bedrock` when the ToS cutover lands. `ANTHROPIC_API_KEY` stays unset in subscription mode. + +### 6.2 Fixer write path (B2 resolved by D2) +Box stays read-only; CI applies the patch and opens the draft PR. The CI workflow + its OIDC role is **new IAM** +and goes through the **mandatory GPT-4.1 cross-review before it is built** (§7, B3). + +### 6.3 step-ca + Roles Anywhere (D5) +New internal CA (step-ca) issues short-lived leaf certs auto-renewed by a systemd timer; the Roles Anywhere +**trust anchor + the read-only AWS role** are **new IAM** and go through the mandatory cross-review before build +(B3). The leaf is short-lived (self-expiring), which is stronger than the box's long-lived GitHub PAT. + +### 6.4 Anti-complacency per role +Each agent ships with its own canary corpus, versioned in the repo. A role with a failing/stale canary is skipped +with a COMPLACENCY ALARM, never run silently degraded. **Canary update process (resolves F2):** canary corpora +are version-controlled **per agent/role** (no shared global corpus, to avoid cross-role confusion); a change to +any canary goes through a PR + review with a revert point, and because the canary runs every night, a canary edit +that silently weakens recall is itself caught on the next run. + +### 6.5 Host capacity +4GB / 2 vCPU / 40GB. Work is I/O-bound, but more report history + step-ca may pressure disk. Re-check headroom +after Phase 1; size up the Hyper-V VM (snapshot first per `feedback_ec2_replacement_snapshot` discipline) if +needed rather than risking the secrev workload. + +### 6.6 Claude budget realism + contention (resolves B1) +A multi-stage pipeline draws far more Claude than a single sweep (clarifier loop + planner + verifier-read per +task), all on the **same Max pool as Adam's interactive Claude Code**. Three controls: +- **Pre-build measurement (gate before P3 builds anything).** Using the existing `telemetry.py` token capture, + measure the per-stage Claude token draw on a representative task plus the worst-case clarifier loop, then + project per-task and daily aggregate at expected task volume. P1/P2 must emit these numbers before P3 + proceeds; the design measures, it does not assume. +- **One shared daily Claude cap across ALL R720 Claude work** (pipeline + Plane-1 sweeps), tracked in the + persistent budget ledger. Non-Claude stages are pushed to GPT/Gemini/DeepSeek (own-account billing) to keep + the Claude draw down. +- **Interactive-first contention rule.** Adam's interactive Claude Code is never blocked. The box keeps a + reserve headroom; before starting a stage it checks remaining headroom, and if below the reserve it **parks + new pipeline tasks and ALARMs** rather than competing for the pool. A task already mid-flight checkpoints and + pauses at the next stage boundary (never killed). A clarifier is capped at `N` turns per task, then + escalates/parks, so an ambiguous task cannot loop-drain the pool. +- **Fairness + no starvation.** The reserve is a fixed configured fraction of the daily cap (not guessed at + runtime). Parked tasks are FIFO-aged with a `MAX_PARK` window; a task that exceeds it escalates (ALARM, and a + Jira ticket per §5) rather than starving silently, and an operator can force-resume or re-prioritize it via the + CLI. A clarifier parked at the turn cap is resumable the same way: Adam adds context and re-opens it, so it is + never an indefinite stall. + +### 6.7 State durability + backup (resolves F1) +All durable state (the LangGraph SQLite checkpoint, the `pending_questions` ledger, the budget ledger, the +Plane-1 rotation/coverage pointer) is written atomically (write-temp-then-rename), integrity-checked on load, and +included in the nightly offsite backup (mode 600). On corruption the coordinator refuses to proceed silently: the +rotation pointer is rebuildable from the report history, and a corrupt task checkpoint parks that task with an +ALARM rather than restarting it blindly. "Integrity-checked" is concrete: schema-version match + a stored content +hash + a logical-consistency check (e.g. no `answered` question whose graph is already past that turn). After a +restore, a reconciliation step re-syncs against external state (in-flight CI runs, current GitHub PR status) +before any task resumes, so a restored backup cannot act on stale external assumptions. + +## 7. Phased rollout (re-sequenced for provisioning order, rollback, cross-review, and docs-as-you-go) + +- **Phase 0 — substrate factoring (with rollback, B4).** Back up `nightly_sweep.sh` (tag a revert point); + extract discovery/mirror/budget-ledger/rotation/Slack/canary into a shared module used by both secrev and the + team. **Gate:** secrev passes its canary + existing behavior after the refactor, else revert. No team behavior + yet. +- **Phase 1 — one checker end to end.** Build `compliance-drift` + its canary + the mode-600 report path + a + **routing dry-run** (F4) for the Slack alarm. Dry-run on the mirrors. Proves the substrate generalizes. + **Create the `project_r720_agent_team` memory now** (B5/docs-as-you-go). +- **Phase 2 — coordinator + second checker.** Add the coordinator (shared budget, dedup, versioned state F1) and + `dependency-cve`. **Run a forced budget-squeeze dry-run** to prove deferral-not-drop + COVERAGE ALARM (F2). +- **Phase 3 — doc-drift + step-ca/Roles Anywhere + aws-posture.** Stand up step-ca and the Roles Anywhere trust + anchor + read-only AWS role; **cross-review the IAM before building** the agent (B3). Wire aws-posture. Add + doc-drift. +- **Phase 4 — planner + confluence-doc.** `plan-groomer` writing into the report only (D3). For + `confluence-doc`: create the `confluence-bot` service account with IT-space-only edit rights (D6), put its + token in `~/secrev.env`; ship the scheduled gap-detection (read-only, recommend) first, then wire the + on-demand SSH-invoked write path. The Mermaid script (`~/.claude/scripts/confluence_mermaid.py`, already + written and offline-tested) must pass a **live dry-run against page 1540098** (verify it lists all 16 weweave + macros and that a no-op set is clean) before any `--apply`. Notion/Jira auto-write stays a later toggle. +- **Phase 5 — fixer (D2).** Build the org CI apply-and-open-draft-PR workflow; **cross-review its OIDC IAM + before building** (B3). Confirm fixer PRs hit pre-push hooks + CI + Claude Code App (F3). Start with the + narrowest fix class (dep bumps). Draft PRs only. +- **Phase 6 — document as standing infra.** Update Confluence (the team **and** the still-undocumented secrev + host) in the IT host/LAN inventory; write the operator runbook. **The runbook must cover incident handling + (resolves F4):** pipeline stalls, stuck/parked tasks, failed human-in-the-loop resumes, budget exhaustion + mid-pipeline, transport outages, and COMPLACENCY/COVERAGE alarms, each with the manual CLI recovery steps + (§3.3.1) and the escalation ladder (§5). + +**Provisioning + rollback + cross-review gate (resolves B3/B5/F5), applied to every phase:** +- **Cross-review is a hard gate, not a note.** Any new IAM role/policy, trust anchor, OIDC role, or permission + change is provisioned AND passes the mandatory GPT-4.1 cross-review (plus `/sh-security-review` where it + touches the CI / untrusted-input surface) **before** any code that depends on it is built. A phase cannot + start its dependent work until that review is recorded. +- **Every stateful phase has a revert point.** Not just Phase 0: before standing up step-ca, the Roles Anywhere + trust anchor + AWS role, or the CI apply-workflow, capture a documented rollback (remove the role/CA, restore + the prior workflow, revert the cert config) and gate the phase on a successful dry-run. The rollback is + **exercised** (sandbox or simulated teardown/re-provision), not merely written, before the phase is accepted. + VM changes snapshot first per `feedback_ec2_replacement_snapshot`. +- **Docs land with the change, enforced as definition-of-done.** A phase is not "done" until its memory entry + and the relevant Confluence page are updated; that update is a checklist item in the phase, not deferred + (Phase 6 is only the final standing-infra writeup). + +### 7.1 Pipeline track (Plane 2) — depends only on Phase 0 substrate + +This track is largely independent of the Plane-1 checker phases (1-6); both build on the Phase 0 substrate. +Given the north star, **Adam may prioritize this track first.** Sequencing within it: + +- **Phase P1 — skeleton + the human gate.** LangGraph graph + SQLite checkpointer on the box; one trivial task + type; Mac SSH-invoke intake; the **clarifier** with `interrupt()`/resume over **one** transport (Slack first). + Stops at an approved plan, no build yet. This proves durable suspend/resume across a real human answer (the + riskiest mechanic) before anything else. **Exit criteria (must demonstrate §3.3.1):** (a) kill the box + mid-wait and have the task resume after restart; (b) submit a duplicate answer and confirm it no-ops; (c) + submit an answer after the deadline expired and confirm it is rejected and the task parked; (d) two tasks + suspended concurrently resume independently to the correct thread. +- **Phase P2 — planner + review loop.** Wire the planner and the GPT-4.1 review loop (reuse `cross_reviewer`), + including loop-back and the escalate-to-Adam path. +- **Phase P3 — builders + verifier via org CI.** Builders emit a candidate diff; the Option-B OIDC workflow + builds/tests/security-reviews; the verifier reads CI results and produces a draft PR. Start with the narrowest + task class (e.g. a dependency bump or a single-file fix), draft PRs only. **Build the §3.3.2 trust boundary:** + split untrusted/privileged CI jobs, diff-hash integrity, the trust-control-surface denylist, and the pure-code + pass/fail gate. The CI apply/verify workflow + its OIDC role go through **`/sh-security-review` AND the + mandatory GPT-4.1 cross-review** before this phase ships (it is IaC/IAM + untrusted-input handling). +- **Phase P4 — more transports + GitHub intake.** Add the ticket-comment and Claude-Code responder adapters + (D10) and GitHub-issue intake (D13). +- **Phase P5 — checker findings as a task source.** Let a confirmed Plane-1 finding open a pipeline task, closing + the loop between the two planes. + +Observability for both planes (D12): LangSmith stays off; the existing local JSONL (`telemetry.py`) covers the +LangChain/LangGraph path, the Agent-SDK path keeps its own run logs, and self-hosted Phoenix is an optional +later add if per-run trace UI is wanted. + +## 8. Open items folded in (no longer blocking) +- Shared vs separate timer: **shared** with secrev (one discovery/mirror pass, one shared budget); error + isolation handled by per-role try/skip + canary, documented in the runbook (N2). +- Read-only PAT sufficiency (Q2): confirm the existing fine-grained PAT covers all mirrors before Phase 1; it + already clones every non-archived org repo for secrev, so this is a verification step, not a change. + +## 9. Obligations on build (per global instructions) +- **Memory:** `project_r720_agent_team` created in Phase 1; cross-link `project_security_review_agent`, + `project_orchestration_migration`, `reference_claude_subscription_billing`, `feedback_cloudwatch_alarms`. +- **Confluence:** document the team (and secrev) as standing infra (always-on VM holding read-only org PAT + now + a Roles Anywhere AWS identity). +- **Handbook/naming:** kebab-case dirs/resources (`agent-team`), snake_case importable Python modules; secrets in + `.env`/Secrets Manager/`~/secrev.env` (mode 600), never committed; CI/CD for the fixer apply-workflow. +- **Cross-review:** the fixer CI OIDC role and the Roles Anywhere trust anchor + AWS read role each go through the + mandatory GPT-4.1 cross-review before their phase builds. +- **Service-account lifecycle (F3):** the `confluence-bot` token is rotated on a schedule (90 days, calendared + like the GitHub PAT), has a documented revocation step, keeps its edits attributable in Confluence page + history, and is decommissioned if the agent is retired. +- **Backups (F1):** the durable state stores (checkpoint, ledgers, rotation pointer) are included in the nightly + offsite backup.