Commit 6a8b3f32 authored by George Papathanail's avatar George Papathanail
Browse files

feat: persist operations + app_instances before publish

parent 75dd3a6b
Loading
Loading
Loading
Loading
Loading
+35 −0
Original line number Diff line number Diff line
@@ -39,8 +39,13 @@ from open_exposure_gateway.domain.edge_application_management import (
    Subject,
)
from open_exposure_gateway.domain.models import (
    AppInstance,
    AppInstanceState,
    AppRegistration,
    AppRegistrationStatus,
    Operation,
    OperationStatus,
    OperationType,
    PackageType,
)
from open_exposure_gateway.ports.database.callbacks import CallbackRegistrationRepository
@@ -204,6 +209,8 @@ class EdgeApplicationManagementService:
    ) -> AppInstanceInfo:
        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:
            raise RuntimeError("Operation/AppInstance repositories are not available")
        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")
@@ -218,6 +225,34 @@ class EdgeApplicationManagementService:
            app_provider_id=app_provider_id,
            correlation_id=correlation_id,
        )

        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,
                app_registration_id=app_registration.app_registration_id,
                metadata={
                    "name": translation.name,
                    "resource_zone_id": translation.resource_zone_id,
                    "compute_domain_id": translation.compute_domain_id,
                },
            )
        )
        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=request.edgeCloudZoneId,
                state=AppInstanceState.INSTANTIATING,
            )
        )

        _command = build_deploy_command(translation, requested_at)
        await self._publish(
            Subject.TASK_DEPLOY,
+24 −1
Original line number Diff line number Diff line
@@ -12,7 +12,7 @@ real hole in the flow, not a broken test.
from collections.abc import Callable
from typing import Any
from unittest.mock import AsyncMock
from uuid import uuid4
from uuid import UUID, uuid4

import pytest
from fastapi.testclient import TestClient
@@ -27,9 +27,12 @@ from open_exposure_gateway.domain.edge_application_management import (
    SRMTerminateCommand,
    Subject,
)
from open_exposure_gateway.domain.models import AppInstanceState, OperationStatus
from tests.unit.fakes import (
    FakeAppInstanceRepository,
    FakeDataBus,
    FakeMsg,
    FakeOperationRepository,
    FakeSRMClient,
    wire_operation_consumer,
)
@@ -93,6 +96,26 @@ class TestCreateAppInstanceFlow:
        assert body["status"] == "instantiating"
        assert "appInstanceId" in body

    def test_persists_operation_and_app_instance_rows(
        self,
        api_client: TestClient,
        live_srm: FakeSRMClient,
        operation_repo: FakeOperationRepository,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        """Proves the DI wiring end-to-end: the router/service must reach the
        real repositories, not just the constructor accepting them."""
        response = api_client.post(f"{EAM_BASE}/appinstances", json=CREATE_INSTANCE_BODY)
        instance_id = response.json()["appInstanceId"]

        operations = list(operation_repo.rows.values())
        assert len(operations) == 1
        assert operations[0].status == OperationStatus.PENDING

        stored_instance = app_instance_repo.rows.get(UUID(instance_id))
        assert stored_instance is not None
        assert stored_instance.state == AppInstanceState.INSTANTIATING

    def test_srm_receives_a_valid_deploy_command(
        self, api_client: TestClient, fake_bus: FakeDataBus, live_srm: FakeSRMClient
    ) -> None:
+71 −2
Original line number Diff line number Diff line
@@ -38,8 +38,11 @@ from open_exposure_gateway.domain.edge_application_management import (
    Subject,
)
from open_exposure_gateway.domain.models import (
    AppInstanceState,
    AppRegistration,
    AppRegistrationStatus,
    OperationStatus,
    OperationType,
    PackageType,
)
from tests.unit.fakes import (
@@ -477,11 +480,77 @@ class TestCreateAppInstance:
        assert result.appId == APP_ID
        assert result.edgeCloudZoneId == ZONE_ID

    async def test_persists_pending_operation_row(
        self,
        service: EdgeApplicationManagementService,
        operation_repo: FakeOperationRepository,
        app_registration_repo: FakeAppRegistrationRepository,
    ) -> None:
        await service.create_app_instance(
            request=self._make_request(),
            tenant_id="tenant-1",
            app_provider_id="provider-1",
        )
        operations = list(operation_repo.rows.values())
        assert len(operations) == 1
        operation = operations[0]
        assert operation.status == OperationStatus.PENDING
        assert operation.operation_type == OperationType.DEPLOY
        assert operation.subject == Subject.TASK_DEPLOY
        assert operation.tenant_id == "tenant-1"
        assert operation.app_provider_id == "provider-1"
        registered = await app_registration_repo.get_by_app_id(APP_ID)
        assert registered is not None
        assert operation.app_registration_id == registered.app_registration_id
        assert operation.metadata["name"] == "myapp_inst"
        assert operation.metadata["resource_zone_id"] == str(ZONE_ID)

    async def test_persists_instantiating_app_instance_row(
        self,
        service: EdgeApplicationManagementService,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        result = await service.create_app_instance(
            request=self._make_request(),
            tenant_id="tenant-1",
            app_provider_id="provider-1",
        )
        stored = await app_instance_repo.get_by_id(result.appInstanceId)
        assert stored is not None
        assert stored.state == AppInstanceState.INSTANTIATING
        assert stored.edge_cloud_zone_id == ZONE_ID

    async def test_raises_when_operation_repo_unavailable(
        self,
        srm_client: AsyncMock,
        app_registration_repo: FakeAppRegistrationRepository,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        service = EdgeApplicationManagementService(
            srm_client=srm_client,
            publisher=AsyncMock(),
            app_registration_repo=app_registration_repo,
            operation_repo=None,
            app_instance_repo=app_instance_repo,
        )
        with pytest.raises(RuntimeError, match="Operation/AppInstance repositories"):
            await service.create_app_instance(
                request=self._make_request(), tenant_id="t", app_provider_id="p"
            )

    async def test_raises_when_publisher_unavailable(
        self, srm_client: AsyncMock, app_registration_repo: FakeAppRegistrationRepository
        self,
        srm_client: AsyncMock,
        app_registration_repo: FakeAppRegistrationRepository,
        operation_repo: FakeOperationRepository,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        service = EdgeApplicationManagementService(
            srm_client=srm_client, publisher=None, app_registration_repo=app_registration_repo
            srm_client=srm_client,
            publisher=None,
            app_registration_repo=app_registration_repo,
            operation_repo=operation_repo,
            app_instance_repo=app_instance_repo,
        )
        with pytest.raises(RuntimeError, match="DataBus publisher is not available"):
            await service.create_app_instance(