From 2c05443c897570b28e25e1d33212e10ab257ae3d Mon Sep 17 00:00:00 2001 From: "Nikita Tyukalov, ASUS, Linux" Date: Mon, 17 Aug 2026 03:33:54 +0300 Subject: [PATCH] Initial commit --- .gitignore | 7 ++ README.md | 80 +++++++++++++++++++ bot.py | 202 +++++++++++++++++++++++++++++++++++++++++++++++ config.py | 53 +++++++++++++ database.py | 0 datatypes.py | 15 ++++ logic.py | 148 ++++++++++++++++++++++++++++++++++ main.py | 52 ++++++++++++ requirements.txt | 3 + util.py | 164 ++++++++++++++++++++++++++++++++++++++ 10 files changed, 724 insertions(+) create mode 100644 .gitignore create mode 100644 README.md create mode 100644 bot.py create mode 100644 config.py create mode 100644 database.py create mode 100644 datatypes.py create mode 100644 logic.py create mode 100644 main.py create mode 100644 requirements.txt create mode 100644 util.py diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..bca140b --- /dev/null +++ b/.gitignore @@ -0,0 +1,7 @@ +.venv/ +session_storage/ +__pycache__/ +*.swp +*.swo +*.vscode +*.json \ No newline at end of file diff --git a/README.md b/README.md new file mode 100644 index 0000000..c0c435d --- /dev/null +++ b/README.md @@ -0,0 +1,80 @@ +# 2026-matrix-csonac + +Реализация CSoNaC (Centralized System of Notification and Control) для Matrix. +Раньше был бот с таким же функционалом, но для Telegram. Telegram больше не в +почёте, и теперь у меня всё в локальном Matrix, поэтому бот тоже перенесён сюда. + +## Подготовка окружения + +1. Клонируйте репозиторий и перейдите в его директорию +```bash +git clone https://git.tyukalov.su/nikita/2026-matrix-csonac.git +cd 2026-matrix-csonac +``` +2. Создайте `venv`, активируйте его +```bash +python3 -m venv .venv +. .venv/bin/activate +``` +5. Запустите приложение чтобы сгененерировать конфиг +```bash +python main.py +``` +6. Отредактируйте конфиг +7. Запустите приложение ещё раз, чтобы авторизоваться +```bash +python main.py +``` + +## Запуск + +Для запуска приложения вы можете либо активировать `venv`, либо просто использовать полный путь до интерпретатора Python: +```bash +~/2026-matrix-downloader/.venv/bin/python main.py +``` +```bash +. .venv/bin/activate +python main.py +``` + +## Как работает этот бот + +Этот бот запускает веб-сервер, используя `web_host` и `web_port` из файла +конфигурации. + +> Не открывайте доступ к серверу напрямую! Используйте reverse proxy, например, +> `nginx`. Это также позволит вам использовать защищённое соединение, что +> исключит возможность применения атаки Man-in-the-Middle для перехвата токена. + +Все запросы к веб-серверу требуют авторизации, используя HTTP заголовок +`Authorization` и схему `Bearer`. Токен авторизации генерируется посредством +взаимодействия с ботом в Matrix. + +> Для каждого сервиса, использующего службу, рекомендуется генерировать свой +> собственный токен. Это позволит отозвать токен только для одной службы, если +> он будет украден. + +Доступные эндпоинты: +- `POST //notify` + - **Описание.** Используется, чтобы отправить уведомление в указанный канал. + Вместо `` указывается код канала, получаемый при помощи команды + `!channel_code`, выполненной в комнате Matrix. + - **Тело запроса.** Тело запроса представляет собой `json` объект: + ```json + { + "service": "<название сервиса, передаётся только если разрешено>", + "subservice": "<название подсервиса>", + "text": "<текст уведомления>", + "urgent": false + } + ``` + - **Тело ответа.** Тело ответа представляет собой `json` объект. В случае + успеха в нём будут все поля, перечисляемые ниже. В случае провала - только + поле `error`, содержащее текстовое описание ошибки. + ```json + { + "error": null, + "notification_id": "<здесь будет Notification ID>" + } + ``` +- `` \ No newline at end of file diff --git a/bot.py b/bot.py new file mode 100644 index 0000000..adae7f1 --- /dev/null +++ b/bot.py @@ -0,0 +1,202 @@ +""" +matrix-nio basics wrapper +""" + +import asyncio +import time +import traceback +from pathlib import Path + +from nio import LoginError, LoginResponse, SyncResponse +from nio import WhoamiResponse +from nio import AsyncClient, AsyncClientConfig + +from datatypes import AppConfig +import util + + + +# +# DATA +# +_app_config: AppConfig +_client: AsyncClient +_task: asyncio.Task | None = None +_stop: asyncio.Event | None = None +_since: str | None = None +_last_since_save_time: float = 0 + + + +# +# CALLBACKS +# +async def _sync_callback(response: SyncResponse): + global _since, _last_since_save_time + t = time.time() + _since = response.next_batch + if t - _last_since_save_time >= 120.0: + await util.set_next_batch(_app_config, _since) + _last_since_save_time = t + + + +# +# PRIVATE +# +async def _bot_login_using_token(token: str, device_id: str) -> bool: + """Tries to login using access_token. Returns True on success.""" + util.log_info("Authorizing using access_token...") + _client.restore_login( + user_id=f"@{_app_config.matrix_user}:{util.get_hostname_from_url(_app_config.matrix_homeserver)}", + device_id=device_id, + access_token=token + ) + result = await _client.whoami() + if type(result) is not WhoamiResponse: + return False + util.log_info(f"Logged in as {result.user_id} using access_token") + return True + +async def _bot_login_using_password(password: str) -> tuple[str, str] | tuple[None, None]: + """Tries to login using password. Returns (access_token, device_id) on success.""" + util.log_info("Authorizing using password...") + result = await _client.login(password=password) + if type(result) is LoginResponse: + util.log_info(f"Authorized using password") + return result.access_token, result.device_id + elif type(result) is LoginError: + util.log_error(f"Failed to authorize using password: {result.message}") + return None, None + else: + raise RuntimeError(f"Invalid login result: {result}") + +async def _bot_login(config: AppConfig) -> bool: + """Tries to login""" + try: + # get the session token and try to use it + session_token, device_id = await util.get_session_data(config) + if session_token is not None and device_id is not None: + if await _bot_login_using_token(session_token, device_id): + return True + await util.set_session_data(config, None) + util.log_warning("Existing access_token is deleted") + # get the password and try to use it + password = await util.get_password() + if password is None: + util.log_error("No password provided (consider using MATRIX_PASSWORD environment variable)") + return False + access_token, device_id = await _bot_login_using_password(password) + if access_token is not None and device_id is not None: + await util.set_session_data(config, (access_token, device_id)) + util.log_warning("Saved new access_token and device_id") + return True + # can't login + return False + except asyncio.CancelledError: + return False + except: + traceback.print_exc() + return False + +async def _bot_loop(config: AppConfig) -> None: + """Bot loop""" + global _client, _since + # app stop task + if _stop is None: + raise RuntimeError("_stop can't be None") + stop_task = asyncio.create_task(_stop.wait()) + # setup the callback for syncing + _client.add_response_callback(_sync_callback, SyncResponse) # type: ignore + # login + login_task = asyncio.create_task(_bot_login(config)) + done, _ = await asyncio.wait( + [login_task, stop_task], + return_when=asyncio.FIRST_COMPLETED + ) + # stopped + if stop_task in done: + login_task.cancel() + return + # failed to login + if login_task.exception() or not login_task.result(): + util.request_app_stop("Can't authorize into matrix") + return + # load initial `next_batch` + _since = await util.get_next_batch(config) + # sync forever + while True: + # sync + sync_task = asyncio.create_task(_client.sync_forever(timeout=5000, since=_since)) + done, _ = await asyncio.wait( + [sync_task, stop_task], + return_when=asyncio.FIRST_COMPLETED + ) + # stopped + if stop_task in done: + util.log_info("Stopping sync_forever...") + _client.stop_sync_forever() + util.log_info("Waiting for sync_forever to quit...") + await sync_task + sync_task.cancel() + break + # something happened + try: + sync_task.result() + except: + traceback.print_exc() + await asyncio.sleep(1) + + + +# +# PUBLIC +# +async def start(config: AppConfig) -> bool: + """Starts the bot""" + global _client + global _task, _stop, _app_config + if _task is not None: + return False + _app_config = config + # create the bot + store_dir = Path.cwd() / config.store_dir + store_dir.mkdir(parents=True, exist_ok=True) + client_config = AsyncClientConfig( + store_name="storefile", + encryption_enabled=True, + store_sync_tokens=False + ) + _client = AsyncClient( + homeserver=config.matrix_homeserver, + user=config.matrix_user, + store_path=str(store_dir), + config=client_config + ) + _stop = asyncio.Event() + _task = asyncio.create_task(_bot_loop(config)) + return True + +def get_client() -> AsyncClient: + return _client + +async def stop() -> None: + """Stop the bot""" + global _task, _stop + if _task is None or _stop is None: + return + _stop.set() + try: + await _task + except asyncio.CancelledError: + pass + except: + traceback.print_exc() + _task = None + _stop = None + try: + await _client.close() + except: + pass + await util.set_next_batch(_app_config, _since) + diff --git a/config.py b/config.py new file mode 100644 index 0000000..b059af0 --- /dev/null +++ b/config.py @@ -0,0 +1,53 @@ +"""Module for generating and reading config""" + +import json +import aiofiles +import os +import traceback + +from datatypes import AppConfig +import util + +DEFAULT_CONFIG = { + "matrix_homeserver": "https://matrix.domain.net", + "matrix_user": "short_username", + "store_dir": "session_storage" +} + + + +# +# PRIVATE +# +async def write_default_config(path: str) -> None: + try: + async with aiofiles.open(path, "w") as f: + await f.write(json.dumps(DEFAULT_CONFIG, indent=4)) + util.log_info(f"Default config is saved to `{path}`") + except: + traceback.print_exc() + + + +# +# PUBLIC +# +async def load_config(path: str = "config.json") -> AppConfig | None: + """Load config from specified path. Creates default config if it does not exist (but still returns None).""" + # no file + if not os.path.isfile(path): + await write_default_config(path) + return None + # open and read + try: + async with aiofiles.open(path, "r") as f: + j = json.loads((await f.read()).strip()) + except: + traceback.print_exc() + return None + # parse + try: + cfg = AppConfig(**j) + return cfg + except: + return None \ No newline at end of file diff --git a/database.py b/database.py new file mode 100644 index 0000000..e69de29 diff --git a/datatypes.py b/datatypes.py new file mode 100644 index 0000000..838a7b9 --- /dev/null +++ b/datatypes.py @@ -0,0 +1,15 @@ +""" +Types used across the application +""" + +from dataclasses import dataclass +from enum import Enum + +@dataclass +class AppConfig: + matrix_homeserver: str + matrix_user: str + store_dir: str + +class MessageType(Enum): + TEXT = "m.text" \ No newline at end of file diff --git a/logic.py b/logic.py new file mode 100644 index 0000000..f89506e --- /dev/null +++ b/logic.py @@ -0,0 +1,148 @@ +"""Main logic implementation""" + +import traceback +from typing import Any +import html + +from nio import AsyncClient +from nio import JoinResponse, RoomSendResponse, RoomSendError +from nio import MatrixInvitedRoom, InviteMemberEvent +from nio import MatrixRoom, RoomMessageText + +from nio import OlmUnverifiedDeviceError + +from datatypes import * +import util + + +# +# DATA +# +_client: AsyncClient + + + +# +# PRIVATE +# +async def _verify_all_devices() -> None: + """Verifies all known devices""" + for user_id in _client.device_store.users: + for device_id, olm_device in _client.device_store[user_id].items(): + # can't trust ourselves + if device_id == _client.device_id and user_id == _client.user_id: + continue + # they are already verified + if olm_device.verified: + continue + # verify them + _client.verify_device(olm_device) + +def _handle_html_in_kwargs(kwargs: dict[str, Any]) -> None: + """Modified `kwargs` in-place so that `formatted_body` appears if needed""" + if "formatted_body" in kwargs or "body" not in kwargs: + return + is_html, text_without_html = util.check_and_remove_html(kwargs["body"]) + if not is_html: + return + kwargs["format"] = "org.matrix.custom.html" + kwargs["formatted_body"] = kwargs["body"] + kwargs["body"] = text_without_html + +async def _send_message_to(room_id: str, message_type: MessageType, **kwargs) -> str: + """Sends a message to the room and returns event_id. + + This function automatically detects `body` key in `kwargs` and checks + if it is a valid HTML. If it is a valid HTML, it will send it as such. + Moreover, `body` attribute will be cleaned from any HTML tags, so that + the text will be looking well. `formatted_body` attribute is added + automatically and you should not add it manually. + """ + try: + # handle HTML + _handle_html_in_kwargs(kwargs) + # try to send the message + result = await _client.room_send( + room_id=room_id, + message_type="m.room.message", + content={ + "msgtype": message_type.value, + **kwargs + } + ) + # success + if type(result) is RoomSendResponse: + return result.event_id + # error + elif type(result) is RoomSendError: + raise Exception(result) + # unknown error + else: + raise RuntimeError() + except OlmUnverifiedDeviceError as e: + # verify everyone and retry + await _verify_all_devices() + return await _send_message_to(room_id, message_type, **kwargs) + except: + raise + +async def _send_text_to(room_id: str, text: str) -> str: + """Sends a text message to the room. `text` may be HTML""" + return await _send_message_to( + room_id=room_id, + message_type=MessageType.TEXT, + body=text + ) + + + +# +# CALLBACKS +# +async def _message_callback(room: MatrixRoom, event: RoomMessageText) -> None: + """Handle commands received from Matrix""" + try: + # do not process messages sent by ourselves + if event.sender == _client.user_id: + return + # prepare response + response = "Получено сообщение" + response += f"

Room ID: {html.escape(room.room_id)}" + response += f"
Sender: {html.escape(event.sender)}" + print(await _send_text_to(room.room_id, response)) + except: + traceback.print_exc() + +async def _invite_callback(room: MatrixInvitedRoom, event: InviteMemberEvent) -> None: + """Happens when the bot is invited to somewhere""" + try: + result = await _client.join(room.room_id) + if type(result) is JoinResponse: + util.log_info(f"Joined the room {room.room_id}") + else: + util.log_error(f"Can't join room {room.room_id}") + except: + traceback.print_exc() + +async def _generic_test_callback(*args, **kwargs) -> None: + """Use this callback to check argument types""" + print("GENERIC TEST CALLBACK") + for a in args: + print(f" - {type(a)}") + for k in kwargs: + print(f" * {k} = {kwargs[k]}") + + + +# +# PUBLIC +# +async def setup(client: AsyncClient) -> None: + global _client + _client = client + client.add_event_callback(_message_callback, RoomMessageText) # type: ignore + client.add_event_callback(_invite_callback, InviteMemberEvent) # type: ignore + +async def stop() -> None: + """Stop all ongoing processes""" + pass diff --git a/main.py b/main.py new file mode 100644 index 0000000..561ea80 --- /dev/null +++ b/main.py @@ -0,0 +1,52 @@ +""" +2026-matrix-downloader entry point +""" + +import asyncio +import traceback +import signal + +import config +import util +import bot +import logic +from datatypes import AppConfig + +async def main() -> None: + """Entry point""" + # setup signal handler + util.setup_app_stop_event() + loop = asyncio.get_running_loop() + for sig in (signal.SIGINT, signal.SIGTERM): + loop.add_signal_handler(sig, lambda: util.request_app_stop("Signal caught")) + + # setup + util.setup_logging() + util.log_info("2026-matrix-csonac") + + # load the config + cfg: AppConfig = await config.load_config("config.json") # type: ignore + if cfg is None: + util.log_error("Could't load config") + return + + # start the bot + if not await bot.start(cfg): + util.log_error("Could't start the bot") + return + await logic.setup(bot.get_client()) + + # wait for stop + await util.get_app_stop_event().wait() + + # stop + await logic.stop() + await bot.stop() + +if __name__ == "__main__": + try: + asyncio.run(main()) + except SystemExit: + raise + except: + traceback.print_exc() diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..14573dc --- /dev/null +++ b/requirements.txt @@ -0,0 +1,3 @@ +aioconsole +aiofiles +matrix-nio[e2e] diff --git a/util.py b/util.py new file mode 100644 index 0000000..b26c8a3 --- /dev/null +++ b/util.py @@ -0,0 +1,164 @@ +""" +Utilities +""" + +import asyncio +import sys +import os +import json +import logging +import traceback +from html.parser import HTMLParser +from pathlib import Path +from urllib.parse import urlparse + +import aiofiles +import aioconsole + +from datatypes import AppConfig + + + +# +# PRIVATE +# +_stop_event: asyncio.Event + + + +# +# PUBLIC +# +def setup_logging() -> None: + """Setup logging""" + logging.basicConfig( + format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", + datefmt="%Y-%m-%d %H:%M:%S", + level=logging.INFO + ) + logging.getLogger("nio").setLevel(logging.CRITICAL + 1) + +def setup_app_stop_event() -> None: + """Setup app stop asyncio event""" + global _stop_event + _stop_event = asyncio.Event() + +def request_app_stop(reason: str) -> None: + try: + log_info(f"{reason}. The application is stopping.") + _stop_event.set() + except: + traceback.print_exc() + +def get_app_stop_event() -> asyncio.Event: + return _stop_event + +def is_terminal_interactive() -> bool: + """Returns True if the terminal is interactive""" + return sys.stdin.isatty() and sys.stdout.isatty() + +def log_info(text: str) -> None: + logging.info(text) + +def log_warning(text: str) -> None: + logging.warning(text) + +def log_error(text: str) -> None: + logging.error(text) + +async def ainput(text: str = "") -> str: + # simulate empty input if not TTY + if not is_terminal_interactive(): + return "" + return await aioconsole.ainput(text) + +async def set_next_batch(config: AppConfig, next_batch: str | None) -> bool: + """Save `next_batch` to session directory.""" + try: + next_batch_file_path = Path.cwd() / config.store_dir / "next_batch.txt" + if next_batch is None: + next_batch_file_path.unlink(True) + return True + async with aiofiles.open(next_batch_file_path, "w") as f: + await f.write(next_batch) + return True + except: + return False + +async def get_next_batch(config: AppConfig) -> str | None: + """Get `next_batch`""" + try: + next_batch_file_path = Path.cwd() / config.store_dir / "next_batch.txt" + async with aiofiles.open(next_batch_file_path, "r") as f: + next_batch = (await f.read()).strip() + if not next_batch: + next_batch = None + return next_batch + except: + return None + +async def get_session_data(config: AppConfig) -> tuple[str, str] | tuple[None, None]: + """ + Get (access_token, device_id) or (None, None) + """ + try: + token_file_path = Path.cwd() / config.store_dir / "auth.json" + async with aiofiles.open(token_file_path, "r") as f: + j = json.loads((await f.read()).strip()) + return j["access_token"], j["device_id"] + except: + return None, None + +async def set_session_data(config: AppConfig, token_device_pair: tuple[str, str] | None) -> bool: + """ + Set new (access_token, device_id) pair; use None to remove it. + + Returns: + True on success + """ + try: + token_file_path = Path.cwd() / config.store_dir / "auth.json" + if token_device_pair is None: + token_file_path.unlink(True) + return True + async with aiofiles.open(token_file_path, "w") as f: + await f.write(json.dumps({"access_token": token_device_pair[0], "device_id": token_device_pair[1]})) + return True + except: + return False + +async def get_password() -> str | None: + if "MATRIX_PASSWORD" in os.environ: + return os.environ["MATRIX_PASSWORD"] + if not is_terminal_interactive(): + return None + return await ainput("Matrix password: ") + +def get_hostname_from_url(url: str) -> str | None: + """Returns `matrix.domain.net` for `https://matrix.domain.net/bla/bla/bla`""" + try: + return urlparse(url).hostname + except: + return None + +def check_and_remove_html(possible_html: str) -> tuple[bool, str]: + """Checks if `possible_html` is a valid HTML text and returns (is_html, text_without_tags)""" + has_tags = False + text_fragments = [] + class Extractor(HTMLParser): + def handle_starttag(self, tag, attrs): + nonlocal has_tags + has_tags = True + def handle_data(self, data): + text_fragments.append(data) + + parser = Extractor(convert_charrefs=True) + parser.feed(possible_html) + + try: + if has_tags: + return (True, " ".join("".join(text_fragments).split())) + except: + traceback.print_exc() + + return (False, possible_html) \ No newline at end of file