From eb137762ff99f16d798eea0907d25e56b32824ba Mon Sep 17 00:00:00 2001 From: Simon van Laak Date: Thu, 13 Aug 2026 10:18:17 -0700 Subject: [PATCH] feat(slack): render tool progress as native plan/task cards (opt-in) Adds platforms.slack.extra.native_task_cards: when enabled, live tool calls render as Slack-native plan/task cards via chat.startStream / chat.appendStream (task_display_mode: plan, task_update chunks) instead of text/edit progress bubbles. ID-bearing tool_start/tool_complete callbacks correlate concurrent same-name tool calls correctly; any native API failure falls back to one continuously edited text update. The stream is stopped exactly once when the turn finalizes. Salvaged from PR #29496 onto current main (TurnRunner/TurnContext seam); closes #29483. --- cli-config.yaml.example | 5 + gateway/run.py | 315 ++++++++++++++++++++- gateway/turn_context.py | 9 + plugins/platforms/slack/adapter.py | 201 +++++++++++++ website/docs/user-guide/messaging/slack.md | 7 + 5 files changed, 533 insertions(+), 4 deletions(-) diff --git a/cli-config.yaml.example b/cli-config.yaml.example index fb36868ce8..4ec0a24a40 100644 --- a/cli-config.yaml.example +++ b/cli-config.yaml.example @@ -1060,6 +1060,11 @@ platform_toolsets: # priority_mode: prepend # priority: # - my_plugin_command +# slack: +# extra: +# # Render live tool calls as Slack-native plan/task cards. This explicit +# # opt-in works even though Slack text tool_progress defaults to off. +# native_task_cards: false # webhook: # extra: # # Route scripts default to a 30 second timeout. Scripts must live under diff --git a/gateway/run.py b/gateway/run.py index e819c43d4a..d589999701 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -3932,6 +3932,16 @@ class TurnRunner: ctx.progress_queue.put(msg) return + # Native task cards consume the authoritative ID-bearing + # tool_start/tool_complete callbacks instead. Do not also enqueue + # name-correlated text events, which would duplicate cards and + # mispair concurrent calls to the same tool. + if ctx._native_slack_task_cards and event_type in { + "tool.started", + "tool.completed", + }: + return + # If tool_progress is off, only _thinking passes through (above). # Regular tool calls are suppressed. if not ctx.tool_progress_enabled: @@ -4108,6 +4118,188 @@ class TurnRunner: ctx.progress_queue.put(msg) + async def _send_native_task_card_progress(self, adapter) -> None: + """Drain the progress queue into Slack-native plan/task cards (#29483). + + Consumes the ID-bearing lifecycle dicts queued by + native_tool_start_callback / native_tool_complete_callback and renders + them through the adapter's chat.startStream plan/task-card stream. + On any native failure, falls back to an editable in-thread text + message so progress stays live for the rest of the turn. + """ + ctx = self._ctx + tasks: Dict[str, Dict[str, str]] = {} + task_order: List[str] = [] + fallback_msg_id: Optional[str] = None + native_failed = False + anonymous_seq = 0 + + def _compact(value: Any, limit: int = 120) -> str: + text = re.sub(r"\s+", " ", str(value or "")).strip() + if len(text) <= limit: + return text + return text[: limit - 3].rstrip() + "..." + + def _visible_tasks() -> List[Dict[str, str]]: + return [tasks[task_id] for task_id in task_order[-8:]] + + def _fallback_text() -> str: + labels = { + "in_progress": "running", + "complete": "complete", + "error": "error", + } + lines = [ + f"- {task['title']} - {labels.get(task['status'], task['status'])}" + for task in _visible_tasks() + ] + return "Hermes is working\n" + "\n".join(lines) + + def _apply_native_event(raw: Any) -> bool: + nonlocal anonymous_seq + if not isinstance(raw, dict): + return False + event_type = raw.get("type") + if event_type not in {"tool.started", "tool.completed"}: + return False + call_id = str(raw.get("tool_call_id") or "") + if not call_id: + anonymous_seq += 1 + call_id = f"anonymous_{anonymous_seq}" + tool_name = str(raw.get("tool_name") or "tool") + + if event_type == "tool.started": + title = tool_name + preview = _compact(raw.get("preview"), 64) + if preview: + title = f"{tool_name} - {preview}" + if call_id not in tasks: + task_order.append(call_id) + tasks[call_id] = { + "id": call_id, + "title": _compact(title), + "status": "in_progress", + } + return True + + task = tasks.get(call_id) + if task is None: + # Completion-only events are rare but valid on some + # runtimes. Keep their real ID instead of guessing a + # same-name pending call. + task = { + "id": call_id, + "title": _compact(tool_name), + "status": "in_progress", + } + tasks[call_id] = task + task_order.append(call_id) + task["status"] = "error" if raw.get("is_error") else "complete" + return True + + async def _send_or_edit_fallback() -> None: + nonlocal fallback_msg_id + text = _fallback_text() + if fallback_msg_id: + result = await adapter.edit_message( + chat_id=ctx.source.chat_id, + message_id=fallback_msg_id, + content=text, + metadata=ctx._progress_metadata, + ) + if getattr(result, "success", False): + return + result = await adapter.send( + chat_id=ctx.source.chat_id, + content=text, + reply_to=ctx._progress_reply_to, + metadata=ctx._progress_metadata, + ) + if getattr(result, "success", False) and getattr( + result, "message_id", None + ): + fallback_msg_id = str(result.message_id) + if ctx._cleanup_progress: + ctx._cleanup_msg_ids.append(fallback_msg_id) + + async def _publish_native_progress() -> None: + nonlocal native_failed + if not tasks: + return + if not native_failed: + result = await adapter.send_native_task_card_progress( + chat_id=ctx.source.chat_id, + tasks=_visible_tasks(), + title="Hermes is working", + reply_to=ctx._progress_reply_to, + metadata=ctx._progress_metadata, + fallback_text=_fallback_text(), + ) + if getattr(result, "success", False): + return + native_failed = True + logger.warning( + "Slack native task-card progress failed; falling back " + "to an editable text update: %s", + getattr(result, "error", "unknown error"), + ) + # Once the native rail fails, every later lifecycle event + # edits the same fallback message so progress remains live. + await _send_or_edit_fallback() + + def _drain_native_queue() -> bool: + changed = False + while True: + try: + changed = _apply_native_event( + ctx.progress_queue.get_nowait() + ) or changed + except queue.Empty: + return changed + except Exception: + logger.debug( + "Slack native progress queue drain failed", + exc_info=True, + ) + return changed + + def _agent_interrupted() -> bool: + try: + _agent = ctx.agent_holder[0] if ctx.agent_holder else None + return bool( + _agent is not None and getattr(_agent, "is_interrupted", False) + ) + except Exception: + return False + + try: + while True: + if not ctx._run_still_current(): + return + try: + raw = ctx.progress_queue.get_nowait() + except queue.Empty: + await asyncio.sleep(0.1) + continue + + if _agent_interrupted(): + continue + + if _apply_native_event(raw): + await _publish_native_progress() + except asyncio.CancelledError: + if _drain_native_queue() and ctx._run_still_current(): + if not _agent_interrupted(): + await _publish_native_progress() + return + finally: + if hasattr(adapter, "stop_native_task_card_progress"): + await adapter.stop_native_task_card_progress( + ctx.source.chat_id, + reply_to=ctx._progress_reply_to, + metadata=ctx._progress_metadata, + ) + async def send_progress_messages(self): ctx = self._ctx if not ctx.progress_queue: @@ -4117,6 +4309,12 @@ class TurnRunner: if not adapter: return + if ctx._native_slack_task_cards and hasattr( + adapter, "send_native_task_card_progress" + ): + await self._send_native_task_card_progress(adapter) + return + # Skip tool progress for platforms that don't support message # editing (e.g. iMessage/BlueBubbles) — each progress update # would become a separate message bubble, which is noisy. @@ -4484,6 +4682,68 @@ class TurnRunner: except Exception as _ack_err: logger.debug("voice ack schedule failed: %s", _ack_err) + # ── Slack-native task cards: ID-bearing lifecycle callbacks (#29483) ── + # These ride agent.tool_start_callback / agent.tool_complete_callback so + # start/completion events correlate by the REAL tool-call id — the + # name-correlated text events in progress_callback would duplicate cards + # and mispair concurrent calls to the same tool. + + def native_tool_start_callback(self, call_id, tool_name, args): + """Queue an ID-correlated native progress start from the agent thread.""" + ctx = self._ctx + if not ctx.progress_queue or not ctx._run_still_current(): + return + try: + _agent = ctx.agent_holder[0] if ctx.agent_holder else None + if _agent is not None and getattr(_agent, "is_interrupted", False): + return + except Exception: + pass + from agent.display import build_tool_preview + + ctx.progress_queue.put( + { + "type": "tool.started", + "tool_call_id": str(call_id or ""), + "tool_name": str(tool_name or "tool"), + "preview": build_tool_preview( + str(tool_name or "tool"), args or {}, max_len=64 + ) + or "", + } + ) + + def native_tool_complete_callback(self, call_id, tool_name, args, result): + """Queue the matching native completion using the real tool-call ID.""" + ctx = self._ctx + if not ctx.progress_queue or not ctx._run_still_current(): + return + try: + _agent = ctx.agent_holder[0] if ctx.agent_holder else None + if _agent is not None and getattr(_agent, "is_interrupted", False): + return + except Exception: + pass + from agent.display import _detect_tool_failure + + is_error, _ = _detect_tool_failure(str(tool_name or "tool"), result) + ctx.progress_queue.put( + { + "type": "tool.completed", + "tool_call_id": str(call_id or ""), + "tool_name": str(tool_name or "tool"), + "is_error": bool(is_error), + } + ) + + def combined_tool_start_callback(self, call_id, tool_name, args): + """Compose the voice ack + native task-card start consumers.""" + ctx = self._ctx + if ctx._voice_ack_guild[0] is not None: + self.voice_ack_callback(call_id, tool_name, args) + if ctx._native_slack_task_cards: + self.native_tool_start_callback(call_id, tool_name, args) + def _step_callback_sync(self, iteration: int, prev_tools: list) -> None: ctx = self._ctx if not ctx._run_still_current(): @@ -5082,10 +5342,23 @@ class TurnRunner: ) else None ) - # Discord voice verbal-ack hook (fires once per turn on first tool - # call; armed only when in a voice channel with the mixer running). + # Compose ID-bearing lifecycle consumers: Discord's one-time voice + # ack and Slack's native task cards both ride the authoritative + # start callback, so neither has to infer identity from tool names. + _combined_start_cb = ctx.native_tool_start_callback or ctx.voice_ack_callback agent.tool_start_callback = ( - ctx.voice_ack_callback if ctx._voice_ack_guild[0] is not None else None + _combined_start_cb + if ( + ctx._voice_ack_guild[0] is not None + or ctx._native_slack_task_cards + ) + else None + ) + agent.tool_complete_callback = ( + ctx.native_tool_complete_callback + if ctx._native_slack_task_cards + and ctx.native_tool_complete_callback is not None + else None ) agent.step_callback = ctx._step_callback_sync if ctx._hooks_ref.loaded_hooks else None agent.stream_delta_callback = _stream_delta_cb @@ -26020,7 +26293,27 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew require_platform_override_for={Platform.MATTERMOST}, ) _thinking_enabled = _thinking_mode != "off" - needs_progress_queue = tool_progress_enabled or _thinking_enabled + # Slack-native task cards (#29483): when the Slack adapter's opt-in + # is set, tool progress renders as native plan/task cards via + # chat.startStream — the progress queue is needed even though Slack + # keeps ordinary text tool_progress off by default (requiring both + # flags would silently leave the native feature inactive). + _progress_adapter_for_native = self._adapter_for_source(source) + _native_slack_task_cards = False + if ( + source.platform == Platform.SLACK + and _progress_adapter_for_native is not None + and hasattr(_progress_adapter_for_native, "native_task_cards_enabled") + ): + try: + _native_slack_task_cards = bool( + _progress_adapter_for_native.native_task_cards_enabled() + ) + except Exception: + logger.debug("Slack native task-card config check failed", exc_info=True) + needs_progress_queue = ( + tool_progress_enabled or _thinking_enabled or _native_slack_task_cards + ) # Queue for progress messages (thread-safe) @@ -26112,6 +26405,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew log_mode_enabled=log_mode_enabled, interim_assistant_messages_enabled=interim_assistant_messages_enabled, needs_progress_queue=needs_progress_queue, + _native_slack_task_cards=_native_slack_task_cards, _voice_ack_fired=_voice_ack_fired, _voice_ack_guild=_voice_ack_guild, _voice_ack_loop=_voice_ack_loop, @@ -26132,6 +26426,10 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew # TurnRunner.progress_callback (bound method, same signature). turn_ctx.progress_callback = turn_runner.progress_callback turn_ctx.voice_ack_callback = turn_runner.voice_ack_callback + turn_ctx.native_tool_start_callback = turn_runner.combined_tool_start_callback + turn_ctx.native_tool_complete_callback = ( + turn_runner.native_tool_complete_callback + ) # Background task to send progress messages # Accumulates tool lines into a single message that gets edited. @@ -26209,6 +26507,15 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew # reply anchor; carry it so progress joins that thread. _progress_metadata = {"reply_to_message_id": event_message_id} _progress_metadata = _non_conversational_metadata(_progress_metadata, platform=source.platform) + if _native_slack_task_cards: + # chat.startStream in channels requires the recipient team/user + # pair; harmless extras elsewhere, so stamp them whenever known. + _progress_metadata = dict(_progress_metadata or {}) + if source.scope_id: + _progress_metadata.setdefault("recipient_team_id", source.scope_id) + _progress_metadata.setdefault("slack_team_id", source.scope_id) + if source.user_id: + _progress_metadata.setdefault("recipient_user_id", source.user_id) _progress_reply_to = ( event_message_id if ( diff --git a/gateway/turn_context.py b/gateway/turn_context.py index ac727e2b4f..302272909c 100644 --- a/gateway/turn_context.py +++ b/gateway/turn_context.py @@ -129,3 +129,12 @@ class TurnContext: _step_callback_sync: Optional[Callable] = None _event_callback_sync: Optional[Callable] = None _status_callback_sync: Optional[Callable] = None + + # --- Slack-native task-card progress (opt-in; #29483) ------------------ + # True when the Slack adapter's ``native_task_cards_enabled()`` opt-in is + # set for this turn's platform. The ID-bearing lifecycle callbacks are + # published by TurnRunner (like voice_ack_callback above) so tool starts + # and completions correlate by real tool-call ID instead of tool name. + _native_slack_task_cards: bool = False + native_tool_start_callback: Optional[Callable] = None + native_tool_complete_callback: Optional[Callable] = None diff --git a/plugins/platforms/slack/adapter.py b/plugins/platforms/slack/adapter.py index 6702c5fea3..691e25cbe4 100644 --- a/plugins/platforms/slack/adapter.py +++ b/plugins/platforms/slack/adapter.py @@ -310,6 +310,18 @@ def slack_deps_present() -> bool: return SLACK_AVAILABLE +@dataclass +class _NativeTaskCardStream: + """Serialized state for one workspace-scoped Slack progress stream.""" + + team_id: str + channel: str + thread_ts: str + stream_ts: str = "" + stopped: bool = False + lock: asyncio.Lock = field(default_factory=asyncio.Lock) + + def check_slack_requirements() -> bool: """Check if Slack dependencies are available. @@ -1033,6 +1045,12 @@ class SlackAdapter(BasePlatformAdapter): # eviction (key[2] is the thread ts). self._active_status_threads: Dict[Tuple[str, str, str], Dict[str, Any]] = {} self._ACTIVE_STATUS_THREADS_MAX = 1000 + # Native progress streams share Slack's workspace/thread isolation. + # Each stream owns a lock so concurrent start/append/stop calls cannot + # race into duplicate streams or append after finalization. + self._native_task_card_streams: Dict[ + Tuple[str, str, str], _NativeTaskCardStream + ] = {} # Best-effort guard so automatic Slack AI thread titles are set once # per visible DM thread instead of on every reply. self._titled_assistant_threads: set = set() @@ -2322,6 +2340,12 @@ class SlackAdapter(BasePlatformAdapter): "[Slack] Watchdog task raised during disconnect", exc_info=True ) + # Finalize native streams while workspace clients are still live. The + # gateway normally stops each turn's stream first; this is the shutdown + # safety net for cancellation/reconnect races. + for key, stream in list(self._native_task_card_streams.items()): + await self._stop_native_task_card_stream(key, stream) + await self._stop_socket_mode_handler() await self._close_workspace_clients() self._app = None @@ -2477,6 +2501,183 @@ class SlackAdapter(BasePlatformAdapter): ignored = self._slack_ignored_channels() return "*" in ignored or parent_channel_id in ignored + @staticmethod + def _truthy_config(value: Any) -> bool: + if isinstance(value, bool): + return value + if isinstance(value, str): + return value.strip().lower() in {"1", "true", "yes", "on"} + return bool(value) + + def native_task_cards_enabled(self) -> bool: + """Return whether Slack-native tool progress is explicitly enabled.""" + extra = self.config.extra if isinstance(self.config.extra, dict) else {} + direct = extra.get("native_task_cards", extra.get("nativeTaskCards")) + if direct is not None: + return self._truthy_config(direct) + + streaming = extra.get("streaming") + if isinstance(streaming, dict): + progress = streaming.get("progress") + if isinstance(progress, dict): + nested = progress.get( + "native_task_cards", progress.get("nativeTaskCards") + ) + if nested is not None: + return self._truthy_config(nested) + return False + + def _native_task_card_key( + self, + chat_id: str, + reply_to: Optional[str], + metadata: Optional[Dict[str, Any]], + ) -> Optional[Tuple[str, str, str]]: + thread_ts = self._resolve_thread_ts(reply_to, metadata) + if not thread_ts: + return None + return self._workspace_thread_key( + self._metadata_team_id(metadata), chat_id, str(thread_ts) + ) + + async def send_native_task_card_progress( + self, + chat_id: str, + tasks: List[Dict[str, str]], + *, + title: str = "Hermes is working", + reply_to: Optional[str] = None, + metadata: Optional[Dict[str, Any]] = None, + fallback_text: Optional[str] = None, + ) -> SendResult: + """Start or update a Slack-native plan/task progress stream.""" + if not self._app: + return SendResult(success=False, error="Not connected") + if not tasks: + return SendResult(success=False, error="No tasks") + + key = self._native_task_card_key(chat_id, reply_to, metadata) + if key is None: + return SendResult(success=False, error="No Slack thread target") + + stream = self._native_task_card_streams.get(key) + if stream is None or stream.stopped: + stream = _NativeTaskCardStream( + team_id=key[0], + channel=chat_id, + thread_ts=key[2], + ) + # There is no await between lookup and assignment, so competing + # coroutines on this event loop will observe this same lock. + self._native_task_card_streams[key] = stream + + async with stream.lock: + if stream.stopped: + return SendResult(success=False, error="Progress stream already stopped") + try: + client = self._get_client(chat_id, team_id=stream.team_id) + if not stream.stream_ts: + start_payload: Dict[str, Any] = { + "channel": chat_id, + "thread_ts": stream.thread_ts, + "task_display_mode": "plan", + } + md = metadata or {} + recipient_team_id = ( + md.get("recipient_team_id") + or md.get("team_id") + or md.get("slack_team_id") + ) + recipient_user_id = md.get("recipient_user_id") or md.get("user_id") + if recipient_team_id: + start_payload["recipient_team_id"] = recipient_team_id + if recipient_user_id: + start_payload["recipient_user_id"] = recipient_user_id + + result = await client.api_call( + "chat.startStream", json=start_payload + ) + if hasattr(result, "get"): + stream.stream_ts = str( + result.get("ts") or result.get("message_ts") or "" + ) + if not stream.stream_ts: + raise RuntimeError("Slack startStream returned no stream timestamp") + + chunks: List[Dict[str, Any]] = [ + {"type": "plan_update", "title": str(title)[:256]} + ] + for task in tasks: + status = str(task.get("status") or "in_progress") + if status not in {"in_progress", "complete", "error"}: + status = "in_progress" + task_id = str(task.get("id") or task.get("task_id") or "task") + chunks.append( + { + "type": "task_update", + "id": task_id, + "title": str(task.get("title") or task_id)[:256], + "status": status, + } + ) + + append_payload: Dict[str, Any] = { + "channel": chat_id, + "ts": stream.stream_ts, + "chunks": chunks, + } + if fallback_text: + append_payload["markdown_text"] = fallback_text + await client.api_call("chat.appendStream", json=append_payload) + return SendResult(success=True, message_id=stream.stream_ts) + except Exception as exc: # pragma: no cover - defensive logging + logger.error( + "[Slack] Native task-card progress error: %s", + exc, + exc_info=True, + ) + return SendResult(success=False, error=str(exc), retryable=True) + + async def _stop_native_task_card_stream( + self, + key: Tuple[str, str, str], + stream: _NativeTaskCardStream, + ) -> None: + async with stream.lock: + if stream.stopped: + return + stream.stopped = True + try: + if self._app and stream.stream_ts: + await self._get_client( + stream.channel, team_id=stream.team_id + ).api_call( + "chat.stopStream", + json={"channel": stream.channel, "ts": stream.stream_ts}, + ) + except Exception as exc: # pragma: no cover - defensive logging + logger.debug( + "[Slack] Native task-card stopStream failed: %s", exc + ) + finally: + if self._native_task_card_streams.get(key) is stream: + self._native_task_card_streams.pop(key, None) + + async def stop_native_task_card_progress( + self, + chat_id: str, + *, + reply_to: Optional[str] = None, + metadata: Optional[Dict[str, Any]] = None, + ) -> None: + """Finalize an active Slack-native progress stream exactly once.""" + key = self._native_task_card_key(chat_id, reply_to, metadata) + if key is None: + return + stream = self._native_task_card_streams.get(key) + if stream is not None: + await self._stop_native_task_card_stream(key, stream) + async def send( self, chat_id: str, diff --git a/website/docs/user-guide/messaging/slack.md b/website/docs/user-guide/messaging/slack.md index 683d1a823c..7b91c9487d 100644 --- a/website/docs/user-guide/messaging/slack.md +++ b/website/docs/user-guide/messaging/slack.md @@ -423,6 +423,12 @@ platforms: # Requires rich_blocks: true. Default: false. feedback_buttons: false + # Render live tool calls as Slack-native plan/task cards. This explicit + # opt-in activates native progress even when text tool_progress is off. + # If Slack rejects the native stream, Hermes keeps one editable text + # fallback current for the rest of the turn. + native_task_cards: false + # Suggested prompts pinned at the top of Agent view's Messages tab. # Either a list of {title, message} rows, or a titled object: # {title: "Start here", prompts: [{title: "Plan", message: "..."}]} @@ -454,6 +460,7 @@ platforms: | `platforms.slack.extra.reply_broadcast` | `false` | When `true`, thread replies are also posted to the main channel. Only the first chunk is broadcast. | | `platforms.slack.extra.rich_blocks` | `false` | When `true`, agent messages are rendered as [Block Kit](https://docs.slack.dev/block-kit/) blocks (headers, dividers, true nested lists, and native tables). A plain-text fallback is always sent. Tables over Slack's limits fall back to aligned monospace. No app reinstall required — it's a send-side change only. | | `platforms.slack.extra.feedback_buttons` | `false` | When `true` with `rich_blocks`, appends Slack-native feedback controls to final replies. | +| `platforms.slack.extra.native_task_cards` | `false` | When `true`, renders live tool calls as Slack-native plan/task cards. This is an explicit progress opt-in independent of Slack's default `tool_progress: off`; native API failures fall back to one continuously edited text update. | | `platforms.slack.extra.suggested_prompts` | `[]` | Up to four `{title, message}` prompts for Agent/Assistant DM entry points; accepts either a list or `{title, prompts}`. | | `platforms.slack.extra.assistant_thread_titles` | `true` | When `true`, names Agent/Assistant DM threads from the first user message. | | `platforms.slack.extra.allow_bots` | `"none"` | Controls messages from other Slack bots: `"none"` ignores them, `"mentions"` accepts a bot message only when **that message itself** @mentions Hermes, and `"all"` accepts all of them. Use `"mentions"` for the safest bot-to-bot collaboration mode. See [Accepting messages from other bots](#accepting-messages-from-other-bots-allow_bots). |