diff --git a/gateway/stream_consumer.py b/gateway/stream_consumer.py index 774fd606ffe87..1843dc0f0e596 100644 --- a/gateway/stream_consumer.py +++ b/gateway/stream_consumer.py @@ -165,11 +165,12 @@ def __init__( self._reset_message_state() # Transports, resolved in run(). Draft: animated frames via adapter.send_draft; - # the final still uses first-send; the first failure disables drafts. Native + # the final still uses first-send; a terminal failure disables drafts. Native # (WeCom msgtype "stream"): the ONLY channel — any failure falls back to edit/send. self._use_draft_streaming = False self._draft_id: Optional[int] = None self._draft_failures = 0 + self._draft_retry_until = 0.0 # TERMINAL authorization refusal for THIS RUN (see _send_draft_frame). # Per-run state, constructed fresh each turn, so a refusal can never # mute a healthy destination on a later turn. @@ -735,6 +736,8 @@ def _should_edit(self, tick: "_Tick") -> bool: if self._use_native_streaming: # No platform edit-rate limit: push every delta immediately. should_edit = bool(self._accumulated) or self._tool_progress_active + elif self._use_draft_streaming: + should_edit = self._should_push_draft() else: elapsed = time.monotonic() - self._last_edit_time # buffer_threshold is a codepoint debounce heuristic, not a diff --git a/gateway/stream_consumer_transport.py b/gateway/stream_consumer_transport.py index 3be6b06a52432..ca075bb3e6b9f 100644 --- a/gateway/stream_consumer_transport.py +++ b/gateway/stream_consumer_transport.py @@ -163,8 +163,15 @@ def _resolve_native_streaming(self) -> bool: logger.debug("supports_native_streaming probe raised", exc_info=True) return False + def _should_push_draft(self) -> bool: + """Client animation consumes snapshots; a full buffer must never bypass cadence.""" + now = time.monotonic() + return bool(self._accumulated) and now >= self._draft_retry_until and ( + self._last_edit_time == 0.0 + or now - self._last_edit_time >= self._current_edit_interval) + async def _send_draft_frame(self, text: str) -> bool: - """Emit one draft frame; any failure permanently disables drafts for this run. + """Emit a draft frame; budget deferrals and flood cooldowns keep drafts enabled. 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()). @@ -178,7 +185,12 @@ async def _send_draft_frame(self, text: str) -> bool: logger.debug("send_draft raised, disabling draft transport for this run: %s", e) else: if getattr(result, "success", False): + raw_response = getattr(result, "raw_response", None) + if isinstance(raw_response, dict) and raw_response.get("skipped"): + self._defer_draft_retry(result) + return True self._last_sent_text = text # parity with the edit-based no-op skip + self._draft_retry_until = 0.0 return True # P5(b): an AUTHORIZATION decline is terminal for the whole run, not # merely "drafts are unusable". Disabling drafts alone routes the @@ -193,12 +205,23 @@ async def _send_draft_frame(self, text: str) -> bool: "is not approved for this connection)" ) self._egress_declined = True + elif self._defer_draft_retry(result): + return True logger.debug("send_draft returned success=False, disabling draft transport: %s", getattr(result, "error", "unknown")) self._draft_failures += 1 self._use_draft_streaming = False return False + def _defer_draft_retry(self, result) -> bool: + """Keep the latest snapshot pending during a server-requested cooldown.""" + retry_after = getattr(result, "retry_after", None) + if not isinstance(retry_after, (int, float)) or retry_after <= 0: + return False + self._draft_retry_until = max( + self._draft_retry_until, time.monotonic() + float(retry_after)) + return True + async def _abandon_native_stream(self) -> None: """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. @@ -434,7 +457,9 @@ async def _draft_push(self, text: str, pre_fence_text: str, *, finalize: bool, stream_is_msg = self._stream_is_message() if finalize and not (stream_is_msg and not is_turn_final): return None - frame_text = pre_fence_text if stream_is_msg else text + preserve_prefix = stream_is_msg or ( + getattr(type(self.adapter), "DRAFT_STREAM_PREFIX_STABLE", False) is True) + frame_text = pre_fence_text if preserve_prefix else text # Strip the cursor: native streams render their own indicator, and # "...text▉" is never a prefix of "...text more▉", which forces the # connector's whole-text re-append on EVERY tick (stacked copies). diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index 5850c2083add7..fe75bf26b793f 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -178,6 +178,7 @@ def _bold_label_html(line: str) -> str: "video file": "platform.telegram.media.kind_video"} from gateway.platforms.event import MessageEvent, MessageType, ProcessingOutcome +from plugins.platforms.telegram.telegram_drafts import TelegramDraftMixin from plugins.platforms.telegram.telegram_entities import expand_link_entities from plugins.platforms.telegram.telegram_held_inbound import TelegramHeldInboundMixin from plugins.platforms.telegram.telegram_ids import normalize_telegram_chat_id @@ -511,7 +512,7 @@ class _PollingStallError(RuntimeError): """ -class TelegramAdapter(TelegramHeldInboundMixin, BasePlatformAdapter): +class TelegramAdapter(TelegramDraftMixin, TelegramHeldInboundMixin, BasePlatformAdapter): """Telegram bot adapter: users/groups, MarkdownV2 replies, forum topics, media.""" # Bound for the per-(chat_id, status_key) status-message cache; FIFO half-trim on overflow. @@ -1567,27 +1568,6 @@ def _should_attempt_rich_draft(self, content: str) -> bool: and not getattr(self, "_rich_draft_disabled", False) and self._rich_content_ok(content)) - async def _try_send_rich_draft(self, chat_id: str, draft_id: int, content: str, metadata: Optional[Dict[str, Any]]) -> bool: - """Emit one ``sendRichMessageDraft`` frame; True on success. Frames are ephemeral, so any failure - returns False and the caller renders the legacy draft; capability failures latch off.""" - payload: Dict[str, Any] = { - "chat_id": normalize_telegram_chat_id(chat_id), "draft_id": int(draft_id), "rich_message": self._rich_message_payload(content)} - payload.update(self._thread_kwargs_for_draft(chat_id, metadata)) - try: - return bool(await _await_with_thread_deadline( - self._bot.do_api_request("sendRichMessageDraft", api_kwargs=payload), - timeout=_TEXT_SEND_DEADLINE, label="telegram-send", dump_on_blocked_loop=False)) - except Exception as exc: - if self._is_rich_capability_error(exc): - self._rich_draft_disabled = True - logger.debug( - "[%s] sendRichMessageDraft unsupported (%s) — using legacy drafts", self.name, _redact_telegram_error_text(exc)) - else: - logger.debug( - "[%s] sendRichMessageDraft transient failure (%s) — legacy draft this frame", self.name, - _redact_telegram_error_text(exc)) - return False - async def _drain_polling_connections(self) -> None: """Reset the httpx pool used for getUpdates polling before a reconnect. @@ -4071,49 +4051,6 @@ def supports_draft_streaming(self, chat_type: Optional[str] = None, metadata: Op return False return (chat_type or "").lower() in {"dm", "private"} - async def send_draft(self, chat_id: str, draft_id: int, content: str, metadata: Optional[Dict[str, Any]] = None) -> SendResult: - """Stream a partial message via ``sendRichMessageDraft`` (when rich is enabled and supported) else - ``sendMessageDraft``; reusing ``draft_id`` animates the preview. The caller sends the final text.""" - if not self._bot: - return SendResult(success=False, error="not_connected") - # Rich draft fast-path; any failure degrades to the plain draft below. Drafts have no message_id. - if self._should_attempt_rich_draft(content) and await self._try_send_rich_draft(chat_id, draft_id, content, metadata): - return SendResult(success=True, message_id=None) - if not hasattr(self._bot, "send_message_draft"): - return SendResult(success=False, error="api_unavailable") - # Drafts share the regular-send UTF-16 length contract. - text = content if len( - content) <= self.MAX_MESSAGE_LENGTH else self.truncate_message(content, self.MAX_MESSAGE_LENGTH, len_fn=utf16_len)[0] - # Same MarkdownV2 conversion as ``send`` (MarkdownV2 then plain) so the draft doesn't snap at the end. Exception: a Rich - # final with rich drafts disabled previews raw — the legacy formatter would turn pipe tables into bullets. - plain_rich_preview = bool( - getattr(self, "_rich_messages_enabled", False) and not getattr(self, "_rich_drafts_enabled", False) - and self._needs_rich_rendering(text)) - draft_thread_kwargs = self._thread_kwargs_for_draft(chat_id, metadata) - for use_markdown in ((False,) if plain_rich_preview else (True, False)): - kwargs: Dict[str, Any] = { - "chat_id": normalize_telegram_chat_id(chat_id), "draft_id": int(draft_id), - "text": self.format_message(text) if use_markdown else text} - if use_markdown: - kwargs["parse_mode"] = ParseMode.MARKDOWN_V2 - kwargs.update(draft_thread_kwargs) - try: - if await _await_with_thread_deadline( - self._bot.send_message_draft(**kwargs), timeout=_TEXT_SEND_DEADLINE, label="telegram-send", dump_on_blocked_loop=False): - return SendResult(success=True, message_id=None) - return SendResult(success=False, error="draft_rejected") - except Exception as e: - # MarkdownV2 parse failure → retry once as plain text; anything else returns to the caller, - # which falls back to edit-based streaming for this response. - if use_markdown and self._is_bad_request_error(e): - logger.debug( - "[%s] sendMessageDraft MarkdownV2 rejected, retrying as plain text (chat=%s draft_id=%s): %s", - self.name, chat_id, draft_id, _redact_telegram_error_text(e)) - continue - logger.debug("[%s] sendMessageDraft failed (chat=%s draft_id=%s): %s", self.name, chat_id, draft_id, e) - return SendResult(success=False, error=_redact_telegram_error_text(e)) - return SendResult(success=False, error="draft_rejected") - async def _send_message_with_thread_fallback(self, **kwargs): """Send a control-style message (approval prompts, pickers), retrying once without message_thread_id on 'Message thread not found' (stale thread_id); ``send`` has its own. @@ -5467,23 +5404,23 @@ def _record_typing_cooldown(self, chat_id: str, exc: Exception) -> None: """Suppress Telegram typing refreshes for this chat after transient failures.""" if not hasattr(self, "_telegram_typing_cooldown_until"): self._telegram_typing_cooldown_until = {} - retry_after = getattr(exc, "retry_after", None) + retry_after = self._record_ai_action_cooldown(chat_id, exc) try: delay = float(retry_after) if retry_after is not None else self._telegram_typing_cooldown_seconds except (TypeError, ValueError): delay = self._telegram_typing_cooldown_seconds - self._telegram_typing_cooldown_until[str(chat_id)] = asyncio.get_running_loop().time() + max(1.0, min(delay, 300.0)) + self._telegram_typing_cooldown_until[str(normalize_telegram_chat_id(chat_id))] = asyncio.get_running_loop().time() + max(1.0, min(delay, 300.0)) def _typing_in_cooldown(self, chat_id: str) -> bool: if not hasattr(self, "_telegram_typing_cooldown_until"): self._telegram_typing_cooldown_until = {} self._telegram_typing_cooldown_seconds = 30.0 - until = self._telegram_typing_cooldown_until.get(str(chat_id)) + until = self._telegram_typing_cooldown_until.get(str(normalize_telegram_chat_id(chat_id))) if until is None: return False if asyncio.get_running_loop().time() < until: return True - self._telegram_typing_cooldown_until.pop(str(chat_id), None) + self._telegram_typing_cooldown_until.pop(str(normalize_telegram_chat_id(chat_id)), None) return False # --- per-chat send ordering + flood cooldown (#114396) --------------------------------------------- @@ -5561,14 +5498,16 @@ def _hold_chat_outbound_slot(self, chat_id: Any) -> None: async def send_typing(self, chat_id: str, metadata: Optional[Dict[str, Any]] = None) -> None: """Send typing indicator.""" - if not self._bot or self._typing_in_cooldown(chat_id): + if not self._bot or self._typing_in_cooldown(chat_id) or self._native_draft_recent(chat_id): return _is_dm_topic: bool = False message_thread_id: Optional[int] = None async def _action(**kw) -> None: + if self._claim_ai_action_slot(chat_id): + return await self._bot.send_chat_action(chat_id=normalize_telegram_chat_id(chat_id), action="typing", **kw) - self._telegram_typing_cooldown_until.pop(str(chat_id), None) + self._telegram_typing_cooldown_until.pop(str(normalize_telegram_chat_id(chat_id)), None) try: _is_dm_topic = self._dm_topic_fallback(metadata) message_thread_id = self._message_thread_id_for_typing(self._metadata_thread_id(metadata)) @@ -5576,15 +5515,15 @@ async def _action(**kw) -> None: except Exception as e: # DM topic lanes: Telegram may reject message_thread_id — retry without it so the indicator at # least appears in the main DM view. - if _is_dm_topic and message_thread_id is not None: + if self._is_transient_typing_error(e): + self._record_typing_cooldown(chat_id, e) + elif _is_dm_topic and message_thread_id is not None: try: await _action() return except Exception as fallback_exc: if self._is_transient_typing_error(fallback_exc): self._record_typing_cooldown(chat_id, fallback_exc) - elif self._is_transient_typing_error(e): - self._record_typing_cooldown(chat_id, e) logger.debug("[%s] Failed to send Telegram typing indicator: %s", self.name, _redact_telegram_error_text(e), exc_info=True) async def get_chat_info(self, chat_id: str) -> Dict[str, Any]: diff --git a/plugins/platforms/telegram/telegram_drafts.py b/plugins/platforms/telegram/telegram_drafts.py new file mode 100644 index 0000000000000..de1d7d4df09b6 --- /dev/null +++ b/plugins/platforms/telegram/telegram_drafts.py @@ -0,0 +1,160 @@ +"""Prefix-stable native drafts and their shared Telegram typing-action budget.""" + +from __future__ import annotations + +import asyncio +import logging +import math +from collections import deque +from typing import Any, Optional + +from gateway.platforms.base import SendResult, _prefix_within_utf16_limit +from plugins.platforms.telegram.telegram_ids import normalize_telegram_chat_id + +logger = logging.getLogger(__name__) + + +def _ai_action_now() -> float: + return asyncio.get_running_loop().time() + + +class TelegramDraftMixin: + DRAFT_STREAM_PREFIX_STABLE = True + + def _claim_ai_action_slot(self, chat_id: Any) -> float: + """Reserve one actual draft/typing request, or return its non-blocking retry delay. + + Telegram shares 20 calls/5s and 40 calls/30s per peer across sendMessageDraft and + sendChatAction: https://core.telegram.org/api/bots/ai. Reserving before the await + prevents concurrent topics or typing refreshes from spending the same slot. + """ + key = str(normalize_telegram_chat_id(chat_id)) + now = _ai_action_now() + cooldowns = self.__dict__.setdefault("_telegram_ai_action_cooldown_until", {}) + remaining = cooldowns.get(key, now) - now + if remaining > 0: + return remaining + cooldowns.pop(key, None) + calls = self.__dict__.setdefault("_telegram_ai_action_calls", {}).setdefault(key, deque()) + while calls and now - calls[0] >= 30.0: + calls.popleft() + short_delay = calls[-20] + 5.0 - now if len(calls) >= 20 else 0.0 + long_delay = calls[0] + 30.0 - now if len(calls) >= 40 else 0.0 + delay = max(short_delay, long_delay, 0.0) + if delay == 0: + calls.append(now) + return delay + + def _record_ai_action_cooldown(self, chat_id: Any, exc: Exception) -> Optional[float]: + """Respect RetryAfter for drafts and typing without delaying completed messages.""" + retry_after = getattr(exc, "retry_after", None) + if retry_after is None: + return None + if hasattr(retry_after, "total_seconds"): + retry_after = retry_after.total_seconds() + try: + delay = float(retry_after) + except (TypeError, ValueError): + return None + if not math.isfinite(delay): + return None + delay = max(1.0, delay) + key = str(normalize_telegram_chat_id(chat_id)) + cooldowns = self.__dict__.setdefault("_telegram_ai_action_cooldown_until", {}) + cooldowns[key] = max(cooldowns.get(key, 0.0), _ai_action_now() + delay) + return delay + + @staticmethod + def _draft_skipped(delay: float) -> SendResult: + return SendResult(success=True, raw_response={"skipped": True}, retry_after=delay) + + def _native_draft_recent(self, chat_id: Any) -> bool: + """A native draft is itself a typing action; skip redundant indicator refreshes.""" + key = str(normalize_telegram_chat_id(chat_id)) + sent_at = self.__dict__.setdefault("_telegram_native_draft_sent_at", {}) + last_sent = sent_at.get(key) + if last_sent is not None and _ai_action_now() - last_sent < 4.0: + return True + sent_at.pop(key, None) + return False + + def _record_native_draft(self, chat_id: Any) -> SendResult: + sent_at = self.__dict__.setdefault("_telegram_native_draft_sent_at", {}) + sent_at[str(normalize_telegram_chat_id(chat_id))] = _ai_action_now() + return SendResult(success=True) + + async def _try_send_rich_draft( + self, chat_id: str, draft_id: int, content: str, metadata: Optional[dict[str, Any]], + ) -> Optional[SendResult]: + """Preserve explicit rich drafts, falling back only outside a flood cooldown.""" + from plugins.platforms.telegram.adapter import ( + _TEXT_SEND_DEADLINE, _await_with_thread_deadline, _redact_telegram_error_text, + ) + + delay = self._claim_ai_action_slot(chat_id) + if delay: + return self._draft_skipped(delay) + payload = { + "chat_id": normalize_telegram_chat_id(chat_id), "draft_id": int(draft_id), + "rich_message": self._rich_message_payload(content), + } + payload.update(self._thread_kwargs_for_draft(chat_id, metadata)) + try: + accepted = await _await_with_thread_deadline( + self._bot.do_api_request("sendRichMessageDraft", api_kwargs=payload), + timeout=_TEXT_SEND_DEADLINE, label="telegram-send", dump_on_blocked_loop=False, + ) + except Exception as exc: # health: allow BLE001 -- API boundary; raw tracebacks may expose bot-token URLs, so log redacted text. + delay = self._record_ai_action_cooldown(chat_id, exc) + if delay is not None: + return self._draft_skipped(delay) + if self._is_rich_capability_error(exc): + self._rich_draft_disabled = True + logger.debug( + "[%s] sendRichMessageDraft rejected; using legacy draft: %s", + self.name, _redact_telegram_error_text(exc), + ) + return None + return self._record_native_draft(chat_id) if accepted else None + + async def send_draft( + self, chat_id: str, draft_id: int, content: str, metadata: Optional[dict[str, Any]] = None, + ) -> SendResult: + """Send an append-only plain preview; persistent finals retain their normal formatting.""" + from plugins.platforms.telegram.adapter import ( + _TEXT_SEND_DEADLINE, _await_with_thread_deadline, _redact_telegram_error_text, + ) + + if not self._bot: + return SendResult(success=False, error="not_connected") + if self._should_attempt_rich_draft(content): + rich_result = await self._try_send_rich_draft(chat_id, draft_id, content, metadata) + if rich_result is not None: + return rich_result + if not hasattr(self._bot, "send_message_draft"): + return SendResult(success=False, error="api_unavailable") + delay = self._claim_ai_action_slot(chat_id) + if delay: + return self._draft_skipped(delay) + # The Markdown renderer and chunker can rewrite earlier characters as fences/tables grow. + # A raw UTF-16-bounded prefix lets Telegram animate only newly appended characters. + kwargs = { + "chat_id": normalize_telegram_chat_id(chat_id), "draft_id": int(draft_id), + "text": _prefix_within_utf16_limit(content, self.MAX_MESSAGE_LENGTH), + **self._thread_kwargs_for_draft(chat_id, metadata), + } + try: + accepted = await _await_with_thread_deadline( + self._bot.send_message_draft(**kwargs), timeout=_TEXT_SEND_DEADLINE, + label="telegram-send", dump_on_blocked_loop=False, + ) + except Exception as exc: # health: allow BLE001 -- API boundary; raw tracebacks may expose bot-token URLs, so log redacted text. + delay = self._record_ai_action_cooldown(chat_id, exc) + if delay is not None: + return self._draft_skipped(delay) + safe_error = _redact_telegram_error_text(exc) + logger.debug("[%s] sendMessageDraft failed: %s", self.name, safe_error) + return SendResult(success=False, error=safe_error) + if accepted: + return self._record_native_draft(chat_id) + return SendResult(success=False, error="draft_rejected") diff --git a/tests/gateway/test_stream_consumer_native_draft.py b/tests/gateway/test_stream_consumer_native_draft.py new file mode 100644 index 0000000000000..a8881cd691240 --- /dev/null +++ b/tests/gateway/test_stream_consumer_native_draft.py @@ -0,0 +1,163 @@ +"""Animated draft previews coalesce deltas without rewriting already shown text.""" + +from types import SimpleNamespace + +import pytest + +from gateway.platforms.base import BasePlatformAdapter, SendResult +from gateway.stream_consumer import GatewayStreamConsumer, StreamConsumerConfig + + +class AnimatedDraftAdapter(BasePlatformAdapter): + DRAFT_STREAM_PREFIX_STABLE = True + + async def connect(self): + return True + + async def disconnect(self): + pass + + def supports_draft_streaming(self, chat_type=None, metadata=None, **kwargs): + return chat_type == "dm" + + async def send_draft(self, chat_id, draft_id, content, metadata=None): + self.drafts.append((draft_id, content)) + return self.draft_results.pop(0) if self.draft_results else SendResult(success=True) + + async def send(self, chat_id, content, reply_to=None, metadata=None): + self.messages.append(content) + return SendResult(success=True, message_id=str(len(self.messages))) + + async def edit_message(self, chat_id, message_id, content, **kwargs): + self.messages.append(content) + return SendResult(success=True, message_id=message_id) + + async def send_typing(self, chat_id, metadata=None): + pass + + async def get_chat_info(self, chat_id): + return {"type": "dm"} + + +def make_consumer(monkeypatch, *, cursor="", edit_interval=0.8): + adapter = AnimatedDraftAdapter.__new__(AnimatedDraftAdapter) + adapter.drafts = [] + adapter.draft_results = [] + adapter.messages = [] + clock = SimpleNamespace(now=100.0) + fake_time = SimpleNamespace(monotonic=lambda: clock.now) + monkeypatch.setattr("gateway.stream_consumer.time", fake_time) + monkeypatch.setattr("gateway.stream_consumer_transport.time", fake_time) + consumer = GatewayStreamConsumer( + adapter, "123", StreamConsumerConfig( + transport="draft", chat_type="dm", edit_interval=edit_interval, + buffer_threshold=4, cursor=cursor, + ), + ) + return consumer, adapter, clock + + +async def push_delta(consumer, text): + consumer.on_delta(text) + tick = consumer._drain_queue() + if consumer._should_edit(tick): + await consumer._push_update(tick) + + +@pytest.mark.asyncio +async def test_native_drafts_coalesce_fast_deltas_and_publish_latest_snapshot(monkeypatch): + consumer, adapter, clock = make_consumer(monkeypatch) + await consumer._start_transports() + + await push_delta(consumer, "首个回复") + for delta in (",", "接", "着", "输", "出"): + clock.now += 0.05 + await push_delta(consumer, delta) + assert [text for _, text in adapter.drafts] == ["首个回复"] + + clock.now = 100.81 + await push_delta(consumer, "。") + assert [text for _, text in adapter.drafts] == ["首个回复", "首个回复,接着输出。"] + assert len({draft_id for draft_id, _ in adapter.drafts}) == 1 + + +@pytest.mark.asyncio +async def test_native_draft_code_prefix_has_no_synthetic_fence_or_cursor(monkeypatch): + consumer, adapter, clock = make_consumer(monkeypatch, cursor="▉") + await consumer._start_transports() + + await push_delta(consumer, "代码:\n```python\nprint(") + clock.now += 0.81 + await push_delta(consumer, "'hello')\n```") + + frames = [text for _, text in adapter.drafts] + assert frames == ["代码:\n```python\nprint(", "代码:\n```python\nprint('hello')\n```"] + assert frames[1].startswith(frames[0]) + + +@pytest.mark.asyncio +async def test_skipped_draft_retries_latest_snapshot_without_claiming_delivery(monkeypatch): + consumer, adapter, clock = make_consumer(monkeypatch) + await consumer._start_transports() + adapter.draft_results = [SendResult(success=True, raw_response={"skipped": True})] + + await push_delta(consumer, "首个回复") + assert consumer._last_sent_text == "" + assert consumer.already_sent is False + assert consumer._use_draft_streaming is True + assert adapter.messages == [] + + clock.now += 0.81 + await push_delta(consumer, ",后续内容") + assert consumer._last_sent_text == "首个回复,后续内容" + assert adapter.drafts[-1][1] == "首个回复,后续内容" + assert len({draft_id for draft_id, _ in adapter.drafts}) == 1 + + +@pytest.mark.asyncio +@pytest.mark.parametrize("success", [False, True]) +async def test_draft_cooldown_keeps_preview_transport_and_full_final_is_immediate( + monkeypatch, success, +): + consumer, adapter, clock = make_consumer(monkeypatch) + await consumer._start_transports() + adapter.draft_results = [SendResult( + success=success, error=None if success else "flood_control:9", retry_after=9, + raw_response={"skipped": True} if success else None, + )] + + await push_delta(consumer, "首个回复") + clock.now += 1.0 + await push_delta(consumer, ",新内容") + assert len(adapter.drafts) == 1 + assert consumer._use_draft_streaming is True + assert consumer._last_sent_text == "" + assert adapter.messages == [] + + # Completion bypasses preview cooldown and carries post-stream augmentation. + final_text = "首个回复,新内容。\n\n完整的最终结论。" + consumer.finish(final_text) + tick = consumer._drain_queue() + await consumer._push_update(tick) + await consumer._finalize_turn(tick) + assert adapter.messages == [final_text] + assert consumer.delivered_final_matches(final_text) is True + + +@pytest.mark.asyncio +async def test_draft_resumes_after_cooldown_with_latest_accumulated_text(monkeypatch): + consumer, adapter, clock = make_consumer(monkeypatch) + await consumer._start_transports() + adapter.draft_results = [SendResult( + success=True, retry_after=9, raw_response={"skipped": True}, + )] + + await push_delta(consumer, "首个回复") + clock.now += 1.0 + await push_delta(consumer, ",新内容") + clock.now += 8.1 + await push_delta(consumer, "。") + + assert [text for _, text in adapter.drafts] == ["首个回复", "首个回复,新内容。"] + assert consumer._last_sent_text == "首个回复,新内容。" + assert adapter.messages == [] diff --git a/tests/gateway/test_telegram_native_draft_budget.py b/tests/gateway/test_telegram_native_draft_budget.py new file mode 100644 index 0000000000000..247754548ceb1 --- /dev/null +++ b/tests/gateway/test_telegram_native_draft_budget.py @@ -0,0 +1,131 @@ +"""Native previews remain prefix-stable and share Telegram's typing-action limits.""" + +from datetime import timedelta +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock + +import pytest + +from gateway.config import PlatformConfig +from gateway.platforms.base import utf16_len +from plugins.platforms.telegram import adapter as tg_mod +from plugins.platforms.telegram.adapter import TelegramAdapter + + +@pytest.fixture +def draft_adapter(monkeypatch): + now = SimpleNamespace(value=1000.0) + monkeypatch.setattr( + tg_mod.asyncio, "get_running_loop", lambda: SimpleNamespace(time=lambda: now.value), + ) + + async def await_request(awaitable, **kwargs): + return await awaitable + + monkeypatch.setattr(tg_mod, "_await_with_thread_deadline", await_request) + adapter = TelegramAdapter(PlatformConfig(enabled=True, token="fake-token")) + adapter._bot = MagicMock() + adapter._bot.send_message_draft = AsyncMock(return_value=True) + adapter._bot.send_chat_action = AsyncMock() + return adapter, now + + +@pytest.mark.asyncio +@pytest.mark.parametrize("content", ["```python\nx = ", "😀" * 2100], ids=["open-code-fence", "astral-limit"]) +async def test_plain_draft_keeps_original_prefix_within_utf16_limit(draft_adapter, content): + adapter, _ = draft_adapter + adapter.format_message = lambda text: f"changed::{text}" + + result = await adapter.send_draft("123", 7, content) + + kwargs = adapter._bot.send_message_draft.call_args.kwargs + assert result.success + assert "parse_mode" not in kwargs + assert content.startswith(kwargs["text"]) + assert utf16_len(kwargs["text"]) <= adapter.MAX_MESSAGE_LENGTH + assert kwargs["text"] == ("😀" * 2048 if content.startswith("😀") else content) + + +@pytest.mark.asyncio +async def test_drafts_and_typing_share_both_windows_for_normalized_peer(draft_adapter): + adapter, now = draft_adapter + for _ in range(10): + await adapter.send_typing("00123") + for index in range(10): + await adapter.send_draft("123", 7, f"first {index}") + + skipped = await adapter.send_draft("00123", 7, "latest first batch") + assert skipped.success and skipped.raw_response["skipped"] + assert adapter._bot.send_message_draft.await_count == 10 + assert adapter._bot.send_chat_action.await_count == 10 + + now.value += 5.01 + for index in range(20): + await adapter.send_draft("123", 7, f"second {index}") + skipped = await adapter.send_draft("123", 7, "latest second batch") + assert skipped.success and skipped.raw_response["skipped"] + assert adapter._bot.send_message_draft.await_count == 30 + + other_peer = await adapter.send_draft("124", 8, "independent chat") + assert other_peer.success and not other_peer.raw_response + now.value = 1030.01 + latest = await adapter.send_draft("00123", 7, "latest, not a queued older frame") + assert latest.success and not latest.raw_response + assert adapter._bot.send_message_draft.call_args.kwargs["text"] == "latest, not a queued older frame" + + +@pytest.mark.asyncio +@pytest.mark.parametrize("rich", [False, True]) +async def test_retry_after_skips_preview_and_typing_then_resumes_native_draft(draft_adapter, rich): + adapter, now = draft_adapter + error = RuntimeError("Flood control exceeded") + error.retry_after = timedelta(seconds=3) + draft_api = adapter._bot.send_message_draft + if rich: + adapter._rich_messages_enabled = adapter._rich_drafts_enabled = True + adapter._allow_cjk_rich_messages = True + adapter._rich_content_ok = lambda content: True + draft_api = adapter._bot.do_api_request = AsyncMock() + draft_api.side_effect = [error, True] + + rejected = await adapter.send_draft("123", 7, "first preview") + assert rejected.success and rejected.raw_response["skipped"] + assert rejected.retry_after == pytest.approx(3) + assert draft_api.await_count == 1 + if rich: + adapter._bot.send_message_draft.assert_not_awaited() + + await adapter.send_typing("00123") + during_cooldown = await adapter.send_draft("00123", 7, "newer preview") + assert during_cooldown.success and during_cooldown.raw_response["skipped"] + adapter._bot.send_chat_action.assert_not_awaited() + assert draft_api.await_count == 1 + + now.value += 3.01 + resumed = await adapter.send_draft("123", 7, "latest preview") + assert resumed.success and not resumed.raw_response + assert draft_api.await_count == 2 + assert not adapter._rich_draft_disabled + + +@pytest.mark.asyncio +async def test_recent_native_draft_replaces_redundant_typing_refresh(draft_adapter): + adapter, now = draft_adapter + await adapter.send_draft("123", 7, "working preview") + await adapter.send_typing("00123") + adapter._bot.send_chat_action.assert_not_awaited() + + now.value += 4.01 + await adapter.send_typing("123") + adapter._bot.send_chat_action.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_unsupported_draft_still_reports_failure_for_edit_fallback(draft_adapter): + adapter, _ = draft_adapter + adapter._bot.send_message_draft.side_effect = RuntimeError("method not found") + + result = await adapter.send_draft("123", 7, "preview") + + assert not result.success + assert "method not found" in result.error diff --git a/tests/gateway/test_telegram_native_stream_pipeline.py b/tests/gateway/test_telegram_native_stream_pipeline.py new file mode 100644 index 0000000000000..e344596bbe7b7 --- /dev/null +++ b/tests/gateway/test_telegram_native_stream_pipeline.py @@ -0,0 +1,129 @@ +"""Native Telegram previews preserve prefixes without changing final delivery. + +Exercise the real consumer and adapter together; only the bot transport is fake. +Queue acknowledgements synchronise input with delivered frames, without asserting +wall-clock latency or claiming to reproduce a Telegram client's animation. +""" + +import asyncio +from types import SimpleNamespace + +import pytest + +from gateway.config import PlatformConfig +from gateway.stream_consumer import GatewayStreamConsumer, StreamConsumerConfig +from plugins.platforms.telegram.adapter import TelegramAdapter + + +class RecordingBot: + def __init__(self, *, reject_drafts=False): + self.reject_drafts = reject_drafts + self.drafts = [] + self.messages = [] + self.edits = [] + self.rich_messages = [] + self.frames = asyncio.Queue() + self.deliveries = asyncio.Queue() + + async def send_message_draft(self, **kwargs): + self.drafts.append(kwargs) + self.frames.put_nowait(kwargs) + return not self.reject_drafts + + async def send_message(self, **kwargs): + self.messages.append(kwargs) + self.deliveries.put_nowait(kwargs) + return SimpleNamespace(message_id=123) + + async def edit_message_text(self, **kwargs): + self.edits.append(kwargs) + return SimpleNamespace(message_id=123) + + async def do_api_request(self, endpoint, *, api_kwargs): + self.rich_messages.append((endpoint, api_kwargs)) + return SimpleNamespace(message_id=123) + + async def send_chat_action(self, **kwargs): + return True + + +def _consumer(bot, *, rich_messages=False): + adapter = TelegramAdapter(PlatformConfig( + enabled=True, token="fake-token", + extra={"rich_messages": rich_messages, "rich_drafts": False}, + )) + adapter._bot = bot + consumer = GatewayStreamConsumer( + adapter, "12345", + StreamConsumerConfig( + transport="draft", chat_type="dm", edit_interval=0.05, + buffer_threshold=1, cursor=" ▉", + ), + ) + return adapter, consumer + + +@pytest.mark.asyncio +async def test_burst_drafts_are_plain_prefixes_and_final_retains_rich_content(): + bot = RecordingBot() + _adapter, consumer = _consumer(bot, rich_messages=True) + first = "**Steps**\n\n```python\nprint(" + middle = "'ready')\n" + tail = "```\n\n| Item | Status |\n|---|---|\n| stream | complete |" + complete = first + middle + tail + final = complete + "\n\nVerification details: " + "detail " * 800 + "complete." + + consumer.on_delta(first) + task = asyncio.create_task(consumer.run()) + try: + await asyncio.wait_for(bot.frames.get(), timeout=5) + # Several producer deltas arrive together: the next preview must catch up + # to all of them instead of slowly replaying a fixed character allowance. + for delta in (middle, tail): + consumer.on_delta(delta) + await asyncio.wait_for(bot.frames.get(), timeout=5) + consumer.finish(final) + await asyncio.wait_for(task, timeout=5) + finally: + if not task.done(): + task.cancel() + await task + + assert [frame["text"] for frame in bot.drafts] == [first, complete] + assert len({frame["draft_id"] for frame in bot.drafts}) == 1 + assert all("parse_mode" not in frame for frame in bot.drafts) + assert bot.messages == [] + assert bot.edits == [] + assert bot.rich_messages == [( + "sendRichMessage", + { + "chat_id": 12345, + "rich_message": {"markdown": final}, + }, + )] + assert consumer.delivered_final_matches(final) + + +@pytest.mark.asyncio +async def test_rejected_native_draft_falls_back_without_losing_final_content(): + bot = RecordingBot(reject_drafts=True) + adapter, consumer = _consumer(bot) + first = "A complete **preview**" + final = first + "\n\nFinal details: [reference](https://example.com)." + + consumer.on_delta(first) + task = asyncio.create_task(consumer.run()) + try: + await asyncio.wait_for(bot.deliveries.get(), timeout=5) + consumer.finish(final) + await asyncio.wait_for(task, timeout=5) + finally: + if not task.done(): + task.cancel() + await task + + assert len(bot.drafts) == 1 + assert len(bot.messages) == 1 + assert bot.rich_messages == [] + assert bot.edits[-1]["text"] == adapter.format_message(final) + assert consumer.delivered_final_matches(final) diff --git a/tests/gateway/test_telegram_send_draft_format.py b/tests/gateway/test_telegram_send_draft_format.py index 91e2a2ec8673c..b9d6d754c0798 100644 --- a/tests/gateway/test_telegram_send_draft_format.py +++ b/tests/gateway/test_telegram_send_draft_format.py @@ -1,60 +1,30 @@ -"""TelegramAdapter.send_draft MarkdownV2 formatting parity. +"""Native drafts preserve text prefixes; final replies retain Markdown formatting.""" -Bot API 9.5 ``sendMessageDraft`` powers the animated streaming preview in -DMs. The regular ``send`` path renders with MarkdownV2, so the draft must -too — otherwise the live preview streams as raw text and the final -``sendMessage`` snaps into formatted output, producing a jarring visual -shift at the end of the response (reported by an external user, May 2026). - -These tests pin: - 1. The happy path passes ``parse_mode=MARKDOWN_V2`` with format_message'd - text (formatting parity with the final message). - 2. A MarkdownV2 BadRequest triggers a single plain-text retry rather than - killing draft streaming for the whole response. - 3. A non-BadRequest failure propagates so the caller falls back to edit. -""" from unittest.mock import AsyncMock, MagicMock import pytest from gateway.config import PlatformConfig -import plugins.platforms.telegram.adapter as tg_mod # noqa: E402 -from plugins.platforms.telegram.adapter import TelegramAdapter # noqa: E402 +from plugins.platforms.telegram.adapter import TelegramAdapter +from telegram.constants import ParseMode -def _make_adapter() -> TelegramAdapter: - adapter = TelegramAdapter(PlatformConfig(enabled=True, token="***")) +@pytest.mark.asyncio +async def test_plain_preview_keeps_markdown_quality_in_persistent_final(): + adapter = TelegramAdapter(PlatformConfig(enabled=True, token="fake-token")) adapter._bot = MagicMock() adapter._bot.send_message_draft = AsyncMock(return_value=True) - return adapter - - -@pytest.mark.asyncio -async def test_send_draft_falls_back_to_plain_text_on_markdownv2_error(): - """A MarkdownV2 BadRequest retries once as plain text (no parse_mode), - instead of aborting draft streaming for the whole response.""" - adapter = _make_adapter() - adapter.format_message = lambda content: f"FMT::{content}" - - # Resolve the BadRequest type the adapter checks via _is_bad_request_error. - from telegram.error import BadRequest # type: ignore - calls = [] - - async def _draft(**kwargs): - calls.append(kwargs) - if "parse_mode" in kwargs: - raise BadRequest("can't parse entities") - return True - - adapter._bot.send_message_draft = AsyncMock(side_effect=_draft) - - result = await adapter.send_draft("123", 9, "weird _text") - - assert result.success is True - # First attempt: MarkdownV2; second attempt: plain text, no parse_mode. - assert len(calls) == 2 - assert "parse_mode" in calls[0] - assert "parse_mode" not in calls[1] - assert calls[1]["text"] == "weird _text" # raw, unformatted - - + adapter._bot.send_message = AsyncMock(return_value=MagicMock(message_id=9)) + adapter._bot.send_chat_action = AsyncMock() + content = "**A bold answer** with `inline_code` and _underscores_" + + preview = await adapter.send_draft("123", 7, content) + final = await adapter.send("123", content) + + assert preview.success and final.success + adapter._bot.send_message_draft.assert_awaited_once_with( + chat_id=123, draft_id=7, text=content, + ) + final_kwargs = adapter._bot.send_message.call_args.kwargs + assert final_kwargs["parse_mode"] == ParseMode.MARKDOWN_V2 + assert final_kwargs["text"] == adapter.format_message(content) diff --git a/website/docs/user-guide/messaging/telegram.md b/website/docs/user-guide/messaging/telegram.md index b318eaafb0f01..548b674d868ae 100644 --- a/website/docs/user-guide/messaging/telegram.md +++ b/website/docs/user-guide/messaging/telegram.md @@ -1047,7 +1047,7 @@ To find a topic's `thread_id`, open the topic in Telegram Web or Desktop and loo - **Bot API 9.4 (Feb 2026):** Private Chat Topics — bots can create forum topics in 1-on-1 DM chats via `createForumTopic`. Hermes uses this for two distinct features: operator-curated [Private Chat Topics](#private-chat-topics-bot-api-94) (config-driven, fixed topic list) and user-driven [Multi-session DM mode](#multi-session-dm-mode-topic) (activated by `/topic`, unlimited user-created topics). - **Privacy policy:** Telegram now requires bots to have a privacy policy. Set one via BotFather with `/setprivacy_policy`, or Telegram may auto-generate a placeholder. This is particularly important if your bot is public-facing. -- **Bot API 9.5 (Mar 2026): Native streaming via `sendMessageDraft`.** Hermes supports Telegram's native streaming-draft API as an opt-in transport for private chats. The default remains the legacy `editMessageText` path because draft previews can visibly collapse and re-render on some Telegram clients. +- **Bot API 9.5 (Mar 2026): Native streaming via `sendMessageDraft`.** Hermes supports Telegram's native streaming-draft API for private chats. Some Telegram clients can visibly collapse or re-render draft previews; select `transport: edit` if this happens on your client. ### Streaming transport (`gateway.streaming.transport`) @@ -1055,7 +1055,7 @@ When streaming is enabled (`gateway.streaming.enabled: true`), Hermes picks one | Value | Behaviour | |---|---| -| `auto` (default) | Native draft streaming on supported chats (currently Telegram DMs); legacy edit-based path otherwise. Falls back gracefully if a draft frame fails. | +| `auto` (default) | Native draft streaming on supported chats (currently Telegram DMs); legacy edit-based path otherwise. Falls back if the draft endpoint is unavailable or rejected. | | `draft` | Force native drafts. Logs a downgrade and falls back to edit if the chat doesn't support drafts (e.g. groups/topics). | | `edit` | Legacy progressive `editMessageText` polling for every chat type. | | `off` | Disable streaming entirely (final reply only, no progressive updates). | @@ -1063,23 +1063,27 @@ When streaming is enabled (`gateway.streaming.enabled: true`), Hermes picks one In `~/.hermes/config.yaml`: ```yaml -gateway: - streaming: - enabled: true - transport: auto # auto | draft | edit | off +streaming: + enabled: true + transport: draft # auto | draft | edit | off + edit_interval: 0.8 ``` -**What you'll see in DMs with `edit` (default)** — the gateway sends a normal preview message and progressively updates it via `editMessageText`, avoiding Telegram's draft-preview collapse/rollback effect. +The nested `gateway.streaming` form is also accepted; a top-level `streaming` block takes precedence. After changing streaming or Telegram rendering settings, run `hermes gateway restart` (or send `/restart` as a gateway administrator). This restarts the background program that handles Telegram messages so it reads the new settings; it does not restart your Telegram app. + +**What you'll see in DMs with `edit`** — the gateway sends a normal preview message and progressively updates it via `editMessageText`, avoiding Telegram's draft-preview collapse/rollback effect. + +**What you'll see in DMs with `auto` or `draft`** — Hermes updates the same ephemeral draft with cumulative text. With the default `rich_drafts: false`, previews append raw text without reformatting unfinished Markdown or adding a cursor. The draft preview is bounded by Telegram's text limit; the complete answer still uses the normal MarkdownV2 or opt-in Rich Message renderer, splitting long messages when required, and stays in your chat history. Preview batching does not change the model, prompt, or completed answer. -**What you'll see in DMs with `auto` or `draft`** — Telegram shows an animated draft preview that updates token-by-token. When the reply finishes, it's delivered as a regular message and the draft preview clears naturally on the client. Drafts have no message id, so the final answer is what stays in your chat history. +**Update cadence and animation are separate.** `edit_interval` controls how often Hermes submits a new preview, not how quickly Telegram paints each character. Between updates, incoming deltas are combined so the next preview contains the latest available text rather than replaying a fixed number of characters. Telegram's client controls the animation of newly added text; its smoothness can differ between clients. Hermes also shares Telegram's per-chat draft/typing allowance (20 calls per 5 seconds and 40 per 30 seconds), so lowering `edit_interval` cannot force unlimited updates. See [Telegram's streaming guidance](https://core.telegram.org/api/bots/ai). **What about groups, supergroups, forum topics?** Telegram restricts `sendMessageDraft` to private chats (DMs). The gateway transparently falls back to the edit-based path for everything else — same UX as before. -**What if a draft frame fails?** Any failure (transient network error, server-side rejection, older python-telegram-bot install) flips that response back to the edit-based path for the rest of the stream. The next response gets a fresh attempt. +**What if a draft frame fails?** A locally rate-limited preview is skipped and the next allowed update uses the latest text. Telegram's `RetryAfter` temporarily pauses draft updates without abandoning native streaming. An unavailable endpoint or a rejected draft falls back to editing a regular message for the rest of that response; the next response gets a fresh attempt. The complete final answer is delivered through the existing final-message path even when intermediate previews are skipped. ## Rendering: Rich Messages, Tables and Link Previews -**Rich Messages (Bot API 10.1).** Final replies that contain constructs the legacy MarkdownV2 path degrades — tables, task lists, collapsible `
`, and block math — are sent with Telegram's native [`sendRichMessage`](https://core.telegram.org/bots/api#sendrichmessage) using the agent's **raw markdown**, so they render natively with no client-side flattening. In DMs, the default `rich_drafts: false` keeps the streaming preview plain — it uses Telegram's ephemeral draft transport with legacy rendering (tables and other rich-only constructs stay as raw markdown in the preview) — then persists the completed response with `sendRichMessage`. Setting `rich_drafts: true` makes the live preview use `sendRichMessageDraft` too. Edit-based streams can finalize an existing preview in place through `editMessageText`'s `rich_message` parameter. Ordinary replies (plain prose, bold/italic, simple lists) stay on the MarkdownV2 path for consistent font weight and spacing across clients. +**Rich Messages (Bot API 10.1).** Final replies that contain constructs the legacy MarkdownV2 path degrades — tables, task lists, collapsible `
`, and block math — are sent with Telegram's native [`sendRichMessage`](https://core.telegram.org/bots/api#sendrichmessage) using the agent's **raw markdown**, so they render natively with no client-side flattening. In DMs, the default `rich_drafts: false` keeps the ephemeral streaming preview as raw plain text, then persists the completed response with `sendRichMessage`. Setting `rich_drafts: true` makes the live preview use `sendRichMessageDraft` too. Edit-based streams can finalize an existing preview in place through `editMessageText`'s `rich_message` parameter. Ordinary replies (plain prose, bold/italic, simple lists) stay on the MarkdownV2 path for consistent font weight and spacing across clients. The rich path is skipped automatically when content exceeds the 32,768-character rich text limit, and any rejection from Telegram (unsupported endpoint on an older `python-telegram-bot`, parser error, oversized blocks/columns) **transparently falls back** to the MarkdownV2 path — your message is never lost. Transient/network errors are *not* silently re-sent (no duplicate final message). diff --git a/website/i18n/zh-Hans/docusaurus-plugin-content-docs/current/user-guide/messaging/telegram.md b/website/i18n/zh-Hans/docusaurus-plugin-content-docs/current/user-guide/messaging/telegram.md index 665dd1a6ed823..5caad92c01db6 100644 --- a/website/i18n/zh-Hans/docusaurus-plugin-content-docs/current/user-guide/messaging/telegram.md +++ b/website/i18n/zh-Hans/docusaurus-plugin-content-docs/current/user-guide/messaging/telegram.md @@ -845,7 +845,7 @@ platforms: - **Bot API 9.4(2026 年 2 月):** 私聊话题——机器人可以通过 `createForumTopic` 在一对一私聊中创建论坛话题。Hermes 将此用于两个不同功能:运营商策划的[私聊话题](#private-chat-topics-bot-api-94)(配置驱动,固定话题列表)和用户驱动的[多会话私聊模式](#multi-session-dm-mode-topic)(通过 `/topic` 激活,用户创建的无限话题)。 - **隐私政策:** Telegram 现在要求机器人有隐私政策。通过 BotFather 的 `/setprivacy_policy` 设置,或 Telegram 可能自动生成占位符。如果你的机器人面向公众,这一点尤为重要。 -- **Bot API 9.5(2026 年 3 月):通过 `sendMessageDraft` 实现原生流式传输。** Hermes 支持 Telegram 的原生流式草稿 API,作为私聊的可选传输方式。默认仍使用旧版 `editMessageText` 路径,因为草稿预览在某些 Telegram 客户端上可能出现明显的折叠和重新渲染。 +- **Bot API 9.5(2026 年 3 月):通过 `sendMessageDraft` 实现原生流式传输。** Hermes 支持 Telegram 的私聊原生流式草稿 API。部分 Telegram 客户端可能出现草稿折叠或重新渲染;如果你的客户端出现这种情况,可以选择 `transport: edit`。 ### 流式传输(`gateway.streaming.transport`) @@ -853,7 +853,7 @@ platforms: | 值 | 行为 | |---|---| -| `auto`(默认) | 在支持的聊天(目前为 Telegram 私聊)上使用原生草稿流式传输;否则使用旧版基于编辑的路径。如果草稿帧失败,会优雅回退。 | +| `auto`(默认) | 在支持的聊天(目前为 Telegram 私聊)上使用原生草稿流式传输;否则使用旧版基于编辑的路径。草稿接口不可用或请求被拒绝时会回退。 | | `draft` | 强制使用原生草稿。如果聊天不支持草稿(例如群组/话题),记录降级日志并回退到编辑方式。 | | `edit` | 对所有聊天类型使用旧版渐进式 `editMessageText` 轮询。 | | `off` | 完全禁用流式传输(仅最终回复,无渐进更新)。 | @@ -861,23 +861,27 @@ platforms: 在 `~/.hermes/config.yaml` 中: ```yaml -gateway: - streaming: - enabled: true - transport: auto # auto | draft | edit | off +streaming: + enabled: true + transport: draft # auto | draft | edit | off + edit_interval: 0.8 ``` -**使用 `edit` 传输时私聊中的效果** — gateway 发送一条普通预览消息,并通过 `editMessageText` 渐进更新,避免 Telegram 草稿预览折叠/回滚效果。 +也支持嵌套的 `gateway.streaming` 写法;同时存在时,顶层 `streaming` 优先。修改流式输出或 Telegram 渲染设置后,运行 `hermes gateway restart`,也可以由网关管理员在聊天中发送 `/restart`。这会重启负责收发 Telegram 消息的后台程序,让它读取新设置,不会重启你的 Telegram 客户端。 + +**使用 `edit` 传输时私聊中的效果**——网关发送一条普通预览消息,并通过 `editMessageText` 渐进更新,避免 Telegram 草稿预览折叠或回滚。 + +**使用 `auto` 或 `draft` 时私聊中的效果**——Hermes 用累计文字更新同一条临时草稿。默认的 `rich_drafts: false` 会追加原始纯文本,不会反复格式化尚未写完的 Markdown,也不会添加光标。草稿预览受 Telegram 文字长度上限限制;完整答案仍按正常的 MarkdownV2 或选择启用的富消息格式送达,必要时拆分长消息,并保留在聊天记录中。预览合并不会改变模型、提示词或最终答案。 -**使用 `auto` 或 `draft` 时私聊中的效果** — Telegram 显示逐 token 更新的动画草稿预览。回复完成后,它作为普通消息投递,草稿预览在客户端自然清除。草稿没有消息 ID,因此最终答案才是保留在聊天历史中的内容。 +**更新频率和文字动画是两件事。** `edit_interval` 决定 Hermes 多久提交一次预览,不决定 Telegram 逐字显示的速度。两次更新之间收到的文字会合并,下一次预览直接包含最新内容,不会按固定字数慢慢回放。新增文字的动画由 Telegram 客户端控制,不同客户端的连贯程度可能不同。Hermes 也会遵守每个聊天中草稿和「正在输入」共同使用的额度:每 5 秒 20 次、每 30 秒 40 次。因此,调低 `edit_interval` 不能让更新次数无限增加。参见 [Telegram 流式输出说明](https://core.telegram.org/api/bots/ai)。 **群组、超级群组、论坛话题怎么办?** Telegram 将 `sendMessageDraft` 限制为私聊(私信)。gateway 对其他所有内容透明地回退到基于编辑的路径——与之前的用户体验相同。 -**如果草稿帧失败怎么办?** 任何失败(瞬时网络错误、服务器端拒绝、旧版 python-telegram-bot 安装)都会将该响应的剩余流切换回基于编辑的路径。下一个响应会重新尝试。 +**如果草稿更新失败怎么办?** 本地额度不足时跳过这次预览,下一次允许更新时发送最新文字。Telegram 返回 `RetryAfter` 时会暂时停止草稿更新,随后继续使用原生草稿。接口不可用或草稿请求被拒绝时,这条回复的剩余部分回退为编辑普通消息;下一条回复会重新尝试。即使中间预览被跳过,完整答案仍通过原有的最终消息流程送达。 ## 渲染:富消息、表格和链接预览 -**富消息(Bot API 10.1)。** 最终回复中那些会被旧版 MarkdownV2 路径降级的结构——表格、任务列表、可折叠的 `
` 以及块级数学公式——会通过 Telegram 原生的 [`sendRichMessage`](https://core.telegram.org/bots/api#sendrichmessage) 发送,使用 Agent 的**原始 markdown**,从而原生渲染、无需客户端展平。在流式传输过程中,最终答案通过 `editMessageText` 的 `rich_message` 参数**就地编辑现有预览**来交付——不发第二条消息、不删除,因此一轮结束时不会出现重复投递的闪烁。在私聊中,实时流式预览也使用 `sendRichMessageDraft`,因此动画草稿与最终的富消息保持一致。普通回复(纯文本、粗体/斜体、简单列表)仍走 MarkdownV2 路径,以在各客户端保持一致的字重和间距。 +**富消息(Bot API 10.1)。** 最终回复中那些会被旧版 MarkdownV2 路径降级的结构——表格、任务列表、可折叠的 `
` 以及块级数学公式——会通过 Telegram 原生的 [`sendRichMessage`](https://core.telegram.org/bots/api#sendrichmessage) 发送,使用 Agent 的**原始 Markdown**,从而原生渲染,无需客户端展平。私聊中默认的 `rich_drafts: false` 会让临时预览保持原始纯文本,完成后再通过 `sendRichMessage` 送达完整答案。设置 `rich_drafts: true` 才会同时用 `sendRichMessageDraft` 渲染实时预览。使用普通消息编辑的流式回复,可以通过 `editMessageText` 的 `rich_message` 参数就地送达富消息格式的最终答案。普通回复(纯文本、粗体、斜体、简单列表)仍走 MarkdownV2 路径,以在各客户端保持一致的字重和间距。 当内容超过 32,768 字符的富文本上限时,富消息路径会自动跳过;Telegram 的任何拒绝(较旧 `python-telegram-bot` 不支持该端点、解析错误、块/列过多)都会**透明回退**到 MarkdownV2 路径——消息绝不会丢失。瞬时/网络错误**不会**被静默重发(不会产生重复的最终消息)。 @@ -1238,4 +1242,4 @@ HERMES_TELEGRAM_NOTIFICATIONS=all 切勿公开分享你的机器人 token。如果泄露,请立即通过 BotFather 的 `/revoke` 命令撤销。 -更多详情,请参阅[安全文档](../security.md)。你也可以使用 [DM 配对](./index.md#dm-pairing-alternative-to-allowlists) 进行更动态的用户授权方式。 \ No newline at end of file +更多详情,请参阅[安全文档](../security.md)。你也可以使用 [DM 配对](./index.md#dm-pairing-alternative-to-allowlists) 进行更动态的用户授权方式。