diff --git a/plugins/hermes-agent/README.md b/plugins/hermes-agent/README.md index bf6266a1f..73165a7e9 100644 --- a/plugins/hermes-agent/README.md +++ b/plugins/hermes-agent/README.md @@ -25,6 +25,7 @@ Supported: - `INLINE_TOKEN`, `INLINE_BOT_TOKEN`, `platforms.inline.token`, and `inline.token` auth paths, including simple `${ENV_NAME}` config references. - Supervised loopback Node sidecar using the Inline realtime SDK. - Realtime inbound messages, catch-up, replies to bot messages, and action callbacks. +- Inbound SDK receipts remain pending until Python handles the event. Temporary routing-metadata and authorization failures retry before deduplication or effects, preserving per-chat order while other chats can progress. Lost acknowledgement responses retry only the acknowledgement. Agent-button preflight retries are bounded, then show a retry prompt so an unavailable button cannot indefinitely block the shared action cursor. Shutdown leaves unresolved receipts available for catch-up; this is an at-least-once handoff, not a durable exactly-once guarantee for model turns or external effects. - Outbound text, Markdown parsing, opt-in edit-message streaming, long-message splitting, edits, deletes, typing, and presence. - Inline reply-thread routing, explicit-request auto mode, `/threads` controls, explicit `/follow` and `/unfollow` dialog relevance controls, parent chat metadata, parent/thread prompt fallback, and thread-specific skill bindings. - Native Hermes `inline` tool for current-chat/thread reads, bounded history and search, exact message lookup, button-message sends, editing/deleting bot-owned messages, reactions, pin/unpin/list pins, reply-thread creation, top-level thread/chat creation outside the current conversation, and avatar presence/status. @@ -33,7 +34,7 @@ Supported: - DM and group policies, user allowlists, group sender allowlists, mention requirements, strict mention mode, allowed chats, and free-response chats. - Native Inline `/` command-menu sync for Hermes slash commands, including `/threads`, `/follow`, `/unfollow`, `/inline_sync`, `/inline_version`, and `/update`; typed slash commands continue to work even if menu sync is disabled or rejected. - Inline-native buttons for clarify prompts, command approvals, slash confirmations, and model selection. -- Agent-created `send_message`/`edit_message` button rows with opaque callback data. A callback is acknowledged immediately and becomes a normal Hermes turn naming the source message and exact action fields. The normal response edits that source message and clears omitted buttons; the agent can instead call `edit_message` with replacement buttons and finish with `NO_REPLY` so the explicit edit remains authoritative. +- Agent-created `send_message`/`edit_message` button rows with opaque callback data. A callback is acknowledged after its target and access checks succeed, then enters Hermes as a normal turn naming the source message and exact action fields. The normal response edits that source message and clears omitted buttons; the agent can instead call `edit_message` with replacement buttons and finish with `NO_REPLY` so the explicit edit remains authoritative. - Outbound local photo, video, voice, and document uploads with configurable size caps. - Inbound photo, video, voice, and document summaries, with URL-backed media cached locally for Hermes when available. - Reactions on bot messages, plus opt-in lifecycle/system events as synthetic Hermes messages. @@ -334,6 +335,22 @@ Access control follows Hermes' native platform model: Policy is evaluated in three ordered stages: **access**, then **wake**, then **delivery**. Access checks the chat and sender policy and is a hard gate; a mention, reply, callback, or command never grants access to a blocked chat or actor. Wake decides whether an allowed group turn invokes Hermes: free-response chats wake normally, while mention-gated chats require an explicit mention unless a configured reply-to-bot or followed-thread exception applies. Delivery keeps existing child-thread conversations in place; top-level `auto` creates a child only for explicit thread intent, `on` always creates one, and `off` stays flat. +`open` controls intake; it does not grant permission to use the bot. After +Inline's sender and group/thread restrictions pass, the adapter asks Hermes' +registered authorization callback before executing local commands, creating +reply threads, downloading media, collecting context, or exposing runtime +settings. This uses the same pairing and profile-aware authorization as native +Hermes adapters. An unapproved sender reaches Hermes' pairing or rejection +flow with a minimal text event and no preceding adapter mutations or media work. +Pairing approval does not override an explicit Inline allowlist, disabled +policy, or excluded group. DMs remain exempt from `allowed_chats`. +Child-thread authorization preserves both the child and parent chat IDs so +Hermes can select the correct profile. If required group metadata is unavailable, +the adapter denies the operation instead of guessing a profile. Controls that +change the parent's reply-thread mode also require access to that parent; +thread-local model and following settings remain available to an authorized +child-thread user. + Equivalent Hermes YAML can use `allow_from`, `allowed_users`, `group_allow_from`, `dm_policy`, `group_policy`, `require_mention`, `strict_mention`, `allowed_chats`, `free_response_chats`, `reply_threads`, @@ -361,11 +378,14 @@ platforms: skills: ["support-triage", "incident-report"] ``` -Inline-native button callbacks, such as approvals and clarify choices, require -the clicking actor to pass an explicit Inline or global Hermes allowlist, or -`INLINE_ALLOW_ALL_USERS=true` / `GATEWAY_ALLOW_ALL_USERS=true`. This includes -model-picker callbacks. The stricter callback gate prevents group-visible -buttons from becoming a bypass when message intake is otherwise `open`. +Inline-native button callbacks, including approvals, clarify choices, +model pickers, and thread controls, use the same local restrictions and Hermes +authorization callback. Pairing-approved users can use them without a second +Inline allowlist. A denial, error, or unknown result from a registered callback +blocks the action. Only standalone adapters without a registered callback use +explicit Inline/global allowlists or allow-all settings as a fallback; `open` +alone never authorizes controls. Settings requests from unauthorized users show +an access guide without runtime or model information. Adapter-owned controls use `system:` action IDs and stay in these deterministic handlers. Agent-authored callbacks use `agent:` IDs and follow the ordinary message intake policy because they are conversational input, not approval or diff --git a/plugins/hermes-agent/plugin/inline/adapter.py b/plugins/hermes-agent/plugin/inline/adapter.py index 37b3fc3c7..e336d0091 100644 --- a/plugins/hermes-agent/plugin/inline/adapter.py +++ b/plugins/hermes-agent/plugin/inline/adapter.py @@ -70,7 +70,10 @@ _DEDUP_WINDOW_SECONDS = 48 * 3600 _CHAT_INFO_CACHE_SECONDS = 10 * 60 _CHAT_INFO_CACHE_MAX_SIZE = 512 -_BOT_SETTINGS_CHAT_INFO_TIMEOUT_SECONDS = 2.0 +_CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS = 2.0 +_INBOUND_ACCESS_LOOKUP_TIMEOUT_SECONDS = 30.0 +_INBOUND_RETRY_INITIAL_SECONDS = 1.0 +_MAX_INBOUND_DELIVERIES = 16 # Provider discovery may invoke host CLIs. Keep repeated panel opens and # unrelated mutations off that slow path while always injecting the live # current/default model into the document below. @@ -191,6 +194,10 @@ class _InlineCommandSpec(NamedTuple): } +class InlineInboundDeferred(RuntimeError): + """Pre-effect dependency unavailable; keep the SDK delivery unacknowledged.""" + + class InlineSidecarError(RuntimeError): def __init__(self, path: str, status_code: int, message: str, error_kind: str = "unknown", raw: Optional[Any] = None): self.path = path @@ -973,6 +980,7 @@ def __init__(self, config: PlatformConfig, *, use_ephemeral_sidecar_port: bool = self._last_catalog_sync: Optional[Dict[str, Any]] = None self._inbound_task: Optional[asyncio.Task] = None self._bot_settings_tasks: set[asyncio.Task] = set() + self._inbound_deliveries: Dict[str, asyncio.Task] = {} self._bot_settings_locks: Dict[str, asyncio.Lock] = {} self._bot_settings_lock_users: Dict[str, int] = {} self._bot_settings_model_catalog_cache: Optional[tuple[float, List[Dict[str, Any]]]] = None @@ -1147,7 +1155,7 @@ async def _bot_settings_context( try: info = await asyncio.wait_for( self._get_chat_info(chat_id), - timeout=_BOT_SETTINGS_CHAT_INFO_TIMEOUT_SECONDS, + timeout=_CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS, ) except asyncio.TimeoutError: logger.warning("[inline] agent settings chat metadata timed out") @@ -1163,13 +1171,15 @@ async def _bot_settings_context( scope_chat_id = parent_chat_id or chat_id is_reply_thread = bool(parent_chat_id) chat_type = self._bot_settings_chat_type(info) - can_inspect = self._allowed(chat_type, actor_id) and self._chat_allowed( - chat_id, - chat_id if is_reply_thread else None, - parent_chat_id, + can_inspect = self._actor_authorization( + chat_type, actor_id, chat_id, + thread_id=chat_id if is_reply_thread else None, + parent_chat_id=parent_chat_id, + is_bot=_inline_sender_profile(event).get("bot") is True, ) - if not can_inspect: - return {"access": "guideOnly", "scope_id": scope_chat_id, "reply_threads": "auto"} + if can_inspect is not True: + return {"access": "guideOnly", "scope_id": scope_chat_id, "reply_threads": "auto", + **({"unavailable_reason": "authorization"} if can_inspect is None else {})} source = self.build_source( chat_id=chat_id, @@ -1180,7 +1190,8 @@ async def _bot_settings_context( parent_chat_id=parent_chat_id, ) runner = self._bot_settings_runner() - can_modify = self._actor_authorized(chat_type, actor_id) + can_modify = can_inspect + can_set_reply_threads = not is_reply_thread or await self._chat_actor_authorized(scope_chat_id, actor_id) context: Dict[str, Any] = { "access": "full" if can_modify else "readOnly", "scope_id": scope_chat_id, @@ -1189,7 +1200,8 @@ async def _bot_settings_context( "chat_type": chat_type, "is_reply_thread": is_reply_thread, "following": None if chat_type == "dm" else self._chat_follow_mode_following(info), - "reply_threads": self._reply_thread_mode_for_chat(scope_chat_id), + "reply_threads": self._reply_thread_mode_for_chat(scope_chat_id) if can_set_reply_threads else None, + "can_set_reply_threads": can_set_reply_threads, "can_set_default_model": can_modify, "source": source, "runner": runner, @@ -1323,7 +1335,7 @@ def _bot_settings_document(self, context: Dict[str, Any]) -> Dict[str, Any]: if context.get("access") == "guideOnly": unavailable_text = ( "Hermes could not verify this chat. Try again when the bot is reachable." - if context.get("unavailable_reason") == "chat_metadata" + if context.get("unavailable_reason") in {"chat_metadata", "authorization"} else "This chat is not allowed by Hermes' access policy." ) return { @@ -1386,15 +1398,26 @@ def _bot_settings_document(self, context: Dict[str, Any]) -> Dict[str, Any]: "id": "following", "label": "Following", "description": "Wake on eligible activity.", **common, "control": {"oneofKind": "toggle", "toggle": {"value": bool(context.get("following"))}}, }]}) - sections.append({"id": "replies", "items": [{ - "id": "reply-threads", "label": "Reply in threads", **common, - "control": {"oneofKind": "select", "select": { - "value": context.get("reply_threads") or "auto", - "options": [ - {"value": "auto", "label": "Auto", "description": "Agent decides.", "disabled": False}, - {"value": "on", "label": "On", "description": "Always use threads.", "disabled": False}, - {"value": "off", "label": "Off", "description": "Stay in chat.", "disabled": False}, - ], + if context.get("can_set_reply_threads"): + sections.append({"id": "replies", "items": [{ + "id": "reply-threads", "label": "Reply in threads", **common, + "control": {"oneofKind": "select", "select": { + "value": context.get("reply_threads") or "auto", + "options": [ + {"value": "auto", "label": "Auto", "description": "Agent decides.", "disabled": False}, + {"value": "on", "label": "On", "description": "Always use threads.", "disabled": False}, + {"value": "off", "label": "Off", "description": "Stay in chat.", "disabled": False}, + ], + }}, + }]}) + else: + sections.append({"id": "replies", "items": [{ + "id": "reply-threads-access", "label": "Reply in threads", "disabled": False, + "control": {"oneofKind": "info", "info": { + "text": ("Parent chat access is temporarily unavailable. Try again." + if context.get("can_set_reply_threads") is None + else "Changing reply mode requires access to the parent chat."), + "tone": _INLINE_BOT_SETTINGS_TONE_WARNING, }}, }]}) if disabled: @@ -1492,6 +1515,8 @@ async def _apply_bot_setting(self, event: Dict[str, Any], context: Dict[str, Any self._invalidate_chat_info(context["chat_id"]) return if item_id == "reply-threads" and value_body.get("oneofKind") == "stringValue": + if not context.get("can_set_reply_threads"): + raise ValueError("parent chat thread settings are not authorized") mode = str(value_body.get("stringValue") or "") if mode not in {"auto", "on", "off"}: raise ValueError("invalid reply mode") @@ -1768,6 +1793,7 @@ async def _handle_thread_command( *, chat_id: str, msg_id: str, + from_id: str, text: str, chat_type: str, thread_id: Optional[str], @@ -1779,6 +1805,9 @@ async def _handle_thread_command( metadata = {"thread_id": thread_id} if thread_id else None target_chat_id = parent_chat_id or chat_id + if target_chat_id != chat_id and not await self._chat_actor_authorized(target_chat_id, from_id): + await self.send(chat_id, "Parent chat access could not be confirmed. Check access and try again.", reply_to=msg_id, metadata=metadata) + return True if action == "on": self._set_reply_threads_for_chat(target_chat_id, "on") elif action == "off": @@ -1935,6 +1964,12 @@ async def disconnect(self) -> None: except Exception: pass self._inbound_task = None + deliveries = list(self._inbound_deliveries.values()) + for task in deliveries: + task.cancel() + if deliveries: + await asyncio.gather(*deliveries, return_exceptions=True) + self._inbound_deliveries.clear() for task in list(self._bot_settings_tasks): task.cancel() if self._bot_settings_tasks: @@ -2247,6 +2282,78 @@ async def _on_inbound(self, line: str) -> None: event = json.loads(line) except json.JSONDecodeError: return + delivery_id = event.get("_inlineDeliveryId") + if not isinstance(delivery_id, str) or not delivery_id: + await self._dispatch_inbound(event) + return + # The SDK bounds outstanding receipts and orders each chat. Retain only + # their active coroutines; reconnect replay must not run one twice. + if delivery_id in self._inbound_deliveries: + return + # An ACK may reach the sidecar while its response is lost, releasing an + # SDK slot before this coroutine finishes. Bound those extra ACK waiters + # too, using stream backpressure rather than another event queue. + while len(self._inbound_deliveries) >= _MAX_INBOUND_DELIVERIES: + await asyncio.wait(tuple(self._inbound_deliveries.values()), return_when=asyncio.FIRST_COMPLETED) + task = asyncio.create_task(self._consume_inbound_delivery(delivery_id, event)) + self._inbound_deliveries[delivery_id] = task + + def finished(completed: asyncio.Task) -> None: + self._inbound_deliveries.pop(delivery_id, None) + if not completed.cancelled(): + try: + completed.result() + except Exception as exc: + self._report_error("inbound.delivery", exc) + task.add_done_callback(finished) + + async def _consume_inbound_delivery(self, delivery_id: str, event: Dict[str, Any]) -> None: + try: + await self._handle_inbound_delivery(delivery_id, event) + except Exception as exc: + # Do not ACK an unexpected failure or silently pin its chat forever. + # Native notification shields runner teardown/reconnect from the + # cancellation of this task during adapter.disconnect(). + self._report_error("inbound.delivery", exc) + self._set_fatal_error("INBOUND_FAILED", str(exc), retryable=True) + await self._notify_fatal_error() + + async def _handle_inbound_delivery(self, delivery_id: str, event: Dict[str, Any]) -> None: + delay = _INBOUND_RETRY_INITIAL_SECONDS + attempts = 0 + while True: + attempts += 1 + try: + await self._dispatch_inbound(event) + break + except InlineSidecarError as exc: + if exc.error_kind not in {"forbidden", "not_found"}: + raise + logger.info("[inline] inbound target is no longer accessible; retiring delivery") + break + except InlineInboundDeferred as exc: + # Action events share the SDK's user-cursor barrier. A persistent + # outage must not hold every chat behind a stale button press. + if event.get("kind") == "message.action.invoke" and attempts >= 2: + await self._answer_action(str(event.get("interactionId") or ""), + "Temporarily unavailable. Please try again.") + break + logger.warning("[inline] inbound preflight unavailable; retrying in %.1fs: %s", delay, exc) + await asyncio.sleep(delay) + delay = min(delay * 2, 30.0) + # A lost ACK response must retry only the ACK, never local commands or + # model delivery. The endpoint is idempotent, including after reconnect. + delay = _INBOUND_RETRY_INITIAL_SECONDS + while True: + try: + await self._sidecar_call("/inbound/ack", {"deliveryId": delivery_id}) + return + except Exception: + logger.warning("[inline] inbound acknowledgement unavailable; retrying in %.1fs", delay) + await asyncio.sleep(delay) + delay = min(delay * 2, 30.0) + + async def _dispatch_inbound(self, event: Dict[str, Any]) -> None: self._me_id = str(event.get("meId") or self._me_id or "") or None self._me_username = _normalize_inline_username(event.get("meUsername") or self._me_username) # Chat snapshots include mutable dialog and routing fields such as @@ -2257,10 +2364,16 @@ async def _on_inbound(self, line: str) -> None: self._invalidate_chat_info(event.get("chatId")) kind = event.get("kind") if kind == "bot.chatSettings.request": - self._schedule_bot_settings_task(self._handle_bot_settings_request(event)) + if event.get("_inlineDeliveryId"): + await self._handle_bot_settings_request(event) + else: + self._schedule_bot_settings_task(self._handle_bot_settings_request(event)) return if kind == "bot.chatSettings.item.invoke": - self._schedule_bot_settings_task(self._handle_bot_settings_item(event)) + if event.get("_inlineDeliveryId"): + await self._handle_bot_settings_item(event) + else: + self._schedule_bot_settings_task(self._handle_bot_settings_item(event)) return if kind == "message.action.invoke": if await self._handle_action(event): @@ -2349,10 +2462,6 @@ async def _dispatch_message(self, event: Dict[str, Any], *, edit: bool = False) dedup_key = f"new:{chat_id}:msg:{msg_id}:date:{message_date}" else: dedup_key = f"new:{chat_id}:msg:{msg_id}" - if self._is_duplicate(dedup_key): - return - if not edit and self._is_duplicate_message_instance(chat_id, msg, event): - return from_id = str(msg.get("fromId") or "") if msg.get("out") or (self._me_id and from_id == self._me_id): return @@ -2361,18 +2470,7 @@ async def _dispatch_message(self, event: Dict[str, Any], *, edit: bool = False) raw_message_text = str(msg.get("message") or "") text = raw_message_text.strip() - media_text, media_urls, media_types, message_type = await self._normalize_media(msg) - if media_text: - text = f"{text}\n{media_text}".strip() if text else media_text - if not text and not media_urls: - text = "[Inline message with no text]" explicitly_mentions_me = self._message_entity_mentions_me(msg) - if event.get("_inlineSenderProvenanceVerified") is False and not explicitly_mentions_me: - self._remember_observed_context(chat_id, msg, text) - return - if sender_profile.get("bot") is True and not explicitly_mentions_me: - self._remember_observed_context(chat_id, msg, text) - return chat_type = self._chat_type_from_message(msg) thread_id = self._thread_id_from_message(msg) @@ -2381,7 +2479,20 @@ async def _dispatch_message(self, event: Dict[str, Any], *, edit: bool = False) chat_name = chat_id chat_info: Dict[str, Any] = {} if chat_type == "group": - chat_info = await self._get_chat_info(chat_id) + try: + chat_info = await asyncio.wait_for( + self._get_chat_info(chat_id, required=True), + timeout=_CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS if agent_action else _INBOUND_ACCESS_LOOKUP_TIMEOUT_SECONDS, + ) + except asyncio.TimeoutError as exc: + raise InlineInboundDeferred("chat metadata timed out") from exc + except InlineSidecarError as exc: + if exc.error_kind in {"forbidden", "not_found"}: + logger.info("[inline] chat is no longer accessible; retiring inbound delivery") + return + raise + if not chat_info: + raise InlineInboundDeferred("chat metadata unavailable") chat_name = self._chat_title_from_info(chat_info) or chat_id info_parent_chat_id = self._chat_info_id(chat_info, "parentChatId") info_parent_message_id = self._chat_info_id(chat_info, "parentMessageId") @@ -2402,9 +2513,52 @@ async def _dispatch_message(self, event: Dict[str, Any], *, edit: bool = False) ) if has_command_target and not command_addressed_to_me: return + actor_authorized = self._actor_authorization( + chat_type, from_id, chat_id, thread_id=thread_id, + parent_chat_id=parent_chat_id, is_bot=sender_profile.get("bot") is True, + ) + if actor_authorized is None: + raise InlineInboundDeferred("host authorization unavailable") + if self._is_duplicate(dedup_key): + return + if not edit and self._is_duplicate_message_instance(chat_id, msg, event): + return + context_only = not explicitly_mentions_me and ( + event.get("_inlineSenderProvenanceVerified") is False or sender_profile.get("bot") is True + ) + if not actor_authorized: + if context_only: + return + # Let Hermes pair/ignore/decline the sender, without first downloading + # media, collecting context, changing settings, or creating a thread. + source = self.build_source( + chat_id=chat_id, chat_name=chat_name, chat_type=chat_type, + user_id=from_id, user_name=sender_name or None, + thread_id=thread_id, parent_chat_id=parent_chat_id, message_id=msg_id, + is_bot=sender_profile.get("bot") is True, + ) + await self.handle_message(MessageEvent( + text=text or "[Inline message with no text]", message_type=MessageType.TEXT, + source=source, raw_message=event, + message_id=build_inline_agent_action_turn_id( + agent_action.get("messageId"), agent_action.get("interactionId"), + ) if agent_action else msg_id, + timestamp=self._timestamp(event.get("date") or msg.get("date")), + allow_gateway_control=not bool(agent_action), + )) + return + media_text, media_urls, media_types, message_type = await self._normalize_media(msg) + if media_text: + text = f"{text}\n{media_text}".strip() if text else media_text + if not text and not media_urls: + text = "[Inline message with no text]" + if context_only: + self._remember_observed_context(chat_id, msg, text) + return if not agent_action and await self._handle_thread_command( chat_id=chat_id, msg_id=msg_id, + from_id=from_id, text=text, chat_type=chat_type, thread_id=thread_id, @@ -2584,6 +2738,7 @@ async def _dispatch_message(self, event: Dict[str, Any], *, edit: bool = False) chat_type=chat_type, user_id=from_id, user_name=sender_name or None, + is_bot=sender_profile.get("bot") is True, thread_id=thread_id, parent_chat_id=parent_chat_id, message_id=msg_id, @@ -2618,7 +2773,6 @@ async def _dispatch_message(self, event: Dict[str, Any], *, edit: bool = False) async def _dispatch_agent_action(self, event: Dict[str, Any]) -> None: interaction_id = str(event.get("interactionId") or "") - await self._answer_action(interaction_id, "") chat_id = str(event.get("chatId") or "") message_id = str(event.get("messageId") or "") @@ -2626,13 +2780,19 @@ async def _dispatch_agent_action(self, event: Dict[str, Any]) -> None: if not chat_id or not message_id or not interaction_id or not actor_user_id: logger.warning("[inline] ignored incomplete agent action event") return - target = await self._fetch_message(chat_id, message_id) + try: + target = await asyncio.wait_for( + self._fetch_message(chat_id, message_id, required=True), + timeout=_CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS, + ) + except asyncio.TimeoutError as exc: + raise InlineInboundDeferred("action target lookup timed out") from exc if not target: - logger.info("[inline] ignored agent action for unavailable message %s", message_id) + await self._answer_action(interaction_id, "Action expired or no longer accessible.") return synthetic_message = dict(target) - for stale_key in ("actions", "attachments", "entities", "media", "reactions", "replies"): + for stale_key in ("actions", "attachments", "entities", "media", "reactions", "replies", "sender"): synthetic_message.pop(stale_key, None) synthetic_message.update({ "id": message_id, @@ -2653,6 +2813,7 @@ async def _dispatch_agent_action(self, event: Dict[str, Any]) -> None: "message": synthetic_message, "_inlineAgentAction": dict(event), }) + await self._answer_action(interaction_id, "") def _message_explicitly_mentions_me(self, msg: Dict[str, Any]) -> bool: if not self._me_id: @@ -2758,13 +2919,19 @@ async def _dispatch_reaction(self, event: Dict[str, Any], *, added: bool) -> Non if self._me_id and user_id == self._me_id: return key = f"{event.get('kind')}:{chat_id}:{message_id}:{user_id}:{emoji}:{event.get('seq') or ''}" - if self._is_duplicate(key): - return target = await self._fetch_message(chat_id, message_id) chat_type = self._chat_type_from_message(target or {"peerId": {"type": {"oneofKind": "chat"}}}) if not self._allowed(chat_type, user_id): return + scope = await self._message_scope(chat_id, target or {}, required=True) + if scope is None: + return + chat_type, thread_id, parent_chat_id = scope + if chat_type == "group" and not self._chat_allowed(chat_id, thread_id, parent_chat_id): + return + if self._is_duplicate(key): + return target_text = str((target or {}).get("message") or "") or None target_author = str((target or {}).get("fromId") or "") or None target_is_own = bool(self._me_id and target_author == self._me_id) @@ -2777,6 +2944,8 @@ async def _dispatch_reaction(self, event: Dict[str, Any], *, added: bool) -> Non chat_type=chat_type, user_id=user_id, user_name=_inline_sender_identity(_inline_sender_profile(event))[0] or None, + thread_id=thread_id, + parent_chat_id=parent_chat_id, message_id=key, ) await self.handle_message(MessageEvent( @@ -2813,7 +2982,13 @@ async def _dispatch_system_event(self, event: Dict[str, Any]) -> None: text = "message:history_cleared" if not chat_id or not text: return - if self._group_policy == "disabled": + if self._group_policy == "disabled" or (user_id and not self._allowed("group", user_id)): + return + scope = await self._message_scope(chat_id, {}, required=True) + if scope is None: + return + _, thread_id, parent_chat_id = scope + if not self._chat_allowed(chat_id, thread_id, parent_chat_id): return key = f"{kind}:{chat_id}:{user_id}:{event.get('seq') or ''}:{text}" if self._is_duplicate(key): @@ -2824,8 +2999,12 @@ async def _dispatch_system_event(self, event: Dict[str, Any]) -> None: chat_type="group", user_id=user_id or None, user_name=_inline_sender_identity(_inline_sender_profile(event))[0] or None, + thread_id=thread_id, + parent_chat_id=parent_chat_id, message_id=key, ) + # Delete/history events may have no actor. Preserve the host's own + # authorization of these events rather than inventing a sender grant. await self.handle_message(MessageEvent( text=text, message_type=MessageType.TEXT, @@ -3665,7 +3844,7 @@ def _reply_to_for_target(self, reply_to: Optional[str], target: Dict[str, str]) return None return str(reply_to) - async def _get_chat_info(self, chat_id: str) -> Dict[str, Any]: + async def _get_chat_info(self, chat_id: str, *, required: bool = False) -> Dict[str, Any]: target = _target_from_chat_id(chat_id) normalized = str(target.get("chatId") or "").strip() if not normalized: @@ -3679,7 +3858,16 @@ async def _get_chat_info(self, chat_id: str) -> Dict[str, Any]: data = await self._sidecar_call("/chat", {"target": {"chatId": normalized}}) result = data.get("result") if isinstance(data, dict) else None info = result if isinstance(result, dict) else {} - except Exception: + except Exception as exc: + if required: + if isinstance(exc, InlineSidecarError) and exc.error_kind in {"forbidden", "not_found"}: + raise + raise InlineInboundDeferred("chat metadata request failed") from exc + return {} + resolved_id = self._chat_info_id(info, "id") + if not info or (resolved_id is not None and resolved_id != normalized): + if required: + raise InlineInboundDeferred("chat metadata unavailable or mismatched") return {} self._chat_info_cache[normalized] = (now, info) self._chat_info_cache.move_to_end(normalized) @@ -3779,12 +3967,14 @@ def _allowed(self, chat_type: str, from_id: str) -> bool: return self._id_allowed(self._group_allow_from, from_id) return True - async def _fetch_message(self, chat_id: str, message_id: str) -> Optional[Dict[str, Any]]: + async def _fetch_message(self, chat_id: str, message_id: str, *, required: bool = False) -> Optional[Dict[str, Any]]: try: data = await self._sidecar_call("/messages", {"target": _target_from_chat_id(chat_id), "messageIds": [message_id]}) messages = (data.get("result") or {}).get("messages") or [] return messages[0] if messages else None - except Exception: + except Exception as exc: + if required and not (isinstance(exc, InlineSidecarError) and exc.error_kind in {"forbidden", "not_found"}): + raise InlineInboundDeferred("action target temporarily unavailable") from exc return None async def _handle_action(self, event: Dict[str, Any]) -> bool: @@ -3825,43 +4015,117 @@ async def _handle_action(self, event: Dict[str, Any]) -> bool: async def _action_allowed(self, event: Dict[str, Any]) -> bool: actor_id = str(event.get("actorUserId") or "").strip() interaction_id = str(event.get("interactionId") or "") - chat_type = await self._action_chat_type(event) - if self._actor_authorized(chat_type, actor_id): - return True - await self._answer_action(interaction_id, "Not authorized") - logger.info("[inline] blocked action actor=%s chat_type=%s action=%s", actor_id or "unknown", chat_type or "unknown", event.get("actionId") or "") + chat_id = str(event.get("chatId") or "") + message_id = str(event.get("messageId") or "") + msg = await self._fetch_message(chat_id, message_id) if chat_id and message_id else None + scope = await self._message_scope(chat_id, msg) if msg else None + verdict = None + if scope is not None: + chat_type, thread_id, parent_chat_id = scope + verdict = self._actor_authorization( + chat_type, actor_id, chat_id, thread_id=thread_id, parent_chat_id=parent_chat_id, + is_bot=_inline_sender_profile(event).get("bot") is True, + ) + if verdict is True: + return True + await self._answer_action(interaction_id, "Access check temporarily unavailable. Try again." if verdict is None else "Not authorized") + logger.info("[inline] blocked action actor=%s chat=%s action=%s", actor_id or "unknown", chat_id or "unknown", event.get("actionId") or "") return False - def _actor_authorized(self, chat_type: Optional[str], actor_id: str) -> bool: - if not actor_id or not chat_type: + async def _message_scope( + self, chat_id: str, msg: Dict[str, Any], *, required: bool = False, + ) -> Optional[tuple[str, Optional[str], Optional[str]]]: + """Resolve the same parent/thread identity used by normal message intake.""" + chat_type = self._chat_type_from_message(msg) + thread_id = self._thread_id_from_message(msg) + parent_chat_id = self._parent_chat_id_from_message(msg) + if chat_type == "group" and not thread_id: + if not parent_chat_id: + try: + info = await asyncio.wait_for( + self._get_chat_info(chat_id, required=required), + timeout=_INBOUND_ACCESS_LOOKUP_TIMEOUT_SECONDS if required else _CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS, + ) + except asyncio.TimeoutError as exc: + if required: + raise InlineInboundDeferred("chat metadata timed out") from exc + return None + if not info: + if required: + raise InlineInboundDeferred("chat metadata unavailable") + return None + parent_chat_id = self._chat_info_id(info, "parentChatId") + if parent_chat_id: + thread_id = chat_id + return chat_type, thread_id, parent_chat_id + + async def _chat_actor_authorized(self, chat_id: str, actor_id: str) -> Optional[bool]: + """Check the actual target when a child control changes parent-wide state.""" + try: + info = await asyncio.wait_for( + self._get_chat_info(chat_id), timeout=_CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS, + ) + except asyncio.TimeoutError: + return None + if not info: + return None + parent_chat_id = self._chat_info_id(info, "parentChatId") + return self._actor_authorization( + self._bot_settings_chat_type(info), actor_id, chat_id, + thread_id=chat_id if parent_chat_id else None, parent_chat_id=parent_chat_id, + ) + + def _actor_authorized( + self, chat_type: Optional[str], actor_id: str, chat_id: Optional[str] = None, *, + thread_id: Optional[str] = None, parent_chat_id: Optional[str] = None, is_bot: bool = False, + ) -> bool: + return self._actor_authorization( + chat_type, actor_id, chat_id, thread_id=thread_id, + parent_chat_id=parent_chat_id, is_bot=is_bot, + ) is True + + def _actor_authorization( + self, chat_type: Optional[str], actor_id: str, chat_id: Optional[str] = None, *, + thread_id: Optional[str] = None, parent_chat_id: Optional[str] = None, is_bot: bool = False, + ) -> Optional[bool]: + # The host may trust an adapter's allowlist policy on the assumption that + # intake already checked it. Always enforce our restrictions before asking. + if not actor_id or chat_type not in {"dm", "group"} or not self._allowed(chat_type, actor_id): return False + if chat_type == "group" and not self._chat_allowed(chat_id or "", thread_id, parent_chat_id): + return False + if getattr(self, "_authorization_check", None) is not None: + # Hermes owns pairing, profile routing, and runtime authorization. + # Unknown/error is unavailable, never a reason to bypass a wired checker. + runner = getattr(self, "gateway_runner", None) + if parent_chat_id and runner is not None: + # The three-argument native callback drops parent_chat_id. Inline + # child chats need it for parent-routed profiles, so ask the same + # host predicate with the complete source, retaining transport identity. + try: + source = self.build_source( + chat_id=chat_id or "", chat_type=chat_type, user_id=actor_id, + thread_id=thread_id, parent_chat_id=parent_chat_id, is_bot=is_bot, + ) + if getattr(source, "profile_route_rejected", False): + return False + verdict = runner._is_user_authorized_for_source(source) + return verdict if isinstance(verdict, bool) else None + except Exception: + logger.warning("[inline] thread authorization unavailable", exc_info=True) + return None + return self._is_sender_authorized( + actor_id, chat_type, chat_id, is_bot=is_bot, thread_id=thread_id, + ) + # Standalone adapters have no runner. Only explicit grants are usable; + # an open intake policy alone must not authorize credentialed side effects. if self._allow_all or _truthy(os.getenv("GATEWAY_ALLOW_ALL_USERS"), False): return True if self._id_allowed(self._parse_id_set(os.getenv("GATEWAY_ALLOWED_USERS")), actor_id): return True - if chat_type == "dm": - if self._dm_policy == "disabled": - return False - if self._dm_policy == "allowlist" or self._allow_from: - return self._id_allowed(self._allow_from, actor_id) - return False - if self._group_policy == "disabled": - return False - if self._id_allowed(self._group_allow_from, actor_id): - return True - if self._id_allowed(self._allow_from, actor_id): - return True - return False - - async def _action_chat_type(self, event: Dict[str, Any]) -> Optional[str]: - chat_id = str(event.get("chatId") or "") - message_id = str(event.get("messageId") or "") - if not chat_id or not message_id: - return None - msg = await self._fetch_message(chat_id, message_id) - if not msg: - return None - return self._chat_type_from_message(msg) + return self._id_allowed(self._allow_from, actor_id) or ( + chat_type == "group" and self._id_allowed(self._group_allow_from, actor_id) + ) async def _handle_clarify_action(self, event: Dict[str, Any]) -> bool: action_id = str(event.get("actionId") or "") @@ -3990,7 +4254,12 @@ async def _handle_thread_action(self, event: Dict[str, Any]) -> bool: if not target_chat_id: await self._answer_action(interaction_id, "Thread controls expired") return True - if not await self._thread_action_allowed(event, state): + if not await self._action_allowed(event): + return True + if target_chat_id != chat_id and not await self._chat_actor_authorized( + target_chat_id, str(event.get("actorUserId") or ""), + ): + await self._answer_action(interaction_id, "Parent chat access could not be confirmed. Check access and try again.") return True try: if choice == "reset": @@ -4012,24 +4281,6 @@ async def _handle_thread_action(self, event: Dict[str, Any]) -> bool: await self._answer_action(interaction_id, "Thread setting failed") return True - async def _thread_action_allowed(self, event: Dict[str, Any], state: Dict[str, Any]) -> bool: - actor_id = str(event.get("actorUserId") or "").strip() - interaction_id = str(event.get("interactionId") or "") - chat_type = await self._action_chat_type(event) - if self._actor_authorized(chat_type, actor_id): - return True - display_chat_id = self._chat_key(state.get("display_chat_id")) - target_chat_id = self._chat_key(state.get("target_chat_id")) - if chat_type and actor_id and self._allowed(chat_type, actor_id): - if chat_type == "dm": - return True - thread_id = display_chat_id if target_chat_id and display_chat_id != target_chat_id else None - if self._chat_allowed(display_chat_id or str(event.get("chatId") or ""), thread_id, target_chat_id): - return True - await self._answer_action(interaction_id, "Not authorized") - logger.info("[inline] blocked thread action actor=%s chat_type=%s action=%s", actor_id or "unknown", chat_type or "unknown", event.get("actionId") or "") - return False - @staticmethod def _is_model_picker_action(action_id: str) -> bool: return action_id.startswith(("mp:", "mpg:", "mm:", "mc:", "mg:", "mb:", "mx:")) diff --git a/plugins/hermes-agent/plugin/inline/cli.py b/plugins/hermes-agent/plugin/inline/cli.py index d29e061cf..ffef0a1d1 100644 --- a/plugins/hermes-agent/plugin/inline/cli.py +++ b/plugins/hermes-agent/plugin/inline/cli.py @@ -514,15 +514,17 @@ def _compatibility_status() -> dict: return {"ok": False, "reason": "plugin_load_failed"} if "inline" not in loaded.tools_registered: return {"ok": False, "reason": "tool_not_registered"} - # Older supported Hermes releases predate the moved-import scanner. - # Their loader still checks real imports and plugin registration above. + # The moved-import scanner is optional: hosts may predate it or retain + # only updater stubs after removing the compat layer. Actual loading and + # this installation's tool registration above remain mandatory. try: - from hermes_cli.plugin_compat import plugin_hits + import hermes_cli.plugin_compat as plugin_compat except ModuleNotFoundError as error: if error.name != "hermes_cli.plugin_compat": raise else: - if plugin_hits(loaded.manifest): + plugin_hits = getattr(plugin_compat, "plugin_hits", None) + if plugin_hits is not None and plugin_hits(loaded.manifest): return {"ok": False, "reason": "deprecated_imports"} return {"ok": True, "reason": "loaded", "pluginPath": str(plugin_dir)} except Exception: diff --git a/plugins/hermes-agent/plugin/inline/sidecar/index.mjs b/plugins/hermes-agent/plugin/inline/sidecar/index.mjs index 25e9bf64e..3719d443d 100644 --- a/plugins/hermes-agent/plugin/inline/sidecar/index.mjs +++ b/plugins/hermes-agent/plugin/inline/sidecar/index.mjs @@ -5599,10 +5599,13 @@ var require_websocket_server = __commonJS(function(exports, module) { import http from "node:http"; // src/sidecar/inbound-stream.ts +import { randomUUID } from "node:crypto"; + class InboundStream { consumer = null; stopped = false; changed = new Set; + pending = new Map; attach(consumer) { if (this.stopped) { consumer.end(); @@ -5621,6 +5624,13 @@ class InboundStream { this.wake(); previous?.end(); } + acknowledge(deliveryId) { + const receipt = this.pending.get(deliveryId); + if (!receipt) + return; + receipt.acknowledged = true; + this.wake(); + } close() { this.stopped = true; const previous = this.consumer; @@ -5629,54 +5639,41 @@ class InboundStream { previous?.end(); } async deliver(event) { - const line = JSON.stringify(event) + ` + if (!event || typeof event !== "object" || Array.isArray(event)) { + throw new Error("Inbound delivery requires an event object"); + } + const deliveryId = randomUUID(); + const line = JSON.stringify({ ...event, _inlineDeliveryId: deliveryId }) + ` `; - while (!this.stopped) { - const owner = this.consumer; - if (!owner) { - await new Promise((resolve) => this.changed.add(resolve)); - continue; - } - let cleanup = () => {}; - const drained = new Promise((resolve) => { - const finish = (ok) => { - cleanup(); - resolve(ok); - }; - const onDrain = () => finish(true); - const onChange = () => finish(false); - cleanup = () => { - owner.off("drain", onDrain); - owner.off("close", onChange); - owner.off("error", onChange); - this.changed.delete(onChange); - }; - owner.once("drain", onDrain); - owner.once("close", onChange); - owner.once("error", onChange); - this.changed.add(onChange); - }); - try { - if (owner.destroyed || owner.writableEnded) { - if (this.consumer === owner) - this.consumer = null; + const receipt = { acknowledged: false }; + this.pending.set(deliveryId, receipt); + let writtenTo = null; + try { + while (!this.stopped) { + if (receipt.acknowledged) + return; + const owner = this.consumer; + if (owner && owner !== writtenTo) { + try { + if (owner.destroyed || owner.writableEnded) { + if (this.consumer === owner) + this.consumer = null; + continue; + } + writtenTo = owner; + owner.write(line); + } catch { + if (this.consumer === owner) + this.consumer = null; + } continue; } - const accepted = owner.write(line); - if (this.consumer !== owner || this.stopped) - continue; - if (accepted) - return; - if (await drained && this.consumer === owner && !this.stopped) - return; - } catch { - if (this.consumer === owner) - this.consumer = null; - } finally { - cleanup(); + await new Promise((resolve) => this.changed.add(resolve)); } + throw new Error("Inbound stream closed before delivery completed"); + } finally { + this.pending.delete(deliveryId); } - throw new Error("Inbound stream closed before delivery completed"); } wake() { for (const resolve of this.changed) @@ -46757,7 +46754,7 @@ function readUsers(result, kind) { } // src/sidecar/telemetry.ts -import { randomUUID } from "node:crypto"; +import { randomUUID as randomUUID2 } from "node:crypto"; import { readFile as readFile2 } from "node:fs/promises"; var TELEMETRY_TIMEOUT_MS = 2000; var TELEMETRY_DEDUP_MS = 5 * 60000; @@ -46818,7 +46815,7 @@ function buildHermesTelemetryEvent(operation, error, options = {}) { const exception = error instanceof Error ? error : new Error(String(error ?? "Unknown error")); const frames = parseStack(exception, env); return { - event_id: randomUUID().replaceAll("-", ""), + event_id: randomUUID2().replaceAll("-", ""), timestamp: new Date().toISOString(), platform: "javascript", level: "error", @@ -47090,6 +47087,13 @@ async function handleRequest(req, res) { writeJson(res, 405, { ok: false, error: "method not allowed", errorKind: "bad_format" }); return; } + if (url.pathname === "/inbound/ack") { + const body2 = asRecord(await readJsonBody(req)); + const deliveryId = readRequiredString(body2, "deliveryId"); + inboundStream.acknowledge(deliveryId); + writeJson(res, 200, { ok: true, result: {} }); + return; + } if (!connected && url.pathname !== "/shutdown") { writeJson(res, 503, { ok: false, diff --git a/plugins/hermes-agent/src/sidecar/inbound-stream.test.ts b/plugins/hermes-agent/src/sidecar/inbound-stream.test.ts index a272bc3b2..c48a73464 100644 --- a/plugins/hermes-agent/src/sidecar/inbound-stream.test.ts +++ b/plugins/hermes-agent/src/sidecar/inbound-stream.test.ts @@ -7,47 +7,108 @@ const flush = async () => { for (let i = 0; i < 10; i++) await Promise.resolve() } +function readEvents(consumer: PassThrough): Array<{ seq: number; _inlineDeliveryId: string }> { + const events: Array<{ seq: number; _inlineDeliveryId: string }> = [] + let buffered = "" + consumer.on("data", (chunk) => { + buffered += chunk.toString() + let newline: number + while ((newline = buffered.indexOf("\n")) >= 0) { + events.push(JSON.parse(buffered.slice(0, newline))) + buffered = buffered.slice(newline + 1) + } + }) + return events +} + describe("InboundStream lifecycle", () => { - it.each(["close", "error", "replacement"])("replays pending delivery after old consumer %s", async (cause) => { + it("holds the receipt after socket write until handling is acknowledged", async () => { + const stream = new InboundStream() + const consumer = new PassThrough() + const events = readEvents(consumer) + stream.attach(consumer) + let complete = false + const pending = stream.deliver({ seq: 1, _inlineDeliveryId: "untrusted" }).then(() => { complete = true }) + await flush() + expect(events).toHaveLength(1) + expect(events[0]._inlineDeliveryId).not.toBe("untrusted") + expect(complete).toBe(false) + stream.acknowledge("unknown") + await flush() + expect(complete).toBe(false) + stream.acknowledge(events[0]._inlineDeliveryId) + await pending + // Lost HTTP responses may cause the same acknowledgement to be retried. + stream.acknowledge(events[0]._inlineDeliveryId) + expect(complete).toBe(true) + stream.close() + consumer.destroy() + }) + + it.each(["close", "error", "replacement"])("replays the same unacknowledged identity after consumer %s", async (cause) => { const stream = new InboundStream() - const old = new PassThrough({ highWaterMark: 1 }) + const old = new PassThrough() + const original = readEvents(old) stream.attach(old) let complete = false - const pending = stream.deliver({ seq: 1 }).then(() => { - complete = true - }) + const pending = stream.deliver({ seq: 1 }).then(() => { complete = true }) await flush() expect(complete).toBe(false) const replacement = new PassThrough() - let output = "" - const received = new Promise((resolve) => - replacement.on("data", (chunk) => { - output += chunk.toString() - if (output.endsWith('{"seq":2}\n')) resolve() - }) - ) + const replay = readEvents(replacement) if (cause === "close") { const closed = once(old, "close") old.destroy() await closed + } else if (cause === "error") { + old.emit("error", new Error("consumer failed")) } stream.attach(replacement) - if (cause === "error") old.emit("error", new Error("retired stream")) + await flush() + expect(replay).toEqual(original) + expect(complete).toBe(false) + stream.acknowledge(replay[0]._inlineDeliveryId) await pending - await stream.deliver({ seq: 2 }) - await received - expect(output).toBe('{"seq":1}\n{"seq":2}\n') + // A stale retired consumer must not disconnect the replacement. + old.emit("error", new Error("retired stream")) + const next = stream.deliver({ seq: 2 }) + await flush() + expect(replay.map((event) => event.seq)).toEqual([1, 2]) + expect(replay[1]._inlineDeliveryId).not.toBe(replay[0]._inlineDeliveryId) + stream.acknowledge(replay[1]._inlineDeliveryId) + await next expect(old.listenerCount("drain")).toBe(0) stream.close() old.destroy() replacement.destroy() }) - it("shutdown releases both backpressure and absent-consumer waits without acknowledgement", async () => { - for (const withConsumer of [false, true]) { + it("acknowledges unrelated concurrent deliveries independently", async () => { + const stream = new InboundStream() + const consumer = new PassThrough() + const events = readEvents(consumer) + stream.attach(consumer) + let firstComplete = false + const first = stream.deliver({ seq: 1 }).then(() => { firstComplete = true }) + const second = stream.deliver({ seq: 2 }) + await flush() + expect(events).toHaveLength(2) + expect(new Set(events.map((event) => event._inlineDeliveryId)).size).toBe(2) + stream.acknowledge(events[1]._inlineDeliveryId) + await second + expect(firstComplete).toBe(false) + stream.acknowledge(events[0]._inlineDeliveryId) + await first + stream.close() + consumer.destroy() + }) + + it("shutdown rejects pending receipts with absent, writable, or backpressured consumers", async () => { + for (const mode of ["absent", "writable", "backpressured"]) { const stream = new InboundStream() - const consumer = new PassThrough({ highWaterMark: 1 }) - if (withConsumer) stream.attach(consumer) + const consumer = new PassThrough({ highWaterMark: mode === "backpressured" ? 1 : 16384 }) + if (mode === "writable") readEvents(consumer) + if (mode !== "absent") stream.attach(consumer) const pending = stream.deliver({ seq: 1 }) const rejected = expect(pending).rejects.toThrow("closed") await flush() diff --git a/plugins/hermes-agent/src/sidecar/inbound-stream.ts b/plugins/hermes-agent/src/sidecar/inbound-stream.ts index 3fa678515..014d22e7e 100644 --- a/plugins/hermes-agent/src/sidecar/inbound-stream.ts +++ b/plugins/hermes-agent/src/sidecar/inbound-stream.ts @@ -1,12 +1,15 @@ +import { randomUUID } from "node:crypto" import type { Writable } from "node:stream" -/** Retains a pending line across consumer replacement. A close after write is - * uncertain, so replay keeps the original event identity for downstream dedup. +/** Holds the SDK receipt until Python acknowledges handling. Consumer replacement + * replays the same delivery identity; this is an at-least-once process handoff, + * not a durable exactly-once transaction with the gateway's effects. */ export class InboundStream { private consumer: Writable | null = null private stopped = false private readonly changed = new Set<() => void>() + private readonly pending = new Map() attach(consumer: Writable): void { if (this.stopped) { @@ -27,6 +30,16 @@ export class InboundStream { previous?.end() } + /** Repeated/unknown acknowledgements are harmless, including a retry after the + * first acknowledgement succeeded but its HTTP response was lost. + */ + acknowledge(deliveryId: string): void { + const receipt = this.pending.get(deliveryId) + if (!receipt) return + receipt.acknowledged = true + this.wake() + } + close(): void { this.stopped = true const previous = this.consumer @@ -36,48 +49,40 @@ export class InboundStream { } async deliver(event: unknown): Promise { - const line = JSON.stringify(event) + "\n" - while (!this.stopped) { - const owner = this.consumer - if (!owner) { - await new Promise((resolve) => this.changed.add(resolve)) - continue - } - let cleanup = () => {} - const drained = new Promise((resolve) => { - const finish = (ok: boolean) => { - cleanup() - resolve(ok) - } - const onDrain = () => finish(true) - const onChange = () => finish(false) - cleanup = () => { - owner.off("drain", onDrain) - owner.off("close", onChange) - owner.off("error", onChange) - this.changed.delete(onChange) - } - owner.once("drain", onDrain) - owner.once("close", onChange) - owner.once("error", onChange) - this.changed.add(onChange) - }) - try { - if (owner.destroyed || owner.writableEnded) { - if (this.consumer === owner) this.consumer = null + if (!event || typeof event !== "object" || Array.isArray(event)) { + throw new Error("Inbound delivery requires an event object") + } + const deliveryId = randomUUID() + const line = JSON.stringify({ ...event, _inlineDeliveryId: deliveryId }) + "\n" + const receipt = { acknowledged: false } + this.pending.set(deliveryId, receipt) + let writtenTo: Writable | null = null + try { + while (!this.stopped) { + if (receipt.acknowledged) return + const owner = this.consumer + if (owner && owner !== writtenTo) { + try { + if (owner.destroyed || owner.writableEnded) { + if (this.consumer === owner) this.consumer = null + continue + } + // Socket acceptance/drain is not handling acknowledgement. Each + // SDK-owned pending delivery writes only once per consumer, and the + // SDK bounds concurrent receipts/backpressure upstream. + writtenTo = owner + owner.write(line) + } catch { + if (this.consumer === owner) this.consumer = null + } continue } - const accepted = owner.write(line) - if (this.consumer !== owner || this.stopped) continue - if (accepted) return - if ((await drained) && this.consumer === owner && !this.stopped) return - } catch { - if (this.consumer === owner) this.consumer = null - } finally { - cleanup() + await new Promise((resolve) => this.changed.add(resolve)) } + throw new Error("Inbound stream closed before delivery completed") + } finally { + this.pending.delete(deliveryId) } - throw new Error("Inbound stream closed before delivery completed") } private wake(): void { diff --git a/plugins/hermes-agent/src/sidecar/index.ts b/plugins/hermes-agent/src/sidecar/index.ts index f3ace9393..26db4e997 100644 --- a/plugins/hermes-agent/src/sidecar/index.ts +++ b/plugins/hermes-agent/src/sidecar/index.ts @@ -261,6 +261,15 @@ async function handleRequest(req: IncomingMessage, res: ServerResponse) { return } + // Receipt completion must remain available while the upstream SDK is offline. + if (url.pathname === "/inbound/ack") { + const body = asRecord(await readJsonBody(req)) + const deliveryId = readRequiredString(body, "deliveryId") + inboundStream.acknowledge(deliveryId) + writeJson(res, 200, { ok: true, result: {} }) + return + } + if (!connected && url.pathname !== "/shutdown") { writeJson(res, 503, { ok: false, diff --git a/plugins/hermes-agent/tests/adapter-python.test.ts b/plugins/hermes-agent/tests/adapter-python.test.ts index 760720d91..0975c7a9e 100644 --- a/plugins/hermes-agent/tests/adapter-python.test.ts +++ b/plugins/hermes-agent/tests/adapter-python.test.ts @@ -89,6 +89,19 @@ class BasePlatformAdapter: self.platform = platform self.name = str(platform) self.connected = False + self._authorization_check = None + + def set_authorization_check(self, callback): + self._authorization_check = callback + + def _is_sender_authorized(self, user_id, chat_type=None, chat_id=None, **kwargs): + if self._authorization_check is None: + return None + try: + result = self._authorization_check(user_id, chat_type, chat_id, **kwargs) + return result if result is True or result is False else None + except Exception: + return None def truncate_message(self, text, max_len): return [text[i:i + max_len] for i in range(0, len(text), max_len)] or [""] @@ -306,6 +319,12 @@ from inline.message_actions import ( ) base_extra = {"token": "fake", "context_history_limit": 0} +# Behavioral fixtures explicitly trust their actors; authorization cases use base_extra. +trusted_extra = {**base_extra, "allow_all": True} + +async def root_chat_info(chat_id, **kwargs): + return {"id": chat_id, "peer": {"type": {"oneofKind": "chat"}}} + assert resolve_inline_message_action_ownership("agent:1:2").owner == "agent" assert resolve_inline_message_action_ownership("system:cl:abc:0").native_action_id == "cl:abc:0" assert resolve_inline_message_action_ownership("legacy").explicit is False @@ -1740,7 +1759,7 @@ assert _target_from_chat_id("chat:55") == {"chatId": "55"} async def assert_thread_bindings(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, "channel_prompts": {"thread:99": "Thread prompt", "10": "Parent prompt"}, "channel_skill_bindings": [ @@ -1748,6 +1767,7 @@ async def assert_thread_bindings(): {"id": "10", "skill": "parent"}, ], })) + adapter._get_chat_info = root_chat_info events = [] async def fake_handle_message(event): @@ -1788,7 +1808,7 @@ asyncio.run(assert_thread_bindings()) async def assert_reply_thread_chat_metadata(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, "channel_prompts": {"123": "Parent prompt"}, "channel_skill_bindings": [{"id": "123", "skill": "parent-skill"}], @@ -1798,7 +1818,7 @@ async def assert_reply_thread_chat_metadata(): async def fake_handle_message(event): events.append(event) - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): if chat_id == "456": return { "chatId": "456", @@ -1840,7 +1860,7 @@ asyncio.run(assert_reply_thread_chat_metadata()) async def assert_default_auto_reply_threads_keep_fresh_parent_messages_flat(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, })) events = [] @@ -1848,7 +1868,7 @@ async def assert_default_auto_reply_threads_keep_fresh_parent_messages_flat(): async def fake_handle_message(event): events.append(event) - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): assert chat_id == "10" return {"chatId": "10", "title": "New thread", "lastMsgId": "9001"} @@ -1879,7 +1899,7 @@ asyncio.run(assert_default_auto_reply_threads_keep_fresh_parent_messages_flat()) async def assert_mentioned_agent_projection(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, "channel_skill_bindings": [{"id": "10", "skills": ["triage", "analysis"]}], })) @@ -1889,7 +1909,7 @@ async def assert_mentioned_agent_projection(): async def fake_handle_message(event): events.append(event) - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): return {"chatId": chat_id, "title": "Planning"} async def fake_sidecar_call(path, body): @@ -1947,9 +1967,10 @@ asyncio.run(assert_mentioned_agent_projection()) async def assert_ambient_bot_message_is_context_only(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, })) + adapter._get_chat_info = root_chat_info adapter._me_id = "20" events = [] @@ -1976,9 +1997,10 @@ asyncio.run(assert_ambient_bot_message_is_context_only()) async def assert_unverified_sender_provenance_is_context_only(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, })) + adapter._get_chat_info = root_chat_info adapter._me_id = "20" events = [] @@ -2005,7 +2027,7 @@ asyncio.run(assert_unverified_sender_provenance_is_context_only()) async def assert_activated_agent_avoids_lookup(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, })) adapter._me_id = "20" @@ -2014,7 +2036,7 @@ async def assert_activated_agent_avoids_lookup(): async def fake_handle_message(event): events.append(event) - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): return {"chatId": chat_id, "title": "Planning"} async def fake_sidecar_call(path, body): @@ -2053,7 +2075,7 @@ asyncio.run(assert_activated_agent_avoids_lookup()) async def assert_forced_reply_thread_creation(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "reply_threads": "on", "require_mention": False, "channel_prompts": {"99": "Thread prompt", "10": "Parent prompt"}, @@ -2065,7 +2087,7 @@ async def assert_forced_reply_thread_creation(): async def fake_handle_message(event): events.append(event) - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): assert chat_id == "10" return {"chatId": "10", "title": "Parent room"} @@ -2079,7 +2101,7 @@ async def assert_forced_reply_thread_creation(): return {"ok": True, "result": {}} raise AssertionError(f"unexpected sidecar path {path}") - async def fake_fetch_message(chat_id, msg_id): + async def fake_fetch_message(chat_id, msg_id, **kwargs): assert chat_id == "10" assert msg_id == "6" return {"id": "6", "chatId": "10", "fromId": "u2", "message": "parent quote"} @@ -2154,7 +2176,8 @@ async def assert_forced_reply_thread_creation(): asyncio.run(assert_forced_reply_thread_creation()) async def assert_default_dm_reply_thread_creation(): - adapter = InlineAdapter(PlatformConfig(extra=base_extra)) + adapter = InlineAdapter(PlatformConfig(extra=trusted_extra)) + adapter._get_chat_info = root_chat_info events = [] calls = [] @@ -2210,7 +2233,7 @@ asyncio.run(assert_default_dm_reply_thread_creation()) async def assert_reply_threads_disabled_preserves_existing_threads(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, "reply_threads": False, })) @@ -2220,7 +2243,7 @@ async def assert_reply_threads_disabled_preserves_existing_threads(): async def fake_handle_message(event): events.append(event) - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): return {"chatId": chat_id, "title": f"Chat {chat_id}"} async def fake_sidecar_call(path, body): @@ -2263,7 +2286,7 @@ asyncio.run(assert_reply_threads_disabled_preserves_existing_threads()) async def assert_inline_entity_context(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, "reply_threads": False, })) @@ -2272,7 +2295,7 @@ async def assert_inline_entity_context(): async def fake_handle_message(event): events.append(event) - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): return {"chatId": chat_id, "title": f"Chat {chat_id}"} adapter.handle_message = fake_handle_message @@ -2330,7 +2353,7 @@ asyncio.run(assert_inline_entity_context()) async def assert_inline_thread_context_history(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, "reply_threads": False, "context_backfill": "selective", @@ -2342,7 +2365,7 @@ async def assert_inline_thread_context_history(): async def fake_handle_message(event): events.append(event) - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): if chat_id == "456": return { "chatId": "456", @@ -2434,7 +2457,7 @@ asyncio.run(assert_inline_thread_context_history()) async def assert_inline_reply_context_window(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": False, "reply_threads": False, "context_backfill": "selective", @@ -2446,13 +2469,13 @@ async def assert_inline_reply_context_window(): async def fake_handle_message(event): events.append(event) - async def fake_fetch_message(chat_id, message_id): + async def fake_fetch_message(chat_id, message_id, **kwargs): assert chat_id == "10" assert message_id == "50" return {"id": "50", "chatId": "10", "fromId": "u2", "message": "Can we ship this?"} - async def fake_get_chat_info(chat_id): - return {} + async def fake_get_chat_info(chat_id, **kwargs): + return {"id": chat_id} async def fake_sidecar_call(path, body): calls.append((path, body)) @@ -2498,13 +2521,14 @@ asyncio.run(assert_inline_reply_context_window()) async def assert_observed_context_buffer(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": True, "reply_threads": False, "context_backfill": "off", "observe_unmentioned_messages": True, "observed_context_limit": 2, })) + adapter._get_chat_info = root_chat_info events = [] async def fake_handle_message(event): @@ -2545,7 +2569,7 @@ asyncio.run(assert_observed_context_buffer()) async def assert_explicit_addressing_precedence(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "require_mention": True, "reply_threads": False, "context_backfill": "off", @@ -2560,10 +2584,10 @@ async def assert_explicit_addressing_precedence(): async def fake_handle_message(event): events.append(event) - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): return {"chatId": chat_id, "dialogFollowMode": 1} - async def fake_fetch_message(chat_id, message_id): + async def fake_fetch_message(chat_id, message_id, **kwargs): return {"id": message_id, "chatId": chat_id, "fromId": "999", "message": "bot reply"} async def fake_send(chat_id, content, reply_to=None, metadata=None, actions=None): @@ -2689,7 +2713,7 @@ async def assert_reply_thread_slash_command(): with tempfile.TemporaryDirectory() as tmp: settings_path = Path(tmp) / "settings.json" adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "settings_path": str(settings_path), "require_mention": True, })) @@ -2716,12 +2740,12 @@ async def assert_reply_thread_slash_command(): sidecar_calls.append((path, body)) return {"ok": True, "result": {}} - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): if chat_id == "99": return {"chatId": "99", "title": "Child thread", "parentChatId": "10"} return {"chatId": chat_id, "title": f"Chat {chat_id}"} - async def fake_fetch_message(chat_id, message_id): + async def fake_fetch_message(chat_id, message_id, **kwargs): if chat_id == "20": return {"peerId": {"type": {"oneofKind": "user", "user": {"userId": chat_id}}}} return {"peerId": {"type": {"oneofKind": "chat", "chat": {"chatId": chat_id}}}} @@ -2973,7 +2997,8 @@ async def assert_new_message_delivery_dedup(): await adapter._dispatch_message(delivery) return [item.text for item in events] - sequenced = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": False})) + sequenced = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": False})) + sequenced._get_chat_info = root_chat_info assert await capture_with(sequenced, [ event(20, 100, "original"), event(20, 100, "duplicate delivery"), @@ -2981,7 +3006,8 @@ async def assert_new_message_delivery_dedup(): event(22, 100, "replacement"), ]) == ["original", "replacement"] - legacy = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": False})) + legacy = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": False})) + legacy._get_chat_info = root_chat_info assert await capture_with(legacy, [ event(None, 100, "legacy original"), event(None, 100, "legacy duplicate"), @@ -3001,7 +3027,7 @@ async def assert_group_room_controls(): async def capture(event): events.append(event) - async def fetch_message(chat_id, message_id): + async def fetch_message(chat_id, message_id, **kwargs): return reply adapter.handle_message = capture @@ -3017,34 +3043,38 @@ async def assert_group_room_controls(): "peerId": {"peer": {"oneofKind": "chat"}}, } - restricted = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": False, "allowed_chats": "99"})) + restricted = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": False, "allowed_chats": "99"})) + restricted._get_chat_info = root_chat_info assert await run(restricted, base_msg) == [] - allowed = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": False, "allowed_chats": "10"})) + allowed = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": False, "allowed_chats": "10"})) + allowed._get_chat_info = root_chat_info assert len(await run(allowed, base_msg)) == 1 - thread_allowed = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": False, "allowed_chats": "99"})) + thread_allowed = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": False, "allowed_chats": "99"})) + thread_allowed._get_chat_info = root_chat_info thread_msg = {**base_msg, "replies": {"chatId": "99"}} assert len(await run(thread_allowed, thread_msg)) == 1 - async def child_thread_info(chat_id): + async def child_thread_info(chat_id, **kwargs): if chat_id == "456": return {"chatId": "456", "title": "Child thread", "parentChatId": "10"} if chat_id == "10": return {"chatId": "10", "title": "Parent room"} raise AssertionError(f"unexpected chat info {chat_id}") - parent_allowed = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": False, "allowed_chats": "10"})) + parent_allowed = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": False, "allowed_chats": "10"})) parent_allowed._get_chat_info = child_thread_info child_events = await run(parent_allowed, {**base_msg, "id": "room-msg-child", "chatId": "456"}) assert len(child_events) == 1 assert child_events[0].source.thread_id == "456" assert child_events[0].source.parent_chat_id == "10" - free = InlineAdapter(PlatformConfig(extra={**base_extra, "free_response_chats": "10"})) + free = InlineAdapter(PlatformConfig(extra={**trusted_extra, "free_response_chats": "10"})) + free._get_chat_info = root_chat_info assert len(await run(free, base_msg)) == 1 - async def followed_info(chat_id): + async def followed_info(chat_id, **kwargs): return { "chatId": chat_id, "title": f"Followed {chat_id}", @@ -3052,13 +3082,13 @@ async def assert_group_room_controls(): "followModeMentionEligible": True, } - followed = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": True})) + followed = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": True})) followed._get_chat_info = followed_info followed_events = await run(followed, base_msg) assert len(followed_events) == 1 assert followed_events[0].source.chat_id == "10" - async def followed_large_info(chat_id): + async def followed_large_info(chat_id, **kwargs): return { "chatId": chat_id, "title": f"Large followed {chat_id}", @@ -3069,15 +3099,16 @@ async def assert_group_room_controls(): "followModeMentionEligible": False, } - followed_large = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": True})) + followed_large = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": True})) followed_large._get_chat_info = followed_large_info assert len(await run(followed_large, base_msg)) == 1 - strict_followed = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": True, "strict_mention": True})) + strict_followed = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": True, "strict_mention": True})) strict_followed._get_chat_info = followed_info assert await run(strict_followed, base_msg) == [] - strict = InlineAdapter(PlatformConfig(extra={**base_extra, "strict_mention": True})) + strict = InlineAdapter(PlatformConfig(extra={**trusted_extra, "strict_mention": True})) + strict._get_chat_info = root_chat_info strict._me_id = "bot" own_reply = {"id": "parent", "fromId": "bot", "message": "answer"} assert await run(strict, {**base_msg, "replyToMsgId": "parent"}, reply=own_reply) == [] @@ -3086,7 +3117,7 @@ async def assert_group_room_controls(): asyncio.run(assert_group_room_controls()) async def assert_chat_info_cache_invalidation(): - adapter = InlineAdapter(PlatformConfig(extra={**base_extra, "require_mention": True})) + adapter = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": True})) adapter._me_id = "bot" following = {"value": True} calls = [] @@ -3176,7 +3207,8 @@ async def assert_chat_info_cache_invalidation(): asyncio.run(assert_chat_info_cache_invalidation()) async def assert_action_thread_targets(): - adapter = InlineAdapter(PlatformConfig(extra=base_extra)) + adapter = InlineAdapter(PlatformConfig(extra=trusted_extra)) + adapter._get_chat_info = root_chat_info calls = [] async def fake_send_sidecar(path, body): @@ -3406,6 +3438,7 @@ asyncio.run(assert_model_picker_flow()) async def assert_choice_picker_flow(): adapter = InlineAdapter(PlatformConfig(extra={**base_extra, "allow_all": True})) + adapter._get_chat_info = root_chat_info calls = [] answers = [] selected = [] @@ -3414,7 +3447,7 @@ async def assert_choice_picker_flow(): calls.append((path, body)) return SendResult(success=True, message_id=body.get("messageId") or "choice-message", raw_response=body) - async def fake_fetch_message(chat_id, message_id): + async def fake_fetch_message(chat_id, message_id, **kwargs): return {"peerId": {"type": {"oneofKind": "chat", "chat": {"chatId": chat_id}}}} async def fake_answer_action(interaction_id, toast): @@ -3615,6 +3648,7 @@ async def assert_choice_picker_flow(): assert selected == [("chat:10", "high"), ("user:123", "medium")] failed = InlineAdapter(PlatformConfig(extra=base_extra)) + failed._get_chat_info = root_chat_info async def failed_send_sidecar(path, body): return SendResult(success=False, error="send failed") @@ -3637,6 +3671,7 @@ asyncio.run(assert_choice_picker_flow()) async def assert_update_prompt_flow(): adapter = InlineAdapter(PlatformConfig(extra={**base_extra, "allow_all": True})) + adapter._get_chat_info = root_chat_info calls = [] answers = [] @@ -3644,7 +3679,7 @@ async def assert_update_prompt_flow(): calls.append((path, body)) return SendResult(success=True, message_id=body.get("messageId") or "update-message", raw_response=body) - async def fake_fetch_message(chat_id, message_id): + async def fake_fetch_message(chat_id, message_id, **kwargs): return {"peerId": {"type": {"oneofKind": "chat", "chat": {"chatId": chat_id}}}} async def fake_answer_action(interaction_id, toast): @@ -3769,6 +3804,7 @@ async def assert_update_prompt_flow(): hermes_constants.get_hermes_home = original_get_hermes_home failed = InlineAdapter(PlatformConfig(extra=base_extra)) + failed._get_chat_info = root_chat_info async def failed_send_sidecar(path, body): return SendResult(success=False, error="prompt delivery failed") @@ -4173,19 +4209,399 @@ async def assert_processing_reactions(): asyncio.run(assert_processing_reactions()) + +async def assert_host_authorization_boundaries(): + # No registered host and no explicit grant is default-deny for local work. + adapter = InlineAdapter(PlatformConfig(extra=base_extra)) + assert not adapter._actor_authorized("dm", "u1", "10") + adapter = InlineAdapter(PlatformConfig(extra={**base_extra, "allow_all": True})) + for verdict in (False, None, "yes", 1): + adapter.set_authorization_check(lambda *args, **kwargs: verdict) + assert not adapter._actor_authorized("dm", "u1", "10"), verdict + def broken_check(*args, **kwargs): + raise RuntimeError("authorization unavailable") + adapter.set_authorization_check(broken_check) + assert not adapter._actor_authorized("dm", "u1", "10") + checks = [] + def host_grant(user_id, chat_type, chat_id, **kwargs): + checks.append((user_id, chat_type, chat_id, kwargs)) + return True + adapter.set_authorization_check(host_grant) + adapter._dm_policy = "disabled" + assert not adapter._actor_authorized("dm", "u1", "10") + adapter._group_policy = "disabled" + assert not adapter._actor_authorized("group", "u1", "10") + assert checks == [] + adapter._group_policy = "open" + adapter._allowed_chats = {"parent"} + assert not adapter._actor_authorized("group", "u1", "blocked") + assert adapter._actor_authorized("group", "u1", "child", thread_id="child", parent_chat_id="parent", is_bot=True) + assert checks[-1] == ("u1", "group", "child", {"is_bot": True, "thread_id": "child"}) + + with tempfile.TemporaryDirectory() as tmp: + settings = Path(tmp) / "settings.json" + adapter = InlineAdapter(PlatformConfig(extra={ + **base_extra, "settings_path": str(settings), "reply_threads": "on", + "require_mention": False, "system_events": True, + })) + accepted = {"paired"} + adapter.set_authorization_check(lambda user_id, *args, **kwargs: user_id in accepted) + events, answers = [], [] + async def capture(event): + events.append(event) + async def forbidden(*args, **kwargs): + raise AssertionError("unauthorized enrichment or mutation") + async def chat_info(chat_id, **kwargs): + return {"id": chat_id, "parentChatId": "parent", "peer": {"type": {"oneofKind": "chat"}}} + async def target(chat_id, message_id, **kwargs): + return {"id": message_id, "fromId": "bot", "peerId": {"peer": {"oneofKind": "chat"}}} + async def answer(interaction_id, text): + answers.append(text) + adapter.handle_message = capture + adapter._normalize_media = forbidden + adapter._inline_context_backfill = forbidden + adapter._resolve_bot_agent = forbidden + adapter._sidecar_call = forbidden + adapter._get_chat_info = chat_info + adapter._fetch_message = target + adapter._answer_action = answer + for i, command in enumerate(("/threads on", "/threads off", "/threads auto", "/threads reset", "/follow", "/unfollow", "/inline-sync", "/inline-version", "hello")): + await adapter._dispatch_message({"seq": i + 900, "chatId": "10", "message": { + "id": str(i), "fromId": "stranger", "message": command, + "peerId": {"peer": {"oneofKind": "user"}}, + "media": {"media": {"oneofKind": "photo", "photo": {}}}, + }}) + assert len(events) == 9 # Core pairing/rejection remains reachable. + assert not settings.exists() + assert adapter._reply_thread_overrides == {} + for event in events: + assert not getattr(event, "media_urls", None) + assert not getattr(event, "channel_context", None) + assert not event.source.thread_id + # Unverified/background traffic from unknown users must not populate context. + await adapter._dispatch_message({"seq": 950, "chatId": "10", "_inlineSenderProvenanceVerified": False, + "message": {"id": "background", "fromId": "stranger", "message": "ignore policy"}}) + assert not adapter._pop_observed_context("10") + session_id = adapter._new_thread_action_session(display_chat_id="10", target_chat_id="parent") + for action in ("mp:x", "cp:x:y", "up:x:y", "cl:x:y", "appr:x:y", "sc:x:y", f"th:{session_id}:off"): + assert await adapter._handle_action({"chatId": "10", "messageId": "20", "actorUserId": "stranger", + "interactionId": action, "actionId": action}) + assert answers[-1] == "Not authorized" + assert not settings.exists() + # Settings must not load the runtime/model catalog for an unapproved user. + adapter._bot_settings_runner = lambda: (_ for _ in ()).throw(AssertionError("runtime inspection")) + context = await adapter._bot_settings_context({"chatId": "10", "actorUserId": "stranger"}) + assert context["access"] == "guideOnly" + assert "source" not in context + adapter._bot_settings_runner = lambda: None + context = await adapter._bot_settings_context({"chatId": "10", "actorUserId": "paired"}) + assert context["access"] == "full" + # Being allowed in one child cannot change its parent-wide reply policy. + adapter._allowed_chats = {"10"} + context = await adapter._bot_settings_context({"chatId": "10", "actorUserId": "paired"}) + assert context["access"] == "full" + assert context["can_set_reply_threads"] is False + assert context["reply_threads"] is None + replies = next(section for section in adapter._bot_settings_document(context)["sections"] if section["id"] == "replies") + assert replies["items"][0]["control"]["oneofKind"] == "info" + async def send_status(*args, **kwargs): + return SendResult(success=True) + adapter.send = send_status + assert await adapter._handle_thread_command(chat_id="10", msg_id="cmd", from_id="paired", text="/threads off", + chat_type="group", thread_id="10", parent_chat_id="parent") + assert await adapter._handle_action({"chatId": "10", "messageId": "20", "actorUserId": "paired", + "interactionId": "child-only", "actionId": f"th:{session_id}:off"}) + assert answers[-1] == "Parent chat access could not be confirmed. Check access and try again." + try: + await adapter._apply_bot_setting({"itemId": "reply-threads", "value": {"value": { + "oneofKind": "stringValue", "stringValue": "off", + }}}, context) + except ValueError: + pass + else: + raise AssertionError("child settings changed parent-wide state") + assert not settings.exists() + # Parent group grants apply to child-thread buttons, but never excluded rooms. + adapter._allowed_chats = {"parent"} + button = {"chatId": "10", "messageId": "20", "actorUserId": "paired", "interactionId": "paired"} + assert await adapter._action_allowed(button) + adapter._allowed_chats = {"other"} + assert not await adapter._action_allowed(button) + before = len(events) + adapter._me_id = "bot" + await adapter._dispatch_reaction({"chatId": "10", "messageId": "20", "userId": "paired", "emoji": "ok"}, added=True) + await adapter._dispatch_system_event({"kind": "chat.participant.add", "chatId": "10", "userId": "paired"}) + assert len(events) == before + adapter._allowed_chats = {"parent"} + # Actorless lifecycle input retains core authorization after local scope, + # including sender-allowlisted groups (the host may grant by chat). + adapter._group_policy = "allowlist" + adapter._group_allow_from = {"paired"} + await adapter._dispatch_system_event({"kind": "message.delete", "chatId": "10", "messageIds": ["20"]}) + assert len(events) == before + 1 + assert events[-1].source.user_id is None + assert events[-1].source.parent_chat_id == "parent" + adapter._group_policy = "open" + # Opted-in reactions survive an unavailable target; core decides admission. + async def missing_target(*args, **kwargs): + return None + adapter._fetch_message = missing_target + await adapter._dispatch_reaction({"chatId": "10", "messageId": "gone", "userId": "paired", "emoji": "ok"}, added=True) + assert len(events) == before + 2 + adapter._fetch_message = target + # Parent resolution is bounded and never grants access after timeout. + original_timeout = inline_adapter_module._CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS + async def slow_chat(*args, **kwargs): + await asyncio.sleep(1) + adapter._get_chat_info = slow_chat + try: + inline_adapter_module._CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS = 0.01 + assert not await adapter._action_allowed(button) + finally: + inline_adapter_module._CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS = original_timeout + adapter._get_chat_info = chat_info + async def missing_chat(*args, **kwargs): + return {} + adapter._get_chat_info = missing_chat + assert not await adapter._action_allowed(button) + before = len(events) + try: + await adapter._dispatch_message({"seq": 999, "chatId": "10", "message": { + "id": "unknown-parent", "fromId": "paired", "message": "/follow", + "peerId": {"peer": {"oneofKind": "chat"}}, + }}) + except inline_adapter_module.InlineInboundDeferred: + pass + else: + raise AssertionError("missing routing metadata was not deferred") + assert len(events) == before + adapter._get_chat_info = chat_info + accepted.clear() + assert not await adapter._action_allowed(button) + context = await adapter._bot_settings_context({"chatId": "10", "actorUserId": "paired"}) + assert context["access"] == "guideOnly" + +asyncio.run(assert_host_authorization_boundaries()) + +async def assert_inbound_receipt_recovery(): + original_delay = inline_adapter_module._INBOUND_RETRY_INITIAL_SECONDS + original_timeout = inline_adapter_module._INBOUND_ACCESS_LOOKUP_TIMEOUT_SECONDS + inline_adapter_module._INBOUND_RETRY_INITIAL_SECONDS = 0.001 + inline_adapter_module._INBOUND_ACCESS_LOOKUP_TIMEOUT_SECONDS = 0.01 + def event(chat, seq, delivery): + return {"kind": "message.new", "chatId": chat, "seq": seq, "_inlineDeliveryId": delivery, + "message": {"id": str(seq), "fromId": "paired", "message": "hello", "mentioned": True, + "peerId": {"type": {"oneofKind": "chat"}}, + "media": {"media": {"oneofKind": "photo", "photo": {}}}}} + try: + adapter = InlineAdapter(PlatformConfig(extra={**trusted_extra, "require_mention": False, + "reply_threads": "off", "context_backfill": "off"})) + available = asyncio.Event() + delivered, normalized, acknowledgements = [], [], [] + ack_failed = False + async def sidecar(path, body): + nonlocal ack_failed + if path == "/chat": + chat = body["target"]["chatId"] + if chat == "10" and not available.is_set(): + await available.wait() + return {"ok": True, "result": {"id": chat}} + assert path == "/inbound/ack", path + acknowledgements.append(body["deliveryId"]) + if body["deliveryId"] == "held" and not ack_failed: + ack_failed = True + raise RuntimeError("ACK response lost") + return {"ok": True} + async def media(msg): + normalized.append(msg["id"]) + return ("[photo]", ["/tmp/receipt-test-photo"], ["image/jpeg"], MessageType.PHOTO) + async def capture(ev): + delivered.append(ev) + adapter._sidecar_call = sidecar + adapter._normalize_media = media + adapter.handle_message = capture + held = event("10", 100, "held") + await adapter._on_inbound(json.dumps(held)) + first = adapter._inbound_deliveries["held"] + await adapter._on_inbound(json.dumps(held)) + assert adapter._inbound_deliveries["held"] is first + await adapter._on_inbound(json.dumps(event("20", 200, "healthy"))) + await asyncio.wait_for(adapter._inbound_deliveries["healthy"], 1) + await asyncio.sleep(0.015) + assert [e.source.chat_id for e in delivered] == ["20"] + assert normalized == ["200"] + assert "held" not in acknowledgements + assert not any("10:" in key for key in adapter._seen_messages) + available.set() + await asyncio.wait_for(first, 1) + assert [e.source.chat_id for e in delivered] == ["20", "10"] + assert delivered[-1].media_urls == ["/tmp/receipt-test-photo"] + assert delivered[-1].message_id == "100" + assert normalized == ["200", "100"] + assert acknowledgements.count("held") == 2 + # A fresh receipt for the same update still deduplicates completed work. + held["_inlineDeliveryId"] = "replayed" + await adapter._on_inbound(json.dumps(held)) + await asyncio.wait_for(adapter._inbound_deliveries["replayed"], 1) + assert len(delivered) == 2 + # Cancellation before authorization does not burn either dedup identity. + available.clear() + cancelled = event("10", 101, "cancelled") + await adapter._on_inbound(json.dumps(cancelled)) + task = adapter._inbound_deliveries["cancelled"] + await asyncio.sleep(0) + task.cancel() + await asyncio.gather(task, return_exceptions=True) + await asyncio.sleep(0) + assert "cancelled" not in acknowledgements + available.set() + await adapter._on_inbound(json.dumps(cancelled)) + await asyncio.wait_for(adapter._inbound_deliveries["cancelled"], 1) + assert [e.message_id for e in delivered] == ["200", "100", "101"] + # Native unknown is temporary, not a denial that strips incoming media. + verdicts = iter([None, True]) + adapter.set_authorization_check(lambda *a, **kw: next(verdicts)) + await adapter._on_inbound(json.dumps(event("30", 300, "auth-recovery"))) + await asyncio.wait_for(adapter._inbound_deliveries["auth-recovery"], 1) + assert delivered[-1].media_urls == ["/tmp/receipt-test-photo"] + assert normalized.count("300") == 1 + # Host exceptions preserve a local command instead of sending it to the model. + auth_calls, commands_seen = [], [] + def interrupted_auth(*args, **kwargs): + auth_calls.append(True) + if len(auth_calls) == 1: + raise RuntimeError("temporary host outage") + return True + async def text_only(msg): + return ("", [], [], MessageType.TEXT) + async def command(**kwargs): + commands_seen.append(kwargs["text"]) + return True + adapter.set_authorization_check(interrupted_auth) + adapter._normalize_media = text_only + original_command = adapter._handle_thread_command + adapter._handle_thread_command = command + local_command = event("30", 301, "command-recovery") + local_command["message"]["message"] = "/threads on" + await adapter._on_inbound(json.dumps(local_command)) + await asyncio.wait_for(adapter._inbound_deliveries["command-recovery"], 1) + assert commands_seen == ["/threads on"] + assert len(delivered) == 4 + adapter._handle_thread_command = original_command + # A transient callback-target fetch must not consume the button press. + adapter.set_authorization_check(lambda *a, **kw: True) + target_calls, action_answers = [], [] + async def action_sidecar(path, body): + if path == "/messages": + target_calls.append(True) + if len(target_calls) == 1: + raise RuntimeError("target lookup offline") + return {"ok": True, "result": {"messages": [{"id": "20", "fromId": "bot", + "message": "choose", "peerId": {"type": {"oneofKind": "user"}}}]}} + if path == "/answer-action": + action_answers.append(body) + elif path == "/inbound/ack": + acknowledgements.append(body["deliveryId"]) + else: + raise AssertionError(path) + return {"ok": True} + adapter._sidecar_call = action_sidecar + await adapter._on_inbound(json.dumps({"kind": "message.action.invoke", "chatId": "60", + "messageId": "20", "actorUserId": "paired", "interactionId": "click", + "actionId": "agent:1:1", "_inlineDeliveryId": "callback-recovery"})) + await asyncio.wait_for(adapter._inbound_deliveries["callback-recovery"], 1) + assert len(target_calls) >= 2 + assert len(action_answers) == 1 + assert delivered[-1].message_id == "inline-agent-action:20:click" + assert delivered[-1].allow_gateway_control is False + assert acknowledgements.count("callback-recovery") == 1 + # A persistently unavailable button must release the SDK user barrier. + blocked = InlineAdapter(PlatformConfig(extra=trusted_extra)) + retries, toasts, retired = [], [], [] + async def always_unavailable(event): + retries.append(event) + raise inline_adapter_module.InlineInboundDeferred("offline") + async def toast(interaction, text): + toasts.append((interaction, text)) + async def retire(path, body): + assert path == "/inbound/ack" + retired.append(body["deliveryId"]) + blocked._dispatch_inbound = always_unavailable + blocked._answer_action = toast + blocked._sidecar_call = retire + await blocked._on_inbound(json.dumps({"kind": "message.action.invoke", "interactionId": "blocked", + "_inlineDeliveryId": "blocked"})) + await asyncio.wait_for(blocked._inbound_deliveries["blocked"], 1) + assert len(retries) == 2 + assert toasts == [("blocked", "Temporarily unavailable. Please try again.")] + assert retired == ["blocked"] + # Lost ACK responses can outlive SDK slots; Python admission is bounded. + bounded = InlineAdapter(PlatformConfig(extra=trusted_extra)) + ack_response = asyncio.Event() + async def handled(event): + pass + async def lost_response(path, body): + await ack_response.wait() + bounded._dispatch_inbound = handled + bounded._sidecar_call = lost_response + old_limit = inline_adapter_module._MAX_INBOUND_DELIVERIES + inline_adapter_module._MAX_INBOUND_DELIVERIES = 2 + try: + for receipt in ("one", "two"): + await bounded._on_inbound(json.dumps({"_inlineDeliveryId": receipt})) + third = asyncio.create_task(bounded._on_inbound(json.dumps({"_inlineDeliveryId": "three"}))) + await asyncio.sleep(0) + assert not third.done() + assert len(bounded._inbound_deliveries) == 2 + await bounded._on_inbound(json.dumps({"_inlineDeliveryId": "one"})) + ack_response.set() + await asyncio.wait_for(third, 1) + await asyncio.gather(*list(bounded._inbound_deliveries.values())) + finally: + inline_adapter_module._MAX_INBOUND_DELIVERIES = old_limit + # Empty/mismatched successful responses must not poison the cache. + responses = iter([{}, {"id": "wrong"}, {"id": "40"}]) + async def snapshots(path, body): + return {"ok": True, "result": next(responses)} + adapter._sidecar_call = snapshots + assert await adapter._get_chat_info("40") == {} + assert await adapter._get_chat_info("40") == {} + assert await adapter._get_chat_info("40") == {"id": "40"} + # Permanent loss of chat access must not poison the same-chat queue. + adapter.set_authorization_check(lambda *a, **kw: True) + async def inaccessible(path, body): + if path == "/chat": + raise inline_adapter_module.InlineSidecarError(path, 403, "forbidden", "forbidden") + acknowledgements.append(body["deliveryId"]) + return {"ok": True} + adapter._sidecar_call = inaccessible + await adapter._on_inbound(json.dumps(event("50", 500, "removed"))) + await asyncio.wait_for(adapter._inbound_deliveries["removed"], 1) + assert "removed" in acknowledgements + assert "500" not in normalized + finally: + inline_adapter_module._INBOUND_RETRY_INITIAL_SECONDS = original_delay + inline_adapter_module._INBOUND_ACCESS_LOOKUP_TIMEOUT_SECONDS = original_timeout + +asyncio.run(assert_inbound_receipt_recovery()) + async def assert_action_authorization(): os.environ.pop("GATEWAY_ALLOWED_USERS", None) os.environ.pop("GATEWAY_ALLOW_ALL_USERS", None) adapter = InlineAdapter(PlatformConfig(extra={**base_extra, "group_policy": "allowlist", "group_allow_from": "u1"})) answers = [] - async def fake_fetch_message(chat_id, message_id): + async def fake_fetch_message(chat_id, message_id, **kwargs): return {"peerId": {"type": {"oneofKind": "chat", "chat": {"chatId": chat_id}}}} + async def fake_chat_info(chat_id, **kwargs): + return {"id": chat_id} + async def fake_answer_action(interaction_id, toast): answers.append((interaction_id, toast)) adapter._fetch_message = fake_fetch_message + adapter._get_chat_info = fake_chat_info adapter._answer_action = fake_answer_action adapter._approval_sessions["approval-1"] = "session-1" @@ -4217,11 +4633,12 @@ async def assert_action_authorization(): "actionId": "appr:approval-1:approve", }) assert not unknown_context - assert answers[-1] == ("interaction-3", "Not authorized") + assert answers[-1] == ("interaction-3", "Access check temporarily unavailable. Try again.") open_adapter = InlineAdapter(PlatformConfig(extra=base_extra)) open_answers = [] open_adapter._fetch_message = fake_fetch_message + open_adapter._get_chat_info = fake_chat_info async def fake_open_answer_action(interaction_id, toast): open_answers.append((interaction_id, toast)) @@ -4239,6 +4656,7 @@ async def assert_action_authorization(): os.environ["GATEWAY_ALLOWED_USERS"] = "u1" gateway_allowed = InlineAdapter(PlatformConfig(extra=base_extra)) gateway_allowed._fetch_message = fake_fetch_message + gateway_allowed._get_chat_info = fake_chat_info assert await gateway_allowed._action_allowed({ "chatId": "10", "messageId": "20", @@ -4250,6 +4668,7 @@ async def assert_action_authorization(): inline_allowed = InlineAdapter(PlatformConfig(extra={**base_extra, "allow_from": "u1"})) inline_allowed._fetch_message = fake_fetch_message + inline_allowed._get_chat_info = fake_chat_info assert await inline_allowed._action_allowed({ "chatId": "10", "messageId": "20", @@ -4355,7 +4774,8 @@ async def assert_callback_state_lifecycle(): asyncio.run(assert_callback_state_lifecycle()) async def assert_agent_action_turn_and_same_message_response(): - adapter = InlineAdapter(PlatformConfig(extra={**base_extra, "allow_all": True})) + adapter = InlineAdapter(PlatformConfig(extra={**trusted_extra, "allow_all": True})) + adapter._get_chat_info = root_chat_info adapter._me_id = "bot" answers = [] events = [] @@ -4366,11 +4786,12 @@ async def assert_agent_action_turn_and_same_message_response(): order.append("ack") answers.append((interaction_id, toast)) - async def fake_fetch_message(chat_id, message_id): + async def fake_fetch_message(chat_id, message_id, **kwargs): return { "id": message_id, "chatId": chat_id, "fromId": "bot", + "sender": {"id": "bot", "bot": True}, "message": "Approve proposal 17?", "out": True, "peerId": {"peer": {"oneofKind": "user", "user": {"userId": "u1"}}}, @@ -4412,7 +4833,7 @@ async def assert_agent_action_turn_and_same_message_response(): })) assert answers == [("30", "")] - assert order == ["ack", "turn"] + assert order == ["turn", "ack"] assert len(events) == 1 action_turn = events[0] assert action_turn.text.startswith("Inline action button pressed on message 20.") @@ -4459,17 +4880,27 @@ async def assert_agent_action_turn_and_same_message_response(): assert answers[-1] == ("31", "Action expired") assert len(events) == 1 + await adapter._on_inbound(json.dumps({ + "kind": "message.action.invoke", "seq": 502, "chatId": "10", "messageId": "20", + "interactionId": "32", "actorUserId": "u1", "actionId": "agent:1:1", + })) + assert len(events) == 2 + assert events[-1].source.user_id == "u1" + assert events[-1].source.is_bot is False + assert events[-1].allow_gateway_control is False + asyncio.run(assert_agent_action_turn_and_same_message_response()) async def assert_inline_lifecycle_events(): - adapter = InlineAdapter(PlatformConfig(extra={**base_extra, "group_policy": "open"})) + adapter = InlineAdapter(PlatformConfig(extra={**trusted_extra, "group_policy": "open"})) + adapter._get_chat_info = root_chat_info adapter._me_id = "bot" events = [] async def capture(event): events.append(event) - async def own_message(chat_id, message_id): + async def own_message(chat_id, message_id, **kwargs): return { "id": message_id, "fromId": "bot", @@ -4496,7 +4927,7 @@ async def assert_inline_lifecycle_events(): assert events[0].source.chat_id == "10" assert events[0].source.user_id == "u1" - async def human_message(chat_id, message_id): + async def human_message(chat_id, message_id, **kwargs): return { "id": message_id, "fromId": "u2", @@ -4514,7 +4945,8 @@ async def assert_inline_lifecycle_events(): })) assert len(events) == 1 - system_adapter = InlineAdapter(PlatformConfig(extra={**base_extra, "system_events": True})) + system_adapter = InlineAdapter(PlatformConfig(extra={**trusted_extra, "system_events": True})) + system_adapter._get_chat_info = root_chat_info system_adapter._me_id = "bot" system_events = [] @@ -4555,7 +4987,7 @@ asyncio.run(assert_inline_lifecycle_events()) async def assert_join_mention_recovery(): adapter = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "group_policy": "open", "require_mention": True, "context_backfill": "off", @@ -4608,7 +5040,7 @@ async def assert_join_mention_recovery(): return {"ok": True, "result": {"messages": [recent, boundary, too_old]}} return {"ok": True, "result": {}} - async def fake_get_chat_info(chat_id): + async def fake_get_chat_info(chat_id, **kwargs): return {"chatId": chat_id, "title": "Project Room"} adapter.handle_message = capture @@ -4634,7 +5066,7 @@ async def assert_join_mention_recovery(): assert [event.message_id for event in events] == ["5000", "5001"] paged = InlineAdapter(PlatformConfig(extra={ - **base_extra, + **trusted_extra, "group_policy": "open", "require_mention": True, "context_backfill": "off", @@ -5038,6 +5470,7 @@ async def assert_bot_settings_document_and_mutation(): "is_reply_thread": False, "following": True, "reply_threads": "auto", + "can_set_reply_threads": True, "runner": None, "source": None, "model_options": [], @@ -5086,12 +5519,12 @@ asyncio.run(assert_bot_settings_document_and_mutation()) async def assert_bot_settings_fail_closed_and_serialized(): adapter = InlineAdapter(PlatformConfig(extra=base_extra)) - async def slow_chat_info(chat_id): + async def slow_chat_info(chat_id, **kwargs): await asyncio.sleep(30) return {"id": chat_id} - original_timeout = inline_adapter_module._BOT_SETTINGS_CHAT_INFO_TIMEOUT_SECONDS - inline_adapter_module._BOT_SETTINGS_CHAT_INFO_TIMEOUT_SECONDS = 0.01 + original_timeout = inline_adapter_module._CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS + inline_adapter_module._CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS = 0.01 try: adapter._get_chat_info = slow_chat_info started_at = time.monotonic() @@ -5100,9 +5533,9 @@ async def assert_bot_settings_fail_closed_and_serialized(): assert context["access"] == "guideOnly" assert context["unavailable_reason"] == "chat_metadata" finally: - inline_adapter_module._BOT_SETTINGS_CHAT_INFO_TIMEOUT_SECONDS = original_timeout + inline_adapter_module._CHAT_ACCESS_LOOKUP_TIMEOUT_SECONDS = original_timeout - async def missing_chat_info(chat_id): + async def missing_chat_info(chat_id, **kwargs): return {} adapter._get_chat_info = missing_chat_info @@ -5112,22 +5545,17 @@ async def assert_bot_settings_fail_closed_and_serialized(): guide = adapter._bot_settings_document(context) assert "could not verify this chat" in guide["sections"][0]["items"][0]["control"]["info"]["text"] - async def allowed_chat_info(chat_id): + async def allowed_chat_info(chat_id, **kwargs): return {"id": chat_id, "peer": {"type": {"oneofKind": "chat"}}} adapter._get_chat_info = allowed_chat_info adapter._allowed = lambda chat_type, actor_id: True adapter._chat_allowed = lambda chat_id, thread_id, parent_chat_id=None: True - adapter._actor_authorized = lambda chat_type, actor_id: False + adapter._actor_authorized = lambda *args, **kwargs: False context = await adapter._bot_settings_context({"chatId": "42", "actorUserId": "u1"}) - assert context["access"] == "readOnly" - readonly = adapter._bot_settings_document(context) - assert all( - item.get("disabled") is True - for section in readonly["sections"][:3] - for item in section["items"] - if item["id"] not in {"runtime-unavailable"} - ) + assert context["access"] == "guideOnly" + guide = adapter._bot_settings_document(context) + assert [section["id"] for section in guide["sections"]] == ["access"] active = 0 max_active = 0 diff --git a/plugins/hermes-agent/tests/inbound-flow.test.ts b/plugins/hermes-agent/tests/inbound-flow.test.ts index 603fd663f..56dab278d 100644 --- a/plugins/hermes-agent/tests/inbound-flow.test.ts +++ b/plugins/hermes-agent/tests/inbound-flow.test.ts @@ -93,15 +93,15 @@ it.each(["recover", "abort"])("real SDK, directory and stream isolate held sende deliver: (event) => stream.deliver(event), }) ) - const message = (chatId: bigint, mention: boolean) => + const message = (chatId: bigint, mention: boolean, seq = 2) => Update.create({ - seq: 2, + seq, date: 100n, update: { oneofKind: "newMessage", newMessage: { message: { - id: 2n, + id: BigInt(seq), chatId, fromId: 42n, peerId: { type: { oneofKind: "chat", chat: { chatId } } }, @@ -119,13 +119,18 @@ it.each(["recover", "abort"])("real SDK, directory and stream isolate held sende body: { oneofKind: "message", message: { - payload: { oneofKind: "update", update: { updates: [message(10n, false), message(20n, true)] } }, + payload: { oneofKind: "update", update: { updates: [message(10n, false), message(20n, true), message(10n, false, 3)] } }, }, }, }) ) await flush() expect(lines.map((line) => JSON.parse(line).chatId)).toEqual(["20"]) + // Bytes reaching Python are not application completion. A healthy chat can + // acknowledge while the other chat's sender/preflight is still pending. + expect(client.exportState().lastSeqByChatId).toEqual({ "10": 1, "20": 1 }) + stream.acknowledge(JSON.parse(lines[0]!)._inlineDeliveryId) + await flush() expect(client.exportState().lastSeqByChatId).toEqual({ "10": 1, "20": 2 }) if (action === "abort") { abort.abort() @@ -145,7 +150,15 @@ it.each(["recover", "abort"])("real SDK, directory and stream isolate held sende }) await flush() expect(lines.map((line) => JSON.parse(line).chatId)).toEqual(["20", "10"]) + expect(client.exportState().lastSeqByChatId).toEqual({ "10": 1, "20": 2 }) + stream.acknowledge(JSON.parse(lines[1]!)._inlineDeliveryId) + await flush() expect(client.exportState().lastSeqByChatId).toEqual({ "10": 2, "20": 2 }) + expect(lines.map((line) => JSON.parse(line).chatId)).toEqual(["20", "10", "10"]) + expect(JSON.parse(lines[2]!).seq).toBe(3) + stream.acknowledge(JSON.parse(lines[2]!)._inlineDeliveryId) + await flush() + expect(client.exportState().lastSeqByChatId).toEqual({ "10": 3, "20": 2 }) } finally { abort.abort() stream.close() diff --git a/plugins/hermes-agent/tests/install.test.ts b/plugins/hermes-agent/tests/install.test.ts index bc44d2599..0a60ab214 100644 --- a/plugins/hermes-agent/tests/install.test.ts +++ b/plugins/hermes-agent/tests/install.test.ts @@ -409,7 +409,9 @@ compat.plugin_hits = lambda _manifest: [] assert cli._compatibility_status() == {"ok": True, "reason": "loaded", "pluginPath": str(plugin_dir)} compat.plugin_hits = lambda _manifest: [object()] assert cli._compatibility_status()["reason"] == "deprecated_imports" -compat.plugin_hits = lambda _manifest: [] +# Upstream main keeps updater stubs but has removed the scanner exports. +del compat.plugin_hits +assert cli._compatibility_status() == {"ok": True, "reason": "loaded", "pluginPath": str(plugin_dir)} loaded.tools_registered = [] assert cli._compatibility_status()["reason"] == "tool_not_registered" loaded.tools_registered = ["inline"] diff --git a/scripts/ci/check-hermes-host.py b/scripts/ci/check-hermes-host.py index 60e37ef6a..d55f83de9 100644 --- a/scripts/ci/check-hermes-host.py +++ b/scripts/ci/check-hermes-host.py @@ -11,6 +11,7 @@ import sys import tempfile from pathlib import Path +from types import SimpleNamespace from unittest.mock import patch import httpx @@ -26,11 +27,15 @@ manifest = yaml.safe_load((plugin / "plugin.yaml").read_text()) assert "inline" in manifest.get("provides_tools", []), "inline tool missing from plugin.yaml" -# Scan even before the removal deadline: deprecated facade imports are future failures. +# Older hosts retain the deprecated-import scanner. After the compat layer's +# removal, plugin_compat remains only as updater stubs; real loading below is +# the authoritative import check on those hosts. if importlib.util.find_spec("hermes_cli.plugin_compat"): - from hermes_cli.plugin_compat import scan_plugin - hits = scan_plugin(plugin) - assert not hits, f"deprecated Hermes imports: {hits}" + from hermes_cli import plugin_compat + scan_plugin = getattr(plugin_compat, "scan_plugin", None) + if scan_plugin is not None: + hits = scan_plugin(plugin) + assert not hits, f"deprecated Hermes imports: {hits}" if importlib.util.find_spec("hermes_cli.plugin_validate"): from hermes_cli.plugin_validate import validate_plugin_dir report = validate_plugin_dir(plugin) @@ -81,6 +86,7 @@ async def reply(event): adapter._http_client = httpx.AsyncClient(transport=httpx.MockTransport(transport)) adapter.set_message_handler(reply) + adapter.set_authorization_check(lambda *_args, **_kwargs: True) event = {"kind": "message.new", "seq": 1, "chatId": "101", "meId": "999", "message": { "id": "201", "chatId": "101", "fromId": "42", "message": "hello Hermes", "peerId": {"peer": {"oneofKind": "user", "user": {"userId": "42"}}}, @@ -124,4 +130,305 @@ def media_client(*args, **kwargs): await adapter._http_client.aclose() asyncio.run(exercise()) -print("Real Hermes admission, registration, inbound/reply delivery, deduplication and media safety passed (offline transport).") + + +async def exercise_authorization(): + from gateway.run import GatewayRunner + from gateway.pairing import PairingStore + + # Construct only the runner state required by its real authorization callback; + # starting a gateway would introduce credentials, networking and provider work. + runner = GatewayRunner.__new__(GatewayRunner) + runner.config = SimpleNamespace(multiplex_profiles=False) + runner.adapters = {adapter.platform: adapter} + runner._profile_adapters = {} + runner._primary_profile_name = "default" + runner.pairing_store = PairingStore() + runner.pairing_stores = {} + adapter.set_authorization_check(runner._make_adapter_auth_check(adapter.platform)) + + def authorized(user="42", chat="101", chat_type="dm"): + return adapter._actor_authorized(chat_type, user, chat) + + # Pairing's allowlist mirroring is intentionally disabled: exercise the real + # persisted pairing grant in isolation without reading or writing env files. + auth_env = {key: "" for key in ( + "INLINE_ALLOWED_USERS", "INLINE_ALLOW_ALL_USERS", "GATEWAY_ALLOWED_USERS", + "GATEWAY_ALLOW_ALL_USERS", + )} + with patch.dict(os.environ, auth_env), patch("gateway.pairing._sync_allowlist_add"), patch("gateway.pairing._sync_allowlist_remove"): + assert authorized() is False, "unknown sender bypassed real host default deny" + code = runner.pairing_store.generate_code("inline", "42", "Offline fixture") + assert code is not None + assert runner.pairing_store.approve_code("inline", code) is not None + assert authorized() is True, "host pairing grant did not reach local authorization" + assert authorized("43") is False + + secondary_store = PairingStore(profile="secondary") + runner.pairing_stores["secondary"] = secondary_store + runner._profile_adapters["secondary"] = {adapter.platform: adapter} + adapter.set_authorization_check(runner._make_adapter_auth_check(adapter.platform, "secondary")) + assert authorized() is False, "primary pairing grant leaked into secondary profile" + secondary_code = secondary_store.generate_code("inline", "42", "Secondary fixture") + assert secondary_code is not None + assert secondary_store.approve_code("inline", secondary_code) is not None + assert authorized() is True, "profile-bound callback missed its pairing grant" + assert secondary_store.revoke("inline", "42") + assert authorized() is False, "secondary pairing revocation was not observed" + adapter.set_authorization_check(runner._make_adapter_auth_check(adapter.platform)) + assert authorized() is True, "secondary revocation removed primary grant" + + # The callback executes the actual runner predicate, which may consult + # adapter policy. A recursion here must fail this positive pairing check. + adapter._dm_policy = "disabled" + assert authorized() is False, "host grant overrode disabled DM policy" + adapter._dm_policy = "open" + adapter._allowed_chats = {"102"} + assert authorized(chat_type="group") is False, "host grant overrode allowed chats" + adapter._allowed_chats = set() + adapter._dm_policy = "allowlist" + adapter._allow_from = {"43"} + assert authorized() is False, "pairing bypassed explicit local sender restriction" + adapter._allow_from = {"42"} + assert authorized() is True + adapter._dm_policy = "open" + adapter._allow_from = set() + + # Exercise local mutations through ingress, with host-approved and unknown + # actors. HTTP and final message delivery stay offline; no LLM is invoked. + delivered, mutations = [], [] + + async def receive(event): + delivered.append(event) + + async def status(**kwargs): + pass + + async def create_thread(*args): + mutations.append(args) + return "301" + + async def dispatch(user, message, seq): + await adapter._dispatch_message({"kind": "message.new", "seq": seq, "chatId": "101", "message": { + "id": str(seq), "chatId": "101", "fromId": user, "message": message, + "peerId": {"peer": {"oneofKind": "user", "user": {"userId": user}}}, + }}) + + with patch.object(adapter, "handle_message", receive), patch.object(adapter, "_send_thread_status", status), patch.object(adapter, "_create_reply_thread", create_thread): + await dispatch("43", "/threads off", 1001) + assert "101" not in adapter._reply_thread_overrides + assert delivered[-1].text == "/threads off", "unknown DM lost host pairing ingress" + assert delivered[-1].source.user_id == "43" + await dispatch("42", "/threads off", 1002) + assert adapter._reply_thread_overrides["101"] == "off" + adapter._set_reply_threads_for_chat("101", "on") + await dispatch("43", "ordinary unknown message", 1003) + assert not mutations, "unknown sender created a reply thread before host authorization" + await dispatch("42", "ordinary paired message", 1004) + assert len(mutations) == 1, "paired sender lost automatic threads" + + assert runner.pairing_store.revoke("inline", "42") + assert authorized() is False, "revocation was hidden by cached adapter authorization" + adapter._dm_policy = "allowlist" + adapter._allow_from = {"42"} + assert authorized() is True, "real host adapter-policy delegation recursed or lost allowlist grant" + adapter._dm_policy = "open" + adapter._allow_from = set() + + # A real host denial must outrank an adapter-local allow-all setting. + adapter._allow_all = True + assert authorized() is False + adapter._allow_all = False + with patch.dict(os.environ, {"GATEWAY_ALLOW_ALL_USERS": "true"}): + assert authorized() is True + adapter._dm_policy = "disabled" + assert authorized() is False + adapter._dm_policy = "open" + adapter._allowed_chats = {"102"} + assert authorized(chat_type="group") is False + adapter._allowed_chats = set() + + from gateway.profile_routing import ProfileRoute + + # Parent-routed Inline threads need the full source: the standard host + # callback carries a child/thread ID but cannot carry parent_chat_id. + runner.config.multiplex_profiles = True + runner.config.profile_routes = [ + ProfileRoute(name="parent-work", platform="inline", profile="work", chat_id="700"), + ] + adapter.gateway_runner = runner + for profile in ("work", "child"): + (home / "profiles" / profile).mkdir(parents=True, exist_ok=True) + (home / "profiles" / profile / "config.yaml").write_text("{}\n") + runner.pairing_stores[profile] = PairingStore(profile=profile) + + def approve(store, user): + code = store.generate_code("inline", user, "Route fixture") + assert code is not None + assert store.approve_code("inline", code) is not None + + approve(runner.pairing_stores["work"], "44") + adapter.set_authorization_check(runner._make_adapter_auth_check(adapter.platform)) + assert adapter._is_sender_authorized("44", "group", "701", thread_id="701") is False + with patch.object(runner, "_is_user_authorized_for_source", wraps=runner._is_user_authorized_for_source) as host_authorize: + assert adapter._actor_authorized("group", "44", "701", thread_id="701", parent_chat_id="700") is True + admitted = host_authorize.call_args.args[0] + assert (admitted.chat_id, admitted.thread_id, admitted.parent_chat_id, admitted.profile) == ("701", "701", "700", "work") + + # A more specific child route wins; replacing the child ID with its + # parent to authorize would accidentally admit the parent-profile user. + runner.config.profile_routes.append( + ProfileRoute(name="child-only", platform="inline", profile="child", thread_id="701"), + ) + # GatewayConfig normally receives this order from parse_profile_routes. + runner.config.profile_routes.sort(key=lambda route: route.specificity, reverse=True) + approve(runner.pairing_stores["child"], "45") + assert adapter._actor_authorized("group", "44", "701", thread_id="701", parent_chat_id="700") is False + with patch.object(runner, "_is_user_authorized_for_source", wraps=runner._is_user_authorized_for_source) as host_authorize: + assert adapter._actor_authorized("group", "45", "701", thread_id="701", parent_chat_id="700") is True + admitted = host_authorize.call_args.args[0] + assert (admitted.chat_id, admitted.thread_id, admitted.profile) == ("701", "701", "child") + + approve(runner.pairing_store, "44") + runner.config.profile_routes = [ + ProfileRoute(name="unserved", platform="inline", profile="missing", chat_id="700"), + ] + rejected = adapter.build_source(chat_id="701", chat_type="group", user_id="44", thread_id="701", parent_chat_id="700") + assert rejected.profile_route_rejected is True + assert adapter._actor_authorized("group", "44", "701", thread_id="701", parent_chat_id="700") is False + + # Missing child metadata must not turn a routed thread into a default + # profile chat and use that profile's otherwise valid pairing grant. + async def missing_info(chat_id, **kwargs): + return {} + + async def child_message(chat_id, message_id): + return {"id": message_id, "peerId": {"type": {"oneofKind": "chat"}}} + + answers = [] + + async def answer(interaction_id, message): + answers.append(message) + + with patch.object(adapter, "_get_chat_info", missing_info), patch.object(adapter, "_fetch_message", child_message), patch.object(adapter, "_answer_action", answer), patch.object(adapter, "handle_message", receive): + assert not await adapter._action_allowed({"chatId": "701", "messageId": "1", "actorUserId": "44", "interactionId": "missing-route"}) + assert answers == ["Access check temporarily unavailable. Try again."] + delivered.clear() + try: + await adapter._dispatch_message({"kind": "message.new", "seq": 2001, "chatId": "701", "message": { + "id": "2001", "fromId": "44", "message": "@bot hello", "mentioned": True, + "peerId": {"type": {"oneofKind": "chat"}}, + }}) + except RuntimeError as exc: + assert type(exc).__name__ == "InlineInboundDeferred" + else: + raise AssertionError("missing routing metadata was not deferred") + assert not delivered, "missing routing metadata borrowed default pairing approval" + + +asyncio.run(exercise_authorization()) +async def exercise_receipt_recovery(): + from gateway.run import GatewayRunner + from gateway.pairing import PairingStore + + recovered = platform_registry.create_adapter("inline", PlatformConfig( + enabled=True, token="offline-test-token", extra={ + "dm_policy": "open", "reply_threads": "off", "context_backfill": "off", + "sync_commands": False, "text_debounce_seconds": 0, + }, + )) + assert recovered is not None + module = sys.modules[type(recovered).__module__] + runner = GatewayRunner.__new__(GatewayRunner) + runner.config = SimpleNamespace(multiplex_profiles=False) + runner.adapters = {recovered.platform: recovered} + runner._profile_adapters = {} + runner._primary_profile_name = "default" + runner.pairing_store = PairingStore() + runner.pairing_stores = {} + recovered.set_authorization_check(runner._make_adapter_auth_check(recovered.platform)) + acknowledgements, received = [], [] + model_started, release_model, replied = asyncio.Event(), asyncio.Event(), asyncio.Event() + + def transport(request): + body = json.loads(request.content) + if request.url.path == "/inbound/ack": + acknowledgements.append(body["deliveryId"]) + elif request.url.path == "/send": + replied.set() + return httpx.Response(200, json={"ok": True, "result": {"messageId": "9002"}}) + + async def model(event): + received.append(event) + model_started.set() + await release_model.wait() + return "Recovered native turn" + + recovered._http_client = httpx.AsyncClient(transport=httpx.MockTransport(transport)) + recovered.set_message_handler(model) + event = {"kind": "message.new", "seq": 9001, "chatId": "901", "meId": "999", + "_inlineDeliveryId": "native-recovery", "message": { + "id": "9001", "chatId": "901", "fromId": "46", "message": "recover native admission", + "peerId": {"peer": {"oneofKind": "user", "user": {"userId": "46"}}}, + }} + auth_env = {key: "" for key in ( + "INLINE_ALLOWED_USERS", "INLINE_ALLOW_ALL_USERS", "GATEWAY_ALLOWED_USERS", "GATEWAY_ALLOW_ALL_USERS", + )} + try: + with patch.dict(os.environ, auth_env), patch("gateway.pairing._sync_allowlist_add"), patch("gateway.pairing._sync_allowlist_remove"): + code = runner.pairing_store.generate_code("inline", "46", "Receipt recovery fixture") + assert code is not None + assert runner.pairing_store.approve_code("inline", code) is not None + native_authorize = runner._is_user_authorized + attempts = 0 + + def transient_authorize(source): + nonlocal attempts + attempts += 1 + if attempts == 1: + raise RuntimeError("offline injected authorization outage") + return native_authorize(source) + + with patch.object(runner, "_is_user_authorized", side_effect=transient_authorize), patch.object(module, "_INBOUND_RETRY_INITIAL_SECONDS", 0.001): + await recovered._on_inbound(json.dumps(event)) + receipt = recovered._inbound_deliveries["native-recovery"] + await asyncio.wait_for(receipt, timeout=10) + await asyncio.wait_for(model_started.wait(), timeout=10) + assert attempts == 2, "native unknown authorization was not retried" + assert acknowledgements == ["native-recovery"] + assert len(received) == 1 and received[0].source.user_id == "46" + assert not release_model.is_set() and not replied.is_set(), "receipt waited for model completion" + assert recovered._active_sessions, "real host did not retain the background turn" + release_model.set() + await asyncio.wait_for(replied.wait(), timeout=10) + + # Exercise the actual native shielded notification. Teardown cancels its + # carrier delivery task, but the detached fatal handler must still finish. + notified, teardown_finished = asyncio.Event(), asyncio.Event() + + async def fatal_handler(failed): + assert failed is recovered + notified.set() + await failed.disconnect() + await asyncio.sleep(0) + teardown_finished.set() + + async def unexpected_failure(_event): + raise RuntimeError("offline injected dispatch failure") + + recovered.set_fatal_error_handler(fatal_handler) + with patch.object(recovered, "_dispatch_inbound", side_effect=unexpected_failure): + await recovered._on_inbound(json.dumps({**event, "_inlineDeliveryId": "native-fatal"})) + await asyncio.wait_for(notified.wait(), timeout=10) + await asyncio.wait_for(teardown_finished.wait(), timeout=10) + assert recovered.has_fatal_error and recovered.fatal_error_retryable + assert recovered.fatal_error_code == "INBOUND_FAILED" + assert acknowledgements == ["native-recovery"], "unexpected failure acknowledged its SDK receipt" + assert not recovered._inbound_deliveries, "fatal teardown leaked delivery tasks" + finally: + release_model.set() + await recovered.disconnect() + + +asyncio.run(exercise_receipt_recovery()) +print("Real Hermes admission, registration, authorization/pairing, receipt recovery, fatal teardown, local effects, inbound/reply delivery, deduplication and media safety passed (offline transport).") diff --git a/scripts/ci/check-hermes-source.mjs b/scripts/ci/check-hermes-source.mjs index 84779397a..70349f991 100644 --- a/scripts/ci/check-hermes-source.mjs +++ b/scripts/ci/check-hermes-source.mjs @@ -97,6 +97,11 @@ try { assert.ok(connected, `installed source sidecar did not connect to offline mock: ${output}`) assert.equal(health.result.version, expectedVersion, "source sidecar lost plugin version metadata") assert.equal((await post("/healthz", {}, false)).status, 401) + assert.equal((await post("/inbound/ack", { deliveryId: "offline-receipt" }, false)).status, 401) + // ACK retries are harmless even after the original receipt has been retired. + for (let attempt = 0; attempt < 2; attempt++) { + assert.equal((await post("/inbound/ack", { deliveryId: "offline-receipt" })).status, 200) + } const sent = await (await post("/send", { target: { chatId: "123" }, text: "source catalog smoke", parseMarkdown: false })).json() assert.equal(sent.ok, true) assert.ok(sent.result.messageId)