Commit d75989ab authored by George Papathanail's avatar George Papathanail
Browse files

refactor: generalize NatsOperationConsumer over the event model

parent 5dfdba1d
Loading
Loading
Loading
Loading
+9 −4
Original line number Diff line number Diff line
@@ -6,7 +6,7 @@ import nats
import structlog
from nats.aio.client import Client
from nats.aio.subscription import Subscription
from pydantic import ValidationError
from pydantic import BaseModel, ValidationError

from open_exposure_gateway.core.config import NatsSettings
from open_exposure_gateway.domain.srm_events import SRMOperationCompleted
@@ -78,11 +78,13 @@ class NatsOperationConsumer:
        self,
        client: Client,
        subject: str,
        handler: Callable[[SRMOperationCompleted], Awaitable[None]] | None = None,
        handler: Callable[[Any], Awaitable[None]] | None = None,
        event_model: type[BaseModel] = SRMOperationCompleted,
    ) -> None:
        self._client = client
        self._subject = subject
        self._handler = handler
        self._event_model = event_model
        self._subscription: Subscription | None = None

    async def start(self) -> None:
@@ -105,7 +107,7 @@ class NatsOperationConsumer:


        try:
            event = SRMOperationCompleted.model_validate(raw)
            event = self._event_model.model_validate(raw)
        except ValidationError as exc:
            logger.warning("invalid_operation_completed_event", subject=msg.subject, error=str(exc))
            return
@@ -113,4 +115,7 @@ class NatsOperationConsumer:
        try:
            await self._handler(event)
        except Exception:
            logger.exception("operation_completed_handler_failed", operation_id=event.operation_id)
            logger.exception(
                "operation_completed_handler_failed",
                operation_id=getattr(event, "operation_id", None),
            )