Merge branch 'main' into feat/relay-slack-parity
One conflict, gateway/relay/adapter.py send_for_platform: main added the turn-final draft-seal interception (_sfp_metadata with the _interim_send marker stripped, seal-or-fall-through); this branch added format-hint stamping on the same frame. COMPOSED: the plain-send frame now stamps _with_format_hints_for_platform over _sfp_metadata (the stripped copy), so both the seal fall-through contract and the cron-lane block hints hold. Note: the seal frame itself (op:draft final) does not stamp hints — cron sends are never open drafts, so the flagship path is unaffected; noted as a connector-PR follow-up for streamed interactive finals.
This commit is contained in:
+613
-5
@@ -98,6 +98,41 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
# Bounded FIFO seen-set for inbound replay dedupe (finding #3);
|
||||
# dict preserves insertion order, giving cheap oldest-first eviction.
|
||||
self._seen_inbound: Dict[str, None] = {}
|
||||
# chat_id -> draft_id of the currently OPEN native draft stream
|
||||
# (NS-658 live cards). Armed by send_draft on a successful frame;
|
||||
# consumed by send() to convert the turn-final delivery into the
|
||||
# sealing draft(final=true) frame instead of a duplicate post.
|
||||
# Keyed by _draft_key (chat + per-turn identity), NOT bare chat:
|
||||
# parallel turns in one DM are distinct streams (live finding #10 —
|
||||
# per-chat keying collided three concurrent turns: merged task
|
||||
# cards, clobbered seal state, 3x duplicate finals).
|
||||
self._open_draft_by_chat: Dict[str, int] = {}
|
||||
# chat_id -> draft_id of the most recently SEALED stream (gateway
|
||||
# mirror of the connector's sealed-key tombstone): post-seal
|
||||
# straggler frames must neither re-arm interception nor re-open a
|
||||
# stream. One entry per turn key; a NEW turn's fresh draft_id
|
||||
# differs, so it arms normally and writes its own tombstone at its
|
||||
# own seal. Bounded like the sibling caches (see send_draft).
|
||||
self._sealed_draft_by_chat: Dict[str, int] = {}
|
||||
# Stream-is-the-message marker (finding #4): the stream consumer
|
||||
# checks this to keep ONE draft stream per turn instead of bumping
|
||||
# draft_id at tool boundaries (which opens a new Slack message per
|
||||
# segment on native streaming — Telegram-shaped adapters want the
|
||||
# bump, we don't).
|
||||
#
|
||||
# SLACK-ONLY semantic, gated on the negotiated descriptor (review
|
||||
# B4): the base send_draft contract is Telegram-shaped — the draft
|
||||
# clears and the final arrives as a separate real send. Setting
|
||||
# this unconditionally made ANY relay connector that advertises
|
||||
# the draft op (e.g. a Telegram connector) intercept the turn-final
|
||||
# into draft(final=true), so no real history message was ever
|
||||
# posted. A future connector platform whose native streaming is
|
||||
# also stream-is-the-message should advertise it explicitly
|
||||
# (descriptor field within the contract) rather than widening this
|
||||
# platform check by guesswork.
|
||||
self.draft_stream_is_message = (
|
||||
str(getattr(descriptor, "platform", "") or "").lower() == "slack"
|
||||
)
|
||||
# chat_id -> event fired when the entry above lands, so a consumer that
|
||||
# arrives before the send can wait for it instead of polling. See
|
||||
# wait_for_auto_thread_info.
|
||||
@@ -254,8 +289,504 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
self,
|
||||
chat_type: Optional[str] = None,
|
||||
metadata: Optional[Dict[str, Any]] = None,
|
||||
chat_id: Optional[str] = None,
|
||||
) -> bool:
|
||||
return self.descriptor.supports_draft_streaming
|
||||
# Native draft streaming needs BOTH the descriptor flag and the
|
||||
# "draft" op. supported_ops is fail-open for legacy connectors
|
||||
# (empty tuple = pre-contract ops only), but "draft" did not exist
|
||||
# pre-contract, so it must NOT fail open: an explicit advertisement
|
||||
# is required. Without it the stream consumer stays on the
|
||||
# edit-based path exactly as today.
|
||||
#
|
||||
# Per-chat resolution (review r2, finding 2): one adapter fronts N
|
||||
# platforms, and the scalar descriptor only reflects the PRIMARY
|
||||
# identity — a Telegram primary must not starve a secondary Slack
|
||||
# chat of native streaming, nor vice versa. When the caller can
|
||||
# name the chat, resolve through its platform's negotiated
|
||||
# descriptor; the scalar remains the fallback (chat unknown,
|
||||
# single-platform gateways: identical behavior).
|
||||
desc = (
|
||||
self._descriptor_for_chat(str(chat_id))
|
||||
if chat_id is not None
|
||||
else self.descriptor
|
||||
)
|
||||
return (
|
||||
desc.supports_draft_streaming
|
||||
and "draft" in (desc.supported_ops or ())
|
||||
)
|
||||
|
||||
def stream_is_message_for_chat(self, chat_id: str) -> bool:
|
||||
"""Per-chat stream-is-the-message semantic (review r2, finding 2).
|
||||
|
||||
The class-level ``draft_stream_is_message`` can only reflect the
|
||||
primary identity's platform. On a multi-platform relay, a Slack
|
||||
primary must not impose seal semantics on a Telegram chat (its
|
||||
turn-final would become draft(final=true) — no history message),
|
||||
and a Telegram primary must not deny a secondary Slack chat its
|
||||
native streaming. Resolve through the chat's own negotiated
|
||||
descriptor. Platform-name inference is deliberate for now — a
|
||||
descriptor-level field is the eventual contract (gg follow-up)
|
||||
so a future platform can advertise the semantic explicitly.
|
||||
"""
|
||||
return (
|
||||
str(self._descriptor_for_chat(str(chat_id)).platform or "").lower()
|
||||
== "slack"
|
||||
)
|
||||
|
||||
# ── Live cards: native draft streaming + task cards (NS-658) ─────────
|
||||
#
|
||||
# Additive relay ops within contract v1. The gateway side is dumb: it
|
||||
# emits ops when the negotiated descriptor advertises them; the
|
||||
# connector owns the platform API mechanics (chat.startStream et al.),
|
||||
# per-workspace feature-gate caching, and the send+edit fallback.
|
||||
#
|
||||
# Semantic bridge: the base send_draft contract is Telegram-shaped —
|
||||
# the draft clears and the final answer arrives as a separate send().
|
||||
# Slack native streaming makes the stream THE message, sealed once.
|
||||
# The adapter tracks the open draft per chat; the turn-final send()
|
||||
# for that chat converts to draft(final=true) so the connector seals
|
||||
# the stream instead of posting a duplicate message.
|
||||
|
||||
def supports_native_task_cards(self) -> bool:
|
||||
"""Descriptor probe for the TurnRunner's task-card lane.
|
||||
|
||||
Explicit advertisement required — same no-fail-open rule as
|
||||
"draft" (the op did not exist pre-contract).
|
||||
"""
|
||||
return "task_card" in (self.descriptor.supported_ops or ())
|
||||
|
||||
def native_task_cards_enabled(self) -> bool:
|
||||
"""TurnRunner opt-in probe (gateway/run.py) — the card lane calls
|
||||
THIS name (same contract as the native Slack adapter's opt-in);
|
||||
``supports_native_task_cards`` is the descriptor-level capability.
|
||||
Live-canary finding: without this alias the lane silently stays
|
||||
text-mode (hasattr probe fails) even though the connector
|
||||
advertises task_card."""
|
||||
return self.supports_native_task_cards()
|
||||
|
||||
@staticmethod
|
||||
def _draft_key(chat_id: str, metadata: Optional[Dict[str, Any]]) -> str:
|
||||
"""Coordination key for one turn's stream.
|
||||
|
||||
Prefers a PER-TURN identity — the triggering inbound message id
|
||||
(``message_id`` is stamped by the gateway's Slack thread metadata,
|
||||
``reply_to_message_id`` by the consumer's send path; both carry the
|
||||
same event id) — over the thread anchor. Finding #10 keyed on the
|
||||
thread anchor alone, which is simultaneously too coarse and too
|
||||
fragile (review B2 + flat-DM concern):
|
||||
|
||||
- two parallel turns REPLYING INSIDE ONE THREAD share thread_ts,
|
||||
so turn A's final sealed turn B's stream with A's content;
|
||||
- a flat DM whose metadata carries no anchor at all degraded to
|
||||
the bare chat id, re-creating the original #10 collision.
|
||||
|
||||
The thread anchor remains the fallback for callers that only have
|
||||
placement metadata, and the bare chat is the last resort
|
||||
(single-turn semantics).
|
||||
"""
|
||||
md = metadata or {}
|
||||
turn_id = md.get("message_id") or md.get("reply_to_message_id")
|
||||
if turn_id:
|
||||
return f"{chat_id}:turn:{turn_id}"
|
||||
anchor = md.get("thread_ts") or md.get("thread_id") or ""
|
||||
return f"{chat_id}:{anchor}"
|
||||
|
||||
# Cap for the draft/seal coordination dicts, matching the sibling
|
||||
# bounded caches (_auto_thread_by_chat). Entries are per-turn keys;
|
||||
# 512 in-flight-or-recent turns per adapter is far beyond any real
|
||||
# concurrency, and matches the connector's tombstone store size.
|
||||
_DRAFT_STATE_CAP = 512
|
||||
|
||||
@classmethod
|
||||
def _evict_oldest(cls, d: Dict[str, int]) -> None:
|
||||
"""FIFO-bound a coordination dict in place (review M1)."""
|
||||
while len(d) > cls._DRAFT_STATE_CAP:
|
||||
d.pop(next(iter(d)), None)
|
||||
|
||||
@staticmethod
|
||||
def _card_key(
|
||||
reply_to: Optional[str], metadata: Optional[Dict[str, Any]]
|
||||
) -> str:
|
||||
"""Per-turn task-card identity — same precedence as ``_draft_key``.
|
||||
|
||||
``reply_to`` (the triggering message id from the TurnRunner) wins;
|
||||
metadata message ids cover the flat-DM / resolver lanes; the thread
|
||||
anchor is only a fallback because two turns replying inside one
|
||||
thread share ``thread_ts`` and must not share a card (review B2).
|
||||
One derivation for send AND stop, so the stop always hits the
|
||||
stream the send opened.
|
||||
"""
|
||||
md = metadata or {}
|
||||
anchor = (
|
||||
reply_to
|
||||
or md.get("message_id")
|
||||
or md.get("reply_to_message_id")
|
||||
or md.get("thread_ts")
|
||||
or md.get("thread_id")
|
||||
or "root"
|
||||
)
|
||||
return f"turn:{anchor}"
|
||||
|
||||
def _match_open_draft(
|
||||
self, chat_id: str, metadata: Optional[Dict[str, Any]]
|
||||
) -> Optional[str]:
|
||||
"""Resolve which open stream (if any) a turn-final send belongs to.
|
||||
|
||||
Exact key match first. Callers WITHOUT a per-turn message id fall
|
||||
into two classes (review r2, finding 5):
|
||||
|
||||
- placement-only metadata (thread_ts/thread_id but no message
|
||||
id — legacy resolver lanes): the thread anchor is placement
|
||||
info, not turn identity, and streams are keyed per turn — an
|
||||
exact match will practically never fire for them. They may
|
||||
still absorb the final via the single-open-stream fallback.
|
||||
- no metadata at all: same fallback.
|
||||
|
||||
The fallback only fires when the chat has EXACTLY one open
|
||||
stream. With several open, the send stays a plain send: a
|
||||
duplicate message is recoverable, sealing someone else's stream
|
||||
with the wrong content is not (review B2). Callers that DO carry
|
||||
a message id never fall back — their identity is authoritative,
|
||||
and a mismatch means the stream is someone else's.
|
||||
"""
|
||||
key = self._draft_key(str(chat_id), metadata)
|
||||
if key in self._open_draft_by_chat:
|
||||
return key
|
||||
md = metadata or {}
|
||||
# Only a per-turn MESSAGE id is turn identity. Thread anchors are
|
||||
# placement info shared by every turn in the thread — treating
|
||||
# them as identity made the single-open-stream fallback dead for
|
||||
# placement-only callers (probed: plain final beside an open
|
||||
# turn-keyed stream).
|
||||
if md.get("message_id") or md.get("reply_to_message_id"):
|
||||
return None
|
||||
prefix = f"{chat_id}:"
|
||||
candidates = [
|
||||
k for k in self._open_draft_by_chat if k.startswith(prefix)
|
||||
]
|
||||
if len(candidates) == 1:
|
||||
return candidates[0]
|
||||
return None
|
||||
|
||||
async def send_draft(
|
||||
self,
|
||||
chat_id: str,
|
||||
draft_id: int,
|
||||
content: str,
|
||||
metadata: Optional[Dict[str, Any]] = None,
|
||||
) -> SendResult:
|
||||
if not self.supports_draft_streaming(chat_id=str(chat_id)):
|
||||
raise NotImplementedError(
|
||||
"connector does not advertise the 'draft' relay op"
|
||||
)
|
||||
if self._transport is None:
|
||||
return SendResult(success=False, error="no transport")
|
||||
# Audit fix G-D1 + regression fix (2026-08-15): arm optimistically
|
||||
# BEFORE the transport call (lossy ack: a timeout/WS-drop 'failure'
|
||||
# often means delivered), but NEVER for a draft_id that has already
|
||||
# been sealed this chat — the gateway-side mirror of the connector's
|
||||
# sealed-key tombstone. Without this, a straggler frame arriving
|
||||
# after the seal re-armed interception with no live stream, and the
|
||||
# next unrelated send (media follow-up, next-turn text) was wrongly
|
||||
# converted to draft(final=true) on the tombstoned key — clearing
|
||||
# the tombstone, re-opening a stream, and freezing it (the observed
|
||||
# escalating-frozen-prefixes regression).
|
||||
chat_key = self._draft_key(str(chat_id), metadata)
|
||||
if self._sealed_draft_by_chat.get(chat_key) == draft_id:
|
||||
# Post-seal straggler: its content is already in the sealed
|
||||
# message; report success, send nothing, arm nothing.
|
||||
return SendResult(success=True)
|
||||
# Arm seal-interception ONLY for stream-is-the-message chats
|
||||
# (review B4, per-chat in r2 finding 2): on a Telegram-shaped
|
||||
# connector the draft clears client-side and the final MUST go out
|
||||
# as a separate real send — arming here would intercept that final
|
||||
# into draft(final=true) and no history message would ever be
|
||||
# posted. Resolved per chat: one adapter fronts N platforms.
|
||||
if self.stream_is_message_for_chat(str(chat_id)):
|
||||
self._open_draft_by_chat[chat_key] = draft_id
|
||||
self._evict_oldest(self._open_draft_by_chat)
|
||||
try:
|
||||
result = await self._transport.send_outbound(
|
||||
{
|
||||
"op": "draft",
|
||||
"chat_id": chat_id,
|
||||
"draft_id": draft_id,
|
||||
"content": content,
|
||||
"final": False,
|
||||
"metadata": self._with_scope(chat_id, dict(metadata or {})),
|
||||
},
|
||||
platform=self._platform_by_chat.get(str(chat_id)),
|
||||
)
|
||||
except Exception as e:
|
||||
# Ambiguous by definition (stale socket, mid-write drop): the
|
||||
# frame may have been delivered. Keep interception armed.
|
||||
return SendResult(success=False, error=f"draft transport error: {e}")
|
||||
if result.get("success"):
|
||||
return SendResult(success=True)
|
||||
if result.get("ambiguous"):
|
||||
# Ack lost (transport timeout) — the production ws transport
|
||||
# RETURNS this shape rather than raising (PR 85796 review,
|
||||
# round 2): the connector may have applied the frame. Same
|
||||
# contract as the except branch: keep interception armed.
|
||||
return SendResult(
|
||||
success=False, error=str(result.get("error") or "draft ack lost")
|
||||
)
|
||||
# DEFINITE connector rejection (an explicit non-ambiguous result):
|
||||
# disarm interception for this key. The stream consumer disables
|
||||
# the draft transport on this failure and falls back to edit-based
|
||||
# streaming — its turn-final must go out as a REAL send, not get
|
||||
# converted into a seal on a stream the connector just told us is
|
||||
# unusable. (This restores the disarm-on-failure semantics the
|
||||
# G-D1 optimistic-arming change silently dropped; the ambiguity
|
||||
# that motivated G-D1 lives in the except branch above and the
|
||||
# ambiguous-result branch — both keep the key armed.)
|
||||
if self._open_draft_by_chat.get(chat_key) == draft_id:
|
||||
self._open_draft_by_chat.pop(chat_key, None)
|
||||
return SendResult(
|
||||
success=False, error=str(result.get("error") or "draft failed")
|
||||
)
|
||||
|
||||
async def _seal_open_draft(
|
||||
self,
|
||||
chat_id: str,
|
||||
content: str,
|
||||
metadata: Optional[Dict[str, Any]],
|
||||
*,
|
||||
draft_key: Optional[str] = None,
|
||||
) -> SendResult:
|
||||
"""Convert the turn-final send into the sealing draft frame."""
|
||||
if draft_key is None:
|
||||
draft_key = self._draft_key(str(chat_id), metadata)
|
||||
draft_id = self._open_draft_by_chat.pop(draft_key)
|
||||
# Tombstone BEFORE the transport call (regression fix): whatever the
|
||||
# ack says, this draft_id's stream must never be re-armed by a
|
||||
# straggler frame — the connector-side tombstone handles its half.
|
||||
self._sealed_draft_by_chat[draft_key] = draft_id
|
||||
# Bounded like the sibling caches (review M1): the key embeds a
|
||||
# per-turn identity, so an unbounded dict grows one entry per turn
|
||||
# for the life of the process. FIFO eviction matches the
|
||||
# straggler window this tombstone exists for (seconds, not days);
|
||||
# the connector holds its own 512-entry tombstone store.
|
||||
self._evict_oldest(self._sealed_draft_by_chat)
|
||||
if self._transport is None:
|
||||
return SendResult(success=False, error="no transport")
|
||||
seal_frame = {
|
||||
"op": "draft",
|
||||
"chat_id": chat_id,
|
||||
"draft_id": draft_id,
|
||||
"content": content,
|
||||
"final": True,
|
||||
"metadata": self._with_scope(chat_id, dict(metadata or {})),
|
||||
}
|
||||
_seal_platform = self._platform_by_chat.get(str(chat_id))
|
||||
_transport = self._transport # narrowed by the None-guard above
|
||||
|
||||
async def _attempt() -> Optional[Dict[str, Any]]:
|
||||
"""One seal attempt; None means ambiguous (exception or lost ack)."""
|
||||
try:
|
||||
r = await _transport.send_outbound(
|
||||
seal_frame, platform=_seal_platform
|
||||
)
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception as e:
|
||||
logger.warning("relay seal transport error (ambiguous): %s", e)
|
||||
return None
|
||||
if r.get("ambiguous"):
|
||||
# The production ws transport returns this shape on ack
|
||||
# timeout instead of raising (PR 85796 review, round 2):
|
||||
# the connector may have sealed and lost only the ack.
|
||||
logger.warning(
|
||||
"relay seal ack lost (ambiguous): %s", r.get("error")
|
||||
)
|
||||
return None
|
||||
return r
|
||||
|
||||
# Ambiguous outcomes (exception OR timeout-shaped result) retry the
|
||||
# SAME idempotent frame once: the connector's sealed-key tombstone
|
||||
# returns the original stream ts for a repeated final and never
|
||||
# opens a second stream, so the retry can turn "unknown" into a
|
||||
# definite answer for free. Only after BOTH attempts stay ambiguous
|
||||
# do we report failure — the caller's fail-open plain send is a
|
||||
# possible duplicate, but a silent loss is worse, and two
|
||||
# consecutive ack losses on one socket almost always mean the
|
||||
# transport is actually down (so the plain send fails too and the
|
||||
# gateway's fallback owns delivery).
|
||||
#
|
||||
# Cancellation safety (review r2, finding 4): the open entry was
|
||||
# popped and the tombstone written BEFORE the await. If the task is
|
||||
# cancelled mid-seal, CancelledError bypasses the failure handling
|
||||
# and the later abandon pass would find nothing to close — the
|
||||
# connector-side stream stays visibly live until eviction. Restore
|
||||
# the open entry (and drop our premature tombstone) before
|
||||
# re-raising so the abandon path can seal it.
|
||||
try:
|
||||
result = await _attempt()
|
||||
if result is None:
|
||||
result = await _attempt()
|
||||
except asyncio.CancelledError:
|
||||
self._open_draft_by_chat[draft_key] = draft_id
|
||||
if self._sealed_draft_by_chat.get(draft_key) == draft_id:
|
||||
self._sealed_draft_by_chat.pop(draft_key, None)
|
||||
raise
|
||||
if result is None:
|
||||
return SendResult(
|
||||
success=False,
|
||||
error="draft seal ambiguous after retry (transport ack lost)",
|
||||
)
|
||||
if result.get("success"):
|
||||
# The connector returns the stream's ts as the message identity.
|
||||
return SendResult(
|
||||
success=True,
|
||||
message_id=str(result.get("message_id") or "") or None,
|
||||
)
|
||||
return SendResult(
|
||||
success=False, error=str(result.get("error") or "draft seal failed")
|
||||
)
|
||||
|
||||
async def send_native_task_card_progress(
|
||||
self,
|
||||
chat_id: str,
|
||||
tasks: list,
|
||||
*,
|
||||
title: str = "Hermes is working",
|
||||
reply_to: Optional[str] = None,
|
||||
metadata: Optional[Dict[str, Any]] = None,
|
||||
fallback_text: Optional[str] = None,
|
||||
) -> SendResult:
|
||||
"""Relay leg of the #85476 task-card lane: emit one card frame.
|
||||
|
||||
SIGNATURE CONTRACT (live-canary finding): the TurnRunner calls this
|
||||
with the NATIVE Slack adapter's keyword contract (tasks/title/
|
||||
reply_to/metadata/fallback_text) — not a card_id. The card stream
|
||||
key is derived per (chat, reply_to-thread): one card per turn
|
||||
thread, matching the connector's (channel, card_id) keying.
|
||||
``fallback_text``/``title`` are accepted for contract parity; the
|
||||
connector's plan-mode stream renders task chunks, so they are not
|
||||
forwarded.
|
||||
|
||||
``tasks`` are the TurnRunner's normalized task dicts (id/title/
|
||||
status/details/output); the connector maps them onto its
|
||||
workspace-scoped card stream (task_update chunks, 256-char field
|
||||
limits enforced connector-side where the API lives).
|
||||
"""
|
||||
if not self.supports_native_task_cards():
|
||||
return SendResult(
|
||||
success=False, error="connector does not advertise task_card"
|
||||
)
|
||||
if self._transport is None:
|
||||
return SendResult(success=False, error="no transport")
|
||||
# Finding #10 + review B2: one card per TURN. reply_to (the
|
||||
# triggering message id) is already per-turn; when it is absent
|
||||
# (flat DM, resolver lanes) fall back to the same per-turn
|
||||
# identity the draft lane keys on — metadata message ids first,
|
||||
# thread anchor only after that (two turns replying inside one
|
||||
# thread share thread_ts and must not share a card).
|
||||
card_id = self._card_key(reply_to, metadata)
|
||||
merged_meta = dict(metadata or {})
|
||||
if reply_to and "thread_ts" not in merged_meta:
|
||||
# Slack card streams are thread replies (same rule as draft):
|
||||
# anchor on the triggering message when the runner gave us one.
|
||||
merged_meta["thread_ts"] = str(reply_to)
|
||||
try:
|
||||
result = await self._transport.send_outbound(
|
||||
{
|
||||
"op": "task_card",
|
||||
"chat_id": chat_id,
|
||||
"card_id": card_id,
|
||||
"chunks": [dict(t) for t in tasks],
|
||||
"metadata": self._with_scope(chat_id, merged_meta),
|
||||
},
|
||||
platform=self._platform_by_chat.get(str(chat_id)),
|
||||
)
|
||||
except Exception as e:
|
||||
# Progress is advisory: a transport drop must degrade to the
|
||||
# TurnRunner's text fallback (failed SendResult), never raise
|
||||
# into the progress loop / turn-cleanup path (review B7 — an
|
||||
# escaping card exception in cleanup skipped final delivery).
|
||||
return SendResult(
|
||||
success=False, error=f"task_card transport error: {e}"
|
||||
)
|
||||
if result.get("success"):
|
||||
return SendResult(success=True)
|
||||
return SendResult(
|
||||
success=False, error=str(result.get("error") or "task_card failed")
|
||||
)
|
||||
|
||||
async def stop_native_task_card_progress(
|
||||
self,
|
||||
chat_id: str,
|
||||
*,
|
||||
reply_to: Optional[str] = None,
|
||||
metadata: Optional[Dict[str, Any]] = None,
|
||||
) -> SendResult:
|
||||
"""Seal the card stream at turn end (idempotent connector-side).
|
||||
|
||||
Same NATIVE-contract signature as send (canary finding above);
|
||||
card key derived identically so the stop hits the open stream.
|
||||
"""
|
||||
if not self.supports_native_task_cards():
|
||||
return SendResult(
|
||||
success=False, error="connector does not advertise task_card"
|
||||
)
|
||||
if self._transport is None:
|
||||
return SendResult(success=False, error="no transport")
|
||||
# Same per-turn key derivation as send (shared helper) so the stop
|
||||
# hits the open stream.
|
||||
card_id = self._card_key(reply_to, metadata)
|
||||
try:
|
||||
result = await self._transport.send_outbound(
|
||||
{
|
||||
"op": "task_card_stop",
|
||||
"chat_id": chat_id,
|
||||
"card_id": card_id,
|
||||
"metadata": self._with_scope(chat_id, dict(metadata or {})),
|
||||
},
|
||||
platform=self._platform_by_chat.get(str(chat_id)),
|
||||
)
|
||||
except Exception as e:
|
||||
# Best-effort by contract: the stop runs in the progress loop's
|
||||
# finally block on the turn-cleanup path — an escaping transport
|
||||
# exception there skipped final delivery (review B7). The
|
||||
# connector seals orphaned card streams on its own (recycling /
|
||||
# eviction), so a lost stop is cosmetic.
|
||||
return SendResult(
|
||||
success=False, error=f"task_card_stop transport error: {e}"
|
||||
)
|
||||
return SendResult(success=bool(result.get("success")))
|
||||
|
||||
async def abandon_open_draft(
|
||||
self,
|
||||
chat_id: str,
|
||||
content: str,
|
||||
metadata: Optional[Dict[str, Any]] = None,
|
||||
) -> SendResult:
|
||||
"""Seal an orphaned stream when its turn dies (review B8).
|
||||
|
||||
A stopped (/stop, /new) or superseded turn previously left its
|
||||
native stream open forever: the Slack message kept the live
|
||||
streaming indicator and the adapter kept armed interception state,
|
||||
which the NEXT turn's key could inherit. Seal in place with
|
||||
``content`` — the text already on screen (the consumer passes its
|
||||
last delivered frame), so the seal adds nothing and claims
|
||||
nothing: it only ends the stream. Delivery flags are the
|
||||
consumer's business; this never sets any.
|
||||
|
||||
Best-effort by contract: failure is reported, never raised — the
|
||||
connector reaps truly orphaned streams via recycling/eviction.
|
||||
"""
|
||||
draft_key = self._match_open_draft(str(chat_id), metadata)
|
||||
if draft_key is None:
|
||||
return SendResult(success=True) # nothing armed — no-op
|
||||
try:
|
||||
return await self._seal_open_draft(
|
||||
chat_id, content, metadata, draft_key=draft_key
|
||||
)
|
||||
except Exception as e:
|
||||
return SendResult(
|
||||
success=False, error=f"abandon seal transport error: {e}"
|
||||
)
|
||||
|
||||
|
||||
# ── abstract methods (delegated to the transport) ────────────────────
|
||||
async def connect(self, *, is_reconnect: bool = False) -> bool:
|
||||
@@ -408,18 +939,33 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
_SEEN_INBOUND_MAX = 512
|
||||
|
||||
def _inbound_dedupe_key(self, event) -> Optional[str]:
|
||||
"""Stable replay identity: (chat, platform message id).
|
||||
"""Stable replay identity: (platform, chat, platform message id).
|
||||
|
||||
Chat identity lives on ``event.source`` (MessageEvent has no top-level
|
||||
``chat_id``), and this adapter can front SEVERAL platforms over one
|
||||
relay socket (Phase 1.5 multiplex), so the underlying platform joins
|
||||
the key — two platforms' numeric chat/message ids must never collide
|
||||
into one identity.
|
||||
|
||||
Returns None when the event carries no platform message id (synthetic
|
||||
events, some prompt responses) — those never dedupe, fail-open by
|
||||
design: dropping a real user message is strictly worse than rerunning
|
||||
one, so only dedupe when identity is certain.
|
||||
"""
|
||||
source = getattr(event, "source", None)
|
||||
message_id = getattr(event, "message_id", None)
|
||||
chat_id = getattr(event, "chat_id", None)
|
||||
chat_id = getattr(source, "chat_id", None)
|
||||
if not message_id or not chat_id:
|
||||
return None
|
||||
return f"{chat_id}:{message_id}"
|
||||
# Normalize the platform component: production wire decoding always
|
||||
# yields a Platform enum (unknowns canonicalize to Platform.RELAY),
|
||||
# but alternate constructors may carry the plain string. Use the
|
||||
# enum's value when present, the string itself otherwise — both
|
||||
# spellings of one platform must produce ONE key, and two different
|
||||
# string platforms must not collapse into the same empty component.
|
||||
raw_platform = getattr(source, "platform", None)
|
||||
platform = getattr(raw_platform, "value", raw_platform) or ""
|
||||
return f"{platform}:{chat_id}:{message_id}"
|
||||
|
||||
def _relay_slack_extra(self) -> Dict[str, Any]:
|
||||
"""The Slack-behavior subset of the RELAY platform config.
|
||||
@@ -1100,6 +1646,31 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
success=False,
|
||||
error=f"relay does not front platform {platform_value}",
|
||||
)
|
||||
_sfp_metadata = dict(metadata or {})
|
||||
# Gateway-internal interim marker (see send()): strip before the
|
||||
# wire; an interim send through this door also skips interception.
|
||||
_interim = bool(_sfp_metadata.pop("_interim_send", False))
|
||||
# Finding #7 (live canary): the delivery resolver calls THIS method
|
||||
# directly (gateway/delivery.py), bypassing send() — an open native
|
||||
# stream must absorb the turn-final here too, or the stream is left
|
||||
# unsealed (frozen live indicator) and the final posts as a separate
|
||||
# duplicate message.
|
||||
if not _interim:
|
||||
_sfp_key = self._match_open_draft(str(chat_id), _sfp_metadata)
|
||||
else:
|
||||
_sfp_key = None
|
||||
if _sfp_key is not None:
|
||||
seal = await self._seal_open_draft(
|
||||
chat_id, content, _sfp_metadata, draft_key=_sfp_key
|
||||
)
|
||||
if seal.success:
|
||||
return seal
|
||||
# Failed seal falls through to the plain send below (review
|
||||
# finding, PR 85796 point 1): never swallow the turn-final.
|
||||
logger.warning(
|
||||
"relay seal failed (%s); delivering turn-final as plain send",
|
||||
seal.error,
|
||||
)
|
||||
if self._transport is None:
|
||||
return SendResult(success=False, error="no transport")
|
||||
result = await self._transport.send_outbound(
|
||||
@@ -1111,10 +1682,12 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
# format_hints on the explicit-platform lane too: this is the
|
||||
# scheduled/cron delivery path — the in_channel brief itself —
|
||||
# and it must render blocks exactly like an interactive send.
|
||||
# Stamps _sfp_metadata (the interim-marker-stripped copy, per
|
||||
# the seal path above), composing both sides of the merge.
|
||||
"metadata": self._with_scope(
|
||||
chat_id,
|
||||
self._with_format_hints_for_platform(
|
||||
str(platform_value), metadata
|
||||
str(platform_value), _sfp_metadata
|
||||
),
|
||||
),
|
||||
},
|
||||
@@ -1231,6 +1804,41 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
) -> SendResult:
|
||||
send_metadata = dict(metadata or {})
|
||||
explicit_platform = send_metadata.pop("_relay_logical_platform", None)
|
||||
# Consumer-declared interim send (commentary, tail flush): NOT the
|
||||
# turn-final, so it must never trigger seal-interception — sealing
|
||||
# the live stream with interim text orphans the true final into a
|
||||
# plain-send duplicate (live finding, 2026-08-16 canary). The
|
||||
# marker is gateway-internal; strip before the wire.
|
||||
_interim = bool(send_metadata.pop("_interim_send", False))
|
||||
# NS-658 seal-interception — checked BEFORE the explicit-platform
|
||||
# branch (finding #7, live canary): the delivery-resolver lane
|
||||
# (follow-up queue, media-accompanied finals, scheduled sends) routes
|
||||
# through send_for_platform, which posted a plain send while the
|
||||
# native stream stayed open — the user got the stream frozen
|
||||
# mid-word (live indicator, never sealed) PLUS the final as a
|
||||
# separate message. An open stream absorbs the turn-final send no
|
||||
# matter which egress door it arrives through; the stream IS the
|
||||
# message.
|
||||
if not _interim:
|
||||
_send_key = self._match_open_draft(str(chat_id), send_metadata)
|
||||
else:
|
||||
_send_key = None
|
||||
if _send_key is not None:
|
||||
seal = await self._seal_open_draft(
|
||||
chat_id, content, send_metadata, draft_key=_send_key
|
||||
)
|
||||
if seal.success:
|
||||
return seal
|
||||
# Review finding (PR 85796, point 1): a failed seal must NOT
|
||||
# swallow the turn-final — the stream consumer has already
|
||||
# disabled the draft transport, so returning failure here means
|
||||
# the user never gets the answer. Fall through to a plain send
|
||||
# (the orphaned stream is sealed connector-side by recycling /
|
||||
# MAX_OPEN_STREAMS eviction).
|
||||
logger.warning(
|
||||
"relay seal failed (%s); delivering turn-final as plain send",
|
||||
seal.error,
|
||||
)
|
||||
if explicit_platform:
|
||||
return await self.send_for_platform(
|
||||
explicit_platform,
|
||||
|
||||
Reference in New Issue
Block a user