Initial commit
This commit is contained in:
7
.gitignore
vendored
Normal file
7
.gitignore
vendored
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
.venv/
|
||||||
|
session_storage/
|
||||||
|
__pycache__/
|
||||||
|
*.swp
|
||||||
|
*.swo
|
||||||
|
*.vscode
|
||||||
|
*.json
|
||||||
80
README.md
Normal file
80
README.md
Normal file
@@ -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 /<channel>/notify`
|
||||||
|
- **Описание.** Используется, чтобы отправить уведомление в указанный канал.
|
||||||
|
Вместо `<channel>` указывается код канала, получаемый при помощи команды
|
||||||
|
`!channel_code`, выполненной в комнате Matrix.
|
||||||
|
- **Тело запроса.** Тело запроса представляет собой `json` объект:
|
||||||
|
```json
|
||||||
|
{
|
||||||
|
"service": "<название сервиса, передаётся только если разрешено>",
|
||||||
|
"subservice": "<название подсервиса>",
|
||||||
|
"text": "<текст уведомления>",
|
||||||
|
"urgent": false
|
||||||
|
}
|
||||||
|
```
|
||||||
|
- **Тело ответа.** Тело ответа представляет собой `json` объект. В случае
|
||||||
|
успеха в нём будут все поля, перечисляемые ниже. В случае провала - только
|
||||||
|
поле `error`, содержащее текстовое описание ошибки.
|
||||||
|
```json
|
||||||
|
{
|
||||||
|
"error": null,
|
||||||
|
"notification_id": "<здесь будет Notification ID>"
|
||||||
|
}
|
||||||
|
```
|
||||||
|
- ``
|
||||||
202
bot.py
Normal file
202
bot.py
Normal file
@@ -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)
|
||||||
|
|
||||||
53
config.py
Normal file
53
config.py
Normal file
@@ -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
|
||||||
0
database.py
Normal file
0
database.py
Normal file
15
datatypes.py
Normal file
15
datatypes.py
Normal file
@@ -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"
|
||||||
148
logic.py
Normal file
148
logic.py
Normal file
@@ -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 = "<b>Получено сообщение</b>"
|
||||||
|
response += f"<br><br><b>Room ID:</b> {html.escape(room.room_id)}"
|
||||||
|
response += f"<br><b>Sender:</b> {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
|
||||||
52
main.py
Normal file
52
main.py
Normal file
@@ -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()
|
||||||
3
requirements.txt
Normal file
3
requirements.txt
Normal file
@@ -0,0 +1,3 @@
|
|||||||
|
aioconsole
|
||||||
|
aiofiles
|
||||||
|
matrix-nio[e2e]
|
||||||
164
util.py
Normal file
164
util.py
Normal file
@@ -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)
|
||||||
Reference in New Issue
Block a user