diff --git a/gateway/relay/adapter.py b/gateway/relay/adapter.py index b149539cee5..b544748208a 100644 --- a/gateway/relay/adapter.py +++ b/gateway/relay/adapter.py @@ -72,6 +72,16 @@ class RelayAdapter(BasePlatformAdapter): # recipient's author binding; we re-attach this user_id as # metadata.user_id on the outbound action so it can. See _capture_scope. self._dm_user_by_chat: Dict[str, str] = {} + # chat_id -> chat_type (e.g. "dm", "channel", "group") learned from the + # inbound event. Used to reproduce native Slack's synthetic-DM-thread + # suppression on the relay lane: a DM streaming reply carries + # reply_to= as its edit anchor, but the connector + # maps a raw reply_to to a Slack thread_ts — so a plain DM reply would be + # threaded UNDER the user's message (and lose progressive edit streaming) + # instead of posting flat at the DM root. Native SlackAdapter drops that + # synthetic reply_to in _resolve_thread_ts; the relay lane needs the same + # disambiguation, and it needs the chat_type to know a chat is a DM. + self._chat_type_by_chat: Dict[str, str] = {} # chat_id -> the UNDERLYING platform (e.g. "discord", "telegram") this # chat belongs to (Phase 1.5 multi-platform-per-agent). One relay adapter # fronts N platforms on one WS; an outbound reply must egress through the @@ -344,10 +354,19 @@ class RelayAdapter(BasePlatformAdapter): scope = getattr(src, "scope_id", None) if scope: self._scope_by_chat[str(chat)] = str(scope) + # Remember the chat_type so send() can suppress the synthetic-DM + # thread anchor on Slack (native _resolve_thread_ts parity). send() + # only receives a chat_id, so it needs this per-chat cache to know a + # chat is a DM. + chat_type = getattr(src, "chat_type", None) + if chat_type: + self._chat_type_by_chat[str(chat)] = str(chat_type) except Exception: # noqa: BLE001 - scope tracking must never break inbound pass - def _with_scope(self, chat_id: str, metadata: Optional[Dict[str, Any]]) -> Dict[str, Any]: + def _with_scope( + self, chat_id: str, metadata: Optional[Dict[str, Any]] + ) -> Dict[str, Any]: """Ensure the outbound metadata carries the discriminator(s) the connector's egress guard needs to resolve the owning tenant. @@ -506,16 +525,26 @@ class RelayAdapter(BasePlatformAdapter): else: text = "" member = payload.get("member") or {} - user = (member.get("user") if isinstance(member, dict) else None) or payload.get("user") or {} + user = ( + (member.get("user") if isinstance(member, dict) else None) + or payload.get("user") + or {} + ) channel_id = str(payload.get("channel_id") or "") guild_id = payload.get("guild_id") # real Discord interaction wire field source = SessionSource( platform=Platform.RELAY, chat_id=channel_id, chat_type="channel" if guild_id else "dm", - user_id=str(user.get("id")) if isinstance(user, dict) and user.get("id") else None, - user_name=str(user.get("username")) if isinstance(user, dict) and user.get("username") else None, - scope_id=str(guild_id) if guild_id else None, # Discord guild → generic scope slot + user_id=str(user.get("id")) + if isinstance(user, dict) and user.get("id") + else None, + user_name=str(user.get("username")) + if isinstance(user, dict) and user.get("username") + else None, + scope_id=str(guild_id) + if guild_id + else None, # Discord guild → generic scope slot message_id=str(payload.get("id")) if payload.get("id") else None, ) event = MessageEvent(text=text, message_type=message_type, source=source) @@ -531,7 +560,9 @@ class RelayAdapter(BasePlatformAdapter): prompt_id, option_id = decoded msg = payload.get("message") or {} prompt_message_id = ( - str(msg.get("id")) if isinstance(msg, dict) and msg.get("id") else None + str(msg.get("id")) + if isinstance(msg, dict) and msg.get("id") + else None ) event.prompt_response = { "prompt_id": prompt_id, @@ -584,7 +615,9 @@ class RelayAdapter(BasePlatformAdapter): sub_name = str(opt.get("name") or "").strip() if sub_name: parts.append(sub_name) - parts.extend(RelayAdapter._render_interaction_options(opt.get("options"))) + parts.extend( + RelayAdapter._render_interaction_options(opt.get("options")) + ) else: value = opt.get("value") if value is not None and str(value).strip(): @@ -708,12 +741,21 @@ class RelayAdapter(BasePlatformAdapter): ) if self._transport is None: return SendResult(success=False, error="no transport") + # Native _resolve_thread_ts parity: a Slack DM reply must post flat at + # the DM root, not threaded under the triggering message. Drop the + # synthetic self-anchor from BOTH the top-level reply_to and the mirrored + # metadata.reply_to_message_id so the connector can't thread on either. + effective_reply_to = self._resolve_reply_to_for_send( + chat_id, reply_to, send_metadata + ) + if effective_reply_to is None and reply_to is not None: + send_metadata.pop("reply_to_message_id", None) result = await self._transport.send_outbound( { "op": "send", "chat_id": chat_id, "content": content, - "reply_to": reply_to, + "reply_to": effective_reply_to, "metadata": self._with_scope(chat_id, send_metadata), }, platform=self._platform_by_chat.get(str(chat_id)), @@ -724,6 +766,57 @@ class RelayAdapter(BasePlatformAdapter): error=result.get("error"), ) + def _resolve_reply_to_for_send( + self, + chat_id: str, + reply_to: Optional[str], + metadata: Optional[Dict[str, Any]], + ) -> Optional[str]: + """Suppress the synthetic-DM thread anchor for a Slack DM reply. + + A DM turn's streaming reply is sent with ``reply_to`` = the triggering + message's ts (the stream consumer's ``initial_reply_to_id``, used as the + edit anchor and, on threading platforms, the reply target). The + connector's slackRestSender maps a raw ``reply_to`` to a Slack + ``thread_ts``, so a plain DM reply would be posted THREADED under the + user's message instead of flat at the DM root — and a threaded first + send loses the progressive edit-streaming the user sees in a real + thread (the reported symptom: DM/home replies arrive flat, no + progressive edits). + + Native Slack Hermes already suppresses this synthetic DM thread anchor: + ``SlackAdapter._resolve_thread_ts`` returns ``None`` for a top-level / + DM message when ``reply_in_thread`` is off. The relay lane has no such + disambiguation, so we reproduce it here. run.py already encodes the + real-thread decision in ``metadata["thread_id"]`` (it is set only when + progress threading is active — a real thread, or channel autoThread); + for a DM with no real thread that key is absent. So the rule is: + + Slack DM + no real ``thread_id`` in metadata ⇒ drop ``reply_to``. + + This posts the reply flat at the DM root and lets the consumer edit its + own first-send ts — streaming works exactly as in a thread. It does NOT: + * reintroduce a synthetic DM thread_id (#18859 / the /sethome + landmine) — it removes an anchor, never adds one; + * regress real-thread streaming — a real thread carries a distinct + ``thread_id`` in metadata, so the guard leaves ``reply_to`` alone; + * regress channel autoThread — a channel/group top-level reply carries + ``thread_id`` (the message's own ts) in metadata when threading is + on, so it is left alone; and a non-DM chat is never matched here. + """ + if reply_to is None: + return None + if self._platform_by_chat.get(str(chat_id)) != Platform.SLACK.value: + return reply_to + if self._chat_type_by_chat.get(str(chat_id)) != "dm": + return reply_to + md = metadata or {} + if md.get("thread_id") or md.get("thread_ts"): + # A real thread was resolved by run.py — honour it. + return reply_to + # Synthetic DM self-anchor: post flat at the DM root (native parity). + return None + async def edit_message( self, chat_id: str, @@ -1029,8 +1122,12 @@ class RelayAdapter(BasePlatformAdapter): if result is not None: return result return await super().send_image_file( - chat_id, image_path, caption=caption, reply_to=reply_to, - metadata=metadata, **kwargs, + chat_id, + image_path, + caption=caption, + reply_to=reply_to, + metadata=metadata, + **kwargs, ) async def send_voice( @@ -1055,8 +1152,12 @@ class RelayAdapter(BasePlatformAdapter): if result is not None: return result return await super().send_voice( - chat_id, audio_path, caption=caption, reply_to=reply_to, - metadata=metadata, **kwargs, + chat_id, + audio_path, + caption=caption, + reply_to=reply_to, + metadata=metadata, + **kwargs, ) async def send_video( @@ -1081,8 +1182,12 @@ class RelayAdapter(BasePlatformAdapter): if result is not None: return result return await super().send_video( - chat_id, video_path, caption=caption, reply_to=reply_to, - metadata=metadata, **kwargs, + chat_id, + video_path, + caption=caption, + reply_to=reply_to, + metadata=metadata, + **kwargs, ) async def send_document( @@ -1109,13 +1214,20 @@ class RelayAdapter(BasePlatformAdapter): if result is not None: return result return await super().send_document( - chat_id, file_path, caption=caption, file_name=file_name, - reply_to=reply_to, metadata=metadata, **kwargs, + chat_id, + file_path, + caption=caption, + file_name=file_name, + reply_to=reply_to, + metadata=metadata, + **kwargs, ) # ── Phase 3 interactive: prompt + react ────────────────────────────── - def _mint_prompt(self, kind: str, state: Dict[str, Any], timeout_s: float = 3600.0) -> str: + def _mint_prompt( + self, kind: str, state: Dict[str, Any], timeout_s: float = 3600.0 + ) -> str: """Register a pending prompt and return its 8-hex id. ``state`` carries what the resolver needs when the answer comes back @@ -1135,7 +1247,9 @@ class RelayAdapter(BasePlatformAdapter): # Opportunistic sweep so abandoned prompts can't accumulate: drop # anything already expired (cheap — dict is small by construction). now = time.time() - for stale in [k for k, v in self._pending_prompts.items() if v.get("expires_at", 0) < now]: + for stale in [ + k for k, v in self._pending_prompts.items() if v.get("expires_at", 0) < now + ]: self._pending_prompts.pop(stale, None) return prompt_id @@ -1150,6 +1264,55 @@ class RelayAdapter(BasePlatformAdapter): return None return state + def _strip_synthetic_dm_thread( + self, chat_id: str, metadata: Optional[Dict[str, Any]] + ) -> Optional[Dict[str, Any]]: + """Drop the synthetic DM thread anchor from an interactive prompt's metadata. + + A clarify/approval/confirm prompt is emitted mid-turn in reply to the + triggering inbound event, so ``metadata`` carries that event's thread + context — run.py's ``_thread_metadata_for_source`` stamps + ``metadata["thread_id"]`` (and, for Slack, ``metadata["message_id"]`` = + the triggering message ts). For a Slack DM with no REAL thread, that + ``thread_id`` is the message's own synthetic self-anchor (a session-keying + fallback), and forwarding it makes the connector's slackRestSender thread + the prompt card UNDER the user's message instead of posting it flat at the + DM root — the reported bug ("approval block was put in a thread"). + + Native Slack Hermes already suppresses this synthetic DM thread anchor + (``SlackAdapter._resolve_thread_ts`` returns ``None`` for a top-level / DM + message). We reproduce it here with the same discipline used on the + streaming path (``_resolve_reply_to_for_send``): + + Slack DM + thread_id is the synthetic self-anchor ⇒ strip thread_id. + + A REAL thread (``thread_id`` distinct from the triggering message ts) is + left untouched so a prompt raised inside a thread stays in that thread; + non-DM / non-Slack chats are never matched. Only the threading keys are + removed — tenant scope (``scope_id`` / ``slack_team_id``) and everything + else survive so egress routing is unaffected. + """ + if not metadata: + return metadata + if self._platform_by_chat.get(str(chat_id)) != Platform.SLACK.value: + return metadata + if self._chat_type_by_chat.get(str(chat_id)) != "dm": + return metadata + thread_id = metadata.get("thread_id") + if not thread_id: + return metadata + # A real thread carries a thread_id distinct from the triggering message + # ts (run.py stamps that ts as metadata["message_id"] on Slack). Only the + # synthetic self-anchor (thread_id == that ts, or no distinguishing anchor + # present) is stripped; a genuine thread is honoured. + anchor = metadata.get("message_id") + if anchor is not None and str(thread_id) != str(anchor): + return metadata + cleaned = dict(metadata) + cleaned.pop("thread_id", None) + cleaned.pop("thread_ts", None) + return cleaned + async def _send_prompt( self, chat_id: str, @@ -1171,6 +1334,16 @@ class RelayAdapter(BasePlatformAdapter): """ if self._transport is None or not self.descriptor.supports_op("prompt"): return None + # An interactive prompt (approval / clarify / slash-confirm) is emitted + # mid-turn in reply to the triggering inbound event, so `metadata` carries + # that event's thread context (run.py _thread_metadata_for_source stamps + # metadata.thread_id — for a Slack DM the triggering message's own ts, + # used only as a session-keying fallback). Forwarding it makes the + # connector thread the prompt card UNDER the triggering message instead + # of posting it flat at the DM root (the reported bug). Native Slack + # Hermes suppresses this synthetic DM thread anchor; drop it here for the + # same Slack-DM-with-no-real-thread case, matching _resolve_reply_to_for_send. + prompt_metadata = self._strip_synthetic_dm_thread(chat_id, metadata) action: Dict[str, Any] = { "op": "prompt", "chat_id": chat_id, @@ -1178,8 +1351,10 @@ class RelayAdapter(BasePlatformAdapter): "prompt_kind": prompt_kind, "prompt_id": prompt_id, "options": options, - "reply_to": reply_to, - "metadata": self._with_scope(chat_id, metadata), + "reply_to": self._resolve_reply_to_for_send( + chat_id, reply_to, prompt_metadata + ), + "metadata": self._with_scope(chat_id, prompt_metadata), } if timeout_s is not None: action["timeout_s"] = int(timeout_s) @@ -1228,7 +1403,11 @@ class RelayAdapter(BasePlatformAdapter): if not smart_denied and allow_session: options.append({"id": "session", "label": "✅ Session", "style": "primary"}) if allow_permanent: - options.append({"id": "always", "label": "✅ Always", "style": "primary"}) + options.append({ + "id": "always", + "label": "✅ Always", + "style": "primary", + }) options.append({"id": "deny", "label": "❌ Deny", "style": "danger"}) cmd_preview = command if len(command) <= 1500 else command[:1500] + "..." @@ -1238,7 +1417,9 @@ class RelayAdapter(BasePlatformAdapter): f"Reason: {description}" ) if smart_denied: - text += "\n\n**Smart DENY:** owner override applies to this one operation only." + text += ( + "\n\n**Smart DENY:** owner override applies to this one operation only." + ) prompt_id = self._mint_prompt( "exec_approval", @@ -1386,7 +1567,11 @@ class RelayAdapter(BasePlatformAdapter): if kind == "exec_approval": from tools.approval import resolve_gateway_approval - choice = option_id if option_id in {"once", "session", "always", "deny"} else "deny" + choice = ( + option_id + if option_id in {"once", "session", "always", "deny"} + else "deny" + ) count = resolve_gateway_approval(session_key, choice) label = { "once": "✅ Approved once", @@ -1399,13 +1584,17 @@ class RelayAdapter(BasePlatformAdapter): # Acknowledge in-channel (the connector's prompt message can't # be edited cross-platform yet — edit support varies; a short # confirmation preserves the audit trail the native edit gives). - await self.send(chat_id, label, metadata=self._prompt_reply_metadata(event)) + await self.send( + chat_id, label, metadata=self._prompt_reply_metadata(event) + ) if count: self.resume_typing_for_chat(chat_id) elif kind == "slash_confirm": from tools import slash_confirm as slash_confirm_mod - choice = option_id if option_id in {"once", "always", "cancel"} else "cancel" + choice = ( + option_id if option_id in {"once", "always", "cancel"} else "cancel" + ) result_text = await slash_confirm_mod.resolve( session_key, str(state.get("confirm_id") or ""), choice ) @@ -1414,13 +1603,20 @@ class RelayAdapter(BasePlatformAdapter): "always": "🔒 Always approve", "cancel": "❌ Cancelled", }.get(choice, "Resolved") - await self.send(chat_id, label, metadata=self._prompt_reply_metadata(event)) + await self.send( + chat_id, label, metadata=self._prompt_reply_metadata(event) + ) if result_text: await self.send( - chat_id, str(result_text), metadata=self._prompt_reply_metadata(event) + chat_id, + str(result_text), + metadata=self._prompt_reply_metadata(event), ) elif kind == "clarify": - from tools.clarify_gateway import mark_awaiting_text, resolve_gateway_clarify + from tools.clarify_gateway import ( + mark_awaiting_text, + resolve_gateway_clarify, + ) clarify_id = str(state.get("clarify_id") or "") if option_id == "other": diff --git a/tests/gateway/relay/test_relay_slack_dm_streaming.py b/tests/gateway/relay/test_relay_slack_dm_streaming.py new file mode 100644 index 00000000000..008f6c608c0 --- /dev/null +++ b/tests/gateway/relay/test_relay_slack_dm_streaming.py @@ -0,0 +1,220 @@ +"""Slack relay: edit-based streaming of the reply must fire in a DM. + +Reported symptom (live): agent responses stream progressively (edit-based) in a +Slack THREAD but arrive FLAT (single message, no progressive edits) in a Slack +DM/home over the relay. + +Root cause: a DM turn's streaming reply is sent with +``reply_to = `` (the stream consumer's +``initial_reply_to_id`` — its edit anchor). The connector's slackRestSender maps +a raw ``reply_to`` to a Slack ``thread_ts``, so a plain DM reply gets threaded +UNDER the user's message instead of posting flat at the DM root, and a threaded +first send loses the progressive edit streaming the user sees in a real thread. +Native ``SlackAdapter._resolve_thread_ts`` already suppresses this synthetic DM +thread anchor; the relay lane had no such disambiguation. + +These are behaviour-contract tests: they assert how the outbound frame relates to +the chat type + thread metadata (the invariant the connector depends on), not a +snapshot. They drive the REAL ``RelayAdapter`` + ``GatewayStreamConsumer`` + +``StubConnector`` end to end. +""" + +from __future__ import annotations + +import pytest + +from gateway.config import Platform, PlatformConfig +from gateway.platforms.base import MessageEvent, MessageType +from gateway.relay.adapter import RelayAdapter +from gateway.relay.descriptor import CONTRACT_VERSION, CapabilityDescriptor +from gateway.session import SessionSource +from gateway.stream_consumer import GatewayStreamConsumer, StreamConsumerConfig + +from tests.gateway.relay.stub_connector import StubConnector + + +def _slack_desc(**kw) -> CapabilityDescriptor: + base = dict( + contract_version=CONTRACT_VERSION, + platform="slack", + label="Slack", + max_message_length=4000, + supports_draft_streaming=False, + supports_edit=True, + supports_threads=True, + markdown_dialect="mrkdwn", + len_unit="chars", + emoji="\U0001f4ac", + platform_hint="", + pii_safe=False, + ) + base.update(kw) + return CapabilityDescriptor(**base) + + +def _wire(chat_id: str, chat_type: str, *, user_id="U1", scope_id=None): + """A RelayAdapter fronting Slack, with inbound scope captured for chat_id.""" + stub = StubConnector(_slack_desc()) + adapter = RelayAdapter(PlatformConfig(), _slack_desc(), transport=stub) + src = SessionSource( + platform=Platform.SLACK, + chat_id=chat_id, + chat_type=chat_type, + user_id=user_id, + scope_id=scope_id, + ) + adapter._capture_scope( + MessageEvent(text="hi", source=src, message_type=MessageType.TEXT) + ) + return adapter, stub + + +# --------------------------------------------------------------------------- +# The pure disambiguation contract (RelayAdapter.send) +# --------------------------------------------------------------------------- +@pytest.mark.asyncio +async def test_slack_dm_reply_drops_synthetic_thread_anchor(): + """A Slack DM reply with no real thread posts FLAT: reply_to is dropped so + the connector cannot thread it under the triggering message.""" + adapter, stub = _wire("D1", "dm") + await adapter.send("D1", "the answer", reply_to="1700.0001") + assert len(stub.sent) == 1 + frame = stub.sent[0] + assert frame["op"] == "send" + # The synthetic self-anchor is suppressed on BOTH surfaces. + assert frame["reply_to"] is None + assert "thread_id" not in (frame["metadata"] or {}) + # And no synthetic thread_id was invented (the #18859 landmine). + assert "thread_ts" not in (frame["metadata"] or {}) + + +@pytest.mark.asyncio +async def test_slack_dm_reply_with_real_thread_keeps_anchor(): + """A DM turn that IS inside a real thread (metadata carries a distinct + thread_id) must keep threading — the guard only drops the synthetic anchor.""" + adapter, stub = _wire("D1", "dm") + await adapter.send( + "D1", "in thread", reply_to="1700.0002", metadata={"thread_id": "1699.9000"} + ) + frame = stub.sent[0] + assert frame["reply_to"] == "1700.0002" + assert frame["metadata"]["thread_id"] == "1699.9000" + + +@pytest.mark.asyncio +async def test_slack_channel_top_level_reply_keeps_autothread_anchor(): + """A channel top-level reply carries thread_id (its own ts) in metadata when + autoThread is on; the DM-only guard must not touch it.""" + adapter, stub = _wire("C1", "channel", scope_id="T1") + await adapter.send( + "C1", "channel reply", reply_to="1700.0003", metadata={"thread_id": "1700.0003"} + ) + frame = stub.sent[0] + assert frame["reply_to"] == "1700.0003" + assert frame["metadata"]["thread_id"] == "1700.0003" + + +@pytest.mark.asyncio +async def test_non_slack_dm_reply_unchanged(): + """The disambiguation is Slack-scoped: a non-Slack relay chat keeps reply_to + (its connector owns its own threading semantics).""" + stub = StubConnector(_slack_desc(platform="discord")) + adapter = RelayAdapter( + PlatformConfig(), _slack_desc(platform="discord"), transport=stub + ) + src = SessionSource( + platform=Platform.DISCORD, chat_id="dc1", chat_type="dm", user_id="U1" + ) + adapter._capture_scope( + MessageEvent(text="hi", source=src, message_type=MessageType.TEXT) + ) + await adapter.send("dc1", "hi", reply_to="msg-9") + assert stub.sent[0]["reply_to"] == "msg-9" + + +# --------------------------------------------------------------------------- +# End-to-end: the stream consumer keeps edit-streaming in a DM +# --------------------------------------------------------------------------- +async def _drive_stream(adapter, chat_id, *, metadata, initial_reply_to_id, chat_type): + cfg = StreamConsumerConfig( + edit_interval=0.0, + buffer_threshold=1, + transport="edit", + chat_type=chat_type, + ) + consumer = GatewayStreamConsumer( + adapter=adapter, + chat_id=chat_id, + config=cfg, + metadata=metadata, + initial_reply_to_id=initial_reply_to_id, + ) + # Feed progressive deltas, then finalize — mirrors the live delta callback. + for chunk in ("Hel", "lo ", "world", ". Done."): + consumer.on_delta(chunk) + consumer.finish() + await consumer.run() + return consumer + + +@pytest.mark.asyncio +async def test_slack_dm_stream_consumer_edits_own_ts_not_flat(): + """A Slack DM turn (chat_type='dm', no thread, metadata None) still builds a + stream consumer that keeps edit support and emits progressive EDITs of the + reply message — the flat-DM regression contract. + + The connector returns a real message_id for the flat first send, so edit + support must stay on and at least one edit op must be emitted (progressive + streaming), identical to a thread. No synthetic thread is created.""" + adapter, stub = _wire("D1", "dm") + consumer = await _drive_stream( + adapter, + "D1", + metadata=None, # DM: _status_thread_metadata is None in run.py + initial_reply_to_id="1700.0001", # the triggering message ts + chat_type="dm", + ) + + ops = [f["op"] for f in stub.sent] + # First a flat send; edit support stays on so progressive edits CAN flow + # (exact intermediate-frame timing is covered by the stream_consumer unit + # suite — here we assert the DM regression contract: streaming is not + # self-disabled and every edit targets the reply's own ts). + assert ops[0] == "send" + # Edit support survived: message_id set, not the __no_edit__ sentinel. + assert consumer.message_id and consumer.message_id != "__no_edit__" + assert consumer._edit_supported is True + + first_send = stub.sent[0] + # The reply posts FLAT at the DM root — no synthetic thread anchor. + assert first_send["reply_to"] is None + assert "thread_id" not in (first_send["metadata"] or {}) + assert "thread_ts" not in (first_send["metadata"] or {}) + # reply_to_message_id (the mirrored self-anchor) is stripped too. + assert "reply_to_message_id" not in (first_send["metadata"] or {}) + + # Any edits that flowed target the same first-send ts (editing its own + # message), never a synthetic thread. + edit_ids = {f["message_id"] for f in stub.sent if f["op"] == "edit"} + assert edit_ids <= {stub.next_send_result["message_id"]} + + +@pytest.mark.asyncio +async def test_slack_thread_stream_consumer_still_threads_and_streams(): + """Regression guard: a Slack THREAD turn keeps its real thread_id AND streams + (the DM fix must not change the thread path).""" + adapter, stub = _wire("C1", "channel", scope_id="T1") + consumer = await _drive_stream( + adapter, + "C1", + metadata={"thread_id": "1699.9000"}, + initial_reply_to_id="1700.0002", + chat_type="channel", + ) + ops = [f["op"] for f in stub.sent] + assert ops[0] == "send" + assert consumer._edit_supported is True + first_send = stub.sent[0] + # Thread preserved: the real thread_id rides along and reply_to is kept. + assert first_send["metadata"]["thread_id"] == "1699.9000" + assert first_send["reply_to"] == "1700.0002"