Skip to content
Merged
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
116 changes: 16 additions & 100 deletions plugins/platforms/telegram/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -412,7 +412,6 @@ def __init__(self, config: PlatformConfig):
)
self._pending_text_batches: Dict[str, MessageEvent] = {}
self._pending_text_batch_tasks: Dict[str, asyncio.Task] = {}
self._drop_delayed_deliveries = False
self._polling_error_task: Optional[asyncio.Task] = None
self._polling_conflict_count: int = 0
self._polling_network_error_count: int = 0
Expand Down Expand Up @@ -500,27 +499,6 @@ def __init__(self, config: PlatformConfig):
# same key edit the same message instead of appending new ones (#30045).
self._status_message_ids: Dict[tuple, str] = {}

def _mark_connected(self) -> None:
self._drop_delayed_deliveries = False
super()._mark_connected()

def _mark_disconnected(self) -> None:
self._drop_delayed_deliveries = True
super()._mark_disconnected()

def _set_fatal_error(self, code: str, message: str, *, retryable: bool) -> None:
self._drop_delayed_deliveries = True
super()._set_fatal_error(code, message, retryable=retryable)

def _should_drop_delayed_delivery(self) -> bool:
"""True once teardown/fatal-error started — delayed flushes must drop.

Buffered text/photo/media-group flushes sit behind an asyncio.sleep().
If disconnect wins the race, dispatching them spawns an agent on a
torn-down session, producing stale/duplicate deliveries.
"""
return bool(getattr(self, "_drop_delayed_deliveries", False))

def _notification_kwargs(
self, metadata: Optional[Dict[str, Any]]
) -> Dict[str, Any]:
Expand Down Expand Up @@ -2937,60 +2915,8 @@ async def _set_status_indicator(self, online: bool) -> None:
self.name, text, e,
)

async def _cancel_pending_delivery_tasks(self) -> None:
"""Cancel every delayed-delivery task family before disconnect completes.

Covers media-group, photo-batch and text-batch flush tasks plus the
polling-error recovery task. Each sits behind an ``asyncio.sleep()``;
if teardown leaves them running they dispatch ``handle_message`` into a
torn-down session. Skips the current task so the coroutine driving
teardown does not cancel itself.
"""
current_task = asyncio.current_task()
pending_tasks: list[asyncio.Task] = []
awaitable_tasks: list[asyncio.Task] = []
seen: set[int] = set()

def collect(task: Optional[asyncio.Task]) -> None:
if not task or task.done() or task is current_task:
return
marker = id(task)
if marker in seen:
return
seen.add(marker)
pending_tasks.append(task)
if asyncio.isfuture(task) or asyncio.iscoroutine(task):
awaitable_tasks.append(task)

for task in list(self._media_group_tasks.values()):
collect(task)
for task in list(self._pending_photo_batch_tasks.values()):
collect(task)
for task in list(self._pending_text_batch_tasks.values()):
collect(task)
collect(self._polling_error_task)

for task in pending_tasks:
task.cancel()
if awaitable_tasks:
await asyncio.gather(*awaitable_tasks, return_exceptions=True)

self._media_group_tasks.clear()
self._media_group_events.clear()
self._pending_photo_batch_tasks.clear()
self._pending_photo_batches.clear()
self._pending_text_batch_tasks.clear()
self._pending_text_batches.clear()
if self._polling_error_task is not current_task:
self._polling_error_task = None

async def disconnect(self) -> None:
"""Stop polling/webhook, cancel pending delayed deliveries, and disconnect."""
# Mark disconnected first so the drop guard short-circuits any flush
# that wins the race against teardown and prevents new delayed tasks
# from being scheduled by late update handlers.
self._mark_disconnected()

"""Stop polling/webhook, cancel pending album flushes, and disconnect."""
# Cancel the heartbeat before tearing down the app so the probe task
# cannot fire get_me() into a half-shutdown bot client.
if self._polling_heartbeat_task and not self._polling_heartbeat_task.done():
Expand All @@ -3011,7 +2937,13 @@ async def disconnect(self) -> None:
except Exception:
pass

await self._cancel_pending_delivery_tasks()
pending_media_group_tasks = list(self._media_group_tasks.values())
for task in pending_media_group_tasks:
task.cancel()
if pending_media_group_tasks:
await asyncio.gather(*pending_media_group_tasks, return_exceptions=True)
self._media_group_tasks.clear()
self._media_group_events.clear()

if self._app:
try:
Expand All @@ -3025,6 +2957,13 @@ async def disconnect(self) -> None:
logger.warning("[%s] Error during Telegram disconnect: %s", self.name, e, exc_info=True)
self._release_platform_lock()

for task in self._pending_photo_batch_tasks.values():
if task and not task.done():
task.cancel()
self._pending_photo_batch_tasks.clear()
self._pending_photo_batches.clear()
Comment on lines 2940 to +2964

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 Text-batch flush tasks not cancelled during disconnect — can dispatch messages into torn-down session (bug)

The PR refactored disconnect() in plugins/platforms/telegram/adapter.py, removing the centralized _cancel_pending_delivery_tasks() method and replacing it with inline cancellation. The inline code at lines 2940–2964 only handles _media_group_tasks and _pending_photo_batch_tasks. _pending_text_batch_tasks receives no cancellation, awaiting, or clearing. Likewise, _pending_text_batches is never cleared. The _should_drop_delayed_delivery() guard that previously short-circuited _flush_text_batch and _enqueue_text_event was also removed entirely.

A _flush_text_batch task (line 6901) sleeping behind await asyncio.sleep(delay) at line 6932 can wake up after disconnect() has shut down the app (line 2955), released the platform lock (line 2958), and set self._app = None / self._bot = None (lines 2967–2968). It then pops its event from _pending_text_batches and calls self.handle_message(event) at line 6940, spawning agent logic on a torn-down adapter. The old code's _should_drop_delayed_delivery docstring explicitly warned: "If disconnect wins the race, dispatching them spawns an agent on a torn-down session, producing stale/duplicate deliveries."

💡 Suggestion: Add text-batch task cancellation inside disconnect() before or after the photo-batch cancellation block. Either restore the _should_drop_delayed_delivery() guard in _flush_text_batch as defense-in-depth, or ensure text-batch tasks are always cancelled and their dicts cleared.

📋 Prompt for AI Agents

In plugins/platforms/telegram/adapter.py, in the disconnect() method, add text-batch task cleanup after the photo-batch cancellation block (around line 2964):

# Cancel text-batch flush tasks to prevent stale deliveries
for task in self._pending_text_batch_tasks.values():
    if task and not task.done():
        task.cancel()
self._pending_text_batch_tasks.clear()
self._pending_text_batches.clear()


self._mark_disconnected()
self._app = None
self._bot = None
logger.info("[%s] Disconnected from Telegram", self.name)
Expand Down Expand Up @@ -6935,10 +6874,6 @@ def _enqueue_text_event(self, event: MessageEvent) -> None:
concatenates them and waits for a short quiet period before
dispatching the combined message.
"""
if self._should_drop_delayed_delivery():
logger.debug("[Telegram] Dropping text batch enqueue after disconnect started")
return

key = self._text_batch_key(event)
existing = self._pending_text_batches.get(key)
chunk_len = len(event.text or "")
Expand Down Expand Up @@ -6998,9 +6933,6 @@ async def _flush_text_batch(self, key: str) -> None:
event = self._pending_text_batches.pop(key, None)
if not event:
return
if self._should_drop_delayed_delivery():
logger.debug("[Telegram] Dropping text batch flush after disconnect started")
return
logger.info(
"[Telegram] Flushing text batch %s (%d chars)",
key, len(event.text or ""),
Expand Down Expand Up @@ -7035,9 +6967,6 @@ async def _flush_photo_batch(self, batch_key: str) -> None:
event = self._pending_photo_batches.pop(batch_key, None)
if not event:
return
if self._should_drop_delayed_delivery():
logger.debug("[Telegram] Dropping photo batch flush after disconnect started")
return
logger.info("[Telegram] Flushing photo batch %s with %d image(s)", batch_key, len(event.media_urls))
await self.handle_message(event)
finally:
Expand All @@ -7046,10 +6975,6 @@ async def _flush_photo_batch(self, batch_key: str) -> None:

def _enqueue_photo_event(self, batch_key: str, event: MessageEvent) -> None:
"""Merge photo events into a pending batch and schedule flush."""
if self._should_drop_delayed_delivery():
logger.debug("[Telegram] Dropping photo batch enqueue after disconnect started")
return

existing = self._pending_photo_batches.get(batch_key)
if existing is None:
self._pending_photo_batches[batch_key] = event
Expand Down Expand Up @@ -7354,10 +7279,6 @@ async def _queue_media_group_event(self, media_group_id: str, event: MessageEven
new user message and interrupts the first. We debounce briefly and merge the
attachments into a single MessageEvent.
"""
if self._should_drop_delayed_delivery():
logger.debug("[Telegram] Dropping media group enqueue after disconnect started")
return

existing = self._media_group_events.get(media_group_id)
if existing is None:
self._media_group_events[media_group_id] = event
Expand All @@ -7376,20 +7297,15 @@ async def _queue_media_group_event(self, media_group_id: str, event: MessageEven
)

async def _flush_media_group_event(self, media_group_id: str) -> None:
current_task = asyncio.current_task()
try:
await asyncio.sleep(self.MEDIA_GROUP_WAIT_SECONDS)
event = self._media_group_events.pop(media_group_id, None)
if event is not None:
if self._should_drop_delayed_delivery():
logger.debug("[Telegram] Dropping media group flush after disconnect started")
return
await self.handle_message(event)
except asyncio.CancelledError:
return
finally:
if self._media_group_tasks.get(media_group_id) is current_task:
self._media_group_tasks.pop(media_group_id, None)
self._media_group_tasks.pop(media_group_id, None)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 Media-group flush unconditionally pops _media_group_tasks entry, erasing replacement task (bug)

In _flush_media_group_event (line 7299–7308), the finally block was changed from a guarded pop (if self._media_group_tasks.get(media_group_id) is current_task) to an unconditional pop (self._media_group_tasks.pop(media_group_id, None)).

When a follow-up album update arrives for the same media_group_id, _queue_media_group_event (line 7274) cancels the prior flush task and creates a replacement stored under the same key — all synchronously, before the cancelled task's CancelledError is delivered. The cancelled original's finally block then unconditionally pops the dict entry, removing the replacement task, not the original.

Consequences:

  1. The replacement task is no longer tracked in _media_group_tasks, so disconnect() cannot find and cancel it.
  2. Subsequent album updates won't see the existing task to cancel/replace it, creating multiple concurrent flushes for the same album.
  3. An orphaned task may still dispatch handle_message with partially-aggregated events.

This same guarded-pop pattern is correctly preserved in _flush_text_batch (line 6942) and _flush_photo_batch (line 6973), confirming the removal was unintended.

💡 Suggestion: Restore the current_task identity check in the finally block of _flush_media_group_event, matching the pattern still present in _flush_text_batch (line 6942) and _flush_photo_batch (line 6973).

📋 Prompt for AI Agents

In plugins/platforms/telegram/adapter.py, in _flush_media_group_event (line 7299), add back current_task = asyncio.current_task() as the first line of the method (before the try block). Then change line 7308 from:

self._media_group_tasks.pop(media_group_id, None)

to:

if self._media_group_tasks.get(media_group_id) is current_task:
    self._media_group_tasks.pop(media_group_id, None)

This matches the identical guarded-pop pattern used in _flush_text_batch (line 6942) and _flush_photo_batch (line 6973).


async def _handle_sticker(self, msg: Message, event: "MessageEvent") -> None:
"""
Expand Down
Loading