diff --git a/gateway/stream_consumer.py b/gateway/stream_consumer.py index 8811db6436..e021b7919c 100644 --- a/gateway/stream_consumer.py +++ b/gateway/stream_consumer.py @@ -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) diff --git a/plugins/platforms/wecom/adapter.py b/plugins/platforms/wecom/adapter.py index 7757843dc7..929178df53 100644 --- a/plugins/platforms/wecom/adapter.py +++ b/plugins/platforms/wecom/adapter.py @@ -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, diff --git a/tests/gateway/test_stream_consumer_wecom_native.py b/tests/gateway/test_stream_consumer_wecom_native.py index cd691eeb25..875ab3b149 100644 --- a/tests/gateway/test_stream_consumer_wecom_native.py +++ b/tests/gateway/test_stream_consumer_wecom_native.py @@ -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 diff --git a/tests/gateway/test_wecom.py b/tests/gateway/test_wecom.py index abfe5d5a35..a46a1caded 100644 --- a/tests/gateway/test_wecom.py +++ b/tests/gateway/test_wecom.py @@ -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 seed. + """First frame for a chat sends 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"] == "" + 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: diff --git a/tests/gateway/test_wecom_stream_dup_fix.py b/tests/gateway/test_wecom_stream_dup_fix.py new file mode 100644 index 0000000000..7329b2a145 --- /dev/null +++ b/tests/gateway/test_wecom_stream_dup_fix.py @@ -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 (""); 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()