Commit 94657429 authored by Dimitrios Gogos's avatar Dimitrios Gogos
Browse files

Merge branch 'refactor/nats-adapter' into 'refactor/dev'

SRM Refactor: Implementation of NATS adapter

See merge request !15
parents 40cd9a8b 62768689
Loading
Loading
Loading
Loading
Loading
+5 −0
Original line number Diff line number Diff line
@@ -5,3 +5,8 @@ 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_ATTEMPTS = 3
NATS_SETTINGS__DRAIN_TIMEOUT = 30
+8 −20
Original line number Diff line number Diff line
default:
  tags:
    - docker
  image: python:3.12-slim
  cache:
    paths:
@@ -19,7 +17,6 @@ stages:
  - format
  - test
  - build
  - push

variables:
  UV_CACHE_DIR: "$CI_PROJECT_DIR/.cache/uv"
@@ -57,28 +54,19 @@ test:

build:
  stage: build
  tags:
    - shell
  before_script:
    - docker info
  script:
    - export TEST_IMAGE_TAG="ci-${CI_COMMIT_REF_SLUG}-${CI_COMMIT_SHORT_SHA}"
    - docker build --network=host -t "$CI_REGISTRY_IMAGE:$TEST_IMAGE_TAG" .
  rules:
    - if: '$CI_COMMIT_BRANCH'

push:
  stage: push
  tags:
    - shell
  needs:
    - build
  image: docker:cli
  variables:
    DOCKER_HOST: tcp://docker:2375
    DOCKER_TLS_CERTDIR: ""
  services:
    - docker:dind
  before_script:
    - docker info
  script:
    - export TEST_IMAGE_TAG="ci-${CI_COMMIT_REF_SLUG}-${CI_COMMIT_SHORT_SHA}"
    - echo "$CI_REGISTRY_PASSWORD" | docker login -u "$CI_REGISTRY_USER" "$CI_REGISTRY" --password-stdin
    - docker build --network=host -t "$CI_REGISTRY_IMAGE:$TEST_IMAGE_TAG" .
    - docker push "$CI_REGISTRY_IMAGE:$TEST_IMAGE_TAG"
    - docker logout "$CI_REGISTRY"
  rules:
    - if: '$CI_COMMIT_BRANCH'
    - if: '$CI_COMMIT_BRANCH == "main" || $CI_COMMIT_BRANCH == "develop"'
+1 −0
Original line number Diff line number Diff line
@@ -27,6 +27,7 @@ repos:
        - "docker>=7.0.0"
        - "fastapi[standard]>=0.135.1"
        - "httpx>=0.27.0"
        - "nats-py>=2.10.0"
        - "pydantic-settings>=2.13.1"
        - "pytest>=9.0.2"
        - "pytest-asyncio>=0.24"
+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]
+76 −0
Original line number Diff line number Diff line
import asyncio

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
        self._connect_lock = asyncio.Lock()

    @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)

        async with self._connect_lock:
            if self.is_connected:
                logger.info("nats_already_connected")
                return

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

    async def close(self) -> None:
        if self._client is None:
            return

        client = self._client
        try:
            # drain() is bounded by the drain_timeout passed to connect(); on timeout it
            # reports through error_cb and closes anyway, so shutdown cannot hang here.
            await client.drain()
        except Exception as e:
            logger.warning("nats_drain_failed", error=str(e))
            await client.close()
        finally:
            self._client = None
Loading