From ff26aa924797738c4cbde3dacce41239a78bc490 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 20:24:52 -0700 Subject: [PATCH] =?UTF-8?q?refactor(gateway):=20stream=20transport/fallbac?= =?UTF-8?q?k/events=20=E2=80=94=20docstrings=20to=20summary=20+=20invarian?= =?UTF-8?q?t?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- gateway/stream_consumer_fallback.py | 57 +++++++----------- gateway/stream_consumer_transport.py | 87 +++++++++------------------- gateway/stream_dispatch.py | 26 ++++----- gateway/stream_events.py | 38 ++++-------- 4 files changed, 69 insertions(+), 139 deletions(-) diff --git a/gateway/stream_consumer_fallback.py b/gateway/stream_consumer_fallback.py index a77971aaae..a6068a8475 100644 --- a/gateway/stream_consumer_fallback.py +++ b/gateway/stream_consumer_fallback.py @@ -76,10 +76,8 @@ class StreamFallbackMixin: def _truncate_for_stream( self, text: str, limit: int, len_fn: "Callable[[str], int]", ) -> list[str]: - """Split via the adapter's canonical truncate_message (platform-specific rules). - - Non-base test doubles / legacy adapters keep the two-argument call shape. - """ + """Split via the adapter's canonical truncate_message (platform-specific rules); + non-base test doubles / legacy adapters keep the two-argument call shape.""" truncate = getattr(self.adapter, "truncate_message", None) if not callable(truncate): return self._split_text_chunks(text, limit, len_fn) @@ -92,10 +90,8 @@ class StreamFallbackMixin: return list(chunks) async def _send_fallback_final(self, text: str) -> None: - """Send the final continuation after streaming edits stop working. - - Retries each chunk once on flood-control failures with a short delay. - """ + """Send the final continuation after streaming edits stop working (one flood retry + per chunk).""" # Balance fences BEFORE computing the continuation so the closing fence # reaches the user even when only the tail is delivered. final_text = ensure_closed_code_fences(self._clean_for_display(text)) @@ -154,11 +150,9 @@ class StreamFallbackMixin: self._fallback_preserve_partial_messages = False async def _fallback_when_nothing_unseen(self, final_text: str) -> Optional[str]: - """Fallback entered but the visible prefix already covers ``final_text``. - - Returns the continuation to send (the whole final when the prefix is from a - *previous* segment) or None when the turn is settled here. - """ + """Fallback entered but the visible prefix already covers ``final_text``: returns the + continuation to send (the whole final when the prefix is from a *previous* segment) + or None when the turn is settled here.""" visible = self._visible_prefix() # Telegram clients can lose (part of) a streamed preview after a failed # final edit, so opt-in adapters commit a fresh final send. @@ -219,11 +213,8 @@ class StreamFallbackMixin: return _len_fn, raw_limit async def _send_with_flood_retry(self, *, content: str, retry_log: str, reply_to=None): - """adapter.send(final metadata) with ONE bounded retry on flood control. - - Exceptions propagate (callers decide whether a raise means "ambiguous"). - Returns the last SendResult (success or not). - """ + """adapter.send(final metadata) with ONE bounded flood retry; returns the last + SendResult. Exceptions propagate (callers decide whether a raise is "ambiguous").""" kwargs = dict( chat_id=self.chat_id, content=content, metadata=self._metadata_for_send(final=True), ) @@ -242,11 +233,9 @@ class StreamFallbackMixin: return result async def _send_empty_fallback_final(self, final_text: str) -> str: - """Commit a completed answer after Telegram finalization fails. - - Returns "delivered", "failed" (gateway may retry), "ambiguous" (a timeout may - have landed) or "preview" (flood control; the complete preview is authoritative). - """ + """Commit a completed answer after Telegram finalization fails: "delivered", "failed" + (gateway may retry), "ambiguous" (a timeout may have landed) or "preview" (flood + control; the complete preview is authoritative).""" # Segment-scoped only: never delete an earlier finalized preamble. stale_ids = self._stale_preview_ids(segment_only=True) try: @@ -314,10 +303,8 @@ class StreamFallbackMixin: return "flood" in err_lower or "retry after" in err_lower or "rate" in err_lower async def _flush_segment_tail_on_edit_failure(self) -> None: - """Send the unseen tail after the delivered prefix as a new message before a segment reset. - - Also best-effort strips the stuck cursor from the partial message. - """ + """Before a segment reset, send the unseen tail as a new message (and best-effort + strip the stuck cursor from the partial).""" if not self._fallback_final_send: await self._try_strip_cursor() visible = self._fallback_prefix or self._visible_prefix() @@ -381,11 +368,9 @@ class StreamFallbackMixin: return False def _raw_message_limit(self) -> int: - """Per-chat length budget (adapter ``message_len_fn`` units) before overflow splits. - - Rich-capable adapters may raise it via ``streaming_overflow_limit`` so a - reply that fits one rich message isn't fragmented at the edit limit. - """ + """Per-chat length budget (``message_len_fn`` units) before overflow splits; rich + adapters may raise it via ``streaming_overflow_limit`` so a reply that fits one rich + message isn't fragmented at the edit limit.""" base = getattr(self.adapter, "MAX_MESSAGE_LENGTH", 4096) # isinstance gate keeps MagicMock adapters (mock attrs, not ints) on base. if isinstance(self.adapter, _BasePlatformAdapter): @@ -403,11 +388,9 @@ class StreamFallbackMixin: return base async def _suppress_silence_marker(self) -> None: - """Retract any streamed preview when the final reply is a bare silence marker. - - Flags stay False: nothing was delivered, and the gateway's whole-response - filter turns the marker into "" so no fallback send happens either. - """ + """Retract any streamed preview when the final reply is a bare silence marker. Flags + stay False: the gateway's whole-response filter turns the marker into "" so no + fallback send happens either.""" # A native-stream bubble isn't a deletable message — close an open one # (e.g. from an eager re-seed) with an empty finalize so it doesn't hang. if self._native_stream_opened: diff --git a/gateway/stream_consumer_transport.py b/gateway/stream_consumer_transport.py index c9620d03a5..c96945e0da 100644 --- a/gateway/stream_consumer_transport.py +++ b/gateway/stream_consumer_transport.py @@ -49,10 +49,8 @@ class StreamTransportMixin: ) async def _try_seed_frame(self, fail_log: str, *, exc_info: bool = False) -> bool: - """_send_seed_frame() as a bool; a raise logs ``fail_log`` at DEBUG and reads as False. - - ``exc_info`` logs the traceback instead of formatting the error into ``fail_log``. - """ + """_send_seed_frame() as a bool; a raise logs ``fail_log`` at DEBUG (with the error + formatted in, or the traceback when ``exc_info``) and reads as False.""" try: return bool(await self._send_seed_frame()) except Exception as e: @@ -87,11 +85,8 @@ class StreamTransportMixin: self._reopen_seeded_eagerly = False def _degrade_native_to_buffered_send(self) -> None: - """Leave native mode; post-boundary output goes out as ONE send() at got_done. - - buffer_only avoids mid-stream flushes that create multiple messages on - non-editable platforms. - """ + """Leave native mode; buffer_only so post-boundary output is ONE send() at got_done + (mid-stream flushes would create multiple messages on non-editable platforms).""" self._use_native_streaming = False self._close_native_state() self.cfg.buffer_only = True @@ -105,10 +100,7 @@ class StreamTransportMixin: return md or None def _stale_preview_ids(self, *, segment_only: bool = False) -> set: - """Preview message ids a fresh final replaces. - - ``segment_only``: never delete an earlier finalized preamble. - """ + """Preview ids a fresh final replaces; ``segment_only`` spares finalized preambles.""" stale_ids = set( self._segment_preview_message_ids if segment_only else self._preview_message_ids ) @@ -136,11 +128,8 @@ class StreamTransportMixin: logger.debug("%s preview cleanup failed (%s): %s", label, stale_id, e) def _resolve_draft_streaming(self) -> bool: - """Whether this run should use draft streaming per ``cfg.transport``. - - "edit"/"off" → False. "draft"/"auto" → the adapter's supports_draft_streaming - probe (chat type, platform-version gates); "draft" logs the downgrade. - """ + """cfg.transport "draft"/"auto" → the adapter's supports_draft_streaming probe + ("draft" logs the downgrade); "edit"/"off" → False.""" transport = (self.cfg.transport or "edit").lower() if transport in ("edit", "off"): return False @@ -169,11 +158,8 @@ class StreamTransportMixin: return bool(supported) def _resolve_native_streaming(self) -> bool: - """Whether to use native streaming (adapter.send_stream_frame for ALL frames). - - Requires a BasePlatformAdapter subclass with class-level - SUPPORTS_NATIVE_STREAMING and a truthy supports_native_streaming probe. - """ + """Native streaming (send_stream_frame for ALL frames): a BasePlatformAdapter with + class-level SUPPORTS_NATIVE_STREAMING and a truthy supports_native_streaming probe.""" if not ( isinstance(self.adapter, _BasePlatformAdapter) and getattr(type(self.adapter), "SUPPORTS_NATIVE_STREAMING", False) @@ -190,9 +176,7 @@ class StreamTransportMixin: async def _send_draft_frame(self, text: str) -> bool: """Emit one draft frame; any failure permanently disables drafts for this run. - - Drafts have no message_id and clear on the client when the final sendMessage lands. - """ + Drafts have no message_id and clear on the client when the final send lands.""" if self._draft_id is None: # Should never happen (set in tandem with _use_draft_streaming in run()). self._use_draft_streaming = False @@ -219,11 +203,9 @@ class StreamTransportMixin: return False async def _abandon_native_stream(self) -> None: - """Seal an orphaned draft stream in place on turn death (stale exit / cancel). - - Else the live indicator stays forever and the adapter's armed interception - state leaks into the next turn. Never sets delivery flags. - """ + """Seal an orphaned draft stream on turn death (stale exit / cancel): else the live + indicator stays forever and armed interception state leaks into the next turn. + Never sets delivery flags.""" if not self._use_draft_streaming: return if getattr(type(self.adapter), "abandon_open_draft", None) is None: @@ -267,10 +249,8 @@ class StreamTransportMixin: self._track_preview_id(mid) def _adapter_prefers_fresh_final(self, text: str) -> bool: - """Adapter's prefers_fresh_final_streaming hook (e.g. Telegram's richer send path). - - False when there's no real preview, no hook, or on any error. - """ + """Adapter's prefers_fresh_final_streaming hook (Telegram's richer send path); + False without a real preview / hook, or on any error.""" fn = getattr(self.adapter, "prefers_fresh_final_streaming", None) if fn is None or not self._has_real_preview(): return False @@ -292,11 +272,8 @@ class StreamTransportMixin: return result is True async def _try_fresh_final(self, text: str, *, is_turn_final: bool = True) -> bool: - """Send ``text`` as a fresh message and best-effort delete the preview(s). - - False on any failure so the caller falls back to edit. ``is_turn_final=False`` - (interim segment) leaves the final-delivery flag unset. - """ + """Send ``text`` fresh and best-effort delete the preview(s); False on any failure so + the caller falls back to edit. ``is_turn_final=False`` leaves the delivery flag unset.""" # Replacing every preview is only sound while ``text`` holds the whole answer; # after a split, deleting sealed heads would erase delivered text. if self._turn_split_delivery: @@ -335,12 +312,9 @@ class StreamTransportMixin: async def _send_or_edit( self, text: str, *, finalize: bool = False, is_turn_final: bool = True, ) -> bool: - """Send or edit the streaming message; True if delivered. - - ``finalize`` marks the last edit of a streaming sequence. Transport order: - native frame → draft frame → edit existing → first send; a transport returns - None to fall through to the next. - """ + """Send or edit the streaming message; True if delivered. ``finalize`` marks the + last edit. Transport order: native frame → draft frame → edit existing → first + send; a transport returns None to fall through to the next.""" text = self._clean_for_display(text) # Stream-is-the-message draft frames must stay prefix-stable: a closing ``` # on a mid-code-block frame makes frame N not a prefix of N+1 and the @@ -395,11 +369,8 @@ class StreamTransportMixin: async def _native_push( self, text: str, *, finalize: bool, is_turn_final: bool, ) -> Optional[bool]: - """Native streaming: every frame goes through send_stream_frame(). - - Lazy re-seed here after a boundary closed the stream. Returns None when - native was disabled (seed/frame failure) so the caller falls through. - """ + """Native streaming: every frame goes through send_stream_frame(); lazy re-seed after + a boundary. None when native was disabled (seed/frame failure) → caller falls through.""" if not self._native_stream_opened and text: if not await self._try_seed_frame("Re-seed failed, disabling native streaming: %s"): self._use_native_streaming = False @@ -462,11 +433,9 @@ class StreamTransportMixin: self, text: str, pre_fence_text: str, *, finalize: bool, is_turn_final: bool, ) -> Optional[bool]: """Draft frame while no message_id exists; None = not applicable / drafts just failed. - - Drafts are skipped when finalizing (the real send clears the draft), EXCEPT - stream-is-the-message adapters keep ONE stream per turn: a segment-break - finalize must not become a real send (it would seal at every tool boundary). - """ + Skipped when finalizing (the real send clears the draft), EXCEPT stream-is-the-message + adapters keep ONE stream per turn: a segment-break finalize must not become a real + send (it would seal at every tool boundary).""" stream_is_msg = self._stream_is_message() if finalize and not (stream_is_msg and not is_turn_final): return None @@ -561,10 +530,8 @@ class StreamTransportMixin: async def _on_edit_failure( self, result, text: str, *, finalize: bool, is_turn_final: bool, ) -> bool: - """Classify a failed edit: partial overflow, flood backoff, or fallback mode. - - Always False; the caller's finalize path may still deliver the tail. - """ + """Classify a failed edit: partial overflow, flood backoff, or fallback mode. Always + False; the caller's finalize path may still deliver the tail.""" turn_final = finalize and is_turn_final if ( turn_final diff --git a/gateway/stream_dispatch.py b/gateway/stream_dispatch.py index f2f7d4f3d9..140dd11b05 100644 --- a/gateway/stream_dispatch.py +++ b/gateway/stream_dispatch.py @@ -1,11 +1,9 @@ -"""Adapter-driven dispatch of structured stream events to a delivery sink. +"""Adapter-driven dispatch of structured stream events (gateway/stream_events.py). -``GatewayEventDispatcher`` routes each typed event (gateway/stream_events.py) -through the adapter's render hooks: message events flow into the consumer; tool -events are formatted by the adapter — which may return None to *eat* them on -platforms without tool chrome — and enqueued onto the same tool-progress queue the -gateway drains, so the two paths never race. Synchronous, no asyncio: callable -from the agent's worker thread. +Message events flow into the consumer; tool events are formatted by the adapter (None += eat it on platforms without tool chrome) and enqueued onto the same tool-progress +queue the gateway drains, so the two paths never race. Synchronous: callable from +the agent's worker thread. """ from __future__ import annotations @@ -23,15 +21,11 @@ logger = logging.getLogger("gateway.stream_events") class GatewayEventDispatcher: """Route typed stream events through an adapter onto a delivery sink. - adapter: provides ``render_message_event`` / ``format_tool_event``. - sink: the GatewayStreamConsumer; None when streaming is disabled (message - events are dropped — the final response still goes out normally). - enqueue_tool_line: puts a rendered tool-progress line on the gateway's - progress queue; None when tool progress is disabled for the channel. - tool_mode: "all" / "new" / "verbose" / "off". preview_max_len: resolved - ``tool_preview_length`` (0 = no cap in verbose). - on_long_tool / on_notice: optional hooks so the gateway owns the - "should I surface this here?" decision. + sink: the GatewayStreamConsumer, or None when streaming is disabled (message + events dropped; the final still goes out normally). enqueue_tool_line: None + when tool progress is disabled. tool_mode: "all"/"new"/"verbose"/"off"; + preview_max_len: ``tool_preview_length`` (0 = no cap in verbose). + on_long_tool / on_notice: the gateway owns the "surface this here?" decision. """ def __init__( diff --git a/gateway/stream_events.py b/gateway/stream_events.py index 86156a2cc3..20db054757 100644 --- a/gateway/stream_events.py +++ b/gateway/stream_events.py @@ -1,10 +1,8 @@ """Structured streaming events — the agent→gateway delivery contract. -A small typed vocabulary naming *what happened* without prescribing *how it is -delivered*: the agent emits these from its worker thread, ``GatewayStreamConsumer`` -is the single sink and the platform adapter decides rendering. Plain frozen -dataclasses: no behavior, no I/O, safe across the thread/async boundary. Events -describe *transport*, never *context* — whatever the gateway "eats" must never +Typed *what happened* events (frozen dataclasses, no I/O) emitted from the agent's +worker thread; ``GatewayStreamConsumer`` is the sink and the adapter decides rendering. +Events describe *transport*, never *context*: whatever the gateway "eats" must never diverge from the agent-owned message history. """ @@ -22,12 +20,8 @@ class MessageChunk: @dataclass(frozen=True) class MessageStop: - """The current assistant text segment is complete. - - ``final`` is True only for the terminal stop of the turn; an intermediate stop - (text → tool call → more text) makes the consumer start a fresh segment below - tool chrome. - """ + """Assistant text segment complete. ``final`` only for the turn's terminal stop; an + intermediate stop (text → tool → text) starts a fresh segment below tool chrome.""" final: bool = False @@ -43,16 +37,13 @@ class ToolCallChunk: tool_name: str preview: Optional[str] = None args: Optional[Dict[str, Any]] = None - # Monotonic per-turn index: correlates a finish with its start. - index: int = 0 + index: int = 0 # monotonic per-turn index: correlates a finish with its start @dataclass(frozen=True) class ToolCallFinished: - """A tool invocation completed. Tool *output* never travels here (it is history). - - Drives progress-bubble settling and one-time onboarding hints (LongToolHint). - """ + """A tool invocation completed (drives bubble settling + LongToolHint). Tool *output* + never travels here — it is history.""" tool_name: str duration: float = 0.0 # wall-clock seconds ok: bool = True # returned without raising @@ -61,21 +52,16 @@ class ToolCallFinished: @dataclass(frozen=True) class LongToolHint: - """One-shot onboarding nudge when a tool runs longer than the threshold. - - The gateway gates it on platform capability (/verbose usable) and first-time use. - """ + """One-shot onboarding nudge for a long tool run; the gateway gates it on platform + capability (/verbose usable) and first-time use.""" tool_name: str = "" duration: float = 0.0 @dataclass(frozen=True) class GatewayNotice: - """A gateway-originated control message. - - ``kind`` is a stable string adapters switch on (``"restart"`` / ``"online"`` / - ``"long_run"`` / …); ``text`` is the default rendering. - """ + """Gateway-originated control message; ``kind`` is a stable string adapters switch on + (``"restart"`` / ``"online"`` / ``"long_run"`` / …), ``text`` the default rendering.""" kind: str text: str = "" extra: Dict[str, Any] = field(default_factory=dict)