diff --git a/plugins/platforms/slack/adapter.py b/plugins/platforms/slack/adapter.py index 42dc734217..456c7979a5 100644 --- a/plugins/platforms/slack/adapter.py +++ b/plugins/platforms/slack/adapter.py @@ -36,22 +36,10 @@ from agent.secret_scope import UnscopedSecretError, get_secret from gateway.config import Platform, PlatformConfig from gateway.platforms.helpers import MessageDeduplicator from gateway.platforms.base import ( - gateway_trust_env, - BasePlatformAdapter, - MessageEvent, - MessageType, - ProcessingOutcome, - SendResult, - SUPPORTED_DOCUMENT_TYPES, - SUPPORTED_VIDEO_TYPES, - _TEXT_INJECT_EXTENSIONS, - is_host_excluded_by_no_proxy, - resolve_proxy_url, - safe_url_for_log, - _ssrf_redirect_guard, - cache_document_from_bytes, - cache_video_from_bytes, -) + gateway_trust_env, BasePlatformAdapter, MessageEvent, MessageType, ProcessingOutcome, + SendResult, SUPPORTED_DOCUMENT_TYPES, SUPPORTED_VIDEO_TYPES, _TEXT_INJECT_EXTENSIONS, + is_host_excluded_by_no_proxy, resolve_proxy_url, safe_url_for_log, _ssrf_redirect_guard, + cache_document_from_bytes, cache_video_from_bytes) try: # sibling module; support both package and flat plugin-dir import from .block_kit import render_blocks, sanitize_blocks @@ -91,8 +79,7 @@ def _slack_unfurl_kwargs(extra: Optional[Dict[str, Any]]) -> Dict[str, bool]: async def _read_error_text_limited( - response: Any, *, limit: int = _SLACK_ERROR_BODY_LIMIT_BYTES -) -> str: + response: Any, *, limit: int = _SLACK_ERROR_BODY_LIMIT_BYTES) -> str: content = getattr(response, "content", None) read = getattr(content, "read", None) if callable(read): @@ -191,9 +178,7 @@ def _align_table(rows: List[str]) -> List[str]: widths = [max(_disp_width(r[c]) for r in parsed) for c in range(n_cols)] out: List[str] = [] for idx, row in enumerate(parsed): - cells = [ - "-" * widths[c] if idx == 1 else _pad(row[c], widths[c]) for c in range(n_cols) - ] + cells = ["-" * widths[c] if idx == 1 else _pad(row[c], widths[c]) for c in range(n_cols)] out.append("| " + " | ".join(cells) + " |") return out @@ -214,8 +199,7 @@ def _wrap_markdown_tables(text: str) -> str: not in_fence and "|" in line and i + 1 < len(lines) - and _TABLE_SEPARATOR_RE.match(lines[i + 1]) - ): + and _TABLE_SEPARATOR_RE.match(lines[i + 1])): block = [line, lines[i + 1]] j = i + 2 while j < len(lines) and _is_table_row(lines[j]): @@ -235,8 +219,7 @@ def _wrap_markdown_tables(text: str) -> str: # match the right stashed response_url under concurrent slashes. ContextVars # propagate to child asyncio.Tasks, so the background processing task sees it. _slash_user_id: contextvars.ContextVar[Optional[str]] = contextvars.ContextVar( - "_slash_user_id", default=None -) + "_slash_user_id", default=None) @dataclass @@ -292,8 +275,7 @@ def check_slack_requirements() -> bool: "AsyncSocketModeHandler": AsyncSocketModeHandler, "AsyncWebClient": AsyncWebClient, "aiohttp": aiohttp, - "SLACK_AVAILABLE": True, - } + "SLACK_AVAILABLE": True} from tools.lazy_deps import ensure_and_bind @@ -395,8 +377,7 @@ _INLINE_ENTITY_FORMATS = { "usergroup": ("", "usergroup_id", ""), "team": ("", "team_id", ""), "emoji": (":{}:", "name", ""), - "broadcast": ("", "range", "here"), -} + "broadcast": ("", "range", "here")} def _render_slack_inline_element(el: dict) -> str: @@ -460,10 +441,8 @@ def _extract_text_from_slack_blocks(blocks: list) -> str: _render_inline_elements( child.get("elements", []) if child.get("type", "") == "rich_text_section" - else [child] - ) - for child in elem.get("elements", []) - ] + else [child]) + for child in elem.get("elements", [])] code_text = "\n".join(line for line in code_lines if line) if code_text: lang = elem.get("language", "") @@ -565,14 +544,12 @@ def _normalize_slack_text_for_dedupe(text: str, bot_uid: str = "") -> str: def _extract_additional_text_from_slack_blocks( - blocks: list, primary_text: str, bot_uid: str = "" -) -> str: + blocks: list, primary_text: str, bot_uid: str = "") -> str: """Render rich-text content not already represented by primary_text.""" primary = _normalize_slack_text_for_dedupe(primary_text, bot_uid) primary_fenced = { _normalize_slack_text_for_dedupe(match.group(0), bot_uid) - for match in _SLACK_FENCED_CODE_RE.finditer(primary_text or "") - } + for match in _SLACK_FENCED_CODE_RE.finditer(primary_text or "")} parts: list[str] = [] for block in blocks or []: if (block or {}).get("type") != "rich_text": @@ -580,8 +557,7 @@ def _extract_additional_text_from_slack_blocks( for element in block.get("elements", []): element_type = element.get("type", "") rendered = _extract_text_from_slack_blocks( - [{"type": "rich_text", "elements": [element]}] - ).strip() + [{"type": "rich_text", "elements": [element]}]).strip() if not rendered: continue normalized = _normalize_slack_text_for_dedupe(rendered, bot_uid) @@ -597,12 +573,10 @@ def _extract_additional_text_from_slack_blocks( # Block Kit keys kept in the agent-facing payload dump (scalars copied; containers recursed). _BLOCK_SCALAR_KEYS = frozenset( - "type block_id action_id style dispatch_action optional multiple emoji".split() -) + "type block_id action_id style dispatch_action optional multiple emoji".split()) _BLOCK_RECURSIVE_KEYS = frozenset( "text title description label placeholder accessory fields elements options " - "option_groups confirm submit close hint".split() -) + "option_groups confirm submit close hint".split()) def _serialize_slack_blocks_for_agent(blocks: list, max_chars: int = 6000) -> str: @@ -615,8 +589,7 @@ def _serialize_slack_blocks_for_agent(blocks: list, max_chars: int = 6000) -> st def _sanitize(value): if isinstance(value, list): return [ - item for item in (_sanitize(v) for v in value) if item not in (None, {}, [], "") - ] + item for item in (_sanitize(v) for v in value) if item not in (None, {}, [], "")] if isinstance(value, dict): sanitized = {} for key, item in value.items(): @@ -653,10 +626,10 @@ def _extract_urls_from_slack_blocks(blocks: list) -> list[str]: if isinstance(node, dict): for key in ("url", "image_url", "external_url"): value = node.get(key) - if isinstance(value, str) and value.startswith(("http://", "https://")): - if value not in seen: - seen.add(value) - found.append(value) + is_url = isinstance(value, str) and value.startswith(("http://", "https://")) + if is_url and value not in seen: + seen.add(value) + found.append(value) for value in node.values(): _walk(value) elif isinstance(node, list): @@ -703,8 +676,7 @@ async def _cancel_socket_tasks(tasks: Any) -> None: for task in tasks if task is not None and callable(getattr(task, "cancel", None)) - and not (callable(getattr(task, "done", None)) and task.done()) - ] + and not (callable(getattr(task, "done", None)) and task.done())] for task in live: task.cancel() pending = set(live) @@ -718,10 +690,8 @@ async def _cancel_socket_tasks(tasks: Any) -> None: logger.debug("[Slack] Socket Mode task failed while stopping", exc_info=True) if still_running: # pragma: no cover - defensive logging logger.warning( - "[Slack] %d Socket Mode task(s) did not stop within %.1fs", - len(still_running), - _SOCKET_TASK_CANCEL_TIMEOUT_S, - ) + "[Slack] %d Socket Mode task(s) did not stop within %.1fs", len(still_running), + _SOCKET_TASK_CANCEL_TIMEOUT_S) _SLACK_PROXY_HOSTS = ("slack.com", "files.slack.com", "wss-primary.slack.com") @@ -736,8 +706,7 @@ def _resolve_slack_proxy_url() -> Optional[str]: if not normalized.startswith(("http://", "https://")): logger.info( "[Slack] Ignoring unsupported proxy scheme for Slack transport: %s", - safe_url_for_log(proxy_url), - ) + safe_url_for_log(proxy_url)) return None if any(is_host_excluded_by_no_proxy(host) for host in _SLACK_PROXY_HOSTS): logger.info("[Slack] NO_PROXY bypasses Slack proxy configuration") @@ -767,14 +736,12 @@ _SLACK_AUDIO_MIME_TO_EXT = { "audio/ogg": ".ogg", "audio/opus": ".ogg", "audio/mpeg": ".mp3", "audio/mp3": ".mp3", "audio/wav": ".wav", "audio/x-wav": ".wav", "audio/webm": ".webm", "audio/mp4": ".m4a", "audio/x-m4a": ".m4a", "audio/m4a": ".m4a", "audio/aac": ".m4a", "audio/flac": ".flac", - "audio/x-flac": ".flac", -} + "audio/x-flac": ".flac"} # Extensions OpenAI/Whisper-family STT backends accept (kept in sync with # tools/transcription_tools.SUPPORTED_FORMATS). _SLACK_STT_SUPPORTED_EXTS = frozenset( - {".mp3", ".mp4", ".mpeg", ".mpga", ".m4a", ".wav", ".webm", ".ogg", ".aac", ".flac"} -) + {".mp3", ".mp4", ".mpeg", ".mpga", ".m4a", ".wav", ".webm", ".ogg", ".aac", ".flac"}) # Cached extension → reported ``audio/*`` mimetype for ``video/mp4``-mislabeled # voice clips, so media_type stays coherent with the cached bytes (STT gate keys @@ -782,8 +749,7 @@ _SLACK_STT_SUPPORTED_EXTS = frozenset( _SLACK_EXT_TO_AUDIO_MIME = { ".mp4": "audio/mp4", ".m4a": "audio/mp4", ".mp3": "audio/mpeg", ".mpeg": "audio/mpeg", ".mpga": "audio/mpeg", ".wav": "audio/wav", ".webm": "audio/webm", ".ogg": "audio/ogg", - ".aac": "audio/aac", ".flac": "audio/flac", -} + ".aac": "audio/aac", ".flac": "audio/flac"} def _resolve_slack_audio_ext(file_obj: Dict[str, Any], mimetype: str) -> str: @@ -818,8 +784,7 @@ _IMAGE_CT_EXTS = (("jpeg", "jpg"), ("jpg", "jpg"), ("gif", "gif"), ("webp", "web _TRANSIENT_UPLOAD_MARKERS = ( "rate_limited", "ratelimited", "429", "connection reset", "service unavailable", - "temporarily unavailable", -) + "temporarily unavailable") _SLACK_PERMISSION_ERRORS = frozenset( @@ -841,18 +806,14 @@ def _is_transient_transport_error(e: BaseException) -> bool: error_type for error_type in ( getattr(aiohttp_module, "ClientSSLError", None), - getattr(aiohttp_module, "ServerFingerprintMismatch", None), - ) - if isinstance(error_type, type) - ) + getattr(aiohttp_module, "ServerFingerprintMismatch", None)) + if isinstance(error_type, type)) is_permanent_tls_error = bool(permanent_tls_error_types) and isinstance( - e, permanent_tls_error_types - ) + e, permanent_tls_error_types) return isinstance(e, TimeoutError) or ( isinstance(connection_error_type, type) and isinstance(e, connection_error_type) - and not is_permanent_tls_error - ) + and not is_permanent_tls_error) def _extra_or_env_flag_getter(key: str, env_var: str, *, strip: bool = False) -> Callable[..., bool]: @@ -866,8 +827,7 @@ def _extra_or_env_flag_getter(key: str, env_var: str, *, strip: bool = False) -> def _extra_or_env_channel_set_getter( - key: str, env_var: str, *, coerce_scalar: bool = False -) -> Callable[..., set]: + key: str, env_var: str, *, coerce_scalar: bool = False) -> Callable[..., set]: """Method factory: ``self._extra_or_env_channel_set(key, env_var, coerce_scalar=...)``.""" def getter(self) -> set: @@ -999,8 +959,7 @@ class SlackAdapter(BasePlatformAdapter): """Close any Slack SDK clients that may own aiohttp sessions.""" primary_client = getattr(self._app, "client", None) if self._app is not None else None clients = ([primary_client] if primary_client is not None else []) + list( - self._team_clients.values() - ) + self._team_clients.values()) seen_ids: set[int] = set() for client in clients: if id(client) in seen_ids: @@ -1025,8 +984,7 @@ class SlackAdapter(BasePlatformAdapter): @classmethod def _discard_oldest_by_thread_ts( - cls, entries: Any, count: int, ts_getter: Callable[[Any], Any] = lambda e: e - ) -> None: + cls, entries: Any, count: int, ts_getter: Callable[[Any], Any] = lambda e: e) -> None: """Discard the *count* entries (set or dict keys) with the oldest embedded Slack ts. Sets iterate in arbitrary order, so ``list(entries)[:count]`` could evict the most ACTIVE entry; sort chronologically by the embedded ts instead.""" @@ -1038,8 +996,7 @@ class SlackAdapter(BasePlatformAdapter): remove(entry) def _evict_oldest_by_ts( - self, entries: Any, cap: int, ts_getter: Callable[[Any], Any] = lambda e: e - ) -> None: + self, entries: Any, cap: int, ts_getter: Callable[[Any], Any] = lambda e: e) -> None: """Once ``entries`` exceeds ``cap``, drop oldest-ts-first down to half the cap.""" if len(entries) > cap: self._discard_oldest_by_thread_ts(entries, len(entries) - cap // 2, ts_getter) @@ -1050,8 +1007,7 @@ class SlackAdapter(BasePlatformAdapter): def _trim_mentioned_threads(self) -> None: if len(self._mentioned_threads) > self._MENTIONED_THREADS_MAX: self._discard_oldest_by_thread_ts( - self._mentioned_threads, self._MENTIONED_THREADS_MAX // 2 - ) + self._mentioned_threads, self._MENTIONED_THREADS_MAX // 2) @staticmethod def _trim_oldest_dict_entries(mapping: Dict[Any, Any], max_size: int) -> None: @@ -1114,15 +1070,13 @@ class SlackAdapter(BasePlatformAdapter): self._handler = self._socket_mode_task = None client = getattr(handler, "client", None) await _cancel_socket_tasks( - [task] + [getattr(client, attr, None) for attr in _SOCKET_CLIENT_TASK_ATTRS] - ) + [task] + [getattr(client, attr, None) for attr in _SOCKET_CLIENT_TASK_ATTRS]) if handler is not None: try: await handler.close_async() except Exception as e: # pragma: no cover - defensive logging logger.warning( - "[Slack] Error while closing Socket Mode handler: %s", e, exc_info=True - ) + "[Slack] Error while closing Socket Mode handler: %s", e, exc_info=True) async def _socket_transport_connected(self) -> Optional[bool]: """Best-effort check of current Socket Mode transport state.""" @@ -1199,8 +1153,7 @@ class SlackAdapter(BasePlatformAdapter): raise except Exception: # pragma: no cover - defensive logging logger.warning( - "[Slack] Socket Mode watchdog iteration failed; continuing", exc_info=True - ) + "[Slack] Socket Mode watchdog iteration failed; continuing", exc_info=True) def _on_socket_watchdog_done(self, task: asyncio.Task) -> None: if task is not self._socket_watchdog_task: @@ -1262,8 +1215,7 @@ class SlackAdapter(BasePlatformAdapter): loop.create_task(self._restart_socket_mode("socket task exited")) def _describe_slack_api_error( - self, response: Any, *, file_obj: Optional[Dict[str, Any]] = None - ) -> Optional[str]: + self, response: Any, *, file_obj: Optional[Dict[str, Any]] = None) -> Optional[str]: """Convert Slack API auth/permission failures into actionable user-facing text.""" if response is None or not hasattr(response, "get"): return None @@ -1278,8 +1230,7 @@ class SlackAdapter(BasePlatformAdapter): provided_hint = f" Current bot scopes: {provided}." if provided else "" return ( f"Slack attachment access failed for {file_label}. {needed_hint}{provided_hint}" - " Update the Slack app scopes/settings and reinstall the app to the workspace." - ) + " Update the Slack app scopes/settings and reinstall the app to the workspace.") if error in {"not_authed", "invalid_auth", "account_inactive", "token_revoked"}: return f"Slack attachment access failed for {file_label} because the bot token is not authorized ({error}). Refresh the token/reinstall the app." if error in {"file_not_found", "file_deleted"}: @@ -1289,8 +1240,7 @@ class SlackAdapter(BasePlatformAdapter): return None def _describe_slack_download_failure( - self, exc: Exception, *, file_obj: Optional[Dict[str, Any]] = None - ) -> Optional[str]: + self, exc: Exception, *, file_obj: Optional[Dict[str, Any]] = None) -> Optional[str]: """Translate Slack download exceptions into user-facing attachment diagnostics.""" file_label = _attachment_label(file_obj) response = getattr(exc, "response", None) @@ -1313,8 +1263,7 @@ class SlackAdapter(BasePlatformAdapter): if "Slack returned HTML instead of media" in message or "non-image data" in message: return ( f"Slack attachment access failed for {file_label}: Slack returned an HTML/login or non-media response. " - "This usually means a scope, auth, or file-permission problem." - ) + "This usually means a scope, auth, or file-permission problem.") return None # ------------------------------------------------------------------ @@ -1347,8 +1296,7 @@ class SlackAdapter(BasePlatformAdapter): for k in [ k for k, v in self._slash_command_contexts.items() - if now - v["ts"] > self._SLASH_CTX_TTL - ]: + if now - v["ts"] > self._SLASH_CTX_TTL]: self._slash_command_contexts.pop(k, None) def _format_chunks(self, content: str) -> List[str]: @@ -1371,8 +1319,7 @@ class SlackAdapter(BasePlatformAdapter): chunks = chunks[:5] chunks[-1] = ( chunks[-1].rstrip() + f"\n\n_[Reply truncated: {dropped} more part(s) exceeded " - "Slack's ephemeral reply limit.]_" - ) + "Slack's ephemeral reply limit.]_") try: async with aiohttp.ClientSession(trust_env=gateway_trust_env()) as session: for idx, chunk in enumerate(chunks): @@ -1387,16 +1334,14 @@ class SlackAdapter(BasePlatformAdapter): "[Slack] response_url POST returned %s: %s", resp.status, body[:200] ) return SendResult( - success=False, error=f"response_url POST returned {resp.status}" - ) + success=False, error=f"response_url POST returned {resp.status}") return SendResult(success=True, message_id=None) except Exception as e: logger.warning("[Slack] response_url POST failed: %s", e) return SendResult(success=False, error=str(e)) async def _post_ephemeral_fallback( - self, chat_id: str, ctx: Dict[str, Any], content: str - ) -> "SendResult": + self, chat_id: str, ctx: Dict[str, Any], content: str) -> "SendResult": """Deliver a slash reply via ``chat.postEphemeral`` when ``response_url`` fails. Keeps the reply private (a public channel post must never happen for an ephemeral reply). Cannot ``replace_original``, so the ack stays; no 5-POST cap applies here.""" @@ -1441,8 +1386,7 @@ class SlackAdapter(BasePlatformAdapter): "'mpim:history' (and 'mpim:read') to bot scopes, add 'message.mpim' to event " "subscriptions, then REINSTALL the app to the workspace. Regenerating the app " "from `hermes slack` produces a manifest with these already included.", - team_key or "this workspace", - ) + team_key or "this workspace") except Exception: # pragma: no cover - diagnostics must never break connect pass @@ -1480,30 +1424,12 @@ class SlackAdapter(BasePlatformAdapter): "OF THAT PERSON will be misrouted as mentions of the bot (the bot replies to " "messages merely addressed to them). Use the 'Bot User OAuth Token' " "(xoxb-...) from your Slack app's 'OAuth & Permissions' page in " - "SLACK_BOT_TOKEN.", - team_key or "this workspace", - user_id, - ) + "SLACK_BOT_TOKEN.", team_key or "this workspace", user_id) except Exception: # pragma: no cover - diagnostics must never break connect pass - def _load_bot_tokens(self, raw_token: str) -> List[str]: - """Comma-separated SLACK_BOT_TOKEN plus saved OAuth tokens (slack_tokens.json).""" - return _load_slack_bot_tokens(raw_token, quiet=False) - - @staticmethod - def _bolt_event_listener(handler): - """Wrap ``handler(event, body)`` as a Bolt listener taking ``(event, say, body)``.""" - - async def _listener(event, say, body): - await handler(event, body) - - return _listener - def _register_bolt_handlers(self) -> None: """Wire every Bolt listener onto ``self._app``; must run before Socket Mode starts.""" - import re as _re - # Bolt injects listener args by NAME and passes None for unknown ones, so every # handler takes (event, say, body). When Slack fires BOTH message and # app_mention they share an event ts and the deduplicator drops the second. @@ -1517,20 +1443,22 @@ class SlackAdapter(BasePlatformAdapter): return _handler + def _listener_for(handler): + async def _listener(event, say, body): + await handler(event, body) + + return _listener + for event_type, handler in ( - ("message", self._handle_slack_message), - ("app_mention", self._handle_slack_message), + ("message", self._handle_slack_message), ("app_mention", self._handle_slack_message), ("app_home_opened", self._handle_app_home_opened), ("app_context_changed", self._handle_app_context_changed), - ("file_shared", self._handle_slack_file_shared), - ("file_created", _noop), - ("file_change", _noop), - ("reaction_added", _reaction(False)), + ("file_shared", self._handle_slack_file_shared), ("file_created", _noop), + ("file_change", _noop), ("reaction_added", _reaction(False)), ("reaction_removed", _reaction(True)), ("assistant_thread_started", self._handle_assistant_thread_lifecycle_event), - ("assistant_thread_context_changed", self._handle_assistant_thread_lifecycle_event), - ): - self._app.event(event_type)(self._bolt_event_listener(handler)) + ("assistant_thread_context_changed", self._handle_assistant_thread_lifecycle_event)): + self._app.event(event_type)(_listener_for(handler)) # Catch-all ack: unacked envelopes count as failures and past 95%/60-min Slack # auto-disables Event Subscriptions (killing ALL inbound). Must be registered # AFTER every named handler — bolt dispatches to the first match. @@ -1540,8 +1468,7 @@ class SlackAdapter(BasePlatformAdapter): "[Slack] Ignoring unhandled event type=%s (no listener registered; subscribed " "events not handled by Hermes can be removed from the Slack app manifest via " "`hermes slack manifest`)", - (event or {}).get("type", (body or {}).get("event", {}).get("type", "unknown")), - ) + (event or {}).get("type", (body or {}).get("event", {}).get("type", "unknown"))) # Every COMMAND_REGISTRY command is a native slash (no /hermes prefix), # dispatched via one regex matcher. Slash commands must ALSO be declared @@ -1551,11 +1478,10 @@ class SlackAdapter(BasePlatformAdapter): _slash_names = [name for name, _d, _h in slack_native_slashes()] if _slash_names: - _slash_pattern = _re.compile( - r"^/(?:" + "|".join(_re.escape(n) for n in _slash_names) + r")$" - ) + _slash_pattern = re.compile( + r"^/(?:" + "|".join(re.escape(n) for n in _slash_names) + r")$") else: # pragma: no cover - registry always non-empty - _slash_pattern = _re.compile(r"^/hermes$") + _slash_pattern = re.compile(r"^/hermes$") @self._app.command(_slash_pattern) async def handle_hermes_command(ack, command): @@ -1574,7 +1500,7 @@ class SlackAdapter(BasePlatformAdapter): self._app.action("hermes_feedback")(self._handle_feedback_action) # Clarify buttons (tools/clarify_gateway.py); indexed action IDs because # Block Kit requires unique IDs within an actions block. - self._app.action(_re.compile(r"^hermes_clarify_choice_\d+$"))(self._handle_clarify_action) + self._app.action(re.compile(r"^hermes_clarify_choice_\d+$"))(self._handle_clarify_action) self._app.action("hermes_clarify_other")(self._handle_clarify_action) # Plugin action handlers (``ctx.register_slack_action_handler``), wired # before Socket Mode dispatches. Each is wrapped so a plugin exception @@ -1595,11 +1521,8 @@ class SlackAdapter(BasePlatformAdapter): await cb(ack, body, action) except Exception as exc: # pragma: no cover - defensive logger.error( - "[Slack] Plugin '%s' action handler raised: %s", - plugin_name, - exc, - exc_info=True, - ) + "[Slack] Plugin '%s' action handler raised: %s", plugin_name, exc, + exc_info=True) # Best-effort ack so Slack doesn't retry the click. try: await ack() @@ -1611,8 +1534,7 @@ class SlackAdapter(BasePlatformAdapter): for _action_id, _cb, _plugin_name in _plugin_handlers: self._app.action(_action_id)(_make_wrapper(_cb, _plugin_name)) logger.debug( - "[Slack] Registered plugin action handler %s (from %s)", _action_id, _plugin_name - ) + "[Slack] Registered plugin action handler %s (from %s)", _action_id, _plugin_name) if _plugin_handlers: logger.info("[Slack] Wired %d plugin action handler(s)", len(_plugin_handlers)) # ctx.register_platform_handler("slack", ...) factories get the full @@ -1642,8 +1564,7 @@ class SlackAdapter(BasePlatformAdapter): if self._bot_display_name is None: self._bot_display_name = bot_name logger.info( - "[Slack] Authenticated as @%s in workspace %s (team: %s)", bot_name, team_name, team_id - ) + "[Slack] Authenticated as @%s in workspace %s (team: %s)", bot_name, team_name, team_id) self._warn_if_missing_group_dm_scopes(auth_response, team_name) self._warn_if_not_bot_token(auth_response, team_name) self._warn_if_inchannel_without_flat_reply(team_name) @@ -1667,21 +1588,17 @@ class SlackAdapter(BasePlatformAdapter): logger.error( "[Slack] %s not set — this is a permanent config error; set %s via `hermes " "gateway setup` or in the active profile's ~/.hermes/.env file, then restart the " - "gateway.", - env_name, - env_name, - ) + "gateway.", env_name, env_name) self._set_fatal_error( f"missing_{env_name.lower()}", f"{env_name} not configured. Use `hermes gateway setup` " "or add it to your active profile's ~/.hermes/.env file, then restart the gateway.", - retryable=False, - ) + retryable=False) return False proxy_url = _resolve_slack_proxy_url() if proxy_url: logger.info("[Slack] Using proxy for Slack transport: %s", safe_url_for_log(proxy_url)) - bot_tokens = self._load_bot_tokens(raw_token) + bot_tokens = _load_slack_bot_tokens(raw_token, quiet=False) lock_acquired = False try: if not self._acquire_platform_lock("slack-app-token", app_token, "Slack app token"): @@ -1701,10 +1618,8 @@ class SlackAdapter(BasePlatformAdapter): self._bot_user_id = self._bot_display_name = None self._team_clients, self._team_bot_user_ids, self._team_bot_names = {}, {}, {} self._app = AsyncApp( - token=bot_tokens[0], - client=self._new_web_client(bot_tokens[0], proxy_url), - before_authorize=_slack_per_request_proxy_middleware(proxy_url), - ) + token=bot_tokens[0], client=self._new_web_client(bot_tokens[0], proxy_url), + before_authorize=_slack_per_request_proxy_middleware(proxy_url)) _apply_slack_proxy(self._app.client, proxy_url) for token in bot_tokens: await self._authenticate_workspace(token, proxy_url) @@ -1733,9 +1648,7 @@ class SlackAdapter(BasePlatformAdapter): "appropriate (run 'hermes slack manifest' if unsure), and (b) the other bot's " "Slack user id is in SLACK_ALLOWED_USERS or GATEWAY_ALLOW_ALL_USERS=true. " "Without these, bot events are silently dropped upstream of the allow_bots " - "gate.", - _allow_bots_cfg, - ) + "gate.", _allow_bots_cfg) return True except Exception as e: # pragma: no cover - defensive logging logger.error("[Slack] Connection failed: %s", e, exc_info=True) @@ -1760,11 +1673,8 @@ class SlackAdapter(BasePlatformAdapter): return str(ts) if ts else None except Exception as exc: logger.warning( - "[%s] Handoff thread: seed-post failed for channel %s: %s", - self.name, - parent_chat_id, - exc, - ) + "[%s] Handoff thread: seed-post failed for channel %s: %s", self.name, + parent_chat_id, exc) return None async def disconnect(self) -> None: @@ -1794,8 +1704,7 @@ class SlackAdapter(BasePlatformAdapter): if not metadata: return "" found = _first_truthy( - metadata, ("scope_id", "slack_team_id", "team_id", "team", "guild_id", "workspace_id") - ) + metadata, ("scope_id", "slack_team_id", "team_id", "team", "guild_id", "workspace_id")) if found: return str(found) source = metadata.get("source") @@ -1864,23 +1773,18 @@ class SlackAdapter(BasePlatformAdapter): if dm_id: self._dm_conversation_cache[cache_key] = dm_id self._trim_oldest_dict_entries( - self._dm_conversation_cache, self._DM_CONVERSATION_CACHE_MAX - ) + self._dm_conversation_cache, self._DM_CONVERSATION_CACHE_MAX) if team_id: self._remember_channel_team(dm_id, team_id) return dm_id except Exception as e: logger.warning( "[Slack] conversations.open failed for user target %s: %s " - "(check the bot's im:write scope)", - cid, - e, - ) + "(check the bot's im:write scope)", cid, e) return chat_id async def _clear_thread_status_quietly( - self, chat_id: str, metadata: Optional[Dict[str, Any]] = None - ) -> None: + self, chat_id: str, metadata: Optional[Dict[str, Any]] = None) -> None: """Best-effort assistant-status clear for send() paths that skip the normal clear. Otherwise the thread stays stuck "is thinking..." on empty responses, ephemeral slash replies, or exceptions before ``thread_ts`` resolved. ``stop_typing`` is idempotent; errors @@ -1928,15 +1832,9 @@ class SlackAdapter(BasePlatformAdapter): 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: + 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") @@ -1959,13 +1857,11 @@ class SlackAdapter(BasePlatformAdapter): start_payload: Dict[str, Any] = { "channel": chat_id, "thread_ts": stream.thread_ts, - "task_display_mode": "plan", - } + "task_display_mode": "plan"} md = metadata or {} recipients = ( ("recipient_team_id", ("recipient_team_id", "team_id", "slack_team_id")), - ("recipient_user_id", ("recipient_user_id", "user_id")), - ) + ("recipient_user_id", ("recipient_user_id", "user_id"))) for key, sources in recipients: value = _first_truthy(md, sources) if value: @@ -1978,10 +1874,7 @@ class SlackAdapter(BasePlatformAdapter): chunks: List[Dict[str, Any]] = [{"type": "plan_update", "title": str(title)[:256]}] chunks.extend(self._task_update_chunk(task) for task in tasks) append_payload: Dict[str, Any] = { - "channel": chat_id, - "ts": stream.stream_ts, - "chunks": chunks, - } + "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) @@ -2001,12 +1894,10 @@ class SlackAdapter(BasePlatformAdapter): "type": "task_update", "id": task_id, "title": str(task.get("title") or task_id)[:256], - "status": status, - } + "status": status} async def _stop_native_task_card_stream( - self, key: Tuple[str, str, str], stream: _NativeTaskCardStream - ) -> None: + self, key: Tuple[str, str, str], stream: _NativeTaskCardStream) -> None: async with stream.lock: if stream.stopped: return @@ -2014,8 +1905,7 @@ class SlackAdapter(BasePlatformAdapter): 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} - ) + "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: @@ -2023,12 +1913,8 @@ class SlackAdapter(BasePlatformAdapter): 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: + 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: @@ -2053,8 +1939,7 @@ class SlackAdapter(BasePlatformAdapter): return None async def _call_with_block_fallback( - self, client_fn: Callable[[], Any], method: str, kwargs: Dict[str, Any], verb: str - ) -> Any: + self, client_fn: Callable[[], Any], method: str, kwargs: Dict[str, Any], verb: str) -> Any: """``client_fn().(**kwargs)``; on a Block Kit rejection retry once without ``blocks`` (an edit sends ``blocks=[]`` so the message drops its stale layout). The client is re-resolved for the retry.""" @@ -2068,18 +1953,13 @@ class SlackAdapter(BasePlatformAdapter): else: retry_kwargs.pop("blocks", None) logger.info( - "[Slack] Block Kit payload rejected; retrying %s without blocks: %s", verb, e - ) + "[Slack] Block Kit payload rejected; retrying %s without blocks: %s", verb, e) return await getattr(client_fn(), method)(**retry_kwargs) raise async def send( - self, - chat_id: str, - content: str, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: + self, chat_id: str, content: str, reply_to: Optional[str] = None, + metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Send a message to a Slack channel or DM.""" blocked = self._outbound_blocked(chat_id, "outbound generic send to") if blocked: @@ -2115,9 +1995,7 @@ class SlackAdapter(BasePlatformAdapter): kwargs = { "channel": chat_id, "text": chunk, - "mrkdwn": True, - **_slack_unfurl_kwargs(self.config.extra), - } + "mrkdwn": True, **_slack_unfurl_kwargs(self.config.extra)} if blocks and i == 0: kwargs["blocks"] = blocks if thread_ts: @@ -2127,8 +2005,7 @@ class SlackAdapter(BasePlatformAdapter): kwargs["reply_broadcast"] = True client_fn = lambda: self._get_client(chat_id, team_id=team_id) # noqa: E731 last_result = await self._call_with_block_fallback( - client_fn, "chat_postMessage", kwargs, "send" - ) + client_fn, "chat_postMessage", kwargs, "send") # Clear Slack Assistant status as soon as the final message is posted. if thread_ts: await self.stop_typing(chat_id, metadata=metadata) @@ -2148,11 +2025,8 @@ class SlackAdapter(BasePlatformAdapter): logger.error("[Slack] Send error: %s", e, exc_info=True) _retryable = self._is_retryable_upload_error(e) return SendResult( - success=False, - error=str(e), - retryable=_retryable, - retry_after=self._retry_after_from_exc(e) if _retryable else None, - ) + success=False, error=str(e), retryable=_retryable, + retry_after=self._retry_after_from_exc(e) if _retryable else None) @staticmethod def _retry_after_from_exc(e: BaseException) -> Optional[float]: @@ -2167,12 +2041,8 @@ class SlackAdapter(BasePlatformAdapter): return None async def _send_slash_reply( - self, - chat_id: str, - slash_ctx: Dict[str, Any], - content: str, - metadata: Optional[Dict[str, Any]], - ) -> SendResult: + self, chat_id: str, slash_ctx: Dict[str, Any], content: str, + metadata: Optional[Dict[str, Any]]) -> SendResult: """Deliver a slash-command reply ephemerally, replacing the "Running /cmd…" ack. response_url first, then chat.postEphemeral; NEVER a public post — a reply the user expects to be private must not leak because a path failed. Ephemeral replies don't auto-clear the @@ -2183,8 +2053,7 @@ class SlackAdapter(BasePlatformAdapter): return ephemeral_result logger.warning( "[Slack] response_url slash reply failed (%s); retrying via chat.postEphemeral", - ephemeral_result.error, - ) + ephemeral_result.error) fallback_result = await self._post_ephemeral_fallback(chat_id, slash_ctx, content) if fallback_result.success: await self._clear_thread_status_quietly(chat_id, metadata) @@ -2192,19 +2061,12 @@ class SlackAdapter(BasePlatformAdapter): # The user still has the ack; the error is returned so the gateway can react. logger.error( "[Slack] Ephemeral slash reply failed on both response_url and chat.postEphemeral " - "(%s); dropping rather than posting publicly", - fallback_result.error, - ) + "(%s); dropping rather than posting publicly", fallback_result.error) return fallback_result async def send_private_notice( - self, - chat_id: str, - user_id: str, - content: str, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: + self, chat_id: str, user_id: str, content: str, reply_to: Optional[str] = None, + metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Send a Slack ephemeral message visible only to one user.""" blocked = self._outbound_blocked(chat_id, "outbound generic ephemeral notice to") if blocked: @@ -2219,22 +2081,15 @@ class SlackAdapter(BasePlatformAdapter): kwargs["thread_ts"] = thread_ts result = await self._client_for(chat_id, metadata).chat_postEphemeral(**kwargs) return SendResult( - success=True, - message_id=result.get("message_ts") or result.get("ts"), - raw_response=result, - ) + success=True, message_id=result.get("message_ts") or result.get("ts"), + raw_response=result) except Exception as e: # pragma: no cover - defensive logging logger.error("[Slack] Ephemeral send error: %s", e, exc_info=True) return SendResult(success=False, error=str(e)) async def send_or_update_status( - self, - chat_id: str, - status_key: str, - content: str, - *, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: + self, chat_id: str, status_key: str, content: str, *, + metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Send a status message, or edit the previous one with the same key. Keyed by (channel, thread, status_key) so repeated progress callbacks edit one bubble via ``chat.update`` instead of spamming the thread. If the edit fails (deleted, too old) the @@ -2244,8 +2099,7 @@ class SlackAdapter(BasePlatformAdapter): cached_id = self._status_message_ids.get(key) if cached_id is not None: result = await self.edit_message( - chat_id, cached_id, content, finalize=False, metadata=metadata - ) + chat_id, cached_id, content, finalize=False, metadata=metadata) if result.success: if result.message_id: self._status_message_ids[key] = str(result.message_id) @@ -2262,14 +2116,8 @@ class SlackAdapter(BasePlatformAdapter): return result async def edit_message( - self, - chat_id: str, - message_id: str, - content: str, - *, - finalize: bool = False, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: + self, chat_id: str, message_id: str, content: str, *, finalize: bool = False, + metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Edit a previously sent Slack message.""" blocked = self._outbound_blocked(chat_id, "message edit in") if blocked: @@ -2283,10 +2131,7 @@ class SlackAdapter(BasePlatformAdapter): chunks = self.truncate_message(formatted, self.MAX_MESSAGE_LENGTH) formatted = chunks[0] if chunks else formatted update_kwargs: Dict[str, Any] = { - "channel": chat_id, - "ts": message_id, - "text": formatted, - } + "channel": chat_id, "ts": message_id, "text": formatted} # Only render Block Kit on the FINAL edit. Intermediate streaming # edits stay plain mrkdwn — re-deriving a full block layout on every # progressive flush would be wasteful and jittery. ``text`` is kept @@ -2296,8 +2141,7 @@ class SlackAdapter(BasePlatformAdapter): if blocks: update_kwargs["blocks"] = blocks await self._call_with_block_fallback( - lambda: self._client_for(chat_id, metadata), "chat_update", update_kwargs, "edit" - ) + lambda: self._client_for(chat_id, metadata), "chat_update", update_kwargs, "edit") if finalize: await self._clear_thread_status_quietly(chat_id, metadata) return SendResult(success=True, message_id=message_id) @@ -2310,21 +2154,12 @@ class SlackAdapter(BasePlatformAdapter): # failure as permanent makes every later tool update a new post. logger.error( "[Slack] transient chat.update failure on message %s in channel %s: %s", - message_id, - chat_id, - e, - exc_info=True, - ) + message_id, chat_id, e, exc_info=True) return SendResult( - success=False, error=str(e), retryable=True, error_kind="transient" - ) + success=False, error=str(e), retryable=True, error_kind="transient") logger.error( - "[Slack] Failed to edit message %s in channel %s: %s", - message_id, - chat_id, - e, - exc_info=True, - ) + "[Slack] Failed to edit message %s in channel %s: %s", message_id, chat_id, e, + exc_info=True) return SendResult(success=False, error=str(e)) async def delete_message(self, chat_id: str, message_id: str) -> bool: @@ -2336,16 +2171,12 @@ class SlackAdapter(BasePlatformAdapter): if hasattr(response, "get") and response.get("ok") is False: logger.debug( "[Slack] chat.delete returned ok=false for message %s in channel %s: %s", - message_id, - chat_id, - response.get("error", "unknown"), - ) + message_id, chat_id, response.get("error", "unknown")) return False return True except Exception as e: # pragma: no cover - best-effort cleanup logger.debug( - "[Slack] Failed to delete message %s in channel %s: %s", message_id, chat_id, e - ) + "[Slack] Failed to delete message %s in channel %s: %s", message_id, chat_id, e) return False # ── Native streaming (chat.startStream / appendStream / stopStream) ── @@ -2359,12 +2190,10 @@ class SlackAdapter(BasePlatformAdapter): _STREAM_CURSOR_GLYPHS = ("\u2589", "▍", "▌", "…") _NATIVE_STREAM_UNSUPPORTED_MARKERS = ( "not_allowed", "missing_scope", "feature_not_enabled", "invalid_method", "unknown_method", - "method_deprecated", "not_authed", "streaming_not_allowed", - ) + "method_deprecated", "not_authed", "streaming_not_allowed") def supports_draft_streaming( - self, chat_type: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None - ) -> bool: + self, chat_type: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None) -> bool: """Return whether Slack's native stream can preserve configured behavior.""" if self._native_stream_unsupported: return False @@ -2421,11 +2250,7 @@ class SlackAdapter(BasePlatformAdapter): if not ts: raise RuntimeError("chat.startStream returned no ts") self._active_streams[chat_id] = { - "ts": str(ts), - "draft_id": draft_id, - "sent": text, - "started": time.time(), - } + "ts": str(ts), "draft_id": draft_id, "sent": text, "started": time.time()} self._bot_message_ts.add(str(ts)) return SendResult(success=True, message_id=str(ts)) sent = stream.get("sent", "") @@ -2451,20 +2276,14 @@ class SlackAdapter(BasePlatformAdapter): logger.warning( "[Slack] Native streaming unavailable (%s). Falling back to edit-based " "streaming. To enable native streaming, turn on the Agents & AI Apps feature " - "for this Slack app (and ensure the assistant:write scope).", - err, - ) + "for this Slack app (and ensure the assistant:write scope).", err) else: logger.debug("[Slack] Native stream frame failed: %s", err) return SendResult(success=False, error=err) async def _seal_stream( - self, - chat_id: str, - stream: Dict[str, Any], - final_text: Optional[str] = None, - blocks: Optional[list] = None, - ) -> bool: + self, chat_id: str, stream: Dict[str, Any], final_text: Optional[str] = None, + blocks: Optional[list] = None) -> bool: """Best-effort chat.stopStream for an open stream. ``final_text`` is the complete final content; only the unsent delta is passed to stopStream (append-only API). Returns True on success.""" @@ -2480,8 +2299,7 @@ class SlackAdapter(BasePlatformAdapter): return True except Exception as e: # pragma: no cover - defensive logger.debug( - "[Slack] chat.stopStream failed for %s/%s: %s", chat_id, stream.get("ts"), e - ) + "[Slack] chat.stopStream failed for %s/%s: %s", chat_id, stream.get("ts"), e) return False async def _try_finalize_stream(self, chat_id: str, content: str) -> Optional[SendResult]: @@ -2510,12 +2328,10 @@ class SlackAdapter(BasePlatformAdapter): if blocks: try: await self._get_client(chat_id).chat_update( - channel=chat_id, ts=ts, text=self.format_message(text), blocks=blocks - ) + channel=chat_id, ts=ts, text=self.format_message(text), blocks=blocks) except Exception as e: logger.debug( - "[Slack] Post-stream Block Kit update failed (markdown fallback stands): %s", e - ) + "[Slack] Post-stream Block Kit update failed (markdown fallback stands): %s", e) await self.stop_typing(chat_id) return SendResult(success=True, message_id=ts) @@ -2532,8 +2348,7 @@ class SlackAdapter(BasePlatformAdapter): # metadata.thread_id is the message's own ts, and setStatus on it # would open an assistant thread before the reply is sent. thread_ts = self._resolve_thread_ts( - reply_to=metadata.get("message_id"), metadata=metadata - ) + reply_to=metadata.get("message_id"), metadata=metadata) if not thread_ts: return # Can only set status in a thread context team_id = self._metadata_team_id(metadata) or self._channel_team.get(chat_id, "") @@ -2551,24 +2366,24 @@ class SlackAdapter(BasePlatformAdapter): self._active_status_threads[status_key] = { "thread_ts": str(thread_ts), "team_id": str(team_id) if team_id else "", - "started": _status_started, - } + "started": _status_started} # Evict oldest-thread-first (key[2] is the thread ts) so the newest survives. self._evict_oldest_by_ts( - self._active_status_threads, self._ACTIVE_STATUS_THREADS_MAX, lambda k: k[2] - ) + self._active_status_threads, self._ACTIVE_STATUS_THREADS_MAX, lambda k: k[2]) + # May lack assistant:write scope or assistant context; reactions still work. + _status = getattr(self, "_status_text", {}).get(str(chat_id)) or getattr( + self.config, "typing_status_text", None) + _status = _status or self._default_status_text(_status_started) + await self._set_thread_status(chat_id, team_id, thread_ts, _status, "failed") + + async def _set_thread_status( + self, chat_id: str, team_id: str, thread_ts: str, status: str, fail_label: str) -> None: + """``assistant.threads.setStatus`` (empty ``status`` clears); failures are debug-logged.""" try: - _status = getattr(self, "_status_text", {}).get(str(chat_id)) or getattr( - self.config, "typing_status_text", None - ) - if not _status: - _status = self._default_status_text(_status_started) await self._get_client(chat_id, team_id=team_id).assistant_threads_setStatus( - channel_id=chat_id, thread_ts=thread_ts, status=_status - ) + channel_id=chat_id, thread_ts=thread_ts, status=status) except Exception as e: - # May lack assistant:write scope or assistant context; reactions still work. - logger.debug("[Slack] assistant.threads.setStatus failed: %s", e) + logger.debug("[Slack] assistant.threads.setStatus %s: %s", fail_label, e) @staticmethod def _default_status_text(started: Optional[float]) -> str: @@ -2604,8 +2419,7 @@ class SlackAdapter(BasePlatformAdapter): key for key in self._active_status_threads if key[1] == str(chat_id) - and (not requested_thread_ts or key[2] == requested_thread_ts) - ] + and (not requested_thread_ts or key[2] == requested_thread_ts)] if len(matching_keys) == 1: active = self._active_status_threads.pop(matching_keys[0], None) ambiguous_tracked = bool(requested_thread_ts) and len(matching_keys) > 1 @@ -2619,12 +2433,7 @@ class SlackAdapter(BasePlatformAdapter): thread_ts = requested_thread_ts if not thread_ts: return - try: - await self._get_client(chat_id, team_id=team_id).assistant_threads_setStatus( - channel_id=chat_id, thread_ts=thread_ts, status="" - ) - except Exception as e: - logger.debug("[Slack] assistant.threads.setStatus clear failed: %s", e) + await self._set_thread_status(chat_id, team_id, thread_ts, "", "clear failed") def _dm_top_level_threads_as_sessions(self) -> bool: """Each top-level DM reply thread is its own session (default True; set @@ -2650,17 +2459,14 @@ class SlackAdapter(BasePlatformAdapter): session.""" try: if self._cron_continuable_surface() == "in_channel" and self.config.extra.get( - "reply_in_thread", True - ): + "reply_in_thread", True): logger.warning( "[Slack] %s: cron_continuable_surface=in_channel is set WITHOUT " "reply_in_thread=false. A continuable in-channel cron job will deliver flat, " "but the bot will still reply to your continuation in a thread — so it falls " "back to a threaded continuation (\u2248 default behaviour), not the flat " "channel session you asked for. Set platforms.slack.extra.reply_in_thread: " - "false to pair them.", - team_name, - ) + "false to pair them.", team_name) except Exception: pass @@ -2685,8 +2491,7 @@ class SlackAdapter(BasePlatformAdapter): raw = os.getenv("SLACK_API_HUMAN_USERS", "") parts = raw if isinstance(raw, (list, tuple, set)) else str(raw).split(",") cached = self._api_human_users_cache = frozenset( - str(p).strip() for p in parts if str(p).strip() - ) + str(p).strip() for p in parts if str(p).strip()) return cached def _event_declares_bot_sender(self, event: dict) -> bool: @@ -2728,30 +2533,17 @@ class SlackAdapter(BasePlatformAdapter): return reply_to async def _upload_with_retry( - self, - chat_id: str, - file_path: Optional[str], - filename: str, - caption: Optional[str], - thread_ts: Optional[str], - metadata: Optional[Dict[str, Any]], - label: str = "Upload", - *, - content: Optional[bytes] = None, - attempts: int = 3, - ) -> SendResult: + self, chat_id: str, file_path: Optional[str], filename: str, caption: Optional[str], + thread_ts: Optional[str], metadata: Optional[Dict[str, Any]], label: str = "Upload", *, + content: Optional[bytes] = None, attempts: int = 3) -> SendResult: """``files_upload_v2`` of a local path (or in-memory ``content``) with up to ``attempts`` tries on transient errors; re-raises otherwise.""" source = {"file": file_path} if content is None else {"content": content} for attempt in range(attempts): try: result = await self._client_for(chat_id, metadata).files_upload_v2( - channel=chat_id, - **source, - filename=filename, - initial_comment=caption or "", - thread_ts=thread_ts, - ) + channel=chat_id, **source, filename=filename, initial_comment=caption or "", + thread_ts=thread_ts) self._record_uploaded_file_thread(chat_id, thread_ts, metadata) return SendResult(success=True, raw_response=result) except Exception as exc: @@ -2761,13 +2553,8 @@ class SlackAdapter(BasePlatformAdapter): await asyncio.sleep(1.5 * (attempt + 1)) async def _upload_file( - self, - chat_id: str, - file_path: str, - caption: Optional[str] = None, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: + self, chat_id: str, file_path: str, caption: Optional[str] = None, + reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Upload a local file to Slack (raises FileNotFoundError when missing).""" blocked = self._outbound_blocked(chat_id, "file upload in") if blocked: @@ -2777,21 +2564,12 @@ class SlackAdapter(BasePlatformAdapter): chat_id = await self._dm_target(chat_id, metadata) thread_ts = self._resolve_thread_ts(reply_to, metadata) return await self._upload_with_retry( - chat_id, file_path, os.path.basename(file_path), caption, thread_ts, metadata - ) + chat_id, file_path, os.path.basename(file_path), caption, thread_ts, metadata) async def _send_local_file( - self, - chat_id: str, - file_path: str, - caption: Optional[str], - reply_to: Optional[str], - metadata: Optional[Dict[str, Any]], - kind: str, - filename: str, - not_found_error: str, - failure_notice: str, - ) -> SendResult: + self, chat_id: str, file_path: str, caption: Optional[str], reply_to: Optional[str], + metadata: Optional[Dict[str, Any]], kind: str, filename: str, not_found_error: str, + failure_notice: str) -> SendResult: """Shared body of ``send_video``/``send_document``: upload with retry, text notice on failure. The host-local path is never echoed into chat; ``failure_notice`` is posted instead.""" if not self._app: @@ -2803,22 +2581,16 @@ class SlackAdapter(BasePlatformAdapter): thread_ts = self._resolve_thread_ts(reply_to, metadata) label = f"{kind.capitalize()} upload" return await self._upload_with_retry( - chat_id, file_path, filename, caption, thread_ts, metadata, label - ) + chat_id, file_path, filename, caption, thread_ts, metadata, label) except Exception as e: # pragma: no cover - defensive logging logger.error( - "[%s] Failed to send %s %s: %s", self.name, kind, file_path, e, exc_info=True - ) + "[%s] Failed to send %s %s: %s", self.name, kind, file_path, e, exc_info=True) text = f"{caption}\n{failure_notice}" if caption else failure_notice return await self.send(chat_id, text, reply_to=reply_to, metadata=metadata) async def send_multiple_images( - self, - chat_id: str, - images: List[Tuple[str, str]], - metadata: Optional[Dict[str, Any]] = None, - human_delay: float = 0.0, - ) -> None: + self, chat_id: str, images: List[Tuple[str, str]], + metadata: Optional[Dict[str, Any]] = None, human_delay: float = 0.0) -> None: """Send a batch of images as a single Slack message with multiple file uploads. Uses ``files_upload_v2`` with its ``file_uploads`` parameter so all images show up attached to one ``initial_comment`` message instead of N separate messages. Falls back to the base @@ -2845,35 +2617,23 @@ class SlackAdapter(BasePlatformAdapter): await asyncio.sleep(human_delay) try: file_uploads, initial_comment_parts = await self._collect_image_uploads( - chunk, _unquote, _is_safe_url, create_ssrf_safe_async_client - ) + chunk, _unquote, _is_safe_url, create_ssrf_safe_async_client) if not file_uploads: continue initial_comment = "\n".join(initial_comment_parts) if initial_comment_parts else "" logger.info( "[Slack] Sending %d image(s) in single files_upload_v2 (chunk %d/%d)", - len(file_uploads), - chunk_idx + 1, - len(chunks), - ) + len(file_uploads), chunk_idx + 1, len(chunks)) await self._client_for(chat_id, metadata).files_upload_v2( - channel=chat_id, - file_uploads=file_uploads, - initial_comment=initial_comment, - thread_ts=thread_ts, - ) + channel=chat_id, file_uploads=file_uploads, initial_comment=initial_comment, + thread_ts=thread_ts) self._record_uploaded_file_thread(chat_id, thread_ts, metadata) except Exception as e: logger.warning( "[Slack] Multi-image files_upload_v2 failed (chunk %d/%d), falling back to per-image: %s", - chunk_idx + 1, - len(chunks), - e, - exc_info=True, - ) + chunk_idx + 1, len(chunks), e, exc_info=True) await super().send_multiple_images( - chat_id, chunk, metadata, human_delay=human_delay - ) + chat_id, chunk, metadata, human_delay=human_delay) @staticmethod async def _collect_image_uploads( @@ -2896,8 +2656,7 @@ class SlackAdapter(BasePlatformAdapter): logger.warning("[Slack] Skipping missing image: %s", local_path) continue file_uploads.append( - {"file": local_path, "filename": os.path.basename(local_path)} - ) + {"file": local_path, "filename": os.path.basename(local_path)}) continue if not is_safe_url_fn(image_url): logger.warning("[Slack] Blocked unsafe image URL in batch") @@ -2908,13 +2667,11 @@ class SlackAdapter(BasePlatformAdapter): ct = response.headers.get("content-type", "") ext = next((e for k, e in _IMAGE_CT_EXTS if k in ct), "png") file_uploads.append({ - "content": response.content, - "filename": f"image_{len(file_uploads)}.{ext}", + "content": response.content, "filename": f"image_{len(file_uploads)}.{ext}" }) except Exception as dl_err: logger.warning( - "[Slack] Download failed for %s: %s", safe_url_for_log(image_url), dl_err - ) + "[Slack] Download failed for %s: %s", safe_url_for_log(image_url), dl_err) return file_uploads, initial_comment_parts def _record_uploaded_file_thread( @@ -2935,8 +2692,7 @@ class SlackAdapter(BasePlatformAdapter): body = " ".join( str(part) for part in (exc, getattr(exc, "message", ""), getattr(exc, "response", None)) - if part - ).lower() + if part).lower() if any(m in body for m in _TRANSIENT_UPLOAD_MARKERS): return True return self._is_retryable_error(body) @@ -2986,16 +2742,11 @@ class SlackAdapter(BasePlatformAdapter): "positive_button": { "text": {"type": "plain_text", "text": "Good Response"}, "accessibility_label": ("Submit positive feedback on this response"), - "value": "positive", - }, + "value": "positive"}, "negative_button": { "text": {"type": "plain_text", "text": "Bad Response"}, "accessibility_label": ("Submit negative feedback on this response"), - "value": "negative", - }, - } - ], - } + "value": "negative"}}]} def _append_feedback_block(self, blocks: Optional[list]) -> Optional[list]: """Append response feedback controls when enabled and block budget allows.""" @@ -3062,19 +2813,16 @@ class SlackAdapter(BasePlatformAdapter): return _ph(f"<{url}|{label}>") text = re.sub( - r"(?\n]+>)", lambda m: _ph(m.group(1)), text - ) + r"(<(?:[@#!]|(?:https?|mailto|tel):)[^>\n]+>)", lambda m: _ph(m.group(1)), text) # 5) Protect blockquote markers text = re.sub(r"^(>+\s)", lambda m: _ph(m.group(0)), text, flags=re.MULTILINE) # 6) Escape control chars. Unescape first in ONE regex pass (sequential # replaces would decode "&lt;" twice) to avoid double-escaping. text = re.sub( - r"&(amp|lt|gt);", lambda m: {"amp": "&", "lt": "<", "gt": ">"}[m.group(1)], text - ) + r"&(amp|lt|gt);", lambda m: {"amp": "&", "lt": "<", "gt": ">"}[m.group(1)], text) text = text.replace("&", "&").replace("<", "<").replace(">", ">") # 7) Headers → *bold* def _convert_header(m): @@ -3097,8 +2845,7 @@ class SlackAdapter(BasePlatformAdapter): # 10) *text* → _text_ only when non-whitespace touches both delimiters # ("a * b * c" is preserved). text = re.sub( - r"(? bool: + self, channel: str, timestamp: str, emoji: str, team_id: str, *, remove: bool) -> bool: """reactions.add / reactions.remove; True on success. Failures (already reacted, missing scope) are debug-logged only.""" if not self._app: @@ -3122,18 +2868,15 @@ class SlackAdapter(BasePlatformAdapter): return True except Exception as e: logger.debug( - "[Slack] reactions.%s failed (%s): %s", "remove" if remove else "add", emoji, e - ) + "[Slack] reactions.%s failed (%s): %s", "remove" if remove else "add", emoji, e) return False async def _add_reaction( - self, channel: str, timestamp: str, emoji: str, team_id: str = "" - ) -> bool: + self, channel: str, timestamp: str, emoji: str, team_id: str = "") -> bool: return await self._react(channel, timestamp, emoji, team_id, remove=False) async def _remove_reaction( - self, channel: str, timestamp: str, emoji: str, team_id: str = "" - ) -> bool: + self, channel: str, timestamp: str, emoji: str, team_id: str = "") -> bool: return await self._react(channel, timestamp, emoji, team_id, remove=True) def _reactions_enabled(self) -> bool: @@ -3220,8 +2963,7 @@ class SlackAdapter(BasePlatformAdapter): return channel_id try: resp = await self._get_client(channel_id, team_id=team_id or None).conversations_info( - channel=channel_id - ) + channel=channel_id) payload = _slack_response_payload(resp) ch = payload.get("channel") or {} if not payload.get("ok"): @@ -3231,8 +2973,7 @@ class SlackAdapter(BasePlatformAdapter): name = ( await self._resolve_user_name(peer_user, chat_id=channel_id, team_id=team_id) if peer_user - else channel_id - ) + else channel_id) else: name = ch.get("name") or ch.get("name_normalized") or channel_id except Exception as e: @@ -3265,8 +3006,7 @@ class SlackAdapter(BasePlatformAdapter): unaffected). Gives the agent a "that's me" anchor to distinguish its own mention from similarly named participants.""" name = ( - (team_id and self._team_bot_names.get(team_id)) or self._bot_display_name or "" - ).strip() + (team_id and self._team_bot_names.get(team_id)) or self._bot_display_name or "").strip() if not name: return "" return ( @@ -3278,12 +3018,10 @@ class SlackAdapter(BasePlatformAdapter): f'because "@{name}" is absent. In messages, each line is prefixed ' f"with the sender's name, and visible mentions are shown as " f"@DisplayName; a mention of any other participant is not a " - f"mention of you, even if their name is similar." - ) + f"mention of you, even if their name is similar.") async def _resolve_user_is_bot( - self, user_id: str, chat_id: str = "", team_id: str = "" - ) -> bool: + self, user_id: str, chat_id: str = "", team_id: str = "") -> bool: """Resolve whether a Slack user ID is a bot account, with caching. Workspace-scoped like :meth:`_resolve_user_name` — Slack user IDs are team-local, so the cache key includes the team.""" @@ -3327,25 +3065,18 @@ class SlackAdapter(BasePlatformAdapter): is_bot = bool( user.get("is_bot") or user.get("is_workflow_bot") - or (isinstance(profile, dict) and profile.get("bot_id")) - ) + or (isinstance(profile, dict) and profile.get("bot_id"))) name = ( profile.get("display_name") or profile.get("real_name") or user.get("real_name") or user.get("name") - or user_id - ) + or user_id) return name, is_bot async def send_image_file( - self, - chat_id: str, - image_path: str, - caption: Optional[str] = None, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: + self, chat_id: str, image_path: str, caption: Optional[str] = None, + reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Send a local image file to Slack by uploading it.""" try: return await self._upload_file(chat_id, image_path, caption, reply_to, metadata) @@ -3361,13 +3092,8 @@ class SlackAdapter(BasePlatformAdapter): return await self.send(chat_id, text, reply_to=reply_to, metadata=metadata) async def send_image( - self, - chat_id: str, - image_url: str, - caption: Optional[str] = None, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: + self, chat_id: str, image_url: str, caption: Optional[str] = None, + reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Send an image to Slack by uploading the URL as a file.""" if not self._app: return SendResult(success=False, error="Not connected") @@ -3376,8 +3102,7 @@ class SlackAdapter(BasePlatformAdapter): if not is_safe_url(image_url): logger.warning("[Slack] Blocked unsafe image URL (SSRF protection)") return await super().send_image( - chat_id, image_url, caption, reply_to, metadata=metadata - ) + chat_id, image_url, caption, reply_to, metadata=metadata) try: async def _ssrf_redirect_guard(response): @@ -3390,8 +3115,7 @@ class SlackAdapter(BasePlatformAdapter): # Download the image first async with create_ssrf_safe_async_client( - timeout=30.0, - follow_redirects=True, + timeout=30.0, follow_redirects=True, event_hooks={"response": [_ssrf_redirect_guard]}, ) as client: response = await client.get(image_url) @@ -3399,36 +3123,20 @@ class SlackAdapter(BasePlatformAdapter): thread_ts = self._resolve_thread_ts(reply_to, metadata) chat_id = await self._dm_target(chat_id, metadata) return await self._upload_with_retry( - chat_id, - None, - "image.png", - caption, - thread_ts, - metadata, - content=response.content, - attempts=1, - ) + chat_id, None, "image.png", caption, thread_ts, metadata, content=response.content, + attempts=1) except Exception as e: # pragma: no cover - defensive logging logger.warning( "[Slack] Failed to upload image from URL %s, falling back to text: %s", - safe_url_for_log(image_url), - e, - exc_info=True, - ) + safe_url_for_log(image_url), e, exc_info=True) # Fall back to sending the URL as text text = f"{caption}\n{image_url}" if caption else image_url return await self.send( - chat_id=chat_id, content=text, reply_to=reply_to, metadata=metadata - ) + chat_id=chat_id, content=text, reply_to=reply_to, metadata=metadata) async def send_voice( - self, - chat_id: str, - audio_path: str, - caption: Optional[str] = None, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, - **kwargs, + self, chat_id: str, audio_path: str, caption: Optional[str] = None, + reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, **kwargs, ) -> SendResult: """Send an audio file to Slack.""" try: @@ -3440,28 +3148,17 @@ class SlackAdapter(BasePlatformAdapter): return SendResult(success=False, error=str(e)) async def send_video( - self, - chat_id: str, - video_path: str, - caption: Optional[str] = None, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: + self, chat_id: str, video_path: str, caption: Optional[str] = None, + reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Send a video file to Slack.""" return await self._send_local_file( chat_id, video_path, caption, reply_to, metadata, "video", os.path.basename(video_path), - f"Video file not found: {video_path}", "⚠️ Couldn't deliver the video attachment.", - ) + f"Video file not found: {video_path}", "⚠️ Couldn't deliver the video attachment.") async def send_document( - self, - chat_id: str, - file_path: str, - caption: Optional[str] = None, - file_name: Optional[str] = None, - reply_to: Optional[str] = None, - metadata: Optional[Dict[str, Any]] = None, - ) -> SendResult: + self, chat_id: str, file_path: str, caption: Optional[str] = None, + file_name: Optional[str] = None, reply_to: Optional[str] = None, + metadata: Optional[Dict[str, Any]] = None) -> SendResult: """Send a document/file attachment to Slack. Only ``display_name`` (never the host-local path) goes in the failure notice.""" display_name = file_name or os.path.basename(file_path) @@ -3487,8 +3184,7 @@ class SlackAdapter(BasePlatformAdapter): @staticmethod def _workspace_thread_key( - team_id: str, channel_id: str, thread_ts: str - ) -> Optional[Tuple[str, str, str]]: + team_id: str, channel_id: str, thread_ts: str) -> Optional[Tuple[str, str, str]]: """Return a workspace-scoped key for thread-local state. Slack Connect can expose the same channel/thread IDs in several workspaces.""" if not channel_id or not thread_ts: @@ -3514,13 +3210,11 @@ class SlackAdapter(BasePlatformAdapter): contexts[key] = { field: value for field, value in metadata.items() - if field in {"channel_id", "context_channel_id", "team_id", "user_id"} and value - } + if field in {"channel_id", "context_channel_id", "team_id", "user_id"} and value} self._trim_oldest_dict_entries(contexts, self._AGENT_VIEW_CONTEXTS_MAX) def _agent_view_context_for_event( - self, event: dict, team_id: str, user_id: str - ) -> Dict[str, str]: + self, event: dict, team_id: str, user_id: str) -> Dict[str, str]: """Read Slack's inline Agent context, falling back to lifecycle state.""" context = event.get("app_context") or event.get("context") or {} context_channel_id = self._context_channel_id(context) @@ -3530,8 +3224,7 @@ class SlackAdapter(BasePlatformAdapter): return { "context_channel_id": context_channel_id or cached.get("context_channel_id", ""), "team_id": team_id, - "user_id": user_id, - } + "user_id": user_id} def _remember_processed_message_ts(self, ts: str) -> None: """Mark a message ts as claimed, for the ``message_changed`` guard. @@ -3580,31 +3273,26 @@ class SlackAdapter(BasePlatformAdapter): return "" def _extract_assistant_thread_metadata( - self, event: dict, body: Optional[dict] = None - ) -> Dict[str, str]: + self, event: dict, body: Optional[dict] = None) -> Dict[str, str]: """Extract Slack Assistant thread identity data from an event payload.""" assistant_thread = event.get("assistant_thread") or {} context = ( assistant_thread.get("context") or event.get("app_context") or event.get("context") - or {} - ) + or {}) channel_id = ( assistant_thread.get("channel_id") or event.get("channel") or context.get("channel_id") - or "" - ) + or "") thread_ts = ( assistant_thread.get("thread_ts") or event.get("thread_ts") or event.get("message_ts") - or "" - ) + or "") user_id = ( - assistant_thread.get("user_id") or event.get("user") or context.get("user_id") or "" - ) + assistant_thread.get("user_id") or event.get("user") or context.get("user_id") or "") team_id = self._event_team_id(event, body) or str(assistant_thread.get("team_id") or "") context_channel_id = self._context_channel_id(context) return { @@ -3612,8 +3300,7 @@ class SlackAdapter(BasePlatformAdapter): "thread_ts": _str_or_empty(thread_ts), "user_id": _str_or_empty(user_id), "team_id": _str_or_empty(team_id), - "context_channel_id": _str_or_empty(context_channel_id), - } + "context_channel_id": _str_or_empty(context_channel_id)} def _cache_assistant_thread_metadata(self, metadata: Dict[str, str]) -> None: """Remember workspace-local assistant identity for later message events.""" @@ -3630,14 +3317,8 @@ class SlackAdapter(BasePlatformAdapter): self._remember_channel_team(channel_id, team_id) def _lookup_assistant_thread_metadata( - self, - event: dict, - *, - channel_id: str = "", - thread_ts: str = "", - team_id: str = "", - body: Optional[dict] = None, - ) -> Dict[str, str]: + self, event: dict, *, channel_id: str = "", thread_ts: str = "", team_id: str = "", + body: Optional[dict] = None) -> Dict[str, str]: """Load workspace-scoped assistant metadata for the current event.""" metadata = self._extract_assistant_thread_metadata(event, body) if channel_id and not metadata.get("channel_id"): @@ -3647,10 +3328,8 @@ class SlackAdapter(BasePlatformAdapter): if team_id and not metadata.get("team_id"): metadata["team_id"] = str(team_id) key = self._workspace_thread_key( - metadata.get("team_id", ""), - metadata.get("channel_id", ""), - metadata.get("thread_ts", ""), - ) + metadata.get("team_id", ""), metadata.get("channel_id", ""), + metadata.get("thread_ts", "")) cached = self._assistant_threads.get(key, {}) if key else {} if cached: return {**cached, **{k: v for k, v in metadata.items() if v}} @@ -3678,8 +3357,7 @@ class SlackAdapter(BasePlatformAdapter): return title, prompts async def _set_assistant_suggested_prompts( - self, channel_id: str, *, team_id: str = "", thread_ts: str = "" - ) -> None: + self, channel_id: str, *, team_id: str = "", thread_ts: str = "") -> None: """Best-effort Slack AI suggested prompts setup.""" if not self._app or not channel_id: return @@ -3702,16 +3380,14 @@ class SlackAdapter(BasePlatformAdapter): return bool(raw) async def _set_assistant_thread_title( - self, channel_id: str, thread_ts: str, title_source: str, *, team_id: str = "" - ) -> None: + self, channel_id: str, thread_ts: str, title_source: str, *, team_id: str = "") -> None: """Best-effort title for visible Slack AI DM threads.""" if ( not self._app or not channel_id or not thread_ts or not title_source - or not self._assistant_thread_title_enabled() - ): + or not self._assistant_thread_title_enabled()): return key = self._workspace_thread_key(team_id, channel_id, thread_ts) if not key or key in self._titled_assistant_threads: @@ -3722,16 +3398,14 @@ class SlackAdapter(BasePlatformAdapter): title = title[:77].rstrip() + "..." if len(title) > 80 else title try: await self._get_client(channel_id, team_id=team_id).assistant_threads_setTitle( - channel_id=channel_id, thread_ts=thread_ts, title=title - ) + channel_id=channel_id, thread_ts=thread_ts, title=title) except Exception as e: logger.debug("[Slack] assistant.threads.setTitle failed: %s", e) return self._titled_assistant_threads.add(key) # Evict oldest thread_ts first so recently titled threads keep their guard. self._evict_oldest_by_ts( - self._titled_assistant_threads, self._TITLED_ASSISTANT_THREADS_MAX, lambda e: e[2] - ) + self._titled_assistant_threads, self._TITLED_ASSISTANT_THREADS_MAX, lambda e: e[2]) def _seed_dm_session( self, metadata: Dict[str, str], *, thread_ts: Optional[str], fail_log: Tuple[Any, ...] @@ -3744,22 +3418,19 @@ class SlackAdapter(BasePlatformAdapter): source = self.build_source( chat_id=channel_id, chat_name=self._channel_name_cache.get( - (str(metadata.get("team_id") or ""), channel_id), channel_id - ), + (str(metadata.get("team_id") or ""), channel_id), channel_id), chat_type="dm", user_id=user_id, thread_id=thread_ts, chat_topic=metadata.get("context_channel_id") or None, - scope_id=metadata.get("team_id") or None, - ) + scope_id=metadata.get("team_id") or None) try: session_store.get_or_create_session(source) except Exception: logger.debug(*fail_log, exc_info=True) async def _handle_assistant_thread_lifecycle_event( - self, event: dict, body: Optional[dict] = None - ) -> None: + self, event: dict, body: Optional[dict] = None) -> None: """Handle Slack Assistant lifecycle events that carry user/thread identity.""" metadata = self._extract_assistant_thread_metadata(event, body) self._cache_assistant_thread_metadata(metadata) @@ -3770,15 +3441,10 @@ class SlackAdapter(BasePlatformAdapter): thread_ts=thread_ts, fail_log=( "[Slack] Failed to seed assistant thread session for %s/%s", - metadata.get("channel_id", ""), - thread_ts, - ), - ) + metadata.get("channel_id", ""), thread_ts)) await self._set_assistant_suggested_prompts( - metadata.get("channel_id", ""), - team_id=metadata.get("team_id", ""), - thread_ts=metadata.get("thread_ts", ""), - ) + metadata.get("channel_id", ""), team_id=metadata.get("team_id", ""), + thread_ts=metadata.get("thread_ts", "")) async def _handle_app_context_changed(self, event: dict, body: Optional[dict] = None) -> None: """Cache the current Agent-view context without entering the agent loop.""" @@ -3796,8 +3462,7 @@ class SlackAdapter(BasePlatformAdapter): return { "context_channel_id": self._context_channel_id(context), "user_id": _str_or_empty(user_id), - "team_id": _str_or_empty(team_id), - } + "team_id": _str_or_empty(team_id)} async def _handle_app_home_opened(self, event: dict, body: Optional[dict] = None) -> None: """Handle Slack Agent DM-open lifecycle events without producing replies.""" @@ -3811,8 +3476,7 @@ class SlackAdapter(BasePlatformAdapter): "channel_id": _str_or_empty(channel_id), "user_id": fields["user_id"], "team_id": fields["team_id"], - "context_channel_id": fields["context_channel_id"], - } + "context_channel_id": fields["context_channel_id"]} self._cache_agent_view_context(metadata) # ``app_home_opened`` (tab == "messages") replaces ``assistant_thread_started`` in # Slack's Agent experience; lifecycle only (no welcome message, no agent loop). @@ -3820,13 +3484,9 @@ class SlackAdapter(BasePlatformAdapter): metadata, thread_ts=None, fail_log=( - "[Slack] Failed to seed agent DM session for %s", - metadata.get("channel_id", ""), - ), - ) + "[Slack] Failed to seed agent DM session for %s", metadata.get("channel_id", ""))) await self._set_assistant_suggested_prompts( - metadata["channel_id"], team_id=metadata["team_id"] - ) + metadata["channel_id"], team_id=metadata["team_id"]) # Common reaction names → unicode emoji. Used by ``_handle_slack_reaction`` # so skills that match on ``text`` see the same character whether the user @@ -3834,8 +3494,7 @@ class SlackAdapter(BasePlatformAdapter): _REACTION_EMOJI_MAP: ClassVar[Dict[str, str]] = { "thumbsup": "👍", "+1": "👍", "thumbsdown": "👎", "-1": "👎", "white_check_mark": "✅", "heavy_check_mark": "✅", "x": "❌", "no_entry": "⛔", "warning": "⚠️", "rotating_light": "🚨", - "eyes": "👀", "rocket": "🚀", "tada": "🎉", "fire": "🔥", "wave": "👋", - } + "eyes": "👀", "rocket": "🚀", "tada": "🎉", "fire": "🔥", "wave": "👋"} async def _handle_slack_reaction(self, event: dict, removed: bool = False) -> None: """Forward reaction events through the normal message pipeline. @@ -3876,8 +3535,7 @@ class SlackAdapter(BasePlatformAdapter): "message_ts": msg_ts, "team_id": team_id, "event_ts": event.get("event_ts"), - "raw_event": event, - }) + "raw_event": event}) except Exception: # pragma: no cover - hook contract is non-blocking logger.debug("[Slack] reaction hook forwarding failed", exc_info=True) # None → routing disabled; empty set → all emoji; non-empty → allowlist. @@ -3888,8 +3546,7 @@ class SlackAdapter(BasePlatformAdapter): if explicit_allowlist and reaction_name.strip(":") not in triggers: return thread_ts = await self._reaction_thread_ts( - client, channel_id, msg_ts, event, team_id, explicit_allowlist - ) + client, channel_id, msg_ts, event, team_id, explicit_allowlist) if thread_ts is None: return emoji_text = self._REACTION_EMOJI_MAP.get(reaction_name, reaction_name) @@ -3909,9 +3566,7 @@ class SlackAdapter(BasePlatformAdapter): "name": reaction_name, "action": action, "reacted_to_ts": msg_ts, - "event_ts": event.get("event_ts"), - }, - } + "event_ts": event.get("event_ts")}} if team_id: synthetic["team"] = team_id # Optional handoff target channel/thread. A channel-only target is a @@ -3929,14 +3584,8 @@ class SlackAdapter(BasePlatformAdapter): await self._handle_slack_message(synthetic) async def _reaction_thread_ts( - self, - client, - channel_id: str, - msg_ts: str, - event: dict, - team_id: str, - explicit_allowlist: bool, - ) -> Optional[str]: + self, client, channel_id: str, msg_ts: str, event: dict, team_id: str, + explicit_allowlist: bool) -> Optional[str]: """Thread to route a reaction into, or None when it must be dropped. Looks up the reacted-to message for its thread and author; on lookup failure @@ -3950,8 +3599,7 @@ class SlackAdapter(BasePlatformAdapter): if client is not None: try: history = await client.conversations_replies( - channel=channel_id, ts=msg_ts, limit=1, inclusive=True - ) + channel=channel_id, ts=msg_ts, limit=1, inclusive=True) messages = (history or {}).get("messages") or [] if messages: first = messages[0] @@ -4025,8 +3673,7 @@ class SlackAdapter(BasePlatformAdapter): channel_id = event.get("channel_id") or event.get("channel") or "" if self._is_ignored_channel(channel_id): logger.info( - "[Slack] Ignoring file_shared event in configured ignored channel %s", channel_id - ) + "[Slack] Ignoring file_shared event in configured ignored channel %s", channel_id) return file_id = event.get("file_id") or (event.get("file") or {}).get("id") or "" if not channel_id or not file_id: @@ -4037,17 +3684,14 @@ class SlackAdapter(BasePlatformAdapter): info_resp = await (client or self._get_client(channel_id)).files_info(file=file_id) except Exception as exc: detail = self._describe_slack_api_error( - getattr(exc, "response", None), file_obj={"id": file_id} - ) + getattr(exc, "response", None), file_obj={"id": file_id}) logger.warning("[Slack] files.info error for file_shared %s: %s", file_id, detail or exc) return if not info_resp.get("ok"): detail = self._describe_slack_api_error(info_resp, file_obj={"id": file_id}) logger.warning( - "[Slack] files.info failed for file_shared %s: %s", - file_id, - detail or info_resp.get("error"), - ) + "[Slack] files.info failed for file_shared %s: %s", file_id, + detail or info_resp.get("error")) return file_obj = info_resp.get("file") or {} if not str(file_obj.get("mimetype", "")).startswith("video/"): @@ -4069,8 +3713,7 @@ class SlackAdapter(BasePlatformAdapter): "channel_type": "im" if channel_id.startswith("D") else "channel", "team": team_id, "ts": "", # already recorded above; avoid tripping our own dedup guard - "files": [file_obj], - } + "files": [file_obj]} if thread_ts and thread_ts != ts: fallback_event["thread_ts"] = thread_ts await self._handle_slack_message(fallback_event) @@ -4085,8 +3728,7 @@ class SlackAdapter(BasePlatformAdapter): self._trim_mentioned_threads() async def _bot_authored_thread_root( - self, channel_id: str, thread_ts: str, team_id: str = "" - ) -> bool: + self, channel_id: str, thread_ts: str, team_id: str = "") -> bool: """Return True when the thread root was authored by this bot. Catches roots posted via direct chat.postMessage (outside send(), so not in _bot_message_ts); derived from the Slack API, so it survives restarts. Checks @@ -4103,23 +3745,15 @@ class SlackAdapter(BasePlatformAdapter): for cached_key, cached_entry in self._thread_context_cache.items(): if cached_key.startswith(f"{channel_id}:{thread_ts}:"): return bool( - cached_entry.parent_user_id and cached_entry.parent_user_id == bot_uid - ) + cached_entry.parent_user_id and cached_entry.parent_user_id == bot_uid) if attempt == 0: await self._fetch_thread_context( - channel_id=channel_id, thread_ts=thread_ts, current_ts="", team_id=team_id - ) + channel_id=channel_id, thread_ts=thread_ts, current_ts="", team_id=team_id) return False async def _should_wake_on_unmentioned_message( - self, - event_thread_ts, - channel_id: str, - user_id: str, - is_thread_reply: bool, - team_id: str = "", - chat_type: str = "group", - ) -> bool: + self, event_thread_ts, channel_id: str, user_id: str, is_thread_reply: bool, + team_id: str = "", chat_type: str = "group") -> bool: """Return True if the bot should wake on an un-mentioned message. Checks, in order: root sent via send() (_bot_message_ts); thread previously @-mentioned; active session; bot-authored root via raw chat.postMessage; thread parent @-mentioned the @@ -4130,22 +3764,16 @@ class SlackAdapter(BasePlatformAdapter): # Check scoped marker AND bare ts: entries recorded before team_id was # known are bare, and a scoped-vs-bare mismatch must not silence the bot. if is_thread_reply and ( - thread_marker in self._bot_message_ts or event_thread_ts in self._bot_message_ts - ): + thread_marker in self._bot_message_ts or event_thread_ts in self._bot_message_ts): return True if thread_marker in self._mentioned_threads or event_thread_ts in self._mentioned_threads: return True if is_thread_reply and self._has_active_session_for_thread( - channel_id=channel_id, - thread_ts=event_thread_ts, - user_id=user_id, - team_id=team_id, - chat_type=chat_type, - ): + channel_id=channel_id, thread_ts=event_thread_ts, user_id=user_id, team_id=team_id, + chat_type=chat_type): return True if is_thread_reply and await self._bot_authored_thread_root( - channel_id=channel_id, thread_ts=event_thread_ts, team_id=team_id - ): + channel_id=channel_id, thread_ts=event_thread_ts, team_id=team_id): return True # Thread PARENT @-mentioned the bot but the mention predates this process # (restart) or asked for a follow-up: a bare reply like "run" is for us. @@ -4153,11 +3781,8 @@ class SlackAdapter(BasePlatformAdapter): bot_uid = self._team_bot_user_ids.get(team_id, self._bot_user_id) if bot_uid: parent_text = await self._fetch_thread_parent_text( - channel_id=channel_id, - thread_ts=event_thread_ts, - team_id=team_id, - strip_bot_mention=False, - ) + channel_id=channel_id, thread_ts=event_thread_ts, team_id=team_id, + strip_bot_mention=False) if parent_text and f"<@{bot_uid}>" in parent_text: # Remember so later replies skip the fetch. if not self._slack_strict_mention(): @@ -4173,9 +3798,7 @@ class SlackAdapter(BasePlatformAdapter): if stripped_blocks: logger.debug( "Slack: extracted additional text from blocks " - "(likely quoted/forwarded content; chars=%d)", - len(stripped_blocks), - ) + "(likely quoted/forwarded content; chars=%d)", len(stripped_blocks)) text = (text.strip() + "\n" + stripped_blocks).strip() blocks_payload = _serialize_slack_blocks_for_agent(blocks) if blocks_payload: @@ -4197,8 +3820,7 @@ class SlackAdapter(BasePlatformAdapter): decision = ( self._is_sender_authorized(user_id, chat_type, channel_id) if user_id and getattr(self, "_authorization_check", None) is not None - else None - ) + else None) auth_fn = self._runner_auth_fn() if decision is None and user_id and callable(auth_fn): source = self.build_source( @@ -4207,25 +3829,14 @@ class SlackAdapter(BasePlatformAdapter): decision = bool(auth_fn(source)) if decision is False: logger.warning( - "[Slack] Early reject of unauthorized user %s in channel %s", user_id, channel_id - ) + "[Slack] Early reject of unauthorized user %s in channel %s", user_id, channel_id) return True return False async def _channel_gate_allows( - self, - *, - channel_id: str, - routing_text: str, - bot_uid: str, - is_mentioned: bool, - is_thread_reply: bool, - event_thread_ts, - user_id: str, - team_id: str, - is_dm: bool, - force_process: bool, - ) -> bool: + self, *, channel_id: str, routing_text: str, bot_uid: str, is_mentioned: bool, + is_thread_reply: bool, event_thread_ts, user_id: str, team_id: str, is_dm: bool, + force_process: bool) -> bool: """Channel/MPIM routing gate: respond if free-response channel (still gated by ``thread_require_mention``), @mentioned, or a wake check passes. Always silent outside ``allowed_channels`` or when addressed to another @@ -4240,39 +3851,30 @@ class SlackAdapter(BasePlatformAdapter): self._slack_ignore_other_user_mentions() and not is_mentioned and not self._slack_message_mentions_self(routing_text, self_uids) - and self._slack_message_addressed_to_other_user(routing_text, self_uids) - ): + and self._slack_message_addressed_to_other_user(routing_text, self_uids)): logger.debug( - "[Slack] Ignoring message addressed to another user in channel %s", channel_id - ) + "[Slack] Ignoring message addressed to another user in channel %s", channel_id) return False thread_gated = self._slack_thread_require_mention() and is_thread_reply and not is_mentioned if force_process: return True free_channel = channel_id not in self._slack_require_mention_channels() and ( - channel_id in self._slack_free_response_channels() or not self._slack_require_mention() - ) + channel_id in self._slack_free_response_channels() or not self._slack_require_mention()) if not free_channel and self._slack_strict_mention() and not is_mentioned: return False # Strict mode: ignore until @-mentioned again if thread_gated: logger.debug( "[Slack] Ignoring thread reply without mention " - "(thread_require_mention=true): channel=%s thread_ts=%s", - channel_id, - event_thread_ts, - ) + "(thread_require_mention=true): channel=%s thread_ts=%s", channel_id, + event_thread_ts) return False if free_channel: return True if not is_mentioned: return await self._should_wake_on_unmentioned_message( - event_thread_ts=event_thread_ts, - channel_id=channel_id, - user_id=user_id, - team_id=team_id, - is_thread_reply=is_thread_reply, - chat_type="dm" if is_dm else "group", - ) + event_thread_ts=event_thread_ts, channel_id=channel_id, user_id=user_id, + team_id=team_id, is_thread_reply=is_thread_reply, + chat_type="dm" if is_dm else "group") return True def _normalize_changed_message(self, event: dict) -> Optional[dict]: @@ -4291,8 +3893,7 @@ class SlackAdapter(BasePlatformAdapter): changed_event_ts = ( str(event.get("event_ts") or edited_ts or "") or (outer_event_ts if outer_event_ts != original_message_ts else "") - or (f"{original_message_ts}:changed" if original_message_ts else "") - ) + or (f"{original_message_ts}:changed" if original_message_ts else "")) normalized_event = dict(updated_message) for key in ("channel", "channel_type", "team", "team_id"): if not normalized_event.get(key) and event.get(key): @@ -4340,8 +3941,7 @@ class SlackAdapter(BasePlatformAdapter): return text def _session_thread_ts( - self, event: dict, ts: str, is_dm: bool, assistant_meta: Dict[str, str] - ) -> Optional[str]: + self, event: dict, ts: str, is_dm: bool, assistant_meta: Dict[str, str]) -> Optional[str]: """Resolve the thread_ts used for session keying. DMs: each top-level thread is its own session unless ``dm_top_level_threads_as_sessions: false``. Reaction handoffs (``_hermes_no_thread_response``) reply top-level, never under the @@ -4363,16 +3963,8 @@ class SlackAdapter(BasePlatformAdapter): return None async def _hydrate_thread_context( - self, - *, - channel_id: str, - event_thread_ts, - ts: str, - user_id: str, - team_id: str, - is_thread_reply: bool, - is_mentioned: bool, - is_dm: bool, + self, *, channel_id: str, event_thread_ts, ts: str, user_id: str, team_id: str, + is_thread_reply: bool, is_mentioned: bool, is_dm: bool, ) -> Tuple[Optional[str], List[str], List[str]]: """Return ``(channel_context, root_media_urls, root_media_types)`` for a thread reply. No session: hydrate full thread + root images once, set watermark. Session + @mention: delta @@ -4385,56 +3977,40 @@ class SlackAdapter(BasePlatformAdapter): if not is_thread_reply: return channel_context, thread_root_media_urls, thread_root_media_types has_active_thread_session = self._has_active_session_for_thread( - channel_id=channel_id, - thread_ts=event_thread_ts, - user_id=user_id, - team_id=team_id, - chat_type="dm" if is_dm else "group", - ) + channel_id=channel_id, thread_ts=event_thread_ts, user_id=user_id, team_id=team_id, + chat_type="dm" if is_dm else "group") async def _fetch(**kw) -> None: nonlocal channel_context thread_context = await self._fetch_thread_context( - channel_id=channel_id, - thread_ts=event_thread_ts, - current_ts=ts, - team_id=team_id, - **kw, - ) + channel_id=channel_id, thread_ts=event_thread_ts, current_ts=ts, team_id=team_id, + **kw) if thread_context: channel_context = thread_context def _advance_watermark() -> None: self._set_thread_watermark( - channel_id=channel_id, - thread_ts=event_thread_ts, - user_id=user_id, - watermark_ts=ts, - team_id=team_id, - ) + channel_id=channel_id, thread_ts=event_thread_ts, user_id=user_id, watermark_ts=ts, + team_id=team_id) def _mark_checked() -> None: self._mark_thread_rehydration_checked(channel_id, event_thread_ts, user_id, team_id) def _watermark() -> str: return self._get_thread_watermark( - channel_id=channel_id, thread_ts=event_thread_ts, user_id=user_id, team_id=team_id - ) + channel_id=channel_id, thread_ts=event_thread_ts, user_id=user_id, team_id=team_id) if not has_active_thread_session: await _fetch() ( - thread_root_media_urls, - thread_root_media_types, + thread_root_media_urls, thread_root_media_types, ) = await self._collect_thread_root_images( - channel_id=channel_id, thread_ts=event_thread_ts, team_id=team_id - ) + channel_id=channel_id, thread_ts=event_thread_ts, team_id=team_id) elif is_mentioned: await _fetch(after_ts=_watermark(), force_refresh=True) else: rehydration_key = self._thread_rehydration_key( - channel_id, event_thread_ts, user_id, team_id - ) + channel_id, event_thread_ts, user_id, team_id) if rehydration_key in self._thread_rehydration_checked: _advance_watermark() return channel_context, thread_root_media_urls, thread_root_media_types @@ -4451,10 +4027,8 @@ class SlackAdapter(BasePlatformAdapter): if not media_types: return MessageType.TEXT for prefix, kind in ( - ("image/", MessageType.PHOTO), - ("video/", MessageType.VIDEO), - ("audio/", MessageType.VOICE), - ): + ("image/", MessageType.PHOTO), ("video/", MessageType.VIDEO), + ("audio/", MessageType.VOICE)): if any(m.startswith(prefix) for m in media_types): return kind return MessageType.DOCUMENT @@ -4470,8 +4044,7 @@ class SlackAdapter(BasePlatformAdapter): channel_prompt = ( f"{identity_prompt}\n\n{channel_prompt}".strip() if channel_prompt - else identity_prompt - ) + else identity_prompt) return channel_prompt def _track_reacting_message(self, team_id: str, ts: str) -> None: @@ -4496,10 +4069,7 @@ class SlackAdapter(BasePlatformAdapter): _claims.pop(_ts, None) logger.warning( "[%s] handler failed after claiming ts=%s; claim released " - "so a retry or edit can re-drive the turn", - self.name, - _ts, - ) + "so a retry or edit can re-drive the turn", self.name, _ts) raise async def _drop_bot_sender(self, event: dict) -> bool: @@ -4514,10 +4084,8 @@ class SlackAdapter(BasePlatformAdapter): sender_is_bot = self._event_declares_bot_sender(event) if not sender_is_bot and msg_user and not event.get("client_msg_id"): sender_is_bot = await self._resolve_user_is_bot( - msg_user, - chat_id=event.get("channel", ""), - team_id=str(event.get("team") or event.get("team_id") or ""), - ) + msg_user, chat_id=event.get("channel", ""), + team_id=str(event.get("team") or event.get("team_id") or "")) if not sender_is_bot: return False allow_bots = self._slack_allow_bots() @@ -4529,50 +4097,96 @@ class SlackAdapter(BasePlatformAdapter): if self._bot_user_id and f"<@{self._bot_user_id}>" not in text_check: logger.debug( "[Slack] Dropping bot message under allow_bots=mentions: " - "no <@%s> mention in flat text or blocks", - self._bot_user_id, - ) + "no <@%s> mention in flat text or blocks", self._bot_user_id) return True return bool(msg_user and self._bot_user_id and msg_user == self._bot_user_id) - async def _handle_slack_message_impl(self, event: dict, payload: Optional[dict] = None) -> None: - """Handle an incoming Slack message event.""" + async def _prefilter_inbound( + self, event: dict, payload: Optional[dict]) -> Optional[Tuple[dict, str, str]]: + """Normalize edits, then drop replays / ignored channels / bot posts / deletions. + Returns ``(event, team_id, channel_id)`` for messages the handler should consider.""" # Entry log BEFORE any filtering so operators can tell "dropped here" # from "never subscribed in the manifest". Metadata only, never text. if logger.isEnabledFor(logging.DEBUG): _bot_profile = event.get("bot_profile") or {} logger.debug( "[Slack] event received type=%s subtype=%s user=%s bot_id=%s bot_name=%s " - "channel=%s ts=%s thread_ts=%s", - event.get("type"), - event.get("subtype"), - event.get("user", "") or "", - event.get("bot_id", "") or "", + "channel=%s ts=%s thread_ts=%s", event.get("type"), event.get("subtype"), + event.get("user", "") or "", event.get("bot_id", "") or "", (_bot_profile.get("name") if isinstance(_bot_profile, dict) else "") or "", - event.get("channel", ""), - event.get("ts", ""), - event.get("thread_ts", ""), - ) + event.get("channel", ""), event.get("ts", ""), event.get("thread_ts", "")) if event.get("subtype") == "message_changed": event = self._normalize_changed_message(event) if event is None: - return + return None # Socket Mode redelivers after reconnects. Scope by workspace: ts is # only unique per team. event_ts = event.get("_slack_changed_event_ts") or event.get("ts", "") dedup_team_id = self._event_team_id(event, payload) if event_ts and self._dedup.is_duplicate(self._workspace_event_id(dedup_team_id, event_ts)): - return + return None channel_id = event.get("channel", "") if self._is_ignored_channel(channel_id): logger.info("[Slack] Ignoring message in configured ignored channel %s", channel_id) - return + return None if await self._drop_bot_sender(event): - return + return None # Edits were normalized above so an @mention added by edit can wake the bot once. - subtype = event.get("subtype") - if subtype == "message_deleted": + if event.get("subtype") == "message_deleted": + return None + return event, dedup_team_id, channel_id + + async def _peer_bot_drop( + self, event: dict, user_id: str, bot_uid: Optional[str], channel_id: str, team_id: str, + is_mentioned: bool) -> bool: + """True when a bot *user* post (peer agent: no bot_id/subtype) must be dropped. + Such posts would otherwise re-trigger via old thread mentions or active sessions and cause + agent-agent loops. Under ``mentions`` only the current text counts as a summons.""" + if not user_id or user_id == bot_uid: + return False + sender_is_bot_user = self._event_declares_bot_sender(event) + if not sender_is_bot_user: + sender_is_bot_user = await self._resolve_user_is_bot( + user_id, chat_id=channel_id, team_id=team_id) + if not sender_is_bot_user: + return False + allow_bots = self._slack_allow_bots() + return allow_bots == "none" or (allow_bots == "mentions" and not is_mentioned) + + def _apply_bot_mention( + self, text: str, original_text: str, command_probe_text: str, is_command_text: bool, + bot_uid: str, thread_ts: Optional[str], team_id: str) -> Tuple[str, str, str, bool]: + """Strip our mention, re-probe for a command hidden behind it, remember the thread. + Returns updated ``(text, original_text, command_probe_text, is_command_text)``.""" + text = text.replace(f"<@{bot_uid}>", "").strip() + # Re-probe commands on the canonical text (not block-augmented, which + # would leak quoted text into arguments): handles ``@bot !cmd`` and + # ``@bot /cmd``, where the command token hid behind the mention. + mention_stripped = original_text.replace(f"<@{bot_uid}>", "").strip() + command_text = ( + mention_stripped + if mention_stripped.startswith("/") + else _rewrite_known_bang_command(mention_stripped)) + if command_text.startswith("/"): + original_text = text = command_probe_text = command_text + is_command_text = True + # Remember the thread so follow-ups auto-trigger. Skipped under + # strict_mention/thread_require_mention (would defeat them and + # re-enable ack loops). Uses session-scoped ``thread_ts`` because a + # top-level @mention STARTS a thread whose replies must trigger too. + if ( + thread_ts + and not self._slack_strict_mention() + and not self._slack_thread_require_mention()): + self._register_mentioned_thread(thread_ts, team_id=team_id) + return text, original_text, command_probe_text, is_command_text + + async def _handle_slack_message_impl(self, event: dict, payload: Optional[dict] = None) -> None: + """Handle an incoming Slack message event.""" + accepted = await self._prefilter_inbound(event, payload) + if accepted is None: return + event, dedup_team_id, channel_id = accepted original_text = event.get("text", "") # Slack rejects slash commands inside threads, so a leading ``!`` is an # alternate prefix rewritten to ``/`` — only when the first token is a @@ -4588,18 +4202,13 @@ class SlackAdapter(BasePlatformAdapter): blocks = event.get("blocks") if blocks and not is_command_text: text = self._append_block_text( - text, blocks, self._team_bot_user_ids.get(dedup_team_id, self._bot_user_id) or "" - ) + text, blocks, self._team_bot_user_ids.get(dedup_team_id, self._bot_user_id) or "") text = self._append_link_unfurls(text, event.get("attachments") or []) ts = event.get("ts", "") outer_team_id = dedup_team_id assistant_meta = self._lookup_assistant_thread_metadata( - event, - channel_id=channel_id, - thread_ts=event.get("thread_ts", ""), - team_id=outer_team_id, - body=payload, - ) + event, channel_id=channel_id, thread_ts=event.get("thread_ts", ""), + team_id=outer_team_id, body=payload) user_id = event.get("user") or assistant_meta.get("user_id", "") if not channel_id: channel_id = assistant_meta.get("channel_id", "") @@ -4608,8 +4217,7 @@ class SlackAdapter(BasePlatformAdapter): outer_team_id or assistant_meta.get("team_id", "") or self._channel_team.get(channel_id, "") ) agent_context = self._agent_view_context_for_event( - event, str(team_id or ""), str(user_id or "") - ) + event, str(team_id or ""), str(user_id or "")) if team_id and channel_id: self._remember_channel_team(channel_id, team_id) channel_type = event.get("channel_type", "") or ("im" if channel_id.startswith("D") else "") @@ -4617,9 +4225,7 @@ class SlackAdapter(BasePlatformAdapter): if is_dm and self._slack_disable_dms(): logger.info( "[Slack] Ignoring DM because Slack DMs are disabled: channel=%s user=%s", - channel_id, - user_id, - ) + channel_id, user_id) return # Only a 1:1 IM earns DM exemptions (no mention needed, free reactions). # An MPIM is a shared surface and obeys channel gating like any channel; @@ -4638,45 +4244,22 @@ class SlackAdapter(BasePlatformAdapter): routing_text = _slack_mention_detection_text(event) or original_text or "" is_mentioned = bool( (bot_uid and f"<@{bot_uid}>" in routing_text) - or self._slack_message_matches_mention_patterns(routing_text) - ) + or self._slack_message_matches_mention_patterns(routing_text)) event_thread_ts = event.get("thread_ts") is_thread_reply = bool(event_thread_ts and event_thread_ts != ts) # Internal triggers (reactions) skip the mention requirement but NOT # allowed_channels or user authorization. force_process = bool(event.get("_hermes_force_process")) - # Bot posts with only a bot *user* id (no bot_id/subtype — peer Hermes - # agents look like this) would otherwise re-trigger via old thread - # mentions or active sessions and cause agent-agent loops. Under - # "mentions" only the current text counts as a summons. - if user_id and user_id != bot_uid: - sender_is_bot_user = self._event_declares_bot_sender(event) - if not sender_is_bot_user: - sender_is_bot_user = await self._resolve_user_is_bot( - user_id, chat_id=channel_id, team_id=team_id - ) - if sender_is_bot_user: - allow_bots = self._slack_allow_bots() - if allow_bots == "none": - return - if allow_bots == "mentions" and not is_mentioned: - return + if await self._peer_bot_drop(event, user_id, bot_uid, channel_id, team_id, is_mentioned): + return if ( not is_one_to_one_dm and bot_uid and not await self._channel_gate_allows( - channel_id=channel_id, - routing_text=routing_text, - bot_uid=bot_uid, - is_mentioned=is_mentioned, - is_thread_reply=is_thread_reply, - event_thread_ts=event_thread_ts, - user_id=user_id, - team_id=team_id, - is_dm=is_dm, - force_process=force_process, - ) - ): + channel_id=channel_id, routing_text=routing_text, bot_uid=bot_uid, + is_mentioned=is_mentioned, is_thread_reply=is_thread_reply, + event_thread_ts=event_thread_ts, user_id=user_id, team_id=team_id, is_dm=is_dm, + force_process=force_process)): return # Claim the message ts HERE: a link unfurl emits `message_changed` for # the same message with a different event ts, so only the @@ -4688,48 +4271,20 @@ class SlackAdapter(BasePlatformAdapter): if _claim_ts: self._remember_processed_message_ts(_claim_ts) if is_mentioned: - text = text.replace(f"<@{bot_uid}>", "").strip() - # Re-probe commands on the canonical text (not block-augmented, which - # would leak quoted text into arguments): handles ``@bot !cmd`` and - # ``@bot /cmd``, where the command token hid behind the mention. - mention_stripped = original_text.replace(f"<@{bot_uid}>", "").strip() - command_text = ( - mention_stripped - if mention_stripped.startswith("/") - else _rewrite_known_bang_command(mention_stripped) + text, original_text, command_probe_text, is_command_text = self._apply_bot_mention( + text, original_text, command_probe_text, is_command_text, bot_uid, thread_ts, + team_id, ) - if command_text.startswith("/"): - original_text = text = command_probe_text = command_text - is_command_text = True - # Remember the thread so follow-ups auto-trigger. Skipped under - # strict_mention/thread_require_mention (would defeat them and - # re-enable ack loops). Uses session-scoped ``thread_ts`` because a - # top-level @mention STARTS a thread whose replies must trigger too. - if ( - thread_ts - and not self._slack_strict_mention() - and not self._slack_thread_require_mention() - ): - self._register_mentioned_thread(thread_ts, team_id=team_id) # Thread history stays out of ``text``: prepending would push a command off char zero. ( - channel_context, - thread_root_media_urls, - thread_root_media_types, + channel_context, thread_root_media_urls, thread_root_media_types, ) = await self._hydrate_thread_context( - channel_id=channel_id, - event_thread_ts=event_thread_ts, - ts=ts, - user_id=user_id, - team_id=team_id, - is_thread_reply=is_thread_reply, - is_mentioned=is_mentioned, - is_dm=is_dm, - ) + channel_id=channel_id, event_thread_ts=event_thread_ts, ts=ts, user_id=user_id, + team_id=team_id, is_thread_reply=is_thread_reply, is_mentioned=is_mentioned, + is_dm=is_dm) # Thread-root media is delivered ahead of the trigger message's own files. media_urls, media_types, text = await self._collect_inbound_media( - event, channel_id, team_id, text, thread_root_media_urls, thread_root_media_types - ) + event, channel_id, team_id, text, thread_root_media_urls, thread_root_media_types) # Commands are restored from canonical authored input: the gateway # parser needs the token at char zero and enrichment (blocks, unfurls, # file text, history) must never mutate command arguments. @@ -4741,8 +4296,7 @@ class SlackAdapter(BasePlatformAdapter): # Best-effort: title the DM thread from the prompt for Slack's AI Agent Messages tab. if is_dm and thread_ts and msg_type != MessageType.COMMAND: await self._set_assistant_thread_title( - channel_id, thread_ts, original_text or text, team_id=team_id - ) + channel_id, thread_ts, original_text or text, team_id=team_id) source = self.build_source( chat_id=channel_id, chat_name=channel_name, @@ -4754,8 +4308,7 @@ class SlackAdapter(BasePlatformAdapter): # Workflow/app posts have user=None; flag them so the gateway # SLACK_ALLOW_BOTS bypass can authorize them. Same predicate as # the drop gate above (api_human_users stay human). - is_bot=self._event_declares_bot_sender(event), - ) + is_bot=self._event_declares_bot_sender(event)) from gateway.platforms.base import resolve_channel_skills _channel_prompt = self._channel_prompt_with_identity(channel_id, team_id) @@ -4780,9 +4333,7 @@ class SlackAdapter(BasePlatformAdapter): metadata={ "slack_team_id": team_id, "slack_channel_id": channel_id, - "slack_thread_ts": thread_ts, - }, - ) + "slack_thread_ts": thread_ts}) # React only when directly addressed; MPIMs are shared, so they need a # mention like any channel. if (is_one_to_one_dm or is_mentioned) and self._reactions_enabled(): @@ -4794,20 +4345,14 @@ class SlackAdapter(BasePlatformAdapter): if context_channel_id and context_channel_id != channel_id and not is_command_text: msg_event.text = ( f"[Slack app context: user is viewing channel {context_channel_id}]\n\n" - f"{msg_event.text}" - ) + f"{msg_event.text}") if ts: self._remember_processed_message_ts(ts) await self.handle_message(msg_event) def _note_attachment_failure( - self, - notices: List[str], - detail: Optional[str], - fallback_msg: str, - *fallback_args: Any, - exc_info: bool = False, - ) -> None: + self, notices: List[str], detail: Optional[str], fallback_msg: str, *fallback_args: Any, + exc_info: bool = False) -> None: """Record a user-facing attachment diagnostic, else log the raw failure.""" if detail: notices.append(detail) @@ -4837,12 +4382,8 @@ class SlackAdapter(BasePlatformAdapter): if notices is not None: detail = self._describe_slack_api_error(info_resp, file_obj=f) self._note_attachment_failure( - notices, - detail, - "[Slack] files.info failed for %s: %s", - file_id, - info_resp.get("error"), - ) + notices, detail, "[Slack] files.info failed for %s: %s", file_id, + info_resp.get("error")) return None @staticmethod @@ -4896,8 +4437,7 @@ class SlackAdapter(BasePlatformAdapter): return None raw_bytes = await self._download_slack_file_bytes(url, team_id=team_id) cached_path = cache_document_from_bytes( - raw_bytes, original_filename or f"document{ext or '.bin'}" - ) + raw_bytes, original_filename or f"document{ext or '.bin'}") doc_mime = SUPPORTED_DOCUMENT_TYPES.get(ext, mimetype or "application/octet-stream") logger.debug("[Slack] Cached user document: %s (%s)", cached_path, doc_mime) injection = "" @@ -4913,13 +4453,8 @@ class SlackAdapter(BasePlatformAdapter): return cached_path, doc_mime, injection async def _collect_inbound_media( - self, - event: dict, - channel_id: str, - team_id: str, - text: str, - thread_root_media_urls: List[str], - thread_root_media_types: List[str], + self, event: dict, channel_id: str, team_id: str, text: str, + thread_root_media_urls: List[str], thread_root_media_types: List[str], ) -> Tuple[List[str], List[str], str]: """Download/cache ``event["files"]`` and return ``(media_urls, media_types, text)``. Thread-root images lead the lists. Small text-like docs are injected into ``text`` (gated on @@ -4949,13 +4484,8 @@ class SlackAdapter(BasePlatformAdapter): text = f"{injection}\n\n{text}" if text else injection except Exception as e: # pragma: no cover - defensive logging self._note_attachment_failure( - notices, - self._describe_slack_download_failure(e, file_obj=f), - f"[Slack] Failed to cache {kind} from %s: %s", - url, - e, - exc_info=True, - ) + notices, self._describe_slack_download_failure(e, file_obj=f), + f"[Slack] Failed to cache {kind} from %s: %s", url, e, exc_info=True) if notices: notice_block = "[Slack attachment notice]\n" + "\n".join(f"- {n}" for n in notices) text = f"{notice_block}\n\n{text}" if text else notice_block @@ -4965,8 +4495,7 @@ class SlackAdapter(BasePlatformAdapter): @staticmethod def _button( - text: str, action_id: str, value: str, *, style: str = "", emoji: bool = False - ) -> dict: + text: str, action_id: str, value: str, *, style: str = "", emoji: bool = False) -> dict: """Block Kit button element; ``style``/``emoji`` keys only present when set.""" text_obj: Dict[str, Any] = {"type": "plain_text", "text": text} if emoji: @@ -4979,45 +4508,55 @@ class SlackAdapter(BasePlatformAdapter): return btn async def _post_interactive_blocks( - self, - chat_id: str, - text: str, - blocks: list, - metadata: Optional[Dict[str, Any]], - *, - sanitize: bool = True, - team_scoped: bool = True, - ): + self, chat_id: str, text: str, blocks: list, metadata: Optional[Dict[str, Any]], *, + sanitize: bool = True, team_scoped: bool = True): """chat.postMessage with ``blocks`` (threaded via metadata); returns the raw response.""" kwargs: Dict[str, Any] = { "channel": chat_id, "text": text, - "blocks": sanitize_blocks(blocks) if sanitize else blocks, - } + "blocks": sanitize_blocks(blocks) if sanitize else blocks} thread_ts = self._resolve_thread_ts(None, metadata) if thread_ts: kwargs["thread_ts"] = thread_ts team_id = self._metadata_team_id(metadata) if team_scoped else None return await self._get_client(chat_id, team_id=team_id).chat_postMessage(**kwargs) - async def send_exec_approval( - self, - chat_id: str, - command: str, - session_key: str, - description: str = "dangerous command", - metadata: Optional[Dict[str, Any]] = None, - allow_permanent: bool = True, - allow_session: bool = True, - smart_denied: bool = False, - ) -> SendResult: - """Send a Block Kit approval prompt with interactive buttons. - The buttons call ``resolve_gateway_approval()`` to unblock the waiting agent thread — same - mechanism as the text ``/approve`` flow.""" + async def _send_interactive_prompt( + self, chat_id: str, metadata: Optional[Dict[str, Any]], + build: Callable[[], Tuple[str, list]], label: str, *, + resolved: Optional[Dict[Any, bool]] = None, resolved_max: int = 0, + team_scoped_key: bool = True, sanitize: bool = True) -> SendResult: + """Shared body of the Block Kit prompt senders: DM-resolve, ``build()`` -> ``(fallback + text, blocks)``, post, then mark the message unresolved in ``resolved`` (double-click + guard). Any failure is logged as ``