Newer
Older
gn-synapse / app / worker / tasks.py
"""Задачи воркера: 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}"