fix(relay): resolve fresh-final unfurl decision per chat, not per primary identity (#99206)

The stream consumer called prefers_fresh_final_streaming(text,
metadata=...) only, and no metadata producer stamps a platform key — so
RelayAdapter's hook always fell back to the PRIMARY descriptor's
platform (the scalar-vs-per-chat capability seam, third occurrence).
Two failure directions on multiplexed relays with
platforms.relay.extra.slack.unfurl_links/media: true (#97957):

- Slack primary fronting Telegram/Discord: every link-bearing streamed
  final on the non-Slack chats finalized as a fresh send with no delete
  op advertised -> the answer delivered TWICE (orphaned preview).
- Non-Slack primary fronting Slack: the hook returned False, leaving
  the force-on unfurl feature dark on exactly the chats it shipped for.

Pass chat_id=self.chat_id from the consumer; the relay hook already
accepted it and resolves via _platform_by_chat + the per-platform
negotiated descriptor. Graduated TypeError fallback keeps the
single-platform hook signatures (Telegram, base class) and legacy test
doubles working unchanged.

Both regression tests verified RED against the unfixed consumer, GREEN
with the fix; single-platform relays are unaffected (#97957's own 30
tests unchanged-green).
This commit is contained in:
Ben Barclay
2026-08-31 16:28:21 +10:00
committed by GitHub
parent 6681f9ebc3
commit 1f99a4b2f2
2 changed files with 103 additions and 4 deletions
+16 -4
View File
@@ -2826,11 +2826,23 @@ class GatewayStreamConsumer:
return False
try:
try:
result = fn(text, metadata=self.metadata)
# Pass the chat id so multi-platform adapters (relay) resolve
# the decision through THIS chat's negotiated platform, not
# the primary identity's. Without it a Slack-primary relay
# with unfurl force-on misroutes a fronted Telegram/Discord
# chat's final through the fresh-send lane (duplicate
# delivery: those descriptors advertise no ``delete`` op),
# and the mirror posture leaves fronted Slack chats dark.
result = fn(text, metadata=self.metadata, chat_id=self.chat_id)
except TypeError:
# Adapter / test double whose hook doesn't accept the metadata
# keyword — fall back to the positional-only form.
result = fn(text)
try:
# Single-platform hook signature (Telegram, base class):
# (content, metadata=None) — no chat_id keyword.
result = fn(text, metadata=self.metadata)
except TypeError:
# Adapter / test double whose hook doesn't accept the
# metadata keyword — fall back to the positional-only form.
result = fn(text)
except Exception as e:
logger.debug("prefers_fresh_final_streaming check failed: %s", e)
return False
@@ -412,6 +412,93 @@ class TestConsumerRoutesForceOnFinalAsFreshSend:
assert ops[-1] == "edit", f"ops={ops}"
class TestMultiplexPerChatFreshFinalSeam:
"""The consumer must resolve the fresh-final decision through the CHAT's
negotiated platform, not the relay's primary identity (the scalar-vs-
per-chat descriptor seam). Two failure directions:
- Slack PRIMARY + force-on unfurl: a fronted TELEGRAM chat must keep the
edit lane — its descriptor advertises no ``delete`` op, so a fresh
final would deliver the answer twice (orphaned preview).
- Non-Slack primary: a fronted SLACK chat must still get the fresh
stamped final, or the shipped feature is dark on exactly the chats it
was built for.
"""
class _MultiplexTransport(_RecordingTransport):
def __init__(self):
super().__init__()
self.descriptors = {
"telegram": make_desc(
platform="telegram",
supports_edit=True,
supported_ops=("send", "edit", "typing"),
),
"slack": make_desc(
platform="slack",
supports_edit=True,
supported_ops=("send", "edit", "typing", "delete"),
),
}
def descriptor_for_platform(self, platform):
return self.descriptors.get(platform)
def _consumer(self, adapter, chat_id):
from gateway.stream_consumer import (
GatewayStreamConsumer,
StreamConsumerConfig,
)
return GatewayStreamConsumer(
adapter=adapter, chat_id=chat_id, config=StreamConsumerConfig(),
)
@pytest.mark.asyncio
async def test_telegram_chat_on_slack_primary_keeps_edit_lane(self):
transport = self._MultiplexTransport()
a = RelayAdapter(
PlatformConfig(extra={"slack": {"unfurl_links": True}}),
make_desc(platform="slack", supports_edit=True), # Slack PRIMARY
transport=transport,
)
# Relay learned this chat is Telegram from inbound traffic.
a._platform_by_chat["TG1"] = "telegram"
consumer = self._consumer(a, "TG1")
await consumer._send_or_edit("Working on it…")
await consumer._send_or_edit("see https://studiotwin.ai", finalize=True)
ops = [x.get("op") for x in transport.actions]
assert ops[-1] == "edit", (
f"telegram chat misrouted through fresh-final: ops={ops}"
)
@pytest.mark.asyncio
async def test_slack_chat_on_telegram_primary_gets_fresh_final(self):
transport = self._MultiplexTransport()
a = RelayAdapter(
PlatformConfig(extra={"slack": {"unfurl_links": True}}),
make_desc(platform="telegram", supports_edit=True), # non-Slack primary
transport=transport,
)
a._platform_by_chat["D1"] = "slack"
consumer = self._consumer(a, "D1")
await consumer._send_or_edit("Working on it…")
await consumer._send_or_edit("see https://studiotwin.ai", finalize=True)
ops = [x.get("op") for x in transport.actions]
# Fresh stamped final followed by cleanup of the sealed preview
# (Slack's descriptor advertises delete) — no duplicate.
assert "delete" in ops, f"preview not cleaned up: ops={ops}"
sends = [x for x in transport.actions if x.get("op") == "send"]
assert len(sends) == 2, f"slack chat on non-slack primary left dark: ops={ops}"
final = sends[-1]
assert final["content"] == "see https://studiotwin.ai"
assert final["metadata"]["unfurl_links"] is True
class TestDeleteOpForFreshFinalCleanup:
"""Relay delete_message: emitted only when the negotiated descriptor
advertises the additive `delete` op; older connectors degrade to the