"""Синхронный клиент 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))