Newer
Older
gn-synapse / app / worker / tasks.py
"""Задачи воркера: ping, ingest (Routing Engine), retry_due (ретраи),
expire_events (ретеншн, beat ежечасно).

Приём (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).
- истёкшие ttl_seconds события не маршрутизируются (status=expired) и
  удаляются beat-задачей synapse.expire_events вместе с событиями старше
  RETENTION_DAYS (см. «Ретеншн» в docs/05).
"""

import uuid
from datetime import UTC, datetime, timedelta

from celery import shared_task
from jinja2 import Template
from sqlalchemy import delete, select

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.settings_store import get_setting
from app.worker.senders import send_s2s, send_push

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

# Целевые каналы без цели: адресат — payload.user_id (docs/05, docs/06).
_USER_CHANNELS = ("user", "push")


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 пуст)"
        if target.deleted_at is not None:
            # ретрай по цели, заархивированной после маршрутизации
            return False, "цель в архиве — restore в админке вернёт её в строй"
        return send_s2s(target, event, source_name)
    if delivery.channel == "push":
        # Web-push: подписки получателя хранятся в push_subscriptions (docs/06).
        return send_push(db, delivery, 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 _user_recipient(event: Event) -> tuple[str | None, str | None]:
    """payload.user_id → (адресат, причина пропуска).

    Смысловое значение поля — sub gnexus-auth: непустая строка (число
    допускаем — приводится к строке). Отсутствие поля — не ошибка источника
    (None, None); тип-не-строка — ошибка источника, причина уйдёт в аудит
    skipped-доставки, иначе доставка ушла бы «в никуда» молча.
    """
    raw = (event.payload or {}).get("user_id")
    if raw is None:
        return None, None
    if isinstance(raw, str):
        rid = raw.strip()[:64]
        return (rid, None) if rid else (None, None)
    if isinstance(raw, int) and not isinstance(raw, bool):
        rid = str(raw)[:64]
        return (rid, None) if rid else (None, None)
    return None, f"payload.user_id имеет тип {type(raw).__name__} — не похоже на sub gnexus-auth"


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

        now = datetime.now(UTC)
        # ttl_seconds: событие бесполезно после срока (docs/05) — не маршрутизируем,
        # удалять будет beat-задача expire_events.
        if event.expires_at is not None and event.expires_at < now:
            event.status = "expired"
            db.commit()
            return f"skip: event {event_id} истёк ttl"

        source = db.get(Source, event.source_id)
        event.status = "processing"

        rules = db.execute(
            select(RoutingRule).where(
                RoutingRule.enabled.is_(True),
                RoutingRule.deleted_at.is_(None),
            ).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_setting(db, "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/push адресуют пользователя из payload.user_id
                # (конвенция, docs/05); без валидного sub доставить некому.
                recipient, user_reason = None, None
                if act.channel in _USER_CHANNELS:
                    recipient, user_reason = _user_recipient(event)

                # Архивная цель — «сюда больше не ходим»: правило живо (или
                # цель заархивировали позже правила), доставка уходит в
                # skipped с аудитом (аналог user без адресата).
                if act.channel not in _USER_CHANNELS and act.target_id is not None:
                    t_row = db.get(ChannelTarget, act.target_id)
                    if t_row is None or t_row.deleted_at is not None:
                        db.add(
                            Delivery(
                                event_id=event.id,
                                rule_id=rule.id,
                                channel=act.channel,
                                channel_target_id=act.target_id,
                                status="skipped",
                                attempts=0,
                                rendered_message=_render(
                                    act.template or rule.template or "", event, source.name
                                ),
                                last_error="цель в архиве — restore в админке вернёт её в строй",
                            )
                        )
                        n_skipped += 1
                        continue

                # шумодав: срезаем (с аудит-записью), 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 in _USER_CHANNELS:
                    # Лично-адресованные каналы: user — личный лог (исполняется
                    # сразу), push — отправка в подписки пользователя (через
                    # _attempt_send/ретраи). Без адресата — skipped с аудитом.
                    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=user_reason or "payload.user_id не задан — некому адресовать",
                            )
                        )
                        n_skipped += 1
                        continue

                    message = _render(
                        act.template or rule.template or "", event, source.name
                    )
                    if act.channel == "user":
                        db.add(
                            Delivery(
                                event_id=event.id,
                                rule_id=rule.id,
                                channel=act.channel,
                                recipient_user_id=recipient,
                                status="delivered",
                                attempts=1,
                                rendered_message=message,
                                delivered_at=now,
                            )
                        )
                        continue

                    # push: pending-доставка с готовым сообщением; отправка
                    # во все подписки получателя — в _attempt_send.
                    delivery = Delivery(
                        event_id=event.id,
                        rule_id=rule.id,
                        channel=act.channel,
                        recipient_user_id=recipient,
                        status="pending",
                        attempts=0,
                        rendered_message=message,
                    )
                    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
                    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}"


@celery_app.task(name="synapse.expire_events")
def expire_events() -> str:
    """Ретеншн (beat, ежечасно): удаление событий с истёкшим ttl_seconds
    и событий старше RETENTION_DAYS; доставки падают каскадом — в том
    числе записи личного лога канала user (/api/v1/me/events читает их).
    RETENTION_DAYS <= 0 отключает удаление по возрасту (ttl продолжит).
    """
    now = datetime.now(UTC)
    ids: list[uuid.UUID] = []
    with SessionLocal() as db:
        ids += db.execute(
            select(Event.id).where(Event.expires_at.is_not(None), Event.expires_at < now)
        ).scalars().all()
        if get_setting(db, "retention_days") > 0:
            cutoff = now - timedelta(days=get_setting(db, "retention_days"))
            ids += db.execute(
                select(Event.id).where(Event.created_at < cutoff).limit(5000)
            ).scalars().all()
        ids = list(dict.fromkeys(ids))  # уникальные, порядок неважен
        if not ids:
            return "expire_events: удалять нечего"
        db.execute(delete(Event).where(Event.id.in_(ids)))  # deliveries — каскадом
        db.commit()
    return f"expire_events: удалено {len(ids)} событий"