From 1e97dd9fe5a6f524fde6237450a443ce7da17d77 Mon Sep 17 00:00:00 2001 From: 3252a8 <3252a8@proton.me> Date: Fri, 15 May 2026 23:59:49 +0300 Subject: [PATCH] refactor: add per-chat throttle in message queue to avoid 429 on rapid same-chat sends --- bot/utils/message_queue.py | 53 ++++++++++++++++++++++++++------------ 1 file changed, 37 insertions(+), 16 deletions(-) diff --git a/bot/utils/message_queue.py b/bot/utils/message_queue.py index d0088c7..28fea39 100644 --- a/bot/utils/message_queue.py +++ b/bot/utils/message_queue.py @@ -27,6 +27,11 @@ class QueuedMessage: class MessageQueue: """Message queue with rate limiting for Telegram API""" + # Telegram allows ~1 message/sec to the same chat before returning 429. + PER_CHAT_MIN_INTERVAL_SECONDS = 1.0 + # Drop per-chat timestamps older than this to keep the dict bounded. + PER_CHAT_TTL_SECONDS = 5.0 + def __init__(self, messages_per_second: float, burst_size: int = 5): self.messages_per_second = messages_per_second self.burst_size = burst_size @@ -36,6 +41,7 @@ class MessageQueue: self.delay_between_messages = 1.0 / messages_per_second self.total_sent = 0 self.total_failed = 0 + self._chat_last_sent: Dict[int, datetime] = {} async def add_message(self, message: QueuedMessage) -> None: """Add message to queue""" @@ -52,14 +58,14 @@ class MessageQueue: try: while self.queue: - # Check if we need to wait - await self._wait_if_needed() + # Peek at the head to honor per-chat throttling before popping. + message = self.queue[0] + await self._wait_if_needed(message.chat_id) + self.queue.popleft() - # Get and process next message - message = self.queue.popleft() try: await self._send_message(message) - self._record_send_time() + self._record_send_time(message.chat_id) except TelegramBadRequest as exc: fallback_message = self._build_profile_link_fallback(message, exc) @@ -72,7 +78,7 @@ class MessageQueue: ) try: await self._send_message(fallback_message) - self._record_send_time() + self._record_send_time(message.chat_id) continue except Exception as retry_exc: self.total_failed += 1 @@ -91,19 +97,27 @@ class MessageQueue: finally: self.is_processing = False - async def _wait_if_needed(self) -> None: - """Wait if we need to respect rate limits""" - if not self.last_send_times: - return + async def _wait_if_needed(self, chat_id: Optional[int] = None) -> None: + """Wait if we need to respect global and per-chat rate limits.""" + now = datetime.now() + waits: list[float] = [] - # Calculate time since last message - time_since_last = (datetime.now() - self.last_send_times[-1]).total_seconds() + if self.last_send_times: + time_since_last = (now - self.last_send_times[-1]).total_seconds() + if time_since_last < self.delay_between_messages: + waits.append(self.delay_between_messages - time_since_last) - if time_since_last < self.delay_between_messages: - wait_time = self.delay_between_messages - time_since_last - await asyncio.sleep(wait_time) + if chat_id is not None: + last_chat = self._chat_last_sent.get(chat_id) + if last_chat is not None: + time_since_chat = (now - last_chat).total_seconds() + if time_since_chat < self.PER_CHAT_MIN_INTERVAL_SECONDS: + waits.append(self.PER_CHAT_MIN_INTERVAL_SECONDS - time_since_chat) - def _record_send_time(self) -> None: + if waits: + await asyncio.sleep(max(waits)) + + def _record_send_time(self, chat_id: Optional[int] = None) -> None: """Track sent message timestamps and purge old entries for rate limiting.""" now = datetime.now() self.last_send_times.append(now) @@ -113,6 +127,13 @@ class MessageQueue: while self.last_send_times and self.last_send_times[0] < cutoff_time: self.last_send_times.popleft() + if chat_id is not None: + self._chat_last_sent[chat_id] = now + chat_cutoff = now - timedelta(seconds=self.PER_CHAT_TTL_SECONDS) + stale = [cid for cid, ts in self._chat_last_sent.items() if ts < chat_cutoff] + for cid in stale: + self._chat_last_sent.pop(cid, None) + def _build_profile_link_fallback( self, message: QueuedMessage, exc: Exception ) -> Optional[QueuedMessage]: