diff --git a/.env.example b/.env.example index ad9a50c..f7049e7 100644 --- a/.env.example +++ b/.env.example @@ -34,6 +34,10 @@ # Режим правил: all (все подошедшие) | first (только первое по weight). ROUTING_MATCH_MODE=all +# --- Секреты s2s-целей (docs/05, HMAC-подпись доставок) --- +# channel_targets.config несёт token_ref ("navi-rei"); значение секрета здесь: +# S2S_SECRET_NAVI_REI= + # --- Misc --- # Порт API на хосте. SYNAPSE_PORT=8013 \ No newline at end of file diff --git a/README.md b/README.md index 017b8bb..95a868f 100644 --- a/README.md +++ b/README.md @@ -75,7 +75,7 @@ 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, раздел «Типы и правила». +Ретраи (#30): провал доставки → пауза 30 с → 2 м → 10 м → 30 м, 5 попыток → `failed`; скан due-доставок — beat воркера раз в минуту (`synapse.retry_due`). Шумодав — `--throttle N` у правила (не чаще одной доставки в цель за N с, critical проходит всегда); все подошедшие правила применяются (режим `all`, `ROUTING_MATCH_MODE=first` — только первое). Подробности — docs/05, раздел «Типы и правила». ## Вход в админку @@ -102,7 +102,8 @@ models/ схема БД (#34): источники, ключи, типы, правила, события, доставки auth/ SSO-валидация + гейт ролей admin/superadmin api/routes.py healthz, readyz, admin/me - worker/ Celery: celery_app, tasks + signature.py HMAC подпись/проверка вебхуков (s2s + входящие gnexus-auth) + worker/ Celery: celery_app, tasks, senders (s2s) alembic/ миграции (env.py читает DATABASE_URL из .env) docs/ docs/04 — схема БД, docs/05 — контракт Ingestion API frontend/ Vue 3 SPA админки (сборка кладётся в spa_static/ образа) diff --git a/app/signature.py b/app/signature.py new file mode 100644 index 0000000..33d18e7 --- /dev/null +++ b/app/signature.py @@ -0,0 +1,38 @@ +"""HMAC-подпись вебхуков — одна схема на всю экосистему (docs/05). + +Зеркало gnexus-auth `app/Domain/Webhooks/WebhookSignature.php`: + + sig = "t=,v1=" + hex(hmac_sha256(".", secret)) + +verify — константное сравнение (hmac.compare_digest, аналог hash_equals) ++ окно свежести для защиты от replay. Референс-реализация для +отправки (`app/worker/senders.py`) и приёма (входящие вебхуки gnexus-auth). +""" + +import hashlib +import hmac +import time + + +def make_signature(body: bytes, secret: str, timestamp: int | None = None) -> str: + """Подпись по «сырому» телу запроса (байт в байт, до любого парсинга).""" + if timestamp is None: + timestamp = int(time.time()) + digest = hmac.new( + secret.encode(), f"{timestamp}.".encode() + body, hashlib.sha256 + ).hexdigest() + return f"t={timestamp},v1={digest}" + + +def verify_signature(body: bytes, header: str, secret: str, max_skew: int = 300) -> bool: + """Проверка заголовка «t=…,v1=…»: свежесть + константное сравнение.""" + try: + parts = dict(part.split("=", 1) for part in header.split(",")) + timestamp = int(parts["t"].strip()) + except (TypeError, AttributeError, ValueError): + return False + if abs(int(time.time()) - timestamp) > max_skew: + return False + expected = make_signature(body, secret, timestamp).split("v1=", 1)[1] + supplied = parts.get("v1", "").strip() + return hmac.compare_digest(expected, supplied) \ No newline at end of file diff --git a/app/worker/senders.py b/app/worker/senders.py new file mode 100644 index 0000000..2ea3888 --- /dev/null +++ b/app/worker/senders.py @@ -0,0 +1,82 @@ +"""Отправка в системные webhooks (s2s): подпись app/signature.py, docs/05. + +Секрет цели — per-target: channel_targets.config = {"endpoint": …, "token_ref": +"navi-rei"}; значение секрета — в окружении (S2S_SECRET_NAVI_REI, .env / +gnexus-creds), в БД только ссылка. Ротация секрета цели не трогает правила. +""" + +import json +import os +import time + +import httpx + +from app.models import ChannelTarget, Event +from app.signature import make_signature + +#: Референс-получателю достаточно этих трёх; остальное — настройка цели. +S2S_TIMEOUT = httpx.Timeout(10.0) + + +def secret_env_name(token_ref: str) -> str: + """token_ref «navi-rei» → переменная окружения S2S_SECRET_NAVI_REI.""" + cleaned = token_ref.upper().replace("-", "_").replace(".", "_") + return f"S2S_SECRET_{cleaned}" + + +def resolve_secret(token_ref: str) -> str | None: + return os.environ.get(secret_env_name(token_ref)) or None + + +def send_s2s(target: ChannelTarget, event: Event, source_name: str) -> tuple[bool, str | None]: + """POST конверта события в endpoint цели, подписанный по схеме gnexus-auth. + + Тело — конверт без изменений + event_id (см. docs/05): + описание источника в него не входит. Успех — любой 2xx; прочее — + провал попытки (воркер поставит ретрай). + """ + cfg = target.config or {} + endpoint = cfg.get("endpoint") + token_ref = cfg.get("token_ref") + if not endpoint: + return False, "s2s: в config цели нет endpoint" + secret = resolve_secret(token_ref) if token_ref else None + if not secret: + return False, ( + f"s2s: секрет цели не задан в окружении " + f"({secret_env_name(token_ref) if token_ref else 'нет token_ref в config'})" + ) + + envelope = { + "event_id": str(event.id), + "source": source_name, + "subject": event.subject, + "action": event.action, + "priority": event.priority, + "payload": dict(event.payload or {}), + } + body = json.dumps(envelope, ensure_ascii=False, separators=(",", ":")).encode() + headers = { + "Content-Type": "application/json", + "X-Gnexus-Event-Id": str(event.id), + "X-Gnexus-Event-Type": f"{source_name}.{event.subject}.{event.action}", + "X-Gnexus-Event-Timestamp": str(int(time.time())), + "X-Gnexus-Signature": make_signature(body, secret), + "X-Synapse-Source": source_name, + "User-Agent": "Synapse/0.1 (gnexus notification hub)", + } + try: + response = httpx.post( + endpoint, + content=body, + headers=headers, + timeout=S2S_TIMEOUT, + verify=cfg.get("verify_tls", True), + follow_redirects=False, + ) + except httpx.HTTPError as err: + return False, f"s2s: сеть: {err}" + + if 200 <= response.status_code < 300: + return True, None + return False, f"s2s: HTTP {response.status_code} · {response.text[:200]}" \ No newline at end of file diff --git a/app/worker/tasks.py b/app/worker/tasks.py index a70e1b5..9ae0869 100644 --- a/app/worker/tasks.py +++ b/app/worker/tasks.py @@ -3,8 +3,9 @@ Приём (api) кладёт только конверт в БД; здесь живёт Routing Engine: - подбор правил по условиям (source/subjects/actions/priority_min/payload); - шумодав правила (throttle_seconds, critical проходит всегда); -- фактическая отправка каналов — задачи #29/#28; internal_log исполняется - сразу (delivery и есть запись лога), прочие каналы pending с ретраями. +- отправка: internal_log сразу (delivery и есть запись лога), s2s — подписано + (senders.py, docs/05); Telegram/Email — заглушки под #29; провалы + — ретрай-цикл через beat (synapse.retry_due). """ import uuid @@ -17,6 +18,7 @@ from app.config import get_settings from app.database import SessionLocal from app.models import ( + ChannelTarget, Delivery, Event, RoutingRule, @@ -24,14 +26,37 @@ 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/#28) — честный фейл с ретраем. -_CHANNELS_NOT_IMPL = ("telegram", "email", "s2s") +# Каналы, чью отправку воркер ещё не умеет (#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: @@ -196,19 +221,27 @@ ) 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)", - ) + # Фактическая отправка (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() @@ -217,17 +250,6 @@ 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). @@ -259,7 +281,8 @@ delivery.status = "failed" delivery.last_error = "событие удалено вместе с доставкой" continue - ok, error = _attempt_send(delivery, event) + 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 diff --git a/docs/05-ingestion-api.md b/docs/05-ingestion-api.md index 5beda8b..25482f8 100644 --- a/docs/05-ingestion-api.md +++ b/docs/05-ingestion-api.md @@ -136,7 +136,7 @@ 2. Сравнить через **константное время** (PHP `hash_equals`, Python `hmac.compare_digest`) — не через `==`. 3. Свежесть: `|now - t|` в пределах допуска (gnexus-auth — 5 минут) → защита от replay. -Секрет — **свой у каждой цели** (per-target, поле `token_ref` в `channel_targets.config`), значения в `.env`/gnexus-creds, в БД только ссылки. Ротация секрета цели не затрагивает правила маршрутизации. +Секрет — **свой у каждой цели** (per-target, поле `token_ref` в `channel_targets.config`), значения в `.env`/gnexus-creds (переменная `S2S_SECRET_`), в БД только ссылки. Ротация секрета цели не затрагивает правила маршрутизации. Цепочка доверия: получатель проверил подпись ⇒ целостность и авторство Synapse. Синтезировать чужое событие Synapse не может — на приёме его ключ и `source` сверились бы с реестром (403), а ретрансляция чужого ключом источника невозможна, поскольку событие должно нести `source`, совпадающий с ключом (анти-спуфинг). @@ -153,6 +153,8 @@ `description` — атрибут записи источника в реестре Synapse (админка), а не поля события: конверт и s2s-вебхук его не несут. Получатель видит имя источника в `X-Synapse-Source`; зачем оно — смотрит в админке Synapse, где у каждого источника/цели прописано описание («что это за сервис» для человека и ИИ-агента, не гадать по названию). +Отправка: первая попытка — сразу при маршрутизации события в воркере; провал → общий ретрай-цикл (30 с → 2 м → 10 м → 30 м). Референс-код Synapse: `app/signature.py` (подпись/проверка — зеркало php-класса выше), `app/worker/senders.py` (отправщик). Эту пару модулей можно заимствовать в принимающем сервисе (только Python; Navi-инстансы на другом стеке — по этой же спецификации). + ## Таблицы БД (постановка #34) `api_keys` (хэш, привязка к source) · `notification_types(source, subject, action)` · `routing_rules` (условия, действия, шаблон) · `events` (конверт + payload + статус) · `deliveries` (канал, цель, статус, попытки, ошибки) · `channel_targets`. \ No newline at end of file