Первоначальная версия VidConf
This commit is contained in:
250
workers/tasks/pipeline.py
Normal file
250
workers/tasks/pipeline.py
Normal file
@@ -0,0 +1,250 @@
|
||||
"""Оркестрация 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
|
||||
Reference in New Issue
Block a user