Commit 0e2c35df authored by George Papathanail's avatar George Papathanail
Browse files

feat: add EdgeApplicationManagementService.create_app_deployment

parent 082fb9e0
Loading
Loading
Loading
Loading
+130 −0
Original line number Diff line number Diff line
@@ -7,10 +7,12 @@ from pydantic import BaseModel

from open_exposure_gateway.adapters.errors import DuplicateAppRegistrationError
from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas import (
    AppDeploymentId,
    AppInstanceInfo,
    AppInstanceStatus,
    AppManifest,
    AppManifestEnvelope,
    CreateAppDeploymentRequest,
    CreateAppInstanceRequest,
    EdgeCloudZone,
    KubernetesResources,
@@ -24,7 +26,9 @@ from open_exposure_gateway.application.mappers.edge_application_mapper import (
    build_app_registration_translation,
    build_catalog_payload,
    build_deploy_command,
    build_deploy_targets,
    build_edge_cloud_zone,
    build_multi_zone_deploy_command,
    build_submitted_app,
    build_terminate_instance_command,
    to_app_instance_status,
@@ -42,6 +46,8 @@ from open_exposure_gateway.domain.edge_application_management import (
    Subject,
)
from open_exposure_gateway.domain.models import (
    AppDeployment,
    AppDeploymentState,
    AppInstance,
    AppInstanceState,
    AppRegistration,
@@ -58,6 +64,7 @@ from open_exposure_gateway.ports.database.callbacks import (
    CallbackDeliveryRepository,
    CallbackRegistrationRepository,
)
from open_exposure_gateway.ports.database.deployments import AppDeploymentRepository
from open_exposure_gateway.ports.database.instances import AppInstanceRepository
from open_exposure_gateway.ports.database.operations import OperationRepository
from open_exposure_gateway.ports.database.registration import AppRegistrationRepository
@@ -117,6 +124,7 @@ class EdgeApplicationManagementService:
        callback_registration_repo: CallbackRegistrationRepository | None = None,
        callback_delivery_port: CallbackDeliveryPort | None = None,
        callback_delivery_repo: CallbackDeliveryRepository | None = None,
        app_deployment_repo: AppDeploymentRepository | None = None,
    ) -> None:
        self.srm_client = srm_client
        self._publisher = publisher
@@ -126,6 +134,7 @@ class EdgeApplicationManagementService:
        self._callback_registration_repo = callback_registration_repo
        self._callback_delivery_port = callback_delivery_port
        self._callback_delivery_repo = callback_delivery_repo
        self._app_deployment_repo = app_deployment_repo

    async def get_edge_cloud_zones(
        self,
@@ -409,6 +418,127 @@ class EdgeApplicationManagementService:
            edgeCloudZoneId=UUID(translation.zone_id),
        )

    async def create_app_deployment(
        self,
        request: CreateAppDeploymentRequest,
        tenant_id: str,
        app_provider_id: str,
        x_correlator: Optional[str] = None,
        idempotency_key: Optional[str] = None,
    ) -> AppDeploymentId:
        if self._app_registration_repo is None:
            raise RuntimeError("AppRegistrationRepository is not available")
        if (
            self._operation_repo is None
            or self._app_instance_repo is None
            or self._app_deployment_repo is None
        ):
            raise RuntimeError("Operation/AppInstance/AppDeployment repositories are not available")

        if idempotency_key is not None:
            existing_operation = await self._operation_repo.get_by_idempotency_key(
                tenant_id, idempotency_key
            )
            if existing_operation is not None:
                existing_deployment = await self._app_deployment_repo.get_by_operation_id(
                    existing_operation.operation_id
                )
                if existing_deployment is None:
                    raise RuntimeError(
                        "operations row found for idempotency_key but its app_deployments "
                        "row is missing"
                    )
                return AppDeploymentId(appDeploymentId=existing_deployment.app_deployment_id)

        app_registration = await self._app_registration_repo.get_by_app_id(request.appId)
        if app_registration is None:
            raise BadRequestException(message=f"App {request.appId} is not registered")

        operation_id, correlation_id, requested_at = self._new_operation_metadata(x_correlator)
        app_deployment_id = uuid4()
        app_instance_ids = [uuid4() for _ in request.edgeCloudZones]

        try:
            targets = build_deploy_targets(request, app_instance_ids)
        except ValueError as exc:
            raise BadRequestException(message=str(exc)) from exc

        await self._operation_repo.save(
            Operation(
                operation_id=operation_id,
                correlation_id=correlation_id,
                tenant_id=tenant_id,
                app_provider_id=app_provider_id,
                operation_type=OperationType.DEPLOY,
                status=OperationStatus.PENDING,
                subject=Subject.TASK_DEPLOY,
                idempotency_key=idempotency_key,
                app_registration_id=app_registration.app_registration_id,
                metadata={
                    "app_deployment_name": request.appDeploymentName,
                    "edge_cloud_zones": [str(zone_id) for zone_id in request.edgeCloudZones],
                },
            )
        )
        await self._app_deployment_repo.save(
            AppDeployment(
                app_deployment_id=app_deployment_id,
                operation_id=operation_id,
                app_registration_id=app_registration.app_registration_id,
                app_deployment_name=request.appDeploymentName,
                edge_cloud_zones=request.edgeCloudZones,
                state=AppDeploymentState.INSTANTIATING,
            )
        )
        for zone_id, app_instance_id in zip(request.edgeCloudZones, app_instance_ids, strict=True):
            await self._app_instance_repo.save(
                AppInstance(
                    app_instance_id=app_instance_id,
                    operation_id=operation_id,
                    app_registration_id=app_registration.app_registration_id,
                    edge_cloud_zone_id=zone_id,
                    state=AppInstanceState.INSTANTIATING,
                    app_deployment_id=app_deployment_id,
                )
            )

        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_multi_zone_deploy_command(
            app_registration_id=app_registration.app_registration_id,
            targets=targets,
            name=request.appDeploymentName,
            operation_id=operation_id,
            correlation_id=correlation_id,
            requested_at=requested_at,
            app_provider_id=app_provider_id,
        )
        await self._publish(
            Subject.TASK_DEPLOY,
            command,
            "Failed to publish app deployment task",
        )
        return AppDeploymentId(appDeploymentId=app_deployment_id)

    async def get_app_instances(
        self,
        app_id: Optional[UUID] = None,
+22 −0
Original line number Diff line number Diff line
@@ -50,6 +50,7 @@ from open_exposure_gateway.domain.location_retrieval import (
    SRMLocationResult,
)
from open_exposure_gateway.domain.models import (
    AppDeployment,
    AppInstance,
    AppInstanceState,
    AppRegistration,
@@ -71,6 +72,7 @@ from open_exposure_gateway.ports.database.callbacks import (
    CallbackDeliveryRepository,
    CallbackRegistrationRepository,
)
from open_exposure_gateway.ports.database.deployments import AppDeploymentRepository
from open_exposure_gateway.ports.database.instances import AppInstanceRepository
from open_exposure_gateway.ports.database.operations import OperationRepository
from open_exposure_gateway.ports.database.qod_sessions import QodSessionRepository
@@ -476,6 +478,26 @@ class FakeAppInstanceRepository(AppInstanceRepository):
        return stored.model_copy(deep=True)


class FakeAppDeploymentRepository(AppDeploymentRepository):
    def __init__(self) -> None:
        self.rows: dict[UUID, AppDeployment] = {}

    async def get_by_id(self, app_deployment_id: UUID) -> AppDeployment | None:
        found = self.rows.get(app_deployment_id)
        return found.model_copy(deep=True) if found is not None else None

    async def get_by_operation_id(self, operation_id: UUID) -> AppDeployment | None:
        for row in self.rows.values():
            if row.operation_id == operation_id:
                return row.model_copy(deep=True)
        return None

    async def save(self, app_deployment: AppDeployment) -> AppDeployment:
        stored = app_deployment.model_copy(deep=True)
        self.rows[stored.app_deployment_id] = stored
        return stored.model_copy(deep=True)


class FakeQodSessionRepository(QodSessionRepository):
    def __init__(self) -> None:
        self.rows: dict[UUID, QodSession] = {}
+197 −0
Original line number Diff line number Diff line
@@ -11,6 +11,7 @@ from open_exposure_gateway.api.camara.edge_application_management.vwip.schemas i
    ApplicationResources,
    AppManifest,
    AppRepo,
    CreateAppDeploymentRequest,
    CreateAppInstanceRequest,
    KubernetesResources,
    SubscriptionConfig,
@@ -49,6 +50,7 @@ from open_exposure_gateway.domain.edge_application_management import (
    Subject,
)
from open_exposure_gateway.domain.models import (
    AppDeploymentState,
    AppInstance,
    AppInstanceState,
    AppRegistration,
@@ -60,6 +62,7 @@ from open_exposure_gateway.domain.models import (
    PackageType,
)
from tests.unit.fakes import (
    FakeAppDeploymentRepository,
    FakeAppInstanceRepository,
    FakeAppRegistrationRepository,
    FakeCallbackDeliveryPort,
@@ -206,6 +209,11 @@ def app_instance_repo() -> FakeAppInstanceRepository:
    return FakeAppInstanceRepository()


@pytest.fixture()
def app_deployment_repo() -> FakeAppDeploymentRepository:
    return FakeAppDeploymentRepository()


@pytest.fixture()
def callback_registration_repo() -> FakeCallbackRegistrationRepository:
    return FakeCallbackRegistrationRepository()
@@ -231,6 +239,7 @@ def service(
    callback_registration_repo: FakeCallbackRegistrationRepository,
    callback_delivery_port: FakeCallbackDeliveryPort,
    callback_delivery_repo: FakeCallbackDeliveryRepository,
    app_deployment_repo: FakeAppDeploymentRepository,
) -> EdgeApplicationManagementService:
    return EdgeApplicationManagementService(
        srm_client=srm_client,
@@ -241,6 +250,7 @@ def service(
        callback_delivery_port=callback_delivery_port,
        callback_delivery_repo=callback_delivery_repo,
        callback_registration_repo=callback_registration_repo,
        app_deployment_repo=app_deployment_repo,
    )


@@ -1252,6 +1262,193 @@ class TestCreateAppInstanceUnregisteredApp:
            )


class TestCreateAppDeployment:
    ZONE_A = UUID("aaaaaaaa-1111-4000-8000-000000000001")
    ZONE_B = UUID("aaaaaaaa-2222-4000-8000-000000000002")

    def _make_request(
        self,
        zones: list[UUID] | None = None,
        kubernetes_cluster_refs: list[UUID] | None = None,
    ) -> CreateAppDeploymentRequest:
        return CreateAppDeploymentRequest(
            appDeploymentName="video_analytics_eu",
            appId=APP_ID,
            edgeCloudZones=zones or [self.ZONE_A, self.ZONE_B],
            kubernetesClusterRefs=kubernetes_cluster_refs,
        )

    @pytest.fixture(autouse=True)
    async def _registered_app(self, app_registration_repo: FakeAppRegistrationRepository) -> None:
        await app_registration_repo.save(
            AppRegistration(
                app_registration_id=APP_REGISTRATION_ID,
                app_id=APP_ID,
                tenant_id="tenant-1",
                name="myvideoapp",
                version="1.0.0",
                package_type=PackageType.HELM,
                status=AppRegistrationStatus.REGISTERED,
            )
        )

    async def test_returns_app_deployment_id(
        self, service: EdgeApplicationManagementService
    ) -> None:
        result = await service.create_app_deployment(
            request=self._make_request(),
            tenant_id="tenant-1",
            app_provider_id="provider-1",
        )
        assert result.appDeploymentId is not None

    async def test_publishes_one_deploy_command_with_n_targets(
        self, service: EdgeApplicationManagementService, publisher: AsyncMock
    ) -> None:
        await service.create_app_deployment(
            request=self._make_request(),
            tenant_id="tenant-1",
            app_provider_id="provider-1",
        )
        publisher.publish.assert_called_once()
        subject, payload = publisher.publish.call_args.args
        assert subject == Subject.TASK_DEPLOY
        assert len(payload["targets"]) == 2
        assert {t["zone_id"] for t in payload["targets"]} == {str(self.ZONE_A), str(self.ZONE_B)}

    async def test_persists_one_app_deployment_row(
        self,
        service: EdgeApplicationManagementService,
        app_deployment_repo: FakeAppDeploymentRepository,
    ) -> None:
        result = await service.create_app_deployment(
            request=self._make_request(),
            tenant_id="tenant-1",
            app_provider_id="provider-1",
        )
        assert len(app_deployment_repo.rows) == 1
        stored = await app_deployment_repo.get_by_id(result.appDeploymentId)
        assert stored is not None
        assert stored.app_deployment_name == "video_analytics_eu"
        assert stored.app_registration_id == APP_REGISTRATION_ID
        assert set(stored.edge_cloud_zones) == {self.ZONE_A, self.ZONE_B}
        assert stored.state == AppDeploymentState.INSTANTIATING

    async def test_persists_n_app_instance_rows_linked_to_deployment(
        self,
        service: EdgeApplicationManagementService,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        result = await service.create_app_deployment(
            request=self._make_request(),
            tenant_id="tenant-1",
            app_provider_id="provider-1",
        )
        instances = list(app_instance_repo.rows.values())
        assert len(instances) == 2
        assert {i.edge_cloud_zone_id for i in instances} == {self.ZONE_A, self.ZONE_B}
        for instance in instances:
            assert instance.app_deployment_id == result.appDeploymentId
            assert instance.state == AppInstanceState.INSTANTIATING

    async def test_persists_pending_deploy_operation_row(
        self,
        service: EdgeApplicationManagementService,
        operation_repo: FakeOperationRepository,
    ) -> None:
        await service.create_app_deployment(
            request=self._make_request(),
            tenant_id="tenant-1",
            app_provider_id="provider-1",
        )
        (operation,) = list(operation_repo.rows.values())
        assert operation.status == OperationStatus.PENDING
        assert operation.operation_type == OperationType.DEPLOY
        assert operation.subject == Subject.TASK_DEPLOY
        assert operation.app_registration_id == APP_REGISTRATION_ID

    async def test_retried_request_with_same_key_does_not_republish(
        self, service: EdgeApplicationManagementService, publisher: AsyncMock
    ) -> None:
        first = await service.create_app_deployment(
            request=self._make_request(),
            tenant_id="tenant-1",
            app_provider_id="provider-1",
            idempotency_key="retry-key-1",
        )
        second = await service.create_app_deployment(
            request=self._make_request(),
            tenant_id="tenant-1",
            app_provider_id="provider-1",
            idempotency_key="retry-key-1",
        )
        assert second.appDeploymentId == first.appDeploymentId
        publisher.publish.assert_called_once()

    async def test_retried_request_does_not_duplicate_rows(
        self,
        service: EdgeApplicationManagementService,
        operation_repo: FakeOperationRepository,
        app_deployment_repo: FakeAppDeploymentRepository,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        for _ in range(2):
            await service.create_app_deployment(
                request=self._make_request(),
                tenant_id="tenant-1",
                app_provider_id="provider-1",
                idempotency_key="retry-key-1",
            )
        assert len(operation_repo.rows) == 1
        assert len(app_deployment_repo.rows) == 1
        assert len(app_instance_repo.rows) == 2

    async def test_raises_when_cluster_refs_longer_than_zones(
        self, service: EdgeApplicationManagementService, publisher: AsyncMock
    ) -> None:
        request = self._make_request(
            zones=[self.ZONE_A], kubernetes_cluster_refs=[uuid4(), uuid4()]
        )
        with pytest.raises(BadRequestException):
            await service.create_app_deployment(
                request=request, tenant_id="tenant-1", app_provider_id="provider-1"
            )
        publisher.publish.assert_not_called()

    async def test_cluster_refs_longer_than_zones_does_not_persist_anything(
        self,
        service: EdgeApplicationManagementService,
        operation_repo: FakeOperationRepository,
        app_deployment_repo: FakeAppDeploymentRepository,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        request = self._make_request(
            zones=[self.ZONE_A], kubernetes_cluster_refs=[uuid4(), uuid4()]
        )
        with pytest.raises(BadRequestException):
            await service.create_app_deployment(
                request=request, tenant_id="tenant-1", app_provider_id="provider-1"
            )
        assert operation_repo.rows == {}
        assert app_deployment_repo.rows == {}
        assert app_instance_repo.rows == {}


class TestCreateAppDeploymentUnregisteredApp:
    async def test_raises_bad_request_when_app_not_registered(
        self, service: EdgeApplicationManagementService
    ) -> None:
        request = CreateAppDeploymentRequest(
            appDeploymentName="video_analytics_eu",
            appId=APP_ID,
            edgeCloudZones=[uuid4()],
        )
        with pytest.raises(BadRequestException, match=str(APP_ID)):
            await service.create_app_deployment(
                request=request, tenant_id="tenant-1", app_provider_id="provider-1"
            )


class TestGetAppInstances:
    async def test_returns_mapped_instances(
        self,