diff --git a/.env.example b/.env.example index de1dde3..57f3f00 100644 --- a/.env.example +++ b/.env.example @@ -22,6 +22,10 @@ TGCLIENT_WRITE_LIMIT_N=20 # мутации через MCP: столько TGCLIENT_WRITE_LIMIT_WINDOW=60 # ... за столько секунд, на аккаунт +# Пуш-события сообщений в Synapse (tg_message_received/edited/deleted) — +# payload несёт тексты переписки; включать осознанно (=1) +TGCLIENT_NOTIFY_MESSAGES=0 + # gnexus-auth (SSO): задан client_id — браузерная авторизация через SSO; # пусто — сервис открыт, все акки уходят локальному служебному юзеру TGCLIENT_AUTH_BASE_URL= diff --git a/README.md b/README.md index 7839e06..c7f43da 100644 --- a/README.md +++ b/README.md @@ -66,8 +66,14 @@ | `tg_flood_wait` | api/flood_wait | high | день+аккаунт | длинный FloodWait (>60 с) в Telethon | | `mcp_token_issued` | token/issued | normal | — | выпуск персонального MCP-ключа | | `mcp_token_revoked` | token/revoked | normal | — | ревок ключа (свой или админский) | +| `tg_message_received` | message/received | normal | — | новое сообщение (text + media-мета: voice/round/duration/waveform) | +| `tg_message_edited` | message/edited | normal | — | правка сообщения | +| `tg_message_deleted` | message/deleted | low | — | удаление (ids; ttl 1 ч) | -Plaintext MCP-ключа никогда не уходит в конверты (только hint и метаданные). +Пуш-события сообщений — то, на что реагирует ИИ-агент/скрипт: подписка через +**Synapse Target** (webhook) — MCP подписок не даёт (stateless request/response). +Включается env `TGCLIENT_NOTIFY_MESSAGES=1` (payload несёт личные переписки — +по умолчанию выключено). Plaintext MCP-ключа в конверты не уходит. ## Разработка (без docker) diff --git a/backend/app/config.py b/backend/app/config.py index 1d69344..b97033d 100644 --- a/backend/app/config.py +++ b/backend/app/config.py @@ -29,6 +29,9 @@ # Лимиты max_active_accounts: int = 30 # аккаунтов на пользователя + # Пуш-события сообщений в Synapse (tg_message_*): по умолчанию выключены — + # payload несёт личные данные; включается осознанно (TGCLIENT_NOTIFY_MESSAGES=1) + notify_messages: bool = False media_max_bytes: int = 20 * 1024 * 1024 # cap download/upload через MCP write_limit_n: int = 20 # мутаций на аккаунт за окно write_limit_window: float = 60.0 diff --git a/backend/app/synapse_report.py b/backend/app/synapse_report.py index 0f105d0..9b898c5 100644 --- a/backend/app/synapse_report.py +++ b/backend/app/synapse_report.py @@ -17,6 +17,13 @@ | tg_flood_wait | api/flood_wait | high | д.+акк | call_with_flood_guard (длинный FloodWait) | | mcp_token_issued | token/issued | normal | нет | POST /me/mcp_tokens (plaintext НЕ уходит) | | mcp_token_revoked | token/revoked | normal | нет | ревок (свой и админский) | +| tg_message_received | message/received | normal | нет | новое сообщение (tg/events.py, TGCLIENT_NOTIFY_MESSAGES=1) | +| tg_message_edited | message/edited | normal | нет | правка сообщения | +| tg_message_deleted | message/deleted | low | нет | удаление (ids; ttl 1 ч) | + +Пуш-события сообщений — для ИИ-агентов/скриптов: подписка по канону +notifications.md — сервис-получатель регистрируется в Synapse как Target +(webhook); MCP подписок не даёт (stateless request/response). В конверте нет получателей/каналов/топиков — кому доставить, решает маршрутизация Synapse. Secrets: plaintext mcp_* только в ответе 201 и @@ -48,6 +55,9 @@ "tg_flood_wait": Sink("api", "flood_wait", "high", True), "mcp_token_issued": Sink("token", "issued", "normal", False), "mcp_token_revoked": Sink("token", "revoked", "normal", False), + "tg_message_received": Sink("message", "received", "normal", False), + "tg_message_edited": Sink("message", "edited", "normal", False), + "tg_message_deleted": Sink("message", "deleted", "low", False), } _client = None @@ -79,11 +89,12 @@ return datetime.now(timezone.utc).strftime("%Y%m%d") -def report(event_type: str, payload: dict) -> None: +def report(event_type: str, payload: dict, *, ttl_seconds: int | None = None) -> None: """Конверт v1 в Synapse, fire-and-forget: падения доставки не роняют код. Приоритет high/critical уходит send-ом (ошибки видны в логе), остальное — тихим emit. Отчёт по имени из SINKS; неизвестное имя — программная ошибка. + ttl_seconds — переопределение ttl (пуш-события сообщений — 1 ч). """ sink = SINKS.get(event_type) if sink is None: @@ -100,7 +111,7 @@ priority=sink.priority, payload=payload, dedup_key=dedup, - ttl_seconds=TTL_SECONDS, + ttl_seconds=ttl_seconds if ttl_seconds is not None else TTL_SECONDS, ) ) diff --git a/backend/app/tg/events.py b/backend/app/tg/events.py new file mode 100644 index 0000000..531d983 --- /dev/null +++ b/backend/app/tg/events.py @@ -0,0 +1,84 @@ +"""Пуш-события 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()) \ No newline at end of file diff --git a/backend/app/tg/manager.py b/backend/app/tg/manager.py index 20e2ad1..56becea 100644 --- a/backend/app/tg/manager.py +++ b/backend/app/tg/manager.py @@ -68,6 +68,11 @@ raise DomainError(409, f"account #{account_id} is not authorized — login again") self.clients[account_id] = client await self._touch(account_id) + if get_settings().notify_messages: + # пуш-события сообщений → Synapse (агенты/скрипты подписаны в Synapse) + from app.tg.events import register_listeners + + register_listeners(client, account_id) return client async def drop(self, account_id: int) -> None: