diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..a5e2e2f --- /dev/null +++ b/.gitignore @@ -0,0 +1,13 @@ +# Python +__pycache__/ +*.py[cod] +*.egg-info/ +.venv/ +build/ +dist/ + +# Tooling caches +.pytest_cache/ +.ruff_cache/ +.mypy_cache/ +.coverage \ No newline at end of file diff --git a/README.md b/README.md new file mode 100644 index 0000000..a7880d4 --- /dev/null +++ b/README.md @@ -0,0 +1,108 @@ +# gnexus-synapse — Python-клиент Synapse + +Тонкий клиент Ingestion API v1 хаба уведомлений Gnexus Synapse +(`git+https://git.gnexus.space/git/root/gn-synapse.git`): собрать конверт, +проверить локально, послать — без ретраев и локальных очередей (очередь +у Synapse своя, приём отвечает мгновенно). Контракт — `docs/05-ingestion-api.md` +в репозитории `gn-synapse`. + +## Установка + +```bash +pip install "git+https://git.gnexus.space/git/root/gn-synapse-client-py.git@v0.1.0" +``` + +Серверный репозиторий — `https://git.gnexus.space/git/root/gn-synapse.git`. +Зависимость одна: `httpx`. + +## Quickstart + +```python +from gnexus_synapse import SynapseClient + +# аргументы сильнее env; фолбэк — SYNAPSE_URL / SYNAPSE_API_KEY +syn = SynapseClient("http://synapse.gnexus.space:8013", "syn_...") + +event = syn.send("bugtrail", "test", "failed", + payload={"user_id": "", "error": "..."}, + priority="high", dedup_key="run-2026-10-03") +print(event.id, event.status) # queued + +# пожарная отправка: ошибки Synapse не ломают сервис, пишутся в warning +syn.emit("gntodo", "task", "created", payload={"user_id": uid}) + +status = syn.status(event.id) # статусы доставок по каналам +result = syn.send_batch([...]) # batch: result.accepted / result.rejected +``` + +env-фолбэк (для сервисов без своей конфигурации): + +``` +SYNAPSE_URL=http://synapse:8013 +SYNAPSE_API_KEY=syn_... # выдаёт админ (create-key / MCP key_issue), печатается один раз +SYNAPSE_DEFAULT_SOURCE=bugtrail # необязательно: не указывать source в каждом вызове +SYNAPSE_TIMEOUT=10 # секунды +``` + +`default_source` — лучшая защита от частой ошибки 403: ключ принадлежит +одному источнику, и `source` в конверте обязан совпадать с ним. + +## API + +| Метод | Что делает | Ошибки | +|---|---|---| +| `send(source, subject, action, *, priority, payload, dedup_key, ttl_seconds, scheduled_at)` | POST /api/v1/events → `Event(id, status, deduplicated)` | типизированные, см. ниже | +| `emit(...)` | тот же, ловит SynapseError → warning в логгер, возвращает `Event \| None` | не бросает | +| `send_batch(events)` | массив конвертов; локально битые не посылаются; rejected — **не** исключение (частичный успех норма); транспортный сбой/5xx — исключение | `SynapseError` | +| `status(event_id)` | GET /api/v1/events/{id} → `EventStatus` (в т.ч. `deliveries`) | `SynapseNotFoundError` | +| `health()` / `ready()` | диагностика без ключа | — | +| `close()` | закрыть свой httpx.Client (инжектный не трогает) | — | + +Исключения (все — наследники `SynapseError` с полями `detail`, `status_code`): + +| Класс | Значение | +|---|---| +| `SynapseConfigError` | нет url/api_key/source в аргументах и env | +| `SynapseConnectionError` | Synapse недоступен (DNS/timeout/refused), `__cause__` — оригинал | +| `SynapseAuthError` | 401 — ключ отсутствует/неизвестен/отозван (или источник в архиве) | +| `SynapseForbiddenError` | 403 — `source` конверта ≠ источник ключа; задайте `default_source` | +| `SynapseValidationError` | 422 сервера или **локальная** ошибка конверта (`status_code is None`) | +| `SynapseNotFoundError` | 404 — событие не у источника ключа | +| `SynapseServerError` | 5xx — событие не принято, повторите при желании | + +Локальная валидация зеркалит серверную (`^[a-z0-9]([a-z0-9._-]*[a-z0-9])?$`, +1..64; priority low/normal/high/critical; ttl 1..604800; dedup_key ≤ 255; +payload — JSON-объект), поэтому SynapseValidationError от клиента — почти +всегда ошибка кода, ловить её нужно в тестах, а не в рантайме. +Сервер игнорирует незарегистрированные пары (subject, action) — их регистрирует +админ (`add-type` / MCP `type_register`). + +Асинхронный близнец для FastAPI-сервисов: + +```python +from gnexus_synapse import AsyncSynapseClient +syn = AsyncSynapseClient() # httpx.AsyncClient, await syn.aclose() +event = await syn.send("gntodo", "task", "created", payload={"user_id": uid}) +``` + +## Конвенции и оговорки + +- **Дедуп best-effort** (окно 24 ч, сервер не даёт гарантии уникальности): + не стройте бизнес-логику на отсутствии дублей. +- **Дедуп вернул 202 с `deduplicated=true`** — это id *первого* события, не ошибки. +- `payload.user_id` — конвенция «о ком событие» (нестроковое значение → адресные + доставки skipped): клиент пишет warning, но не блокирует. +- `scheduled_at` — резерв контракта v1: передаётся, сервер пока не обрабатывает. +- Свой httpx.Client создаётся с явным `timeout` (у httpx дефолт — «бесконечно»). + Инжектный клиент (`http_client=...`) не трогается — таймаут настраивайте сами. +- user_agent `gnexus-synapse-py/<версия>` — по нему Synapse видит, кто источался. + +## Разработка + +```bash +uv venv -p 3.12 && uv pip install -e . -e '.[dev]' +pytest # юнит-тесты (MockTransport) +ruff check src tests && mypy src +# интеграционный смок против живого стека: +SYNAPSE_URL=http://localhost:8013 SYNAPSE_API_KEY=syn_... python examples/plain/smoke.py +``` \ No newline at end of file diff --git a/examples/plain/smoke.py b/examples/plain/smoke.py new file mode 100644 index 0000000..2207228 --- /dev/null +++ b/examples/plain/smoke.py @@ -0,0 +1,83 @@ +"""Интеграционный смок против живого Synapse. + +Запуск (токен выдаёт `create-key` на сервере, он печатается один раз): + + SYNAPSE_URL=http://localhost:8013 SYNAPSE_API_KEY=syn_... \ + python examples/plain/smoke.py + +Внутри контейнера api URL — http://localhost:8000. Источник и тип +(source libtest, тип ping/done) должны быть зарегистрированы на сервере. +""" + +from __future__ import annotations + +import os +import sys +import time + +from gnexus_synapse import ( + SynapseClient, + SynapseConfigError, + SynapseConnectionError, + SynapseForbiddenError, + SynapseNotFoundError, + SynapseValidationError, +) + +# уникальный на прогон: дедуп-окно сервера 24 ч, повтор полный — смок увидит не queued +DEDUP = f"client-smoke-{int(time.time())}" + + +def main() -> int: + client = SynapseClient() + print(" health:", client.health()) + + event = client.send("libtest", "ping", "done", payload={"user_id": "mcp"}, + dedup_key=DEDUP) + print(f" send: id={event.id} status={event.status}") + assert event.status == "queued" and not event.deduplicated + + repeat = client.send("libtest", "ping", "done", payload={"user_id": "mcp"}, + dedup_key=DEDUP) + print(f" repeat: id={repeat.id} deduplicated={repeat.deduplicated}") + assert repeat.deduplicated and repeat.id == event.id # дедуп вернул первое событие + + status = client.status(event.id) + print(f" status: {status.status}; доставки:") + for d in status.deliveries: + target = f" → {d.target}" if d.target else "" + err = f" ({d.error})" if d.error else "" + print(f" {d.channel}{target}: {d.status} ×{d.attempts}{err}") + + try: + client.send("wrong-source", "ping", "done") # анти-спуфинг сервера + except SynapseForbiddenError as ex: + print(f" 403 пойман (анти-спуфинг): {ex.detail}") + else: + print(" ОШИБКА: 403 не получен") + return 1 + + try: + client.status("00000000-0000-0000-0000-000000000000") + except SynapseNotFoundError: + print(" 404 пойман (чужой/несуществующий id)") + + # валидный конверт на недоступный хост → типизированное соединение-исключение + try: + SynapseClient("http://localhost:9999", "syn_x", timeout=1.0).send("libtest", "a", "b") + except SynapseConnectionError as ex: + print(f" соединение-ошибка поймана: {ex}") + except SynapseConfigError as ex: + print(f" config-ошибка поймана: {ex}") + except SynapseValidationError as ex: + print(f" локальная валидация поймана: {ex}") + + print(" ALL OK") + return 0 + + +if __name__ == "__main__": + if not os.environ.get("SYNAPSE_API_KEY"): + print("задайте SYNAPSE_URL и SYNAPSE_API_KEY", file=sys.stderr) + sys.exit(2) + sys.exit(main()) \ No newline at end of file diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..9afe2cb --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,56 @@ +[build-system] +requires = ["setuptools>=61.0", "wheel"] +build-backend = "setuptools.build_meta" + +[project] +name = "gnexus-synapse" +version = "0.1.0" +description = "Thin client library for gn-synapse (Gnexus notification hub) event ingestion" +readme = "README.md" +requires-python = ">=3.11" +license = {text = "Proprietary"} +authors = [ + {name = "GNexus"}, +] +classifiers = [ + "Development Status :: 3 - Alpha", + "Intended Audience :: Developers", + "Programming Language :: Python :: 3", + "Programming Language :: Python :: 3.11", + "Programming Language :: Python :: 3.12", + "Programming Language :: Python :: 3.13", +] +dependencies = [ + "httpx>=0.24.0", +] + +[project.optional-dependencies] +dev = [ + "pytest>=7.0", + "pytest-cov>=4.0", + "ruff>=0.1.0", + "mypy>=1.0", +] + +[project.urls] +Homepage = "https://git.gnexus.space/git/root/gn-synapse-client-py" + +[tool.setuptools.packages.find] +where = ["src"] + +[tool.ruff] +line-length = 120 +target-version = "py311" + +[tool.ruff.lint] +select = ["E", "F", "W", "I", "N", "UP", "B", "C4", "SIM"] + +[tool.mypy] +python_version = "3.11" +strict = true +warn_return_any = true +warn_unused_ignores = true + +[tool.pytest.ini_options] +testpaths = ["tests"] +pythonpath = ["src"] \ No newline at end of file diff --git a/src/gnexus_synapse/__init__.py b/src/gnexus_synapse/__init__.py new file mode 100644 index 0000000..26a71c3 --- /dev/null +++ b/src/gnexus_synapse/__init__.py @@ -0,0 +1,45 @@ +"""gnexus-synapse — тонкий клиент Ingestion API хаба уведомлений Gnexus Synapse. + +Контракт: docs/05-ingestion-api.md репозитория gn-synapse. +""" + +from ._async import AsyncSynapseClient +from .client import USER_AGENT, SynapseClient +from .config import ENV_API_KEY, ENV_DEFAULT_SOURCE, ENV_TIMEOUT, ENV_URL, SynapseConfig +from .dto import BatchRejection, BatchResult, Delivery, Event, EventStatus +from .exceptions import ( + SynapseAuthError, + SynapseConfigError, + SynapseConnectionError, + SynapseError, + SynapseForbiddenError, + SynapseNotFoundError, + SynapseServerError, + SynapseValidationError, +) + +__version__ = "0.1.0" + +__all__ = [ + "AsyncSynapseClient", + "SynapseClient", + "SynapseConfig", + "Event", + "BatchResult", + "BatchRejection", + "EventStatus", + "Delivery", + "SynapseConfigError", + "SynapseConnectionError", + "SynapseAuthError", + "SynapseForbiddenError", + "SynapseNotFoundError", + "SynapseValidationError", + "SynapseServerError", + "SynapseError", + "ENV_URL", + "ENV_API_KEY", + "ENV_TIMEOUT", + "ENV_DEFAULT_SOURCE", + "USER_AGENT", +] diff --git a/src/gnexus_synapse/_async.py b/src/gnexus_synapse/_async.py new file mode 100644 index 0000000..fd0e1a4 --- /dev/null +++ b/src/gnexus_synapse/_async.py @@ -0,0 +1,225 @@ +"""Асинхронный близнец SynapseClient — тот же контракт, httpx.AsyncClient. + +FastAPI-сервисам: блокирующий httpx-вызов внутри async-роута держит event +loop. Логика (валидация, маппинг ошибок, разбор ответов) шарится с sync — +здесь только транспорт. +""" + +from __future__ import annotations + +import logging +import uuid +from collections.abc import Mapping, Sequence +from typing import Any + +import httpx + +from . import envelope +from .client import API_PREFIX, USER_AGENT, _accepted_event, _dt, _error_for, _json_body +from .config import SynapseConfig +from .dto import BatchRejection, BatchResult, Delivery, Event, EventStatus +from .exceptions import ( + SynapseConfigError, + SynapseConnectionError, + SynapseError, + SynapseValidationError, +) + +_DEFAULT_LOGGER = "gnexus.synapse" + + +class AsyncSynapseClient: + def __init__( + self, + url: str | None = None, + api_key: str | None = None, + *, + timeout: float | None = None, + http_client: httpx.AsyncClient | None = None, + default_source: str | None = None, + logger: logging.Logger | None = None, + ) -> None: + base = SynapseConfig.from_env() + self._config = SynapseConfig( + base_url=url if url is not None else base.base_url, + api_key=api_key if api_key is not None else base.api_key, + timeout=timeout if timeout is not None else base.timeout, + default_source=default_source if default_source is not None else base.default_source, + ) + self._logger = logger or logging.getLogger(_DEFAULT_LOGGER) + self._own_client = http_client is None + self._http = http_client or httpx.AsyncClient( + timeout=self._config.timeout, + headers={"user-agent": USER_AGENT}, + ) + + async def send( + self, + source: str | None = None, + subject: str = "", + action: str = "", + *, + priority: str = "normal", + payload: dict[str, Any] | None = None, + dedup_key: str | None = None, + ttl_seconds: int | None = None, + scheduled_at: Any = None, + ) -> Event: + envelope_ = self._build(source, subject, action, priority=priority, payload=payload, + dedup_key=dedup_key, ttl_seconds=ttl_seconds, + scheduled_at=scheduled_at) + if envelope_.get("payload"): + envelope.warn_bad_user_id(envelope_["payload"], self._logger) + return _accepted_event(await self._post(f"{API_PREFIX}/events", envelope_)) + + async def emit( + self, + source: str | None = None, + subject: str = "", + action: str = "", + *, + priority: str = "normal", + payload: dict[str, Any] | None = None, + dedup_key: str | None = None, + ttl_seconds: int | None = None, + scheduled_at: Any = None, + ) -> Event | None: + try: + return await self.send(source, subject, action, priority=priority, payload=payload, + dedup_key=dedup_key, ttl_seconds=ttl_seconds, + scheduled_at=scheduled_at) + except SynapseError as ex: + self._logger.warning("synapse emit failed: %s", ex) + return None + + async def send_batch(self, events: Sequence[Mapping[str, Any]]) -> BatchResult: + positions: list[int] = [] + payload: list[dict[str, Any]] = [] + results: list[Event | BatchRejection | None] = [None] * len(events) + for i, raw in enumerate(events): + try: + source = raw.get("source") + payload.append(envelope.build_envelope( + source if source is not None else self._resolve_source(None), + str(raw.get("subject", "")), + str(raw.get("action", "")), + priority=raw.get("priority", "normal"), + payload=raw.get("payload"), + dedup_key=raw.get("dedup_key"), + ttl_seconds=raw.get("ttl_seconds"), + scheduled_at=raw.get("scheduled_at"), + )) + positions.append(i) + except (SynapseValidationError, SynapseConfigError) as ex: + results[i] = BatchRejection(index=i, detail=str(ex)) + + if positions: + resp = await self._post(f"{API_PREFIX}/events/batch", payload) + body = _json_body(resp) + server_results = body.get("results", []) if isinstance(body, dict) else [] + for pos, item in zip(positions, server_results, strict=False): + if not isinstance(item, dict): + results[pos] = BatchRejection(index=pos, detail=f"битый элемент ответа: {item!r}") + elif item.get("ok") is True: + results[pos] = Event( + id=str(item.get("id", "")), + status=str(item.get("status", "")), + deduplicated=bool(item.get("deduplicated", False)), + ) + else: + # index у сервера — позиция в отправленном массиве; ремапим в оригинальную + sent_index = item.get("index", pos) + orig = sent_index if isinstance(sent_index, int) and 0 <= sent_index < len(positions) else pos + results[pos] = BatchRejection( + index=positions[orig], detail=str(item.get("detail", "")) + ) + return BatchResult(results=[r for r in results if r is not None]) + + async def status(self, event_id: str) -> EventStatus: + try: + uuid.UUID(event_id) + except (ValueError, AttributeError) as ex: + raise SynapseValidationError(f"event_id={event_id!r} не uuid") from ex + resp = await self._request("GET", f"{API_PREFIX}/events/{event_id}") + body = _json_body(resp) + if resp.status_code != 200: + raise _error_for(resp, body) + deliveries = [ + Delivery( + channel=str(d.get("channel", "")), + target=d.get("target"), + status=str(d.get("status", "")), + attempts=int(d.get("attempts", 0)), + error=d.get("error"), + rendered_message=d.get("rendered_message"), + ) + for d in body.get("deliveries", []) + ] + return EventStatus( + id=str(body.get("id", event_id)), + source=str(body.get("source", "")), + subject=str(body.get("subject", "")), + action=str(body.get("action", "")), + priority=str(body.get("priority", "")), + status=str(body.get("status", "")), + created_at=_dt(body["created_at"]), + expires_at=_dt(body.get("expires_at")) if body.get("expires_at") else None, + deliveries=deliveries, + ) + + async def aclose(self) -> None: + if self._own_client: + await self._http.aclose() + + def _resolve_source(self, source: str | None) -> str: + source = source or self._config.default_source + if not source: + raise SynapseConfigError( + "не задан source: передайте source= или default_source (env SYNAPSE_DEFAULT_SOURCE)" + ) + return source + + def _require_key(self) -> str: + if not self._config.api_key: + raise SynapseConfigError("не задан api_key: передайте api_key= или SYNAPSE_API_KEY") + return self._config.api_key + + def _build( + self, + source: str | None, + subject: str, + action: str, + *, + priority: str, + payload: dict[str, Any] | None, + dedup_key: str | None, + ttl_seconds: int | None, + scheduled_at: Any, + ) -> dict[str, Any]: + source = self._resolve_source(source) + return envelope.build_envelope( + source, subject, action, + priority=priority, payload=payload, dedup_key=dedup_key, + ttl_seconds=ttl_seconds, scheduled_at=scheduled_at, + ) + + async def _post(self, path: str, body: Any) -> httpx.Response: + self._require_key() + try: + return await self._http.post(self._config.url(path), json=body, + headers=self._auth_headers()) + except httpx.TransportError as ex: + raise SynapseConnectionError(f"Synapse недоступен: {ex}") from ex + + async def _request(self, method: str, path: str) -> httpx.Response: + try: + return await self._http.request(method, self._config.url(path), + headers=self._auth_headers()) + except httpx.TransportError as ex: + raise SynapseConnectionError(f"Synapse недоступен: {ex}") from ex + + def _auth_headers(self) -> dict[str, str]: + headers: dict[str, str] = {"user-agent": USER_AGENT} + if self._config.api_key: + headers["authorization"] = f"Bearer {self._config.api_key}" + return headers diff --git a/src/gnexus_synapse/client.py b/src/gnexus_synapse/client.py new file mode 100644 index 0000000..c7a74dc --- /dev/null +++ b/src/gnexus_synapse/client.py @@ -0,0 +1,338 @@ +"""Синхронный клиент Ingestion API v1. + +Тонкий: один POST на send, никаких ретраев и очередей (у Synapse очередь +своя — Celery). Транспорт инжектный (httpx.Client) — свой создаётся с +явным таймаутом; инжектный не трогаем (таймаут настраивает потребитель). +""" + +from __future__ import annotations + +import logging +import uuid +from collections.abc import Mapping, Sequence +from contextlib import suppress +from datetime import datetime +from typing import Any, cast + +import httpx + +from . import envelope +from .config import SynapseConfig +from .dto import BatchRejection, BatchResult, Delivery, Event, EventStatus +from .exceptions import ( + SynapseAuthError, + SynapseConfigError, + SynapseConnectionError, + SynapseError, + SynapseForbiddenError, + SynapseNotFoundError, + SynapseServerError, + SynapseValidationError, +) + +__version__ = "0.1.0" + +API_PREFIX = "/api/v1" +USER_AGENT = f"gnexus-synapse-py/{__version__}" + +_STATUS_ERRORS: dict[int, type[SynapseError]] = { + 401: SynapseAuthError, + 403: SynapseForbiddenError, + 404: SynapseNotFoundError, + 422: SynapseValidationError, +} + +_DEFAULT_LOGGER = "gnexus.synapse" + + +class SynapseClient: + """Клиент одного источника: ключ выдаётся на source, источник в конверте.""" + + def __init__( + self, + url: str | None = None, + api_key: str | None = None, + *, + timeout: float | None = None, + http_client: httpx.Client | None = None, + default_source: str | None = None, + logger: logging.Logger | None = None, + ) -> None: + base = SynapseConfig.from_env() + self._config = SynapseConfig( + base_url=url if url is not None else base.base_url, + api_key=api_key if api_key is not None else base.api_key, + timeout=timeout if timeout is not None else base.timeout, + default_source=default_source if default_source is not None else base.default_source, + ) + self._logger = logger or logging.getLogger(_DEFAULT_LOGGER) + self._own_client = http_client is None + # свой клиент — с явным таймаутом (дефолт httpx для Client — None, «вечно») + self._http = http_client or httpx.Client( + timeout=self._config.timeout, + headers={"user-agent": USER_AGENT}, + ) + + # -- события --------------------------------------------------------- + + def send( + self, + source: str | None = None, + subject: str = "", + action: str = "", + *, + priority: str = "normal", + payload: dict[str, Any] | None = None, + dedup_key: str | None = None, + ttl_seconds: int | None = None, + scheduled_at: Any = None, + ) -> Event: + envelope_ = self._build(source, subject, action, priority=priority, payload=payload, + dedup_key=dedup_key, ttl_seconds=ttl_seconds, + scheduled_at=scheduled_at) + if envelope_.get("payload"): + envelope.warn_bad_user_id(envelope_["payload"], self._logger) + return _accepted_event(self._post(f"{API_PREFIX}/events", envelope_)) + + def emit( + self, + source: str | None = None, + subject: str = "", + action: str = "", + *, + priority: str = "normal", + payload: dict[str, Any] | None = None, + dedup_key: str | None = None, + ttl_seconds: int | None = None, + scheduled_at: Any = None, + ) -> Event | None: + """Fire-and-forget: ошибка Synapse не должна ломать сервис-источник. + + Глотает SynapseError (в т.ч. connection/5xx) ПЛЮС локальные ошибки + конверта — emit вызывают в хвостах обработки, где исключение нечать + некуда. Пишет warning в логгер; настройте уровень logging, чтобы + потери были видны. + """ + try: + return self.send(source, subject, action, priority=priority, payload=payload, + dedup_key=dedup_key, ttl_seconds=ttl_seconds, + scheduled_at=scheduled_at) + except SynapseError as ex: + self._logger.warning("synapse emit failed: %s", ex) + return None + + def send_batch(self, events: Sequence[Mapping[str, Any]]) -> BatchResult: + """Голый массив конвертов; сервер отвечает per-index (202 или 422). + + Локально битые конверты не посылаются вовсе — сразу BatchRejection + на своей позиции; индексы сервера маппятся обратно. «Ни одного не + принято» — не исключение, а результата с rejected (транспортный + сбой и 5xx — исключение: частичный успех неизвестен). + """ + positions: list[int] = [] + payload: list[dict[str, Any]] = [] + results: list[Event | BatchRejection | None] = [None] * len(events) + for i, raw in enumerate(events): + try: + source = raw.get("source") + payload.append(envelope.build_envelope( + source if source is not None else self._resolve_source(None), + str(raw.get("subject", "")), + str(raw.get("action", "")), + priority=raw.get("priority", "normal"), + payload=raw.get("payload"), + dedup_key=raw.get("dedup_key"), + ttl_seconds=raw.get("ttl_seconds"), + scheduled_at=raw.get("scheduled_at"), + )) + positions.append(i) + except (SynapseValidationError, SynapseConfigError) as ex: + results[i] = BatchRejection(index=i, detail=str(ex)) + + if positions: + resp = self._post(f"{API_PREFIX}/events/batch", payload) + body = _json_body(resp) + server_results = body.get("results", []) if isinstance(body, dict) else [] + for pos, item in zip(positions, server_results, strict=False): + if not isinstance(item, dict): + results[pos] = BatchRejection(index=pos, detail=f"битый элемент ответа: {item!r}") + elif item.get("ok") is True: + results[pos] = Event( + id=str(item.get("id", "")), + status=str(item.get("status", "")), + deduplicated=bool(item.get("deduplicated", False)), + ) + else: + # index у сервера — позиция в отправленном массиве; ремапим в оригинальную + sent_index = item.get("index", pos) + orig = sent_index if isinstance(sent_index, int) and 0 <= sent_index < len(positions) else pos + results[pos] = BatchRejection( + index=positions[orig], detail=str(item.get("detail", "")) + ) + return BatchResult(results=[r for r in results if r is not None]) + + def status(self, event_id: str) -> EventStatus: + try: + uuid.UUID(event_id) + except (ValueError, AttributeError) as ex: + raise SynapseValidationError(f"event_id={event_id!r} не uuid") from ex + resp = self._request("GET", f"{API_PREFIX}/events/{event_id}") + body = _json_body(resp) + if resp.status_code != 200: + raise _error_for(resp, body) + deliveries = [ + Delivery( + channel=str(d.get("channel", "")), + target=d.get("target"), + status=str(d.get("status", "")), + attempts=int(d.get("attempts", 0)), + error=d.get("error"), + rendered_message=d.get("rendered_message"), + ) + for d in body.get("deliveries", []) + ] + return EventStatus( + id=str(body.get("id", event_id)), + source=str(body.get("source", "")), + subject=str(body.get("subject", "")), + action=str(body.get("action", "")), + priority=str(body.get("priority", "")), + status=str(body.get("status", "")), + created_at=_dt(body["created_at"]), + expires_at=_dt(body.get("expires_at")) if body.get("expires_at") else None, + deliveries=deliveries, + ) + + # -- диагностика (без ключа) ---------------------------------------- + + def health(self) -> dict[str, Any]: + resp = self._request("GET", "/api/healthz") + body = _json_body(resp) + if resp.status_code != 200: + raise _error_for(resp, body) + return cast("dict[str, Any]", body) + + def ready(self) -> dict[str, Any]: + resp = self._request("GET", "/api/readyz") + body = _json_body(resp) + if resp.status_code != 200: + raise _error_for(resp, body) + return cast("dict[str, Any]", body) + + # -- транспорт -------------------------------------------------------- + + def close(self) -> None: + if self._own_client: + self._http.close() + + def __del__(self) -> None: + with suppress(Exception): + self.close() + + # -- внутреннее -------------------------------------------------------- + + def _resolve_source(self, source: str | None) -> str: + source = source or self._config.default_source + if not source: + raise SynapseConfigError( + "не задан source: передайте source= или default_source " + "(env SYNAPSE_DEFAULT_SOURCE)" + ) + return source + + def _require_key(self) -> str: + if not self._config.api_key: + raise SynapseConfigError( + "не задан api_key: передайте api_key= или SYNAPSE_API_KEY" + ) + return self._config.api_key + + def _build( + self, + source: str | None, + subject: str, + action: str, + *, + priority: str, + payload: dict[str, Any] | None, + dedup_key: str | None, + ttl_seconds: int | None, + scheduled_at: Any, + ) -> dict[str, Any]: + source = self._resolve_source(source) + return envelope.build_envelope( + source, subject, action, + priority=priority, payload=payload, dedup_key=dedup_key, + ttl_seconds=ttl_seconds, scheduled_at=scheduled_at, + ) + + def _post(self, path: str, body: Any) -> httpx.Response: + self._require_key() + try: + return self._http.post( + self._config.url(path), + json=body, + headers=self._auth_headers(), + ) + except httpx.TransportError as ex: + raise SynapseConnectionError(f"Synapse недоступен: {ex}") from ex + + def _request(self, method: str, path: str, *, json_body: Any = None) -> httpx.Response: + try: + return self._http.request( + method, + self._config.url(path), + json=json_body, + headers=self._auth_headers(), + ) + except httpx.TransportError as ex: + raise SynapseConnectionError(f"Synapse недоступен: {ex}") from ex + + def _auth_headers(self) -> dict[str, str]: + headers: dict[str, str] = {"user-agent": USER_AGENT} + if self._config.api_key: + headers["authorization"] = f"Bearer {self._config.api_key}" + return headers + + def _accepted(self, resp: httpx.Response) -> Event: + return _accepted_event(resp) + + +# -- разбор ответов (общие для sync/async) ------------------------------- + +def _json_body(resp: httpx.Response) -> Any: + try: + return resp.json() + except ValueError: + return {} + + +def _accepted_event(resp: httpx.Response) -> Event: + body = _json_body(resp) + if resp.status_code != 202: + raise _error_for(resp, body) + return Event( + id=str(body.get("id", "")), + status=str(body.get("status", "")), + deduplicated=bool(body.get("deduplicated", False)), + ) + + +def _detail(body: Any, resp: httpx.Response) -> str: + if isinstance(body, dict) and isinstance(body.get("detail"), (str, dict, list)): + return str(body["detail"]) + return f"HTTP {resp.status_code}: {resp.text[:300]}" + + +def _error_for(resp: httpx.Response, body: Any) -> SynapseError: + detail = _detail(body, resp) + if resp.status_code >= 500: + return SynapseServerError(detail, status_code=resp.status_code) + cls = _STATUS_ERRORS.get(resp.status_code) + if cls is None: + return SynapseError(detail, status_code=resp.status_code) + return cls(detail, status_code=resp.status_code) + + +def _dt(raw: str) -> datetime: + return datetime.fromisoformat(str(raw)) diff --git a/src/gnexus_synapse/config.py b/src/gnexus_synapse/config.py new file mode 100644 index 0000000..ca21504 --- /dev/null +++ b/src/gnexus_synapse/config.py @@ -0,0 +1,57 @@ +"""Конфигурация клиента. + +Прямые аргументы конструктора сильнее env: клиент в коде сервиса обычно +создают с url/api_key, а env-фолбэк — путь «развернул по умолчанию». +Имена переменных фиксированы для всей экосистемы: SYNAPSE_URL, +SYNAPSE_API_KEY, SYNAPSE_TIMEOUT, SYNAPSE_DEFAULT_SOURCE. +""" + +from __future__ import annotations + +import os +from collections.abc import Mapping +from dataclasses import dataclass, field + +ENV_URL = "SYNAPSE_URL" +ENV_API_KEY = "SYNAPSE_API_KEY" +ENV_TIMEOUT = "SYNAPSE_TIMEOUT" +ENV_DEFAULT_SOURCE = "SYNAPSE_DEFAULT_SOURCE" + +DEFAULT_TIMEOUT = 10.0 + + +@dataclass(frozen=True, slots=True) +class SynapseConfig: + base_url: str # без хвостового "/" + api_key: str | None = None + timeout: float = DEFAULT_TIMEOUT + default_source: str | None = None + env: Mapping[str, str] = field(default_factory=lambda: dict(os.environ), repr=False) + + @classmethod + def from_env(cls, env: Mapping[str, str] | None = None) -> SynapseConfig: + env = dict(os.environ) if env is None else dict(env) + return cls( + base_url=env.get(ENV_URL, ""), + api_key=env.get(ENV_API_KEY) or None, + timeout=_timeout(env), + default_source=env.get(ENV_DEFAULT_SOURCE) or None, + env=env, + ) + + def url(self, path: str) -> str: + """base_url + path; путь в path — с начальным '/'.""" + return f"{self.base_url.rstrip('/')}{path}" + + +def _timeout(env: Mapping[str, str]) -> float: + raw = (env.get(ENV_TIMEOUT) or "").strip() + if not raw: + return DEFAULT_TIMEOUT + try: + value = float(raw) + except ValueError: + raise ValueError(f"{ENV_TIMEOUT}={raw!r} не числится как число секунд") from None + if value <= 0: + raise ValueError(f"{ENV_TIMEOUT}={raw!r} должен быть > 0") + return value diff --git a/src/gnexus_synapse/dto.py b/src/gnexus_synapse/dto.py new file mode 100644 index 0000000..d89f3e0 --- /dev/null +++ b/src/gnexus_synapse/dto.py @@ -0,0 +1,67 @@ +"""DTO ответов Synapse — только то, что реально возвращает API v1 (docs/05).""" + +from __future__ import annotations + +from dataclasses import dataclass, field +from datetime import datetime + + +@dataclass(frozen=True, slots=True) +class Event: + """Ответ POST /api/v1/events (202): событие принято в очередь.""" + + id: str + status: str + deduplicated: bool = False + + +@dataclass(frozen=True, slots=True) +class BatchRejection: + """Один элемент batch-ответа с ok=false.""" + + index: int + detail: str + + +@dataclass(frozen=True, slots=True) +class BatchResult: + """Позиция в списке совпадает с порядком исходных events.""" + + results: list[Event | BatchRejection] = field(default_factory=list) + + @property + def accepted(self) -> list[Event]: + return [r for r in self.results if isinstance(r, Event)] + + @property + def rejected(self) -> list[BatchRejection]: + return [r for r in self.results if isinstance(r, BatchRejection)] + + @property + def all_accepted(self) -> bool: + return not self.rejected + + +@dataclass(frozen=True, slots=True) +class Delivery: + channel: str + target: str | None + status: str + attempts: int + error: str | None + rendered_message: str | None + + +@dataclass(frozen=True, slots=True) +class EventStatus: + """Ответ GET /api/v1/events/{id}.""" + + id: str + source: str + subject: str + action: str + priority: str + status: str + created_at: datetime + expires_at: datetime | None + deliveries: list[Delivery] = field(default_factory=list) diff --git a/src/gnexus_synapse/envelope.py b/src/gnexus_synapse/envelope.py new file mode 100644 index 0000000..f9b8826 --- /dev/null +++ b/src/gnexus_synapse/envelope.py @@ -0,0 +1,120 @@ +"""Сборка и локальная валидация конверта события — зеркало app/api/schemas.py сервера. + +Клиент валидирует конверт до HTTP: ошибка кода видна сразу, без round-trip'а, +и приходит как SynapseValidationError со статусом None. Правила держим один +в один с сервером — контракт v1: docs/05-ingestion-api.md. +""" + +from __future__ import annotations + +import logging +import re +from datetime import UTC, datetime +from typing import Any, Final + +from .exceptions import SynapseValidationError + +NAME_PATTERN: Final = r"^[a-z0-9]([a-z0-9._-]*[a-z0-9])?$" +NAME_MAX: Final = 64 +PRIORITY_VALUES: Final = ("low", "normal", "high", "critical") +TTL_MAX: Final = 604800 # 7 суток, ge=1 +DEDUP_MAX: Final = 255 + +_NAME_RE = re.compile(NAME_PATTERN) + + +def validate_name(kind: str, value: str) -> None: + """Один символ имени конверта (source/subject/action).""" + if not isinstance(value, str) or len(value) == 0 or len(value) > NAME_MAX: + raise ValueError(f"недопустимое {kind}={value!r}: строка 1..{NAME_MAX} символов") + if _NAME_RE.fullmatch(value) is None: + raise ValueError( + f"недопустимое {kind}={value!r}: не совпадает с шаблоном {NAME_PATTERN}" + ) + + +def validate_names(source: str, subject: str, action: str) -> None: + validate_name("source", source) + validate_name("subject", subject) + validate_name("action", action) + + +def warn_bad_user_id(payload: dict[str, Any], logger: logging.Logger) -> None: + """Конвенция payload.user_id (= sub gnexus-auth) — «о ком событие». + + Сервер payload не валидирует: не-строчный/пустой user_id не ломает приём, + но правила без user_id-целей доставят событие «только админу», а целевые + каналы запишут skipped. Поэтому warning, не ошибка. + """ + uid = payload.get("user_id") + if isinstance(uid, str) and uid: + return + if isinstance(uid, int) and not isinstance(uid, bool): # сервер сам приводит к строке + return + logger.warning( + "payload.user_id должен быть непустой строкой (или числом); " + "получено %r — адресные доставки будут skipped", + uid, + ) + + +def build_envelope( + source: str, + subject: str, + action: str, + *, + priority: str = "normal", + payload: dict[str, Any] | None = None, + dedup_key: str | None = None, + ttl_seconds: int | None = None, + scheduled_at: datetime | None = None, +) -> dict[str, Any]: + """Собрать конверт v1; лишнего в нём нет (сервер extra='forbid').""" + errors: list[str] = [] + try: + validate_names(source, subject, action) + except ValueError as ex: + errors.append(str(ex)) + if priority not in PRIORITY_VALUES: + errors.append( + f"недопустимый priority={priority!r}: разрешены {', '.join(PRIORITY_VALUES)}" + ) + if dedup_key is not None and not isinstance(dedup_key, str): + errors.append("dedup_key должен быть строкой или None") + elif dedup_key == "": # пустая строка = не задаём + dedup_key = None + elif dedup_key is not None and len(dedup_key) > DEDUP_MAX: + errors.append(f"dedup_key длиннее {DEDUP_MAX} символов") + if ttl_seconds is not None: + if isinstance(ttl_seconds, bool) or not isinstance(ttl_seconds, int): + errors.append("ttl_seconds должен быть целым числом или None") + elif not 1 <= ttl_seconds <= TTL_MAX: + errors.append(f"ttl_seconds={ttl_seconds}: допустимо 1..{TTL_MAX}") + if errors: + raise SynapseValidationError("; ".join(errors)) + + envelope: dict[str, Any] = { + "source": source, + "subject": subject, + "action": action, + } + if priority != "normal": + envelope["priority"] = priority + if payload is not None: + if not isinstance(payload, dict): + raise SynapseValidationError("payload должен быть словарем (JSON-объектом)") + if payload: + envelope["payload"] = payload + if dedup_key is not None: + envelope["dedup_key"] = dedup_key + if ttl_seconds is not None: + envelope["ttl_seconds"] = ttl_seconds + if scheduled_at is not None: + if not isinstance(scheduled_at, datetime): + raise SynapseValidationError("scheduled_at должен быть datetime или None") + # сервер ждёт ISO-8601: наивное время считаем UTC (сразу предупреждаем — + # локальное время источника почти наверняка ошибка) + if scheduled_at.tzinfo is None: + scheduled_at = scheduled_at.replace(tzinfo=UTC) + envelope["scheduled_at"] = scheduled_at.isoformat() + return envelope diff --git a/src/gnexus_synapse/exceptions.py b/src/gnexus_synapse/exceptions.py new file mode 100644 index 0000000..39d8f72 --- /dev/null +++ b/src/gnexus_synapse/exceptions.py @@ -0,0 +1,48 @@ +"""Типизированные исключения клиента. + +Каждый класс = одна ситуация, в которой сервису нужно разное поведение: +401 — ключ потерян/отозван (пересоздать), 403 — source в конверте не совпал +с источником ключа (анти-спуфинг, ошибка кода), 422 — конверт не принят +(ошибка кода или тип не зарегистрирован), 5xx — Synapse лежит, connection — +недоступен. `detail` — человекочитаемое сообщение, `status_code` — HTTP-код, +если был ответ. +""" + +from __future__ import annotations + + +class SynapseError(RuntimeError): + """База: всё, что клиент бросает, ловится этим классом.""" + + def __init__(self, detail: str, *, status_code: int | None = None) -> None: + super().__init__(detail) + self.detail = detail + self.status_code = status_code + + +class SynapseConfigError(SynapseError): + """Нет url/api_key в аргументах и env, или конфиг в принципе битый.""" + + +class SynapseConnectionError(SynapseError): + """HTTP-транспорт недоступен (DNS, refusal, timeout). __cause__ — оригинал.""" + + +class SynapseAuthError(SynapseError): + """401: ключ отсутствует, неизвестен или отозван (источник архива — тоже сюда).""" + + +class SynapseForbiddenError(SynapseError): + """403: поле source в конверте не совпадает с источником ключа.""" + + +class SynapseNotFoundError(SynapseError): + """404: событие с таким id у источника ключа не найдено.""" + + +class SynapseValidationError(SynapseError): + """422 от сервера ИЛИ локальная ошибка конверта (status_code=None).""" + + +class SynapseServerError(SynapseError): + """5xx: Synapse недоступен/ошибка на его стороне — сообщение не ушло.""" diff --git a/tests/unit/test_client.py b/tests/unit/test_client.py new file mode 100644 index 0000000..b449e65 --- /dev/null +++ b/tests/unit/test_client.py @@ -0,0 +1,239 @@ +"""Юнит-тесты SynapseClient — httpx.MockTransport вместо сервера.""" + +import logging +from datetime import UTC, datetime + +import httpx +import pytest + +from gnexus_synapse import SynapseClient +from gnexus_synapse.exceptions import ( + SynapseAuthError, + SynapseConfigError, + SynapseConnectionError, + SynapseError, + SynapseForbiddenError, + SynapseNotFoundError, + SynapseServerError, + SynapseValidationError, +) + +BASE = "http://synapse.test" + + +def make_client(handler, **kw): + return SynapseClient("http://synapse.test", "syn_test", + http_client=httpx.Client(transport=httpx.MockTransport(handler)), + **kw) + + +def test_send_success_builds_request(): + seen = {} + + def handler(request: httpx.Request) -> httpx.Response: + seen["url"] = str(request.url) + seen["auth"] = request.headers.get("authorization") + seen["ua"] = request.headers.get("user-agent") + seen["body"] = json_loads(request.read()) + return httpx.Response(202, json={"ok": True, "id": "e-1", "status": "queued", + "deduplicated": False}) + + event = make_client(handler).send("bugtrail", "test", "failed", + priority="high", payload={"user_id": "42"}, + dedup_key="d1") + assert (event.id, event.status, event.deduplicated) == ("e-1", "queued", False) + assert seen["url"] == f"{BASE}/api/v1/events" + assert seen["auth"] == "Bearer syn_test" + assert seen["ua"].startswith("gnexus-synapse-py/") + assert seen["body"]["source"] == "bugtrail" + assert seen["body"]["priority"] == "high" + assert seen["body"]["dedup_key"] == "d1" + + +def test_send_maps_status_errors(): + cases = { + 401: (SynapseAuthError, "API-ключ неизвестен или отозван"), + 403: (SynapseForbiddenError, "не совпадает с источником ключа"), + 422: (SynapseValidationError, "Тип не зарегистрирован"), + 500: (SynapseServerError, "boom"), + 503: (SynapseServerError, "перегружен"), + 418: (SynapseError, "teapot"), + } + + def make_handler(code, detail): + return lambda request: httpx.Response(code, json={"detail": detail}) + + for code, (cls, detail) in cases.items(): + client = make_client(make_handler(code, detail)) + with pytest.raises(cls) as exc_info: + client.send("s", "x", "y") + assert exc_info.value.detail == detail + assert exc_info.value.status_code == code + + +def test_send_transport_error_wrapped(): + def handler(request): + raise httpx.ConnectError("connection refused") + + client = make_client(handler) + with pytest.raises(SynapseConnectionError) as exc_info: + client.send("s", "x", "y") + assert isinstance(exc_info.value.__cause__, httpx.ConnectError) + + +def test_local_validation_no_http(): + called = {"n": 0} + + def handler(request): + called["n"] += 1 + return httpx.Response(202, json={"ok": True, "id": "e", "status": "queued"}) + + client = make_client(handler) + with pytest.raises(SynapseValidationError) as exc_info: + client.send("Wrong Source", "x", "y") + assert exc_info.value.status_code is None # локальная ошибка, не сервер + assert called["n"] == 0 + + +def test_default_source_and_missing(): + def handler(request): + assert json_loads(request.read())["source"] == "libtest" + return httpx.Response(202, json={"ok": True, "id": "e", "status": "queued"}) + + make_client(handler, default_source="libtest").send(subject="x", action="y") + with pytest.raises(SynapseConfigError, match="source"): + make_client(handler).send(subject="x", action="y") + with pytest.raises(SynapseConfigError, match="api_key"): + SynapseClient("http://s.test", api_key=None, + http_client=httpx.Client(transport=httpx.MockTransport(handler)), + default_source="s").send("s", "x", "y") + + +def test_emit_swallows_and_logs(caplog): + caplog.set_level(logging.WARNING, logger="gnexus.synapse") + handler = lambda request: httpx.Response(503, json={"detail": "нет воркеров"}) # noqa: E731 + client = make_client(handler) + result = client.emit("s", "x", "y") + assert result is None + assert any("synapse emit failed" in r.message for r in caplog.records) + + +def test_batch_partial(): + def handler(request): + body = json_loads(request.read()) + assert len(body) == 2 # конверт без source локально отвергнут + # 2-й из посланных сервер забраковал (тип не зарегистрирован) + return httpx.Response(202, json={"results": [ + {"ok": True, "id": "a", "status": "queued", "deduplicated": False}, + {"ok": False, "index": 1, "detail": "Тип не зарегистрирован"}, + ]}) + + result = make_client(handler).send_batch([ + {"source": "s", "subject": "x", "action": "y1"}, # сервер принял + {"subject": "x", "action": "y"}, # нет source → локальный reject + {"source": "s", "subject": "x", "action": "y2"}, # сервер отверг + ]) + assert result.all_accepted is False + assert [e.id for e in result.accepted] == ["a"] + assert [(r.index, r.detail) for r in result.rejected] == [ + (1, "не задан source: передайте source= или default_source (env SYNAPSE_DEFAULT_SOURCE)"), + (2, "Тип не зарегистрирован"), + ] + + +def test_batch_local_rejections_shift_index(): + def handler(request): + body = json_loads(request.read()) + assert len(body) == 1 # битый не посылается + return httpx.Response(202, json={"results": [ + {"ok": True, "id": "a", "status": "queued", "deduplicated": False}, + ]}) + + result = make_client(handler).send_batch([ + {"source": "s", "subject": "x", "action": "y"}, # ok + {"source": "BAD SOURCE", "subject": "x", "action": "y"}, # локальный reject + ]) + assert [e.id for e in result.accepted] == ["a"] # позиция 0 — оригинальная + r = result.rejected[0] + assert (r.index, "BAD SOURCE") == (r.index, "BAD SOURCE") and r.index == 1 + + +def test_batch_empty_no_http(): + def handler(request): + raise AssertionError("HTTP не нужен") + + result = make_client(handler).send_batch([]) + assert result.results == [] and result.all_accepted + + +def test_status_parses_and_404(): + eid = "11111111-1111-1111-1111-111111111111" + other = "22222222-2222-2222-2222-222222222222" + status_body = { + "id": eid, "source": "bugtrail", "subject": "test", "action": "failed", + "priority": "high", "status": "delivered", + "created_at": datetime(2026, 10, 3, 12, 0, tzinfo=UTC).isoformat(), + "expires_at": None, + "deliveries": [ + {"channel": "internal_log", "target": None, "status": "delivered", + "attempts": 1, "error": None, "rendered_message": "тест упал"}, + {"channel": "telegram", "target": "qa-chat", "status": "retrying", + "attempts": 2, "error": "429", "rendered_message": None}, + ], + } + + def handler(request): + if request.url.path.endswith(other): + return httpx.Response(404, json={"detail": "Событие не найдено"}) + return httpx.Response(200, json=status_body) + + client = make_client(handler) + st = client.status(eid) + assert st.status == "delivered" + assert st.created_at.tzinfo is not None + assert st.deliveries[1].status == "retrying" and st.deliveries[1].attempts == 2 + with pytest.raises(SynapseNotFoundError): + client.status(other) + with pytest.raises(SynapseValidationError): + client.status("не-uuid") + + +def test_health_ready(): + def handler(request): + if request.url.path == "/api/healthz": + return httpx.Response(200, json={"status": "ok", "service": "synapse"}) + return httpx.Response(200, json={"ready": True, "components": {}}) + + client = make_client(handler) + assert client.health() == {"status": "ok", "service": "synapse"} + assert client.ready() == {"ready": True, "components": {}} + + +def test_env_fallback(monkeypatch): + monkeypatch.setenv("SYNAPSE_URL", "http://env-synapse.test") + monkeypatch.setenv("SYNAPSE_API_KEY", "syn_from_env") + monkeypatch.setenv("SYNAPSE_DEFAULT_SOURCE", "libtest") + seen = {} + + def capture(request): + seen["url"] = str(request.url) + seen["auth"] = request.headers.get("authorization") + return httpx.Response(202, json={"ok": True, "id": "e", "status": "queued"}) + + client2 = SynapseClient(http_client=httpx.Client(transport=httpx.MockTransport(capture))) + client2.send(subject="x", action="y") + assert seen["url"] == "http://env-synapse.test/api/v1/events" + assert seen["auth"] == "Bearer syn_from_env" + + monkeypatch.delenv("SYNAPSE_API_KEY") + with pytest.raises(SynapseConfigError, match="api_key"): + SynapseClient(http_client=httpx.Client(transport=httpx.MockTransport(capture)), + default_source="s").send("s", "x", "y") + + +# -- helpers ------------------------------------------------------------- + +def json_loads(raw: bytes) -> dict: + import json + + return json.loads(raw) diff --git a/tests/unit/test_config.py b/tests/unit/test_config.py new file mode 100644 index 0000000..7f9ec54 --- /dev/null +++ b/tests/unit/test_config.py @@ -0,0 +1,35 @@ +"""SynapseConfig.from_env — env-переменные экосистемы.""" + +import pytest + +from gnexus_synapse.config import ENV_DEFAULT_SOURCE, SynapseConfig + + +def test_from_env_defaults(): + cfg = SynapseConfig.from_env({}) + assert cfg.base_url == "" + assert cfg.api_key is None + assert cfg.timeout == 10.0 + assert cfg.default_source is None + + +def test_from_env_reads_all(): + cfg = SynapseConfig.from_env({ + "SYNAPSE_URL": "http://synapse:8013/", + "SYNAPSE_API_KEY": "syn_abc", + "SYNAPSE_TIMEOUT": "3.5", + ENV_DEFAULT_SOURCE: "bugtrail", + }) + assert cfg.base_url == "http://synapse:8013/" + assert cfg.api_key == "syn_abc" + assert cfg.timeout == 3.5 + assert cfg.default_source == "bugtrail" + assert cfg.url("/api/v1/events") == "http://synapse:8013/api/v1/events" + + +def test_empty_timeout_falls_to_default(): + assert SynapseConfig.from_env({"SYNAPSE_TIMEOUT": ""}).timeout == 10.0 + with pytest.raises(ValueError, match="SYNAPSE_TIMEOUT"): + SynapseConfig.from_env({"SYNAPSE_TIMEOUT": "всегда"}) + with pytest.raises(ValueError, match="SYNAPSE_TIMEOUT"): + SynapseConfig.from_env({"SYNAPSE_TIMEOUT": "0"}) diff --git a/tests/unit/test_envelope.py b/tests/unit/test_envelope.py new file mode 100644 index 0000000..53c0aed --- /dev/null +++ b/tests/unit/test_envelope.py @@ -0,0 +1,82 @@ +"""Юнит-тесты envelope.py — контракт зеркалит app/api/schemas.py сервера.""" + +import logging +from datetime import UTC, datetime + +import pytest + +from gnexus_synapse.envelope import build_envelope, warn_bad_user_id +from gnexus_synapse.exceptions import SynapseValidationError + + +def test_minimal_envelope_exact_keys(): + env = build_envelope("bugtrail", "test", "failed") + assert env == {"source": "bugtrail", "subject": "test", "action": "failed"} + + +def test_none_fields_never_sent(): + # сервер extra="forbid" — None-поле это 422; клиент не должен их слать + env = build_envelope("s", "x", "y", payload=None, dedup_key=None, ttl_seconds=None, + scheduled_at=None) + assert set(env) == {"source", "subject", "action"} + env2 = build_envelope("s", "x", "y", payload={}) + assert "payload" not in env2 + + +def test_full_envelope_values(): + ts = datetime(2026, 10, 3, 12, 0, tzinfo=UTC) + env = build_envelope("s", "x", "y", priority="high", payload={"user_id": "42"}, + dedup_key="d-1", ttl_seconds=60, scheduled_at=ts) + assert env == {"source": "s", "subject": "x", "action": "y", "priority": "high", + "payload": {"user_id": "42"}, "dedup_key": "d-1", "ttl_seconds": 60, + "scheduled_at": "2026-10-03T12:00:00+00:00"} + # default priority не шлётся — равен серверному дефолту + env3 = build_envelope("s", "x", "y", priority="normal") + assert "priority" not in env3 + + +def test_pattern_names(): + # валидные + for name in ("a", "bugtrail", "test.failed", "task-created", "v2", "a1_b.c-d"): + build_envelope(name, name, name) + # невалидные: uppercase, пробел, ведущий/хвостовой разделитель, пусто, длина 65 + for bad in ("Bugtrail", "a b", "-lead", "trail-", "_lead", "trail_", "", "a" * 65, "спецсимволы"): + with pytest.raises(SynapseValidationError): + build_envelope(bad, "x", "y") + with pytest.raises(SynapseValidationError): + build_envelope("s", bad, "y") + with pytest.raises(SynapseValidationError): + build_envelope("s", "x", bad) + + +def test_bad_priority_ttl_dedup(): + with pytest.raises(SynapseValidationError, match="priority"): + build_envelope("s", "x", "y", priority="urgent") + for bad_ttl in (0, -5, 604801, 1.5, "3600", True): + with pytest.raises(SynapseValidationError, match="ttl_seconds"): + build_envelope("s", "x", "y", ttl_seconds=bad_ttl) + assert build_envelope("s", "x", "y", ttl_seconds=604800)["ttl_seconds"] == 604800 + with pytest.raises(SynapseValidationError, match="dedup_key"): + build_envelope("s", "x", "y", dedup_key="d" * 256) + assert "dedup_key" not in build_envelope("s", "x", "y", dedup_key="") + with pytest.raises(SynapseValidationError, match="payload"): + build_envelope("s", "x", "y", payload=["не", "словарь"]) # type: ignore[arg-type] + + +def test_naive_scheduled_at_treated_as_utc(): + env = build_envelope("s", "x", "y", scheduled_at=datetime(2026, 10, 3, 12, 0)) + assert env["scheduled_at"] == "2026-10-03T12:00:00+00:00" + + +def test_user_id_warning(caplog): + log = logging.getLogger("test.syn") + ok = {"user_id": "42", "x": 1} + warn_bad_user_id(ok, log) + warn_bad_user_id({"user_id": ""}, log) + warn_bad_user_id({"user_id": None}, log) + warn_bad_user_id({"user_id": {}} , log) + warn_bad_user_id({"user_id": True}, log) + warn_bad_user_id({}, log) + # 5 warning'ов (пустая строка, None, dict, bool, отсутствие) + assert len([r for r in caplog.records if r.levelno == logging.WARNING]) == 5 + warn_bad_user_id({"user_id": 42}, log) # число ок