fix(wecom): eliminate duplicate + split bubbles in native streaming

Fixes two related native-streaming bubble defects surfaced in production:

1. Duplicate bubble on long turns — when keep-alive already refreshed the
   6-min reply window, the Layer-2 clock fallback still declined the finalize
   frame and forced a proactive send(), duplicating the message. Skip the
   clock fallback while keep-alive is active; intermediate-frame failures are
   now fully fire-and-forget (only a failed FINAL frame falls back to send()).

2. Split / mini bubbles ('Cla' + 'ude ...') — two compounding root causes:
   a) In native streaming a mid-turn commentary (e.g. a Hindsight recall
      notice) called _reset_segment_state(), clearing the cumulative
      _accumulated so the next delta + finalize frame carried only the few
      chars accumulated after the reset. Native streaming now skips that
      reset (commentary still posts as its own message via send()).
   b) The adapter-side _BlockChunker.update() 'only grow' guard silently
      dropped any cumulative snapshot shorter than its high-water mark, so
      after a baseline reset the leading characters were stranded before
      _emitted_len. Removed the _BlockChunker sentence-alignment + idle-flush
      layer entirely; intermediate frames are pure identity-dedup, matching
      the fire-and-forget model.

Also removes ~232 lines of now-dead code (_BlockChunker class, idle-flush
machinery, block-stream constants) and aligns the test suite with the
fire-and-forget frame model, including a regression test that locks the
native-commentary-no-reset behavior.

Tests: 177 passed, 3 skipped (wecom + stream_consumer suites).
This commit is contained in:
wansui
2026-08-27 19:46:58 +08:00
committed by Teknium
parent 2ecb544551
commit 9faa953cd1
5 changed files with 530 additions and 452 deletions
+10
View File
@@ -1881,6 +1881,16 @@ class GatewayStreamConsumer:
if _stream_is_msg_c and self._use_draft_streaming:
await self._send_commentary(commentary_text)
self._last_edit_time = time.monotonic()
elif self._use_native_streaming:
# Native streaming (WeCom): commentary is sent as an
# independent message via adapter.send(), but we must
# NOT reset _accumulated — the native stream is
# cumulative and a reset would lose all pre-commentary
# text. Subsequent frames must still carry the full
# accumulated content. Same rationale as segment-break
# no-op for native streaming.
await self._send_commentary(commentary_text)
self._last_edit_time = time.monotonic()
else:
self._reset_segment_state()
await self._send_commentary(commentary_text)
+74 -288
View File
@@ -175,27 +175,6 @@ STREAM_SAFE_DURATION_SECONDS = 330.0 # 5.5 min — Layer 2 clock fallback
STREAM_KEEPALIVE_INTERVAL_SECONDS = 120.0 # 2 min — Layer 1 heartbeat cadence
STREAM_KEEPALIVE_ENABLED_DEFAULT = False # Layer 1 off unless config opts in
# ── Block-streaming parameters (aligned with the official wecom-openclaw-plugin) ──
# The official plugin (wecom-openclaw-plugin/src/webhook/helpers.ts) coalesces
# incoming LLM tokens into sentence-aligned blocks before sending each frame,
# rather than echoing every char-delta. This produces ~1 frame per natural
# sentence boundary instead of one per 60-char delta, which means:
# - frame cadence drops from "every 200ms" to "every 0.5-2s typical",
# - the 5s ack timeout almost never fires (frames are spaced wider apart),
# - and the user sees content land in complete thought-units rather than
# mid-sentence cursors moving every tick.
#
# Values copied verbatim from the official plugin so we match its behaviour
# in WeCom's rate-limit budget (≈30 frames/min/chat).
BLOCK_STREAM_MIN_CHARS = 120 # Don't emit a frame below this size (unless forced)
BLOCK_STREAM_MAX_CHARS = 360 # Force a break above this size
BLOCK_STREAM_IDLE_FLUSH = 0.25 # 250ms — flush partial buffer if no new tokens
# Sentence terminators recognised by the block chunker. Includes the most
# common Chinese full-stop / exclamation / question forms; matches the
# official plugin's "sentence" break preference for the WeCom channel.
_SENTENCE_TERMINATORS = ".!?。!?"
IMAGE_MAX_BYTES = 10 * 1024 * 1024
VIDEO_MAX_BYTES = 10 * 1024 * 1024
VOICE_MAX_BYTES = 2 * 1024 * 1024
@@ -285,124 +264,6 @@ class ReplyQueue:
self.pending_ack: Optional[ReplyFrame] = None
class _BlockChunker:
"""Coalesce streaming text into sentence-aligned blocks.
The consumer feeds us **cumulative** text (every frame is the full
response-so-far). We track the high-water mark and only emit when one
of these conditions holds for the *new* tail:
* new content is at least ``BLOCK_STREAM_MIN_CHARS`` AND ends on a
safe break (sentence terminator or blank-line paragraph boundary);
* new content has grown to ``BLOCK_STREAM_MAX_CHARS`` (hard cap —
force a break);
* ``force=True`` is passed (finalize, segment break, idle timer).
Returns the **cumulative text up to the emit point**, matching the
semantics of WeCom's stream protocol where every frame carries the
full content and the client renders the diff. This is consistent
with the official wecom-openclaw-plugin's flow: each ``deliver``
callback there carries the new block, and the plugin sends the
accumulated text via ``sendWeComReply`` under the same streamId.
"""
__slots__ = (
"_cumulative",
"_min_chars",
"_max_chars",
"_terminators",
"_emitted_len",
)
def __init__(
self,
*,
min_chars: int = BLOCK_STREAM_MIN_CHARS,
max_chars: int = BLOCK_STREAM_MAX_CHARS,
terminators: str = _SENTENCE_TERMINATORS,
) -> None:
self._cumulative: str = ""
self._min_chars = min_chars
self._max_chars = max_chars
self._terminators = terminators
# Length of the prefix that has already been signalled as ready to
# emit at least once. We compare new growth against this watermark.
self._emitted_len = 0
def update(self, cumulative_text: str) -> None:
"""Replace the high-water mark with a fresh cumulative snapshot."""
# Only grow — the consumer never shrinks its accumulator, but defensive
# against accidental resets so we don't lose ground.
if len(cumulative_text) > len(self._cumulative):
self._cumulative = cumulative_text
def has_pending(self) -> bool:
"""True when there is buffered content past the last emit point."""
return len(self._cumulative) > self._emitted_len
def drain(self, *, force: bool = False) -> Optional[str]:
"""Decide whether to emit a frame, return the cumulative text if so.
Returns ``None`` when the chunker wants to keep buffering, or a
cumulative string (full content up to the chosen break) when a
frame is ready to go on the wire.
"""
if not self.has_pending():
return None
new_len = len(self._cumulative) - self._emitted_len
# Hard cap — force a break even mid-sentence.
if new_len >= self._max_chars:
self._emitted_len = len(self._cumulative)
return self._cumulative
if force:
# Finalize / segment-break / idle-timer path: emit everything we
# have. Empty/whitespace-only tails are still allowed to pass —
# the caller (send_stream_frame) handles empty-text finalize.
self._emitted_len = len(self._cumulative)
return self._cumulative
if new_len < self._min_chars:
return None
# Look for a safe break inside the new tail. We scan backwards from
# the end of cumulative — the rightmost terminator that leaves a
# block of at least ``min_chars`` becomes the break point.
# Sentence terminators must be followed by whitespace / end-of-text
# so we don't split inside e.g. "v1.2" or "192.168.1.1".
new_tail_start = self._emitted_len
for idx in range(len(self._cumulative) - 1, new_tail_start - 1, -1):
ch = self._cumulative[idx]
if ch not in self._terminators:
continue
# Tail-end terminator is always safe.
if idx == len(self._cumulative) - 1:
next_ok = True
else:
next_ch = self._cumulative[idx + 1]
next_ok = next_ch.isspace() or next_ch in self._terminators
if not next_ok:
continue
break_at = idx + 1
if break_at - new_tail_start < self._min_chars:
# Sentence is too short; keep scanning leftward looking for a
# later terminator that satisfies min_chars. Once we fall
# below min_chars we know no further matches will satisfy it
# either (we're moving left, blocks shrink).
break
self._emitted_len = break_at
return self._cumulative[:break_at]
# Also accept paragraph boundary ("\n\n") inside the new tail.
para_idx = self._cumulative.rfind("\n\n", new_tail_start)
if para_idx >= 0 and (para_idx + 2 - new_tail_start) >= self._min_chars:
break_at = para_idx + 2
self._emitted_len = break_at
return self._cumulative[:break_at]
return None
class StreamTurn:
"""Per-turn stream state to avoid global state conflicts.
@@ -426,14 +287,8 @@ class StreamTurn:
# MAX_INTERMEDIATE_FRAMES to leave room for the finalize frame).
self._last_frame_sent_at: float = 0.0
self._intermediate_frames_sent: int = 0
# Block-stream chunker — coalesces consumer's cumulative cursor into
# sentence-aligned blocks before each frame goes on the wire. See
# ``_BlockChunker`` for the algorithm; lazily created on first append
# so per-turn cost is paid only for turns that actually stream.
self.chunker: Optional["_BlockChunker"] = None
# Idle flush handle — set when a partial buffer is waiting on the
# 250ms idle timer. Cleared when the timer fires, when a force-flush
# happens, or on cleanup.
# Idle flush handle — retained for _cancel_idle_flush() compatibility
# (called in finalize/boundary paths; always None in fire-and-forget).
self.idle_flush_handle: Optional[asyncio.TimerHandle] = None
# Keep-alive handle (Layer 1) — set when the stream-level keep-alive
# timer is armed. Structurally identical to idle_flush_handle: a
@@ -1983,99 +1838,6 @@ class WeComAdapter(BasePlatformAdapter):
pass
turn.idle_flush_handle = None
def _arm_idle_flush(
self,
turn: StreamTurn,
*,
turn_id: Optional[str],
) -> None:
"""Arm the 250ms idle-flush timer so a partial buffer still ships.
When the LLM pauses (model thinking, network slow), the chunker's
``min_chars`` gate would hold a half-sentence forever. The official
plugin's ``blockStreamingCoalesce.idleMs = 250`` solves this by
flushing whatever is buffered after 250ms of silence.
The timer is per-turn and idempotent — re-arming while one is
pending no-ops; a new timer is only scheduled after the previous
one has fired or been cancelled.
"""
if turn.idle_flush_handle is not None:
return # already armed
try:
loop = asyncio.get_running_loop()
except RuntimeError:
return # not inside a loop (defensive — should never happen here)
handle = loop.call_later(
BLOCK_STREAM_IDLE_FLUSH,
self._on_idle_flush_fire,
turn,
turn_id,
)
turn.idle_flush_handle = handle
def _on_idle_flush_fire(
self,
turn: StreamTurn,
turn_id: Optional[str],
) -> None:
"""Loop callback — dispatch an async flush task without blocking."""
turn.idle_flush_handle = None
if turn.finalized or turn.expired:
return
if turn.chunker is None or not turn.chunker.has_pending():
return
# Schedule on the running loop; swallow scheduling errors so a stale
# timer firing during shutdown can never bubble up.
try:
asyncio.ensure_future(self._idle_flush_send(turn, turn_id))
except RuntimeError:
pass
async def _idle_flush_send(
self,
turn: StreamTurn,
turn_id: Optional[str],
) -> None:
"""Drain the chunker's pending tail and send one intermediate frame.
Called from the idle-flush timer when the consumer hasn't pushed a
fresh cumulative snapshot for ``BLOCK_STREAM_IDLE_FLUSH`` seconds.
Force-drains the chunker (force=True) so even a sub-min_chars tail
reaches the user.
"""
if turn.finalized or turn.expired:
return
if turn._intermediate_frames_sent >= MAX_INTERMEDIATE_FRAMES:
return
chunker = turn.chunker
if chunker is None or not chunker.has_pending():
return
drained = chunker.drain(force=True)
if drained is None or not drained:
return
try:
await self._send_stream_reply(
turn.req_id,
turn.stream_id,
drained,
finish=False,
)
except WeComStreamExpiredError:
turn.expired = True
self._retire_turn(turn, turn_id)
self._stream_expired_chats.add(turn.chat_id)
return
except Exception as exc:
logger.debug(
"[%s] idle-flush send failed (chat=%s, turn=%s): %s",
self.name, turn.chat_id, turn.stream_id, exc,
)
return
turn._last_frame_sent_at = time.monotonic()
turn._intermediate_frames_sent += 1
turn.last_sent_content = drained
# ── Stream-level keep-alive (Layer 1) ─────────────────────────────────
# Structurally mirrors the idle-flush timer: a per-turn asyncio TimerHandle
# stored on the StreamTurn, cancelled on every turn-exit path. The only
@@ -3257,31 +3019,35 @@ class WeComAdapter(BasePlatformAdapter):
# False finalize). Zero new uplink frames; safe for groups in
# the sense that it does not make delivery any worse than the
# current 846608 fallback.
stream_age = time.monotonic() - turn.start_time
if stream_age >= self._stream_safe_duration_seconds:
logger.info(
"[%s] Stream age %.0fs >= safe duration %.0fs for chat "
"%s — declining finalize frame, falling back to "
"proactive send (Layer 2 clock fallback).",
self.name, stream_age,
self._stream_safe_duration_seconds, chat,
)
turn.expired = True
self._retire_turn(turn, turn_id)
self._stream_expired_chats.add(chat)
return False
#
# SKIP entirely when Layer 1 keep-alive is enabled: the
# heartbeat has been refreshing the stream window every
# STREAM_KEEPALIVE_INTERVAL_SECONDS, so an old stream_age does
# NOT mean the stream is dead. Declining a still-live stream on
# a blind clock read would force the consumer's send() fallback
# to re-deliver content the intermediate frames already put on
# screen — the exact duplicate-bubble bug this guards against.
# If the stream truly HAS expired, the _send_stream_reply(
# finish=True) below will hit 846604/846608 and raise
# WeComStreamExpiredError, which the except block turns into the
# real (finalize-only) fallback.
if not self._stream_keepalive_enabled:
stream_age = time.monotonic() - turn.start_time
if stream_age >= self._stream_safe_duration_seconds:
logger.info(
"[%s] Stream age %.0fs >= safe duration %.0fs for chat "
"%s — declining finalize frame, falling back to "
"proactive send (Layer 2 clock fallback).",
self.name, stream_age,
self._stream_safe_duration_seconds, chat,
)
turn.expired = True
self._retire_turn(turn, turn_id)
self._stream_expired_chats.add(chat)
return False
# Force-drain any buffered tail through the chunker so the
# final frame carries the latest content — this is what the
# official plugin does in `finishThinkingStream`. The
# chunker call also cancels the idle-flush timer.
self._cancel_idle_flush(turn)
self._cancel_keepalive(turn)
if turn.chunker is not None:
turn.chunker.update(text)
drained = turn.chunker.drain(force=True)
if drained is not None:
text = drained
# WeCom may silently drop (no ack) a final frame whose content
# is identical to the preceding intermediate frame — it treats
@@ -3306,53 +3072,59 @@ class WeComAdapter(BasePlatformAdapter):
else:
self._cleanup_stream_turn(chat, turn.req_id)
else:
# Block-streaming gate (aligned with wecom-openclaw-plugin):
# the consumer hands us a *cumulative* text snapshot every
# tick. Feed it to the chunker and only put a frame on the
# wire when a complete sentence-aligned block is ready.
# Otherwise schedule a 250ms idle flush so partial buffers
# don't sit forever when the LLM slows down.
if turn.chunker is None:
turn.chunker = _BlockChunker()
turn.chunker.update(text)
# Fire-and-forget: gateway already decides when to push
# (pure identity-dedup in stream_consumer.py). No adapter-
# side buffering — send immediately when content differs
# from the last pushed frame. This removes the _BlockChunker
# sentence-alignment layer whose "only grow" guard in
# update() silently dropped frames whenever gateway-side
# _accumulated was reset (commentary, boundary) and the new
# cumulative text was shorter than the chunker's high-water
# mark — the root cause of the "Cla"/"ude" split-bubble bug.
turn.accumulated_text = text
if turn._intermediate_frames_sent >= MAX_INTERMEDIATE_FRAMES:
# Frame cap reached — drop intermediates, keep accumulating.
# The finalize path will drain whatever is left.
turn.accumulated_text = text
return True
drained = turn.chunker.drain(force=False)
if drained is None:
# Not enough content for a block yet. Make sure the
# idle-flush timer is armed so a partial buffer still
# reaches the user when the LLM pauses.
self._arm_idle_flush(turn, turn_id=turn_id)
turn.accumulated_text = text
# Pure dedup: skip if content is identical to last sent frame.
if text == turn.last_sent_content:
return True
# We have a block — cancel any pending idle flush, send the
# frame, then re-arm only if there's still pending content
# in the chunker (i.e., another block partway grown).
self._cancel_idle_flush(turn)
await self._send_stream_reply(
turn.req_id,
turn.stream_id,
drained,
text,
finish=False,
)
turn._last_frame_sent_at = time.monotonic()
turn._intermediate_frames_sent += 1
turn.accumulated_text = text
turn.last_sent_content = drained
if turn.chunker.has_pending():
self._arm_idle_flush(turn, turn_id=turn_id)
turn.last_sent_content = text
return True
except WeComStreamExpiredError:
# Intermediate frames (finalize=False) are fire-and-forget: a later
# cumulative frame — or the finalize frame — carries the full text
# and overwrites whatever this one would have shown. A transient
# failure here must NOT flip the turn expired or trip the consumer's
# send() fallback; doing so re-delivers content the stream will
# replace anyway (duplicate bubble). The stream is still alive
# (keep-alive is refreshing it), so leave the turn intact and report
# success so the consumer keeps streaming. Only a FINAL frame's
# expiry means the screen is genuinely missing this content and the
# consumer must fall back.
if not finalize:
logger.info(
"[%s] Intermediate stream frame expired (errcode=%d) for "
"chat %s — dropping frame, stream stays live",
self.name, STREAM_EXPIRED_ERRCODE, chat,
)
return True
logger.info(
"[%s] Stream expired (errcode=%d) for chat %s — switching to proactive send",
self.name, STREAM_EXPIRED_ERRCODE, chat,
@@ -3367,6 +3139,20 @@ class WeComAdapter(BasePlatformAdapter):
self._stream_expired_chats.add(chat)
return False
except Exception as exc:
# Same intermediate/final split as the expired path above: a single
# intermediate frame failing is transient and self-healing (the next
# cumulative frame overwrites it), so swallow it and keep the turn
# (and its keep-alive) alive. A final-frame failure genuinely leaves
# the screen short of the answer, so retire the turn and let the
# consumer's send() fallback deliver it.
if not finalize:
logger.info(
"[%s] Intermediate stream frame failed (chat=%s): %s — "
"dropping frame, stream stays live",
self.name, chat, exc,
)
return True
logger.warning(
"[%s] Stream frame failed (chat=%s): %s",
self.name, chat, exc,
@@ -1075,3 +1075,61 @@ class TestClarifyEagerReseed:
consumer.finish()
await task
class TestNativeCommentaryPreservesAccumulated:
"""Regression lock for root cause #1 of the "Cla"/"ude" split-bubble bug.
In native streaming mode a mid-turn commentary (e.g. a Hindsight
"recalled N memories" notice) must be delivered as its own message via
``send()`` WITHOUT calling ``_reset_segment_state()``. The native stream
is cumulative: resetting ``_accumulated`` mid-turn dropped all
pre-commentary text, so the following delta (and the finalize frame)
carried only the few characters accumulated *after* the reset — producing
a mini-bubble like ``Cla`` and a body missing its first characters.
Drives the exact ``on_delta -> on_commentary -> on_delta`` sequence on a
real BasePlatformAdapter subclass so the
``isinstance(BasePlatformAdapter)`` + ``_use_native_streaming`` gate is
satisfied and the new ``elif self._use_native_streaming`` branch runs.
"""
async def _wait_until(self, predicate, timeout: float = 1.0) -> bool:
deadline = asyncio.get_event_loop().time() + timeout
while asyncio.get_event_loop().time() < deadline:
if predicate():
return True
await asyncio.sleep(0.01)
return predicate()
@pytest.mark.asyncio
async def test_commentary_does_not_reset_accumulated_in_native(self):
adapter = _make_native_streaming_adapter()
consumer = GatewayStreamConsumer(
adapter,
"chat_cla",
StreamConsumerConfig(edit_interval=0.01, buffer_threshold=5),
)
task = asyncio.create_task(consumer.run())
consumer.on_delta("Claude Code ")
await self._wait_until(lambda: consumer._native_stream_opened)
# Mid-turn commentary (Hindsight recall) — MUST NOT reset _accumulated.
consumer.on_commentary("🔮 recalled 10 memories")
await self._wait_until(lambda: adapter.send.await_count >= 1)
consumer.on_delta("finished (exit 0).")
consumer.finish()
await task
# 1) Commentary went out as its own proactive message (not a frame).
assert adapter.send.await_count == 1
# 2) The finalize frame carries the FULL cumulative body — no
# first-character loss ("Claude Code finished", never "ude ...").
finalize_frames = [f for f in adapter.frames if f["finalize"]]
assert finalize_frames, "expected a finalize frame"
final_text = finalize_frames[-1]["text"]
assert final_text.startswith("Claude Code "), (
f"pre-commentary prefix lost (the bug): {final_text!r}"
)
assert "finished (exit 0)." in final_text
+58 -164
View File
@@ -957,15 +957,14 @@ class TestSendStreamFrame:
@pytest.mark.asyncio
async def test_first_call_seeds_thinking_frame_then_returns_true(self):
"""First frame for a chat sends <think></think> seed.
"""First frame for a chat sends <think></think> seed, then the
content frame.
Bypasses ack tracking so both seed and content are sent in same call.
Uses a payload above ``BLOCK_STREAM_MIN_CHARS`` ending on a sentence
boundary so the block chunker emits immediately (otherwise the
chunker buffers waiting for either ``min_chars`` + a safe break or
the 250ms idle flush).
Fire-and-forget: intermediate frames are pushed immediately (pure
identity-dedup), so any non-empty payload produces a content frame
right after the seed — no min_chars / sentence-boundary gating.
"""
from plugins.platforms.wecom.adapter import WeComAdapter, BLOCK_STREAM_MIN_CHARS
from plugins.platforms.wecom.adapter import WeComAdapter
adapter = WeComAdapter(PlatformConfig(enabled=True))
adapter._last_chat_req_ids["chat-1"] = "req-1"
@@ -973,9 +972,7 @@ class TestSendStreamFrame:
# Mock _send_reply_queued to bypass ack tracking
self._mock_send_json_with_immediate_ack(adapter)
# Build a sentence-terminated payload above min_chars.
payload = ("hello world. " * 12).strip()
assert len(payload) >= BLOCK_STREAM_MIN_CHARS
payload = "hello world"
ok = await adapter.send_stream_frame(payload, chat_id="chat-1")
assert ok is True
@@ -994,11 +991,11 @@ class TestSendStreamFrame:
async def test_first_and_second_call_share_stream_id(self):
"""Successive frames use the same stream_id.
Both frames carry sentence-terminated payloads above min_chars so the
block chunker emits each one — the test is about stream_id continuity
across frames, not about chunker thresholds.
Fire-and-forget pushes each distinct cumulative payload immediately,
so this exercises stream_id continuity across frames, not chunker
thresholds.
"""
from plugins.platforms.wecom.adapter import WeComAdapter, BLOCK_STREAM_MIN_CHARS
from plugins.platforms.wecom.adapter import WeComAdapter
adapter = WeComAdapter(PlatformConfig(enabled=True))
adapter._last_chat_req_ids["chat-1"] = "req-1"
@@ -1006,10 +1003,8 @@ class TestSendStreamFrame:
# Immediate ack so all frames are sent (no pending-skip)
self._mock_send_json_with_immediate_ack(adapter)
first = ("alpha sentence. " * 10).strip()
second = first + " " + ("beta sentence. " * 10).strip()
assert len(first) >= BLOCK_STREAM_MIN_CHARS
assert len(second) - len(first) >= BLOCK_STREAM_MIN_CHARS
first = "alpha"
second = "alpha beta" # cumulative growth — differs from `first`
await adapter.send_stream_frame(first, chat_id="chat-1")
await adapter.send_stream_frame(second, chat_id="chat-1")
@@ -1212,17 +1207,20 @@ class TestSendStreamFrameFailures:
assert "chat-1" not in adapter._stream_expired_chats
@pytest.mark.asyncio
async def test_generic_transport_error_resets_state(self):
from plugins.platforms.wecom.adapter import WeComAdapter
async def test_generic_transport_error_on_intermediate_is_fire_and_forget(self):
"""A generic transport error on an INTERMEDIATE frame is fire-and-forget.
adapter = WeComAdapter(PlatformConfig(enabled=True))
adapter._last_chat_req_ids["chat-1"] = "req-1"
adapter._send_reply_request = AsyncMock(
side_effect=RuntimeError("ws disconnected"),
)
The seed frame here fails with a generic RuntimeError. An intermediate
frame failing is transient and self-healing — a later cumulative frame
(or the finalize frame) re-carries the full text — so the turn must stay
live (keep-alive keeps refreshing it) and the call returns True. It
must NOT retire the turn or trip the consumer's send() fallback, which
would re-deliver content the stream will overwrite (duplicate bubble).
@pytest.mark.asyncio
async def test_generic_transport_error_resets_state(self):
Contrast with the finalize-frame failure paths, which still return False
and retire so the consumer can fall back (see the double_send /
stream_dup_fix suites).
"""
from plugins.platforms.wecom.adapter import WeComAdapter
adapter = WeComAdapter(PlatformConfig(enabled=True))
@@ -1234,10 +1232,10 @@ class TestSendStreamFrameFailures:
turn_id = "test-turn-3"
ok = await adapter.send_stream_frame("hi", chat_id="chat-1", turn_id=turn_id)
assert ok is False
# Generic error cleans up this specific turn but does NOT mark chat as expired
assert ok is True
# Intermediate failure keeps the turn alive and leaves the chat usable.
turn_key = "chat-1:test-turn-3"
assert turn_key not in adapter._stream_turns
assert turn_key in adapter._stream_turns
assert "chat-1" not in adapter._stream_expired_chats
# === SEND_TYPING TESTS PLACEHOLDER ===
@@ -1383,100 +1381,11 @@ class TestSendClosesActiveStream:
class TestBlockChunker:
"""Unit tests for the sentence-aligned block chunker.
The chunker buffers cumulative text and only emits when a sentence
boundary lands above ``min_chars``, or when a hard ``max_chars`` cap is
reached, or when the caller forces a drain (finalize / idle flush).
"""
def test_emits_nothing_below_min_chars(self):
from plugins.platforms.wecom.adapter import _BlockChunker
chunker = _BlockChunker(min_chars=120, max_chars=360)
chunker.update("Hello world.")
assert chunker.drain(force=False) is None
# State is preserved across drain attempts.
assert chunker.has_pending() is True
def test_emits_on_sentence_boundary_above_min_chars(self):
from plugins.platforms.wecom.adapter import _BlockChunker
chunker = _BlockChunker(min_chars=20, max_chars=360)
chunker.update("This is sentence one. This is two.")
out = chunker.drain(force=False)
assert out is not None
# Cumulative — full content up to the break point.
assert out.endswith(".")
# Nothing pending after a clean emit.
assert chunker.has_pending() is False
def test_emits_on_paragraph_boundary(self):
from plugins.platforms.wecom.adapter import _BlockChunker
chunker = _BlockChunker(min_chars=20, max_chars=360)
text = "Part one of the answer text here\n\nPart two of the answer"
chunker.update(text)
out = chunker.drain(force=False)
assert out is not None
assert out.endswith("\n\n")
assert chunker.has_pending() is True # part two still buffered
def test_hard_cap_forces_break_mid_sentence(self):
from plugins.platforms.wecom.adapter import _BlockChunker
chunker = _BlockChunker(min_chars=20, max_chars=50)
# Single run-on sentence with no terminators in the first 50 chars.
text = "a" * 80
chunker.update(text)
out = chunker.drain(force=False)
assert out is not None
# Hard cap fires — entire cumulative is emitted (the chunker doesn't
# split mid-buffer; it leaves the next slice to the next update).
assert out == text
def test_force_emits_remaining_buffer(self):
from plugins.platforms.wecom.adapter import _BlockChunker
chunker = _BlockChunker(min_chars=120, max_chars=360)
chunker.update("Tiny tail")
out = chunker.drain(force=True)
assert out == "Tiny tail"
assert chunker.has_pending() is False
def test_no_pending_after_complete_drain(self):
from plugins.platforms.wecom.adapter import _BlockChunker
chunker = _BlockChunker(min_chars=10, max_chars=360)
chunker.update("hello there world.")
assert chunker.drain(force=False) == "hello there world."
assert chunker.drain(force=True) is None # nothing to emit
def test_does_not_break_on_decimal_or_ip_address(self):
"""Sentence break requires whitespace after terminator — '1.2' must not split."""
from plugins.platforms.wecom.adapter import _BlockChunker
chunker = _BlockChunker(min_chars=10, max_chars=360)
chunker.update("version 1.2.3 is out and not yet released")
# No safe sentence terminator → wait.
assert chunker.drain(force=False) is None
def test_cumulative_update_does_not_shrink(self):
from plugins.platforms.wecom.adapter import _BlockChunker
chunker = _BlockChunker()
chunker.update("long initial cumulative content " * 5)
before = chunker._cumulative
chunker.update("short") # caller bug — must not shrink high-water mark
assert chunker._cumulative == before
class TestBlockStreamingFrameFlow:
"""Integration: send_stream_frame uses the chunker to gate intermediate frames."""
class TestFireAndForgetFrameFlow:
"""Integration: send_stream_frame pushes each distinct cumulative payload
immediately (pure identity-dedup), with no sentence/min-chars buffering."""
def _mock_send_json_with_immediate_ack(self, adapter):
from plugins.platforms.wecom.adapter import APP_CMD_RESPONSE
sent_frames = []
async def mock_send(reply_req_id, body, **kwargs):
@@ -1492,8 +1401,9 @@ class TestBlockStreamingFrameFlow:
adapter._sent_frames = sent_frames
@pytest.mark.asyncio
async def test_short_text_buffered_no_frame_sent(self):
"""A few tokens below min_chars do not produce an intermediate frame."""
async def test_short_text_sent_immediately(self):
"""Fire-and-forget: even a short body ships right after the seed —
there is no min_chars buffering anymore."""
from plugins.platforms.wecom.adapter import WeComAdapter
adapter = WeComAdapter(PlatformConfig(enabled=True))
@@ -1503,13 +1413,18 @@ class TestBlockStreamingFrameFlow:
ok = await adapter.send_stream_frame("Hello.", chat_id="chat-1")
assert ok is True
# Only the seed frame was sent — the 6-char body is below min_chars.
assert len(adapter._sent_frames) == 1
# seed + content = 2 frames (the 6-char body is NOT buffered).
assert len(adapter._sent_frames) == 2
assert adapter._sent_frames[0]["body"]["stream"]["content"] == "<think></think>"
assert adapter._sent_frames[1]["body"]["stream"]["content"] == "Hello."
assert adapter._sent_frames[1]["body"]["stream"]["finish"] is False
@pytest.mark.asyncio
async def test_finalize_force_drains_buffered_tail(self):
"""Finalize emits whatever the chunker has, even sub-min_chars."""
async def test_finalize_sends_accumulated_tail(self):
"""Finalize emits the accumulated text with finish=true.
With no chunker, finalize uses the caller's cumulative text directly.
"""
from plugins.platforms.wecom.adapter import WeComAdapter
adapter = WeComAdapter(PlatformConfig(enabled=True))
@@ -1517,42 +1432,26 @@ class TestBlockStreamingFrameFlow:
adapter._ws = MagicMock(closed=False)
self._mock_send_json_with_immediate_ack(adapter)
# 1: buffer small text (no intermediate frame sent — only seed).
# 1: intermediate frame (seed + content).
await adapter.send_stream_frame("Short.", chat_id="chat-1")
# 2: finalize with the same tiny tail — force-drains the chunker.
# 2: finalize with the same text. Content equals last_sent_content, so
# the adapter appends a zero-width space to force a distinct final frame.
ok = await adapter.send_stream_frame(
"Short.", chat_id="chat-1", finalize=True,
)
assert ok is True
# seed + finalize = 2 frames.
assert len(adapter._sent_frames) == 2
# seed + content + finalize = 3 frames.
assert len(adapter._sent_frames) == 3
final_frame = adapter._sent_frames[-1]
assert final_frame["body"]["stream"]["finish"] is True
# Content survives the force-drain.
# Content survives the finalize (zero-width space appended when it
# matched the previous frame verbatim).
assert final_frame["body"]["stream"]["content"].startswith("Short.")
@pytest.mark.asyncio
async def test_idle_flush_emits_partial_buffer(self):
"""A buffered sub-min tail still ships after the idle-flush timer fires."""
from plugins.platforms.wecom.adapter import WeComAdapter, BLOCK_STREAM_IDLE_FLUSH
adapter = WeComAdapter(PlatformConfig(enabled=True))
adapter._last_chat_req_ids["chat-1"] = "req-1"
adapter._ws = MagicMock(closed=False)
self._mock_send_json_with_immediate_ack(adapter)
await adapter.send_stream_frame("Half a thought", chat_id="chat-1")
# Wait past the idle window so the timer fires and drains.
await asyncio.sleep(BLOCK_STREAM_IDLE_FLUSH + 0.1)
# seed + idle-flushed partial = 2 frames.
assert len(adapter._sent_frames) == 2
flushed = adapter._sent_frames[-1]
assert flushed["body"]["stream"]["content"] == "Half a thought"
assert flushed["body"]["stream"]["finish"] is False
@pytest.mark.asyncio
async def test_block_emits_on_sentence_boundary(self):
"""Cumulative text crossing a sentence boundary triggers an immediate frame."""
async def test_duplicate_intermediate_content_is_deduped(self):
"""Identical cumulative content skips the send (pure identity-dedup),
but still returns success."""
from plugins.platforms.wecom.adapter import WeComAdapter
adapter = WeComAdapter(PlatformConfig(enabled=True))
@@ -1560,17 +1459,12 @@ class TestBlockStreamingFrameFlow:
adapter._ws = MagicMock(closed=False)
self._mock_send_json_with_immediate_ack(adapter)
# Build cumulative text that ends on a sentence boundary above min_chars.
text = "This is a complete first sentence with enough text. " * 4
await adapter.send_stream_frame(text, chat_id="chat-1")
# seed + sentence block = 2 frames.
await adapter.send_stream_frame("same text", chat_id="chat-1")
ok = await adapter.send_stream_frame("same text", chat_id="chat-1")
assert ok is True
# seed + first content only; the identical repeat was deduped.
assert len(adapter._sent_frames) == 2
block = adapter._sent_frames[-1]
# The chunker breaks immediately after the rightmost safe sentence
# terminator — i.e., right after the last ".", stripping the trailing
# space that the caller's cumulative blob included.
assert block["body"]["stream"]["content"] == text.rstrip()
assert block["body"]["stream"]["finish"] is False
assert adapter._sent_frames[-1]["body"]["stream"]["content"] == "same text"
class TestFinalFrameAckTimeoutSemantics:
+330
View File
@@ -0,0 +1,330 @@
"""Behavior-contract tests for the WeCom native-streaming "重复气泡" fix.
Two production bug classes are covered, each toggle-validated so the test
reproduces the duplicate/mis-decline when the fix is disabled and passes when
it is enabled — proving the assertions track behavior, not a frozen snapshot.
Fix A — Layer 2 clock fallback must NOT decline a finalize frame while Layer 1
keep-alive is enabled. Keep-alive refreshes the stream window every ~2min, so
a large ``stream_age`` does not mean the stream is dead; declining it on a blind
clock read forces the consumer's send() fallback to re-deliver content the
intermediate frames already put on screen (the duplicate bubble). Toggle =
``adapter._stream_keepalive_enabled``.
Fix B — an intermediate frame (finalize=False) failing/expiring must be
fire-and-forget (return True, turn stays live, no consumer fallback), because a
later cumulative frame overwrites it. Only a FINAL frame (finalize=True)
failure means the screen is genuinely missing content and must return False to
trip the consumer's send() fallback. Toggle = the ``finalize`` argument.
These drive the REAL ``WeComAdapter._send_stream_frame_inner`` with only the
byte-level ``_send_stream_reply`` seam faked, so the actual finalize / except
control flow runs. Assertions read observable adapter state: the return value
(what the consumer keys its fallback on), whether the turn survived in
``_stream_turns``, whether the chat was marked expired, and how many finalize
frames actually reached ``_send_stream_reply``.
"""
from __future__ import annotations
import asyncio
from unittest.mock import AsyncMock
import pytest
from gateway.config import PlatformConfig
from plugins.platforms.wecom.adapter import (
WeComAdapter,
WeComStreamExpiredError,
STREAM_EXPIRED_ERRCODE,
)
CHAT_ID = "chat-dup"
REQ_ID = "req-dup"
TURN_ID = "turn-dup"
# Fire-and-forget intermediate frames are pushed as soon as the cumulative
# text differs from the last sent frame (pure identity-dedup, no chunker /
# min-chars gate). A non-trivial body is used so the intermediate-failure
# tests exercise a real content frame rather than only the seed frame.
BLOCK_TEXT = (
"This is a complete sentence used to fill the block chunker past its "
"minimum character threshold so it actually drains a content frame. "
"Here is a second sentence to be safely over the limit."
)
def _make_adapter(*, keepalive_enabled: bool) -> WeComAdapter:
"""Real adapter with only the stream byte-writer faked.
``_send_stream_reply`` is the seam between per-turn logic and the wire.
Faking it here lets each test dictate per-frame success/expiry while the
real finalize / except branches run.
"""
extra = {"stream_keepalive_enabled": keepalive_enabled}
adapter = WeComAdapter(PlatformConfig(enabled=True, extra=extra))
adapter._last_chat_req_ids[CHAT_ID] = REQ_ID
return adapter
def _finalize_calls(mock: AsyncMock) -> list:
"""finish=True calls that reached ``_send_stream_reply``."""
return [c for c in mock.await_args_list if c.kwargs.get("finish") is True]
# ===========================================================================
# Fix A — keep-alive suppresses the Layer 2 clock decline
# ===========================================================================
class TestKeepaliveSuppressesClockDecline:
"""A long-lived, keep-alive-refreshed stream must still finalize natively;
the clock fallback must not decline it and force a duplicate send()."""
@pytest.mark.asyncio
async def test_keepalive_on_old_stream_finalizes_natively(self):
"""FIX ENABLED: keep-alive on + stream_age >> safe_duration + finalize.
Post-fix contract: the finalize frame is sent on the wire (finish=True),
finalize returns True, the turn is finalized+cleaned, and the chat is
NOT marked expired — so the consumer suppresses its send() fallback and
no duplicate bubble is produced.
"""
adapter = _make_adapter(keepalive_enabled=True)
try:
reply = AsyncMock(return_value={"errcode": 0})
adapter._send_stream_reply = reply
# Open the turn (seed + first intermediate).
await adapter._send_stream_frame_inner(
"partial answer", chat=CHAT_ID, finalize=False, turn_id=TURN_ID,
)
turn = adapter._stream_turns[f"{CHAT_ID}:{TURN_ID}"]
# Age the stream far beyond the Layer 2 safe duration.
turn.start_time -= adapter._stream_safe_duration_seconds + 500
ok = await adapter._send_stream_frame_inner(
"the complete final answer",
chat=CHAT_ID, finalize=True, turn_id=TURN_ID,
)
# Native finalize succeeded — consumer will suppress fallback.
assert ok is True
assert len(_finalize_calls(reply)) == 1, (
"keep-alive on: finalize frame must reach the wire, not be "
"declined by the Layer 2 clock fallback"
)
assert CHAT_ID not in adapter._stream_expired_chats
assert f"{CHAT_ID}:{TURN_ID}" not in adapter._stream_turns
finally:
await adapter.disconnect()
@pytest.mark.asyncio
async def test_keepalive_off_old_stream_declines_finalize(self):
"""FIX DISABLED (toggle): keep-alive off preserves original Layer 2.
With keep-alive off there is nothing refreshing the window, so a
stream older than safe_duration is genuinely doomed — the original
clock decline must still fire: no finalize frame on the wire, return
False, turn retired, chat marked expired (consumer takes over via
send()). This proves Fix A is gated on the toggle, not unconditional.
"""
adapter = _make_adapter(keepalive_enabled=False)
try:
reply = AsyncMock(return_value={"errcode": 0})
adapter._send_stream_reply = reply
await adapter._send_stream_frame_inner(
"partial answer", chat=CHAT_ID, finalize=False, turn_id=TURN_ID,
)
turn = adapter._stream_turns[f"{CHAT_ID}:{TURN_ID}"]
turn.start_time -= adapter._stream_safe_duration_seconds + 500
ok = await adapter._send_stream_frame_inner(
"the complete final answer",
chat=CHAT_ID, finalize=True, turn_id=TURN_ID,
)
assert ok is False, "keep-alive off: old stream must decline finalize"
assert len(_finalize_calls(reply)) == 0, (
"keep-alive off: no finalize frame should reach the wire"
)
assert CHAT_ID in adapter._stream_expired_chats
assert f"{CHAT_ID}:{TURN_ID}" not in adapter._stream_turns
finally:
await adapter.disconnect()
@pytest.mark.asyncio
async def test_keepalive_on_truly_expired_stream_falls_back(self):
"""Even with the clock decline skipped, a genuinely dead stream is safe.
Keep-alive on skips the blind clock check, but if the stream really has
expired the finalize ``_send_stream_reply(finish=True)`` hits 846608 and
raises ``WeComStreamExpiredError`` — the existing except path then
retires the turn and returns False for the (correct, finalize-only)
fallback. This shows Fix A does not lose the real-expiry safety net.
"""
adapter = _make_adapter(keepalive_enabled=True)
try:
async def _reply(req_id, stream_id, content, finish=False):
if finish:
raise WeComStreamExpiredError(errcode=STREAM_EXPIRED_ERRCODE)
return {"errcode": 0}
adapter._send_stream_reply = AsyncMock(side_effect=_reply)
await adapter._send_stream_frame_inner(
"partial answer", chat=CHAT_ID, finalize=False, turn_id=TURN_ID,
)
turn = adapter._stream_turns[f"{CHAT_ID}:{TURN_ID}"]
turn.start_time -= adapter._stream_safe_duration_seconds + 500
ok = await adapter._send_stream_frame_inner(
"the complete final answer",
chat=CHAT_ID, finalize=True, turn_id=TURN_ID,
)
assert ok is False, "real 846608 on finalize must fall back"
assert CHAT_ID in adapter._stream_expired_chats
assert f"{CHAT_ID}:{TURN_ID}" not in adapter._stream_turns
finally:
await adapter.disconnect()
# ===========================================================================
# Fix B — intermediate failures are fire-and-forget; only final falls back
# ===========================================================================
class TestIntermediateFrameFailureIsFireAndForget:
"""A single intermediate frame failing must not trip the consumer fallback
or kill the turn; only a final-frame failure does."""
@pytest.mark.asyncio
async def test_intermediate_expired_returns_true_keeps_turn(self):
"""FIX ENABLED: intermediate finish=False hits 846608.
Contract: return True (no fallback), turn survives in the registry
(keep-alive keeps refreshing), chat NOT marked expired. A later
cumulative frame will overwrite the dropped one.
"""
adapter = _make_adapter(keepalive_enabled=True)
try:
# Seed succeeds; the next intermediate content frame expires.
calls = {"n": 0}
async def _reply(req_id, stream_id, content, finish=False):
calls["n"] += 1
# First call is the seed ("<think></think>"); let it succeed so
# the turn opens, then expire the real content frame.
if calls["n"] >= 2 and not finish:
raise WeComStreamExpiredError(errcode=STREAM_EXPIRED_ERRCODE)
return {"errcode": 0}
adapter._send_stream_reply = AsyncMock(side_effect=_reply)
# Send a non-trivial body so a real content frame is drained
# (fire-and-forget: any content differing from the last frame).
ok = await adapter._send_stream_frame_inner(
BLOCK_TEXT,
chat=CHAT_ID, finalize=False, turn_id=TURN_ID,
)
assert ok is True, "intermediate expiry must be fire-and-forget"
assert f"{CHAT_ID}:{TURN_ID}" in adapter._stream_turns, (
"intermediate failure must NOT retire the turn — keep-alive is "
"still refreshing the live stream"
)
turn = adapter._stream_turns[f"{CHAT_ID}:{TURN_ID}"]
assert turn.expired is False
assert CHAT_ID not in adapter._stream_expired_chats
finally:
await adapter.disconnect()
@pytest.mark.asyncio
async def test_intermediate_generic_exception_returns_true_keeps_turn(self):
"""Same fire-and-forget contract for a generic (non-expiry) exception on
an intermediate frame — the whole except class is fixed, not just the
WeComStreamExpiredError path."""
adapter = _make_adapter(keepalive_enabled=True)
try:
calls = {"n": 0}
async def _reply(req_id, stream_id, content, finish=False):
calls["n"] += 1
if calls["n"] >= 2 and not finish:
raise RuntimeError("transient wire error")
return {"errcode": 0}
adapter._send_stream_reply = AsyncMock(side_effect=_reply)
ok = await adapter._send_stream_frame_inner(
BLOCK_TEXT,
chat=CHAT_ID, finalize=False, turn_id=TURN_ID,
)
assert ok is True
assert f"{CHAT_ID}:{TURN_ID}" in adapter._stream_turns
assert CHAT_ID not in adapter._stream_expired_chats
finally:
await adapter.disconnect()
@pytest.mark.asyncio
async def test_final_expired_returns_false_falls_back(self):
"""FIX-INVARIANT (toggle via finalize flag): a final finish=True frame
that expires MUST return False, retire the turn, and mark the chat
expired — the screen is genuinely missing this content, so the consumer
must run its send() fallback. Contrast with the intermediate case above
proves the fix discriminates on finalize, not blanket-swallows."""
adapter = _make_adapter(keepalive_enabled=True)
try:
async def _reply(req_id, stream_id, content, finish=False):
if finish:
raise WeComStreamExpiredError(errcode=STREAM_EXPIRED_ERRCODE)
return {"errcode": 0}
adapter._send_stream_reply = AsyncMock(side_effect=_reply)
await adapter._send_stream_frame_inner(
"partial answer", chat=CHAT_ID, finalize=False, turn_id=TURN_ID,
)
ok = await adapter._send_stream_frame_inner(
"the complete final answer",
chat=CHAT_ID, finalize=True, turn_id=TURN_ID,
)
assert ok is False, "final-frame expiry MUST fall back (return False)"
assert CHAT_ID in adapter._stream_expired_chats
assert f"{CHAT_ID}:{TURN_ID}" not in adapter._stream_turns
finally:
await adapter.disconnect()
@pytest.mark.asyncio
async def test_final_generic_exception_returns_false_retires(self):
"""A generic exception on the final frame also returns False + retires
the turn (consumer fallback)."""
adapter = _make_adapter(keepalive_enabled=True)
try:
async def _reply(req_id, stream_id, content, finish=False):
if finish:
raise RuntimeError("wire down on finalize")
return {"errcode": 0}
adapter._send_stream_reply = AsyncMock(side_effect=_reply)
await adapter._send_stream_frame_inner(
"partial answer", chat=CHAT_ID, finalize=False, turn_id=TURN_ID,
)
ok = await adapter._send_stream_frame_inner(
"the complete final answer",
chat=CHAT_ID, finalize=True, turn_id=TURN_ID,
)
assert ok is False
assert f"{CHAT_ID}:{TURN_ID}" not in adapter._stream_turns
finally:
await adapter.disconnect()