Commit 6bdd1f71 authored by Sergio Gimenez's avatar Sergio Gimenez
Browse files

feat(fm): NATS JetStream command publisher

parent 54a69e11
Loading
Loading
Loading
Loading
+7 −0
Original line number Diff line number Diff line
@@ -16,5 +16,12 @@ services:
    volumes:
      - fm-pgdata:/var/lib/postgresql/data

  nats:
    image: nats:2.10-alpine
    container_name: fm-nats
    command: ["-js"]
    ports:
      - "4222:4222"

volumes:
  fm-pgdata:
+1 −0
Original line number Diff line number Diff line
@@ -11,6 +11,7 @@ requires-python = ">=3.12"
dependencies = [
    "asyncpg>=0.30",
    "fastapi[standard]>=0.115",
    "nats-py>=2.6",
    "pydantic-settings>=2.4",
    "sqlalchemy[asyncio]>=2.0",
]
+0 −0

Empty file added.

+47 −0
Original line number Diff line number Diff line
import json

import nats
from nats.aio.client import Client
from nats.js import JetStreamContext
from nats.js.api import RetentionPolicy, StreamConfig

from federation_manager.contracts.srm import TASK_STREAM


class NatsCommandPublisher:
    def __init__(self, url: str) -> None:
        self._url = url
        self._nc: Client | None = None
        self._js: JetStreamContext | None = None

    async def connect(self) -> None:
        self._nc = await nats.connect(self._url)
        self._js = self._nc.jetstream()

    async def close(self) -> None:
        if self._nc is not None:
            await self._nc.drain()
            self._nc = None
            self._js = None

    async def ensure_task_stream(self) -> None:
        js = self._require_js()
        try:
            await js.stream_info(TASK_STREAM)
        except Exception:
            await js.add_stream(
                StreamConfig(
                    name=TASK_STREAM,
                    subjects=["command.srm.>"],
                    retention=RetentionPolicy.WORK_QUEUE,
                )
            )

    async def publish(self, subject: str, payload: dict[str, object]) -> None:
        js = self._require_js()
        await js.publish(subject, json.dumps(payload).encode())

    def _require_js(self) -> JetStreamContext:
        if self._js is None:
            raise RuntimeError("NATS publisher is not connected")
        return self._js
+4 −0
Original line number Diff line number Diff line
@@ -7,6 +7,10 @@ from pydantic import BaseModel, ConfigDict, Field

Source = Literal["nbi_camara", "nbi_tmf", "operator_portal", "federation"]

TASK_STREAM = "OOP_TASKS"
SUBJECT_DEPLOY = "command.srm.service.deploy"
SUBJECT_TERMINATE = "command.srm.service.terminate"


class CommandEnvelopeV1(BaseModel):
    schema_version: str = "1.0"
Loading