"""Задачи воркера: ping, ingest (Routing Engine), retry_due (ретраи доставок).
Приём (api) кладёт только конверт в БД; здесь живёт Routing Engine:
- подбор правил по условиям (source/subjects/actions/priority_min/payload);
- шумодав правила (throttle_seconds, critical проходит всегда);
- отправка: internal_log сразу (delivery и есть запись лога), user — личный
лог пользователя (payload.user_id, см. docs/05), s2s — подписано
(senders.py, docs/05); Telegram/Email — заглушки под #29; провалы
— ретрай-цикл через beat (synapse.retry_due).
"""
import uuid
from datetime import UTC, datetime, timedelta
from celery import shared_task
from jinja2 import Template
from sqlalchemy import select
from app.config import get_settings
from app.database import SessionLocal
from app.models import (
ChannelTarget,
Delivery,
Event,
RoutingRule,
RoutingRuleAction,
Source,
)
from app.worker.celery_app import celery_app
from app.worker.senders import send_s2s
PRIORITY_ORDER = {"low": 0, "normal": 1, "high": 2, "critical": 3}
# Ретраи доставки: 1+4 попытки, пауза растёт (минуты → часы), после — failed.
RETRY_DELAYS = (30, 120, 600, 1800)
# Каналы, чью отправку воркер ещё не умеет (#29).
_CHANNELS_NOT_IMPL = ("telegram", "email")
def _attempt_send(
db, delivery: Delivery, event: Event, source_name: str
) -> tuple[bool, str | None]:
"""Одна попытка отправки. Возвращает (ok, ошибка).
Устанавливает результат прямо в delivery (статус/attempts/next_retry_at);
ретраи — через synapse.retry_due (beat). Заглушки под #29: telegram/email.
"""
if delivery.channel == "s2s":
target = (
db.get(ChannelTarget, delivery.channel_target_id)
if delivery.channel_target_id
else None
)
if target is None:
return False, "s2s: у доставки нет цели (target_id пуст)"
return send_s2s(target, event, source_name)
if delivery.channel in _CHANNELS_NOT_IMPL:
return False, "канал ещё не реализован (#29)"
return False, f"неизвестный канал {delivery.channel}"
def _payload_matches(actual: dict, wanted: dict) -> bool:
"""Условия payload: точное совпадение по top-level ключам; список = any-of."""
for key, expected in wanted.items():
allowed = expected if isinstance(expected, list) else [expected]
if actual.get(key) not in allowed:
return False
return True
def _rule_matches(
rule: RoutingRule,
*,
source_name: str,
subject: str,
action: str,
priority: str,
payload: dict,
) -> bool:
cond = rule.conditions or {}
if cond.get("source") is not None and cond["source"] != source_name:
return False
subjects = cond.get("subjects") or []
if subjects and subject not in subjects:
return False
actions = cond.get("actions") or []
if actions and action not in actions:
return False
min_priority = cond.get("priority_min")
if min_priority and PRIORITY_ORDER.get(priority, 0) < PRIORITY_ORDER.get(min_priority, 1):
return False
payload_cond = cond.get("payload")
if payload_cond and not _payload_matches(payload, payload_cond):
return False
return True
def _render(template_text: str, event: Event, source_name: str) -> str:
"""Jinja2-рендер (`{{ payload.x }}`); сломанный шаблон не крушит доставку."""
try:
return Template(template_text).render(
payload=event.payload,
id=str(event.id),
source=source_name,
subject=event.subject,
action=event.action,
priority=event.priority,
)
except Exception: # noqa: BLE001 — шаблон не должен ронять событие
return template_text
def _throttled(
db, rule: RoutingRule, channel: str, target_id: int | None, now: datetime,
recipient: str | None = None,
) -> bool:
"""Шумодав: была ли создана доставка в ту же цель в окне rule.throttle_seconds.
Считаются pending/delivered (в полёте тоже держат канал занятым);
skipped окно не продлевает — иначе непрерывный шторм никогда не проходит.
critical не доходит до этого вызова (обход на месте вызова).
Канал user — цель это recipient_user_id.
"""
if not rule.throttle_seconds or (target_id is None and recipient is None):
return False
recent = (
db.execute(
select(Delivery.id).where(
Delivery.channel == channel,
# Канал user адресует пользователя, остальные — ChannelTarget.
Delivery.recipient_user_id == recipient if recipient is not None else
Delivery.channel_target_id == target_id,
Delivery.status.in_(("pending", "delivered")),
Delivery.created_at > now - timedelta(seconds=rule.throttle_seconds),
)
.limit(1)
)
).scalar_one_or_none()
return recent is not None
@celery_app.task(name="synapse.ping")
def ping() -> str:
"""Проверка связки api -> очередь -> воркер."""
return "pong"
@celery_app.task(name="synapse.ingest")
def ingest(event_id: str) -> str:
"""Маршрутизация события: подобрать правила, создать deliveries.
Режим all (дефолт) — применяются все подошедшие правила; first —
только первое по weight. Если правил нет — событие завершается
(конверт остаётся в events как история).
"""
with SessionLocal() as db:
event = db.get(Event, uuid.UUID(event_id))
if event is None:
return f"skip: event {event_id} не найден"
if event.status != "queued":
return f"skip: event {event_id} в статусе {event.status}"
source = db.get(Source, event.source_id)
event.status = "processing"
rules = db.execute(
select(RoutingRule).where(RoutingRule.enabled.is_(True)).order_by(RoutingRule.weight)
).scalars().all()
matched = [
rule for rule in rules
if _rule_matches(
rule,
source_name=source.name if source else "",
subject=event.subject,
action=event.action,
priority=event.priority,
payload=dict(event.payload or {}),
)
]
if get_settings().routing_match_mode == "first":
matched = matched[:1]
now = datetime.now(UTC)
n_skipped = 0
for rule in matched:
actions = db.execute(
select(RoutingRuleAction).where(RoutingRuleAction.rule_id == rule.id)
).scalars().all()
for act in actions:
# Канал user адресует пользователя из payload.user_id
# (конвенция, docs/05); без него доставить некому.
recipient = None
if act.channel == "user":
rid = str((event.payload or {}).get("user_id") or "").strip()
recipient = rid[:64] or None
# шумодав: срезаем (с аудит-записью), critical проходит
if event.priority != "critical" and _throttled(
db, rule, act.channel, act.target_id, now, recipient=recipient
):
db.add(
Delivery(
event_id=event.id,
rule_id=rule.id,
channel=act.channel,
channel_target_id=act.target_id,
recipient_user_id=recipient,
status="skipped",
attempts=0,
rendered_message=(
_render(act.template or rule.template or "", event, source.name)
),
last_error=f"шумодав: доставка в эту цель уже была за "
f"{rule.throttle_seconds} с",
)
)
n_skipped += 1
continue
if act.channel == "internal_log":
# internal_log исполняется сразу: delivery — и есть запись лога.
db.add(
Delivery(
event_id=event.id,
rule_id=rule.id,
channel=act.channel,
channel_target_id=act.target_id,
status="delivered",
attempts=1,
rendered_message=_render(
act.template or rule.template or "", event, source.name
),
delivered_at=now,
)
)
continue
if act.channel == "user":
# Личный лог пользователя: также исполняется сразу.
if recipient is None:
# Аудит: правило направило, но адресовать некому —
# админ увидит причину в /admin/deliveries.
db.add(
Delivery(
event_id=event.id,
rule_id=rule.id,
channel=act.channel,
status="skipped",
attempts=0,
rendered_message=_render(
act.template or rule.template or "", event, source.name
),
last_error="payload.user_id не задан — некому адресовать",
)
)
n_skipped += 1
continue
db.add(
Delivery(
event_id=event.id,
rule_id=rule.id,
channel=act.channel,
recipient_user_id=recipient,
status="delivered",
attempts=1,
rendered_message=_render(
act.template or rule.template or "", event, source.name
),
delivered_at=now,
)
)
continue
# Фактическая отправка (internal_log уже ушёл выше): первая
# попытка сразу на маршрутизации; провал → ретрай-цикл.
delivery = Delivery(
event_id=event.id,
rule_id=rule.id,
channel=act.channel,
channel_target_id=act.target_id,
status="pending",
attempts=0,
)
db.add(delivery)
db.flush()
ok, error = _attempt_send(db, delivery, event, source.name)
if ok:
delivery.status = "delivered"
delivery.delivered_at = now
delivery.last_error = None
else:
delivery.last_error = error
delivery.next_retry_at = now + timedelta(seconds=RETRY_DELAYS[0])
delivery.attempts = 1
event.status = "done"
db.commit()
rules_note = ",".join(str(r.id) for r in matched) or "без правила"
skipped_note = f", skipped={n_skipped}" if n_skipped else ""
return f"event {event_id}: rules={rules_note}{skipped_note}"
@celery_app.task(name="synapse.retry_due")
def retry_due() -> str:
"""Пауза/попытка для pending-доставок с наступившим next_retry_at (beat).
При успехе — delivered; при провале attempts+1, следующая попытка по
RETRY_DELAYS, после последней — failed.
"""
now = datetime.now(UTC)
processed = 0
with SessionLocal() as db:
due = (
db.execute(
select(Delivery)
.where(
Delivery.status == "pending",
Delivery.next_retry_at.is_not(None),
Delivery.next_retry_at <= now,
)
.order_by(Delivery.next_retry_at)
.limit(50)
)
.scalars()
.all()
)
for delivery in due:
processed += 1
event = db.get(Event, delivery.event_id)
if event is None:
delivery.status = "failed"
delivery.last_error = "событие удалено вместе с доставкой"
continue
source = db.get(Source, event.source_id)
ok, error = _attempt_send(db, delivery, event, source.name if source else "")
if ok:
delivery.status = "delivered"
delivery.delivered_at = now
delivery.last_error = None
delivery.next_retry_at = None
else:
delivery.attempts += 1
delivery.last_error = error
if delivery.attempts >= len(RETRY_DELAYS) + 1:
delivery.status = "failed"
delivery.next_retry_at = None
else:
# RETRY_DELAYS[attempts-1]: после 1-й (стартовой) — пауза[0]
delivery.next_retry_at = now + timedelta(
seconds=RETRY_DELAYS[min(delivery.attempts - 1, len(RETRY_DELAYS) - 1)]
)
db.commit()
return f"retry_due: обработано {processed}"