Files
vidconf/workers/tasks/summarize.py

217 lines
11 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Celery-задача суммаризации сеанса.
Завершает AI-пайплайн после реконструкции фраз
(`workers.tasks.pipeline.run_pipeline`): собирает транскрипт сеанса из `phrases` и имён участников, вызывает
активный плагин `Summarizer` (`config/plugins.yaml`) и записывает результат в
`conference_sessions.summary_data`. Статус `pipeline_status` ОСТАЁТСЯ
`summarizing` после записи саммари; готовность к уведомлению определяется
парой `pipeline_status='summarizing'` И `summary_data IS NOT NULL`.
Значение `notified` ставит `workers.tasks.notify.
notify_session` — постановка по имени задачи через общий
`workers.tasks.dispatch.send_task_with_retry` сразу после коммита `summary_data`
(тот же паттерн защиты от потери шага при недоступности брокера, что у
`run_pipeline` при постановке этой самой задачи, см. `workers.tasks.pipeline`
и `workers.tasks.dispatch`). Если задача завершается no-op БЕЗ записи
`summary_data` (summarizer отключён, нет фраз, уже заполнено) — `notify_session`
не ставится.
"""
import logging
import uuid
from typing import Any, Protocol
from celery.exceptions import MaxRetriesExceededError
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from core.plugins.config import InstanceConfig
from core.plugins.factory import create_summarizer
from core.summarization.llm_client import LlmUnavailableError
from models.guest import GuestAccess
from models.participant import ConferenceParticipant
from models.phrase import Phrase
from models.session import ConferenceSession
from models.user import User
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.summarizer.transcript import TranscriptLine, build_transcript
from workers.tasks.dispatch import send_task_with_retry
logger = logging.getLogger(__name__)
RETRY_COUNTDOWN_BASE_S = 60
"""Базовая пауза (сек) перед повтором при обрыве LLM; фактический countdown —
`RETRY_COUNTDOWN_BASE_S * (attempt + 1)` (нарастающий backoff)."""
_DEFAULT_SPEAKER_NAME = "Участник"
"""Резервное имя говорящего на случай отсутствия и `users.name_user`, и
`guest_access.display_name` (не должно происходить при действующем
CHECK-constraint `conference_participants`, но защищает транскрипт от падения)."""
class RetryableTask(Protocol):
"""Минимальный протокол bound-задачи Celery, достаточный для тестируемости.
В отличие от одноимённого протокола в `workers.tasks.pipeline`, здесь
дополнительно нужен `request.retries` — номер уже сделанной попытки,
чтобы считать нарастающий countdown ретрая. Реальный `celery.Task`
(bound self) удовлетворяет протоколу структурно.
"""
request: Any
def retry(self, countdown: int | None = None) -> None: ...
@app.task(
name="workers.tasks.summarize.summarize_session",
bind=True,
max_retries=5,
acks_late=True,
)
def summarize_session(self: RetryableTask, session_id: str) -> None:
"""Точка входа Celery — синхронная обёртка над асинхронной логикой суммаризации."""
run_async(lambda: summarize_session_async(self, uuid.UUID(session_id)))
async def summarize_session_async(
task: RetryableTask,
session_id: uuid.UUID,
*,
plugins_config: InstanceConfig | None = None,
) -> None:
"""Суммаризировать сеанс `session_id`: собрать транскрипт → саммари → `summary_data`.
Guard'ы (строго по порядку):
1. Сеанс не найден либо ещё не завершён (`t_end IS NULL`) — выход.
2. `pipeline_status != 'summarizing'` — no-op (идемпотентность повторной
доставки задачи либо запуска раньше срока).
3. `summary_data IS NOT NULL` — no-op (сеанс уже готов к уведомлению).
4. `summarizer.enabled=false` — выход без изменений (пайплайн
задокументированно завершается после реконструкции фраз).
5. Фраз нет — выход, `summary_data` остаётся `NULL`.
Обрыв LLM (`LlmUnavailableError`) — `task.retry` с нарастающим countdown;
исчерпание попыток — `pipeline_status='failed'`. Один шаг — один commit.
"""
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(
"summarize_session: сеанс %s не найден либо ещё не завершён (t_end IS NULL)",
session_id,
)
return
if session_record.pipeline_status != "summarizing":
logger.info(
"summarize_session: сеанс %s не на шаге суммаризации "
"(pipeline_status=%s) — no-op",
session_id,
session_record.pipeline_status,
)
return
if session_record.summary_data is not None:
logger.info(
"summarize_session: сеанс %s уже готов к уведомлению "
"(summary_data заполнен) — no-op",
session_id,
)
return
cfg = plugins_config or await load_effective_config(session)
if not cfg.summarizer.enabled:
logger.info(
"summarize_session: summarizer отключён (enabled=false) — сеанс %s пропущен",
session_id,
)
return
lines = await _fetch_transcript_lines(session, session_record)
if not lines:
logger.warning(
"summarize_session: у сеанса %s нет фраз — summary_data остаётся NULL",
session_id,
)
return
transcript = build_transcript(lines)
summarizer = create_summarizer(cfg.summarizer)
try:
summary = summarizer.summarize(transcript)
except LlmUnavailableError as exc:
attempt = getattr(task.request, "retries", 0)
countdown = RETRY_COUNTDOWN_BASE_S * (attempt + 1)
try:
task.retry(countdown=countdown)
except MaxRetriesExceededError:
session_record.pipeline_status = "failed"
await session.commit()
logger.warning(
"summarize_session: исчерпаны попытки суммаризации сеанса %s "
"(LLM недоступен: %s) — pipeline failed",
session_id,
exc,
)
# Реальный `Task.retry()` сам бросает исключение Retry (не
# возвращает управление) — до сюда доходим только с
# моком/заглушкой `task.retry` в тестах.
return
session_record.summary_data = summary
await session.commit()
logger.info(
"summarize_session: сеанс %s — саммари сохранено (%d симв.)",
session_id,
len(summary),
)
# По имени задачи, без импорта `workers.tasks.notify` — разграничение
# процессов такое же, как у `run_pipeline` → `summarize_session`.
# Постановка защищена retry на сбой брокера
# (`send_task_with_retry`); неудача не роняет `summarize_session` и не
# трогает уже сохранённый `summary_data` — восстановление берёт на
# себя `workers.tasks.maintenance.recover_stuck_notifications`
# (уровень 2 защиты).
sent = await send_task_with_retry(
"workers.tasks.notify.notify_session", args=[str(session_id)]
)
logger.info(
"summarize_session: постановка задачи уведомления сеанса %s %s",
session_id,
"выполнена" if sent else "не удалась (см. предыдущий error) — ждём recovery",
)
async def _fetch_transcript_lines(
session: AsyncSession, session_record: ConferenceSession
) -> list[TranscriptLine]:
"""Собрать фразы сеанса с именами говорящих для построения транскрипта.
Имя говорящего — `users.name_user` для зарегистрированного участника либо
`guest_access.display_name` для гостя: ровно одно из `user_id`/`guest_id`
заполнено у `conference_participants` (CHECK-constraint, ADR-001 п.6).
`offset_s` — смещение начала фразы от `t_start` сеанса, как ожидает
`build_transcript`.
"""
stmt = (
select(Phrase.t_start, Phrase.data, User.name_user, GuestAccess.display_name)
.join(ConferenceParticipant, Phrase.participant_id == ConferenceParticipant.id)
.outerjoin(User, ConferenceParticipant.user_id == User.id)
.outerjoin(GuestAccess, ConferenceParticipant.guest_id == GuestAccess.id)
.where(Phrase.session_id == session_record.id)
.order_by(Phrase.t_start)
)
rows = (await session.execute(stmt)).all()
return [
TranscriptLine(
speaker=name_user or display_name or _DEFAULT_SPEAKER_NAME,
offset_s=(t_start - session_record.t_start).total_seconds(),
text=text,
)
for t_start, text, name_user, display_name in rows
]