refactor(gateway): stream consumer — overflow predicates, native re-seed early-exit, cursor strip gate

This commit is contained in:
Teknium
2026-09-02 20:09:53 -07:00
parent 917c897e81
commit 2c5dc2ceee
3 changed files with 21 additions and 26 deletions
+8 -10
View File
@@ -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:
+1 -3
View File
@@ -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)
+12 -13
View File
@@ -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.