Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion gateway/stream_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down
29 changes: 27 additions & 2 deletions gateway/stream_consumer_transport.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()).
Expand All @@ -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
Expand All @@ -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.
Expand Down Expand Up @@ -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).
Expand Down
87 changes: 13 additions & 74 deletions plugins/platforms/telegram/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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) ---------------------------------------------
Expand Down Expand Up @@ -5561,30 +5498,32 @@ 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))
await _action(message_thread_id=message_thread_id)
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]:
Expand Down
Loading