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
Port four Chat lifecycle fixes (incl. the core half of Telegram's polling fix): init retries only after a failed state connection; dedupe TTL 5 → 10 min; opt-in WebhookOptions.propagate_handler_errors (default: wait_until gets a task that swallows handler errors); WebhookOptions.deduplicate=False bypasses chat-level dedupe. Python diverges on all four; the wait_until change is consumer-visible.
Upstream changes
0b63791b fix(slack): process Socket Mode retry envelopes instead of dropping them (#667) — chat@4.33.0 — core half only: DEDUPE_TTL_MS 5 → 10 min so the entry outlives Slack's ~+5 min Events API retry.
f233ffe8 fix(chat): retry initialization after a failed attempt (#924) — chat@4.41.0 — ensureInitialized clears initPromise only when state.connect() rejects and only if still the same attempt (this.initPromise === attempt); adapter initialize rejections stay cached until shutdown().
c21ccbc0 fix(chat): propagate handler errors through waitUntil (#943) — chat@4.41.0 — adds WebhookOptions.propagateHandlerErrors; processMessage/processAction/processSlashCommand call waitUntil(propagate ? task : tracked) where tracked logs and swallows. Teams action path now spreads ...webhookOptions.
91683e52 fix(telegram): wait for polling handlers before advancing offset (#942) — chat@4.41.0 — core half only: adds WebhookOptions.deduplicate; handleIncomingMessage splits into routeIncomingMessage(…, deduplicate = true) (self-filter, then dedupe) and dispatchIncomingMessage (history, lock, strategy); processReaction returns its task and always hands waitUntil the tracked version; processAction/processSlashCommand return the raw task.
Current Python behavior
Dedupe TTL:src/chat_sdk/chat.py:83DEDUPE_TTL_MS = 5 * 60 * 1000; types.py:1706dedupe_ttl_ms: int = 300000 always wins, so the constant is dead. chat.py:305config.dedupe_ttl_ms or DEDUPE_TTL_MS turns 0 into the default. tests/test_chat_faithful.py:460test_should_use_default_dedupe_ttl_of_5_minutes never asserts the TTL.
Init retry:chat.py:542-555_ensure_initialized awaits the attempt under _init_lock and resets _init_promise = None on any exception, incl. adapter failures (_do_initialize, :557-572). A retry re-runs state.connect() and every adapter initialize(); asyncio.gather (:567) doesn't cancel siblings, and Teams re-registers handlers on each initialize (adapters/teams/adapter.py:408).
wait_until:chat.py:986-987, :1005-1006, :1023-1024, :1179-1180 pass the raw task, as do :1120-1121 (modal-submit callback task) and the lifecycle process_* at :1161, :1200, :1221, :1242, :1263. Awaiting it re-raises handler errors — effectively upstream's opt-in propagate=True — although the :966 docstring says "swallowed". process_reaction/process_action/process_slash_command return None (:994, :1012, :1168), as do their ChatInstance stubs (types.py:1762-1764). types.py:1151-1155WebhookOptions has only wait_until; grep -rn 'propagate_handler_errors\|deduplicate' src/chat_sdk is empty.
Teams streaming gate (adapters/teams/adapter.py:821-891, docs/UPSTREAM_SYNC.md:655): hooks add_done_callbackonly if handed an asyncio.Task (:828); anything else resolves processing_done immediately, closing the streamer early. It also builds WebhookOptions(wait_until=_chained_wait_until) from scratch (:875), dropping any new caller option — upstream spreads ...baseOptions.
Scope
Dedupe TTL: DEDUPE_TTL_MS = 10 * 60 * 1000; make ChatConfig.dedupe_ttl_ms: int | None = None resolved with is not None (recommended), or default 600_000 — change both constant and default; update the types.py:1705 comment.
Init retry: create the attempt under _init_lock, await it outside the lock (concurrent callers share one task). The attempt awaits state.connect(); on failure clear _init_promise only if self._init_promise is attempt, re-raise. Adapter failures keep the failed task cached (re-raised until shutdown() clears it).
_tracked(task) -> asyncio.Task that awaits the handler task and swallows its exception (the done-callback already logs). wait_until(task if options.propagate_handler_errors else _tracked(task)) for message/action/slash-command; _tracked(task) unconditionally for reaction and every lifecycle process_* that hands a task to wait_until (upstream's lifecycle tasks are already .catched).
process_reaction/process_action/process_slash_command return the raw task (None without a loop); update ChatInstance return types to asyncio.Task[None] | None.
Split handle_incoming_message into _route_incoming_message(adapter, thread_id, message, deduplicate=True) and _dispatch_incoming_message (history, lock key, strategy); keep handle_incoming_message public. process_message calls _route_incoming_message(..., deduplicate=False) when options.deduplicate is False.
Teams _handle_message_activity: build chained options with dataclasses.replace(options or WebhookOptions(), wait_until=_chained_wait_until) so propagate_handler_errors/deduplicate survive; re-verify the processing_done gate with the wrapper Task.
Wrapper must be a Task (a dropped coroutine warns "never awaited" and fails the Teams gate): create via _create_task(..., self._active_tasks); await asyncio.shield(task) so cancelling the wrapper spares the handler; treat an inner cancelled task as completion, re-raise a CancelledError aimed at the wrapper.
Default flip: hosts that awaited the wait_until awaitable to observe errors stop seeing them unless they pass WebhookOptions(propagate_handler_errors=True). The directly returned task still raises.
Init retry — recommended default: adopt upstream. Python's broad retry re-inits running adapters (double-registers Teams handlers, can restart Slack Socket Mode). Trade-off: a transient adapter-init failure wedges the instance until shutdown() (log at error, note in CHANGELOG). If the maintainer keeps the broad retry, record a divergence and baseline-skip the 6 tests below.
Shutdown during init: keep shutdown() clearing _init_promise; the identity check stops an old failing connect from clearing a newer attempt.
Tests
From packages/chat/src/chat.test.ts (fidelity-mapped → tests/test_chat_faithful.py):
describe("Chat initialization retry (#922)"): "does not restart an initialized adapter when another adapter fails"; "does not overlap adapter initialization after a sibling fails"; "recovers through a webhook after repeated state connection failures"; "keeps a newer attempt when a pre-shutdown state connection rejects"; "retries initialization after a failed attempt once the state recovers"; "still shares one attempt between concurrent callers, including a failing one".
"should use default dedupe TTL of 10 minutes" — renametest_should_use_default_dedupe_ttl_of_5_minutes and assert set_if_not_exists is called with 600_000 via an AsyncMock spy (MockStateAdapter.set_if_not_exists doesn't record TTL).
"should optionally propagate handler errors through waitUntil" (message/action/slash-command × propagate False/True: direct result, background outcome, error log).
"lets transports own deduplication when retrying admission".
Python-specific:
ChatConfig(dedupe_ttl_ms=0) is honoured.
Teams DM streaming gate still waits for handler completion when wait_until receives the wrapper; propagate_handler_errors=True passed to the Teams adapter reaches Chat.process_message.
Cancelling the wrapper leaves the handler running; shutdown() cancels both.
Acceptance criteria
Full validation command from CLAUDE.md passes.
docs/UPSTREAM_SYNC.md: init-retry decision recorded; Teams processing_done row (:655) updated (wrapper Task, options spread); Teams dialog-open N/A row added.
CHANGELOG under "Unreleased (4.41 wave)", consumer-visible: wait_until now swallows by default + new opt-in; the three process_* methods return tasks; 10-min dedupe TTL; adapter-init failure cached until shutdown().
Teams DM native streaming proven unchanged by an adapter-level test.
Consumer impact: high — wait_until semantics change for every host, and the Teams DM streaming gate depends on the object wait_until receives. Slack/Teams deployments that awaited wait_until to surface errors must opt in.
Summary
Port four
Chatlifecycle fixes (incl. the core half of Telegram's polling fix): init retries only after a failed state connection; dedupe TTL 5 → 10 min; opt-inWebhookOptions.propagate_handler_errors(default:wait_untilgets a task that swallows handler errors);WebhookOptions.deduplicate=Falsebypasses chat-level dedupe. Python diverges on all four; thewait_untilchange is consumer-visible.Upstream changes
0b63791bfix(slack): process Socket Mode retry envelopes instead of dropping them (#667) — chat@4.33.0 — core half only:DEDUPE_TTL_MS5 → 10 min so the entry outlives Slack's ~+5 min Events API retry.f233ffe8fix(chat): retry initialization after a failed attempt (#924) — chat@4.41.0 —ensureInitializedclearsinitPromiseonly whenstate.connect()rejects and only if still the same attempt (this.initPromise === attempt); adapterinitializerejections stay cached untilshutdown().c21ccbc0fix(chat): propagate handler errors through waitUntil (#943) — chat@4.41.0 — addsWebhookOptions.propagateHandlerErrors;processMessage/processAction/processSlashCommandcallwaitUntil(propagate ? task : tracked)wheretrackedlogs and swallows. Teams action path now spreads...webhookOptions.91683e52fix(telegram): wait for polling handlers before advancing offset (#942) — chat@4.41.0 — core half only: addsWebhookOptions.deduplicate;handleIncomingMessagesplits intorouteIncomingMessage(…, deduplicate = true)(self-filter, then dedupe) anddispatchIncomingMessage(history, lock, strategy);processReactionreturns its task and always handswaitUntilthetrackedversion;processAction/processSlashCommandreturn the raw task.Current Python behavior
src/chat_sdk/chat.py:83DEDUPE_TTL_MS = 5 * 60 * 1000;types.py:1706dedupe_ttl_ms: int = 300000always wins, so the constant is dead.chat.py:305config.dedupe_ttl_ms or DEDUPE_TTL_MSturns0into the default.tests/test_chat_faithful.py:460test_should_use_default_dedupe_ttl_of_5_minutesnever asserts the TTL.chat.py:542-555_ensure_initializedawaits the attempt under_init_lockand resets_init_promise = Noneon any exception, incl. adapter failures (_do_initialize,:557-572). A retry re-runsstate.connect()and every adapterinitialize();asyncio.gather(:567) doesn't cancel siblings, and Teams re-registers handlers on eachinitialize(adapters/teams/adapter.py:408).wait_until:chat.py:986-987,:1005-1006,:1023-1024,:1179-1180pass the raw task, as do:1120-1121(modal-submit callback task) and the lifecycleprocess_*at:1161,:1200,:1221,:1242,:1263. Awaiting it re-raises handler errors — effectively upstream's opt-inpropagate=True— although the:966docstring says "swallowed".process_reaction/process_action/process_slash_commandreturnNone(:994,:1012,:1168), as do theirChatInstancestubs (types.py:1762-1764).types.py:1151-1155WebhookOptionshas onlywait_until;grep -rn 'propagate_handler_errors\|deduplicate' src/chat_sdkis empty.chat.py:1933-1992handle_incoming_messagealways dedupes (set_if_not_exists,:1961).adapters/teams/adapter.py:821-891,docs/UPSTREAM_SYNC.md:655): hooksadd_done_callbackonly if handed anasyncio.Task(:828); anything else resolvesprocessing_doneimmediately, closing the streamer early. It also buildsWebhookOptions(wait_until=_chained_wait_until)from scratch (:875), dropping any new caller option — upstream spreads...baseOptions.Scope
DEDUPE_TTL_MS = 10 * 60 * 1000; makeChatConfig.dedupe_ttl_ms: int | None = Noneresolved withis not None(recommended), or default600_000— change both constant and default; update thetypes.py:1705comment._init_lock, await it outside the lock (concurrent callers share one task). The attempt awaitsstate.connect(); on failure clear_init_promiseonlyif self._init_promise is attempt, re-raise. Adapter failures keep the failed task cached (re-raised untilshutdown()clears it).WebhookOptions: addpropagate_handler_errors: bool = False,deduplicate: bool | None = None._tracked(task) -> asyncio.Taskthat awaits the handler task and swallows its exception (the done-callback already logs).wait_until(task if options.propagate_handler_errors else _tracked(task))for message/action/slash-command;_tracked(task)unconditionally for reaction and every lifecycleprocess_*that hands a task towait_until(upstream's lifecycle tasks are already.catched).process_reaction/process_action/process_slash_commandreturn the raw task (Nonewithout a loop); updateChatInstancereturn types toasyncio.Task[None] | None.handle_incoming_messageinto_route_incoming_message(adapter, thread_id, message, deduplicate=True)and_dispatch_incoming_message(history, lock key, strategy); keephandle_incoming_messagepublic.process_messagecalls_route_incoming_message(..., deduplicate=False)whenoptions.deduplicate is False._handle_message_activity: build chained options withdataclasses.replace(options or WebhookOptions(), wait_until=_chained_wait_until)sopropagate_handler_errors/deduplicatesurvive; re-verify theprocessing_donegate with the wrapper Task.Out of scope
deduplicate=False) → [4.41/TG3] Telegram polling & media: await handlers before advancing offset, media groups, multi-file uploads #227.handleDialogOpenparts of c21ccbc0/91683e52: N/A (dialog-open inbound not ported) — record in UPSTREAM_SYNC.md.deduplicate=Falsedispatch → [4.41/C3] Conversation context + AI tool scoping (read & write guards, strict_scope) #195.Porting notes
_create_task(..., self._active_tasks); awaitasyncio.shield(task)so cancelling the wrapper spares the handler; treat an inner cancelled task as completion, re-raise aCancelledErroraimed at the wrapper.wait_untilawaitable to observe errors stop seeing them unless they passWebhookOptions(propagate_handler_errors=True). The directly returned task still raises.options.deduplicate is False(None= dedupe); replacededupe_ttl_ms or DEDUPE_TTL_MSwithis not None(hazard security: Fix all critical and high findings from security audit #1).shutdown()(log at error, note in CHANGELOG). If the maintainer keeps the broad retry, record a divergence and baseline-skip the 6 tests below.shutdown()clearing_init_promise; the identity check stops an old failing connect from clearing a newer attempt.Tests
From
packages/chat/src/chat.test.ts(fidelity-mapped →tests/test_chat_faithful.py):describe("Chat initialization retry (#922)"): "does not restart an initialized adapter when another adapter fails"; "does not overlap adapter initialization after a sibling fails"; "recovers through a webhook after repeated state connection failures"; "keeps a newer attempt when a pre-shutdown state connection rejects"; "retries initialization after a failed attempt once the state recovers"; "still shares one attempt between concurrent callers, including a failing one".test_should_use_default_dedupe_ttl_of_5_minutesand assertset_if_not_existsis called with600_000via an AsyncMock spy (MockStateAdapter.set_if_not_existsdoesn't record TTL).Python-specific:
ChatConfig(dedupe_ttl_ms=0)is honoured.wait_untilreceives the wrapper;propagate_handler_errors=Truepassed to the Teams adapter reachesChat.process_message.shutdown()cancels both.Acceptance criteria
docs/UPSTREAM_SYNC.md: init-retry decision recorded; Teamsprocessing_donerow (:655) updated (wrapper Task, options spread); Teams dialog-open N/A row added.wait_untilnow swallows by default + new opt-in; the threeprocess_*methods return tasks; 10-min dedupe TTL; adapter-init failure cached untilshutdown().Dependencies
Blocked by #190. Blocks #192, #227, #203.
Metadata
wait_untilsemantics change for every host, and the Teams DM streaming gate depends on the objectwait_untilreceives. Slack/Teams deployments that awaitedwait_untilto surface errors must opt in.sync/4.41-c1bPart of #184.