From 2c5dc2ceeec1fab3660d2c123449581452130cd8 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 20:09:53 -0700 Subject: [PATCH] =?UTF-8?q?refactor(gateway):=20stream=20consumer=20?= =?UTF-8?q?=E2=80=94=20overflow=20predicates,=20native=20re-seed=20early-e?= =?UTF-8?q?xit,=20cursor=20strip=20gate?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- gateway/stream_consumer.py | 18 ++++++++---------- gateway/stream_consumer_fallback.py | 4 +--- gateway/stream_consumer_transport.py | 25 ++++++++++++------------- 3 files changed, 21 insertions(+), 26 deletions(-) diff --git a/gateway/stream_consumer.py b/gateway/stream_consumer.py index cd18ab04e7..9b6a0fca97 100644 --- a/gateway/stream_consumer.py +++ b/gateway/stream_consumer.py @@ -620,11 +620,7 @@ class GatewayStreamConsumer( ): # Overflow split. Native streaming bypasses this: the adapter # truncates against the stream protocol's own limit. - if ( - not self._use_native_streaming - and self._len_fn(self._accumulated) > self._safe_limit - and self._message_id is None - ): + if not self._use_native_streaming and self._first_send_overflows(): if await self._split_first_send(tick) == "return": return continue @@ -866,13 +862,15 @@ class GatewayStreamConsumer( self._signal_flush(tick.flush_event) return "continue" + def _overflows(self) -> bool: + return self._len_fn(self._accumulated) > self._safe_limit + + def _first_send_overflows(self) -> bool: + return self._message_id is None and self._overflows() + async def _seal_overflow_heads(self) -> None: """Existing message overflowing: seal it with the head, start a new message for the rest.""" - while ( - self._len_fn(self._accumulated) > self._safe_limit - and self._message_id is not None - and self._edit_supported - ): + while self._overflows() and self._message_id is not None and self._edit_supported: cp_budget = _custom_unit_to_cp(self._accumulated, self._safe_limit, self._len_fn) split_at = self._accumulated.rfind("\n", 0, cp_budget) if split_at < cp_budget // 2: diff --git a/gateway/stream_consumer_fallback.py b/gateway/stream_consumer_fallback.py index 2236d4cd72..a77971aaae 100644 --- a/gateway/stream_consumer_fallback.py +++ b/gateway/stream_consumer_fallback.py @@ -339,10 +339,8 @@ class StreamFallbackMixin: async def _try_strip_cursor(self) -> None: """Best-effort edit removing a stuck cursor when entering fallback mode.""" - if not self._message_id or self._message_id == "__no_edit__": - return prefix = self._visible_prefix() - if not prefix.strip(): + if not self._has_real_preview() or not prefix.strip(): return with contextlib.suppress(Exception): # never block the fallback path result = await self._edit_message(message_id=self._message_id, content=prefix) diff --git a/gateway/stream_consumer_transport.py b/gateway/stream_consumer_transport.py index 121a06aeb2..49d8cf6de8 100644 --- a/gateway/stream_consumer_transport.py +++ b/gateway/stream_consumer_transport.py @@ -388,22 +388,21 @@ class StreamTransportMixin: """ if not self._native_stream_opened and text: try: - if await self._send_seed_frame(): - self._native_stream_opened = True - self._awaiting_reopen_after_boundary = False - # Paired with the boundary-finalize INFO: typing-reappear latency. - logger.info( - "[latency] Re-opened native stream after boundary " - "(turn=%s, waited for first delta)", - self._turn_id, - ) - else: - self._use_native_streaming = False + seeded = await self._send_seed_frame() except Exception as e: logger.debug("Re-seed failed, disabling native streaming: %s", e) + seeded = False + if not seeded: self._use_native_streaming = False - if not self._use_native_streaming: - return None + return None + self._native_stream_opened = True + self._awaiting_reopen_after_boundary = False + # Paired with the boundary-finalize INFO: typing-reappear latency. + logger.info( + "[latency] Re-opened native stream after boundary " + "(turn=%s, waited for first delta)", + self._turn_id, + ) # WeCom renders each finalize as a separate bubble: only the turn-final and # boundaries close the stream, not segment breaks.