"""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, )