"""Синхронный клиент Ingestion API v1.
Тонкий: один POST на send, никаких ретраев и очередей (у Synapse очередь
своя — Celery). Транспорт инжектный (httpx.Client) — свой создаётся с
явным таймаутом; инжектный не трогаем (таймаут настраивает потребитель).
"""
from __future__ import annotations
import logging
import time
import uuid
from collections.abc import Iterable, 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,
SynapseStatusTimeoutError,
SynapseValidationError,
)
__version__ = "0.1.2"
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,
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, 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_))
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,
user_id: str | None = 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, user_id=user_id)
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 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]:
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,
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,
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))