commit 2ae1e3e0fbe7ce3c6548ca8ba6dd032b5beb8f16 Author: nikita Date: Thu Sep 24 19:32:57 2026 +0300 Initial commit (codex) diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..1a451a6 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,8 @@ +.git +.venv +__pycache__ +*.py[cod] +config.json +runtime +session_storage +work diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..0376696 --- /dev/null +++ b/.gitignore @@ -0,0 +1,7 @@ +.venv/ +__pycache__/ +*.py[cod] +config.json +runtime/ +session_storage/ +work/ diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..e494319 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,27 @@ +FROM python:3.13-slim + +ENV PYTHONDONTWRITEBYTECODE=1 \ + PYTHONUNBUFFERED=1 \ + PIP_DISABLE_PIP_VERSION_CHECK=1 + +RUN apt-get update \ + && apt-get install --no-install-recommends -y \ + ca-certificates \ + docker-cli \ + git \ + libmagic1 \ + libolm3 \ + && rm -rf /var/lib/apt/lists/* + +WORKDIR /usr/src/app + +COPY requirements.txt ./ +RUN python -m pip install --no-cache-dir -r requirements.txt + +COPY bot.py config.py main.py run_sumka.sh ./ +RUN chmod +x /usr/src/app/run_sumka.sh + +WORKDIR /runtime +STOPSIGNAL SIGTERM + +CMD ["python", "/usr/src/app/main.py"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..ef39168 --- /dev/null +++ b/README.md @@ -0,0 +1,90 @@ +# 2026-matrix-sumka + +Matrix-бот для запуска +[`2026-linux-sumka`](https://git.tyukalov.su/nikita/2026-linux-sumka.git). +Бот принимает поддерживаемые аудио- и видеофайлы, запускает `sumka` в отдельном +Docker-контейнере и отправляет пользователю готовый `output.pdf`. + +Одновременно выполняется только одно задание. Новый файл во время обработки +сразу отклоняется и не ставится в очередь. + +Поддерживаемые расширения: `.mp4`, `.mkv`, `.avi`, `.mp3`, `.m4a`, `.wav`. +Файл можно отправить в Matrix как видео, аудио или обычный файл. + +## Как это работает + +Для каждого принятого видео бот: + +1. очищает `runtime/work`; +2. скачивает файл как `runtime/work/input.<расширение>`; +3. запускает контейнер `sumka` через `run_sumka.sh`; +4. читает этапы обработки из `report.txt`; +5. проверяет код выхода процесса; +6. отправляет `output.pdf` в исходную Matrix-комнату. + +Вывод контейнера сохраняется в `runtime/work/sumka.log`. + +## Конфигурация + +При первом запуске автоматически создаётся `runtime/config.json`: + +```json +{ + "matrix_homeserver": "https://matrix.example.org", + "matrix_user": "sumka-bot", + "store_dir": "session_storage", + "sumka_image": "sumka:amd", + "ai_api_key": "change-me", + "sumka_args": [] +} +``` + +`matrix_user` — локальная часть Matrix ID без `@` и имени сервера. +В `sumka_args` можно передать аргументы командной строки `sumka`, например имена +моделей. Для NVIDIA укажите образ с `nvidia` в теге, например `sumka:nvidia`; +для AMD — `sumka:amd`. + +## Запуск в Docker + +На машине заранее должны быть собраны нужный образ `sumka` и образ бота. + +```bash +mkdir -p runtime +docker compose build +docker compose run --rm bot +``` + +Первый запуск создаст конфиг и завершится. Заполните +`runtime/config.json`, затем повторите команду для первой авторизации в Matrix. +После успешного входа остановите процесс через `Ctrl+C`; сессия сохранится в +`runtime/session_storage`. + +Пароль также можно передать через окружение: + +```bash +MATRIX_PASSWORD='matrix-password' \ + docker compose run --rm -e MATRIX_PASSWORD bot +``` + +Постоянный запуск: + +```bash +docker compose up -d +``` + +Бот использует Docker daemon хоста через `/var/run/docker.sock`; вложенный +Docker daemon не запускается. + +## Запуск через `.venv` + +Нужны Python 3.12+, Docker CLI, `libmagic` и `libolm`. + +```bash +python3 -m venv .venv +. .venv/bin/activate +python -m pip install -r requirements.txt +cp config.example.json config.json +python main.py +``` + +Используется стабильная версия `mab` `v0.5.2`. diff --git a/bot.py b/bot.py new file mode 100644 index 0000000..4a54d1c --- /dev/null +++ b/bot.py @@ -0,0 +1,335 @@ +"""Matrix callbacks and the single-job sumka supervisor.""" + +from __future__ import annotations + +import asyncio +import json +import logging +import os +import shutil +from pathlib import Path +from typing import Any + +from mab import ( + CTX_FILE_NAME, + EventContext, + MatrixBot, + MessageHasFile, + MessageType, + MessageTypeFilter, + NewMessageFilter, + SenderIsBotFilter, +) + +from config import AppConfig + + +LOGGER = logging.getLogger(__name__) + +STEP_NAMES = { + "media_separation": "подготовка аудио и видео", + "voice_recognition": "распознавание речи", + "asr_filter": "обработка распознанного текста", + "asr_events": "построение событий по речи", + "video_references": "поиск ссылок на видеоряд", + "reference_resolver": "анализ кадров", + "structure_builder": "построение структуры конспекта", + "structure_refiner": "уточнение структуры конспекта", + "markdown_builder": "сборка Markdown", + "pdf_builder": "сборка PDF", +} + +TERMINAL_REPORT_STATUSES = {"succeeded", "failed", "cancelled"} +SUPPORTED_INPUT_EXTENSIONS = {".mp4", ".mkv", ".avi", ".mp3", ".m4a", ".wav"} + + +class SumkaBotController: + """Accept one supported media file at a time and supervise its conversion.""" + + def __init__(self, bot: MatrixBot, config: AppConfig) -> None: + self._bot = bot + self._config = config + self._state_lock = asyncio.Lock() + self._active_task: asyncio.Task[None] | None = None + self._stopping = False + + def setup_callback(self) -> None: + media_filter = ( + ~SenderIsBotFilter() + & NewMessageFilter() + & MessageTypeFilter( + [MessageType.VIDEO, MessageType.AUDIO, MessageType.FILE] + ) + & MessageHasFile() + ) + self._bot.add_callback(media_filter, self.on_media) + + async def on_media(self, context: EventContext) -> None: + """Handle supported media without ever queueing it behind another job.""" + filename = context[CTX_FILE_NAME] + extension = Path(filename).suffix.lower() + if extension not in SUPPORTED_INPUT_EXTENSIONS: + supported = ", ".join(sorted(SUPPORTED_INPUT_EXTENSIONS)) + await self._safe_send_text( + context.room.room_id, + f"Неподдерживаемый формат файла. Поддерживаются: {supported}.", + ) + return + + current_task = asyncio.current_task() + if current_task is None: + LOGGER.error("Media callback is running outside an asyncio task") + return + + async with self._state_lock: + busy = self._active_task is not None + stopping = self._stopping + if not busy and not stopping: + self._active_task = current_task + + if stopping: + return + if busy: + await self._safe_send_text( + context.room.room_id, + "Сейчас уже обрабатывается другой файл. " + "Новый запрос не поставлен в очередь; отправьте его позже.", + ) + return + + try: + await self._process_media(context, extension) + except asyncio.CancelledError: + LOGGER.info("Active media job was cancelled") + raise + except Exception as error: + LOGGER.exception("Unexpected failure while processing media") + await self._safe_send_text( + context.room.room_id, + f"Не удалось обработать файл: {type(error).__name__}: {error}", + ) + finally: + async with self._state_lock: + if self._active_task is current_task: + self._active_task = None + + async def shutdown(self) -> None: + """Stop accepting jobs and cancel the current runner, if any.""" + async with self._state_lock: + self._stopping = True + task = self._active_task + + current_task = asyncio.current_task() + if task is not None and task is not current_task and not task.done(): + task.cancel() + await asyncio.gather(task, return_exceptions=True) + + async def _process_media( + self, context: EventContext, extension: str + ) -> None: + room_id = context.room.room_id + await asyncio.to_thread(self._recreate_work_directory) + + input_path = self._config.work_dir / f"input{extension}" + await self._safe_send_text(room_id, "Начинаю скачивание файла…") + await self._bot.download_file(context, path=input_path) + await self._safe_send_text( + room_id, + "Файл скачан. Запускаю построение конспекта…", + ) + + exit_code, terminal_report = await self._run_sumka(room_id) + output_path = self._config.work_dir / "output.pdf" + + if exit_code != 0: + details = self._report_error_details(terminal_report) + message = f"Обработка завершилась с ошибкой (код {exit_code})." + if details: + message += f" {details}" + await self._safe_send_text(room_id, message) + return + + if not output_path.is_file(): + await self._safe_send_text( + room_id, + "Обработка завершилась без ошибки, но output.pdf не был создан.", + ) + return + + try: + await self._bot.send_file( + room=room_id, + path=output_path, + filename="output.pdf", + mime_type="application/pdf", + text="Готовый конспект", + is_html=False, + timeout=self._config.upload_timeout_seconds, + ) + except Exception as error: + LOGGER.exception("Could not upload output.pdf") + await self._safe_send_text( + room_id, + "Конспект готов, но отправить output.pdf не удалось: " + f"{type(error).__name__}: {error}", + ) + + def _recreate_work_directory(self) -> None: + work_dir = self._config.work_dir + if work_dir.is_symlink() or work_dir.is_file(): + work_dir.unlink() + elif work_dir.exists(): + shutil.rmtree(work_dir) + work_dir.mkdir(parents=True) + + async def _run_sumka( + self, room_id: str + ) -> tuple[int, dict[str, Any] | None]: + env = os.environ.copy() + env.update( + { + "AI_API_KEY": self._config.ai_api_key, + "SUMKA_IMAGE": self._config.sumka_image, + "SUMKA_CONTAINER_NAME": self._config.sumka_container_name, + } + ) + command = [ + str(self._config.runner_script), + str(self._config.work_dir), + *self._config.sumka_args, + ] + + log_path = self._config.work_dir / "sumka.log" + process: asyncio.subprocess.Process | None = None + with log_path.open("wb", buffering=0) as log_file: + process = await asyncio.create_subprocess_exec( + *command, + cwd=self._config.runtime_dir, + env=env, + stdout=log_file, + stderr=asyncio.subprocess.STDOUT, + ) + try: + return await self._monitor_process(process, room_id) + except asyncio.CancelledError: + await self._terminate_process(process) + raise + + async def _monitor_process( + self, + process: asyncio.subprocess.Process, + room_id: str, + ) -> tuple[int, dict[str, Any] | None]: + report_path = self._config.work_dir / "report.txt" + offset = 0 + seen_steps: set[tuple[str | None, str]] = set() + terminal_report: dict[str, Any] | None = None + + while True: + records, offset = await asyncio.to_thread( + self._read_report_records, report_path, offset + ) + for record in records: + status = record.get("status") + if status in TERMINAL_REPORT_STATUSES: + terminal_report = record + continue + if status != "running": + continue + + step = record.get("step") + if not isinstance(step, str): + continue + step_key = (record.get("run_id"), step) + if step_key in seen_steps: + continue + seen_steps.add(step_key) + + label = STEP_NAMES.get(step, step) + index = record.get("step_index") + total = record.get("steps_total") + if isinstance(index, int) and isinstance(total, int): + text = f"Этап {index}/{total}: {label}." + else: + text = f"Текущий этап: {label}." + await self._safe_send_text(room_id, text) + + if process.returncode is not None: + break + try: + async with asyncio.timeout( + self._config.report_poll_interval_seconds + ): + await process.wait() + except TimeoutError: + pass + + # Consume status records written immediately before process exit. + records, offset = await asyncio.to_thread( + self._read_report_records, report_path, offset + ) + for record in records: + if record.get("status") in TERMINAL_REPORT_STATUSES: + terminal_report = record + + return await process.wait(), terminal_report + + @staticmethod + def _read_report_records( + report_path: Path, offset: int + ) -> tuple[list[dict[str, Any]], int]: + if not report_path.is_file(): + return [], offset + + records: list[dict[str, Any]] = [] + new_offset = offset + with report_path.open("rb") as report_file: + report_file.seek(offset) + while True: + line_start = report_file.tell() + line = report_file.readline() + if not line: + break + if not line.endswith(b"\n"): + new_offset = line_start + break + new_offset = report_file.tell() + try: + record = json.loads(line.decode("utf-8")) + except (UnicodeDecodeError, json.JSONDecodeError): + LOGGER.warning("Ignoring malformed report.txt line") + continue + if isinstance(record, dict): + records.append(record) + return records, new_offset + + @staticmethod + def _report_error_details(report: dict[str, Any] | None) -> str: + if not report: + return "Подробности сохранены в sumka.log." + + error_name = report.get("error") + error_message = report.get("message") + if isinstance(error_name, str) and isinstance(error_message, str): + return f"{error_name}: {error_message}" + if report.get("status") == "cancelled": + return "Задание было отменено." + return "Подробности сохранены в sumka.log." + + @staticmethod + async def _terminate_process(process: asyncio.subprocess.Process) -> None: + if process.returncode is not None: + return + process.terminate() + try: + async with asyncio.timeout(45): + await process.wait() + except TimeoutError: + process.kill() + await process.wait() + + async def _safe_send_text(self, room_id: str, text: str) -> None: + try: + await self._bot.send_text(room_id, text, is_html=False) + except Exception: + LOGGER.exception("Could not send a Matrix status message") diff --git a/config.example.json b/config.example.json new file mode 100644 index 0000000..229763e --- /dev/null +++ b/config.example.json @@ -0,0 +1,8 @@ +{ + "matrix_homeserver": "https://matrix.example.org", + "matrix_user": "sumka-bot", + "store_dir": "session_storage", + "sumka_image": "sumka:amd", + "ai_api_key": "change-me", + "sumka_args": [] +} diff --git a/config.py b/config.py new file mode 100644 index 0000000..193a976 --- /dev/null +++ b/config.py @@ -0,0 +1,136 @@ +"""Loading and validation of the bot's JSON configuration.""" + +from __future__ import annotations + +import json +import os +from dataclasses import dataclass +from pathlib import Path +from typing import Any + + +DEFAULT_CONFIG: dict[str, Any] = { + "matrix_homeserver": "https://matrix.example.org", + "matrix_user": "sumka-bot", + "store_dir": "session_storage", + "sumka_image": "sumka:amd", + "ai_api_key": "change-me", + "sumka_args": [], +} + +SUMKA_CONTAINER_NAME = "matrix-sumka-job" +REPORT_POLL_INTERVAL_SECONDS = 2.0 +UPLOAD_TIMEOUT_SECONDS = 3600.0 + + +class ConfigError(RuntimeError): + """Raised when config.json cannot be used safely.""" + + +@dataclass(frozen=True, slots=True) +class AppConfig: + matrix_homeserver: str + matrix_user: str + store_dir: Path + sumka_image: str + ai_api_key: str + sumka_args: tuple[str, ...] + runtime_dir: Path + + @property + def work_dir(self) -> Path: + return self.runtime_dir / "work" + + @property + def runner_script(self) -> Path: + return Path(__file__).resolve().parent / "run_sumka.sh" + + @property + def sumka_container_name(self) -> str: + return SUMKA_CONTAINER_NAME + + @property + def report_poll_interval_seconds(self) -> float: + return REPORT_POLL_INTERVAL_SECONDS + + @property + def upload_timeout_seconds(self) -> float: + return UPLOAD_TIMEOUT_SECONDS + + +def _non_empty_string(data: dict[str, Any], key: str) -> str: + value = data.get(key) + if not isinstance(value, str) or not value.strip(): + raise ConfigError(f"`{key}` must be a non-empty string") + return value.strip() + + +def _runtime_path(value: str, runtime_dir: Path) -> Path: + path = Path(value).expanduser() + if not path.is_absolute(): + path = runtime_dir / path + return path.resolve() + + +def _parse_config(data: Any, config_path: Path) -> AppConfig: + if not isinstance(data, dict): + raise ConfigError("the config root must be a JSON object") + + expected = set(DEFAULT_CONFIG) + missing = expected - set(data) + unknown = set(data) - expected + if missing: + raise ConfigError(f"missing config keys: {', '.join(sorted(missing))}") + if unknown: + raise ConfigError(f"unknown config keys: {', '.join(sorted(unknown))}") + + args = data.get("sumka_args") + if not isinstance(args, list) or any(not isinstance(arg, str) for arg in args): + raise ConfigError("`sumka_args` must be an array of strings") + + api_key = _non_empty_string(data, "ai_api_key") + if api_key == DEFAULT_CONFIG["ai_api_key"]: + raise ConfigError("replace the placeholder value in `ai_api_key`") + + runtime_dir = config_path.parent.resolve() + store_dir = _runtime_path(_non_empty_string(data, "store_dir"), runtime_dir) + runner_script = Path(__file__).resolve().parent / "run_sumka.sh" + if not runner_script.is_file(): + raise ConfigError(f"runner script does not exist: {runner_script}") + if not os.access(runner_script, os.X_OK): + raise ConfigError(f"runner script is not executable: {runner_script}") + + return AppConfig( + matrix_homeserver=_non_empty_string(data, "matrix_homeserver"), + matrix_user=_non_empty_string(data, "matrix_user"), + store_dir=store_dir, + sumka_image=_non_empty_string(data, "sumka_image"), + ai_api_key=api_key, + sumka_args=tuple(args), + runtime_dir=runtime_dir, + ) + + +def write_default_config(path: str | Path = "config.json") -> Path: + """Write an editable default config and return its absolute path.""" + config_path = Path(path).resolve() + config_path.parent.mkdir(parents=True, exist_ok=True) + config_path.write_text( + json.dumps(DEFAULT_CONFIG, ensure_ascii=False, indent=4) + "\n", + encoding="utf-8", + ) + return config_path + + +def load_config(path: str | Path = "config.json") -> AppConfig | None: + """Load config, or create a template and return None on the first run.""" + config_path = Path(path).resolve() + if not config_path.is_file(): + write_default_config(config_path) + return None + + try: + data = json.loads(config_path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError) as error: + raise ConfigError(f"cannot read {config_path}: {error}") from error + return _parse_config(data, config_path) diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..1f2c09f --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,11 @@ +services: + bot: + build: . + image: 2026-matrix-sumka:local + container_name: matrix-sumka-bot + restart: unless-stopped + stdin_open: true + tty: true + volumes: + - ./runtime:/runtime + - /var/run/docker.sock:/var/run/docker.sock diff --git a/main.py b/main.py new file mode 100644 index 0000000..663573e --- /dev/null +++ b/main.py @@ -0,0 +1,70 @@ +"""2026-matrix-sumka entry point.""" + +from __future__ import annotations + +import asyncio +import logging +import signal +from pathlib import Path + +from mab import MatrixBot, MatrixBotConfig + +from bot import SumkaBotController +from config import ConfigError, load_config + + +LOGGER = logging.getLogger(__name__) + + +async def run() -> int: + config_path = Path("config.json").resolve() + config = load_config(config_path) + if config is None: + LOGGER.error("Created %s; edit it and start the bot again", config_path) + return 2 + + matrix_config = MatrixBotConfig( + matrix_homeserver_url=config.matrix_homeserver, + matrix_username_localpart=config.matrix_user, + storage_directory=config.store_dir, + ) + matrix_bot = MatrixBot(matrix_config) + controller = SumkaBotController(matrix_bot, config) + controller.setup_callback() + + stop_event = asyncio.Event() + loop = asyncio.get_running_loop() + for caught_signal in (signal.SIGINT, signal.SIGTERM): + loop.add_signal_handler(caught_signal, stop_event.set) + + started = False + try: + await matrix_bot.start() + started = True + LOGGER.info("Bot started") + await stop_event.wait() + finally: + await controller.shutdown() + if started: + await matrix_bot.stop() + LOGGER.info("Bot stopped") + return 0 + + +def main() -> int: + logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(levelname)s %(name)s: %(message)s", + ) + logging.getLogger("nio").setLevel(logging.WARNING) + try: + return asyncio.run(run()) + except ConfigError as error: + LOGGER.error("Invalid config: %s", error) + return 2 + except KeyboardInterrupt: + return 130 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..20fe18f --- /dev/null +++ b/requirements.txt @@ -0,0 +1 @@ +mab @ git+https://git.tyukalov.su/nikita/mab@v0.5.2 diff --git a/run_sumka.sh b/run_sumka.sh new file mode 100755 index 0000000..ea9a4a1 --- /dev/null +++ b/run_sumka.sh @@ -0,0 +1,131 @@ +#!/usr/bin/env bash + +set -Eeuo pipefail + +if [[ $# -lt 1 ]]; then + echo "Usage: $0 WORK_DIR [SUMKA_ARGUMENT ...]" >&2 + exit 64 +fi + +work_dir=$1 +shift + +if [[ ! -d "$work_dir" ]]; then + echo "Work directory does not exist: $work_dir" >&2 + exit 66 +fi + +input_found=false +for extension in mp4 mkv avi mp3 m4a wav; do + if [[ -f "$work_dir/input.$extension" ]]; then + input_found=true + break + fi +done +if [[ "$input_found" != true ]]; then + echo "No supported input file found in: $work_dir" >&2 + exit 66 +fi + +: "${AI_API_KEY:?AI_API_KEY is required}" + +sumka_image=${SUMKA_IMAGE:-sumka:amd} +container_name=${SUMKA_CONTAINER_NAME:-matrix-sumka-job} +cache_volume=${SUMKA_CACHE_VOLUME:-whisper-cache} +tmpfs_size=${SUMKA_TMPFS_SIZE:-8g} + +docker_args=( + run + --rm + --name "$container_name" + --env AI_API_KEY + --mount "type=volume,src=${cache_volume},dst=/tmp/sumka-cache/whisper" + --tmpfs "/tmp:rw,size=${tmpfs_size}" +) + +gpu_kind=${SUMKA_GPU:-auto} +if [[ "$gpu_kind" == auto ]]; then + case "$sumka_image" in + *nvidia*) gpu_kind=nvidia ;; + *amd*) gpu_kind=amd ;; + *) gpu_kind=amd ;; + esac +fi + +case "$gpu_kind" in + nvidia) + docker_args+=(--gpus all) + ;; + amd) + docker_args+=( + --device=/dev/kfd + --device=/dev/dri + --group-add video + --security-opt seccomp=unconfined + ) + ;; + *) + echo "Unsupported SUMKA_GPU value: $gpu_kind (expected amd or nvidia)" >&2 + exit 64 + ;; +esac + +volumes_from=${SUMKA_VOLUMES_FROM:-} +if [[ -z "$volumes_from" && -f /.dockerenv ]]; then + # Docker uses the short container ID as the default hostname. Unlike a + # Compose service name, it also works for `docker compose run` containers. + volumes_from=${HOSTNAME:-} +fi + +if [[ -n "$volumes_from" ]]; then + # The sibling sumka container sees the bot's /runtime bind mount at the + # same path, independently of how the bot container itself was named. + docker_args+=( + --volumes-from "${volumes_from}:rw" + --workdir "$work_dir" + ) +else + # Direct host/.venv launch: mount only this job's directory. + work_dir=$(realpath "$work_dir") + docker_args+=( + --mount "type=bind,src=${work_dir},dst=/work" + --workdir /work + ) +fi + +signal_exit_code=0 + +stop_container() { + docker stop --time 10 "$container_name" >/dev/null 2>&1 || true +} + +on_term() { + signal_exit_code=143 + stop_container +} + +on_int() { + signal_exit_code=130 + stop_container +} + +on_hup() { + signal_exit_code=129 + stop_container +} + +trap on_term TERM +trap on_int INT +trap on_hup HUP + +set +e +docker "${docker_args[@]}" "$sumka_image" "$@" & +docker_pid=$! +wait "$docker_pid" +exit_code=$? +set -e + +if [[ $signal_exit_code -ne 0 ]]; then + exit "$signal_exit_code" +fi +exit "$exit_code"