From 3847e43ba32ffa8e70c7a29e48b23e770ecba138 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Tue, 23 Jun 2026 15:11:19 -0400 Subject: [PATCH] =?UTF-8?q?feat(agent-team):=20one=20Slack=20thread=20per?= =?UTF-8?q?=20task=20=E2=80=94=20root=20"Task=20received"=20message=20+=20?= =?UTF-8?q?threaded=20questions/milestones?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- agent-team/agent_team/coordinator.py | 118 ++++++++++++-- agent-team/agent_team/graph.py | 14 ++ agent-team/agent_team/responder.py | 27 +++- agent-team/agent_team/task_model.py | 10 ++ agent-team/agent_team/transport/base.py | 7 + .../agent_team/transport/slack_adapter.py | 29 ++++ .../agent_team/transport/slack_listener.py | 153 ++++++++++++++++-- agent-team/agent_team/transport/slack_live.py | 40 ++++- agent-team/run-team.py | 6 +- agent-team/tests/test_coordinator.py | 101 +++++++++++- agent-team/tests/test_graph.py | 22 +++ agent-team/tests/test_responder.py | 53 +++++- agent-team/tests/test_slack_adapter.py | 40 +++++ .../test_ws0_ws2_ws4_plugin_slack_hook.py | 18 ++- agent-team/tests/test_ws_activation_wiring.py | 19 ++- 15 files changed, 605 insertions(+), 52 deletions(-) diff --git a/agent-team/agent_team/coordinator.py b/agent-team/agent_team/coordinator.py index 3500878..74ffc45 100644 --- a/agent-team/agent_team/coordinator.py +++ b/agent-team/agent_team/coordinator.py @@ -152,7 +152,7 @@ def default_slack_listener_factory( transport: SlackTransport, db_path: Path, enqueue_resume: Callable[[Any], None], - new_task_callback: "Callable[[str, str], str] | None" = None, + new_task_callback: "Callable[[str, str, str], str] | None" = None, ) -> Any: """Build the live :class:`SlackListener` from the coordinator's seams (D-1). @@ -170,13 +170,30 @@ def default_slack_listener_factory( ``new_task_callback`` (WS2) is forwarded to the listener so an allowlisted ``/new-task`` command starts a task. Left ``None``, the listener ignores ``/new-task`` (its built-in default) β€” the ``/new-task`` path still runs - AFTER the AUTHZ-01 owner check regardless. + AFTER the AUTHZ-01 owner check regardless. The listener posts the root + "πŸ“₯ Task received" ack through the transport's own poster (no extra wiring). + + A best-effort πŸ‘ ``reactor`` is built from ``SLACK_BOT_TOKEN`` so the listener + can acknowledge inbound answers with a reaction (requires the + ``reactions:write`` scope). If the SDK/token is unavailable the reactor is + left ``None`` (no reaction attempted) β€” it never blocks listener startup. Imported lazily for the same import-hygiene reason as the clarifier / planner factories (the listener pulls the transport + responder leaves). """ from agent_team.transport.slack_listener import SlackListener + reactor: Any = None + try: + from agent_team.transport.slack_live import build_slack_reactor + + reactor = build_slack_reactor() + except Exception: # noqa: BLE001 - no SDK/token -> run without πŸ‘ reactions + _LOG.info( + "Slack πŸ‘ reactor unavailable (no SDK/token); inbound answers will " + "not be reaction-acknowledged" + ) + return SlackListener( transport, db_path, @@ -184,6 +201,7 @@ def default_slack_listener_factory( app_token=os.environ.get("SLACK_APP_TOKEN") or None, bot_token=os.environ.get("SLACK_BOT_TOKEN") or None, new_task_callback=new_task_callback, + reactor=reactor, ) @@ -449,8 +467,8 @@ class Coordinator: deadline_window: timedelta | None = None, alarm_hook: AlarmHook | None = None, build_listener: ListenerFactory | None = None, - new_task_callback: "Callable[[str, str], str] | None" = None, - notify: "Callable[[str], None] | None" = None, + new_task_callback: "Callable[[str, str, str], str] | None" = None, + notify: "Callable[..., None] | None" = None, ) -> None: self._db_path = Path(db_path) self._transport = transport @@ -505,14 +523,17 @@ class Coordinator: # ------------------------------------------------------------------ # def set_new_task_callback( - self, callback: "Callable[[str, str], str] | None" + self, callback: "Callable[[str, str, str], str] | None" ) -> None: """Wire the WS2 ``/new-task`` handler (call before :meth:`serve`). Avoids the constructor chicken-and-egg of referencing the coordinator's own ``start_task`` at build time: the serve path constructs the - coordinator, then sets ``lambda text, source: self.start_task(...)``. Must - be set before :meth:`_maybe_start_slack_listener` runs (i.e. before + coordinator, then sets + ``lambda text, source, root_ts: self.start_task(...)``. The ``root_ts`` + (the listener's "πŸ“₯ Task received" ack ts) is forwarded into + ``start_task`` so the task threads under it (one-thread-per-task). Must be + set before :meth:`_maybe_start_slack_listener` runs (i.e. before :meth:`serve`); it is read when the listener is built. """ self._new_task_callback = callback @@ -621,7 +642,9 @@ class Coordinator: # Intake. # ------------------------------------------------------------------ # - def start_task(self, *, task_text: str, transport_name: str) -> str: + def start_task( + self, *, task_text: str, transport_name: str, slack_thread_ts: str = "" + ) -> str: """Start one task: run to the first human gate, then notify (Β§3.3, Β§3.3.1). Runs :func:`agent_team.graph.start_task` to the first clarifier @@ -639,6 +662,15 @@ class Coordinator: post-suspend ``update_state`` would clear the pending interrupt and break the human gate). ``intake_node`` returns only a partial state (status/phase), so the seeded ``task`` channel persists into CLARIFY. + + **One-thread-per-task (``slack_thread_ts``).** When the task originates + from a ``/new-task`` slash command, the listener has already posted a + root "πŸ“₯ Task received" message and passes its ``ts`` here. It is seeded + into the graph state (so every later turn can recover it) AND forwarded + to :func:`notify_question` as ``thread_ts``, so the FIRST clarifier + question posts as a threaded reply under that root and its ledger + ``channel_ref`` becomes the root ``ts``. Empty (the default) for a task + with no root post β€” the question posts top-level exactly as before. """ if self._graph is None: raise RuntimeError("Coordinator.start_task called before setup()") @@ -652,7 +684,10 @@ class Coordinator: task_text[:200].replace("\n", "\\n").replace("\r", "\\r"), ) thread_id, _state = graph_mod.start_task( - self._graph, transport=transport_name, task=task_text + self._graph, + transport=transport_name, + task=task_text, + slack_thread_ts=slack_thread_ts, ) question = graph_mod.pending_question(self._graph, thread_id=thread_id) @@ -673,6 +708,7 @@ class Coordinator: self._transport, question_set, deadline=deadline, + thread_ts=slack_thread_ts or None, ) finally: conn.close() @@ -747,15 +783,32 @@ class Coordinator: # Maintenance tick (deadline policy + drain). # ------------------------------------------------------------------ # - def _emit(self, message: str) -> None: + def _emit(self, message: str, *, thread_ts: str | None = None) -> None: """Post a lifecycle status line to the notify sink (never raises). A notification failure must never disturb the pipeline, so the sink call is wrapped; a broken Slack post is logged and swallowed. + + ``thread_ts`` (one-thread-per-task) β€” when set, the milestone is posted + as a threaded reply under that task's root "πŸ“₯ Task received" message so + every notification for a task lands in its one Slack thread. The notify + sink accepts an optional keyword ``thread_ts``; a sink that does not + (an older/simpler sink) is called positionally so the threading hint is + a no-op rather than a crash β€” the call is tried with ``thread_ts`` first + and falls back to the bare message on a ``TypeError``. """ if self._notify is None: return try: + if thread_ts: + try: + self._notify(message, thread_ts=thread_ts) + return + except TypeError: + # The sink does not accept thread_ts; degrade to a top-level + # post rather than dropping the milestone entirely. + self._notify(message) + return self._notify(message) except Exception: # noqa: BLE001 - a notify failure must not break the loop _LOG.warning("notify sink raised; dropping status message", exc_info=True) @@ -792,6 +845,11 @@ class Coordinator: desc = desc[:90] + "…" label = f'"{desc}" (`{short}`)' + # One-thread-per-task: the root "πŸ“₯ Task received" message ts. When + # set, every follow-up question AND every lifecycle milestone for this + # task threads under it. Empty for a non-/new-task origin (top-level). + root_ts = str(values.get("slack_thread_ts") or "") or None + try: question = graph_mod.pending_question(self._graph, thread_id=thread_id) except Exception: # noqa: BLE001 - never let a status check break the loop @@ -800,7 +858,8 @@ class Coordinator: if question is not None: # Multi-turn: a new clarifier question is waiting. Post it to the # transport (the drain path otherwise leaves it unposted) and tell - # the human more input is needed. + # the human more input is needed. Thread it (and its channel_ref) + # under the task's root message so the next reply maps back. try: conn = connect(self._db_path) try: @@ -810,6 +869,7 @@ class Coordinator: question["question_set"], deadline=question.get("deadline") or self._default_deadline(), + thread_ts=root_ts, ) finally: conn.close() @@ -819,7 +879,8 @@ class Coordinator: ) self._emit( f"❓ {label} β€” needs more input. A new clarifying question was " - "posted above; reply in its thread." + "posted above; reply in its thread.", + thread_ts=root_ts, ) continue @@ -844,10 +905,14 @@ class Coordinator: f"β€’ Reached phase: {phase}\n" f"β€’ What's blocking it: {blocker}\n" "β€’ Needs your review β€” the plan could not be auto-approved. " - "Re-assign with more detail, or adjust the requirement to unblock." + "Re-assign with more detail, or adjust the requirement to unblock.", + thread_ts=root_ts, ) else: - self._emit(f"βœ… {label} β€” plan ready for review (phase: {phase}).") + self._emit( + f"βœ… {label} β€” plan ready for review (phase: {phase}).", + thread_ts=root_ts, + ) @staticmethod def _summarize_blocker(values: "dict[str, Any]") -> str: @@ -995,11 +1060,20 @@ class Coordinator: or question.get("deadline") or self._default_deadline() ) + # One-thread-per-task: re-thread the re-posted question under the + # task's root message so its channel_ref is the root ts again and + # the inbound reply still maps back. Read the root ts from the live + # state (robust across the stub + live clarifier node), falling + # back to the interrupt payload. + root_ts = self._task_root_ts(thread_id) or question.get( + "slack_thread_ts" + ) responder_mod.notify_question( conn, self._transport, question_set, deadline=deadline, + thread_ts=str(root_ts) if root_ts else None, ) finally: conn.close() @@ -1321,6 +1395,22 @@ class Coordinator: # Internals. # ------------------------------------------------------------------ # + def _task_root_ts(self, thread_id: str) -> str | None: + """Return the task's Slack root-message ts (one-thread-per-task), or None. + + Reads ``slack_thread_ts`` off the live checkpointed state so milestone / + re-delivery posts can thread under the task's root "πŸ“₯ Task received" + message. Best-effort: any read failure (or an empty value) yields + ``None``, so the caller falls back to a top-level post and a status read + never breaks the loop. + """ + try: + snap = self._graph.get_state(graph_mod.thread_config(thread_id)) + values = getattr(snap, "values", {}) or {} + except Exception: # noqa: BLE001 - a state read must not break delivery + return None + return str(values.get("slack_thread_ts") or "") or None + def _default_deadline(self) -> str: """Compute a fresh ISO deadline from the configured window (Β§3.3.1).""" from datetime import datetime, timezone diff --git a/agent-team/agent_team/graph.py b/agent-team/agent_team/graph.py index 540d547..a724ec5 100644 --- a/agent-team/agent_team/graph.py +++ b/agent-team/agent_team/graph.py @@ -224,6 +224,10 @@ def clarify_node(state: PipelineState) -> PipelineState: turn = len(history) thread_id = state.get("thread_id", "") transport = state.get("transport", "") + # The Slack root-message ts (one-thread-per-task). Empty for non-/new-task + # origins; surfaced in the interrupt payload so the responder can thread the + # question post under the root message (and use it as the channel_ref). + slack_thread_ts = state.get("slack_thread_ts", "") # Stable across the resume replay of this node (see _question_id_for): the # id delivered at suspend == the ledger key == the qa_history entry. @@ -248,6 +252,7 @@ def clarify_node(state: PipelineState) -> PipelineState: "question_set": question_set, "transport": transport, "deadline": deadline, + "slack_thread_ts": slack_thread_ts, } ) @@ -592,6 +597,7 @@ def start_task( thread_id: str | None = None, transport: str = "", task: str = "", + slack_thread_ts: str = "", ) -> tuple[str, PipelineState]: """Start a new pipeline task and run it up to the first human gate (Β§3.3). @@ -608,6 +614,13 @@ def start_task( empty ``task`` (the default) seeds no description β€” the clarifier then asks for one. + ``slack_thread_ts`` is the Slack root-message ``ts`` for a ``/new-task`` task + (the "πŸ“₯ Task received" ack post) β€” when set, every clarifier question and + lifecycle notification for this task threads under it (one-thread-per-task). + Like ``task`` it is seeded into the initial invoke and persists through + INTAKE into CLARIFY (``intake_node`` returns only a partial state). Empty + (the default) for a task with no root post β€” posts are top-level as before. + The graph MUST be compiled with a checkpointer for the suspend to persist; an uncheckpointed graph would run straight through without honouring the interrupt. @@ -619,6 +632,7 @@ def start_task( status=TaskStatus.ACTIVE.value, current_phase=_phase_value(Phase.INTAKE), task=task, + slack_thread_ts=slack_thread_ts, qa_history=[], transport=transport, created_at=now, diff --git a/agent-team/agent_team/responder.py b/agent-team/agent_team/responder.py index 96333ac..997004c 100644 --- a/agent-team/agent_team/responder.py +++ b/agent-team/agent_team/responder.py @@ -156,6 +156,7 @@ def notify_question( *, 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. @@ -177,6 +178,16 @@ def notify_question( 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__ @@ -197,19 +208,29 @@ def notify_question( ) # 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. + # 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: - channel_ref = transport.post_question( + 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. + # 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), diff --git a/agent-team/agent_team/task_model.py b/agent-team/agent_team/task_model.py index 9408ed7..7a26a03 100644 --- a/agent-team/agent_team/task_model.py +++ b/agent-team/agent_team/task_model.py @@ -89,6 +89,9 @@ class TaskRecord: current_phase: Phase # Intake task description (mirrors PipelineState.task). task: str = "" + # Slack root-message ts for one-thread-per-task (mirrors + # PipelineState.slack_thread_ts). Empty for non-/new-task origins. + slack_thread_ts: str = "" qa_history: list[Any] = field(default_factory=list) plan: dict[str, Any] | None = None review_verdicts: list[Any] = field(default_factory=list) @@ -114,6 +117,12 @@ class PipelineState(TypedDict, total=False): # Seeded by graph.start_task and read by the clarifier/planner; a first-class # channel so the seeded value persists across node transitions. task: str + # The Slack root-message ``ts`` for a /new-task task (the "πŸ“₯ Task received" + # ack post). All of the task's clarifier questions and lifecycle milestone + # notifications thread under this ``ts`` so one task maps to one Slack thread. + # Empty/absent for a task that did not originate from /new-task (no root post), + # in which case posts are top-level exactly as before. + slack_thread_ts: str qa_history: list[Any] plan: dict[str, Any] | None review_verdicts: list[Any] @@ -140,6 +149,7 @@ def task_from_dict(data: dict[str, Any]) -> TaskRecord: status=TaskStatus(data["status"]), current_phase=Phase(data["current_phase"]), task=data.get("task", ""), + slack_thread_ts=data.get("slack_thread_ts", ""), qa_history=list(data.get("qa_history", [])), plan=data.get("plan"), review_verdicts=list(data.get("review_verdicts", [])), diff --git a/agent-team/agent_team/transport/base.py b/agent-team/agent_team/transport/base.py index 015cfda..2779816 100644 --- a/agent-team/agent_team/transport/base.py +++ b/agent-team/agent_team/transport/base.py @@ -84,6 +84,7 @@ class Transport(ABC): turn: int, question_set: QuestionSet, deadline: str, + thread_ts: str | None = None, ) -> str: """Deliver ``question_set`` and return its ``channel_ref``. @@ -93,6 +94,12 @@ class Transport(ABC): ``channel_ref`` is the transport's locator for the post (Slack message ``ts`` / issue-comment id / Claude session id) and is stored on the ledger row so reconcile/recovery can act on it (Β§3.3.1). + + ``thread_ts`` is an OPTIONAL transport-specific threading hint (Slack's + one-thread-per-task: post the message as a reply under that root ``ts``). + Transports without native threading may ignore it. The responder only + forwards it when set, so a transport that does not accept it is never + called with it. """ raise NotImplementedError diff --git a/agent-team/agent_team/transport/slack_adapter.py b/agent-team/agent_team/transport/slack_adapter.py index 5c0f7c3..417c33b 100644 --- a/agent-team/agent_team/transport/slack_adapter.py +++ b/agent-team/agent_team/transport/slack_adapter.py @@ -217,6 +217,19 @@ class SlackTransport(Transport): self.channel = channel self._poster: SlackPoster = poster if poster is not None else _default_poster + @property + def poster(self) -> SlackPoster: + """The injected network seam (the ``chat.postMessage`` callable). + + Exposed read-only so collaborators sharing this transport (e.g. the + inbound :class:`~agent_team.transport.slack_listener.SlackListener`) can + post auxiliary messages β€” the root "πŸ“₯ Task received" ack and the πŸ‘ + reaction-bearing posts β€” through the SAME poster the question delivery + uses, instead of constructing a second client. The foundation default + still refuses the network (no poster configured). + """ + return self._poster + def post_question( self, *, @@ -225,12 +238,23 @@ class SlackTransport(Transport): turn: int, question_set: QuestionSet, deadline: str, + thread_ts: str | None = None, ) -> str: """Render + post the question-set; return the Slack ``ts`` channel_ref. Embeds ``question_id`` in the message ``callback_id`` so an inbound answer maps back (Β§3.3.1). On any poster failure raises :class:`SlackPostError` so the ledger row stays ``open`` for reconcile. + + ``thread_ts`` (one-thread-per-task) β€” when set, the message is posted as + a THREADED REPLY under that root ``ts`` (the task's "πŸ“₯ Task received" + ack post), so every clarifier question for a task lands in one Slack + thread. When ``None`` (the default, e.g. a task that did not originate + from ``/new-task``) the message is posted top-level exactly as before. + The returned value is still the POSTED message's own ``ts``; the caller + (:func:`agent_team.responder.notify_question`) is what records the + durable ``channel_ref`` (it uses the root ``thread_ts`` when threading so + an inbound reply's ``thread_ts`` maps back to this question). """ blocks = build_question_blocks(question_set, deadline) message: dict[str, Any] = { @@ -250,6 +274,11 @@ class SlackTransport(Transport): }, }, } + # Thread under the task's root message when one exists (one thread per + # task). Only set the key when non-empty so the top-level-post behavior + # is byte-identical for non-/new-task origins. + if thread_ts: + message["thread_ts"] = thread_ts try: response = self._poster(message) diff --git a/agent-team/agent_team/transport/slack_listener.py b/agent-team/agent_team/transport/slack_listener.py index 7df310d..685b86d 100644 --- a/agent-team/agent_team/transport/slack_listener.py +++ b/agent-team/agent_team/transport/slack_listener.py @@ -82,17 +82,29 @@ from collections.abc import Callable from agent_team.db.schema import connect, find_open_question_by_channel_ref from agent_team.responder import AnswerOutcome, EnqueueResume, submit_answer -from agent_team.transport.slack_adapter import SlackTransport +from agent_team.transport.slack_adapter import SlackPoster, SlackTransport __all__ = [ "NewTaskCallback", + "Reactor", "SlackListener", ] -# Injectable callback for /new-task slash commands: receives (task_text, transport) -# and returns the minted thread_id. Injected at coordinator startup so the -# listener is testable with no coordinator and no graph. -NewTaskCallback = Callable[[str, str], str] +# Injectable callback for /new-task slash commands: receives +# (task_text, transport, slack_thread_ts) and returns the minted thread_id. +# ``slack_thread_ts`` is the root "πŸ“₯ Task received" message ts the listener +# posted before starting the task (one-thread-per-task); the callback seeds it +# into the graph state so every later question/notification threads under it. +# Injected at coordinator startup so the listener is testable with no +# coordinator and no graph. +NewTaskCallback = Callable[[str, str, str], str] + +# A best-effort πŸ‘-reaction adder: given (channel, message_ts), add a reaction so +# the human sees the machine received the inbound message. Returns nothing; any +# failure (missing scope, deleted message, transport error) must be tolerated by +# the caller. Injected so the listener is testable with a fake recorder and no +# live WebClient. +Reactor = Callable[[str, str], None] # The Slack slash command that starts a new pipeline task. _NEW_TASK_COMMAND = "/new-task" @@ -162,10 +174,20 @@ class SlackListener: answer; :meth:`serve` sources it from ``AGENT_TEAM_SLACK_OWNER_IDS`` when not injected. * ``new_task_callback`` β€” optional :data:`NewTaskCallback`; when set, the - listener handles ``/new-task `` slash commands by calling it - with ``(task_text, "slack")`` and returning ``None`` (the task is started; - the owner will receive clarifying questions via the transport). When - ``None`` (the default), ``/new-task`` commands are ignored. + listener handles ``/new-task `` slash commands by posting a + root "πŸ“₯ Task received" ack message (one-thread-per-task) and calling it + with ``(task_text, "slack", root_ts)``, returning ``None`` (the task is + started; the owner will receive clarifying questions threaded under the + root). When ``None`` (the default), ``/new-task`` commands are ignored. + * ``poster`` β€” optional :data:`~agent_team.transport.slack_adapter.SlackPoster` + used ONLY to post the root "πŸ“₯ Task received" ack for ``/new-task``. + Defaults to the injected ``transport``'s own poster so the ack goes through + the same client as the questions. Never used for answers. + * ``reactor`` β€” optional :data:`Reactor`; when set, the listener adds a πŸ‘ + reaction to an inbound message it acted on (a thread-reply answer) AFTER + the AUTHZ-01 owner check passes, best-effort. ``None`` (the default) means + no reaction is attempted. Requires the ``reactions:write`` bot scope; until + that is granted the reactor silently no-ops, which the listener tolerates. The listener never resumes the graph; it only normalizes, submits, and enqueues. See the module SECURITY note for the trust boundary. @@ -181,6 +203,8 @@ class SlackListener: bot_token: str | None = None, owner_ids: set[str] | None = None, new_task_callback: NewTaskCallback | None = None, + poster: SlackPoster | None = None, + reactor: Reactor | None = None, ) -> None: self._transport = transport self._db_path = Path(db_path) @@ -192,6 +216,11 @@ class SlackListener: self._owner_ids: set[str] = set(owner_ids) if owner_ids else set() # Injected /new-task callback (opt-in). None = ignore new-task commands. self._new_task_callback = new_task_callback + # Poster for the root "πŸ“₯ Task received" ack. Defaults to the transport's + # own poster so the ack uses the same client as the question posts. + self._poster: SlackPoster = poster if poster is not None else transport.poster + # πŸ‘-reaction adder (opt-in). None = no reaction attempted. + self._reactor = reactor # The live Socket Mode handler, retained by :meth:`serve` so :meth:`close` # can stop it cleanly on daemon shutdown. ``None`` until ``serve`` opens # the socket. @@ -249,9 +278,17 @@ class SlackListener: # slash command is never misrouted as an answer-to-a-question. Auth has # already cleared above (AUTHZ-01), so only allowlisted owners can start # tasks. The callback is opt-in; if not injected, /new-task is ignored. + # (No πŸ‘ reaction here: a slash command has no reactable message; its ack + # is the "πŸ“₯ Task received" root post instead.) if _is_new_task_command(raw_payload): return self._handle_new_task_command(raw_payload) + # πŸ‘-acknowledge the inbound message the machine is acting on (an answer + # in a task thread). Runs AFTER AUTHZ-01 (a non-owner message above + # already returned None, so this never reacts to an unauthorized sender) + # and is best-effort: a reaction failure must never break handle_event. + self._maybe_react(raw_payload) + # ``submit_answer`` calls ``transport.parse_answer`` internally, which # raises ValueError when no question_id is recoverable. A real free-text # thread reply carries no callback_id / question_id / metadata, so its @@ -315,12 +352,21 @@ class SlackListener: Extracts the task description from ``text`` (the words after the command name). If no ``new_task_callback`` is configured, logs and returns - ``None`` (ignore). Otherwise calls ``new_task_callback(task_text, "slack")`` - and returns ``None`` (the task is started; the owner receives clarifying - questions via the transport; there is no ``AnswerOutcome`` to return here). + ``None`` (ignore). Otherwise: - Any exception raised by the callback is caught and logged; the listen - loop stays alive. + 1. Posts an immediate ROOT "πŸ“₯ Task received" ack message to the channel + (one-thread-per-task) and captures its ``ts`` (``root_ts``). This is + the instant acknowledgement (a slash command has no reactable message, + so this post IS its πŸ‘). A post failure degrades to ``root_ts=""`` so + the task still starts (its questions then post top-level). + 2. Calls ``new_task_callback(task_text, "slack", root_ts)`` so the task is + seeded with the root ts and every clarifier question + lifecycle + notification threads under it. + + Returns ``None`` (the task is started; the owner receives clarifying + questions via the transport; there is no ``AnswerOutcome`` here). Any + exception raised by the callback is caught and logged; the listen loop + stays alive. """ if self._new_task_callback is None: _LOG.debug( @@ -333,11 +379,16 @@ class SlackListener: _LOG.info("/new-task received with empty description; ignoring") return None + # 1. Instant root ack (one-thread-per-task). Best-effort: a failed post + # yields root_ts="" so the task still starts (top-level questions). + root_ts = self._post_task_received(text) + try: - thread_id = self._new_task_callback(text, "slack") + thread_id = self._new_task_callback(text, "slack", root_ts) _LOG.info( - "new task started via /new-task: thread_id=%s task=%r", + "new task started via /new-task: thread_id=%s root_ts=%s task=%r", thread_id, + root_ts or "(none)", text[:80], ) except Exception: # noqa: BLE001 - keep the listen loop alive @@ -348,6 +399,76 @@ class SlackListener: ) return None + def _post_task_received(self, description: str) -> str: + """Post the root "πŸ“₯ Task received" ack to the channel; return its ``ts``. + + The instant acknowledgement for a ``/new-task`` command and the anchor for + one-thread-per-task: every clarifier question and lifecycle notification + threads under the returned ``ts``. Posts through the shared poster (the + same client the question delivery uses) to the transport's channel. + + Best-effort: any failure (no poster configured, transport error, missing + ``ts`` in the response) is swallowed and an empty string returned, so a + post problem never blocks task start β€” the task simply runs with + top-level (un-threaded) questions. + """ + try: + response = self._poster( + { + "channel": self._transport.channel, + "text": ( + f'πŸ“₯ Task received: "{description}" ' + "β€” starting (clarifying first)…" + ), + } + ) + except Exception: # noqa: BLE001 - a failed ack must not block task start + _LOG.warning( + "failed to post 'πŸ“₯ Task received' root ack; task starts un-threaded", + exc_info=True, + ) + return "" + if not isinstance(response, Mapping): + return "" + ts = response.get("ts") + if not ts: + message = response.get("message") + if isinstance(message, Mapping): + ts = message.get("ts") + return str(ts) if ts else "" + + def _maybe_react(self, raw_payload: Mapping[str, Any]) -> None: + """Add a πŸ‘ reaction to the inbound message the machine is acting on. + + Best-effort acknowledgement that the inbound answer was received. Only + fires when a ``reactor`` is configured and the payload is an Events API + message/mention carrying a channel + message ``ts`` (the reactable + thread-reply answer). MUST be called only AFTER AUTHZ-01 has passed (the + caller guarantees this), so a non-owner message is never reacted to. + + Wrapped end-to-end: a missing scope (``reactions:write`` not yet granted), + a deleted message, or any transport error is swallowed β€” a reaction + failure must never break :meth:`handle_event` or the listen loop. + """ + if self._reactor is None: + return + event = _inner_event(raw_payload) + # Only react to a real inbound message/mention (the thread-reply answer + # shape). Interactive/slash payloads have no reactable message ts here. + channel = event.get("channel") + ts = event.get("ts") + if not channel or not ts: + return + try: + self._reactor(str(channel), str(ts)) + except Exception: # noqa: BLE001 - a reaction failure must never break handling + _LOG.debug( + "πŸ‘ reaction add failed (channel=%s ts=%s); ignoring " + "(reactions:write may not be granted yet)", + channel, + ts, + ) + def _is_authorized(self, raw_payload: Mapping[str, Any]) -> bool: """Return ``True`` iff the payload's sender is an allowlisted owner. diff --git a/agent-team/agent_team/transport/slack_live.py b/agent-team/agent_team/transport/slack_live.py index 9f9868e..1a0b64d 100644 --- a/agent-team/agent_team/transport/slack_live.py +++ b/agent-team/agent_team/transport/slack_live.py @@ -38,7 +38,7 @@ do the question-id round-trip the adapter relies on. from __future__ import annotations import os -from collections.abc import Mapping +from collections.abc import Callable, Mapping from typing import Any from agent_team.transport.slack_adapter import SlackPoster, SlackTransport @@ -46,12 +46,15 @@ from agent_team.transport.slack_adapter import SlackPoster, SlackTransport __all__ = [ "build_live_slack_transport", "build_slack_poster", + "build_slack_reactor", ] # Top-level ``chat.postMessage`` keyword arguments the live poster forwards. # ``callback_id`` is deliberately excluded: it is not a postMessage parameter, # and the durable inbound key lives in ``metadata.event_payload`` instead. -_POST_MESSAGE_KEYS = ("channel", "text", "blocks", "metadata") +# ``thread_ts`` IS a postMessage parameter (one-thread-per-task threading) and is +# forwarded when present so a question/notification posts as a threaded reply. +_POST_MESSAGE_KEYS = ("channel", "text", "blocks", "metadata", "thread_ts") def build_slack_poster(token: str | None = None, *, client: Any = None) -> SlackPoster: @@ -86,6 +89,39 @@ def build_slack_poster(token: str | None = None, *, client: Any = None) -> Slack return _poster +def build_slack_reactor( + token: str | None = None, *, client: Any = None +) -> "Callable[[str, str], None]": + """Build a live ``slack_sdk``-backed πŸ‘-reaction adder (one-thread-per-task UX). + + The returned ``(channel, ts) -> None`` callable performs a Slack + ``reactions.add`` (emoji ``thumbsup``) on the message at ``(channel, ts)`` so + a human sees the machine received their inbound answer. It is wired into the + inbound :class:`~agent_team.transport.slack_listener.SlackListener` (which + only calls it AFTER the AUTHZ-01 owner check passes) and is invoked + best-effort β€” the listener swallows any failure. + + Requires the ``reactions:write`` bot scope. Until that scope is granted (the + manifest re-applied + the app reinstalled) ``reactions.add`` fails with a + ``missing_scope`` error; this reactor lets that propagate to the listener, + which swallows it, so the reaction silently no-ops rather than breaking + answer handling. + + ``client`` (optional) injects a pre-built client for testability; any object + exposing ``reactions_add(**kwargs)`` works. When omitted, a + ``slack_sdk.WebClient`` is constructed lazily from ``token`` (falling back to + ``SLACK_BOT_TOKEN``); the deferred-import / missing-token semantics match + :func:`build_slack_poster`. + """ + if client is None: + client = _build_web_client(token) + + def _reactor(channel: str, ts: str) -> None: + client.reactions_add(channel=channel, timestamp=ts, name="thumbsup") + + return _reactor + + def build_live_slack_transport( channel: str, token: str | None = None, *, client: Any = None ) -> SlackTransport: diff --git a/agent-team/run-team.py b/agent-team/run-team.py index 57f9b09..d699217 100644 --- a/agent-team/run-team.py +++ b/agent-team/run-team.py @@ -616,8 +616,10 @@ def _build_coordinator(args: argparse.Namespace) -> Any: # before serve() builds the listener. AUTHZ-01 (owner allowlist) gates this # upstream in the listener; the source label is always "slack". coordinator.set_new_task_callback( - lambda task_text, _source: coordinator.start_task( - task_text=task_text, transport_name="slack" + lambda task_text, _source, slack_thread_ts: coordinator.start_task( + task_text=task_text, + transport_name="slack", + slack_thread_ts=slack_thread_ts, ) ) return coordinator diff --git a/agent-team/tests/test_coordinator.py b/agent-team/tests/test_coordinator.py index 66a6d2b..6de75fc 100644 --- a/agent-team/tests/test_coordinator.py +++ b/agent-team/tests/test_coordinator.py @@ -53,6 +53,9 @@ class FakeTransport(Transport): def __init__(self, *, fail_post: bool = False) -> None: self.posted: list[QuestionSet] = [] + # Records the thread_ts each post was threaded under (None = top-level), + # so one-thread-per-task wiring can be asserted. + self.thread_tss: list[str | None] = [] self.fail_post = fail_post def post_question( @@ -63,10 +66,12 @@ class FakeTransport(Transport): turn: int, question_set: QuestionSet, deadline: str, + thread_ts: str | None = None, ) -> str: if self.fail_post: raise RuntimeError("simulated transport post failure") self.posted.append(question_set) + self.thread_tss.append(thread_ts) return f"fake:{question_id}" def parse_answer(self, raw: Any) -> tuple[str, Any, str]: @@ -201,6 +206,53 @@ def test_start_task_suspends_and_writes_open_ledger_row(db_path: Path) -> None: assert len(transport.posted) == 1 +def test_start_task_threads_first_question_under_root(db_path: Path) -> None: + """One-thread-per-task: a /new-task root ts threads the first question + ref. + + When start_task is given ``slack_thread_ts`` (the listener's "πŸ“₯ Task + received" ack ts), the first clarifier question posts threaded under it AND + the ledger ``channel_ref`` is that ROOT ts β€” so the human's reply in the + thread (thread_ts == root) maps back to this open question with NO change to + the answer-mapping logic. + """ + transport = FakeTransport() + coord = _make_coordinator(db_path, transport=transport) + coord.setup() + + root_ts = "1700000000.ROOT" + thread_id = coord.start_task( + task_text="build a thing", + transport_name="slack", + slack_thread_ts=root_ts, + ) + assert thread_id + + # The question post threaded under the root. + assert transport.thread_tss == [root_ts] + # The ledger channel_ref is the ROOT ts (so an inbound reply's thread_ts maps + # back via find_open_question_by_channel_ref), not the posted reply's ref. + row = _only_open_row(db_path) + assert row["channel_ref"] == root_ts + # And the root ts is seeded on the graph state for later turns/notifications. + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state.get("slack_thread_ts") == root_ts + + +def test_start_task_without_root_posts_top_level(db_path: Path) -> None: + """No slack_thread_ts (e.g. a non-/new-task origin): top-level, ref = post ref.""" + transport = FakeTransport() + coord = _make_coordinator(db_path, transport=transport) + coord.setup() + + thread_id = coord.start_task(task_text="x", transport_name="github") + + assert transport.thread_tss == [None] + row = _only_open_row(db_path) + assert row["channel_ref"] == f"fake:{row['question_id']}" + state = graph_mod.get_pipeline_state(coord.graph, thread_id=thread_id) + assert state.get("slack_thread_ts") == "" + + def test_start_task_lost_post_leaves_open_row_without_ref(db_path: Path) -> None: # A failed transport post is recoverable: the row stays open with no ref. transport = FakeTransport(fail_post=True) @@ -953,7 +1005,9 @@ def test_followups_posts_new_question_and_emits_needs_input( monkeypatch.setattr( coord_mod.responder_mod, "notify_question", - lambda conn, transport, qset, *, deadline: posted.append((qset, deadline)), + lambda conn, transport, qset, *, deadline, thread_ts=None: posted.append( + (qset, deadline, thread_ts) + ), ) coord._post_resume_followups([_resume_result("abc12345deadbeef")]) @@ -962,6 +1016,51 @@ def test_followups_posts_new_question_and_emits_needs_input( assert any("abc12345" in m for m in msgs) +def test_followups_thread_under_task_root_ts(db_path: Path, monkeypatch: Any) -> None: + """A multi-turn follow-up question + milestone thread under the task root ts. + + The task's ``slack_thread_ts`` (seeded by start_task) is read off the live + state and forwarded as ``thread_ts`` to both the follow-up notify_question + and the milestone _emit, so one-thread-per-task holds across turns. + """ + from agent_team import coordinator as coord_mod + + posted: list[Any] = [] + emitted: list[tuple[str, Any]] = [] + coord = _make_coordinator(db_path) + # A notify sink that accepts the optional thread_ts kwarg. + coord._notify = lambda message, *, thread_ts=None: emitted.append( + (message, thread_ts) + ) + coord.setup() + + # A real task carrying a root ts on its state. + root_ts = "1700000000.ROOT" + thread_id = coord.start_task( + task_text="ship it", transport_name="slack", slack_thread_ts=root_ts + ) + + monkeypatch.setattr( + coord_mod.graph_mod, + "pending_question", + lambda _g, *, thread_id: {"question_set": object(), "deadline": "2099-01-01"}, + ) + monkeypatch.setattr( + coord_mod.responder_mod, + "notify_question", + lambda conn, transport, qset, *, deadline, thread_ts=None: posted.append( + thread_ts + ), + ) + + coord._post_resume_followups([_resume_result(thread_id)]) + + # The follow-up question threaded under the root ts. + assert posted == [root_ts] + # The "needs more input" milestone also threaded under the root ts. + assert any(ts == root_ts for _msg, ts in emitted) + + def test_followups_emits_parked(db_path: Path, monkeypatch: Any) -> None: from agent_team import coordinator as coord_mod from agent_team.task_model import TaskStatus diff --git a/agent-team/tests/test_graph.py b/agent-team/tests/test_graph.py index 48d0b97..b4b58be 100644 --- a/agent-team/tests/test_graph.py +++ b/agent-team/tests/test_graph.py @@ -147,6 +147,28 @@ def test_start_task_seeds_task_description_into_state(compiled) -> None: assert state.get("task") == "build a login form" +def test_start_task_seeds_slack_thread_ts_into_state(compiled) -> None: + # One-thread-per-task: the root "πŸ“₯ Task received" message ts must reach the + # graph state so every later question/notification threads under it. Like + # ``task`` it is seeded into the initial invoke and persists through INTAKE + # into the suspended CLARIFY snapshot. + _thread_id, state = start_task( + compiled, transport="slack", slack_thread_ts="1700000000.ROOT" + ) + assert state.get("slack_thread_ts") == "1700000000.ROOT" + # It is also surfaced on the pending interrupt payload so the responder can + # thread the question post. + payload = pending_question(compiled, thread_id=_thread_id) + assert payload["slack_thread_ts"] == "1700000000.ROOT" + + +def test_start_task_default_slack_thread_ts_is_empty(compiled) -> None: + # A task with no root post (e.g. a GitHub-issue origin): slack_thread_ts is + # empty so questions post top-level exactly as before. + _thread_id, state = start_task(compiled, transport="slack") + assert state.get("slack_thread_ts") == "" + + def test_pending_question_carries_foundation_questionset(compiled) -> None: thread_id, _ = start_task(compiled, transport="slack") payload = pending_question(compiled, thread_id=thread_id) diff --git a/agent-team/tests/test_responder.py b/agent-team/tests/test_responder.py index 914b6ae..b541d29 100644 --- a/agent-team/tests/test_responder.py +++ b/agent-team/tests/test_responder.py @@ -71,7 +71,7 @@ class FakeTransport(Transport): self.post_fails = post_fails def post_question( - self, *, thread_id, question_id, turn, question_set, deadline + self, *, thread_id, question_id, turn, question_set, deadline, thread_ts=None ) -> str: if self.post_fails: raise RuntimeError("transport unreachable") @@ -82,6 +82,7 @@ class FakeTransport(Transport): "question_id": question_id, "turn": turn, "deadline": deadline, + "thread_ts": thread_ts, "ref": ref, } ) @@ -197,6 +198,56 @@ def test_notify_lost_post_leaves_open_row_without_ref( assert row["channel_ref"] is None +def test_notify_threads_under_root_and_sets_channel_ref_to_root( + conn: sqlite3.Connection, +) -> None: + """One-thread-per-task: thread_ts is forwarded to post AND becomes channel_ref. + + When ``thread_ts`` (the task's root "πŸ“₯ Task received" ts) is given, the + question posts as a threaded reply under it, and the durable ``channel_ref`` + is set to that ROOT ts (NOT the posted reply's own ts) β€” so an inbound reply + whose ``thread_ts == root_ts`` maps back via + ``find_open_question_by_channel_ref``. + """ + transport = FakeTransport() + qs = _question_set() + root_ts = "1700000000.ROOT" + + ref = notify_question( + conn, + transport, + qs, + deadline="2026-06-18T00:00:00+00:00", + thread_ts=root_ts, + ) + + # The post threaded under the root. + assert transport.posts[0]["thread_ts"] == root_ts + # channel_ref is the ROOT ts, not the posted reply ref ("slack-ts-q1"). + assert ref == root_ts + row = _row(conn, "q1") + assert row["channel_ref"] == root_ts + + +def test_notify_without_thread_ts_uses_posted_ref_as_channel_ref( + conn: sqlite3.Connection, +) -> None: + """No thread_ts (the default) preserves the prior behavior exactly. + + The post is top-level (thread_ts None) and the channel_ref is the posted + message's own ts (the FakeTransport ref). + """ + transport = FakeTransport() + + ref = notify_question( + conn, transport, _question_set(), deadline="2026-06-18T00:00:00+00:00" + ) + + assert transport.posts[0]["thread_ts"] is None + assert ref == "slack-ts-q1" + assert _row(conn, "q1")["channel_ref"] == "slack-ts-q1" + + # --------------------------------------------------------------------------- # submit_answer β€” first-answer-wins (Β§3.3.1). # --------------------------------------------------------------------------- diff --git a/agent-team/tests/test_slack_adapter.py b/agent-team/tests/test_slack_adapter.py index dc9e892..25376e7 100644 --- a/agent-team/tests/test_slack_adapter.py +++ b/agent-team/tests/test_slack_adapter.py @@ -169,6 +169,46 @@ def test_post_question_embeds_question_id_in_callback_id() -> None: assert sent["metadata"]["event_payload"]["thread_id"] == "thread-1" +def test_post_question_threads_under_thread_ts_when_given() -> None: + """One-thread-per-task: a non-empty thread_ts is forwarded as message['thread_ts']. + + The poster receives ``thread_ts`` so Slack posts the question as a threaded + reply under the task's root "πŸ“₯ Task received" message. + """ + poster = _RecordingPoster() + t = SlackTransport(channel="C999", poster=poster) + t.post_question( + thread_id="thread-1", + question_id="q-abc", + turn=0, + question_set=_question_set(), + deadline="2026-06-18T00:00:00Z", + thread_ts="1700000000.ROOT", + ) + assert poster.calls[0]["thread_ts"] == "1700000000.ROOT" + + +def test_post_question_omits_thread_ts_by_default() -> None: + """No thread_ts (the default) => top-level post (no 'thread_ts' key).""" + poster = _RecordingPoster() + t = SlackTransport(channel="C999", poster=poster) + t.post_question( + thread_id="thread-1", + question_id="q-abc", + turn=0, + question_set=_question_set(), + deadline="2026-06-18T00:00:00Z", + ) + assert "thread_ts" not in poster.calls[0] + + +def test_poster_property_exposes_injected_poster() -> None: + """The transport exposes its poster so the listener can post the root ack.""" + poster = _RecordingPoster() + t = SlackTransport(channel="C1", poster=poster) + assert t.poster is poster + + def test_post_question_accepts_nested_message_ts() -> None: poster = _RecordingPoster({"ok": True, "message": {"ts": "1700000000.000300"}}) t = SlackTransport(channel="C1", poster=poster) diff --git a/agent-team/tests/test_ws0_ws2_ws4_plugin_slack_hook.py b/agent-team/tests/test_ws0_ws2_ws4_plugin_slack_hook.py index 4f198be..d14cb40 100644 --- a/agent-team/tests/test_ws0_ws2_ws4_plugin_slack_hook.py +++ b/agent-team/tests/test_ws0_ws2_ws4_plugin_slack_hook.py @@ -157,23 +157,25 @@ def test_is_new_task_command_ignores_non_slash() -> None: def test_new_task_calls_callback_with_text_and_transport() -> None: - calls: list[tuple[str, str]] = [] + calls: list[tuple[str, str, str]] = [] - def cb(task: str, transport: str) -> str: - calls.append((task, transport)) + def cb(task: str, transport: str, root_ts: str) -> str: + calls.append((task, transport, root_ts)) return "thread-abc" listener = _make_listener(new_task_callback=cb) result = listener.handle_event(_new_task_payload("Add OAuth to admin portal")) assert result is None # no AnswerOutcome for new-task - assert calls == [("Add OAuth to admin portal", "slack")] + # The default (non-posting) transport poster raises, so root_ts degrades to + # "" β€” the task still starts, un-threaded. The text + via are forwarded. + assert calls == [("Add OAuth to admin portal", "slack", "")] def test_new_task_unauthorized_sender_rejected() -> None: calls: list[Any] = [] - def cb(task: str, transport: str) -> str: + def cb(task: str, transport: str, root_ts: str) -> str: calls.append(task) return "thread-xyz" @@ -187,7 +189,7 @@ def test_new_task_unauthorized_sender_rejected() -> None: def test_new_task_empty_text_ignored() -> None: calls: list[Any] = [] - def cb(task: str, transport: str) -> str: + def cb(task: str, transport: str, root_ts: str) -> str: calls.append(task) return "thread-123" @@ -205,7 +207,7 @@ def test_new_task_no_callback_configured_is_ignored() -> None: def test_new_task_callback_exception_does_not_crash_listener() -> None: - def boom(task: str, transport: str) -> str: + def boom(task: str, transport: str, root_ts: str) -> str: raise RuntimeError("coordinator exploded") listener = _make_listener(new_task_callback=boom) @@ -222,7 +224,7 @@ def test_other_slash_command_not_intercepted_by_new_task_path( init_db(db_path) calls: list[Any] = [] - def cb(task: str, transport: str) -> str: + def cb(task: str, transport: str, root_ts: str) -> str: calls.append(task) return "thread-xyz" diff --git a/agent-team/tests/test_ws_activation_wiring.py b/agent-team/tests/test_ws_activation_wiring.py index 9988d9c..9e90760 100644 --- a/agent-team/tests/test_ws_activation_wiring.py +++ b/agent-team/tests/test_ws_activation_wiring.py @@ -97,18 +97,27 @@ def test_build_coordinator_wires_new_task_callback_to_start_task( # the real graph. calls: dict[str, Any] = {} - def _fake_start_task(*, task_text: str, transport_name: str) -> str: + def _fake_start_task( + *, task_text: str, transport_name: str, slack_thread_ts: str = "" + ) -> str: calls["task_text"] = task_text calls["transport_name"] = transport_name + calls["slack_thread_ts"] = slack_thread_ts return "thread-xyz" monkeypatch.setattr(coordinator, "start_task", _fake_start_task) cb = coordinator._new_task_callback assert cb is not None - thread_id = cb("fix the flaky test", "slack") + # The 3-arg callback (one-thread-per-task): (task_text, via, root_ts). The + # root_ts is forwarded into start_task so the task threads under the root. + thread_id = cb("fix the flaky test", "slack", "1700000000.000100") assert thread_id == "thread-xyz" - assert calls == {"task_text": "fix the flaky test", "transport_name": "slack"} + assert calls == { + "task_text": "fix the flaky test", + "transport_name": "slack", + "slack_thread_ts": "1700000000.000100", + } def test_set_new_task_callback_overrides() -> None: @@ -124,7 +133,7 @@ def test_set_new_task_callback_overrides() -> None: coord = Coordinator(db_path=":memory:", transport=_T()) assert coord._new_task_callback is None - sentinel = lambda t, s: "tid" # noqa: E731 + sentinel = lambda t, s, r: "tid" # noqa: E731 coord.set_new_task_callback(sentinel) assert coord._new_task_callback is sentinel @@ -150,7 +159,7 @@ def test_slack_listener_factory_forwards_new_task_callback( monkeypatch.setattr(sl_mod, "SlackListener", _FakeListener) - sentinel = lambda t, s: "tid" # noqa: E731 + sentinel = lambda t, s, r: "tid" # noqa: E731 coord_mod.default_slack_listener_factory( transport=MagicMock(), db_path=tmp_path / "x.sqlite",