Loading src/open_exposure_gateway/api/camara/traffic_influence/vwip/router.py +3 −1 Original line number Diff line number Diff line Loading @@ -3,7 +3,7 @@ from uuid import UUID from fastapi import APIRouter, Depends, Query, status from open_exposure_gateway.api.camara.common import XCorrelatorHeader from open_exposure_gateway.api.camara.common import IdempotencyKeyHeader, XCorrelatorHeader from open_exposure_gateway.api.camara.traffic_influence.vwip.schemas import ( PostTrafficInfluence, TrafficInfluence, Loading Loading @@ -87,12 +87,14 @@ async def post_traffic_influence( service: TrafficInfluenceServiceDep, caller: Caller, x_correlator: XCorrelatorHeader = None, idempotency_key: IdempotencyKeyHeader = None, ) -> Any: return await service.create_traffic_influence( request=request, tenant_id=caller.tenant_id, app_provider_id=caller.app_provider_id, x_correlator=x_correlator, idempotency_key=idempotency_key, ) Loading src/open_exposure_gateway/application/services/traffic_influence_service.py +56 −12 Original line number Diff line number Diff line Loading @@ -5,6 +5,7 @@ from uuid import UUID, uuid4 import structlog from pydantic import BaseModel from open_exposure_gateway.adapters.errors import DuplicateOperationError from open_exposure_gateway.api.camara.traffic_influence.vwip.schemas import ( PostTrafficInfluence, ) Loading @@ -18,7 +19,11 @@ from open_exposure_gateway.application.mappers.traffic_influence_mapper import ( build_activate_command, build_deactivate_command, ) from open_exposure_gateway.core.exceptions import DownstreamServiceException, NotFoundException from open_exposure_gateway.core.exceptions import ( ConflictException, DownstreamServiceException, NotFoundException, ) from open_exposure_gateway.domain.models import ( CallbackRegistration, Operation, Loading Loading @@ -97,16 +102,44 @@ class TrafficInfluenceService: return zones[0].resource_zone_id return None async def _replay_create(self, operation: Operation) -> TrafficInfluenceResponse: if operation.operation_type != OperationType.TRAFFIC_INFLUENCE: raise ConflictException( message="Idempotency-Key was already used for a different operation" ) if self._traffic_influence_repo is None: raise RuntimeError("TrafficInfluenceRepository is not available") traffic_influence = await self._traffic_influence_repo.get_by_operation_id( operation.operation_id ) if traffic_influence is None: raise DownstreamServiceException( message="Idempotent replay could not locate the original traffic influence" ) return TrafficInfluenceResponse.model_validate( { **operation.metadata.get("request", {}), "trafficInfluenceID": traffic_influence.traffic_influence_id, "state": traffic_influence.state.value, } ) async def create_traffic_influence( self, request: PostTrafficInfluence, tenant_id: str, app_provider_id: str, x_correlator: Optional[str] = None, idempotency_key: Optional[str] = None, ) -> TrafficInfluenceResponse: if self._operation_repo is None or self._traffic_influence_repo is None: raise RuntimeError("Operation/TrafficInfluence repositories are not available") if idempotency_key is not None: existing = await self._operation_repo.get_by_idempotency_key(tenant_id, idempotency_key) if existing is not None: return await self._replay_create(existing) 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 @@ -127,8 +160,7 @@ class TrafficInfluenceService: resource_zone_id=resource_zone_id, ) await self._operation_repo.save( Operation( operation = Operation( operation_id=operation_id, correlation_id=correlation_id, tenant_id=tenant_id, Loading @@ -136,9 +168,21 @@ class TrafficInfluenceService: operation_type=OperationType.TRAFFIC_INFLUENCE, status=OperationStatus.PENDING, subject=Subject.TASK_ACTIVATE, metadata={"app_id": str(request.appId)}, ) idempotency_key=idempotency_key, metadata={ "app_id": str(request.appId), "request": request.model_dump(mode="json"), }, ) try: await self._operation_repo.save(operation) except DuplicateOperationError: if idempotency_key is None: raise existing = await self._operation_repo.get_by_idempotency_key(tenant_id, idempotency_key) if existing is None: raise return await self._replay_create(existing) source_filter = request.sourceTrafficFilters destination_filter = request.destinationTrafficFilters Loading tests/unit/test_traffic_influence_flows.py +72 −1 Original line number Diff line number Diff line Loading @@ -13,7 +13,12 @@ from fastapi.testclient import TestClient from open_exposure_gateway.api.camara.traffic_influence.vwip.router import ( BASE_PATH as TI_BASE, ) from open_exposure_gateway.domain.models import OperationStatus, TrafficInfluenceState from open_exposure_gateway.domain.models import ( Operation, OperationStatus, OperationType, TrafficInfluenceState, ) from open_exposure_gateway.domain.traffic_influence import ( SRMNetworkCapabilityActivateCommand, SRMNetworkCapabilityDeactivateCommand, Loading Loading @@ -298,3 +303,69 @@ class TestTrafficInfluenceValidation: } response = api_client.post(f"{TI_BASE}/traffic-influences", json=body) assert response.status_code == 400 class TestTrafficInfluenceIdempotency: def test_replayed_post_returns_same_resource_without_republishing( self, api_client: TestClient, fake_bus: FakeDataBus, operation_repo: FakeOperationRepository, traffic_influence_repo: FakeTrafficInfluenceRepository, ) -> None: headers = {"Idempotency-Key": "retry-key-1"} first = api_client.post(f"{TI_BASE}/traffic-influences", json=TI_BODY, headers=headers) second = api_client.post(f"{TI_BASE}/traffic-influences", json=TI_BODY, headers=headers) assert first.status_code == 201 assert second.status_code == 201 assert first.json()["trafficInfluenceID"] == second.json()["trafficInfluenceID"] activates = [p for s, p in fake_bus.published if s == Subject.TASK_ACTIVATE] assert len(activates) == 1 assert len(operation_repo.rows) == 1 assert len(traffic_influence_repo.rows) == 1 def test_distinct_idempotency_keys_create_separate_resources( self, api_client: TestClient, fake_bus: FakeDataBus ) -> None: first = api_client.post( f"{TI_BASE}/traffic-influences", json=TI_BODY, headers={"Idempotency-Key": "key-a"}, ) second = api_client.post( f"{TI_BASE}/traffic-influences", json=TI_BODY, headers={"Idempotency-Key": "key-b"}, ) assert first.json()["trafficInfluenceID"] != second.json()["trafficInfluenceID"] activates = [p for s, p in fake_bus.published if s == Subject.TASK_ACTIVATE] assert len(activates) == 2 def test_idempotency_key_reused_for_different_operation_is_conflict( self, api_client: TestClient, operation_repo: FakeOperationRepository, ) -> None: clashing_operation_id = uuid4() operation_repo.rows[clashing_operation_id] = Operation( operation_id=clashing_operation_id, correlation_id="corr-existing", tenant_id="placeholder", app_provider_id="placeholder", operation_type=OperationType.DEPLOY, status=OperationStatus.COMPLETED, subject="command.srm.service.deploy", idempotency_key="clashing-key", ) response = api_client.post( f"{TI_BASE}/traffic-influences", json=TI_BODY, headers={"Idempotency-Key": "clashing-key"}, ) assert response.status_code == 409 Loading
src/open_exposure_gateway/api/camara/traffic_influence/vwip/router.py +3 −1 Original line number Diff line number Diff line Loading @@ -3,7 +3,7 @@ from uuid import UUID from fastapi import APIRouter, Depends, Query, status from open_exposure_gateway.api.camara.common import XCorrelatorHeader from open_exposure_gateway.api.camara.common import IdempotencyKeyHeader, XCorrelatorHeader from open_exposure_gateway.api.camara.traffic_influence.vwip.schemas import ( PostTrafficInfluence, TrafficInfluence, Loading Loading @@ -87,12 +87,14 @@ async def post_traffic_influence( service: TrafficInfluenceServiceDep, caller: Caller, x_correlator: XCorrelatorHeader = None, idempotency_key: IdempotencyKeyHeader = None, ) -> Any: return await service.create_traffic_influence( request=request, tenant_id=caller.tenant_id, app_provider_id=caller.app_provider_id, x_correlator=x_correlator, idempotency_key=idempotency_key, ) Loading
src/open_exposure_gateway/application/services/traffic_influence_service.py +56 −12 Original line number Diff line number Diff line Loading @@ -5,6 +5,7 @@ from uuid import UUID, uuid4 import structlog from pydantic import BaseModel from open_exposure_gateway.adapters.errors import DuplicateOperationError from open_exposure_gateway.api.camara.traffic_influence.vwip.schemas import ( PostTrafficInfluence, ) Loading @@ -18,7 +19,11 @@ from open_exposure_gateway.application.mappers.traffic_influence_mapper import ( build_activate_command, build_deactivate_command, ) from open_exposure_gateway.core.exceptions import DownstreamServiceException, NotFoundException from open_exposure_gateway.core.exceptions import ( ConflictException, DownstreamServiceException, NotFoundException, ) from open_exposure_gateway.domain.models import ( CallbackRegistration, Operation, Loading Loading @@ -97,16 +102,44 @@ class TrafficInfluenceService: return zones[0].resource_zone_id return None async def _replay_create(self, operation: Operation) -> TrafficInfluenceResponse: if operation.operation_type != OperationType.TRAFFIC_INFLUENCE: raise ConflictException( message="Idempotency-Key was already used for a different operation" ) if self._traffic_influence_repo is None: raise RuntimeError("TrafficInfluenceRepository is not available") traffic_influence = await self._traffic_influence_repo.get_by_operation_id( operation.operation_id ) if traffic_influence is None: raise DownstreamServiceException( message="Idempotent replay could not locate the original traffic influence" ) return TrafficInfluenceResponse.model_validate( { **operation.metadata.get("request", {}), "trafficInfluenceID": traffic_influence.traffic_influence_id, "state": traffic_influence.state.value, } ) async def create_traffic_influence( self, request: PostTrafficInfluence, tenant_id: str, app_provider_id: str, x_correlator: Optional[str] = None, idempotency_key: Optional[str] = None, ) -> TrafficInfluenceResponse: if self._operation_repo is None or self._traffic_influence_repo is None: raise RuntimeError("Operation/TrafficInfluence repositories are not available") if idempotency_key is not None: existing = await self._operation_repo.get_by_idempotency_key(tenant_id, idempotency_key) if existing is not None: return await self._replay_create(existing) 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 @@ -127,8 +160,7 @@ class TrafficInfluenceService: resource_zone_id=resource_zone_id, ) await self._operation_repo.save( Operation( operation = Operation( operation_id=operation_id, correlation_id=correlation_id, tenant_id=tenant_id, Loading @@ -136,9 +168,21 @@ class TrafficInfluenceService: operation_type=OperationType.TRAFFIC_INFLUENCE, status=OperationStatus.PENDING, subject=Subject.TASK_ACTIVATE, metadata={"app_id": str(request.appId)}, ) idempotency_key=idempotency_key, metadata={ "app_id": str(request.appId), "request": request.model_dump(mode="json"), }, ) try: await self._operation_repo.save(operation) except DuplicateOperationError: if idempotency_key is None: raise existing = await self._operation_repo.get_by_idempotency_key(tenant_id, idempotency_key) if existing is None: raise return await self._replay_create(existing) source_filter = request.sourceTrafficFilters destination_filter = request.destinationTrafficFilters Loading
tests/unit/test_traffic_influence_flows.py +72 −1 Original line number Diff line number Diff line Loading @@ -13,7 +13,12 @@ from fastapi.testclient import TestClient from open_exposure_gateway.api.camara.traffic_influence.vwip.router import ( BASE_PATH as TI_BASE, ) from open_exposure_gateway.domain.models import OperationStatus, TrafficInfluenceState from open_exposure_gateway.domain.models import ( Operation, OperationStatus, OperationType, TrafficInfluenceState, ) from open_exposure_gateway.domain.traffic_influence import ( SRMNetworkCapabilityActivateCommand, SRMNetworkCapabilityDeactivateCommand, Loading Loading @@ -298,3 +303,69 @@ class TestTrafficInfluenceValidation: } response = api_client.post(f"{TI_BASE}/traffic-influences", json=body) assert response.status_code == 400 class TestTrafficInfluenceIdempotency: def test_replayed_post_returns_same_resource_without_republishing( self, api_client: TestClient, fake_bus: FakeDataBus, operation_repo: FakeOperationRepository, traffic_influence_repo: FakeTrafficInfluenceRepository, ) -> None: headers = {"Idempotency-Key": "retry-key-1"} first = api_client.post(f"{TI_BASE}/traffic-influences", json=TI_BODY, headers=headers) second = api_client.post(f"{TI_BASE}/traffic-influences", json=TI_BODY, headers=headers) assert first.status_code == 201 assert second.status_code == 201 assert first.json()["trafficInfluenceID"] == second.json()["trafficInfluenceID"] activates = [p for s, p in fake_bus.published if s == Subject.TASK_ACTIVATE] assert len(activates) == 1 assert len(operation_repo.rows) == 1 assert len(traffic_influence_repo.rows) == 1 def test_distinct_idempotency_keys_create_separate_resources( self, api_client: TestClient, fake_bus: FakeDataBus ) -> None: first = api_client.post( f"{TI_BASE}/traffic-influences", json=TI_BODY, headers={"Idempotency-Key": "key-a"}, ) second = api_client.post( f"{TI_BASE}/traffic-influences", json=TI_BODY, headers={"Idempotency-Key": "key-b"}, ) assert first.json()["trafficInfluenceID"] != second.json()["trafficInfluenceID"] activates = [p for s, p in fake_bus.published if s == Subject.TASK_ACTIVATE] assert len(activates) == 2 def test_idempotency_key_reused_for_different_operation_is_conflict( self, api_client: TestClient, operation_repo: FakeOperationRepository, ) -> None: clashing_operation_id = uuid4() operation_repo.rows[clashing_operation_id] = Operation( operation_id=clashing_operation_id, correlation_id="corr-existing", tenant_id="placeholder", app_provider_id="placeholder", operation_type=OperationType.DEPLOY, status=OperationStatus.COMPLETED, subject="command.srm.service.deploy", idempotency_key="clashing-key", ) response = api_client.post( f"{TI_BASE}/traffic-influences", json=TI_BODY, headers={"Idempotency-Key": "clashing-key"}, ) assert response.status_code == 409