refactor(gateway): stream transport/fallback/events — docstrings to summary + invariant
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
+10
-16
@@ -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__(
|
||||
|
||||
+12
-26
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user