Commit 39ed61bc authored by Paris Stentoumis's avatar Paris Stentoumis
Browse files

feat: nats adapter implementation. and some tests

parent 40cd9a8b
Loading
Loading
Loading
Loading
Loading
+4 −0
Original line number Diff line number Diff line
@@ -5,3 +5,7 @@ APP_VERSION="1.5.0"
POSTGRES_SETTINGS__URL = "postgresql+asyncpg://postgres:postgres@localhost:5432/srm"
POSTGRES_SETTINGS__ECHO = true
POSTGRES_SETTINGS__CREATE_SCHEMA_ON_STARTUP = true

NATS_SETTINGS__URL = "nats://localhost:4222"
NATS_SETTINGS__CONNECT_TIMEOUT = 10
NATS_SETTINGS__MAX_RECONNECT_ATTEMPS = 3
+2 −1
Original line number Diff line number Diff line
@@ -20,7 +20,8 @@ dependencies = [
    "structlog>=25.5.0",
    "sunrise6g-opensdk==2.0.0",
    "sqlalchemy>=2.0.48",
    "asyncpg>=0.31.0"
    "asyncpg>=0.31.0",
    "nats-py>=2.10.0",
]

[project.optional-dependencies]

src/srm/adapters/databus/.gitkeep

deleted100644 → 0
+0 −0

Empty file deleted.

+60 −0
Original line number Diff line number Diff line
import nats
import structlog
from nats.aio.client import Client

from srm.config import NatsSettings

logger: structlog.BoundLogger = structlog.get_logger(__name__)


async def init_databus_manager(settings: NatsSettings) -> "NatsConnectionManager":
    connection_manager: NatsConnectionManager = NatsConnectionManager(settings=settings)
    await connection_manager.connect()
    return connection_manager


class NatsConnectionManager:
    def __init__(self, settings: NatsSettings) -> None:
        self._settings = settings
        self._client: Client | None = None

    @property
    def is_connected(self) -> bool:
        return self._client is not None and self._client.is_connected

    @property
    def client(self) -> Client:
        if self._client is None:
            raise RuntimeError("NATS client is not connected")
        return self._client

    async def connect(self) -> None:
        async def _on_error(e: Exception) -> None:
            logger.error("nats_error", error=str(e))

        async def _on_disconnect() -> None:
            logger.warning("nats_disconnected", url=self._settings.url)

        async def _on_reconnect() -> None:
            logger.info("nats_reconnected", url=self._settings.url)

        if not self.is_connected:
            try:
                self._client = await nats.connect(
                    servers=[self._settings.url],
                    connect_timeout=self._settings.connect_timeout,
                    max_reconnect_attempts=self._settings.max_reconnect_attempts,
                    error_cb=_on_error,
                    disconnected_cb=_on_disconnect,
                    reconnected_cb=_on_reconnect,
                )
            except Exception as e:
                logger.error("nats_error", error=str(e))
                raise
        else:
            logger.info("nats_already_connected")

    async def close(self) -> None:
        if self._client is not None:
            await self._client.drain()
            self._client = None
+20 −0
Original line number Diff line number Diff line
import json
from typing import Any

from srm.adapters.databus.nats_connection_manager import NatsConnectionManager
from srm.domain.ports.databus.publisher import DataBusPublisher


class NatsPublisher(DataBusPublisher):
    def __init__(self, connection_manager: NatsConnectionManager) -> None:
        self._connection_manager = connection_manager

    async def publish(
        self,
        subject: str,
        payload: dict[str, Any],
        headers: dict[str, str] | None = None,
    ) -> None:

        body = json.dumps(payload).encode("utf-8")
        await self._connection_manager.client.publish(subject, body, headers=headers)
Loading