diff --git a/README.md b/README.md index 8dbf324..d6eec2c 100644 --- a/README.md +++ b/README.md @@ -10,7 +10,7 @@ ## Установка ```bash -pip install "git+https://git.gnexus.space/git/root/gn-synapse-client-py.git@v0.1.1" +pip install "git+https://git.gnexus.space/git/root/gn-synapse-client-py.git@v0.1.2" ``` Серверный репозиторий — `https://git.gnexus.space/git/root/gn-synapse.git`. @@ -25,7 +25,7 @@ syn = SynapseClient("http://synapse.gnexus.space:8013", "syn_...") event = syn.send("bugtrail", "test", "failed", - payload={"user_id": "", "error": "..."}, + payload={"error": "..."}, user_id="", priority="high", dedup_key="run-2026-10-03") print(event.id, event.status) # queued @@ -52,10 +52,11 @@ | Метод | Что делает | Ошибки | |---|---|---| -| `send(source, subject, action, *, priority, payload, dedup_key, ttl_seconds, scheduled_at)` | POST /api/v1/events → `Event(id, status, deduplicated)` | типизированные, см. ниже | +| `send(source, subject, action, *, priority, payload, user_id, 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` | +| `wait_for_status(event_id, *, statuses=("done", "failed"), timeout=30, interval=1)` | поллинг до целевого статуса (для приёмки/тестов); морг транспорта не прерывает ожидание | `SynapseStatusTimeoutError` | | `health()` / `ready()` | диагностика без ключа | — | | `close()` | закрыть свой httpx.Client (инжектный не трогает) | — | @@ -70,6 +71,7 @@ | `SynapseValidationError` | 422 сервера или **локальная** ошибка конверта (`status_code is None`) | | `SynapseNotFoundError` | 404 — событие не у источника ключа | | `SynapseServerError` | 5xx — событие не принято, повторите при желании | +| `SynapseStatusTimeoutError` | `wait_for_status` не дождался целевого статуса за timeout | | `SynapseWebhookError` | приём s2s: подпись/freshness/не-JSON — не прошло проверку `verify_webhook` | Локальная валидация зеркалит серверную (`^[a-z0-9]([a-z0-9._-]*[a-z0-9])?$`, @@ -105,7 +107,9 @@ 300 с)/не-JSON — всё ловится одним `SynapseWebhookError` (детали в `detail`). Схема подписи — одна на всю экосистему (совместима с gnexus-auth и с `app/signature.py` Synapse); проверка константным сравнением. -`make_signature()` — для тестов приёмника. Секрет per-target — +`make_signature()` — для тестов приёмника; `parse_event_type()` — тройка +из `X-Gnexus-Event-Type` (`("monitoring", "container", "down")`); +`ack(event_id)` — тело ответа 2xx. Секрет per-target — `S2S_SECRET_` (docs/05). ## Конвенции и оговорки diff --git a/examples/plain/smoke.py b/examples/plain/smoke.py index 2207228..8000fe0 100644 --- a/examples/plain/smoke.py +++ b/examples/plain/smoke.py @@ -26,18 +26,19 @@ # уникальный на прогон: дедуп-окно сервера 24 ч, повтор полный — смок увидит не queued DEDUP = f"client-smoke-{int(time.time())}" +# источник теста: из env (SYNAPSE_DEFAULT_SOURCE, см. заголовок) — как и в реальном сервисе +SOURCE = os.environ.get("SYNAPSE_DEFAULT_SOURCE") or "libtest" def main() -> int: client = SynapseClient() print(" health:", client.health()) - event = client.send("libtest", "ping", "done", payload={"user_id": "mcp"}, - dedup_key=DEDUP) + event = client.send(SOURCE, "ping", "done", 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"}, + repeat = client.send(SOURCE, "ping", "done", user_id="mcp", dedup_key=DEDUP) print(f" repeat: id={repeat.id} deduplicated={repeat.deduplicated}") assert repeat.deduplicated and repeat.id == event.id # дедуп вернул первое событие diff --git a/pyproject.toml b/pyproject.toml index 4d03c6d..e4772a8 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ [project] name = "gnexus-synapse" -version = "0.1.1" +version = "0.1.2" description = "Thin client library for gn-synapse (Gnexus notification hub) event ingestion" readme = "README.md" requires-python = ">=3.11" diff --git a/src/gnexus_synapse/__init__.py b/src/gnexus_synapse/__init__.py index 46ea74e..f15f3b5 100644 --- a/src/gnexus_synapse/__init__.py +++ b/src/gnexus_synapse/__init__.py @@ -15,12 +15,22 @@ SynapseForbiddenError, SynapseNotFoundError, SynapseServerError, + SynapseStatusTimeoutError, SynapseValidationError, SynapseWebhookError, ) -from .webhook import DEFAULT_MAX_SKEW, SIGNATURE_HEADER, make_signature, verify_webhook +from .webhook import ( + DEFAULT_MAX_SKEW, + EVENT_TYPE_HEADER, + SIGNATURE_HEADER, + SOURCE_HEADER, + ack, + make_signature, + parse_event_type, + verify_webhook, +) -__version__ = "0.1.1" +__version__ = "0.1.2" __all__ = [ "AsyncSynapseClient", @@ -39,10 +49,15 @@ "SynapseValidationError", "SynapseServerError", "SynapseWebhookError", + "SynapseStatusTimeoutError", "SynapseError", "verify_webhook", "make_signature", + "parse_event_type", + "ack", "SIGNATURE_HEADER", + "EVENT_TYPE_HEADER", + "SOURCE_HEADER", "DEFAULT_MAX_SKEW", "ENV_URL", "ENV_API_KEY", diff --git a/src/gnexus_synapse/_async.py b/src/gnexus_synapse/_async.py index fd0e1a4..87fd9f7 100644 --- a/src/gnexus_synapse/_async.py +++ b/src/gnexus_synapse/_async.py @@ -7,9 +7,10 @@ from __future__ import annotations +import asyncio import logging import uuid -from collections.abc import Mapping, Sequence +from collections.abc import Iterable, Mapping, Sequence from typing import Any import httpx @@ -22,9 +23,11 @@ SynapseConfigError, SynapseConnectionError, SynapseError, + SynapseStatusTimeoutError, SynapseValidationError, ) +# user_agent/version держит client — см. USER_AGENT там же _DEFAULT_LOGGER = "gnexus.synapse" @@ -64,10 +67,11 @@ dedup_key: str | None = None, ttl_seconds: int | None = None, scheduled_at: Any = None, + user_id: str | None = None, ) -> Event: envelope_ = self._build(source, subject, action, priority=priority, payload=payload, dedup_key=dedup_key, ttl_seconds=ttl_seconds, - scheduled_at=scheduled_at) + scheduled_at=scheduled_at, user_id=user_id) if envelope_.get("payload"): envelope.warn_bad_user_id(envelope_["payload"], self._logger) return _accepted_event(await self._post(f"{API_PREFIX}/events", envelope_)) @@ -83,11 +87,12 @@ dedup_key: str | None = None, ttl_seconds: int | None = None, scheduled_at: Any = None, + user_id: str | None = 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) + scheduled_at=scheduled_at, user_id=user_id) except SynapseError as ex: self._logger.warning("synapse emit failed: %s", ex) return None @@ -167,6 +172,35 @@ deliveries=deliveries, ) + async def wait_for_status( + self, + event_id: str, + *, + statuses: Iterable[str] = ("done", "failed"), + timeout: float = 30.0, + interval: float = 1.0, + ) -> EventStatus: + """Поллинг до целевого статуса или таймаут — для приёмки/тестов. + + Морг Synapse (ConnectionError) не прерывает ожидание до дедлайна. + Не достигли → SynapseStatusTimeoutError с последним известным. + """ + wanted = tuple(statuses) + deadline = asyncio.get_event_loop().time() + timeout + while True: + try: + current = await self.status(event_id) + except SynapseConnectionError: + current = None + if current is not None and current.status in wanted: + return current + if asyncio.get_event_loop().time() >= deadline: + seen = current.status if current else "нет ответа" + raise SynapseStatusTimeoutError( + f"событие {event_id} не в {wanted} за {timeout} c (последний статус: {seen})" + ) + await asyncio.sleep(interval) + async def aclose(self) -> None: if self._own_client: await self._http.aclose() @@ -195,8 +229,11 @@ dedup_key: str | None, ttl_seconds: int | None, scheduled_at: Any, + user_id: str | None = None, ) -> dict[str, Any]: source = self._resolve_source(source) + if user_id: + payload = {**(payload or {}), "user_id": user_id} return envelope.build_envelope( source, subject, action, priority=priority, payload=payload, dedup_key=dedup_key, diff --git a/src/gnexus_synapse/client.py b/src/gnexus_synapse/client.py index c7a74dc..4b6d241 100644 --- a/src/gnexus_synapse/client.py +++ b/src/gnexus_synapse/client.py @@ -8,8 +8,9 @@ from __future__ import annotations import logging +import time import uuid -from collections.abc import Mapping, Sequence +from collections.abc import Iterable, Mapping, Sequence from contextlib import suppress from datetime import datetime from typing import Any, cast @@ -27,10 +28,11 @@ SynapseForbiddenError, SynapseNotFoundError, SynapseServerError, + SynapseStatusTimeoutError, SynapseValidationError, ) -__version__ = "0.1.0" +__version__ = "0.1.2" API_PREFIX = "/api/v1" USER_AGENT = f"gnexus-synapse-py/{__version__}" @@ -86,10 +88,11 @@ dedup_key: str | None = None, ttl_seconds: int | None = None, scheduled_at: Any = None, + user_id: str | None = None, ) -> Event: envelope_ = self._build(source, subject, action, priority=priority, payload=payload, dedup_key=dedup_key, ttl_seconds=ttl_seconds, - scheduled_at=scheduled_at) + scheduled_at=scheduled_at, user_id=user_id) if envelope_.get("payload"): envelope.warn_bad_user_id(envelope_["payload"], self._logger) return _accepted_event(self._post(f"{API_PREFIX}/events", envelope_)) @@ -105,6 +108,7 @@ dedup_key: str | None = None, ttl_seconds: int | None = None, scheduled_at: Any = None, + user_id: str | None = None, ) -> Event | None: """Fire-and-forget: ошибка Synapse не должна ломать сервис-источник. @@ -116,7 +120,7 @@ try: return self.send(source, subject, action, priority=priority, payload=payload, dedup_key=dedup_key, ttl_seconds=ttl_seconds, - scheduled_at=scheduled_at) + scheduled_at=scheduled_at, user_id=user_id) except SynapseError as ex: self._logger.warning("synapse emit failed: %s", ex) return None @@ -203,6 +207,37 @@ deliveries=deliveries, ) + def wait_for_status( + self, + event_id: str, + *, + statuses: Iterable[str] = ("done", "failed"), + timeout: float = 30.0, + interval: float = 1.0, + ) -> EventStatus: + """Поллинг до целевого статуса или таймаут — для приёмки/тестов, + не бизнес-кода (события fire-and-forget, docs/09). + + statuses — терминальные статусы; приёмный статус может мигать + Synapse-морг (ConnectionError) — пауза игнорируется и ждём дедлайн. + Не достигли → SynapseStatusTimeoutError с последним известным. + """ + wanted = tuple(statuses) + deadline = time.monotonic() + timeout + while True: + try: + current = self.status(event_id) + except SynapseConnectionError: + current = None + if current is not None and current.status in wanted: + return current + if time.monotonic() >= deadline: + seen = current.status if current else "нет ответа" + raise SynapseStatusTimeoutError( + f"событие {event_id} не в {wanted} за {timeout} c (последний статус: {seen})" + ) + time.sleep(interval) + # -- диагностика (без ключа) ---------------------------------------- def health(self) -> dict[str, Any]: @@ -258,8 +293,13 @@ dedup_key: str | None, ttl_seconds: int | None, scheduled_at: Any, + user_id: str | None = None, ) -> dict[str, Any]: source = self._resolve_source(source) + # user_id — first-class конвенция «о ком событие» (docs/05): мержим в payload, + # явный параметр перебивает payload["user_id"], если тот был + if user_id: + payload = {**(payload or {}), "user_id": user_id} return envelope.build_envelope( source, subject, action, priority=priority, payload=payload, dedup_key=dedup_key, diff --git a/src/gnexus_synapse/exceptions.py b/src/gnexus_synapse/exceptions.py index 08eebd6..7c39de4 100644 --- a/src/gnexus_synapse/exceptions.py +++ b/src/gnexus_synapse/exceptions.py @@ -52,3 +52,8 @@ """s2s-доставка не прошла проверку: нет/битый заголовок подписи, чужой секрет, replay (тело старше окна свежести), не-JSON тело. Локальная ошибка приёма (status_code=None) — это не про отправку.""" + + +class SynapseStatusTimeoutError(SynapseError): + """wait_for_status: событие не достигло целевого статуса за timeout. + Для приёмочных тестов/диагностики, не бизнес-обещаний.""" diff --git a/src/gnexus_synapse/webhook.py b/src/gnexus_synapse/webhook.py index d9f3f0e..29e822c 100644 --- a/src/gnexus_synapse/webhook.py +++ b/src/gnexus_synapse/webhook.py @@ -29,12 +29,46 @@ from .exceptions import SynapseWebhookError -__all__ = ["SIGNATURE_HEADER", "DEFAULT_MAX_SKEW", "verify_webhook", "make_signature"] +__all__ = [ + "SIGNATURE_HEADER", + "EVENT_TYPE_HEADER", + "SOURCE_HEADER", + "DEFAULT_MAX_SKEW", + "verify_webhook", + "make_signature", + "parse_event_type", + "ack", +] SIGNATURE_HEADER = "x-gnexus-signature" +EVENT_TYPE_HEADER = "x-gnexus-event-type" # ".." +SOURCE_HEADER = "x-synapse-source" DEFAULT_MAX_SKEW = 300 # секунд (зеркало gnexus-auth) +def parse_event_type(header: str | None) -> tuple[str, str, str] | None: + """`".."` → кортеж; None — нет/битый заголовок. + + Контекст приёма: первый экран обработки — обычно маршрутизация + по тройке из `X-Gnexus-Event-Type`. Имена источника не валидируются + (это метрика, не конверт): валидирует отправитель на ingestion. + """ + if not header: + return None + parts = header.split(".") + if len(parts) != 3 or not all(parts): + return None + return parts[0], parts[1], parts[2] + + +def ack(event_id: str | None = None) -> dict[str, object]: + """Тело ответа 2xx приёмника: {\"received\": true, \"event_id\": ...}.""" + body: dict[str, object] = {"received": True} + if event_id: + body["event_id"] = event_id + return body + + def make_signature(body: bytes, secret: str, timestamp: int | None = None) -> str: """Подпись «сырое» тело (для тестов и отправщиков, приём — не его работа).""" if timestamp is None: diff --git a/tests/unit/test_client.py b/tests/unit/test_client.py index b449e65..6aa53b2 100644 --- a/tests/unit/test_client.py +++ b/tests/unit/test_client.py @@ -1,5 +1,6 @@ """Юнит-тесты SynapseClient — httpx.MockTransport вместо сервера.""" +import json as _json import logging from datetime import UTC, datetime @@ -237,3 +238,82 @@ import json return json.loads(raw) + + +# -- user_id first-class + wait_for_status (v0.1.2) ----------------- + + + + +def test_user_id_kwarg_merges_into_payload(): + seen = {} + + def handler(request: httpx.Request) -> httpx.Response: + seen["body"] = _json.loads(request.read()) + return httpx.Response(202, json={"ok": True, "id": "e-1", "status": "queued", + "deduplicated": False}) + + make_client(handler).send("bugtrail", "task", "created", + payload={"title": "x"}, user_id="u-42") + assert seen["body"]["payload"] == {"title": "x", "user_id": "u-42"} + + +def test_user_id_overrides_payload_value(): + seen = {} + + def handler(request: httpx.Request) -> httpx.Response: + seen["body"] = _json.loads(request.read()) + return httpx.Response(202, json={"ok": True, "id": "e-1", "status": "queued", + "deduplicated": False}) + + make_client(handler).send("bugtrail", "task", "created", + payload={"user_id": "wrong"}, user_id="u-42") + assert seen["body"]["payload"]["user_id"] == "u-42" + + +def test_user_id_none_leaves_payload_untouched(): + seen = {} + + def handler(request: httpx.Request) -> httpx.Response: + seen["body"] = _json.loads(request.read()) + return httpx.Response(202, json={"ok": True, "id": "e-1", "status": "queued", + "deduplicated": False}) + + make_client(handler).emit("bugtrail", "task", "created") + assert "payload" not in seen["body"] + + +def _status_response(status, eid="11111111-1111-1111-1111-111111111111"): + return httpx.Response(200, json={ + "id": eid, "source": "libtest", "subject": "ping", "action": "done", + "priority": "normal", "status": status, "created_at": "2026-10-03T12:00:00+00:00", + "deliveries": [], + }) + + +def test_wait_for_status_polls_until_terminal(): + calls = {"n": 0} + + def handler(request: httpx.Request) -> httpx.Response: + calls["n"] += 1 + return _status_response("queued" if calls["n"] < 3 else "done") + + client = make_client(handler) + st = client.wait_for_status("11111111-1111-1111-1111-111111111111", + statuses=("done",), timeout=5, interval=0.05) + assert st.status == "done" and calls["n"] == 3 + + +def test_wait_for_status_timeout(): + def handler(request: httpx.Request) -> httpx.Response: + return _status_response("queued") + + try: + make_client(handler).wait_for_status( + "11111111-1111-1111-1111-111111111111", + statuses=("done",), timeout=0.2, interval=0.05) + except SynapseError as ex: + assert type(ex).__name__ == "SynapseStatusTimeoutError" + assert "queued" in str(ex) + else: + raise AssertionError("таймаут не сработал") diff --git a/tests/unit/test_webhook.py b/tests/unit/test_webhook.py index 2a56317..34622d9 100644 --- a/tests/unit/test_webhook.py +++ b/tests/unit/test_webhook.py @@ -106,3 +106,21 @@ pass else: raise AssertionError("JSON-массив принят как конверт") + + +def test_parse_event_type() -> None: + from gnexus_synapse import parse_event_type + + assert parse_event_type("monitoring.container.down") == ("monitoring", "container", "down") + assert parse_event_type("a.b") is None + assert parse_event_type("..") is None + assert parse_event_type(None) is None + + +def test_ack() -> None: + from gnexus_synapse import ack + + assert ack() == {"received": True} + assert ack("11111111-1111-1111-1111-111111111111") == { + "received": True, "event_id": "11111111-1111-1111-1111-111111111111", + }