251 lines
14 KiB
Python
251 lines
14 KiB
Python
"""Оркестрация AI-пайплайна пост-обработки сеанса.
|
||
|
||
Единственная задача-диспетчер `run_pipeline(session_id)`: смотрит текущее
|
||
`pipeline_status` сеанса и эффективную конфигурацию инстанса
|
||
(`services.instance_settings.load_effective_config` — настройки из БД
|
||
поверх дефолтов `config/plugins.yaml`), продолжает работу с
|
||
последнего успешного шага. Ожидание завершения
|
||
записи треков (`egress_ended` приходит позже `room_finished`) — через
|
||
celery-retry с backoff внутри самой задачи. Сама суммаризация здесь
|
||
не выполняется: `run_pipeline` доводит сеанс до `pipeline_status='summarizing'`
|
||
и передаёт эстафету задаче `workers.tasks.summarize.summarize_session` —
|
||
по имени через `app.send_task`, без импорта модуля суммаризации (транскрайбер-
|
||
процесс не должен тянуть его код).
|
||
|
||
Надёжность постановки `summarize_session`: временная
|
||
недоступность брокера Redis в момент `app.send_task` не должна ронять
|
||
`run_pipeline` — фразы к этому моменту уже закоммичены, и `failed` из-за
|
||
одного лишь сбоя постановки задачи стёр бы уже проделанную работу
|
||
транскрибации. Поэтому: (1) общий хелпер `workers.tasks.dispatch.
|
||
send_task_with_retry` делает несколько попыток с коротким backoff
|
||
(переиспользуется и `summarize_session` для постановки `notify_session`);
|
||
(2) если все попытки исчерпаны — ошибка логируется, но `pipeline_status`
|
||
остаётся `summarizing` (не откатывается, не переводится в `failed`), а
|
||
восстановление берёт на себя периодическая задача `workers.tasks.
|
||
maintenance.recover_stuck_summaries` (уровень 2 защиты).
|
||
"""
|
||
|
||
import logging
|
||
import uuid
|
||
from datetime import timedelta
|
||
from pathlib import Path
|
||
from typing import Protocol
|
||
|
||
from celery.exceptions import MaxRetriesExceededError
|
||
from sqlalchemy import delete
|
||
|
||
from core.plugins.config import InstanceConfig
|
||
from core.plugins.factory import create_transcriber
|
||
from core.plugins.transcriber import Segment
|
||
from models.audio_track import SessionAudioTrack
|
||
from models.phrase import Phrase
|
||
from models.session import ConferenceSession
|
||
from repositories.conferences import AudioTrackRepository
|
||
from services.instance_settings import load_effective_config
|
||
from workers.celery_app import app as app
|
||
from workers.db import open_session, run_async
|
||
from workers.tasks.dispatch import send_task_with_retry
|
||
from workers.transcription.phrases import build_phrases
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
RETRY_COUNTDOWN_S = 30
|
||
"""Пауза (сек) перед повторной попыткой, пока треки ещё дописываются egress'ом."""
|
||
|
||
_TRANSCRIBABLE_PIPELINE_STATUSES = frozenset({"recording", "transcribing"})
|
||
"""Статусы сеанса, с которых допустим (повторный) запуск шага транскрибации —
|
||
guard идемпотентности: если пайплайн уже ушёл дальше (`summarizing` и
|
||
позже) или зафиксирован как `failed`, повторный вызов задачи — no-op."""
|
||
|
||
|
||
class RetryableTask(Protocol):
|
||
"""Минимальный протокол объекта задачи с методом `retry` (для тестируемости).
|
||
|
||
Реальный `celery.Task` (bound self) удовлетворяет протоколу структурно —
|
||
отдельный импорт `celery.Task` как типа не нужен.
|
||
"""
|
||
|
||
def retry(self, countdown: int | None = None) -> None: ...
|
||
|
||
|
||
@app.task(name="workers.tasks.pipeline.run_pipeline", bind=True, max_retries=20, acks_late=True)
|
||
def run_pipeline(self: RetryableTask, session_id: str) -> None:
|
||
"""Точка входа Celery — синхронная обёртка над асинхронной логикой диспетчера."""
|
||
run_async(lambda: run_pipeline_async(self, uuid.UUID(session_id)))
|
||
|
||
|
||
async def run_pipeline_async(
|
||
task: RetryableTask,
|
||
session_id: uuid.UUID,
|
||
*,
|
||
plugins_config: InstanceConfig | None = None,
|
||
) -> None:
|
||
"""Диспетчер: продолжить AI-пайплайн сеанса `session_id` с последнего успешного шага.
|
||
|
||
Шаги:
|
||
1. Сеанс не найден либо ещё не завершён (`t_end IS NULL`) — выход.
|
||
2. `transcriber.enabled=false` — выход, `pipeline_status` не меняется.
|
||
3. Пайплайн уже прошёл шаг транскрибации (`pipeline_status` не в
|
||
`{recording, transcribing}`) — выход (идемпотентность повторного вызова).
|
||
4. Есть треки со статусом `recording` (egress ещё пишет) — `task.retry`;
|
||
исчерпание попыток — зависшие треки помечаются `failed`, работа
|
||
продолжается с остальными треками.
|
||
5. `pipeline_status='transcribing'`, commit.
|
||
6. Транскрибация треков со статусом `recorded` и `segments IS NULL`
|
||
(уже транскрибированные при прошлом прогоне — пропускаются); commit
|
||
ПОСЛЕ КАЖДОГО трека — точка возобновления при падении процесса
|
||
посередине (acks_late).
|
||
7. Ни одного трека не транскрибировано (все failed либо треков нет) —
|
||
`pipeline_status='failed'`, выход.
|
||
8. Реконструкция фраз (`build_phrases`, ТЗ §1.3) → DELETE+INSERT `phrases`
|
||
одной транзакцией → `pipeline_status='summarizing'`, commit → постановка
|
||
задачи `workers.tasks.summarize.summarize_session` в очередь по
|
||
умолчанию (её слушает базовый `worker`, а не `worker-transcriber`) —
|
||
через общий `workers.tasks.dispatch.send_task_with_retry` (retry на сбой
|
||
брокера, см. его докстринг и `workers.tasks.maintenance.recover_stuck_summaries`).
|
||
"""
|
||
async with open_session() as session:
|
||
session_record = await session.get(ConferenceSession, session_id)
|
||
if session_record is None or session_record.t_end is None:
|
||
logger.warning(
|
||
"run_pipeline: сеанс %s не найден либо ещё не завершён (t_end IS NULL)",
|
||
session_id,
|
||
)
|
||
return
|
||
|
||
cfg = plugins_config or await load_effective_config(session)
|
||
if not cfg.transcriber.enabled:
|
||
logger.info(
|
||
"run_pipeline: transcriber отключён (enabled=false) — сеанс %s пропущен",
|
||
session_id,
|
||
)
|
||
return
|
||
|
||
if session_record.pipeline_status not in _TRANSCRIBABLE_PIPELINE_STATUSES:
|
||
logger.info(
|
||
"run_pipeline: сеанс %s уже прошёл шаг транскрибации (pipeline_status=%s) — no-op",
|
||
session_id,
|
||
session_record.pipeline_status,
|
||
)
|
||
return
|
||
|
||
track_repo = AudioTrackRepository(session)
|
||
# Сортировка по `started_at` — детерминированный порядок обработки
|
||
# (точка возобновления по-трекового коммита должна быть
|
||
# предсказуемой между прогонами, см. тест идемпотентности №12).
|
||
tracks = sorted(await track_repo.list_by_session(session_id), key=lambda t: t.started_at)
|
||
|
||
still_recording = [track for track in tracks if track.status == "recording"]
|
||
if still_recording:
|
||
try:
|
||
task.retry(countdown=RETRY_COUNTDOWN_S)
|
||
except MaxRetriesExceededError:
|
||
for track in still_recording:
|
||
track.status = "failed"
|
||
await session.commit()
|
||
logger.warning(
|
||
"run_pipeline: исчерпаны попытки ожидания egress для %d треков сеанса %s "
|
||
"— помечены failed",
|
||
len(still_recording),
|
||
session_id,
|
||
)
|
||
else:
|
||
# Реальный `Task.retry()` сам бросает исключение Retry (не
|
||
# возвращает управление) — сюда попадаем только с
|
||
# моком/заглушкой `task.retry` в тестах.
|
||
return
|
||
|
||
session_record.pipeline_status = "transcribing"
|
||
await session.commit()
|
||
|
||
transcriber = create_transcriber(cfg.transcriber)
|
||
for track in tracks:
|
||
if track.status != "recorded" or track.segments is not None:
|
||
continue # уже транскрибирован на прошлом прогоне — идемпотентность
|
||
|
||
if not track.file_path or not Path(track.file_path).exists():
|
||
track.status = "failed"
|
||
logger.warning(
|
||
"run_pipeline: файл трека %s не найден (%s) — трек помечен failed",
|
||
track.id,
|
||
track.file_path,
|
||
)
|
||
await session.commit()
|
||
continue
|
||
|
||
segments = transcriber.transcribe(track.file_path, cfg.transcriber.language)
|
||
track.segments = [
|
||
{"start": segment.start, "end": segment.end, "text": segment.text}
|
||
for segment in segments
|
||
]
|
||
track.status = "transcribed"
|
||
await session.commit()
|
||
|
||
if not any(track.status == "transcribed" for track in tracks):
|
||
session_record.pipeline_status = "failed"
|
||
await session.commit()
|
||
logger.warning(
|
||
"run_pipeline: все треки сеанса %s провалены либо треков нет — pipeline failed",
|
||
session_id,
|
||
)
|
||
return
|
||
|
||
segments_by_participant, track_offsets = _collect_transcribed(tracks, session_record)
|
||
phrases = build_phrases(segments_by_participant, track_offsets)
|
||
|
||
await session.execute(delete(Phrase).where(Phrase.session_id == session_id))
|
||
for phrase in phrases:
|
||
session.add(
|
||
Phrase(
|
||
participant_id=phrase.participant_id,
|
||
session_id=session_id,
|
||
data=phrase.text,
|
||
t_start=session_record.t_start + timedelta(seconds=phrase.start),
|
||
t_end=session_record.t_start + timedelta(seconds=phrase.end),
|
||
)
|
||
)
|
||
session_record.pipeline_status = "summarizing"
|
||
await session.commit()
|
||
|
||
# По имени задачи, без импорта `workers.tasks.summarize` — модуль
|
||
# суммаризации не должен становиться зависимостью
|
||
# транскрайбер-процесса. Отправляем безусловно: если
|
||
# `summarizer.enabled=false`, задача сама завершится по своему guard'у
|
||
# (фразы уже сохранены, `pipeline_status='summarizing'` без summary —
|
||
# задокументированное завершение пайплайна после phrases).
|
||
sent = await send_task_with_retry(
|
||
"workers.tasks.summarize.summarize_session", args=[str(session_id)]
|
||
)
|
||
|
||
logger.info(
|
||
"run_pipeline: сеанс %s — реконструировано %d фраз, статус=summarizing, "
|
||
"постановка задачи суммаризации %s",
|
||
session_id,
|
||
len(phrases),
|
||
"выполнена" if sent else "не удалась (см. предыдущий error) — ждём recovery",
|
||
)
|
||
|
||
|
||
def _collect_transcribed(
|
||
tracks: list[SessionAudioTrack], session_record: ConferenceSession
|
||
) -> tuple[dict[uuid.UUID, list[Segment]], dict[uuid.UUID, float]]:
|
||
"""Собрать сегменты и смещения транскрибированных треков для `build_phrases`.
|
||
|
||
Смещение трека — `(track.started_at - session.t_start).total_seconds()`
|
||
(«Ключевые архитектурные решения», п.5). Ключ обоих словарей
|
||
— `participant_id`: одному участнику соответствует один аудиотрек сеанса
|
||
(одно окно присутствия — один микрофон, ADR-002).
|
||
"""
|
||
segments_by_participant: dict[uuid.UUID, list[Segment]] = {}
|
||
track_offsets: dict[uuid.UUID, float] = {}
|
||
for track in tracks:
|
||
if track.status != "transcribed" or not track.segments:
|
||
continue
|
||
offset = (track.started_at - session_record.t_start).total_seconds()
|
||
track_offsets[track.participant_id] = offset
|
||
segments_by_participant.setdefault(track.participant_id, []).extend(
|
||
Segment(start=raw["start"], end=raw["end"], text=raw["text"])
|
||
for raw in track.segments
|
||
)
|
||
return segments_by_participant, track_offsets
|