Loading src/open_exposure_gateway/application/services/edge_application_management_service.py +22 −0 Original line number Diff line number Diff line Loading @@ -43,6 +43,7 @@ from open_exposure_gateway.domain.models import ( AppInstanceState, AppRegistration, AppRegistrationStatus, CallbackRegistration, Operation, OperationStatus, OperationType, Loading Loading @@ -279,6 +280,27 @@ class EdgeApplicationManagementService: ) ) if request.subscriptionRequest is not None: if self._callback_registration_repo is None: raise RuntimeError("CallbackRegistrationRepository is not available") subscription = request.subscriptionRequest await self._callback_registration_repo.save( CallbackRegistration( id=uuid4(), operation_id=operation_id, tenant_id=tenant_id, api_family="edge-application-management", sink=subscription.sink, event_types=subscription.types, sink_credential_ref=f"secret://oeg/{operation_id}/sink-credential" if subscription.sinkCredential is not None else None, expires_at=subscription.config.subscriptionExpireTime if subscription.config else None, ) ) _command = build_deploy_command(translation, requested_at) await self._publish( Subject.TASK_DEPLOY, Loading tests/unit/test_eam_flows.py +21 −0 Original line number Diff line number Diff line Loading @@ -30,6 +30,7 @@ from open_exposure_gateway.domain.edge_application_management import ( from open_exposure_gateway.domain.models import AppInstanceState, OperationStatus from tests.unit.fakes import ( FakeAppInstanceRepository, FakeCallbackRegistrationRepository, FakeDataBus, FakeMsg, FakeOperationRepository, Loading Loading @@ -116,6 +117,26 @@ class TestCreateAppInstanceFlow: assert stored_instance is not None assert stored_instance.state == AppInstanceState.INSTANTIATING def test_persists_callback_registration_when_subscription_request_present( self, api_client: TestClient, live_srm: FakeSRMClient, callback_registration_repo: FakeCallbackRegistrationRepository, ) -> None: body = { **CREATE_INSTANCE_BODY, "subscriptionRequest": { "sink": "https://client.example.com/callback", "types": [ "org.camaraproject.edge-application-management.v0.app-instance-status-change" ], }, } api_client.post(f"{EAM_BASE}/appinstances", json=body) (callback,) = list(callback_registration_repo.rows.values()) assert callback.sink == "https://client.example.com/callback" def test_srm_receives_a_valid_deploy_command( self, api_client: TestClient, fake_bus: FakeDataBus, live_srm: FakeSRMClient ) -> None: Loading tests/unit/test_eam_service.py +98 −0 Original line number Diff line number Diff line from datetime import datetime, timezone from unittest.mock import AsyncMock from uuid import UUID, uuid4 Loading @@ -10,6 +11,8 @@ from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas i AppRepo, CreateAppInstanceRequest, KubernetesResources, SubscriptionConfig, SubscriptionRequest, VmResources, ) from open_exposure_gateway.application.services.edge_application_management_service import ( Loading Loading @@ -601,6 +604,101 @@ class TestCreateAppInstance: ) assert len(operation_repo.rows) == 2 def _make_request_with_subscription( self, with_credential: bool = False, expires_at: datetime | None = None ) -> CreateAppInstanceRequest: return CreateAppInstanceRequest( name="myapp_inst", appId=APP_ID, edgeCloudZoneId=ZONE_ID, subscriptionRequest=SubscriptionRequest( sink="https://client.example.com/callback", sinkCredential={"credentialType": "ACCESSTOKEN", "accessToken": "raw-token"} if with_credential else None, types=[ "org.camaraproject.edge-application-management.v0.app-instance-status-change" ], config=SubscriptionConfig(subscriptionExpireTime=expires_at) if expires_at else None, ), ) async def test_persists_callback_registration_when_subscription_present( self, service: EdgeApplicationManagementService, operation_repo: FakeOperationRepository, callback_registration_repo: FakeCallbackRegistrationRepository, ) -> None: expires_at = datetime(2027, 1, 17, 13, 18, 23, tzinfo=timezone.utc) await service.create_app_instance( request=self._make_request_with_subscription(expires_at=expires_at), tenant_id="tenant-1", app_provider_id="provider-1", ) (operation,) = list(operation_repo.rows.values()) callback = await callback_registration_repo.get_by_operation_id(operation.operation_id) assert callback is not None assert callback.tenant_id == "tenant-1" assert callback.api_family == "edge-application-management" assert callback.sink == "https://client.example.com/callback" assert callback.event_types == [ "org.camaraproject.edge-application-management.v0.app-instance-status-change" ] assert callback.is_active is True assert callback.expires_at == expires_at assert callback.sink_credential_ref is None async def test_sink_credential_is_stored_as_secret_reference_not_raw( self, service: EdgeApplicationManagementService, callback_registration_repo: FakeCallbackRegistrationRepository, ) -> None: await service.create_app_instance( request=self._make_request_with_subscription(with_credential=True), tenant_id="tenant-1", app_provider_id="provider-1", ) (callback,) = list(callback_registration_repo.rows.values()) assert callback.sink_credential_ref is not None assert callback.sink_credential_ref.startswith("secret://") assert "raw-token" not in callback.sink_credential_ref async def test_no_callback_registration_when_subscription_absent( self, service: EdgeApplicationManagementService, callback_registration_repo: FakeCallbackRegistrationRepository, ) -> None: await service.create_app_instance( request=self._make_request(), tenant_id="tenant-1", app_provider_id="provider-1", ) assert callback_registration_repo.rows == {} async def test_raises_when_callback_registration_repo_unavailable( self, srm_client: AsyncMock, app_registration_repo: FakeAppRegistrationRepository, operation_repo: FakeOperationRepository, app_instance_repo: FakeAppInstanceRepository, ) -> None: service = EdgeApplicationManagementService( srm_client=srm_client, publisher=AsyncMock(), app_registration_repo=app_registration_repo, operation_repo=operation_repo, app_instance_repo=app_instance_repo, callback_registration_repo=None, ) with pytest.raises(RuntimeError, match="CallbackRegistrationRepository"): await service.create_app_instance( request=self._make_request_with_subscription(), tenant_id="t", app_provider_id="p", ) async def test_raises_when_operation_repo_unavailable( self, srm_client: AsyncMock, Loading Loading
src/open_exposure_gateway/application/services/edge_application_management_service.py +22 −0 Original line number Diff line number Diff line Loading @@ -43,6 +43,7 @@ from open_exposure_gateway.domain.models import ( AppInstanceState, AppRegistration, AppRegistrationStatus, CallbackRegistration, Operation, OperationStatus, OperationType, Loading Loading @@ -279,6 +280,27 @@ class EdgeApplicationManagementService: ) ) if request.subscriptionRequest is not None: if self._callback_registration_repo is None: raise RuntimeError("CallbackRegistrationRepository is not available") subscription = request.subscriptionRequest await self._callback_registration_repo.save( CallbackRegistration( id=uuid4(), operation_id=operation_id, tenant_id=tenant_id, api_family="edge-application-management", sink=subscription.sink, event_types=subscription.types, sink_credential_ref=f"secret://oeg/{operation_id}/sink-credential" if subscription.sinkCredential is not None else None, expires_at=subscription.config.subscriptionExpireTime if subscription.config else None, ) ) _command = build_deploy_command(translation, requested_at) await self._publish( Subject.TASK_DEPLOY, Loading
tests/unit/test_eam_flows.py +21 −0 Original line number Diff line number Diff line Loading @@ -30,6 +30,7 @@ from open_exposure_gateway.domain.edge_application_management import ( from open_exposure_gateway.domain.models import AppInstanceState, OperationStatus from tests.unit.fakes import ( FakeAppInstanceRepository, FakeCallbackRegistrationRepository, FakeDataBus, FakeMsg, FakeOperationRepository, Loading Loading @@ -116,6 +117,26 @@ class TestCreateAppInstanceFlow: assert stored_instance is not None assert stored_instance.state == AppInstanceState.INSTANTIATING def test_persists_callback_registration_when_subscription_request_present( self, api_client: TestClient, live_srm: FakeSRMClient, callback_registration_repo: FakeCallbackRegistrationRepository, ) -> None: body = { **CREATE_INSTANCE_BODY, "subscriptionRequest": { "sink": "https://client.example.com/callback", "types": [ "org.camaraproject.edge-application-management.v0.app-instance-status-change" ], }, } api_client.post(f"{EAM_BASE}/appinstances", json=body) (callback,) = list(callback_registration_repo.rows.values()) assert callback.sink == "https://client.example.com/callback" def test_srm_receives_a_valid_deploy_command( self, api_client: TestClient, fake_bus: FakeDataBus, live_srm: FakeSRMClient ) -> None: Loading
tests/unit/test_eam_service.py +98 −0 Original line number Diff line number Diff line from datetime import datetime, timezone from unittest.mock import AsyncMock from uuid import UUID, uuid4 Loading @@ -10,6 +11,8 @@ from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas i AppRepo, CreateAppInstanceRequest, KubernetesResources, SubscriptionConfig, SubscriptionRequest, VmResources, ) from open_exposure_gateway.application.services.edge_application_management_service import ( Loading Loading @@ -601,6 +604,101 @@ class TestCreateAppInstance: ) assert len(operation_repo.rows) == 2 def _make_request_with_subscription( self, with_credential: bool = False, expires_at: datetime | None = None ) -> CreateAppInstanceRequest: return CreateAppInstanceRequest( name="myapp_inst", appId=APP_ID, edgeCloudZoneId=ZONE_ID, subscriptionRequest=SubscriptionRequest( sink="https://client.example.com/callback", sinkCredential={"credentialType": "ACCESSTOKEN", "accessToken": "raw-token"} if with_credential else None, types=[ "org.camaraproject.edge-application-management.v0.app-instance-status-change" ], config=SubscriptionConfig(subscriptionExpireTime=expires_at) if expires_at else None, ), ) async def test_persists_callback_registration_when_subscription_present( self, service: EdgeApplicationManagementService, operation_repo: FakeOperationRepository, callback_registration_repo: FakeCallbackRegistrationRepository, ) -> None: expires_at = datetime(2027, 1, 17, 13, 18, 23, tzinfo=timezone.utc) await service.create_app_instance( request=self._make_request_with_subscription(expires_at=expires_at), tenant_id="tenant-1", app_provider_id="provider-1", ) (operation,) = list(operation_repo.rows.values()) callback = await callback_registration_repo.get_by_operation_id(operation.operation_id) assert callback is not None assert callback.tenant_id == "tenant-1" assert callback.api_family == "edge-application-management" assert callback.sink == "https://client.example.com/callback" assert callback.event_types == [ "org.camaraproject.edge-application-management.v0.app-instance-status-change" ] assert callback.is_active is True assert callback.expires_at == expires_at assert callback.sink_credential_ref is None async def test_sink_credential_is_stored_as_secret_reference_not_raw( self, service: EdgeApplicationManagementService, callback_registration_repo: FakeCallbackRegistrationRepository, ) -> None: await service.create_app_instance( request=self._make_request_with_subscription(with_credential=True), tenant_id="tenant-1", app_provider_id="provider-1", ) (callback,) = list(callback_registration_repo.rows.values()) assert callback.sink_credential_ref is not None assert callback.sink_credential_ref.startswith("secret://") assert "raw-token" not in callback.sink_credential_ref async def test_no_callback_registration_when_subscription_absent( self, service: EdgeApplicationManagementService, callback_registration_repo: FakeCallbackRegistrationRepository, ) -> None: await service.create_app_instance( request=self._make_request(), tenant_id="tenant-1", app_provider_id="provider-1", ) assert callback_registration_repo.rows == {} async def test_raises_when_callback_registration_repo_unavailable( self, srm_client: AsyncMock, app_registration_repo: FakeAppRegistrationRepository, operation_repo: FakeOperationRepository, app_instance_repo: FakeAppInstanceRepository, ) -> None: service = EdgeApplicationManagementService( srm_client=srm_client, publisher=AsyncMock(), app_registration_repo=app_registration_repo, operation_repo=operation_repo, app_instance_repo=app_instance_repo, callback_registration_repo=None, ) with pytest.raises(RuntimeError, match="CallbackRegistrationRepository"): await service.create_app_instance( request=self._make_request_with_subscription(), tenant_id="t", app_provider_id="p", ) async def test_raises_when_operation_repo_unavailable( self, srm_client: AsyncMock, Loading