Improved integration ability
This commit is contained in:
@@ -9,4 +9,5 @@ input.*
|
||||
audio.*
|
||||
video.*
|
||||
output.*
|
||||
report.txt
|
||||
*.json
|
||||
|
||||
1
.gitignore
vendored
1
.gitignore
vendored
@@ -6,6 +6,7 @@ images/
|
||||
debug/
|
||||
output.md
|
||||
output.pdf
|
||||
report.txt
|
||||
*.json
|
||||
|
||||
*.mkv
|
||||
|
||||
56
README.md
56
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 <job-id>` через
|
||||
`subprocess.Popen`/`asyncio.create_subprocess_exec`, параллельно читать новые
|
||||
строки `report.txt` и разбирать их через `json.loads`. Код возврата контейнера
|
||||
остаётся окончательным признаком успеха: `0` — успех, ненулевой — ошибка или
|
||||
отмена. Подробные диагностические логи идут в stdout/stderr контейнера; для
|
||||
отмены задания можно выполнить `docker stop <job-id>`.
|
||||
|
||||
## Архитектура
|
||||
|
||||
Предполагается, что приложение будет запускаться сторонним приложением всякий
|
||||
раз, когда требуется произвести конвертацию медиафайла в конспект (например,
|
||||
Telegram ботом, которому отправили видео).
|
||||
Matrix-ботом, которому отправили видео).
|
||||
|
||||
**Сервис не сохраняет никаких данных между перезапусками**. Всё, что
|
||||
сохраняется - это промежуточные результаты. Например, если сервис сгенерировал
|
||||
файл `asr_events.json`, а после этого его принудительно завершили, то при
|
||||
следующем запуске он будет использовать этот файл, чтобы не повторять дорогие
|
||||
операции. **Поэтому при автоматизации рекомендуется на каждый новый запрос
|
||||
создавать новую промежуточную директорию, а старые директории удалять, когда они
|
||||
становятся не нужны.**
|
||||
**Сервис не хранит состояние вне рабочей директории.** Если сервис сгенерировал
|
||||
промежуточный файл, а после этого его принудительно завершили, то при следующем
|
||||
запуске он использует этот файл, чтобы не повторять дорогую операцию. **Поэтому
|
||||
для каждого нового запроса нужна новая рабочая директория; старые директории
|
||||
можно удалять, когда результат больше не нужен.**
|
||||
|
||||
Алгоритм работы сервиса следующий:
|
||||
1. **Проверить, является входной файл звуком или видео со звуком**
|
||||
|
||||
105
main.py
105
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()
|
||||
|
||||
128
reporting.py
Normal file
128
reporting.py
Normal file
@@ -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
|
||||
Reference in New Issue
Block a user