diff --git a/.env.example b/.env.example index 6872de0..ad9a50c 100644 --- a/.env.example +++ b/.env.example @@ -28,6 +28,12 @@ # Публичный URL SPA для редиректа после callback. Пусто = тот же origin (прод). SPA_PUBLIC_URL=http://localhost:5174 +# --- Routing (docs/05) --- +# Окно дедупликации событий по dedup_key (сек). +DEDUP_WINDOW_SECONDS=86400 +# Режим правил: all (все подошедшие) | first (только первое по weight). +ROUTING_MATCH_MODE=all + # --- Misc --- # Порт API на хосте. SYNAPSE_PORT=8013 \ No newline at end of file diff --git a/README.md b/README.md index f45bfc5..017b8bb 100644 --- a/README.md +++ b/README.md @@ -75,6 +75,8 @@ docker compose exec api celery -A app.worker.celery_app call synapse.ping # возвращает task id ``` +Retraй-механика (#30): провал доставки → пауза 30 с → 2 м → 10 м → 30 м, 5 попыток → `failed`; скан due-доставок — beat воркера раз в минуту (`synapse.retry_due`). Шумодав — `--throttle N` у правила (не чаще одной доставки в цель за N с, critical проходит всегда); все подошедшие правила применяются (режим `all`, `ROUTING_MATCH_MODE=first` — только первое). Подробности — docs/05, раздел «Типы и правила». + ## Вход в админку SSO-поток целиком на сервере (gnexus-auth-client-py, OAuth2 Authorization Code + PKCE): diff --git a/alembic/versions/20261003_0003_routing_full.py b/alembic/versions/20261003_0003_routing_full.py new file mode 100644 index 0000000..0f9a0ed --- /dev/null +++ b/alembic/versions/20261003_0003_routing_full.py @@ -0,0 +1,35 @@ +""" Routing Engine полный: throttle_seconds у правил, статус skipped у доставок + +Revision ID: b7d2f4a19e60 +Revises: a4c1e9b02d77 +Create Date: 2026-10-03 +""" + +from alembic import op +import sqlalchemy as sa + + +revision = "b7d2f4a19e60" +down_revision = "a4c1e9b02d77" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.add_column("routing_rules", sa.Column("throttle_seconds", sa.Integer(), nullable=True)) + op.drop_constraint("ck_delivery_status", "deliveries", type_="check") + op.create_check_constraint( + "ck_delivery_status", + "deliveries", + "status IN ('pending','delivered','failed','skipped')", + ) + + +def downgrade() -> None: + op.drop_constraint("ck_delivery_status", "deliveries", type_="check") + op.create_check_constraint( + "ck_delivery_status", + "deliveries", + "status IN ('pending','delivered','failed')", + ) + op.drop_column("routing_rules", "throttle_seconds") \ No newline at end of file diff --git a/app/cli.py b/app/cli.py index 4d04ac7..59af25b 100644 --- a/app/cli.py +++ b/app/cli.py @@ -3,13 +3,14 @@ В проде всё это будет управляться SPA админ-панелью (постройка UI — #32+); CLI нужен для локальной разработки и smoke-тестов. - python -m app.cli create-source monitoring [--label ...] + python -m app.cli create-source monitoring [--label ...] [--description ...] python -m app.cli create-key monitoring [--name ...] # токен печатается ОДИН раз python -m app.cli add-type monitoring container down - python -m app.cli add-target telegram infra --config '{"chat_id": "@infra"}' + python -m app.cli add-target telegram infra --config '{"chat_id": "@infra"}' [--description ...] python -m app.cli add-rule "контейнеры вниз" --source monitoring \ --subjects container --actions down --priority high \ --template "Контейнер {{ payload.container }} упал" \ + --throttle 300 \ --action telegram:infra --action internal_log """ @@ -128,8 +129,10 @@ "subjects": args.subjects or [], "actions": args.actions or [], "priority_min": args.priority or None, + "payload": json.loads(args.payload) if args.payload else None, }, template=args.template, + throttle_seconds=args.throttle, ) db.add(rule) db.flush() @@ -177,8 +180,10 @@ p.add_argument("--subjects", nargs="*") p.add_argument("--actions", nargs="*") p.add_argument("--priority") + p.add_argument("--payload", help='JSON условий по payload, напр. \'{"container": "api"}\'') p.add_argument("--template") p.add_argument("--weight", type=int, default=0) + p.add_argument("--throttle", type=int, help="шумодав: не чаще одной доставки в эту цель за N с") p.add_argument("--action", action="append", default=[]) # channel[:target] p.set_defaults(func=cmd_add_rule) diff --git a/app/config.py b/app/config.py index 49363ef..569c693 100644 --- a/app/config.py +++ b/app/config.py @@ -38,6 +38,11 @@ # от источника внутри окна вернёт id первого события (deduplicated: true). dedup_window_seconds: int = 24 * 3600 + # Режим маршрутизации (docs/05, «Типы и правила»): + # all — применяются ВСЕ подошедшие правила (каждое добавляет доставки) + # first — только первое подошедшее по weight + routing_match_mode: str = "all" + # SPA лежит в образе рядом с приложением; в dev используется vite-сервер. spa_dist_dir: Path = BASE_DIR / "spa_static" diff --git a/app/models/dicts.py b/app/models/dicts.py index c3e981a..e18f19b 100644 --- a/app/models/dicts.py +++ b/app/models/dicts.py @@ -31,7 +31,7 @@ # Статусы конверта события (events.status): # queued — принято на приёме, ждёт воркера # processing — воркер маршрутизирует/раскидывает доставки -# done — все доставки в терминальном состоянии +# done — маршрутизация завершена; живые статусы — в deliveries # failed — событие не маршрутизировано (после финальных ретраев правил) EVENT_STATUSES = ("queued", "processing", "done", "failed") diff --git a/app/models/events.py b/app/models/events.py index a60088a..a14a85c 100644 --- a/app/models/events.py +++ b/app/models/events.py @@ -22,9 +22,13 @@ from app.database import Base -# Статусы доставки: pending (в работе) → delivered | failed. +# Статусы доставки: +# pending — в работе (квота у отправщика, ждёт next_retry_at) +# delivered — доставлено +# failed — все попытки исчерпаны +# skipped — срезано шумодавом правила (не доставлено намеренно, аудит) # Одна запись на доставку: счётчик attempts и next_retry_at, не новые строки. -DELIVERY_STATUSES = ("pending", "delivered", "failed") +DELIVERY_STATUSES = ("pending", "delivered", "failed", "skipped") class Event(Base): @@ -63,7 +67,9 @@ __tablename__ = "deliveries" __table_args__ = ( CheckConstraint("channel IN ('telegram','email','s2s','internal_log')", name="ck_channel"), - CheckConstraint("status IN ('pending','delivered','failed')", name="ck_delivery_status"), + CheckConstraint( + "status IN ('pending','delivered','failed','skipped')", name="ck_delivery_status" + ), Index("ix_deliveries_retry", "status", "next_retry_at"), ) diff --git a/app/models/rules.py b/app/models/rules.py index dc718c4..11e73f7 100644 --- a/app/models/rules.py +++ b/app/models/rules.py @@ -1,13 +1,12 @@ """Правила маршрутизации (Routing Engine, задача #30). -Условия — JSONB (MVP-набор полей, легко расширять): +Условия — JSONB: {"source": "monitoring" | null, # null = любой источник "subjects": [], # пусто = любые subject "actions": [], # пусто = любые action - "priority_min": "normal"} # в шкале low/normal/high/critical - -Payload-матчинг и теги в v1 контракта нет — условия набора хватает для -принятых примеров (см. docs/05-ingestion-api.md). + "priority_min": "normal", # в шкале low/normal/high/critical + "payload": {"container": "api"} # точное совпадение по top-level ключам + # payload'а; значение-список = any-of} """ from datetime import datetime @@ -31,8 +30,15 @@ class RoutingRule(Base): """Правило: условия → набор действий (каналы/цели/шаблоны). - weight — порядок применения: чем меньше, тем раньше. В MVP применяется - только первое подошедшее правило (later: first/all решается в #30). + weight — порядок применения (важен в режиме first). Режим all/first + задаётся настройкой ROUTING_MATCH_MODE (app/config.py, дефолт all). + + throttle_seconds — шумодав: если доставки в ту же (channel, target) + были созданы не позже N секунд назад (pending/delivered), эта — + записывается со статусом skipped. critical проходит всегда + (обход шумодавов — часть смысла критичности, см. docs/05). + + Теги — позже, когда появится реальное правило, требующее тега. """ __tablename__ = "routing_rules" @@ -42,6 +48,9 @@ enabled: Mapped[bool] = mapped_column(default=True) weight: Mapped[int] = mapped_column(default=0) conditions: Mapped[dict] = mapped_column(JSONB) + # Шумодав: не чаще одной доставки в (channel, target) за N секунд. + # None = без ограничения. critical НЕ обходится автоматически — см. докстринг. + throttle_seconds: Mapped[int | None] = mapped_column(Integer) # Шаблон по умолчанию для доставок правила (Jinja2/`{{ payload.x }}`); # individual action может переопределить. template: Mapped[str | None] = mapped_column(Text) diff --git a/app/worker/celery_app.py b/app/worker/celery_app.py index 4de2145..157c4ec 100644 --- a/app/worker/celery_app.py +++ b/app/worker/celery_app.py @@ -16,7 +16,14 @@ celery_app.conf.update( task_default_queue="synapse", task_track_started=True, - # TODO(#29): политика ретраев и таймаутов — уточнить при реализации доставок. + # Ретраи доставок: beat стартует вместе с воркером (`worker -B`), + # раз в 60 с просмотр pending с наступившим next_retry_at. + beat_schedule={ + "retry-due-deliveries": { + "task": "synapse.retry_due", + "schedule": 60.0, + }, + }, task_acks_late=True, worker_prefetch_multiplier=1, ) \ No newline at end of file diff --git a/app/worker/tasks.py b/app/worker/tasks.py index ca30ba2..a70e1b5 100644 --- a/app/worker/tasks.py +++ b/app/worker/tasks.py @@ -1,21 +1,22 @@ -"""Задачи воркера: ping + ingest (маршрутизация принятых событий). +"""Задачи воркера: ping, ingest (Routing Engine), retry_due (ретраи доставок). -Приём (api) кладёт только конверт в БД; здесь Routing Engine MVP: -первое подходящее enabled-правило по weight → создание Delivery по каждому -действию. Исполнение каналов (Telegram/Email/S2S) — задачи #29/#28; -internal_log исполняется сразу (это запись в эту же таблицу). +Приём (api) кладёт только конверт в БД; здесь живёт Routing Engine: +- подбор правил по условиям (source/subjects/actions/priority_min/payload); +- шумодав правила (throttle_seconds, critical проходит всегда); +- фактическая отправка каналов — задачи #29/#28; internal_log исполняется + сразу (delivery и есть запись лога), прочие каналы pending с ретраями. """ import uuid -from datetime import UTC, datetime +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, @@ -26,8 +27,31 @@ PRIORITY_ORDER = {"low": 0, "normal": 1, "high": 2, "critical": 3} +# Ретраи доставки: 1+4 попытки, пауза растёт (минуты → часы), после — failed. +RETRY_DELAYS = (30, 120, 600, 1800) -def _rule_matches(rule: RoutingRule, *, source_name: str, subject: str, action: str, priority: str) -> bool: +# Каналы, чью отправку воркер ещё не умеет (#29/#28) — честный фейл с ретраем. +_CHANNELS_NOT_IMPL = ("telegram", "email", "s2s") + + +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 @@ -40,6 +64,9 @@ 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 @@ -58,6 +85,29 @@ return template_text +def _throttled(db, rule: RoutingRule, channel: str, target_id: int | None, now: datetime) -> bool: + """Шумодав: была ли создана доставка в ту же цель в окне rule.throttle_seconds. + + Считаются pending/delivered (в полёте тоже держат канал занятым); + skipped окно не продлевает — иначе непрерывный шторм никогда не проходит. + critical не доходит до этого вызова (обход на месте вызова). + """ + if not rule.throttle_seconds or target_id is None: + return False + recent = ( + db.execute( + select(Delivery.id).where( + Delivery.channel == channel, + 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 -> очередь -> воркер.""" @@ -66,11 +116,11 @@ @celery_app.task(name="synapse.ingest") def ingest(event_id: str) -> str: - """Маршрутизация события: найти правило, создать deliveries. + """Маршрутизация события: подобрать правила, создать deliveries. - Первое подошедшее правило (по weight) применяется целиком; прочие - условия/полнота Routing Engine уточняются в #30. Если правил нет — - событие просто завершается (конверт остаётся в events как история). + Режим all (дефолт) — применяются все подошедшие правила; first — + только первое по weight. Если правил нет — событие завершается + (конверт остаётся в events как история). """ with SessionLocal() as db: event = db.get(Event, uuid.UUID(event_id)) @@ -85,50 +135,146 @@ rules = db.execute( select(RoutingRule).where(RoutingRule.enabled.is_(True)).order_by(RoutingRule.weight) ).scalars().all() - - matched: RoutingRule | None = None - for rule in rules: + 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, - ): - matched = rule - break + payload=dict(event.payload or {}), + ) + ] + if get_settings().routing_match_mode == "first": + matched = matched[:1] now = datetime.now(UTC) - if matched is not None: + n_skipped = 0 + for rule in matched: actions = db.execute( - select(RoutingRuleAction).where(RoutingRuleAction.rule_id == matched.id) + select(RoutingRuleAction).where(RoutingRuleAction.rule_id == rule.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)", + # шумодав: срезаем (с аудит-записью), critical проходит + if event.priority != "critical" and _throttled( + db, rule, act.channel, act.target_id, now + ): + 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=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 + + # Фактическая отправка — задачи #29/#28: pending на ретраи. + db.add( + Delivery( + event_id=event.id, + rule_id=rule.id, + channel=act.channel, + channel_target_id=act.target_id, + status="pending", + attempts=1, + next_retry_at=now + timedelta(seconds=RETRY_DELAYS[0]), + last_error="канал ещё не реализован (#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}" \ No newline at end of file + 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}" + + +def _attempt_send(delivery: Delivery, event: Event) -> tuple[bool, str | None]: + """Одна попытка отправки. Возвращает (ok, ошибка). + + Заглушка под провайдеров (#29/#28): internal_log в этом пути не бывает + (доставляется на маршрутизации), остальные каналы пока не умеем. + """ + if delivery.channel in _CHANNELS_NOT_IMPL: + return False, "канал ещё не реализован (#29/#28)" + return False, f"неизвестный канал {delivery.channel}" + + +@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 + ok, error = _attempt_send(delivery, event) + 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}" \ No newline at end of file diff --git a/docker/entrypoint.sh b/docker/entrypoint.sh index 89605f7..b482fbd 100644 --- a/docker/entrypoint.sh +++ b/docker/entrypoint.sh @@ -9,8 +9,8 @@ exec uvicorn app.main:app --host 0.0.0.0 --port 8000 ;; worker) - echo "[synapse] запускаю celery-воркер..." - exec celery -A app.worker.celery_app worker --loglevel=info + echo "[synapse] запускаю celery-воркер (+beat для ретраей)..." + exec celery -A app.worker.celery_app worker -B --loglevel=info ;; *) exec "$@" diff --git a/docs/04-database.md b/docs/04-database.md index 8f34295..9a9f742 100644 --- a/docs/04-database.md +++ b/docs/04-database.md @@ -1,6 +1,6 @@ # 04 — Схема БД (#34), v1 -Реализация — `app/models/` + миграция `alembic/versions/20261003_0001_initial_schema.py` (7 таблиц). Постгрес в докере (postgres:17-alpine, том `pgdata`), Redis 7 (том `redisdata`). +Реализация — `app/models/` + миграции `alembic/versions/` (7 таблиц; `0001` — исходная схема, `0002/0003` — `description` и полный Routing Engine). Постгрес в докере (postgres:17-alpine, том `pgdata`), Redis 7 (том `redisdata`). Два контура данных: **конфигурация** (человек заводит через админку, меняется редко) и **поток событий** (пишется автоматически, растёт всегда). @@ -16,11 +16,11 @@ | Таблица | Что хранит | Ключевые поля | |---|---|---| -| `sources` | Источники событий (регистрирует админ) | `name` (slug, unique: `monitoring`), `label` | +| `sources` | Источники событий (регистрирует админ) | `name` (slug, unique: `monitoring`), `label`, `description` («что это за сервис» для человека и ИИ-агента) | | `api_keys` | Ключи `syn_…` к источникам | **`token_hash`** (sha256, unique), `token_hint` (последние 4 для UI), `revoked_at` (ротация: новый ключ, старая строка маркируется). Ключ — удостоверение; правила матчают `source`, не токен | | `notification_types` | Реестр типов из контракта | тройка `(source_id, subject, action)` уникальна, `payload_schema` JSONB (опц., валидирует воркер) | -| `channel_targets` | Цели каналов: куда доставлять | `channel` ∈ `telegram/email/s2s/internal_log` (CHECK), `config` JSONB (chat_id, адрес, s2s-endpoint). **Креды каналов (токен бота, SMTP) — не здесь, а в .env** | -| `routing_rules` | Правила «условия → действия» | `conditions` JSONB (`source`, `subjects[]`, `actions[]`, `priority_min`), `template` (Jinja2 `{{ payload.x }}`), `weight` (порядок; в MVP применяется первое подошедшее), `enabled` | +| `channel_targets` | Цели каналов: куда доставлять | `channel` ∈ `telegram/email/s2s/internal_log` (CHECK), `config` JSONB (chat_id, адрес, s2s-endpoint; s2s ещё `token_ref` — ссылка на секрет), `description` («что это за клиент-сервис»). **Креды каналов (токен бота, SMTP, HMAC-секреты) — не здесь, а в .env** | +| `routing_rules` | Правила «условия → действия» | `conditions` JSONB (`source`, `subjects[]`, `actions[]`, `priority_min`, `payload` — точный матч по top-level ключам), `template` (Jinja2 `{{ payload.x }}`), `weight` (порядок, важен в режиме first), `throttle_seconds` (шумодав: окно на (channel, target), critical проходит всегда), `enabled` | | `routing_rule_actions` | Действия правила | `rule_id`, `channel` (CHECK), `target_id` → channel_targets (SET NULL при удалении цели), `template` — переопределение шаблона правила | ## Поток событий @@ -31,8 +31,8 @@ | Таблица | Что хранит | Ключевые поля | |---|---|---| -| `events` | Конверт события (контракт docs/05) | UUID pk, `source_id` (RESTRICT — события переживают удаление источника), `subject`/`action`/`priority` (CHECK по шкале), `payload` JSONB, `dedup_key` (+индекс с `created_at` — окно дедупликации), `expires_at` (= created+ttl_seconds), `scheduled_at` (резерв), `status` ∈ `queued/processing/done/failed` | -| `deliveries` | Доставка в конкретную цель | `event_id` (CASCADE), `rule_id` (SET NULL — правило удалим, историю оставим), `channel`, `channel_target_id`, `status` ∈ `pending/delivered/failed`, `attempts`, `last_error`, `next_retry_at` (+индекс `status,next_retry_at` — скан ретраев), `rendered_message` (аудит: что реально ушло в канал) | +| `events` | Конверт события (контракт docs/05) | UUID pk, `source_id` (RESTRICT — события переживают удаление источника), `subject`/`action`/`priority` (CHECK по шкале), `payload` JSONB, `dedup_key` (+индекс с `created_at` — окно дедупликации), `expires_at` (= created+ttl_seconds), `scheduled_at` (резерв), `status` ∈ `queued/processing/done/failed` (done — замаршрутизировано; живые статусы в deliveries) | +| `deliveries` | Доставка в конкретную цель | `event_id` (CASCADE), `rule_id` (SET NULL — правило удалим, историю оставим), `channel`, `channel_target_id`, `status` ∈ `pending/delivered/failed/skipped` (skipped — срезано шумодавом), `attempts`, `last_error`, `next_retry_at` (+индекс `status,next_retry_at` — скан ретраей beat'ом), `rendered_message` (аудит: что реально уйдёт в канал) | ## Что сознательно НЕ в таблицах diff --git a/docs/05-ingestion-api.md b/docs/05-ingestion-api.md index c74c416..5beda8b 100644 --- a/docs/05-ingestion-api.md +++ b/docs/05-ingestion-api.md @@ -71,10 +71,44 @@ Обратный канал (webhook об изменении статуса) позже; контракт поллинга стабилен — подписки докинутся сверху. Для MVP источники либо fire-and-forget, либо поллят этот GET. -## Типы и правила (админка) +## Типы и правила (админка) — Routing Engine -- **Реестр типов**: тройка `(source, subject, action)`, опционально JSON Schema payload'а. Заводит админ Synapse в UI. -- **Правило маршрутизации**: условия (`source`, `subject`, набор `actions`, `priority >= X`, позже — payload-матчинг и теги) → действия (каналы + шаблон + цели: TG-чат, SMTP, s2s-ендпоинт). +### Реестр типов + +Тройка `(source, subject, action)`, опционально JSON Schema payload'а (валидация мягкая, в воркере). Заводит админ Synapse в UI. + +### Правило маршрутизации + +Условия (`conditions` JSONB): + +```json +{ + "source": "monitoring", // null/отсутствует = любой источник + "subjects": ["container"], // [] = любые + "actions": ["down", "restarting"], // [] = любые + "priority_min": "high", // шкала low/normal/high/critical + "payload": {"container": "nomin-web"} // опц.: точный матч по top-level + // ключам payload'а; значение-список + // = any-of («хост melody или pilar») +} +``` + +Действия — по одному на строку: канал + цель (`telegram:infra-alerts`, `internal_log`, позднее `email`, `s2s:navi-rei`) + опциональный шаблон (Jinja2, `{{ payload.x }}`; шаблон действия переопределяет шаблон правила). Шаблон обязан быть устойчив к отсутствию полей payload. + +### Сколько правил применяется + +Задаёт настройка `ROUTING_MATCH_MODE`: + +- `all` (дефолт) — **все** подошедшие правила; каждое добавляет свои доставки. «Контейнеру вниз → TG», «critical → Navi» и «всё от monitoring → лог» работают одновременно. +- `first` — только первое по `weight` (меньше — раньше). Суровые сценарии «одно правило на тип». + +### Шумодав (`throttle_seconds`) + +Правило может ограничить частоту: если в ту же цель `(channel, target)` уже была создана доставка (`pending`/`delivered`) за последние N секунд — новая записывается со статусом **`skipped`** (аудит: `rendered_message` хранит, что отправили бы). Окно продлевают только реальные доставки — непрерывный шторм срезанный не продлевает. **`priority: critical` проходит шумодав всегда** — это часть смысла критичности (обещано в конверте). `low`-интервалы типа «контейнер перезапускается 20 раз» — тот же механизм. + +### Ретраи доставок + +Провал доставки (сеть, 5xx провайдера) → обратно-экспоненциальная пауза: `30 c → 2 м → 10 м → 30 м`, всего 5 попыток → `failed`. Скан due-доставок — Celery beat раз в минуту (`synapse.retry_due`). Пока канал не реализован (#29/#28), доставки честно проходят этот цикл и падают в `failed` — таблица сразу показывает поведение ретраев. ## Доставка s2s: подпись и верификация