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.
354 lines
13 KiB
Python
354 lines
13 KiB
Python
"""Deadline / no-answer timer loop for pending questions (design §3.3.1).
|
|
|
|
The durable human-in-the-loop ledger (``pending_questions``) gives every
|
|
delivered question-set a ``deadline_at``. This module is the **timer loop**
|
|
that the design calls out:
|
|
|
|
> 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``.
|
|
|
|
Design contract this leaf honours:
|
|
|
|
* **Same compare-and-set.** The ``open`` -> ``expired`` flip is the foundation's
|
|
:func:`agent_team.db.schema.expire_question`, imported **verbatim** and not
|
|
redefined here. It runs under ``BEGIN IMMEDIATE`` so the timer and a racing
|
|
responder are serialized: exactly one of "expire" / "answer" wins the row.
|
|
* **Deterministic race.** If :func:`expire_question` returns ``False`` for an
|
|
overdue row, a responder answered it first (or another timer pass already
|
|
expired it); the loop then does **nothing** for that row — it never applies
|
|
the no-answer policy to a question that was actually answered.
|
|
* **Per-question policy.** When the flip wins, the loop applies that question's
|
|
:class:`DeadlinePolicy`: ``PARK`` (park the task + ALARM Adam) or
|
|
``DEFAULT_ANSWER`` (resume the graph with a configured default). Policy is
|
|
resolved per question via an injected resolver so the durable ledger schema
|
|
is not extended by this leaf.
|
|
* **Restart-safe.** All inputs are read from the durable table on each pass; the
|
|
loop holds no in-memory-only state, so a reboot mid-sweep simply re-runs the
|
|
remaining overdue rows on the next pass (each flip is idempotent via the
|
|
compare-and-set).
|
|
|
|
Side effects (ALARM, resume-with-default) are injected as callables. This keeps
|
|
the module pure of transport / graph I/O — exactly what makes the deterministic
|
|
race and the policy branch unit-testable without a live Slack or LangGraph.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import sqlite3
|
|
from collections.abc import Callable
|
|
from dataclasses import dataclass, field
|
|
from datetime import datetime, timezone
|
|
from enum import Enum
|
|
|
|
# Compare-and-set primitive imported VERBATIM from the committed foundation.
|
|
# The SQL / transaction discipline is NOT redefined here.
|
|
from agent_team.db.schema import expire_question
|
|
|
|
__all__ = [
|
|
"DeadlinePolicy",
|
|
"ExpiryAction",
|
|
"ExpiryOutcome",
|
|
"OverdueQuestion",
|
|
"TimerLoopReport",
|
|
"AlarmFn",
|
|
"PolicyResolver",
|
|
"ResumeWithDefaultFn",
|
|
"overdue_open_questions",
|
|
"run_deadline_timer",
|
|
]
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class DeadlinePolicy(Enum):
|
|
"""What to do when a question hits its deadline with no answer (§3.3.1).
|
|
|
|
``PARK`` — park the task and ALARM Adam (the conservative default; a stuck
|
|
task parks rather than spins). ``DEFAULT_ANSWER`` — resume the graph with a
|
|
pre-configured default answer for tasks where a no-answer has a safe,
|
|
defined fallback.
|
|
"""
|
|
|
|
PARK = "park"
|
|
DEFAULT_ANSWER = "default_answer"
|
|
|
|
|
|
class ExpiryAction(Enum):
|
|
"""The action the loop actually took for one overdue row.
|
|
|
|
``PARKED`` / ``DEFAULTED`` follow a *won* expiry flip and the question's
|
|
policy. ``LOST_RACE`` means the compare-and-set returned ``False`` — a
|
|
responder answered (or a prior pass expired) the row first, so no no-answer
|
|
policy was applied. ``ERRORED`` means the flip won but the side effect
|
|
raised; the row is already ``expired`` and the failure is recorded for the
|
|
caller to ALARM/retry.
|
|
"""
|
|
|
|
PARKED = "parked"
|
|
DEFAULTED = "defaulted"
|
|
LOST_RACE = "lost_race"
|
|
ERRORED = "errored"
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class OverdueQuestion:
|
|
"""A minimal read-snapshot of one overdue ``open`` question row (§3.3.1).
|
|
|
|
Only the columns the timer loop needs: identity, ownership (``thread_id`` /
|
|
``turn`` for the turn-guarded resume), transport, and the deadline. Frozen
|
|
because rows are snapshots — mutation goes through the compare-and-set, never
|
|
by editing an instance.
|
|
"""
|
|
|
|
question_id: str
|
|
thread_id: str
|
|
turn: int
|
|
transport: str
|
|
deadline_at: str | None = None
|
|
|
|
@classmethod
|
|
def from_row(cls, row: sqlite3.Row) -> OverdueQuestion:
|
|
"""Build an :class:`OverdueQuestion` from a ``sqlite3.Row``."""
|
|
return cls(
|
|
question_id=row["question_id"],
|
|
thread_id=row["thread_id"],
|
|
turn=int(row["turn"]),
|
|
transport=row["transport"],
|
|
deadline_at=row["deadline_at"],
|
|
)
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class ExpiryOutcome:
|
|
"""The result of processing one overdue question (§3.3.1)."""
|
|
|
|
question_id: str
|
|
thread_id: str
|
|
action: ExpiryAction
|
|
policy: DeadlinePolicy | None = None
|
|
error: str | None = None
|
|
|
|
|
|
@dataclass
|
|
class TimerLoopReport:
|
|
"""Aggregate result of one timer-loop pass (§3.3.1).
|
|
|
|
``outcomes`` is one entry per overdue row examined. The summary counters let
|
|
the coordinator decide whether to ALARM (any ``errored``) without re-walking
|
|
the list.
|
|
"""
|
|
|
|
outcomes: list[ExpiryOutcome] = field(default_factory=list)
|
|
|
|
@property
|
|
def examined(self) -> int:
|
|
"""Number of overdue rows examined this pass."""
|
|
return len(self.outcomes)
|
|
|
|
@property
|
|
def parked(self) -> int:
|
|
"""Rows expired and parked (PARK policy)."""
|
|
return sum(1 for o in self.outcomes if o.action is ExpiryAction.PARKED)
|
|
|
|
@property
|
|
def defaulted(self) -> int:
|
|
"""Rows expired and resumed with a default answer (DEFAULT_ANSWER)."""
|
|
return sum(1 for o in self.outcomes if o.action is ExpiryAction.DEFAULTED)
|
|
|
|
@property
|
|
def lost_race(self) -> int:
|
|
"""Rows that were answered/expired by someone else first."""
|
|
return sum(1 for o in self.outcomes if o.action is ExpiryAction.LOST_RACE)
|
|
|
|
@property
|
|
def errored(self) -> int:
|
|
"""Rows whose flip won but whose side effect raised."""
|
|
return sum(1 for o in self.outcomes if o.action is ExpiryAction.ERRORED)
|
|
|
|
@property
|
|
def expired(self) -> int:
|
|
"""Rows this pass actually flipped ``open`` -> ``expired``.
|
|
|
|
Every parked / defaulted / errored row won its compare-and-set, so the
|
|
row is ``expired`` in the durable table; a ``lost_race`` row did not.
|
|
"""
|
|
return self.examined - self.lost_race
|
|
|
|
|
|
# Injected side-effect / policy seams. Keeping these as callables means the
|
|
# timer loop performs no transport or LangGraph I/O of its own (testable, and
|
|
# faithful to §3.3.1's transport-agnostic core).
|
|
AlarmFn = Callable[[OverdueQuestion], None]
|
|
"""Called once per *expired-and-parked* question to ALARM Adam."""
|
|
|
|
ResumeWithDefaultFn = Callable[[OverdueQuestion], None]
|
|
"""Called once per *expired* question whose policy is DEFAULT_ANSWER, to resume
|
|
the graph with that task's configured default answer."""
|
|
|
|
PolicyResolver = Callable[[OverdueQuestion], DeadlinePolicy]
|
|
"""Resolve the :class:`DeadlinePolicy` for one question. Defaults to PARK."""
|
|
|
|
|
|
def _utc_now_iso() -> str:
|
|
"""Return the current UTC time as an ISO-8601 string."""
|
|
return datetime.now(timezone.utc).isoformat()
|
|
|
|
|
|
def _default_policy_resolver(_question: OverdueQuestion) -> DeadlinePolicy:
|
|
"""Conservative default: park + ALARM on any no-answer (§3.3.1)."""
|
|
return DeadlinePolicy.PARK
|
|
|
|
|
|
def overdue_open_questions(
|
|
conn: sqlite3.Connection,
|
|
*,
|
|
now: str | None = None,
|
|
) -> list[OverdueQuestion]:
|
|
"""Return ``open`` rows whose ``deadline_at`` has passed (§3.3.1).
|
|
|
|
These are the candidates the timer loop flips with the foundation's
|
|
:func:`expire_question`. Rows with a NULL ``deadline_at`` never expire and
|
|
are excluded. ``now`` defaults to the current UTC ISO-8601 time; deadlines
|
|
are compared as ISO-8601 strings, which sort lexicographically iff written
|
|
in a consistent UTC offset — the ledger always writes UTC, so this holds.
|
|
Ordered oldest-deadline-first so the most-overdue questions are handled
|
|
first.
|
|
"""
|
|
cutoff = now or _utc_now_iso()
|
|
rows = conn.execute(
|
|
"SELECT question_id, thread_id, turn, transport, deadline_at "
|
|
"FROM pending_questions "
|
|
"WHERE status='open' AND deadline_at IS NOT NULL AND deadline_at <= ? "
|
|
"ORDER BY deadline_at ASC, question_id ASC",
|
|
(cutoff,),
|
|
).fetchall()
|
|
return [OverdueQuestion.from_row(r) for r in rows]
|
|
|
|
|
|
def run_deadline_timer(
|
|
conn: sqlite3.Connection,
|
|
*,
|
|
on_park: AlarmFn,
|
|
resume_with_default: ResumeWithDefaultFn | None = None,
|
|
policy_resolver: PolicyResolver | None = None,
|
|
now: str | None = None,
|
|
) -> TimerLoopReport:
|
|
"""Run one deadline-timer pass over the ledger (§3.3.1).
|
|
|
|
For every ``open`` question past its ``deadline_at`` as of ``now``:
|
|
|
|
1. Attempt the atomic ``open`` -> ``expired`` flip via the foundation's
|
|
:func:`expire_question` (the same ``BEGIN IMMEDIATE`` compare-and-set
|
|
used by the responder). This is the **deterministic race**: if a
|
|
responder answered the question first, the flip returns ``False`` and the
|
|
loop records :attr:`ExpiryAction.LOST_RACE` and applies **no** policy.
|
|
2. On a winning flip, resolve the question's :class:`DeadlinePolicy` and act:
|
|
* :attr:`DeadlinePolicy.PARK` -> call ``on_park`` (park the task + ALARM
|
|
Adam).
|
|
* :attr:`DeadlinePolicy.DEFAULT_ANSWER` -> call ``resume_with_default``
|
|
(resume the graph with the configured default). If a question resolves
|
|
to ``DEFAULT_ANSWER`` but no ``resume_with_default`` callback was
|
|
supplied, that is a configuration error and raises :class:`ValueError`
|
|
(the row is already ``expired``; failing loud beats silently dropping a
|
|
resume).
|
|
|
|
Side effects are isolated per row: if a callback raises, the row is already
|
|
durably ``expired`` (the flip committed first), so the loop records
|
|
:attr:`ExpiryAction.ERRORED` for that question and continues with the rest of
|
|
the batch rather than aborting the whole pass. The caller ALARMs on any
|
|
``errored`` count.
|
|
|
|
Restart-safety: all inputs are read fresh from the durable table, and each
|
|
flip is idempotent, so re-running the loop after a crash safely processes
|
|
only the rows still ``open`` and overdue.
|
|
|
|
Returns a :class:`TimerLoopReport` describing what happened to each row.
|
|
"""
|
|
resolver = policy_resolver or _default_policy_resolver
|
|
report = TimerLoopReport()
|
|
|
|
for question in overdue_open_questions(conn, now=now):
|
|
# Step 1: deterministic race on flipping ``open`` -> ``expired``.
|
|
won = expire_question(conn, question_id=question.question_id)
|
|
if not won:
|
|
# A responder answered first (or a prior pass expired it). Apply no
|
|
# no-answer policy — the question was actually answered/handled.
|
|
logger.debug(
|
|
"deadline-timer: %s lost the expire race (answered/expired first)",
|
|
question.question_id,
|
|
)
|
|
report.outcomes.append(
|
|
ExpiryOutcome(
|
|
question_id=question.question_id,
|
|
thread_id=question.thread_id,
|
|
action=ExpiryAction.LOST_RACE,
|
|
)
|
|
)
|
|
continue
|
|
|
|
# Step 2: the flip won — the row is durably ``expired``. Apply policy.
|
|
policy = resolver(question)
|
|
outcome = _apply_policy(
|
|
question,
|
|
policy=policy,
|
|
on_park=on_park,
|
|
resume_with_default=resume_with_default,
|
|
)
|
|
report.outcomes.append(outcome)
|
|
|
|
return report
|
|
|
|
|
|
def _apply_policy(
|
|
question: OverdueQuestion,
|
|
*,
|
|
policy: DeadlinePolicy,
|
|
on_park: AlarmFn,
|
|
resume_with_default: ResumeWithDefaultFn | None,
|
|
) -> ExpiryOutcome:
|
|
"""Apply a question's no-answer policy after a winning expiry flip.
|
|
|
|
Isolates the side effect: a callback that raises yields an
|
|
:attr:`ExpiryAction.ERRORED` outcome (the row stays ``expired``) instead of
|
|
crashing the whole sweep. A ``DEFAULT_ANSWER`` policy with no resume callback
|
|
is a configuration error and re-raises :class:`ValueError`.
|
|
"""
|
|
if policy is DeadlinePolicy.DEFAULT_ANSWER and resume_with_default is None:
|
|
raise ValueError(
|
|
f"question {question.question_id!r} resolved to DEFAULT_ANSWER but no "
|
|
"resume_with_default callback was supplied"
|
|
)
|
|
|
|
try:
|
|
if policy is DeadlinePolicy.PARK:
|
|
on_park(question)
|
|
action = ExpiryAction.PARKED
|
|
else: # DeadlinePolicy.DEFAULT_ANSWER (resume callback guaranteed above)
|
|
assert resume_with_default is not None # narrowed by the guard
|
|
resume_with_default(question)
|
|
action = ExpiryAction.DEFAULTED
|
|
except Exception as exc: # noqa: BLE001 — isolate one row's side-effect failure
|
|
logger.exception(
|
|
"deadline-timer: side effect for %s (%s) raised; row remains expired",
|
|
question.question_id,
|
|
policy.value,
|
|
)
|
|
return ExpiryOutcome(
|
|
question_id=question.question_id,
|
|
thread_id=question.thread_id,
|
|
action=ExpiryAction.ERRORED,
|
|
policy=policy,
|
|
error=f"{type(exc).__name__}: {exc}",
|
|
)
|
|
|
|
return ExpiryOutcome(
|
|
question_id=question.question_id,
|
|
thread_id=question.thread_id,
|
|
action=action,
|
|
policy=policy,
|
|
)
|