73 lines
4.3 KiB
Python
73 lines
4.3 KiB
Python
"""Общий хелпер безопасной постановки Celery-задачи по имени.
|
||
|
||
Вынесен из `_send_summarize_task` (`workers.tasks.pipeline`) как общая
|
||
защита от потери шага при недоступности Redis:
|
||
временная недоступность брокера в момент `app.send_task` не должна ронять
|
||
вызывающую задачу — предыдущий шаг (фразы/саммари) уже закоммичен, и
|
||
`failed` из-за одного лишь сбоя постановки СЛЕДУЮЩЕЙ задачи стёр бы уже
|
||
проделанную работу. Используется дважды: `workers.tasks.pipeline.run_pipeline`
|
||
(постановка `summarize_session`) и `workers.tasks.summarize.
|
||
summarize_session` (постановка `notify_session`) — оба вызывающих
|
||
места сами решают, что делать при `False` (не роняются, статус пайплайна не
|
||
откатывается и не проваливается; восстановление — соответствующая beat-задача
|
||
`workers.tasks.maintenance`, уровень 2 защиты).
|
||
"""
|
||
|
||
import asyncio
|
||
import logging
|
||
|
||
from kombu.exceptions import ConnectionError as KombuConnectionError
|
||
from kombu.exceptions import OperationalError as KombuOperationalError
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
SEND_TASK_MAX_ATTEMPTS = 3
|
||
"""Число попыток поставить задачу в очередь при недоступности брокера."""
|
||
|
||
SEND_TASK_BACKOFF_S = 2.0
|
||
"""Базовая пауза (сек) между попытками; растёт линейно с номером попытки."""
|
||
|
||
|
||
async def send_task_with_retry(task_name: str, args: list[str]) -> bool:
|
||
"""Поставить задачу `task_name` в очередь с несколькими попытками при сбое брокера.
|
||
|
||
`args` — позиционные аргументы задачи (как ожидает `Celery.send_task`).
|
||
Возвращает `True`, если постановка удалась хотя бы с одной попытки;
|
||
`False` — все попытки исчерпаны. Вызывающая сторона логирует контекст
|
||
(какой сеанс/задача) и оставляет восстановление уровню 2 защиты — не
|
||
откатывает и не проваливает уже проделанную работу из-за одного лишь
|
||
сбоя постановки следующей задачи.
|
||
"""
|
||
# Отложенный импорт `app`: на верхнем уровне модуля это создало бы цикл
|
||
# (`workers.celery_app` импортирует `workers.tasks.pipeline`/`summarize`,
|
||
# а те — `send_task_with_retry` из этого модуля) — импорт откладывается
|
||
# до первого вызова, когда `workers.celery_app` уже полностью загружен.
|
||
from workers.celery_app import app
|
||
|
||
for attempt in range(1, SEND_TASK_MAX_ATTEMPTS + 1):
|
||
try:
|
||
app.send_task(task_name, args=args)
|
||
except (KombuOperationalError, KombuConnectionError) as exc:
|
||
if attempt < SEND_TASK_MAX_ATTEMPTS:
|
||
logger.warning(
|
||
"send_task_with_retry: сбой постановки %s (попытка %d/%d): %s "
|
||
"— повтор через %.1fс",
|
||
task_name,
|
||
attempt,
|
||
SEND_TASK_MAX_ATTEMPTS,
|
||
exc,
|
||
SEND_TASK_BACKOFF_S * attempt,
|
||
)
|
||
await asyncio.sleep(SEND_TASK_BACKOFF_S * attempt)
|
||
continue
|
||
logger.error(
|
||
"send_task_with_retry: не удалось поставить %s за %d попыток (%s)",
|
||
task_name,
|
||
SEND_TASK_MAX_ATTEMPTS,
|
||
exc,
|
||
)
|
||
return False
|
||
else:
|
||
return True
|
||
return False # недостижимо: цикл либо возвращает, либо продолжает до предела
|