You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
[4.41/TG3] Telegram polling & media: await handlers before advancing offset, media groups, multi-file uploads #227
Telegram long polling advances offsetbefore 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) + 1before:996self.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.
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:
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.
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
That [4.41/C1b] Core lifecycle: init retry, dedupe TTL 10min, propagate_handler_errors, webhook dedupe option #191 landed all three core pieces: WebhookOptions.deduplicate; process_action, process_reaction and process_slash_command returning their task; and the process_message task raising the handler error when awaited. If any piece is missing, add the minimal core change in this PR with a tests/test_chat_faithful.py case and note it on C1b.
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.
Summary
Telegram long polling advances
offsetbefore 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 oneMessageon webhook and polling paths, and outbound posts can send 2–10 files as onesendMediaGroup.Upstream changes
8d7ccdb1feat(telegram): support multiple file and attachment uploads (#605) — chat@4.34.0. Multiplefilesorattachmentsgo out viasendMediaGroup. 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 useattach://multipart parts, and the call returns the last sent message and caches all of them.629e6555fix(telegram): combine incoming media groups (#760) — chat@4.37.0. Webhook album parts are buffered in state undertelegram: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 inmessage_idorder;is_mentionif any part is a mention. Slash routing skips album parts.91683e52fix(telegram): wait for polling handlers before advancing offset (#942) — chat@4.41.0.processUpdatereturns the dispatched tasks. The poller awaits them, grouping bymedia_group_id, and calls withWebhookOptions(deduplicate=False). Failures persist in atelegram:polling:{scope}checkpoint{offset, pending[{update, receivedAt, attempts, retryAt}]}. The backoff ismax(retryDelayMs, 1000) * 2^(attempts-1), capped at 30 s and at leastretry_afteron rate limits. Album parts wait for the settle window before being processed. The core half (thededuplicateoption,processAction/processReactionreturning promises, andprocessMessagereturning 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)::994setsoffset = update.get("update_id", 0) + 1before:996self.process_update(update), which is synchronous. Handler outcomes are never observed and no checkpoint exists.adapter.py:1086-1114(process_update) returnsNoneand discards the task thatchat.process_messagereturns (src/chat_sdk/chat.py:953-988).chat.process_reaction(chat.py:990-1006),process_action(:1008-1024) andprocess_slash_command(:1164) returnNone.src/chat_sdk/types.py:1152WebhookOptionsonly haswait_until, so there is nodeduplicateyet ([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/telegramfinds 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/telegramfinds nothing.Scope
process_updatereturnslist[asyncio.Task]of every dispatched task (message, slash, action, reaction). The webhook path keeps ignoring the list.handle_incoming_message_updateand the slash gate skip parts that have amedia_group_id. Addasync _process_incoming_media_group(msg, thread_id, options)following upstream, usingstate.acquire_lock/release_lockintry/finallyandstate.get/set/delete. Its task is returned fromprocess_updateand handed tooptions.wait_until.polling_loop), following upstreampollingLoop+processPollingUpdates: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).(chat.id, message_thread_id, media_group_id).asyncio.gather(..., return_exceptions=True).retry_at, and advanceoffsetonly past acknowledged updates.getUpdatestimeout while an album is settling.WebhookOptions(deduplicate=False)on polled dispatch._ensure_bot_identityonce updates arrive (from26a06ca5, so the identity recovers after a failed startupgetMe).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_groupand_send_attachment_media_group, plus_validate_media_group_length(2–10) and_validate_attachment_media_group_types. Rejectreply_markup. Build theaiohttp.FormDatawith amediaJSON string andmedia{i}parts. The caption andparse_modego on item 0, with the same MarkdownV2 → plain fallback assend_document. Cache every returned message and return the last.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
deduplicate/ task-returningprocess_*/ propagated handler errors: [4.41/C1b] Core lifecycle: init retry, dedupe TTL 10min, propagate_handler_errors, webhook dedupe option #191.reply_toon the combined message andreply_parameterson media groups: [4.41/TG4] Telegram replies: replied-to context, reply-to-bot as mention, native replies, portable file data #228. Carryreply_tothrough only ifMessage.reply_toalready exists from [4.41/C2a] Core mentions & message model: tri-state is_mention, mention regex, Author.email/is_system, Message.reply_to #192.business_connection_idsegment of the group key: [4.41/DEF1] Deferred upstream 4.32–4.41 features (demand-gated backlog) #189.Porting notes
stop_pollingalready cancels_polling_task,adapter.py:946-960). LetCancelledErrorpropagate from everyawait(never catchBaseException); write the checkpoint with a singlestate.set.receivedAt/retryAtare epoch-ms in state, so storetime.time()*1000there 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.Messageobjects.retry_aftercomes fromAdapterRateLimitError.retry_afterand may beNone; useis not None.FormDataimport stays lazy; themediafield is a JSON string.deduplicate=Falseon polled dispatch; otherwise a retried update is dropped by the 10-minute core dedupe window.Tests
packages/adapter-telegram/src/index.test.tsis not fidelity-mapped (see #78). Port intotests/test_telegram_webhook.py/tests/test_telegram_api.py:Optionally mine
packages/integration-tests/src/polling.test.ts(new in91683e52) for multi-restart checkpoint cases. The corechat.test.tscase "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 updatependingwith backoff scheduled.AsyncMockhandlers, in-memory state, no sleeps.Acceptance criteria
Messagewith N attachments (webhook and polling).post_messagewith 2–10 files or attachments issues onesendMediaGroup.docs/UPSTREAM_SYNC.mdis updated for any divergence (for example, checkpoint clock semantics).Dependencies
Blocked by #191. Soft-depends on #224 (
_ensure_bot_identity/ scope) and #225 (typing and allowlist helpers).Verify first
WebhookOptions.deduplicate;process_action,process_reactionandprocess_slash_commandreturning their task; and theprocess_messagetask raising the handler error when awaited. If any piece is missing, add the minimal core change in this PR with atests/test_chat_faithful.pycase and note it on C1b.Metadata
sendMediaGroupinto its own PR if it runs over)sync/4.41-tg3Part of #184.