This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
orchestrator/agent-team/agent_team/responder.py
Adam Moussa 15a416d31a Add Plane-2 leaf scaffold (pipeline graph, nodes, HITL, transports, CI)
Consolidates the 18 leaf modules from the r720-plane2-scaffold workflow onto
the foundation commit. Full suite: 535 passed, 1 skipped; ruff + format clean.

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

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

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

444 lines
17 KiB
Python

"""Transport-agnostic notify + resume responder (design §3.3, §3.3.1).
This is the Plane-2 leaf that owns the human-in-the-loop **notify+resume seam**.
It sits on top of the committed FOUNDATION contracts and wires them together; it
re-declares none of them:
* :mod:`agent_team.transport.base` — the :class:`~agent_team.transport.base.Transport`
ABC plus :class:`~agent_team.transport.base.QuestionSet` /
:class:`~agent_team.transport.base.NormalizedAnswer` payloads.
* :mod:`agent_team.db.schema` — the ``pending_questions`` ledger and the
``BEGIN IMMEDIATE`` first-answer-wins compare-and-set helpers
(:func:`~agent_team.db.schema.answer_question`,
:func:`~agent_team.db.schema.expire_question`,
:func:`~agent_team.db.schema.supersede_question`).
Three seams, all transport-independent (§3.3.1):
* **Notify (and lost-post).** :func:`notify_question` writes the ledger row
``open`` FIRST, then posts to the chosen transport and stores its
``channel_ref``. If the post raises, the row stays ``open`` with no ref so the
reconcile loop can retry idempotently — the durable ledger is the source of
truth, never an in-memory-only post.
* **Answer (first-answer-wins).** :func:`submit_answer` normalizes the inbound
raw payload via the transport, then runs the single atomic compare-and-set
``UPDATE ... SET status='answered' ... WHERE question_id=? AND status='open'``.
rowcount 1 = first valid answer → enqueue a resume job; rowcount 0 = duplicate
/ late / already-closed → ignored with an "already closed" reply. This one
compare-and-set makes duplicate clicks, transport redelivery, answers via two
channels, and answer-after-timeout all safe.
* **Resume (single-flight, turn-guarded).** :class:`ResumeWorker` serializes per
``thread_id`` and, before resuming, checks the live checkpoint is still
interrupted on this ``turn``. If the graph already advanced (stale or
redelivered job) it marks the question ``superseded`` and skips, so a resume
can never double-apply. The graph is injected as a small structural protocol
so this leaf stays free of a hard LangGraph dependency for pre-deployment
scaffolding.
The deadline timer (overdue ``open`` → ``expired``) and the startup recovery
sweep round out the lifecycle and reuse the same compare-and-set helpers.
No network, no SDK, no provisioning here: the transport, the graph, and the
resume-job queue are all injected, so this module is pure orchestration over the
durable contracts and is unit-testable in isolation.
"""
from __future__ import annotations
import json
import sqlite3
import threading
from dataclasses import dataclass
from datetime import datetime, timezone
from typing import Any, Callable, Protocol, runtime_checkable
from agent_team.db.schema import (
answer_question,
expire_question,
supersede_question,
)
from agent_team.transport.base import QuestionSet, Transport
__all__ = [
"AnswerOutcome",
"GraphHandle",
"ResumeJob",
"ResumeWorker",
"deadline_sweep",
"notify_question",
"recover_open_questions",
"submit_answer",
]
# ---------------------------------------------------------------------------
# Injected collaborators (kept as structural protocols so this leaf has no hard
# LangGraph / queue dependency for pre-deployment scaffolding).
# ---------------------------------------------------------------------------
@runtime_checkable
class GraphHandle(Protocol):
"""The slice of the LangGraph graph the responder needs (§3.3.1).
The real object is a compiled LangGraph graph backed by the SQLite
checkpointer. The responder only needs to (a) ask which ``turn`` a thread is
currently interrupted on and (b) resume it with an answer, so it depends on
this narrow structural protocol rather than importing LangGraph.
"""
def interrupted_turn(self, thread_id: str) -> int | None:
"""Return the ``turn`` the thread is interrupted on, or ``None``.
``None`` means the thread is not currently suspended on an
``interrupt()`` (it already advanced, completed, or never existed). The
turn guard compares this against the question's ``turn``.
"""
...
def resume(self, thread_id: str, answer: Any) -> Any:
"""Resume the thread with ``answer`` (``Command(resume=answer)``).
Drives the graph forward from its checkpointed interrupt. Returns
whatever the graph yields next; the responder does not interpret it.
"""
...
@dataclass(frozen=True)
class ResumeJob:
"""A unit of resume work enqueued after a first-answer-wins compare-and-set.
Carries exactly what the turn-guarded :class:`ResumeWorker` needs:
``thread_id`` (the per-thread single-flight key), ``question_id`` (the
ledger row to supersede if stale), ``turn`` (the turn guard), and the
``answer`` to feed into ``Command(resume=...)``.
"""
thread_id: str
question_id: str
turn: int
answer: Any
# A resume-job enqueue callback. The real queue is durable (re-enqueued on the
# startup sweep, §3.3.1); the responder only needs to hand a job to it.
EnqueueResume = Callable[[ResumeJob], None]
@dataclass(frozen=True)
class AnswerOutcome:
"""Result of :func:`submit_answer` (§3.3.1).
``accepted`` is the rowcount-1 first-answer-wins verdict: ``True`` means this
call recorded the first valid answer and a resume job was enqueued;
``False`` means the question was not ``open`` (already answered / expired /
superseded), so the answer was a duplicate or late and was ignored.
``question_id`` / ``via`` echo the normalized inbound answer for the audit
trail. ``job`` is the enqueued :class:`ResumeJob` iff ``accepted``.
"""
accepted: bool
question_id: str
via: str
job: ResumeJob | None = None
# ---------------------------------------------------------------------------
# Notify (delivery + lost-post) — §3.3.1.
# ---------------------------------------------------------------------------
def notify_question(
conn: sqlite3.Connection,
transport: Transport,
question_set: QuestionSet,
*,
deadline: str,
posted_at: str | None = None,
) -> str | None:
"""Deliver a question-set: write the ledger row ``open`` first, then post.
Implements the §3.3.1 "delivery (and lost-post)" rule precisely:
1. Insert the ``pending_questions`` row as ``open`` with no ``channel_ref``.
The durable ledger is written BEFORE the side-effecting post so a crash
or a failed post never loses the question — recovery sees an ``open`` row
lacking a ref and retries.
2. Post the question-set over ``transport`` (which MUST embed the
``question_id`` so an inbound answer maps back) and capture the returned
``channel_ref``.
3. Persist the ``channel_ref`` (and ``posted_at``) on the row.
Returns the ``channel_ref`` on success, or ``None`` if the post failed (the
row stays ``open`` with no ref for the reconcile loop). The transport
exception is intentionally swallowed: a lost post is a recoverable state in
this design, not a hard error.
The row is inserted with the question's identity (``thread_id``, ``turn``,
``transport`` name) so the turn guard and reconcile can act on it.
"""
stamp = posted_at or _utc_now_iso()
transport_name = type(transport).__name__
# 1. Durable ledger row first (open, no ref).
conn.execute(
"INSERT INTO pending_questions "
"(question_id, thread_id, turn, status, transport, posted_at, deadline_at) "
"VALUES (?, ?, ?, 'open', ?, ?, ?)",
(
question_set.question_id,
question_set.thread_id,
question_set.turn,
transport_name,
stamp,
deadline,
),
)
# 2. Side-effecting post. A failure here is recoverable (row stays open,
# no ref) — do NOT let it bubble up and lose the durable row.
try:
channel_ref = transport.post_question(
thread_id=question_set.thread_id,
question_id=question_set.question_id,
turn=question_set.turn,
question_set=question_set,
deadline=deadline,
)
except Exception:
return None
# 3. Persist the ref so reconcile/recovery can act on the post.
conn.execute(
"UPDATE pending_questions SET channel_ref=? WHERE question_id=?",
(channel_ref, question_set.question_id),
)
return channel_ref
# ---------------------------------------------------------------------------
# Answer (first-answer-wins) — §3.3.1.
# ---------------------------------------------------------------------------
def submit_answer(
conn: sqlite3.Connection,
transport: Transport,
raw: Any,
*,
enqueue_resume: EnqueueResume,
answered_at: str | None = None,
) -> AnswerOutcome:
"""Normalize an inbound answer and run the first-answer-wins compare-and-set.
The transport adapter normalizes ``raw`` to ``(question_id, answer, via)``;
the responder then runs the single atomic statement (via
:func:`agent_team.db.schema.answer_question`, which wraps it in
``BEGIN IMMEDIATE``):
UPDATE pending_questions SET status='answered', answer_json=?,
answered_via=?, answered_at=? WHERE question_id=? AND status='open'
* rowcount 1 → this is the first valid answer: look up the row's ``turn``
and enqueue a :class:`ResumeJob`; return ``accepted=True``.
* rowcount 0 → the question was not ``open`` (already answered / expired /
superseded): the answer is a duplicate or late and is ignored; return
``accepted=False`` so the caller can send an "already closed" reply.
This single compare-and-set is what makes duplicate clicks, transport
redelivery, answers-via-two-channels, and answer-after-timeout all safe —
only one caller can ever flip ``open`` → ``answered``.
"""
question_id, answer, via = transport.parse_answer(raw)
accepted = answer_question(
conn,
question_id=question_id,
answer_json=json.dumps(answer, sort_keys=True),
answered_via=via,
answered_at=answered_at,
)
if not accepted:
# Duplicate / late / already-closed: ignore (caller replies "closed").
return AnswerOutcome(accepted=False, question_id=question_id, via=via)
# First valid answer: enqueue the turn-guarded resume job.
turn = _question_turn(conn, question_id)
job = ResumeJob(
thread_id=_question_thread(conn, question_id),
question_id=question_id,
turn=turn,
answer=answer,
)
enqueue_resume(job)
return AnswerOutcome(accepted=True, question_id=question_id, via=via, job=job)
# ---------------------------------------------------------------------------
# Resume (single-flight, turn-guarded) — §3.3.1.
# ---------------------------------------------------------------------------
class ResumeWorker:
"""Serializes resume work per ``thread_id`` and guards on the turn (§3.3.1).
A resume job for a thread runs under a per-thread lock so two jobs for the
same thread can never resume concurrently (single-flight). Before resuming,
the worker checks the live checkpoint is still interrupted on the job's
``turn``:
* if the graph advanced past this turn (stale or redelivered job) the
question is marked ``superseded`` and the resume is skipped — a resume can
never double-apply;
* otherwise it calls ``graph.resume(thread_id, answer)`` (the injected
``Command(resume=answer)`` seam).
Different threads resume concurrently within the budget cap; only same-thread
work is serialized. The per-thread locks are created lazily under a single
registry lock so this is safe to share across resume threads.
"""
def __init__(self, conn: sqlite3.Connection, graph: GraphHandle) -> None:
self._conn = conn
self._graph = graph
self._registry_lock = threading.Lock()
self._thread_locks: dict[str, threading.Lock] = {}
def _lock_for(self, thread_id: str) -> threading.Lock:
"""Return (creating if needed) the single-flight lock for ``thread_id``."""
with self._registry_lock:
lock = self._thread_locks.get(thread_id)
if lock is None:
lock = threading.Lock()
self._thread_locks[thread_id] = lock
return lock
def run(self, job: ResumeJob) -> bool:
"""Process one resume ``job`` under the per-thread single-flight lock.
Returns ``True`` if the graph was resumed, ``False`` if the job was a
no-op (graph already advanced past the turn → question superseded and
skipped). Idempotent: re-running a job for an already-advanced thread
supersedes-and-skips rather than double-applying.
"""
with self._lock_for(job.thread_id):
live_turn = self._graph.interrupted_turn(job.thread_id)
if live_turn != job.turn:
# Stale / redelivered: the graph already advanced past this turn
# (or is not interrupted at all). Supersede and skip — never
# double-apply a resume.
supersede_question(self._conn, question_id=job.question_id)
return False
self._graph.resume(job.thread_id, job.answer)
return True
# ---------------------------------------------------------------------------
# Deadline / no-answer timer — §3.3.1.
# ---------------------------------------------------------------------------
def deadline_sweep(
conn: sqlite3.Connection,
*,
now: str | None = None,
) -> list[str]:
"""Flip overdue ``open`` questions to ``expired`` (deterministic race).
Selects ``open`` rows whose ``deadline_at`` is non-null and ``<= now`` and
runs the same compare-and-set (:func:`agent_team.db.schema.expire_question`)
on each. Expiry vs answer is a deterministic race on flipping ``open``: an
answer arriving for an already-``expired`` question loses the compare-and-set
and is ignored, and vice versa. Returns the ids actually expired by this
call (rowcount 1), so the caller can apply the task policy (park + ALARM, or
a defined default answer) to exactly those.
"""
moment = now or _utc_now_iso()
rows = conn.execute(
"SELECT question_id FROM pending_questions "
"WHERE status='open' AND deadline_at IS NOT NULL AND deadline_at <= ?",
(moment,),
).fetchall()
expired: list[str] = []
for row in rows:
qid = row["question_id"]
if expire_question(conn, question_id=qid):
expired.append(qid)
return expired
# ---------------------------------------------------------------------------
# Restart recovery — §3.3.1.
# ---------------------------------------------------------------------------
def recover_open_questions(
conn: sqlite3.Connection,
*,
enqueue_resume: EnqueueResume,
) -> list[ResumeJob]:
"""Startup sweep: re-enqueue resume jobs for already-``answered`` questions.
All state is durable, so a reboot converges by replaying the ledger. This
routine handles the ``answered`` slice of that sweep (§3.3.1): for every row
still ``answered`` (i.e. answered before a crash but not yet resumed), it
re-enqueues a :class:`ResumeJob`. The re-enqueue is safe because the
:class:`ResumeWorker` turn guard makes resumes idempotent — a job whose graph
already advanced supersedes-and-skips.
Delivery-retry for ``open`` rows lacking a ref and the deadline policy for
overdue ``open`` rows are the reconcile loop's and
:func:`deadline_sweep`'s jobs respectively; this function owns only the
answered→resume replay. Returns the jobs it enqueued.
"""
rows = conn.execute(
"SELECT question_id, thread_id, turn, answer_json "
"FROM pending_questions WHERE status='answered'"
).fetchall()
jobs: list[ResumeJob] = []
for row in rows:
answer = json.loads(row["answer_json"]) if row["answer_json"] else None
job = ResumeJob(
thread_id=row["thread_id"],
question_id=row["question_id"],
turn=int(row["turn"]),
answer=answer,
)
enqueue_resume(job)
jobs.append(job)
return jobs
# ---------------------------------------------------------------------------
# Internal helpers.
# ---------------------------------------------------------------------------
def _question_turn(conn: sqlite3.Connection, question_id: str) -> int:
"""Return the ``turn`` recorded for ``question_id``."""
row = conn.execute(
"SELECT turn FROM pending_questions WHERE question_id=?",
(question_id,),
).fetchone()
if row is None:
raise KeyError(f"unknown question_id: {question_id!r}")
return int(row["turn"])
def _question_thread(conn: sqlite3.Connection, question_id: str) -> str:
"""Return the ``thread_id`` recorded for ``question_id``."""
row = conn.execute(
"SELECT thread_id FROM pending_questions WHERE question_id=?",
(question_id,),
).fetchone()
if row is None:
raise KeyError(f"unknown question_id: {question_id!r}")
return str(row["thread_id"])
def _utc_now_iso() -> str:
"""Return the current UTC time as an ISO-8601 string."""
return datetime.now(timezone.utc).isoformat()