Loading src/open_exposure_gateway/adapters/databus/nats_adapter.py +39 −13 Original line number Diff line number Diff line Loading @@ -6,10 +6,11 @@ import nats import structlog from nats.aio.client import Client from nats.aio.subscription import Subscription from pydantic import BaseModel, ValidationError from pydantic import ValidationError from open_exposure_gateway.core.config import NatsSettings from open_exposure_gateway.domain.srm_events import SRMOperationCompleted from open_exposure_gateway.domain.edge_application_management import SRMOperationCompleted from open_exposure_gateway.domain.quality_on_demand import SRMOperationStatus from open_exposure_gateway.ports.databus_port import DataBusPort logger: structlog.BoundLogger = structlog.get_logger(__name__) Loading Loading @@ -78,13 +79,11 @@ class NatsOperationConsumer: self, client: Client, subject: str, handler: Callable[[Any], Awaitable[None]] | None = None, event_model: type[BaseModel] = SRMOperationCompleted, handler: Callable[[SRMOperationCompleted], Awaitable[None]], ) -> 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: Loading @@ -101,21 +100,48 @@ class NatsOperationConsumer: logger.warning("invalid_json", subject=msg.subject) return if self._handler is None: logger.info("operation_completed_received", subject=msg.subject, payload=raw) try: event = SRMOperationCompleted.model_validate(raw) except ValidationError as exc: logger.warning("invalid_operation_completed_event", subject=msg.subject, error=str(exc)) return try: await self._handler(event) except Exception: logger.exception("operation_completed_handler_failed", operation_id=event.operation_id) class NatsOperationStatusConsumer: def __init__( self, client: Client, subject: str, handler: Callable[[SRMOperationStatus], Awaitable[None]], ) -> None: self._client = client self._subject = subject self._handler = handler self._subscription: Subscription | None = None async def start(self) -> None: # TODO: same at-most-once caveat as NatsOperationConsumer.start() — see there. self._subscription = await self._client.subscribe(self._subject, cb=self._handle_message) async def _handle_message(self, msg: _Msg) -> None: try: raw: Any = json.loads(msg.data.decode()) except (json.JSONDecodeError, UnicodeDecodeError): logger.warning("invalid_json", subject=msg.subject) return try: event = self._event_model.model_validate(raw) event = SRMOperationStatus.model_validate(raw) except ValidationError as exc: logger.warning("invalid_operation_completed_event", subject=msg.subject, error=str(exc)) logger.warning("invalid_operation_status_event", subject=msg.subject, error=str(exc)) return try: await self._handler(event) except Exception: logger.exception( "operation_completed_handler_failed", operation_id=getattr(event, "operation_id", None), ) logger.exception("operation_status_handler_failed", operation_id=event.operation_id) src/open_exposure_gateway/application/services/quality_on_demand_service.py +2 −2 Original line number Diff line number Diff line Loading @@ -17,6 +17,7 @@ from open_exposure_gateway.application.mappers.quality_on_demand_mapper import ( ) from open_exposure_gateway.core.config import get_settings from open_exposure_gateway.core.exceptions import DownstreamServiceException, NotFoundException from open_exposure_gateway.domain.edge_application_management import SRMOperationCompleted from open_exposure_gateway.domain.models import ( CallbackDelivery, CallbackRegistration, Loading @@ -26,8 +27,7 @@ from open_exposure_gateway.domain.models import ( QodSession, QodSessionState, ) from open_exposure_gateway.domain.quality_on_demand import Subject from open_exposure_gateway.domain.srm_events import SRMOperationCompleted, SRMOperationStatus from open_exposure_gateway.domain.quality_on_demand import SRMOperationStatus, Subject from open_exposure_gateway.ports.database.callbacks import ( CallbackDeliveryRepository, CallbackRegistrationRepository, Loading src/open_exposure_gateway/core/config.py +0 −4 Original line number Diff line number Diff line Loading @@ -21,10 +21,6 @@ class NatsSettings(BaseModel): max_reconnect_attempts: int = 3 class CallbackSettings(BaseModel): timeout: float = 10.0 class ObservabilitySettings(BaseModel): log_level: str = "INFO" Loading src/open_exposure_gateway/domain/edge_application_management.py +42 −9 Original line number Diff line number Diff line from __future__ import annotations from enum import StrEnum from typing import Any from typing import Any, Literal from uuid import UUID from pydantic import BaseModel, Field from open_exposure_gateway.domain.srm_events import ( SRMCompletedInstance, SRMOperationCompleted, ) __all__ = ["SRMCompletedInstance", "SRMOperationCompleted"] from pydantic import BaseModel, Field, model_validator class Subject(StrEnum): Loading Loading @@ -282,6 +275,46 @@ class SRMTerminateCommand(BaseModel): terminate: SRMTerminatePayload class SRMCompletedInstance(BaseModel): service_instance_id: str zone_id: str status: Literal["completed", "failed"] external_ref: str | None = None error: dict[str, Any] | None = None @model_validator(mode="after") def _require_error_when_failed(self) -> SRMCompletedInstance: if self.status == "failed" and self.error is None: raise ValueError("error is required when instance status is failed") return self class SRMOperationCompleted(BaseModel): schema_version: str operation_id: str status: Literal["completed", "partially_completed", "failed"] service_order_id: str | None = None instances: list[SRMCompletedInstance] = [] metadata: dict[str, Any] | None = None error: dict[str, Any] | None = None correlation_id: str completed_at: str @model_validator(mode="after") def _require_error_when_failed(self) -> SRMOperationCompleted: if self.status == "failed" and not self.instances and self.error is None: raise ValueError( "error is required when status is failed and no instances were produced" ) return self @model_validator(mode="after") def _require_instances_unless_failed(self) -> SRMOperationCompleted: if self.status != "failed" and not self.instances: raise ValueError("instances is required when status != failed") return self class SRMCapabilityEndpoint(BaseModel): interface_id: str fqdn: str | None = None Loading src/open_exposure_gateway/domain/quality_on_demand.py +20 −0 Original line number Diff line number Diff line Loading @@ -73,6 +73,26 @@ class SRMNetworkCapabilityDeactivateCommand(BaseModel): network_capability: NetworkCapabilityDeactivateTarget class SRMOperationStatus(BaseModel): """`event.srm.operation.status` (srm/interface-contract.md §C.1). Covers both pre-execution states (`accepted`, `failed_before_start`) and, per §C.3's "backend-driven change" case, post-completion capability condition changes (e.g. a QoD session dropped by the network) reusing the original realization's `operation_id`. Backend-specific status (e.g. `qos_status`) lives in `metadata`, never in `state`. """ schema_version: str operation_id: str service_order_id: str | None = None service_instance_id: str | None = None capability: str | None = None state: Literal["accepted", "failed_before_start", "completed", "failed", "in_progress"] metadata: dict[str, Any] | None = None correlation_id: str emitted_at: str class EventQosStatusChangedData(BaseModel): sessionId: UUID qosStatus: Literal["AVAILABLE", "UNAVAILABLE"] Loading Loading
src/open_exposure_gateway/adapters/databus/nats_adapter.py +39 −13 Original line number Diff line number Diff line Loading @@ -6,10 +6,11 @@ import nats import structlog from nats.aio.client import Client from nats.aio.subscription import Subscription from pydantic import BaseModel, ValidationError from pydantic import ValidationError from open_exposure_gateway.core.config import NatsSettings from open_exposure_gateway.domain.srm_events import SRMOperationCompleted from open_exposure_gateway.domain.edge_application_management import SRMOperationCompleted from open_exposure_gateway.domain.quality_on_demand import SRMOperationStatus from open_exposure_gateway.ports.databus_port import DataBusPort logger: structlog.BoundLogger = structlog.get_logger(__name__) Loading Loading @@ -78,13 +79,11 @@ class NatsOperationConsumer: self, client: Client, subject: str, handler: Callable[[Any], Awaitable[None]] | None = None, event_model: type[BaseModel] = SRMOperationCompleted, handler: Callable[[SRMOperationCompleted], Awaitable[None]], ) -> 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: Loading @@ -101,21 +100,48 @@ class NatsOperationConsumer: logger.warning("invalid_json", subject=msg.subject) return if self._handler is None: logger.info("operation_completed_received", subject=msg.subject, payload=raw) try: event = SRMOperationCompleted.model_validate(raw) except ValidationError as exc: logger.warning("invalid_operation_completed_event", subject=msg.subject, error=str(exc)) return try: await self._handler(event) except Exception: logger.exception("operation_completed_handler_failed", operation_id=event.operation_id) class NatsOperationStatusConsumer: def __init__( self, client: Client, subject: str, handler: Callable[[SRMOperationStatus], Awaitable[None]], ) -> None: self._client = client self._subject = subject self._handler = handler self._subscription: Subscription | None = None async def start(self) -> None: # TODO: same at-most-once caveat as NatsOperationConsumer.start() — see there. self._subscription = await self._client.subscribe(self._subject, cb=self._handle_message) async def _handle_message(self, msg: _Msg) -> None: try: raw: Any = json.loads(msg.data.decode()) except (json.JSONDecodeError, UnicodeDecodeError): logger.warning("invalid_json", subject=msg.subject) return try: event = self._event_model.model_validate(raw) event = SRMOperationStatus.model_validate(raw) except ValidationError as exc: logger.warning("invalid_operation_completed_event", subject=msg.subject, error=str(exc)) logger.warning("invalid_operation_status_event", subject=msg.subject, error=str(exc)) return try: await self._handler(event) except Exception: logger.exception( "operation_completed_handler_failed", operation_id=getattr(event, "operation_id", None), ) logger.exception("operation_status_handler_failed", operation_id=event.operation_id)
src/open_exposure_gateway/application/services/quality_on_demand_service.py +2 −2 Original line number Diff line number Diff line Loading @@ -17,6 +17,7 @@ from open_exposure_gateway.application.mappers.quality_on_demand_mapper import ( ) from open_exposure_gateway.core.config import get_settings from open_exposure_gateway.core.exceptions import DownstreamServiceException, NotFoundException from open_exposure_gateway.domain.edge_application_management import SRMOperationCompleted from open_exposure_gateway.domain.models import ( CallbackDelivery, CallbackRegistration, Loading @@ -26,8 +27,7 @@ from open_exposure_gateway.domain.models import ( QodSession, QodSessionState, ) from open_exposure_gateway.domain.quality_on_demand import Subject from open_exposure_gateway.domain.srm_events import SRMOperationCompleted, SRMOperationStatus from open_exposure_gateway.domain.quality_on_demand import SRMOperationStatus, Subject from open_exposure_gateway.ports.database.callbacks import ( CallbackDeliveryRepository, CallbackRegistrationRepository, Loading
src/open_exposure_gateway/core/config.py +0 −4 Original line number Diff line number Diff line Loading @@ -21,10 +21,6 @@ class NatsSettings(BaseModel): max_reconnect_attempts: int = 3 class CallbackSettings(BaseModel): timeout: float = 10.0 class ObservabilitySettings(BaseModel): log_level: str = "INFO" Loading
src/open_exposure_gateway/domain/edge_application_management.py +42 −9 Original line number Diff line number Diff line from __future__ import annotations from enum import StrEnum from typing import Any from typing import Any, Literal from uuid import UUID from pydantic import BaseModel, Field from open_exposure_gateway.domain.srm_events import ( SRMCompletedInstance, SRMOperationCompleted, ) __all__ = ["SRMCompletedInstance", "SRMOperationCompleted"] from pydantic import BaseModel, Field, model_validator class Subject(StrEnum): Loading Loading @@ -282,6 +275,46 @@ class SRMTerminateCommand(BaseModel): terminate: SRMTerminatePayload class SRMCompletedInstance(BaseModel): service_instance_id: str zone_id: str status: Literal["completed", "failed"] external_ref: str | None = None error: dict[str, Any] | None = None @model_validator(mode="after") def _require_error_when_failed(self) -> SRMCompletedInstance: if self.status == "failed" and self.error is None: raise ValueError("error is required when instance status is failed") return self class SRMOperationCompleted(BaseModel): schema_version: str operation_id: str status: Literal["completed", "partially_completed", "failed"] service_order_id: str | None = None instances: list[SRMCompletedInstance] = [] metadata: dict[str, Any] | None = None error: dict[str, Any] | None = None correlation_id: str completed_at: str @model_validator(mode="after") def _require_error_when_failed(self) -> SRMOperationCompleted: if self.status == "failed" and not self.instances and self.error is None: raise ValueError( "error is required when status is failed and no instances were produced" ) return self @model_validator(mode="after") def _require_instances_unless_failed(self) -> SRMOperationCompleted: if self.status != "failed" and not self.instances: raise ValueError("instances is required when status != failed") return self class SRMCapabilityEndpoint(BaseModel): interface_id: str fqdn: str | None = None Loading
src/open_exposure_gateway/domain/quality_on_demand.py +20 −0 Original line number Diff line number Diff line Loading @@ -73,6 +73,26 @@ class SRMNetworkCapabilityDeactivateCommand(BaseModel): network_capability: NetworkCapabilityDeactivateTarget class SRMOperationStatus(BaseModel): """`event.srm.operation.status` (srm/interface-contract.md §C.1). Covers both pre-execution states (`accepted`, `failed_before_start`) and, per §C.3's "backend-driven change" case, post-completion capability condition changes (e.g. a QoD session dropped by the network) reusing the original realization's `operation_id`. Backend-specific status (e.g. `qos_status`) lives in `metadata`, never in `state`. """ schema_version: str operation_id: str service_order_id: str | None = None service_instance_id: str | None = None capability: str | None = None state: Literal["accepted", "failed_before_start", "completed", "failed", "in_progress"] metadata: dict[str, Any] | None = None correlation_id: str emitted_at: str class EventQosStatusChangedData(BaseModel): sessionId: UUID qosStatus: Literal["AVAILABLE", "UNAVAILABLE"] Loading