Commit 7082a526 authored by George Papathanail's avatar George Papathanail
Browse files

feat: update app_instances rows on operation completion

parent e51075f3
Loading
Loading
Loading
Loading
Loading
+33 −0
Original line number Diff line number Diff line
@@ -71,6 +71,11 @@ _OPERATION_COMPLETION_STATUS_MAP: dict[str, OperationStatus] = {
    "failed": OperationStatus.FAILED,
}

_APP_INSTANCE_COMPLETION_STATE_MAP: dict[str, AppInstanceState] = {
    "completed": AppInstanceState.READY,
    "failed": AppInstanceState.FAILED,
}


class EdgeApplicationManagementService:
    def __init__(
@@ -358,6 +363,8 @@ class EdgeApplicationManagementService:
    async def handle_completed(self, event: SRMOperationCompleted) -> None:
        if self._operation_repo is None:
            raise RuntimeError("OperationRepository is not available")
        if self._app_instance_repo is None:
            raise RuntimeError("AppInstanceRepository is not available")

        operation_id = UUID(event.operation_id)
        operation = await self._operation_repo.get_by_id(operation_id)
@@ -391,3 +398,29 @@ class EdgeApplicationManagementService:
            }
        )
        await self._operation_repo.save(updated)

        if event.instances:
            for instance in event.instances:
                app_instance_id = UUID(instance.service_instance_id)
                app_instance = await self._app_instance_repo.get_by_id(app_instance_id)
                if app_instance is None:
                    logger.warning(
                        "app_instance_completed_for_unknown_instance",
                        app_instance_id=instance.service_instance_id,
                    )
                    continue
                await self._app_instance_repo.save(
                    app_instance.model_copy(
                        update={"state": _APP_INSTANCE_COMPLETION_STATE_MAP[instance.status]}
                    )
                )
        elif status == OperationStatus.FAILED:
            # Total failure carries no instances[] entries to match against.
            # POST /appinstances always creates exactly one app_instances row
            # per operation (ADR-0005), so fall back to that link rather than
            # leaving the pre-created row stuck at instantiating forever.
            app_instance = await self._app_instance_repo.get_by_operation_id(operation_id)
            if app_instance is not None:
                await self._app_instance_repo.save(
                    app_instance.model_copy(update={"state": AppInstanceState.FAILED})
                )
+3 −2
Original line number Diff line number Diff line
@@ -52,14 +52,15 @@ def service_overrides() -> Generator[None, None, None]:
        )
    )
    operation_repo = FakeOperationRepository()
    wire_operation_consumer(bus, operation_repo)
    app_instance_repo = FakeAppInstanceRepository()
    wire_operation_consumer(bus, operation_repo, app_instance_repo)
    wire_srm_worker(bus, srm)
    app.dependency_overrides[get_edge_app_service] = lambda: EdgeApplicationManagementService(
        srm_client=srm,
        publisher=bus,
        app_registration_repo=FakeAppRegistrationRepository(),
        operation_repo=operation_repo,
        app_instance_repo=FakeAppInstanceRepository(),
        app_instance_repo=app_instance_repo,
        callback_registration_repo=FakeCallbackRegistrationRepository(),
    )
    app.dependency_overrides[get_qod_service] = lambda: QualityOnDemandService(srm_client=srm)
+2 −1
Original line number Diff line number Diff line
@@ -69,9 +69,10 @@ def live_srm(
    fake_bus: FakeDataBus,
    fake_srm: FakeSRMClient,
    operation_repo: FakeOperationRepository,
    app_instance_repo: FakeAppInstanceRepository,
) -> FakeSRMClient:
    """Fake SRM with its async side running: consumes commands, publishes completions."""
    wire_operation_consumer(fake_bus, operation_repo)
    wire_operation_consumer(fake_bus, operation_repo, app_instance_repo)
    wire_srm_worker(fake_bus, fake_srm)
    return fake_srm

+9 −4
Original line number Diff line number Diff line
@@ -241,13 +241,18 @@ def wire_srm_worker(bus: FakeDataBus, srm: FakeSRMClient) -> None:


def wire_operation_consumer(
    bus: FakeDataBus, operation_repo: OperationRepository
    bus: FakeDataBus,
    operation_repo: OperationRepository,
    app_instance_repo: AppInstanceRepository,
) -> NatsOperationConsumer:
    """Wires OEG's real completion handler (EdgeApplicationManagementService.handle_completed)
    behind the fake bus, backed by operation_repo -- pass the same instance used to build
    the service under test so a completion event updates the row the test can see."""
    behind the fake bus, backed by operation_repo/app_instance_repo -- pass the same instances
    used to build the service under test so a completion event updates the rows the test can
    see."""
    service = EdgeApplicationManagementService(
        srm_client=AsyncMock(), operation_repo=operation_repo
        srm_client=AsyncMock(),
        operation_repo=operation_repo,
        app_instance_repo=app_instance_repo,
    )
    consumer = NatsOperationConsumer(
        client=AsyncMock(), subject=Subject.OPERATION_COMPLETED, handler=service.handle_completed
+31 −4
Original line number Diff line number Diff line
@@ -212,6 +212,27 @@ class TestCreateAppInstanceFlow:
        assert len(instances) == 1
        assert instances[0]["status"] == "ready"

    def test_operations_and_app_instances_rows_reach_terminal_state_after_completion(
        self,
        api_client: TestClient,
        live_srm: FakeSRMClient,
        operation_repo: FakeOperationRepository,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        """Proves the UC7 completion path end-to-end through real DI wiring --
        not just GET /appinstances (which proxies SRM live), but OEG's own
        operations/app_instances rows, which is what a future retry or an
        internal read would actually see."""
        response = api_client.post(f"{EAM_BASE}/appinstances", json=CREATE_INSTANCE_BODY)
        instance_id = UUID(response.json()["appInstanceId"])

        (operation,) = list(operation_repo.rows.values())
        assert operation.status == OperationStatus.COMPLETED

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


class TestCreateAppInstanceUnregisteredAppFlow:
    """No autouse app registration here — this class exists to prove the
@@ -473,17 +494,23 @@ class TestEdgeCloudZonesFlow:

class TestOperationCompletedHandling:
    async def test_consumer_survives_malformed_event(
        self, fake_bus: FakeDataBus, operation_repo: FakeOperationRepository
        self,
        fake_bus: FakeDataBus,
        operation_repo: FakeOperationRepository,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        consumer = wire_operation_consumer(fake_bus, operation_repo)
        consumer = wire_operation_consumer(fake_bus, operation_repo, app_instance_repo)
        await consumer._handle_message(
            FakeMsg(data=b"not-json", subject=str(Subject.OPERATION_COMPLETED))
        )

    async def test_consumer_survives_event_missing_required_fields(
        self, fake_bus: FakeDataBus, operation_repo: FakeOperationRepository
        self,
        fake_bus: FakeDataBus,
        operation_repo: FakeOperationRepository,
        app_instance_repo: FakeAppInstanceRepository,
    ) -> None:
        consumer = wire_operation_consumer(fake_bus, operation_repo)
        consumer = wire_operation_consumer(fake_bus, operation_repo, app_instance_repo)
        await consumer._handle_message(
            FakeMsg(data=b'{"status": "completed"}', subject=str(Subject.OPERATION_COMPLETED))
        )
Loading