217 lines
11 KiB
Python
217 lines
11 KiB
Python
"""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
|
||
]
|