"""Задачи воркера: ping + ingest (маршрутизация принятых событий).
Приём (api) кладёт только конверт в БД; здесь Routing Engine MVP:
первое подходящее enabled-правило по weight → создание Delivery по каждому
действию. Исполнение каналов (Telegram/Email/S2S) — задачи #29/#28;
internal_log исполняется сразу (это запись в эту же таблицу).
"""
import uuid
from datetime import UTC, datetime
from celery import shared_task
from jinja2 import Template
from sqlalchemy import select
from app.database import SessionLocal
from app.models import (
ChannelTarget,
Delivery,
Event,
RoutingRule,
RoutingRuleAction,
Source,
)
from app.worker.celery_app import celery_app
PRIORITY_ORDER = {"low": 0, "normal": 1, "high": 2, "critical": 3}
def _rule_matches(rule: RoutingRule, *, source_name: str, subject: str, action: str, priority: str) -> 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
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
@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.
Первое подошедшее правило (по weight) применяется целиком; прочие
условия/полнота Routing Engine уточняются в #30. Если правил нет —
событие просто завершается (конверт остаётся в 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: RoutingRule | None = None
for rule in rules:
if _rule_matches(
rule,
source_name=source.name if source else "",
subject=event.subject,
action=event.action,
priority=event.priority,
):
matched = rule
break
now = datetime.now(UTC)
if matched is not None:
actions = db.execute(
select(RoutingRuleAction).where(RoutingRuleAction.rule_id == matched.id)
).scalars().all()
for act in actions:
target_name: str | None = None
if act.target_id:
target = db.get(ChannelTarget, act.target_id)
target_name = target.name if target else None
delivery = Delivery(
event_id=event.id,
rule_id=matched.id,
channel=act.channel,
channel_target_id=act.target_id,
# Фактическая отправка канала — задача #29/#28; пока
# Telegram/Email/S2S лежат pending. internal_log
# исполняется сразу: delivery — и есть запись лога.
status="delivered" if act.channel == "internal_log" else "pending",
attempts=1,
rendered_message=(
_render(act.template or matched.template or "", event, source.name)
if act.channel == "internal_log"
else None
),
delivered_at=now if act.channel == "internal_log" else None,
last_error=None if act.channel == "internal_log" else "канал ещё не реализован (#29/#28)",
)
db.add(delivery)
event.status = "done"
db.commit()
rule_note = f"rule={matched.id}" if matched else "без правила"
return f"event {event_id}: {rule_note}"