commit 2c05443c897570b28e25e1d33212e10ab257ae3d Author: Nikita Tyukalov, ASUS, Linux Date: Mon Aug 17 03:33:54 2026 +0300 Initial commit 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