Skip to content
Draft
Show file tree
Hide file tree
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
3 changes: 3 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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, в секундах |
Expand Down Expand Up @@ -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 |
Expand Down
2 changes: 2 additions & 0 deletions app/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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"),
)
21 changes: 15 additions & 6 deletions app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -87,35 +88,43 @@ 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:
log.info(
"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:
Expand Down
7 changes: 7 additions & 0 deletions app/max_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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] = []
Expand All @@ -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:
Expand Down Expand Up @@ -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:
Expand Down
153 changes: 92 additions & 61 deletions app/max_listener.py
Original file line number Diff line number Diff line change
Expand Up @@ -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__)

Expand Down Expand Up @@ -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<blockquote>{escape(fwd_text)}{attaches_str}</blockquote>{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
Expand Down Expand Up @@ -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("✅ <b>Max:</b> соединение восстановлено")
await sender.send("✅ <b>Max:</b> соединение восстановлено", reply_markup=status_kb)
else:
chat_count = len(resolver.chats)
await sender.send(f"✅ <b>Max:</b> подключён | чатов: {chat_count}")
await sender.send(f"✅ <b>Max:</b> подключён | чатов: {chat_count}", reply_markup=status_kb)
_first_connect = False

@client.on_disconnect
Expand All @@ -325,6 +394,16 @@ async def handle_disconnect():
_last_notif_time = datetime.now()
await sender.send("⚠️ <b>Max:</b> соединение потеряно, переподключение...")

@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(
Expand All @@ -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)
Expand All @@ -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<blockquote>{escape(fwd_text)}{attaches_str}</blockquote>{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
Loading
Loading