Newer
Older
gn-synapse-client-py / src / gnexus_synapse / _async.py
@Eugene Sukhodolskiy Eugene Sukhodolskiy 1 day ago 9 KB Initial client library skeleton (v0.1.0)
"""Асинхронный близнец 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