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 3847e43ba3 feat(agent-team): one Slack thread per task — root "Task received" message + threaded questions/milestones
WS Slack-UX Feature 1. A /new-task task now maps to ONE Slack thread instead of
several top-level messages.

- /new-task posts an immediate root "📥 Task received: …" ack and captures its
  ts (root_ts); this is the instant acknowledgement.
- root_ts is plumbed into start: new PipelineState/TaskRecord channel
  slack_thread_ts, seeded by graph.start_task and threaded through
  Coordinator.start_task. The NewTaskCallback is now (task_text, via, root_ts).
- All clarifier questions for the task post as THREADED REPLIES under root_ts
  (chat.postMessage thread_ts=root_ts), and each question's ledger channel_ref
  is set to root_ts (NOT the reply's own ts). Because answer-mapping resolves a
  reply via find_open_question_by_channel_ref(thread_ts), a reply in the root
  thread (thread_ts==root_ts) maps to the task's currently-open question with NO
  change to the mapping logic or the first-answer-wins CAS. The open-only
  partial-unique index still holds (one open question per task at a time).
- Lifecycle milestones (parked / plan-ready / needs-input) and follow-up
  questions thread under root_ts too; the notify sink gained an optional
  thread_ts kwarg (degrades to top-level on a sink that doesn't accept it).
  notify failures still never break tick.
- SlackTransport.post_question + the live poster accept/forward thread_ts.
- No root_ts (non-/new-task origin) ⇒ top-level posts exactly as before.

AUTHZ-01 (owner-allowlist-first, fail-closed) and the atomic open→answered
compare-and-set are unchanged.

Adds plumbing for the inbound-ack reactor seam used by Feature 2 (dormant until
a reactor is injected). Tests cover thread_ts forwarding, channel_ref=root_ts,
graph seeding, and coordinator threading.
2026-06-23 15:49:40 -04:00

465 lines
18 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,
thread_ts: 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.
``thread_ts`` (one-thread-per-task, Slack) — when set, the question is posted
as a THREADED REPLY under that root message ``ts`` (the task's
"📥 Task received" ack post) AND the row's durable ``channel_ref`` is set to
that SAME root ``ts`` (NOT the posted reply's own ``ts``). This is what makes
the answer-mapping unchanged: a human reply in the root thread carries
``thread_ts == root_ts``, and ``find_open_question_by_channel_ref(thread_ts)``
resolves it to this task's currently-open question. When ``None`` the
question is posted top-level and the ``channel_ref`` is the posted message's
own ``ts`` exactly as before.
"""
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. ``thread_ts``
# is only forwarded when set, so non-threading transports keep their
# existing call shape.
try:
post_kwargs: dict[str, Any] = {}
if thread_ts:
post_kwargs["thread_ts"] = thread_ts
posted_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,
**post_kwargs,
)
except Exception:
return None
# 3. Persist the ref so reconcile/recovery can act on the post. When the post
# threaded under a root message, the durable channel_ref is the ROOT ts
# (so an inbound reply's thread_ts maps back to this question via
# find_open_question_by_channel_ref), NOT the posted reply's own ts.
channel_ref = thread_ts if thread_ts else posted_ref
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()