"""Асинхронный близнец 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