/
azathd
/
mutiagent
Обзор
Документация
Войти
/
azathd
/
mutiagent
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/common/bus.py
109 строк
3 KB
Your Name
init
15 май 2026, 17:47
15 май 2026, 17:47
359c81e
Код
Авторство
О чём код?
"""Async NATS JetStream helper (publish / subscribe).""" from __future__ import annotations import json import os from collections.abc import Awaitable, Callable from typing import Any, TypeAlias from uuid import UUID import nats from nats.aio.client import Client as NATSClient from nats.js import JetStreamContext from nats.js.api import RetentionPolicy, StorageType, StreamConfig from nats.js.errors import NotFoundError from src.common.models import Envelope SubscribeHandler: TypeAlias = Callable[[Envelope], Awaitable[None]] DEFAULT_STREAM = "XRAY_EVENTS" def envelope_to_bytes(env: Envelope) -> bytes: """Serialize ``Envelope`` to JSON bytes.""" def _default(o: Any) -> str: if isinstance(o, UUID): return str(o) raise TypeError(type(o).__name__) return json.dumps(env.model_dump(mode="json"), default=_default).encode("utf-8") def envelope_from_json(data: bytes) -> Envelope: raw: dict[str, Any] = json.loads(data.decode("utf-8")) payload = raw.get("payload") or {} if not isinstance(payload, dict): msg = "payload must be object" raise ValueError(msg) raw["payload"] = payload return Envelope.model_validate(raw) class MessageBus: """Thin JetStream client: guarantees stream exists, publishes with ack.""" def __init__(self, nats_url: str | None = None, stream_name: str = DEFAULT_STREAM) -> None: self._nats_url = nats_url or os.environ.get("NATS_URL", "nats://nats:4222") self._stream_name = stream_name self._nc: NATSClient | None = None self._js: JetStreamContext | None = None @property def js(self) -> JetStreamContext: if self._js is None: msg = "not connected" raise RuntimeError(msg) return self._js async def connect(self) -> None: if self._nc is not None: return nc = await nats.connect(self._nats_url) js = nc.jetstream() try: await js.stream_info(self._stream_name) except NotFoundError: cfg = StreamConfig( name=self._stream_name, subjects=["xray.>"], retention=RetentionPolicy.LIMITS, storage=StorageType.FILE, max_age=86400 * 7, ) await js.add_stream(cfg) self._nc = nc self._js = js async def close(self) -> None: if self._nc: await self._nc.drain() await self._nc.close() self._nc = None self._js = None async def publish(self, subject: str, env: Envelope) -> None: """Publish envelope; awaits server ack.""" await self.js.publish(subject, envelope_to_bytes(env)) async def subscribe_durable(self, subject: str, durable: str, handler: SubscribeHandler) -> None: """Push consumer with durable name; awaits ``handler`` before ack.""" async def _cb(msg: Any) -> None: payload = envelope_from_json(msg.data) await handler(payload) await msg.ack() await self.js.subscribe( subject, durable=durable, cb=_cb, manual_ack=True, stream=self._stream_name, )