diff --git a/.env.example b/.env.example index 63e4ee2..56329f3 100644 --- a/.env.example +++ b/.env.example @@ -33,5 +33,8 @@ PLUGINS_ENABLED=true # Skip muted chats (same as «без звука» / не беспокоить in Max) # SKIP_MUTED=true +# Enable buffered digest for muted chats via /muted and inline button +# MUTED_DIGEST_ENABLED=true + # Log directory (default: logs) # LOG_DIR=logs diff --git a/README.md b/README.md index 7353162..8ee13ea 100644 --- a/README.md +++ b/README.md @@ -79,6 +79,7 @@ cp .env.example .env | `UNREAD_ONLY` | нет | `true` — пересылать только непрочитанные (если прочитали в Max — в TG не придёт) | | `UNREAD_DELAY_SEC` | нет | Задержка в секундах перед проверкой прочитанности (по умолчанию `2`) | | `SKIP_MUTED` | нет | `true` — не пересылать из заглушённых чатов Max («без звука») | +| `MUTED_DIGEST_ENABLED` | нет | `true` — накапливать сообщения из заглушённых чатов и выдавать по `/muted` или кнопке `📭 Заглушённые` | | `LOG_DIR` | нет | Путь к директории логов (по умолчанию `logs`) | | `TG_PROXY` | нет | SOCKS5-прокси для Telegram (`socks5://host:port`) | | `TG_READ_TIMEOUT` | нет | Таймаут чтения HTTP-ответа от Telegram, в секундах | @@ -396,6 +397,7 @@ cp .env.example .env | `UNREAD_ONLY` | no | `true` — forward only unread messages (skip if read in Max) | | `UNREAD_DELAY_SEC` | no | Delay before read check in seconds (default `2`) | | `SKIP_MUTED` | no | `true` — skip muted / do-not-disturb chats in Max | +| `MUTED_DIGEST_ENABLED` | no | `true` — buffer muted-chat messages and flush them via `/muted` or the `📭 Заглушённые` button | | `LOG_DIR` | no | Log directory path (default: `logs`) | | `TG_PROXY` | no | SOCKS5 proxy for Telegram (`socks5://host:port`) | | `TG_READ_TIMEOUT` | no | HTTP read timeout for Telegram responses, in seconds | diff --git a/app/config.py b/app/config.py index a2bcc31..3a43d9e 100644 --- a/app/config.py +++ b/app/config.py @@ -21,6 +21,7 @@ class Settings: unread_only: bool = False unread_delay_sec: float = 2.0 skip_muted: bool = False + muted_digest_enabled: bool = False def load_settings() -> Settings: @@ -58,4 +59,5 @@ def load_settings() -> Settings: unread_only=os.environ.get("UNREAD_ONLY", "").lower() in ("1", "true", "yes"), unread_delay_sec=float(os.environ.get("UNREAD_DELAY_SEC", "2") or "2"), skip_muted=os.environ.get("SKIP_MUTED", "").lower() in ("1", "true", "yes"), + muted_digest_enabled=os.environ.get("MUTED_DIGEST_ENABLED", "").lower() in ("1", "true", "yes"), ) diff --git a/app/main.py b/app/main.py index 82cb08d..ae837d1 100644 --- a/app/main.py +++ b/app/main.py @@ -11,6 +11,7 @@ from app.config import load_settings from app.hooks import hooks from app.max_listener import create_max_client +from app.muted_buffer import MutedMessageBuffer from app.plugins import load_plugins from app.tg_handler import build_tg_app from app.tg_sender import TelegramSender @@ -87,12 +88,16 @@ async def main(): media_write_timeout=settings.tg_media_write_timeout, ) await sender.start() + muted_buffer = MutedMessageBuffer() if settings.muted_digest_enabled else None + skip_muted_enabled = settings.skip_muted or settings.muted_digest_enabled client = create_max_client( settings.max_token, settings.max_device_id, sender, settings.max_chat_ids, debug=settings.debug, reply_enabled=settings.reply_enabled, unread_only=settings.unread_only, unread_delay_sec=settings.unread_delay_sec, - skip_muted=settings.skip_muted, + skip_muted=skip_muted_enabled, + muted_digest_enabled=settings.muted_digest_enabled, + muted_buffer=muted_buffer, ) if settings.unread_only: @@ -100,22 +105,26 @@ async def main(): "Unread-only mode: ON (delay %ss before forward check)", settings.unread_delay_sec, ) - if settings.skip_muted: + if skip_muted_enabled: log.info("Skip-muted mode: ON (no forwards from muted Max chats)") + if settings.muted_digest_enabled: + log.info("Muted digest mode: ON (/muted command enabled)") tg_app = None - if settings.reply_enabled: + if settings.reply_enabled or settings.muted_digest_enabled: tg_app = build_tg_app(settings.tg_bot_token, client, settings.tg_chat_id, - proxy_url=settings.tg_proxy, read_timeout=settings.tg_read_timeout, write_timeout=settings.tg_write_timeout) + proxy_url=settings.tg_proxy, read_timeout=settings.tg_read_timeout, write_timeout=settings.tg_write_timeout, + muted_buffer=muted_buffer, muted_digest_enabled=settings.muted_digest_enabled, + sender=sender, resolver=client.resolver, reply_enabled=settings.reply_enabled) await tg_app.initialize() await tg_app.start() await tg_app.updater.start_polling( drop_pending_updates=True, allowed_updates=Update.ALL_TYPES, ) - log.info("Telegram polling started (reply → Max enabled)") + log.info("Telegram polling started") else: - log.info("Reply to Max disabled (REPLY_ENABLED=false)") + log.info("Telegram polling disabled (REPLY_ENABLED=false and MUTED_DIGEST_ENABLED=false)") log.info("Starting Max listener...") try: diff --git a/app/max_client.py b/app/max_client.py index c2e0d93..ed8e43a 100644 --- a/app/max_client.py +++ b/app/max_client.py @@ -117,6 +117,7 @@ def __init__(self, token: str, device_id: str, chat_ids: str | None = None, debu self._dispatch_counter = 0 self._pending: dict[int, asyncio.Future] = {} self._on_disconnect_cb = None + self._on_mark_cb = None self._auth_pending = False self._auth_timeout_task: asyncio.Task | None = None self.chat_ids: list[int] = [] @@ -137,6 +138,10 @@ def on_disconnect(self, func): self._on_disconnect_cb = func return func + def on_mark(self, func): + self._on_mark_cb = func + return func + # ── transport ────────────────────────────────────────────────── async def _send(self, opcode: int, payload: dict) -> int: @@ -335,6 +340,8 @@ async def _handle(self, data: dict): elif op == OpCode.NOTIF_MARK: if self.read_tracker: self.read_tracker.on_notif_mark(payload) + if self._on_mark_cb: + await self._on_mark_cb(payload) elif op == OpCode.NOTIF_CHAT: if self.mute_tracker: diff --git a/app/max_listener.py b/app/max_listener.py index 198520a..7a2b6e4 100644 --- a/app/max_listener.py +++ b/app/max_listener.py @@ -6,8 +6,9 @@ from app.hooks import HookEvent, hooks from app.max_client import MaxClient, MaxMessage, OpCode +from app.muted_buffer import MutedMessageBuffer from app.resolver import ContactResolver -from app.tg_sender import TelegramSender, reply_keyboard +from app.tg_sender import TelegramSender, muted_digest_keyboard, reply_keyboard log = logging.getLogger(__name__) @@ -252,16 +253,83 @@ def _human_size(n: int) -> str: return f"{n:.1f} ТБ" +async def forward_message( + msg: MaxMessage, + client: MaxClient, + sender: TelegramSender, + resolver: ContactResolver, + reply_enabled: bool = False, +) -> None: + sender_label = escape(await resolver.resolve_user(msg.sender_id)) + is_dm = resolver.is_dm(msg.chat_id) + if len(client.chat_ids) == 1: + chat_label = "" + else: + chat_label = escape(resolver.chat_name(msg.chat_id)) + header_text = _header(msg, sender_label, chat_label, is_dm) + kb = reply_keyboard(msg.chat_id) if reply_enabled else None + + link = msg.link + link_type = link.get("type") if isinstance(link, dict) else None + + if link_type == "FORWARD": + await _handle_forward_message(link, header_text, client, sender, resolver, kb=kb) + if msg.text: + await sender.send(f"{header_text}\n{escape(msg.text)}", reply_markup=kb) + log.info("Forwarded message → TG") + return + + if link_type == "REPLY": + attaches_str, full_header, fwd_text = await _handle_reply_message(link, header_text, resolver) + if msg.text: + await sender.send( + f"{full_header}\n
{escape(fwd_text)}{attaches_str}{escape(msg.text)}", + reply_markup=kb, + ) + log.info("Forwarded reply → TG") + return + + meaningful_attaches = [ + a for a in msg.attaches + if isinstance(a, dict) and a.get("_type") not in ("CONTROL", "WIDGET", "INLINE_KEYBOARD", None) + ] + + if meaningful_attaches: + text_sent = False + for i, attach in enumerate(meaningful_attaches): + if i == 0 and msg.text: + cap = f"{header_text}\n{escape(msg.text)}" + text_sent = True + else: + cap = header_text + await _send_attach(attach, client, sender, cap, msg.chat_id, msg.message_id, kb=kb) + log.info("Forwarded attach _type=%s → TG", attach.get("_type")) + + if msg.text and not text_sent: + await sender.send(f"{header_text}\n{escape(msg.text)}", reply_markup=kb) + return + + if msg.text: + await sender.send(f"{header_text}\n{escape(msg.text)}", reply_markup=kb) + log.info("Forwarded text → TG") + return + + log.warning("Нетекстовое сообщение! %s", msg.attaches) + + def create_max_client( max_token: str, max_device_id: str, sender: TelegramSender, max_chat_ids: str | None = None, debug: bool = False, reply_enabled: bool = False, unread_only: bool = False, unread_delay_sec: float = 2.0, skip_muted: bool = False, + muted_digest_enabled: bool = False, + muted_buffer: MutedMessageBuffer | None = None, ) -> MaxClient: client = MaxClient( token=max_token, device_id=max_device_id, debug=debug, chat_ids=max_chat_ids, unread_only=unread_only, skip_muted=skip_muted, ) resolver = ContactResolver(client=client) + client.resolver = resolver _first_connect = True _notif_count = 0 @@ -307,11 +375,12 @@ async def handle_ready(snapshot: dict): is_reconnect=not _first_connect, ) + status_kb = muted_digest_keyboard() if muted_digest_enabled else None if not _first_connect: - await sender.send("✅ Max: соединение восстановлено") + await sender.send("✅ Max: соединение восстановлено", reply_markup=status_kb) else: chat_count = len(resolver.chats) - await sender.send(f"✅ Max: подключён | чатов: {chat_count}") + await sender.send(f"✅ Max: подключён | чатов: {chat_count}", reply_markup=status_kb) _first_connect = False @client.on_disconnect @@ -325,6 +394,16 @@ async def handle_disconnect(): _last_notif_time = datetime.now() await sender.send("⚠️ Max: соединение потеряно, переподключение...") + @client.on_mark + async def handle_mark(payload: dict) -> None: + if not muted_buffer or not client.read_tracker: + return + chat_id = payload.get("chatId") + if chat_id is None: + return + mark = client.read_tracker.mark_for_chat(chat_id) + await muted_buffer.prune_read(chat_id, mark) + @client.on_message async def handle_message(msg: MaxMessage): log.info( @@ -345,10 +424,6 @@ async def handle_message(msg: MaxMessage): log.info("Message filtered by hook (chat=%s)", msg.chat_id) return - if client.skip_muted and client.mute_tracker and client.mute_tracker.is_muted(msg.chat_id): - log.info("Skipped (muted chat in Max): chat=%s", msg.chat_id) - return - if client.unread_only and client.read_tracker: if unread_delay_sec > 0: await asyncio.sleep(unread_delay_sec) @@ -360,64 +435,20 @@ async def handle_message(msg: MaxMessage): ) return + if client.skip_muted and client.mute_tracker and client.mute_tracker.is_muted(msg.chat_id): + if muted_digest_enabled and muted_buffer: + await muted_buffer.add(msg) + log.info("Buffered (muted chat in Max): chat=%s msg=%s", msg.chat_id, msg.message_id) + else: + log.info("Skipped (muted chat in Max): chat=%s", msg.chat_id) + return + async def _message_sent() -> None: await hooks.emit( HookEvent.ON_MESSAGE_SENT, msg=msg, resolver=resolver, client=client ) - sender_label = escape(await resolver.resolve_user(msg.sender_id)) - is_dm = resolver.is_dm(msg.chat_id) - if len(client.chat_ids) == 1: - chat_label = "" - else: - chat_label = escape(resolver.chat_name(msg.chat_id)) - header_text = _header(msg, sender_label, chat_label, is_dm) - kb = reply_keyboard(msg.chat_id) if reply_enabled else None - - link = msg.link - link_type = link.get("type") if isinstance(link, dict) else None - - if link_type == "FORWARD": - await _handle_forward_message(link, header_text, client, sender, resolver, kb=kb) - if msg.text: - await sender.send(f"{header_text}\n{escape(msg.text)}", reply_markup=kb) - log.info("Forwarded message → TG") - await _message_sent() - return - - if link_type == "REPLY": - attaches_str, full_header, fwd_text = await _handle_reply_message(link, header_text, resolver) - if msg.text: - await sender.send(f"{full_header}\n
{escape(fwd_text)}{attaches_str}{escape(msg.text)}", reply_markup=kb) - log.info("Forwarded reply → TG") - await _message_sent() - return - - meaningful_attaches = [ - a for a in msg.attaches - if isinstance(a, dict) and a.get("_type") not in ("CONTROL", "WIDGET", "INLINE_KEYBOARD", None) - ] - - if meaningful_attaches: - text_sent = False - for i, attach in enumerate(meaningful_attaches): - if i == 0 and msg.text: - cap = f"{header_text}\n{escape(msg.text)}" - text_sent = True - else: - cap = header_text - await _send_attach(attach, client, sender, cap, msg.chat_id, msg.message_id, kb=kb) - log.info("Forwarded attach _type=%s → TG", attach.get("_type")) - - if msg.text and not text_sent: - await sender.send(f"{header_text}\n{escape(msg.text)}", reply_markup=kb) - await _message_sent() - else: - if msg.text: - await sender.send(f"{header_text}\n{escape(msg.text)}", reply_markup=kb) - log.info("Forwarded text → TG") - await _message_sent() - else: - log.warning("Нетекстовое сообщение! %s", msg.attaches) + await forward_message(msg, client, sender, resolver, reply_enabled=reply_enabled) + await _message_sent() return client diff --git a/app/muted_buffer.py b/app/muted_buffer.py new file mode 100644 index 0000000..a84ed8e --- /dev/null +++ b/app/muted_buffer.py @@ -0,0 +1,83 @@ +"""In-memory buffer for messages from muted Max chats.""" + +from __future__ import annotations + +import asyncio +from collections.abc import Iterable +from typing import TYPE_CHECKING, Any + +if TYPE_CHECKING: + from app.max_client import MaxMessage + + +class MutedMessageBuffer: + """Stores muted-chat messages grouped by chat preserving first-seen chat order.""" + + def __init__(self) -> None: + self._lock = asyncio.Lock() + self._chat_order: list[Any] = [] + self._messages: dict[Any, list[MaxMessage]] = {} + + async def add(self, msg: MaxMessage) -> None: + chat_id = msg.chat_id + async with self._lock: + bucket = self._messages.get(chat_id) + if bucket is None: + self._chat_order.append(chat_id) + self._messages[chat_id] = [msg] + return + bucket.append(msg) + + async def get_grouped(self) -> list[tuple[Any, list[MaxMessage]]]: + async with self._lock: + grouped: list[tuple[Any, list[MaxMessage]]] = [] + for chat_id in self._chat_order: + messages = self._messages.get(chat_id) + if messages: + grouped.append((chat_id, list(messages))) + return grouped + + async def pop_grouped(self) -> list[tuple[Any, list[MaxMessage]]]: + async with self._lock: + grouped: list[tuple[Any, list[MaxMessage]]] = [] + for chat_id in self._chat_order: + messages = self._messages.get(chat_id) + if messages: + grouped.append((chat_id, list(messages))) + self._chat_order.clear() + self._messages.clear() + return grouped + + async def clear(self) -> None: + async with self._lock: + self._chat_order.clear() + self._messages.clear() + + async def count(self) -> int: + async with self._lock: + return sum(len(messages) for messages in self._messages.values()) + + async def prune_read(self, chat_id: Any, read_mark: int) -> None: + async with self._lock: + messages = self._messages.get(chat_id) + if not messages: + return + pruned = [msg for msg in messages if not _is_read(msg, read_mark)] + if pruned: + self._messages[chat_id] = pruned + return + self._messages.pop(chat_id, None) + self._chat_order = [cid for cid in self._chat_order if cid != chat_id] + + async def prune_many(self, read_marks: Iterable[tuple[Any, int]]) -> None: + for chat_id, read_mark in read_marks: + await self.prune_read(chat_id, read_mark) + + +def _is_read(msg: MaxMessage, read_mark: int) -> bool: + if msg.timestamp is None: + return False + try: + return int(msg.timestamp) <= int(read_mark) + except (TypeError, ValueError): + return False diff --git a/app/muted_digest.py b/app/muted_digest.py new file mode 100644 index 0000000..03fb68e --- /dev/null +++ b/app/muted_digest.py @@ -0,0 +1,43 @@ +"""Helpers to flush buffered messages from muted chats.""" + +from __future__ import annotations + +from html import escape + +from app.max_client import MaxClient +from app.max_listener import forward_message +from app.muted_buffer import MutedMessageBuffer +from app.resolver import ContactResolver +from app.tg_sender import TelegramSender + + +def _msg_sort_key(msg) -> tuple[int, str]: + ts = msg.timestamp + try: + ts_val = int(ts) + except (TypeError, ValueError): + ts_val = 0 + return ts_val, msg.message_id + + +async def flush_muted_digest( + buffer: MutedMessageBuffer, + client: MaxClient, + sender: TelegramSender, + resolver: ContactResolver, + reply_enabled: bool = False, +) -> int: + grouped = await buffer.pop_grouped() + if not grouped: + await sender.send("📭 Нет новых сообщений из заглушённых чатов.") + return 0 + + total = 0 + for chat_id, messages in grouped: + messages_sorted = sorted(messages, key=_msg_sort_key) + chat_label = escape(resolver.chat_name(chat_id)) + await sender.send(f"🔇 {chat_label} ({len(messages_sorted)} сообщений)") + for msg in messages_sorted: + await forward_message(msg, client, sender, resolver, reply_enabled=reply_enabled) + total += 1 + return total diff --git a/app/read_tracker.py b/app/read_tracker.py index 0d7e5ee..3b07efa 100644 --- a/app/read_tracker.py +++ b/app/read_tracker.py @@ -59,3 +59,6 @@ def is_unread(self, chat_id: Any, message_time: Any) -> bool: except (TypeError, ValueError): return True return msg_ts > self._marks.get(chat_id, 0) + + def mark_for_chat(self, chat_id: Any) -> int: + return self._marks.get(chat_id, 0) diff --git a/app/tg_handler.py b/app/tg_handler.py index 4f90c05..4d64d37 100644 --- a/app/tg_handler.py +++ b/app/tg_handler.py @@ -14,6 +14,10 @@ from app.hooks import HookEvent, hooks from app.max_client import MaxClient +from app.muted_buffer import MutedMessageBuffer +from app.muted_digest import flush_muted_digest +from app.resolver import ContactResolver +from app.tg_sender import TelegramSender log = logging.getLogger(__name__) @@ -21,6 +25,11 @@ PENDING_REPLY_LABEL_KEY = "pending_reply_label" _ALLOWED_CHAT_ID_KEY = "allowed_chat_id" +_MUTED_BUFFER_KEY = "muted_buffer" +_MUTED_DIGEST_ENABLED_KEY = "muted_digest_enabled" +_TG_SENDER_KEY = "tg_sender" +_RESOLVER_KEY = "resolver" +_REPLY_ENABLED_KEY = "reply_enabled" async def _on_reply_button(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: @@ -127,8 +136,62 @@ async def _on_text_reply(update: Update, context: ContextTypes.DEFAULT_TYPE) -> await update.message.reply_text("⚠️ Ошибка при отправке в Max.") +async def _on_muted_digest(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not context.bot_data.get(_MUTED_DIGEST_ENABLED_KEY): + await update.message.reply_text("Функция чтения заглушённых чатов отключена.") + return + buffer: MutedMessageBuffer | None = context.bot_data.get(_MUTED_BUFFER_KEY) + sender: TelegramSender | None = context.bot_data.get(_TG_SENDER_KEY) + resolver: ContactResolver | None = context.bot_data.get(_RESOLVER_KEY) + max_client: MaxClient | None = context.bot_data.get("max_client") + if not buffer or not sender or not resolver or not max_client: + await update.message.reply_text("⚠️ Компоненты для чтения заглушённых чатов недоступны.") + return + + total = await flush_muted_digest( + buffer=buffer, + client=max_client, + sender=sender, + resolver=resolver, + reply_enabled=bool(context.bot_data.get(_REPLY_ENABLED_KEY)), + ) + if total > 0: + await update.message.reply_text(f"✅ Отправлено сообщений из заглушённых чатов: {total}.") + + +async def _on_muted_button(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + query = update.callback_query + await query.answer() + if query.data != "muted:flush": + return + if not context.bot_data.get(_MUTED_DIGEST_ENABLED_KEY): + await query.message.reply_text("Функция чтения заглушённых чатов отключена.") + return + + buffer: MutedMessageBuffer | None = context.bot_data.get(_MUTED_BUFFER_KEY) + sender: TelegramSender | None = context.bot_data.get(_TG_SENDER_KEY) + resolver: ContactResolver | None = context.bot_data.get(_RESOLVER_KEY) + max_client: MaxClient | None = context.bot_data.get("max_client") + if not buffer or not sender or not resolver or not max_client: + await query.message.reply_text("⚠️ Компоненты для чтения заглушённых чатов недоступны.") + return + + total = await flush_muted_digest( + buffer=buffer, + client=max_client, + sender=sender, + resolver=resolver, + reply_enabled=bool(context.bot_data.get(_REPLY_ENABLED_KEY)), + ) + if total > 0: + await query.message.reply_text(f"✅ Отправлено сообщений из заглушённых чатов: {total}.") + + def build_tg_app(token: str, max_client: MaxClient, allowed_chat_id: str, - proxy_url: str | None = None, read_timeout: int | None = None, write_timeout: int | None = None) -> Application: + proxy_url: str | None = None, read_timeout: int | None = None, write_timeout: int | None = None, + muted_buffer: MutedMessageBuffer | None = None, muted_digest_enabled: bool = False, + sender: TelegramSender | None = None, resolver: ContactResolver | None = None, + reply_enabled: bool = False) -> Application: """Build and configure the Telegram Application with handlers.""" builder = Application.builder().token(token) if proxy_url: @@ -140,11 +203,19 @@ def build_tg_app(token: str, max_client: MaxClient, allowed_chat_id: str, app = builder.build() app.bot_data["max_client"] = max_client app.bot_data[_ALLOWED_CHAT_ID_KEY] = int(allowed_chat_id) + app.bot_data[_MUTED_BUFFER_KEY] = muted_buffer + app.bot_data[_MUTED_DIGEST_ENABLED_KEY] = muted_digest_enabled + app.bot_data[_TG_SENDER_KEY] = sender + app.bot_data[_RESOLVER_KEY] = resolver + app.bot_data[_REPLY_ENABLED_KEY] = reply_enabled chat_filter = filters.Chat(chat_id=int(allowed_chat_id)) app.add_handler(CallbackQueryHandler(_on_reply_button, pattern=r"^reply:")) + app.add_handler(CallbackQueryHandler(_on_muted_button, pattern=r"^muted:")) app.add_handler(CommandHandler("cancel", _on_cancel, filters=chat_filter)) + app.add_handler(CommandHandler("muted", _on_muted_digest, filters=chat_filter)) + app.add_handler(CommandHandler("readmuted", _on_muted_digest, filters=chat_filter)) app.add_handler(MessageHandler(filters.TEXT & ~filters.COMMAND & chat_filter, _on_text_reply)) return app diff --git a/app/tg_sender.py b/app/tg_sender.py index 5d136f8..f231d92 100644 --- a/app/tg_sender.py +++ b/app/tg_sender.py @@ -21,6 +21,13 @@ def reply_keyboard(max_chat_id) -> InlineKeyboardMarkup: ]]) +def muted_digest_keyboard() -> InlineKeyboardMarkup: + """Build keyboard with a button to flush muted-chat backlog.""" + return InlineKeyboardMarkup([[ + InlineKeyboardButton("📭 Заглушённые", callback_data="muted:flush") + ]]) + + class TelegramSender: def __init__( self, diff --git a/tests/test_muted_buffer.py b/tests/test_muted_buffer.py new file mode 100644 index 0000000..f6bcf5e --- /dev/null +++ b/tests/test_muted_buffer.py @@ -0,0 +1,43 @@ +from types import SimpleNamespace + +import pytest + +from app.muted_buffer import MutedMessageBuffer + + +@pytest.mark.asyncio +async def test_add_and_get_grouped_preserves_chat_order(): + buffer = MutedMessageBuffer() + await buffer.add(SimpleNamespace(chat_id=2, timestamp=20, message_id="m2")) + await buffer.add(SimpleNamespace(chat_id=1, timestamp=10, message_id="m1")) + await buffer.add(SimpleNamespace(chat_id=2, timestamp=30, message_id="m3")) + + grouped = await buffer.get_grouped() + assert [chat_id for chat_id, _ in grouped] == [2, 1] + assert [m.message_id for m in grouped[0][1]] == ["m2", "m3"] + + +@pytest.mark.asyncio +async def test_prune_read_removes_read_messages_and_chat_bucket(): + buffer = MutedMessageBuffer() + await buffer.add(SimpleNamespace(chat_id=5, timestamp=100, message_id="a")) + await buffer.add(SimpleNamespace(chat_id=5, timestamp=200, message_id="b")) + + await buffer.prune_read(5, 150) + grouped = await buffer.get_grouped() + assert len(grouped) == 1 + assert [m.message_id for m in grouped[0][1]] == ["b"] + + await buffer.prune_read(5, 999) + grouped = await buffer.get_grouped() + assert grouped == [] + + +@pytest.mark.asyncio +async def test_pop_grouped_clears_buffer(): + buffer = MutedMessageBuffer() + await buffer.add(SimpleNamespace(chat_id=7, timestamp=1, message_id="m")) + + grouped = await buffer.pop_grouped() + assert len(grouped) == 1 + assert await buffer.count() == 0 diff --git a/tests/test_muted_digest.py b/tests/test_muted_digest.py new file mode 100644 index 0000000..7f56b6f --- /dev/null +++ b/tests/test_muted_digest.py @@ -0,0 +1,50 @@ +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import pytest + +from app.muted_buffer import MutedMessageBuffer +from app.muted_digest import flush_muted_digest + + +@pytest.mark.asyncio +async def test_flush_muted_digest_empty_buffer_sends_empty_message(): + sender = SimpleNamespace(send=AsyncMock()) + resolver = SimpleNamespace(chat_name=lambda chat_id: f"chat-{chat_id}") + + total = await flush_muted_digest( + buffer=MutedMessageBuffer(), + client=SimpleNamespace(), + sender=sender, + resolver=resolver, + ) + + assert total == 0 + sender.send.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_flush_muted_digest_sorts_within_chat(monkeypatch): + buffer = MutedMessageBuffer() + await buffer.add(SimpleNamespace(chat_id=10, timestamp=20, message_id="2")) + await buffer.add(SimpleNamespace(chat_id=10, timestamp=10, message_id="1")) + await buffer.add(SimpleNamespace(chat_id=20, timestamp=15, message_id="3")) + + sender = SimpleNamespace(send=AsyncMock()) + resolver = SimpleNamespace(chat_name=lambda chat_id: f"chat-{chat_id}") + calls = [] + + async def fake_forward(msg, client, sender, resolver, reply_enabled=False): + calls.append((msg.chat_id, msg.message_id)) + + monkeypatch.setattr("app.muted_digest.forward_message", fake_forward) + + total = await flush_muted_digest( + buffer=buffer, + client=SimpleNamespace(), + sender=sender, + resolver=resolver, + ) + + assert total == 3 + assert calls == [(10, "1"), (10, "2"), (20, "3")] diff --git a/tests/test_tg_handler.py b/tests/test_tg_handler.py index 26222d2..881b07e 100644 --- a/tests/test_tg_handler.py +++ b/tests/test_tg_handler.py @@ -7,6 +7,8 @@ PENDING_REPLY_KEY, PENDING_REPLY_LABEL_KEY, _on_cancel, + _on_muted_button, + _on_muted_digest, _on_reply_button, _on_text_reply, ) @@ -319,3 +321,53 @@ async def test_escapes_html_in_success_label(self): args = update.message.reply_text.call_args[0][0] assert 'evil' not in args assert '<b>evil</b>' in args + + +class TestMutedDigestHandlers: + @pytest.mark.asyncio + async def test_on_muted_digest_disabled(self): + update = _make_message_update("/muted") + ctx = _make_context(bot_data={"muted_digest_enabled": False}) + + await _on_muted_digest(update, ctx) + + update.message.reply_text.assert_called_once() + + @pytest.mark.asyncio + async def test_on_muted_digest_flushes_and_reports_count(self): + update = _make_message_update("/muted") + ctx = _make_context( + bot_data={ + "muted_digest_enabled": True, + "muted_buffer": MagicMock(), + "tg_sender": MagicMock(), + "resolver": MagicMock(), + "max_client": MagicMock(), + "reply_enabled": False, + } + ) + with patch("app.tg_handler.flush_muted_digest", new=AsyncMock(return_value=3)) as flush_mock: + await _on_muted_digest(update, ctx) + + flush_mock.assert_awaited_once() + update.message.reply_text.assert_called_once() + + @pytest.mark.asyncio + async def test_on_muted_button_flushes(self): + query = _make_callback_query("muted:flush") + update = _make_update_with_query(query) + ctx = _make_context( + bot_data={ + "muted_digest_enabled": True, + "muted_buffer": MagicMock(), + "tg_sender": MagicMock(), + "resolver": MagicMock(), + "max_client": MagicMock(), + "reply_enabled": False, + } + ) + with patch("app.tg_handler.flush_muted_digest", new=AsyncMock(return_value=1)) as flush_mock: + await _on_muted_button(update, ctx) + + flush_mock.assert_awaited_once() + query.message.reply_text.assert_called_once()