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.
This commit is contained in:
Simon van Laak
2026-08-13 10:18:17 -07:00
committed by Teknium
parent 47ed5e4964
commit eb137762ff
5 changed files with 533 additions and 4 deletions
+5
View File
@@ -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
+311 -4
View File
@@ -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 (
+9
View File
@@ -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
+201
View File
@@ -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,
@@ -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). |