From d86057da5caaf3f8bf68a566523077fcf3ffb6f7 Mon Sep 17 00:00:00 2001 From: nikita Date: Thu, 24 Sep 2026 18:30:05 +0300 Subject: [PATCH] Improved integration ability --- .dockerignore | 1 + .gitignore | 1 + README.md | 56 ++++++++++++++++++---- main.py | 105 ++++++++++++++++++++++++++++++++--------- reporting.py | 128 ++++++++++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 260 insertions(+), 31 deletions(-) create mode 100644 reporting.py diff --git a/.dockerignore b/.dockerignore index d5f8487..cf0353c 100644 --- a/.dockerignore +++ b/.dockerignore @@ -9,4 +9,5 @@ input.* audio.* video.* output.* +report.txt *.json diff --git a/.gitignore b/.gitignore index 3799ed7..f9d6ea3 100644 --- a/.gitignore +++ b/.gitignore @@ -6,6 +6,7 @@ images/ debug/ output.md output.pdf +report.txt *.json *.mkv diff --git a/README.md b/README.md index afd0d39..49502a6 100644 --- a/README.md +++ b/README.md @@ -68,21 +68,59 @@ python main.py ## Запуск -TODO +Контейнер рассчитан на один запрос. Для каждого задания создайте отдельную +директорию, положите в неё один поддерживаемый файл с именем `input.*` и +смонтируйте её в `/work`: + +```bash +JOB_DIR=$(mktemp -d) +cp input.mkv "$JOB_DIR/input.mkv" + +docker run --rm --name sumka-example --gpus all -e AI_API_KEY \ + --mount type=bind,src="$JOB_DIR",dst=/work \ + --mount type=volume,src=whisper-cache,dst=/tmp/sumka-cache/whisper \ + --tmpfs /tmp:rw,size=8g sumka:nvidia +``` + +Не запускайте два контейнера с одной рабочей директорией. Повторный запуск с +той же директорией продолжит работу по уже созданным промежуточным файлам. +Кэш Whisper можно безопасно разделять между заданиями. + +### Статус и логирование + +Приложение дописывает в `JOB_DIR/report.txt` по одной JSON-записи на строку. +Записи появляются при старте и завершении этапа, а во время долгой операции — +каждые 30 секунд. Интервал меняется через `--report-interval`. + +```json +{"timestamp":"2026-09-24T12:00:00Z","run_id":"...","status":"running","step":"voice_recognition","step_index":2,"steps_total":10} +{"timestamp":"2026-09-24T12:00:30Z","run_id":"...","status":"heartbeat","elapsed_seconds":30,"step":"voice_recognition","step_index":2,"steps_total":10} +{"timestamp":"2026-09-24T12:10:00Z","run_id":"...","status":"succeeded","elapsed_seconds":600,"outputs":["output.md","output.pdf"]} +``` + +`status` принимает значения `started`, `running`, `step_completed`, +`heartbeat`, `succeeded`, `failed` или `cancelled`. Файл не перезаписывается: +повторный запуск добавляет записи с новым `run_id`. Смотреть его вручную можно +через `tail -f "$JOB_DIR/report.txt"`. + +Matrix-боту рекомендуется запускать `docker run --rm --name ` через +`subprocess.Popen`/`asyncio.create_subprocess_exec`, параллельно читать новые +строки `report.txt` и разбирать их через `json.loads`. Код возврата контейнера +остаётся окончательным признаком успеха: `0` — успех, ненулевой — ошибка или +отмена. Подробные диагностические логи идут в stdout/stderr контейнера; для +отмены задания можно выполнить `docker stop `. ## Архитектура Предполагается, что приложение будет запускаться сторонним приложением всякий раз, когда требуется произвести конвертацию медиафайла в конспект (например, -Telegram ботом, которому отправили видео). +Matrix-ботом, которому отправили видео). -**Сервис не сохраняет никаких данных между перезапусками**. Всё, что -сохраняется - это промежуточные результаты. Например, если сервис сгенерировал -файл `asr_events.json`, а после этого его принудительно завершили, то при -следующем запуске он будет использовать этот файл, чтобы не повторять дорогие -операции. **Поэтому при автоматизации рекомендуется на каждый новый запрос -создавать новую промежуточную директорию, а старые директории удалять, когда они -становятся не нужны.** +**Сервис не хранит состояние вне рабочей директории.** Если сервис сгенерировал +промежуточный файл, а после этого его принудительно завершили, то при следующем +запуске он использует этот файл, чтобы не повторять дорогую операцию. **Поэтому +для каждого нового запроса нужна новая рабочая директория; старые директории +можно удалять, когда результат больше не нужен.** Алгоритм работы сервиса следующий: 1. **Проверить, является входной файл звуком или видео со звуком** diff --git a/main.py b/main.py index 9e5b639..a92f27e 100644 --- a/main.py +++ b/main.py @@ -1,12 +1,10 @@ """Application entry point""" -import traceback import argparse import logging import json -import time import os -from dataclasses import asdict +import signal from enum import Enum from typing import Callable @@ -20,6 +18,7 @@ from structure_builder import Structure, StructureBuilder from structure_refiner import StructureRefiner from markdown_builder import MarkdownBuilder from pdf_builder import PdfBuilder +from reporting import StatusReporter from windowizer import Windowizer from agent import Agent @@ -27,6 +26,15 @@ from utils import Timeline, ffmpeg_split_video, ffmpeg_to_mp3 ARGS: argparse.Namespace + +class TerminationRequested(Exception): + """Raised when the container receives SIGTERM.""" + + def __init__(self, signal_number: int) -> None: + super().__init__(f"Received signal {signal_number}") + self.signal_number = signal_number + + class Step(Enum): MEDIA_SEPARATION = "media_separation" VOICE_RECOGNITION = "voice_recognition" @@ -158,8 +166,17 @@ def setup_arguments() -> argparse.Namespace: type=str, default=ai_api_key ) + parser.add_argument( + "--report-interval", + type=float, + default=30.0, + help="seconds between report.txt heartbeat records", + ) parser.add_argument("-v", action='store_true') - return parser.parse_args() + args = parser.parse_args() + if args.report_interval <= 0: + parser.error("--report-interval must be positive") + return args # # Workflow @@ -393,23 +410,28 @@ Step.CODE: ( ) """ -def main() -> None: - """Application entry point""" - global ARGS - check_cuda() - ARGS = setup_arguments() - logging.basicConfig(level=logging.DEBUG if ARGS.v else logging.INFO) + +def run_workflow(reporter: StatusReporter) -> list[str]: + """Execute the workflow and return paths to its final outputs.""" step_to_do = Step.MEDIA_SEPARATION intermediate_result: dict | None = None + step_numbers = { + step: index + for index, step in enumerate(WORKFLOW_DATA, start=1) + } + steps_total = len(WORKFLOW_DATA) + # execute steps while possible while step_to_do: + current_step = step_to_do + step_index = step_numbers[current_step] + reporter.step_started(current_step.value, step_index, steps_total) step_data = WORKFLOW_DATA[step_to_do] output_file_path = step_data[0] func = step_data[1] if not func: - logging.error(f"Step {step_to_do} has no function, stopping") - break - logging.info(f"Executing step {step_to_do}") + raise RuntimeError(f"Step {step_to_do} has no function") + logging.info("Executing step %s", step_to_do.value) step_to_do, ret = func(step_to_do, intermediate_result) if ret is not None: with open(output_file_path, "w", encoding="utf-8") as f: @@ -418,14 +440,53 @@ def main() -> None: else: intermediate_result = ret json.dump(ret, f, ensure_ascii=False, indent=4) - logging.info(f"Intermediate results are saved {output_file_path}") - # final report - logging.info(f"Workflow was interrupted at step {step_to_do}") + logging.info("Intermediate results are saved to %s", output_file_path) + if step_to_do is None and current_step is not Step.PDF_BUILDER: + raise RuntimeError(f"Workflow stopped at step {current_step.value}") + reporter.step_completed(current_step.value, step_index, steps_total) + + return [path for path in ("output.md", "output.pdf") if os.path.isfile(path)] + + +def main() -> None: + """Application entry point.""" + global ARGS + ARGS = setup_arguments() + logging.basicConfig(level=logging.DEBUG if ARGS.v else logging.INFO) + reporter = StatusReporter(heartbeat_interval=ARGS.report_interval) + reporter.start() + + previous_sigterm_handler = signal.getsignal(signal.SIGTERM) + + def handle_sigterm(signal_number: int, _frame: object) -> None: + raise TerminationRequested(signal_number) + + signal.signal(signal.SIGTERM, handle_sigterm) + try: + check_cuda() + outputs = run_workflow(reporter) + except TerminationRequested as error: + logging.warning("Workflow was cancelled by SIGTERM") + reporter.finish("cancelled", signal=error.signal_number) + raise SystemExit(128 + error.signal_number) from None + except KeyboardInterrupt: + logging.warning("Workflow was cancelled by the user") + reporter.finish("cancelled") + raise + except BaseException as error: + logging.error("Workflow failed: %s: %s", type(error).__name__, error) + reporter.finish( + "failed", + error=type(error).__name__, + message=str(error), + ) + raise + else: + reporter.finish("succeeded", outputs=outputs) + logging.info("Workflow completed successfully") + finally: + signal.signal(signal.SIGTERM, previous_sigterm_handler) + if __name__ == "__main__": - try: - main() - except SystemExit: - raise - except: - traceback.print_exc() + main() diff --git a/reporting.py b/reporting.py new file mode 100644 index 0000000..fe1a466 --- /dev/null +++ b/reporting.py @@ -0,0 +1,128 @@ +"""Append-only workflow status reporting for external process supervisors.""" + +from __future__ import annotations + +import json +import threading +import time +import uuid +from datetime import datetime, timezone +from pathlib import Path +from typing import Any, Literal, TextIO + + +TerminalStatus = Literal["succeeded", "failed", "cancelled"] + + +class StatusReporter: + """Write JSON Lines status events and periodic heartbeats.""" + + def __init__( + self, + path: str | Path = "report.txt", + *, + heartbeat_interval: float = 30.0, + ) -> None: + if heartbeat_interval <= 0: + raise ValueError("heartbeat_interval must be positive") + + self._path = Path(path) + self._heartbeat_interval = heartbeat_interval + self._run_id = uuid.uuid4().hex + self._started_at = time.monotonic() + self._step: str | None = None + self._step_index: int | None = None + self._steps_total: int | None = None + self._lock = threading.Lock() + self._stop_event = threading.Event() + self._thread: threading.Thread | None = None + self._file: TextIO | None = None + self._finished = False + + @staticmethod + def _timestamp() -> str: + return datetime.now(timezone.utc).isoformat(timespec="seconds").replace( + "+00:00", + "Z", + ) + + def _write(self, status: str, **fields: Any) -> None: + record = { + "timestamp": self._timestamp(), + "run_id": self._run_id, + "status": status, + **fields, + } + line = json.dumps(record, ensure_ascii=False, separators=(",", ":")) + with self._lock: + if self._file is None: + raise RuntimeError("StatusReporter is not running") + self._file.write(line + "\n") + self._file.flush() + + def start(self) -> None: + """Open the report and start emitting heartbeats.""" + if self._file is not None or self._finished: + raise RuntimeError("StatusReporter cannot be started again") + self._file = self._path.open("a", encoding="utf-8", buffering=1) + self._write("started") + self._thread = threading.Thread( + target=self._heartbeat_loop, + name="status-reporter", + daemon=True, + ) + self._thread.start() + + def _heartbeat_loop(self) -> None: + while not self._stop_event.wait(self._heartbeat_interval): + with self._lock: + step = self._step + step_index = self._step_index + steps_total = self._steps_total + fields: dict[str, Any] = { + "elapsed_seconds": round(time.monotonic() - self._started_at), + } + if step is not None: + fields.update({ + "step": step, + "step_index": step_index, + "steps_total": steps_total, + }) + self._write("heartbeat", **fields) + + def step_started(self, step: str, index: int, total: int) -> None: + with self._lock: + self._step = step + self._step_index = index + self._steps_total = total + self._write( + "running", + step=step, + step_index=index, + steps_total=total, + ) + + def step_completed(self, step: str, index: int, total: int) -> None: + self._write( + "step_completed", + step=step, + step_index=index, + steps_total=total, + ) + + def finish(self, status: TerminalStatus, **fields: Any) -> None: + """Write the terminal event and close the report.""" + if self._finished: + return + self._finished = True + self._stop_event.set() + if self._thread is not None: + self._thread.join() + self._write( + status, + elapsed_seconds=round(time.monotonic() - self._started_at), + **fields, + ) + assert self._file is not None + self._file.close() + self._file = None