Skip to content

[4.41/TG3] Telegram polling & media: await handlers before advancing offset, media groups, multi-file uploads #227

Description

@patrick-chinchill

Summary

Telegram long polling advances offset before any handler runs, so a handler failure or crash loses the update. Upstream 4.41 waits for dispatched handlers to settle, checkpoints failures in state and retries with backoff. Albums (media groups) coalesce into one Message on webhook and polling paths, and outbound posts can send 2–10 files as one sendMediaGroup.

Upstream changes

  • 8d7ccdb1 feat(telegram): support multiple file and attachment uploads (#605) — chat@4.34.0. Multiple files or attachments go out via sendMediaGroup. It requires 2–10 items and rejects an inline keyboard. Photos and videos may mix; documents and audio only group with their own type. The caption goes on the first item, uploads use attach:// multipart parts, and the call returns the last sent message and caches all of them.
  • 629e6555 fix(telegram): combine incoming media groups (#760) — chat@4.37.0. Webhook album parts are buffered in state under telegram:incoming-media-group:{thread}:{group}, guarded by a state lock (5 s TTL, 50 ms retry, 30 s buffer TTL). After a 1 s settle since the newest part, one message is emitted with: the id and raw of the newest part; the text of the first part that has text; all attachments in message_id order; is_mention if any part is a mention. Slash routing skips album parts.
  • 91683e52 fix(telegram): wait for polling handlers before advancing offset (#942) — chat@4.41.0. processUpdate returns the dispatched tasks. The poller awaits them, grouping by media_group_id, and calls with WebhookOptions(deduplicate=False). Failures persist in a telegram:polling:{scope} checkpoint {offset, pending[{update, receivedAt, attempts, retryAt}]}. The backoff is max(retryDelayMs, 1000) * 2^(attempts-1), capped at 30 s and at least retry_after on rate limits. Album parts wait for the settle window before being processed. The core half (the deduplicate option, processAction/processReaction returning promises, and processMessage returning the rejecting task) lands in [4.41/C1b] Core lifecycle: init retry, dedupe TTL 10min, propagate_handler_errors, webhook dedupe option #191.

Current Python behavior

  • src/chat_sdk/adapters/telegram/adapter.py:973-1026 (polling_loop): :994 sets offset = update.get("update_id", 0) + 1 before :996 self.process_update(update), which is synchronous. Handler outcomes are never observed and no checkpoint exists.
  • adapter.py:1086-1114 (process_update) returns None and discards the task that chat.process_message returns (src/chat_sdk/chat.py:953-988). chat.process_reaction (chat.py:990-1006), process_action (:1008-1024) and process_slash_command (:1164) return None.
  • src/chat_sdk/types.py:1152 WebhookOptions only has wait_until, so there is no deduplicate yet ([4.41/C1b] Core lifecycle: init retry, dedupe TTL 10min, propagate_handler_errors, webhook dedupe option #191).
  • grep -rn 'media_group' src/chat_sdk/adapters/telegram finds nothing, so every album part is dispatched as a separate message (adapter.py:1116-1135).
  • adapter.py:1407-1425 (post_message) raises "supports a single file upload per message" / "single attachment upload per message". grep -rn sendMediaGroup src/chat_sdk/adapters/telegram finds nothing.

Scope

  • process_update returns list[asyncio.Task] of every dispatched task (message, slash, action, reaction). The webhook path keeps ignoring the list.
  • Webhook album path: handle_incoming_message_update and the slash gate skip parts that have a media_group_id. Add async _process_incoming_media_group(msg, thread_id, options) following upstream, using state.acquire_lock / release_lock in try/finally and state.get / set / delete. Its task is returned from process_update and handed to options.wait_until.
  • Polling rewrite (polling_loop), following upstream pollingLoop + processPollingUpdates:
    • Key: f"{self._name}:polling:{scope}", with the scope from _ensure_bot_identity ([4.41/TG0] Telegram: require webhook verification by default, dedupe repeated updates #224; if TG0 has not landed, introduce the same helper here).
    • Group by (chat.id, message_thread_id, media_group_id).
    • Await each group's tasks with asyncio.gather(..., return_exceptions=True).
    • Persist pending failures and retry_at, and advance offset only past acknowledged updates.
    • Use a shortened getUpdates timeout while an album is settling.
    • Pass WebhookOptions(deduplicate=False) on polled dispatch.
    • Lazily retry _ensure_bot_identity once updates arrive (from 26a06ca5, so the identity recovers after a failed startup getMe).
  • post_message: remove the single-file and single-attachment errors (keep "mixing file uploads and attachments" and "Attachment upload payload is empty"). Add _send_document_media_group and _send_attachment_media_group, plus _validate_media_group_length (2–10) and _validate_attachment_media_group_types. Reject reply_markup. Build the aiohttp.FormData with a media JSON string and media{i} parts. The caption and parse_mode go on item 0, with the same MarkdownV2 → plain fallback as send_document. Cache every returned message and return the last.
  • Constants from upstream: TELEGRAM_INCOMING_MEDIA_GROUP_BUFFER_TTL_MS = 30_000, _LOCK_TTL_MS = 5_000, _RETRY_MS = 50, _SETTLE_MS = 1_000, TELEGRAM_MEDIA_GROUP_MIN/MAX = 2/10.

Out of scope

Porting notes

  • AbortController → task cancellation (stop_polling already cancels _polling_task, adapter.py:946-960). Let CancelledError propagate from every await (never catch BaseException); write the checkpoint with a single state.set.
  • Settle windows use an injectable monotonic clock and sleep. Tests use a fake clock (receivedAt/retryAt are epoch-ms in state, so store time.time()*1000 there but compare with the same clock).
  • Promise.allSettled → asyncio.gather(*tasks, return_exceptions=True), then re-raise the first exception so the whole group is retried, as upstream does.
  • Checkpoint values must be JSON-serialisable for the Redis and Postgres backends: store raw update dicts, not Message objects.
  • retry_after comes from AdapterRateLimitError.retry_after and may be None; use is not None.
  • aiohttp FormData import stays lazy; the media field is a JSON string.
  • Design decision: upstream deduplicates polled updates by the checkpoint, not by core dedupe. Keep deduplicate=False on polled dispatch; otherwise a retried update is dropped by the 10-minute core dedupe window.

Tests

packages/adapter-telegram/src/index.test.ts is not fidelity-mapped (see #78). Port into tests/test_telegram_webhook.py / tests/test_telegram_api.py:

  • "waits for polled message processing and saves failures before acknowledging updates"
  • "coalesces polled media groups before acknowledging updates"
  • "starts polling, advances offset, and stops cleanly"
  • "combines an incoming media group into one ordered message"
  • "posts multiple files as a Telegram media group"
  • "posts and normalizes mixed image and video attachments as a Telegram media group"
  • "rejects incompatible Telegram media group attachment types"
  • "rejects Telegram media groups with more than 10 files"
  • "rejects mixed file uploads and attachments"
  • "rejects attachments without upload data"

Optionally mine packages/integration-tests/src/polling.test.ts (new in 91683e52) for multi-restart checkpoint cases. The core chat.test.ts case "lets transports own deduplication when retrying admission" belongs to #191.

Python-specific: stop_polling() during a settle returns promptly with a consistent checkpoint; a handler exception keeps the update pending with backoff scheduled. AsyncMock handlers, in-memory state, no sleeps.

Acceptance criteria

  • A handler failure on a polled update does not advance past it without a persisted retry entry, and a restart resumes from the checkpoint.
  • An album of N parts reaches handlers as one Message with N attachments (webhook and polling).
  • post_message with 2–10 files or attachments issues one sendMediaGroup.
  • The full validation command from CLAUDE.md passes.
  • docs/UPSTREAM_SYNC.md is updated for any divergence (for example, checkpoint clock semantics).
  • A CHANGELOG entry under "Unreleased (4.41 wave)" calls out that polling is now at-least-once with retries, albums arrive as one message (handlers stop seeing N messages), and multi-file posts are supported.

Dependencies

Blocked by #191. Soft-depends on #224 (_ensure_bot_identity / scope) and #225 (typing and allowlist helpers).

Verify first

Metadata

  • Effort: L (~1,000–1,400 LOC including tests; split outbound sendMediaGroup into its own PR if it runs over)
  • Consumer impact: low (Telegram only). None for Slack/Teams.
  • Suggested branch: sync/4.41-tg3

Part of #184.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions