Loading src/open_exposure_gateway/adapters/database/mappers.py +2 −2 Original line number Diff line number Diff line Loading @@ -150,7 +150,7 @@ class TrafficInfluenceMapper: return TrafficInfluence( traffic_influence_id=row.traffic_influence_id, operation_id=row.operation_id, app_id=row.app_id, app_registration_id=row.app_registration_id, edge_cloud_zone_id=row.edge_cloud_zone_id, edge_cloud_region=row.edge_cloud_region, source_port=row.source_port, Loading @@ -167,7 +167,7 @@ class TrafficInfluenceMapper: return TrafficInfluenceRow( traffic_influence_id=domain.traffic_influence_id, operation_id=domain.operation_id, app_id=domain.app_id, app_registration_id=domain.app_registration_id, edge_cloud_zone_id=domain.edge_cloud_zone_id, edge_cloud_region=domain.edge_cloud_region, source_port=domain.source_port, Loading src/open_exposure_gateway/adapters/database/repos/traffic_influences.py +5 −3 Original line number Diff line number Diff line Loading @@ -25,10 +25,12 @@ class SqlTrafficInfluenceRepository(TrafficInfluenceRepository): row = await self._session.scalar(stmt) return TrafficInfluenceMapper.to_domain(row) if row is not None else None async def list_by_app_id(self, app_id: UUID | None) -> list[TrafficInfluence]: async def list_by_app_registration_id( self, app_registration_id: UUID | None ) -> list[TrafficInfluence]: stmt = select(TrafficInfluenceRow) if app_id is not None: stmt = stmt.where(TrafficInfluenceRow.app_id == app_id) if app_registration_id is not None: stmt = stmt.where(TrafficInfluenceRow.app_registration_id == app_registration_id) rows = await self._session.scalars(stmt) return [TrafficInfluenceMapper.to_domain(row) for row in rows] Loading src/open_exposure_gateway/adapters/database/sql.py +7 −2 Original line number Diff line number Diff line Loading @@ -207,13 +207,18 @@ class QodSessionRow(AuditedMixin, Base): class TrafficInfluenceRow(AuditedMixin, Base): __tablename__ = "traffic_influences" __table_args__ = (Index("idx_traffic_influences_operation", "operation_id"),) __table_args__ = ( Index("idx_traffic_influences_operation", "operation_id"), Index("idx_traffic_influences_app_registration", "app_registration_id"), ) traffic_influence_id: Mapped[UUID] = mapped_column(PG_UUID(as_uuid=True), primary_key=True) operation_id: Mapped[UUID] = mapped_column( ForeignKey("operations.operation_id"), nullable=False ) app_id: Mapped[UUID] = mapped_column(PG_UUID(as_uuid=True), nullable=False) app_registration_id: Mapped[UUID] = mapped_column( ForeignKey("app_registrations.app_registration_id"), nullable=False ) edge_cloud_zone_id: Mapped[UUID | None] = mapped_column(PG_UUID(as_uuid=True)) edge_cloud_region: Mapped[str | None] = mapped_column(String(256)) source_port: Mapped[int | None] = mapped_column(Integer) Loading src/open_exposure_gateway/adapters/http/srm_client.py +3 −6 Original line number Diff line number Diff line Loading @@ -200,9 +200,6 @@ class SRMClient: ) return SRMNetworkCapability.model_validate(data) # NOTE: /internal/network-capabilities is not documented in # architecture/srm/interface-contract.md §E as a general capability read. # Confirm the real path with the SRM dev before relying on this in production. async def get_traffic_influence( self, traffic_influence_id: str, Loading @@ -216,12 +213,12 @@ class SRMClient: async def get_traffic_influences( self, app_id: UUID | None = None, service_specification_id: UUID | None = None, x_correlator: str | None = None, ) -> list[SRMTrafficInfluenceCapability]: params = {"capability_type": "traffic_influence"} if app_id is not None: params["app_id"] = str(app_id) if service_specification_id is not None: params["service_specification_id"] = str(service_specification_id) headers = {"x-correlator": x_correlator} if x_correlator else None data = await self._request( "GET", "/internal/network-capabilities", params=params, headers=headers Loading src/open_exposure_gateway/application/services/traffic_influence_service.py +30 −7 Original line number Diff line number Diff line Loading @@ -49,6 +49,7 @@ from open_exposure_gateway.ports.database.callbacks import ( CallbackRegistrationRepository, ) from open_exposure_gateway.ports.database.operations import OperationRepository from open_exposure_gateway.ports.database.registration import AppRegistrationRepository from open_exposure_gateway.ports.database.traffic_influences import TrafficInfluenceRepository from open_exposure_gateway.ports.databus_port import DataBusPort from open_exposure_gateway.ports.srm_port import SRMClientPort Loading Loading @@ -89,6 +90,7 @@ class TrafficInfluenceService: callback_registration_repo: Optional[CallbackRegistrationRepository] = None, callback_delivery_repo: Optional[CallbackDeliveryRepository] = None, callback_delivery_port: Optional[TrafficInfluenceCallbackDeliveryPort] = None, app_registration_repo: Optional[AppRegistrationRepository] = None, ) -> None: self.srm_client = srm_client self._publisher = publisher Loading @@ -97,6 +99,7 @@ class TrafficInfluenceService: self._callback_registration_repo = callback_registration_repo self._callback_delivery_repo = callback_delivery_repo self._callback_delivery_port = callback_delivery_port self._app_registration_repo = app_registration_repo def _parse_traffic_influence_id(self, traffic_influence_id: str) -> UUID: try: Loading Loading @@ -152,6 +155,14 @@ class TrafficInfluenceService: return instances[0].zone_id return None async def _resolve_app_registration_id(self, app_id: UUID) -> UUID: if self._app_registration_repo is None: raise RuntimeError("AppRegistrationRepository is not available") app_registration = await self._app_registration_repo.get_by_app_id(app_id) if app_registration is None: raise BadRequestException(message=f"App {app_id} is not registered") return app_registration.app_registration_id async def _replay_create(self, operation: Operation) -> TrafficInfluenceResponse: if operation.operation_type != OperationType.TRAFFIC_INFLUENCE: raise ConflictException( Loading Loading @@ -186,6 +197,8 @@ class TrafficInfluenceService: if existing is not None: return await self._replay_create(existing) app_registration_id = await self._resolve_app_registration_id(request.appId) operation_id, correlation_id, requested_at = self._new_operation_metadata(x_correlator) traffic_influence_id = uuid4() # edgeCloudRegion has no SRM field of its own (srm/canonical-parameters-schema.md Loading @@ -198,10 +211,7 @@ class TrafficInfluenceService: operation_id=operation_id, correlation_id=correlation_id, requested_at=requested_at, # service_specification_id == appId for an already-onboarded app # (srm/interface-contract.md §E.4, ADR-0011) -- unlike QoD, which mints a # throwaway id because its sessions aren't tied to an onboarded catalog entry. service_specification_id=request.appId, service_specification_id=app_registration_id, service_instance_id=traffic_influence_id, app_provider_id=app_provider_id, zone_id=zone_id, Loading Loading @@ -237,7 +247,7 @@ class TrafficInfluenceService: TrafficInfluenceRecord( traffic_influence_id=traffic_influence_id, operation_id=operation_id, app_id=request.appId, app_registration_id=app_registration_id, edge_cloud_zone_id=request.edgeCloudZoneId, edge_cloud_region=request.edgeCloudRegion, source_port=source_filter.sourcePort if source_filter else None, Loading Loading @@ -542,7 +552,20 @@ class TrafficInfluenceService: if self._traffic_influence_repo is None: raise RuntimeError("TrafficInfluence repository is not available") records = await self._traffic_influence_repo.list_by_app_id(app_id) # An unknown appId filters to nothing rather than erroring, matching # EdgeApplicationManagementService.get_app_instances. app_registration_id: UUID | None = None if app_id is not None: if self._app_registration_repo is None: raise RuntimeError("AppRegistrationRepository is not available") app_registration = await self._app_registration_repo.get_by_app_id(app_id) if app_registration is None: return [] app_registration_id = app_registration.app_registration_id records = await self._traffic_influence_repo.list_by_app_registration_id( app_registration_id ) # Best-effort, same rationale as get_traffic_influence. OEG mints # traffic_influence_id and sends it to SRM as service_instance_id at activation Loading @@ -550,7 +573,7 @@ class TrafficInfluenceService: capabilities_by_id: dict[str, SRMTrafficInfluenceCapability] = {} try: capabilities = await self.srm_client.get_traffic_influences( app_id=app_id, service_specification_id=app_registration_id, x_correlator=x_correlator, ) capabilities_by_id = {c.service_instance_id: c for c in capabilities} Loading Loading
src/open_exposure_gateway/adapters/database/mappers.py +2 −2 Original line number Diff line number Diff line Loading @@ -150,7 +150,7 @@ class TrafficInfluenceMapper: return TrafficInfluence( traffic_influence_id=row.traffic_influence_id, operation_id=row.operation_id, app_id=row.app_id, app_registration_id=row.app_registration_id, edge_cloud_zone_id=row.edge_cloud_zone_id, edge_cloud_region=row.edge_cloud_region, source_port=row.source_port, Loading @@ -167,7 +167,7 @@ class TrafficInfluenceMapper: return TrafficInfluenceRow( traffic_influence_id=domain.traffic_influence_id, operation_id=domain.operation_id, app_id=domain.app_id, app_registration_id=domain.app_registration_id, edge_cloud_zone_id=domain.edge_cloud_zone_id, edge_cloud_region=domain.edge_cloud_region, source_port=domain.source_port, Loading
src/open_exposure_gateway/adapters/database/repos/traffic_influences.py +5 −3 Original line number Diff line number Diff line Loading @@ -25,10 +25,12 @@ class SqlTrafficInfluenceRepository(TrafficInfluenceRepository): row = await self._session.scalar(stmt) return TrafficInfluenceMapper.to_domain(row) if row is not None else None async def list_by_app_id(self, app_id: UUID | None) -> list[TrafficInfluence]: async def list_by_app_registration_id( self, app_registration_id: UUID | None ) -> list[TrafficInfluence]: stmt = select(TrafficInfluenceRow) if app_id is not None: stmt = stmt.where(TrafficInfluenceRow.app_id == app_id) if app_registration_id is not None: stmt = stmt.where(TrafficInfluenceRow.app_registration_id == app_registration_id) rows = await self._session.scalars(stmt) return [TrafficInfluenceMapper.to_domain(row) for row in rows] Loading
src/open_exposure_gateway/adapters/database/sql.py +7 −2 Original line number Diff line number Diff line Loading @@ -207,13 +207,18 @@ class QodSessionRow(AuditedMixin, Base): class TrafficInfluenceRow(AuditedMixin, Base): __tablename__ = "traffic_influences" __table_args__ = (Index("idx_traffic_influences_operation", "operation_id"),) __table_args__ = ( Index("idx_traffic_influences_operation", "operation_id"), Index("idx_traffic_influences_app_registration", "app_registration_id"), ) traffic_influence_id: Mapped[UUID] = mapped_column(PG_UUID(as_uuid=True), primary_key=True) operation_id: Mapped[UUID] = mapped_column( ForeignKey("operations.operation_id"), nullable=False ) app_id: Mapped[UUID] = mapped_column(PG_UUID(as_uuid=True), nullable=False) app_registration_id: Mapped[UUID] = mapped_column( ForeignKey("app_registrations.app_registration_id"), nullable=False ) edge_cloud_zone_id: Mapped[UUID | None] = mapped_column(PG_UUID(as_uuid=True)) edge_cloud_region: Mapped[str | None] = mapped_column(String(256)) source_port: Mapped[int | None] = mapped_column(Integer) Loading
src/open_exposure_gateway/adapters/http/srm_client.py +3 −6 Original line number Diff line number Diff line Loading @@ -200,9 +200,6 @@ class SRMClient: ) return SRMNetworkCapability.model_validate(data) # NOTE: /internal/network-capabilities is not documented in # architecture/srm/interface-contract.md §E as a general capability read. # Confirm the real path with the SRM dev before relying on this in production. async def get_traffic_influence( self, traffic_influence_id: str, Loading @@ -216,12 +213,12 @@ class SRMClient: async def get_traffic_influences( self, app_id: UUID | None = None, service_specification_id: UUID | None = None, x_correlator: str | None = None, ) -> list[SRMTrafficInfluenceCapability]: params = {"capability_type": "traffic_influence"} if app_id is not None: params["app_id"] = str(app_id) if service_specification_id is not None: params["service_specification_id"] = str(service_specification_id) headers = {"x-correlator": x_correlator} if x_correlator else None data = await self._request( "GET", "/internal/network-capabilities", params=params, headers=headers Loading
src/open_exposure_gateway/application/services/traffic_influence_service.py +30 −7 Original line number Diff line number Diff line Loading @@ -49,6 +49,7 @@ from open_exposure_gateway.ports.database.callbacks import ( CallbackRegistrationRepository, ) from open_exposure_gateway.ports.database.operations import OperationRepository from open_exposure_gateway.ports.database.registration import AppRegistrationRepository from open_exposure_gateway.ports.database.traffic_influences import TrafficInfluenceRepository from open_exposure_gateway.ports.databus_port import DataBusPort from open_exposure_gateway.ports.srm_port import SRMClientPort Loading Loading @@ -89,6 +90,7 @@ class TrafficInfluenceService: callback_registration_repo: Optional[CallbackRegistrationRepository] = None, callback_delivery_repo: Optional[CallbackDeliveryRepository] = None, callback_delivery_port: Optional[TrafficInfluenceCallbackDeliveryPort] = None, app_registration_repo: Optional[AppRegistrationRepository] = None, ) -> None: self.srm_client = srm_client self._publisher = publisher Loading @@ -97,6 +99,7 @@ class TrafficInfluenceService: self._callback_registration_repo = callback_registration_repo self._callback_delivery_repo = callback_delivery_repo self._callback_delivery_port = callback_delivery_port self._app_registration_repo = app_registration_repo def _parse_traffic_influence_id(self, traffic_influence_id: str) -> UUID: try: Loading Loading @@ -152,6 +155,14 @@ class TrafficInfluenceService: return instances[0].zone_id return None async def _resolve_app_registration_id(self, app_id: UUID) -> UUID: if self._app_registration_repo is None: raise RuntimeError("AppRegistrationRepository is not available") app_registration = await self._app_registration_repo.get_by_app_id(app_id) if app_registration is None: raise BadRequestException(message=f"App {app_id} is not registered") return app_registration.app_registration_id async def _replay_create(self, operation: Operation) -> TrafficInfluenceResponse: if operation.operation_type != OperationType.TRAFFIC_INFLUENCE: raise ConflictException( Loading Loading @@ -186,6 +197,8 @@ class TrafficInfluenceService: if existing is not None: return await self._replay_create(existing) app_registration_id = await self._resolve_app_registration_id(request.appId) operation_id, correlation_id, requested_at = self._new_operation_metadata(x_correlator) traffic_influence_id = uuid4() # edgeCloudRegion has no SRM field of its own (srm/canonical-parameters-schema.md Loading @@ -198,10 +211,7 @@ class TrafficInfluenceService: operation_id=operation_id, correlation_id=correlation_id, requested_at=requested_at, # service_specification_id == appId for an already-onboarded app # (srm/interface-contract.md §E.4, ADR-0011) -- unlike QoD, which mints a # throwaway id because its sessions aren't tied to an onboarded catalog entry. service_specification_id=request.appId, service_specification_id=app_registration_id, service_instance_id=traffic_influence_id, app_provider_id=app_provider_id, zone_id=zone_id, Loading Loading @@ -237,7 +247,7 @@ class TrafficInfluenceService: TrafficInfluenceRecord( traffic_influence_id=traffic_influence_id, operation_id=operation_id, app_id=request.appId, app_registration_id=app_registration_id, edge_cloud_zone_id=request.edgeCloudZoneId, edge_cloud_region=request.edgeCloudRegion, source_port=source_filter.sourcePort if source_filter else None, Loading Loading @@ -542,7 +552,20 @@ class TrafficInfluenceService: if self._traffic_influence_repo is None: raise RuntimeError("TrafficInfluence repository is not available") records = await self._traffic_influence_repo.list_by_app_id(app_id) # An unknown appId filters to nothing rather than erroring, matching # EdgeApplicationManagementService.get_app_instances. app_registration_id: UUID | None = None if app_id is not None: if self._app_registration_repo is None: raise RuntimeError("AppRegistrationRepository is not available") app_registration = await self._app_registration_repo.get_by_app_id(app_id) if app_registration is None: return [] app_registration_id = app_registration.app_registration_id records = await self._traffic_influence_repo.list_by_app_registration_id( app_registration_id ) # Best-effort, same rationale as get_traffic_influence. OEG mints # traffic_influence_id and sends it to SRM as service_instance_id at activation Loading @@ -550,7 +573,7 @@ class TrafficInfluenceService: capabilities_by_id: dict[str, SRMTrafficInfluenceCapability] = {} try: capabilities = await self.srm_client.get_traffic_influences( app_id=app_id, service_specification_id=app_registration_id, x_correlator=x_correlator, ) capabilities_by_id = {c.service_instance_id: c for c in capabilities} Loading