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

527 lines
21 KiB
Python

"""Restart-recovery sweep + post-restore reconciliation (design §3.3.1, §6.7).
All durable human-in-the-loop state lives in two SQLite stores (the LangGraph
``SqliteSaver`` checkpoint and the ``pending_questions`` ledger), so a reboot or
a backup restore must *converge* rather than restart: nothing in the
interaction lifecycle is kept in memory. This module is the startup sweep that
drives that convergence.
Per §3.3.1 ("Restart recovery"), on startup the sweep:
* **redelivers** every ``open`` question whose row has no ``channel_ref`` (its
original post was lost / never landed): it re-posts over the question's
transport idempotently and records the returned ``channel_ref``;
* **re-enqueues a resume** for every ``answered`` question whose graph is still
interrupted on that turn (a crash between the compare-and-set and the resume
worker). The turn guard makes this idempotent — a graph that already advanced
is left alone, the question marked ``superseded``;
* **applies the deadline policy** to every overdue ``open`` question, flipping
it ``expired`` via the same first-answer-wins compare-and-set (an answer that
arrives for an already-expired question loses the race), then parks + ALARMs
or applies a default answer per the task policy.
Per §6.7 ("post-restore reconciliation"), when the sweep runs after a *restore*
(not just a reboot) it first re-syncs against external state (in-flight CI runs,
current GitHub PR status) via an injected reconciler, so a restored backup
cannot resume a task on stale external assumptions.
The module owns no I/O of its own beyond reading the ledger through the
foundation's :func:`agent_team.db.schema.connect` connection: every
side-effecting collaborator (transport delivery, resume enqueue, the live
checkpoint's interrupt/turn probe, the deadline policy, the external
reconciler) is injected as a callable. That keeps the sweep deterministic and
unit-testable, and means it provisions nothing.
It imports the committed foundation contracts verbatim
(:mod:`agent_team.db.schema`, :mod:`agent_team.transport.base`) and does not
redefine them.
"""
from __future__ import annotations
import sqlite3
from collections.abc import Callable
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Protocol
from agent_team.db.schema import (
QUESTION_STATES,
connect,
expire_question,
supersede_question,
)
from agent_team.transport.base import QuestionSet, Transport
__all__ = [
"PendingQuestion",
"RecoveryReport",
"DeadlineOutcome",
"TransportResolver",
"ResumeEnqueuer",
"GraphInterruptProbe",
"DeadlinePolicy",
"ExternalReconciler",
"load_pending_questions",
"redeliver_open_questions",
"reenqueue_answered_resumes",
"apply_deadline_policy",
"run_restart_recovery",
]
# --------------------------------------------------------------------------- #
# Row view #
# --------------------------------------------------------------------------- #
@dataclass(frozen=True)
class PendingQuestion:
"""A read-only view of one ``pending_questions`` ledger row (§3.3.1).
Mirrors the foundation DDL columns. The sweep reads rows through this view
rather than passing raw :class:`sqlite3.Row` objects around, so the
collaborators receive a typed, immutable record.
"""
question_id: str
thread_id: str
turn: int
status: str
transport: str
channel_ref: str | None
posted_at: str | None
deadline_at: str | None
answer_json: str | None
answered_at: str | None
answered_via: str | None
@classmethod
def from_row(cls, row: sqlite3.Row) -> PendingQuestion:
"""Build a :class:`PendingQuestion` from a ``pending_questions`` row."""
return cls(
question_id=row["question_id"],
thread_id=row["thread_id"],
turn=int(row["turn"]),
status=row["status"],
transport=row["transport"],
channel_ref=row["channel_ref"],
posted_at=row["posted_at"],
deadline_at=row["deadline_at"],
answer_json=row["answer_json"],
answered_at=row["answered_at"],
answered_via=row["answered_via"],
)
# --------------------------------------------------------------------------- #
# Injected-collaborator contracts #
# --------------------------------------------------------------------------- #
# A resolver hands the sweep the right Transport adapter for a question's
# ``transport`` string (e.g. "slack"/"github"/"claude_code"). Returning None
# means the transport is currently unreachable/unconfigured; the sweep records
# the redelivery as deferred rather than crashing.
TransportResolver = Callable[[str], Transport | None]
# Enqueue a resume job for (thread_id, question_id, turn). The resume worker
# (§3.3.1) is single-flight per thread and turn-guarded; the sweep only needs
# to (idempotently) put the job on the queue. Returns True if a job was
# enqueued.
ResumeEnqueuer = Callable[[str, str, int], bool]
# Probe the LIVE LangGraph checkpoint: is ``thread_id`` still interrupted on
# exactly ``turn``? Returns True only when the graph is genuinely still waiting
# on this question's turn. A False return means the graph already advanced
# (stale/redelivered) and the question must be superseded, never resumed.
GraphInterruptProbe = Callable[[str, int], bool]
# The deadline policy for an expired question (§3.3.1): park + ALARM, or apply
# a defined default answer. Invoked only after the row is durably flipped to
# ``expired``. Returns the action it took for the report.
DeadlinePolicy = Callable[["PendingQuestion"], "DeadlineOutcome"]
class ExternalReconciler(Protocol):
"""Post-restore external-state reconciliation seam (§6.7).
After a *restore* (not a plain reboot) the sweep must re-sync against
external systems (in-flight CI runs, current GitHub PR status) before any
task resumes, so a restored backup cannot act on stale assumptions. A
concrete reconciler is injected in the leaves; the sweep only invokes it.
"""
def reconcile(self, question: PendingQuestion) -> bool:
"""Reconcile one task's external state.
Return ``True`` when the task is safe to resume, ``False`` when external
state diverged (e.g. the CI run vanished or the PR was closed) and the
task must be held/parked instead of resumed.
"""
...
# --------------------------------------------------------------------------- #
# Report types #
# --------------------------------------------------------------------------- #
@dataclass(frozen=True)
class DeadlineOutcome:
"""What the deadline policy did with one expired question (§3.3.1)."""
question_id: str
action: str # "parked" | "defaulted" | str describing the action taken
detail: str = ""
@dataclass
class RecoveryReport:
"""Structured result of a restart-recovery sweep (§3.3.1, §6.7).
Every list holds ``question_id`` values so the coordinator can ALARM /
report deterministically. ``errors`` carries ``(question_id, message)``
pairs for collaborator failures that were isolated so one bad row cannot
abort the whole sweep.
"""
redelivered: list[str] = field(default_factory=list)
redelivery_deferred: list[str] = field(default_factory=list)
resumes_enqueued: list[str] = field(default_factory=list)
superseded: list[str] = field(default_factory=list)
expired: list[str] = field(default_factory=list)
deadline_outcomes: list[DeadlineOutcome] = field(default_factory=list)
reconcile_held: list[str] = field(default_factory=list)
errors: list[tuple[str, str]] = field(default_factory=list)
@property
def clean(self) -> bool:
"""True when the sweep took no action and hit no errors.
A clean sweep means durable state already matched reality (nothing to
redeliver, resume, expire, or hold) — the §3.3 "clean night posts
nothing" discipline applies to recovery too.
"""
return not any(
(
self.redelivered,
self.redelivery_deferred,
self.resumes_enqueued,
self.superseded,
self.expired,
self.reconcile_held,
self.errors,
)
)
# --------------------------------------------------------------------------- #
# Ledger reads #
# --------------------------------------------------------------------------- #
def load_pending_questions(
conn: sqlite3.Connection,
*,
status: str | None = None,
) -> list[PendingQuestion]:
"""Load ``pending_questions`` rows, optionally filtered by ``status``.
Returned newest-posted-first within a stable secondary key so the sweep is
deterministic. ``status`` must be one of
:data:`agent_team.db.schema.QUESTION_STATES` when given.
"""
if status is not None and status not in QUESTION_STATES:
raise ValueError(
f"unknown status {status!r}; expected one of {QUESTION_STATES}"
)
sql = "SELECT * FROM pending_questions"
params: tuple[Any, ...] = ()
if status is not None:
sql += " WHERE status = ?"
params = (status,)
# Order by question_id as a stable tiebreaker; posted_at may be NULL for a
# never-delivered open row, so it cannot be the sole sort key.
sql += " ORDER BY posted_at IS NULL, posted_at, question_id"
rows = conn.execute(sql, params).fetchall()
return [PendingQuestion.from_row(row) for row in rows]
def _record_channel_ref(
conn: sqlite3.Connection, *, question_id: str, channel_ref: str
) -> None:
"""Persist a freshly obtained ``channel_ref`` for an ``open`` question.
Guarded on ``status='open'`` so a concurrent answer/expire that closed the
row in the meantime is not clobbered — the redelivery simply no-ops on the
ledger if the question is no longer open.
"""
conn.execute("BEGIN IMMEDIATE")
try:
conn.execute(
"UPDATE pending_questions SET channel_ref = ? "
"WHERE question_id = ? AND status = 'open'",
(channel_ref, question_id),
)
conn.execute("COMMIT")
except BaseException:
conn.execute("ROLLBACK")
raise
# --------------------------------------------------------------------------- #
# Sweep step 1 — redeliver lost posts #
# --------------------------------------------------------------------------- #
def redeliver_open_questions(
conn: sqlite3.Connection,
*,
resolve_transport: TransportResolver,
report: RecoveryReport,
) -> None:
"""Re-post every ``open`` question lacking a ``channel_ref`` (§3.3.1).
On interrupt the responder writes the row ``open`` *before* posting, so a
crash (or a failed post) can leave an ``open`` row with no ref. This step
retries delivery idempotently: it resolves the question's transport,
re-posts the question-set, and records the returned ``channel_ref``. If the
transport is unreachable the redelivery is *deferred* (not an error) so a
later sweep / reconcile loop retries — matching the §3.3.1 "reconcile loop
retries idempotently" behaviour.
Rows that already have a ``channel_ref`` are skipped: their post landed.
"""
for question in load_pending_questions(conn, status="open"):
if question.channel_ref:
continue # post already landed; nothing to redeliver.
transport = resolve_transport(question.transport)
if transport is None:
report.redelivery_deferred.append(question.question_id)
continue
try:
channel_ref = transport.post_question(
thread_id=question.thread_id,
question_id=question.question_id,
turn=question.turn,
question_set=_question_set_for(question),
deadline=question.deadline_at or "",
)
except Exception as exc: # isolate one bad transport call.
report.errors.append((question.question_id, f"redeliver: {exc}"))
report.redelivery_deferred.append(question.question_id)
continue
if not channel_ref:
# Transport returned no locator; treat as a deferred retry.
report.redelivery_deferred.append(question.question_id)
continue
_record_channel_ref(
conn, question_id=question.question_id, channel_ref=channel_ref
)
report.redelivered.append(question.question_id)
def _question_set_for(question: PendingQuestion) -> QuestionSet:
"""Build a minimal :class:`QuestionSet` for a redelivery.
The original prompt text lives in the LangGraph checkpoint, not this
ledger; on a redelivery the sweep carries identity (``thread_id`` /
``question_id`` / ``turn``) so the adapter can re-render from the
checkpoint. ``questions`` is left empty here and the adapter fills it from
graph state, keeping the ledger free of duplicated prompt text.
"""
return QuestionSet(
thread_id=question.thread_id,
question_id=question.question_id,
turn=question.turn,
questions=[],
)
# --------------------------------------------------------------------------- #
# Sweep step 2 — re-enqueue resumes for already-answered questions #
# --------------------------------------------------------------------------- #
def reenqueue_answered_resumes(
conn: sqlite3.Connection,
*,
is_interrupted_on_turn: GraphInterruptProbe,
enqueue_resume: ResumeEnqueuer,
report: RecoveryReport,
reconciler: ExternalReconciler | None = None,
) -> None:
"""Re-enqueue resumes for ``answered`` rows whose graph still waits (§3.3.1).
A crash can land between the first-answer-wins compare-and-set (row flipped
``answered``) and the resume worker actually resuming the graph. On startup
every ``answered`` question is checked against the *live* checkpoint:
* still interrupted on this exact turn -> (optionally reconcile external
state per §6.7, then) re-enqueue a resume. The resume worker is itself
turn-guarded and single-flight, so re-enqueuing is idempotent — a
duplicate job no-ops.
* the graph already advanced past this turn -> this is a stale/redelivered
answer; mark the question ``superseded`` (foundation compare-and-set) and
do NOT resume. "A resume can never double-apply."
When a ``reconciler`` is supplied (a restore, not a plain reboot) and it
reports external state diverged, the resume is held rather than enqueued so
the task does not act on stale CI/PR assumptions (§6.7).
"""
for question in load_pending_questions(conn, status="answered"):
try:
still_waiting = is_interrupted_on_turn(question.thread_id, question.turn)
except Exception as exc: # isolate a bad probe.
report.errors.append((question.question_id, f"probe: {exc}"))
continue
if not still_waiting:
# Graph already advanced: the resume already applied (or the turn
# moved on). Supersede so the row can never re-trigger a resume.
if supersede_question(conn, question_id=question.question_id):
report.superseded.append(question.question_id)
continue
if reconciler is not None:
try:
safe = reconciler.reconcile(question)
except Exception as exc: # isolate a bad reconciler.
report.errors.append((question.question_id, f"reconcile: {exc}"))
report.reconcile_held.append(question.question_id)
continue
if not safe:
report.reconcile_held.append(question.question_id)
continue
try:
enqueued = enqueue_resume(
question.thread_id, question.question_id, question.turn
)
except Exception as exc: # isolate a bad enqueue.
report.errors.append((question.question_id, f"resume: {exc}"))
continue
if enqueued:
report.resumes_enqueued.append(question.question_id)
# --------------------------------------------------------------------------- #
# Sweep step 3 — deadline policy for overdue open questions #
# --------------------------------------------------------------------------- #
def apply_deadline_policy(
conn: sqlite3.Connection,
*,
policy: DeadlinePolicy,
report: RecoveryReport,
now: datetime | None = None,
) -> None:
"""Expire overdue ``open`` questions and apply the task policy (§3.3.1).
For each ``open`` question whose ``deadline_at`` is at/before ``now``, the
sweep flips it ``expired`` via :func:`agent_team.db.schema.expire_question`
— the same first-answer-wins compare-and-set the live timer uses, so an
answer racing the same deadline either wins (row already ``answered``, this
call no-ops) or loses (row flipped ``expired``, a late answer is later
ignored). Only after a row is *durably* flipped does the ``policy`` run
(park + ALARM, or apply a default answer), so a crash between the two leaves
an ``expired`` row a later sweep re-processes — the policy must be
idempotent.
Rows with no ``deadline_at`` never expire here (no deadline configured).
``now`` defaults to the current UTC time; it is injectable for tests.
"""
current = now or datetime.now(timezone.utc)
for question in load_pending_questions(conn, status="open"):
if not _is_overdue(question.deadline_at, current):
continue
try:
flipped = expire_question(conn, question_id=question.question_id)
except Exception as exc: # isolate a bad compare-and-set.
report.errors.append((question.question_id, f"expire: {exc}"))
continue
if not flipped:
# Lost the race: the question was answered/closed concurrently.
continue
report.expired.append(question.question_id)
try:
outcome = policy(question)
except Exception as exc: # isolate a bad policy callback.
report.errors.append((question.question_id, f"policy: {exc}"))
continue
report.deadline_outcomes.append(outcome)
def _is_overdue(deadline_at: str | None, now: datetime) -> bool:
"""Return True when ``deadline_at`` (ISO-8601) is at/before ``now``.
A missing or unparseable deadline is treated as "not overdue": the sweep
never expires a question whose deadline it cannot read, it leaves it ``open``
for an operator. ``now`` is timezone-aware (UTC); a naive ``deadline_at`` is
assumed UTC for comparison.
"""
if not deadline_at:
return False
try:
parsed = datetime.fromisoformat(deadline_at)
except ValueError:
return False
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed <= now
# --------------------------------------------------------------------------- #
# Orchestrating sweep #
# --------------------------------------------------------------------------- #
def run_restart_recovery(
db_path: Any,
*,
resolve_transport: TransportResolver,
is_interrupted_on_turn: GraphInterruptProbe,
enqueue_resume: ResumeEnqueuer,
deadline_policy: DeadlinePolicy,
reconciler: ExternalReconciler | None = None,
now: datetime | None = None,
conn: sqlite3.Connection | None = None,
) -> RecoveryReport:
"""Run the full restart-recovery sweep against the ledger (§3.3.1, §6.7).
Convergence order matters and is fixed:
1. **deadline first** — expire overdue ``open`` rows before anything else so
a question past its deadline is never redelivered or resumed as if live;
2. **redeliver** ``open`` rows still lacking a ``channel_ref`` (their post
was lost);
3. **re-enqueue resumes** for ``answered`` rows whose graph still waits,
superseding those whose graph advanced.
When ``reconciler`` is supplied the sweep is treated as a *post-restore*
reconciliation (§6.7): step 3 first re-syncs each task's external state and
holds (does not resume) any task whose external state diverged.
A connection is opened via :func:`agent_team.db.schema.connect` unless one
is injected (tests share an in-memory/temp DB). Every per-row collaborator
failure is isolated into ``report.errors`` so one bad row cannot abort the
sweep — the coordinator decides whether the error set warrants an ALARM.
"""
owns_conn = conn is None
connection = conn if conn is not None else connect(db_path)
report = RecoveryReport()
try:
apply_deadline_policy(
connection, policy=deadline_policy, report=report, now=now
)
redeliver_open_questions(
connection, resolve_transport=resolve_transport, report=report
)
reenqueue_answered_resumes(
connection,
is_interrupted_on_turn=is_interrupted_on_turn,
enqueue_resume=enqueue_resume,
report=report,
reconciler=reconciler,
)
finally:
if owns_conn:
connection.close()
return report