"""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 ]