Newer
Older
tgclient-mcp / backend / app / tg / events.py
"""Пуш-события Telegram → Synapse: события, на которые может реагировать
ИИ-агент или скрипт (например «новое сообщение»).

MCP-протокол подписок не даёт (stateless json_response, запрос-ответ) — канон
10-platform/notifications.md: агент/скрипт регистрируется в Synapse как
Target (webhook) и получает доставки; tgclient пушит конверты о том, что
случилось у него. Маршрутизацию и ретраи делает Synapse.

Включается env TGCLIENT_NOTIFY_MESSAGES=1 (по умолчанию выключено: payload
несёт личные данные — текст переписки). Хендлеры вешаются на каждого
поднятого клиента (AccountManager.get_client), поэтому переживают
автопереподключения Telethon.

События (в app/synapse_report.py SINKS):
- tg_message_received  message/received  normal — новое сообщение в любом
  диалоге аккаунта (текст + media-мета: is_voice/is_round/duration/waveform);
- tg_message_edited    message/edited    normal — правка сообщения;
- tg_message_deleted   message/deleted   low    — удаление (ids; у юзеров
  Telegram присылает только ids, диалог может быть не определён).
"""

from telethon import events

from app.mcp.serializers import marked_peer_id, message_dict
from app.synapse_report import report

MESSAGE_TTL_SECONDS = 3600  # события-триггеры актуальны час; Synapse не доставит старьё


async def _owner_of(account_id: int) -> str | None:
    from app.db import get_db

    cursor = await get_db().execute(
        "SELECT user_id FROM accounts WHERE id = ?", (account_id,)
    )
    row = await cursor.fetchone()
    return row["user_id"] if row else None


async def _report_message(event_type: str, account_id: int, message) -> None:
    """Конверт сообщения: владелец, аккаунт, диалог, сериализация телом тулов."""
    from app.main import get_account_manager

    owner = await _owner_of(account_id)
    if owner is None:
        return
    report(event_type, {
        "user_id": owner,
        "entity": f"account-{account_id}",
        "account_id": account_id,
        "dialog_id": marked_peer_id(getattr(message, "peer_id", None)),
        "message": message_dict(message),
    }, ttl_seconds=MESSAGE_TTL_SECONDS)


async def _on_deleted(account_id: int, event) -> None:
    owner = await _owner_of(account_id)
    if owner is None:
        return
    # Telegram в приватных чатах присылает только удалённые id —
    # какой это диалог, не определяем; агенту достаточно message_ids
    report("tg_message_deleted", {
        "user_id": owner,
        "entity": f"account-{account_id}",
        "account_id": account_id,
        "message_ids": list(getattr(event, "deleted_ids", []) or [])[:50],
    }, ttl_seconds=MESSAGE_TTL_SECONDS)


def register_listeners(client, account_id: int) -> None:
    """Повесить пуш-хендлеры на клиент аккаунта (AccountManager.get_client)."""

    async def on_new(event) -> None:
        await _report_message("tg_message_received", account_id, event.message)

    async def on_edit(event) -> None:
        await _report_message("tg_message_edited", account_id, event.message)

    async def on_delete(event) -> None:
        await _on_deleted(account_id, event)

    client.add_event_handler(on_new, events.NewMessage())
    client.add_event_handler(on_edit, events.MessageEdited())
    client.add_event_handler(on_delete, events.MessageDeleted())