diff --git a/gateway/run.py b/gateway/run.py index a1d2246b73..a11315e825 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -30069,6 +30069,12 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew and source.thread_id and event_message_id ) + or ( + # Buzz has no native thread_id; threading is always via reply-to + # the triggering event id (channel clutter otherwise). + str(getattr(source.platform, "value", source.platform) or "").lower() == "buzz" + and event_message_id + ) or _relay_prospective_thread_id else None ) diff --git a/gateway/stream_consumer.py b/gateway/stream_consumer.py index f11fa43262..4fc33efd20 100644 --- a/gateway/stream_consumer.py +++ b/gateway/stream_consumer.py @@ -2712,11 +2712,12 @@ class GatewayStreamConsumer: # draft(final=true) — that would seal the live stream with # interim text and orphan the true final into a plain-send # duplicate (live finding, 2026-08-16 canary). - _md = dict(self.metadata) if self.metadata else {} + _md = self._metadata_for_send(final=False) or {} _md["_interim_send"] = True result = await self.adapter.send( chat_id=self.chat_id, content=text, + reply_to=self._initial_reply_to_id, metadata=_md, ) # Note: do NOT set _already_sent = True here. diff --git a/plugins/platforms/buzz/adapter.py b/plugins/platforms/buzz/adapter.py index d3ce59aa6e..a60e182c67 100644 --- a/plugins/platforms/buzz/adapter.py +++ b/plugins/platforms/buzz/adapter.py @@ -795,8 +795,13 @@ class BuzzAdapter(BasePlatformAdapter): if not content: return SendResult(success=False, error="Empty message") args = ["messages", "send", "--channel", str(chat_id), "--content", "-"] + # Prefer the stable thread anchor from metadata.thread_id (Slack-style), + # then metadata.reply_to_message_id (gateway stream consumer / + # progress sends), then the explicit reply_to argument. Without + # reply_to_message_id, interim commentary posts flat in the channel. + meta = metadata or {} reply_target = self._resolve_reply_anchor( - (metadata or {}).get("thread_id") or reply_to + meta.get("thread_id") or meta.get("reply_to_message_id") or reply_to ) if reply_target and self._reply_to_mode != "off": args += ["--reply-to", str(reply_target)] diff --git a/tests/gateway/test_buzz_adapter.py b/tests/gateway/test_buzz_adapter.py index 5993bccab9..e3830c54fd 100644 --- a/tests/gateway/test_buzz_adapter.py +++ b/tests/gateway/test_buzz_adapter.py @@ -899,6 +899,29 @@ class TestBuzzAdapterSend: args, _stdin = cli.calls[0] assert args[args.index("--reply-to") + 1] == "buzz-event-123" + @pytest.mark.asyncio + async def test_send_uses_metadata_reply_to_message_id(self): + """Gateway stream/progress pass reply anchors via metadata. + + Without honoring reply_to_message_id, mid-turn commentary posts as + new top-level channel messages instead of thread replies. + """ + adapter = _make_adapter() + adapter._channel_state[CHANNEL] = {"chat_type": "group", "last_ts": 0, "seen": {}} + cli = _ScriptedCli() + cli.script("messages", "send", {"accepted": True, "event_id": "evt-reply", "message": ""}) + adapter._run_cli = cli + + result = await adapter.send( + CHANNEL, + "threaded reply", + metadata={"reply_to_message_id": "root-event-abc"}, + ) + assert result.success is True + args, _stdin = cli.calls[0] + assert "--reply-to" in args + assert args[args.index("--reply-to") + 1] == "root-event-abc" + @pytest.mark.asyncio async def test_send_prefers_stable_thread_root_over_latest_reply(self): adapter = _make_adapter() @@ -917,6 +940,7 @@ class TestBuzzAdapterSend: assert args[args.index("--reply-to") + 1] == "stable-root" + @pytest.mark.asyncio async def test_send_image_local_file_uses_file_flag(self, tmp_path): img = tmp_path / "shot.png" diff --git a/tests/gateway/test_run_progress_topics.py b/tests/gateway/test_run_progress_topics.py index da3dd2efc7..c955605070 100644 --- a/tests/gateway/test_run_progress_topics.py +++ b/tests/gateway/test_run_progress_topics.py @@ -1925,3 +1925,14 @@ class TestSlackReplyInThreadProgressRouting: event_message_id="1700000000.000100", reply_in_thread=False, ) is None + + def test_buzz_uses_event_message_id_as_progress_thread(self): + """Buzz has no native thread_id; progress must reply-to the trigger.""" + from gateway.run import _resolve_progress_thread_id + + assert _resolve_progress_thread_id( + "buzz", + source_thread_id=None, + event_message_id="evt-trigger-001", + reply_in_thread=True, + ) == "evt-trigger-001"