Newer
Older
gn-synapse-client-py / src / gnexus_synapse / client.py
@Eugene Sukhodolskiy Eugene Sukhodolskiy 1 day ago 13 KB Initial client library skeleton (v0.1.0)
"""Синхронный клиент 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))