Commit 7942e015 authored by George Papathanail's avatar George Papathanail Committed by Dimitrios Gogos
Browse files

feat: translate SRM's live capability into TrafficInfluence GET responses

parent 7d9ce6f9
Loading
Loading
Loading
Loading
+45 −13
Original line number Diff line number Diff line
from typing import Any
from uuid import UUID

from pydantic import ValidationError

from open_exposure_gateway.api.camara.traffic_influence.vwip.schemas import PostTrafficInfluence
from open_exposure_gateway.api.camara.traffic_influence.vwip.schemas import (
    TrafficInfluence as TrafficInfluenceResponse,
@@ -14,6 +16,7 @@ from open_exposure_gateway.domain.traffic_influence import (
    NetworkCapabilityTargetApplicationReference,
    SRMNetworkCapabilityActivateCommand,
    SRMNetworkCapabilityDeactivateCommand,
    SRMTrafficInfluenceCapability,
    TrafficFilter,
    TrafficFilters,
)
@@ -86,18 +89,47 @@ def build_deactivate_command(
def build_traffic_influence_response(
    traffic_influence: TrafficInfluenceRecord,
    request_metadata: dict[str, Any],
    capability: SRMTrafficInfluenceCapability | None = None,
) -> TrafficInfluenceResponse:
    """Read-side counterpart of `build_activate_command`. SRM's capability_instance for a
    traffic-influence policy never carries apiConsumerId/appInstanceId/edgeCloudRegion/
    subscriptionRequest (they have no `srm.params/v1` field of their own), so the CAMARA
    response is reconstructed from the original request, stashed on
    `operations.metadata["request"]` at POST time, with the current `trafficInfluenceID`/
    `state` from OEG's own bookkeeping overlaid on top.
    """Read-side counterpart of `build_activate_command` -- TI's equivalent of
    quality_on_demand_mapper.build_session_info. `apiConsumerId`/`appInstanceId`/
    `edgeCloudRegion`/`subscriptionRequest` have no `srm.params/v1` field of their own, so
    they always come from the original request stashed on `operations.metadata["request"]`
    at POST time; `trafficInfluenceID`/`state` always come from OEG's own bookkeeping (the
    latter kept live by handle_completed/handle_status_changed, not by this call). When a
    live `capability` is available, its translated view of `appId`/`edgeCloudZoneId`/the
    traffic filters overrides the stored request for those fields -- most notably
    `edgeCloudZoneId`, which is the only field here SRM can know and OEG's own bookkeeping
    sometimes can't (a request pinned by `edgeCloudRegion` alone, or by neither, only gets
    a concrete zone once SRM realizes it -- srm/interface-contract.md §C.2).
    """
    return TrafficInfluenceResponse.model_validate(
        {
    response_data: dict[str, Any] = {
        **request_metadata,
        "trafficInfluenceID": traffic_influence.traffic_influence_id,
        "state": traffic_influence.state.value,
    }

    if capability is not None:
        response_data["appId"] = (
            capability.parameters_snapshot.target.application_reference.external_app_id
        )
        if capability.zone_id is not None:
            response_data["edgeCloudZoneId"] = capability.zone_id
        try:
            parameters = NetworkCapabilityParameters.model_validate(
                capability.parameters_snapshot.parameters
            )
        except ValidationError:
            parameters = None
        if parameters is not None:
            if parameters.traffic_filters.source is not None:
                response_data["sourceTrafficFilters"] = {
                    "sourcePort": parameters.traffic_filters.source.port,
                }
            if parameters.traffic_filters.destination is not None:
                response_data["destinationTrafficFilters"] = {
                    "destinationPort": parameters.traffic_filters.destination.port,
                    "destinationProtocol": parameters.traffic_filters.destination.protocol,
                }

    return TrafficInfluenceResponse.model_validate(response_data)
+39 −11
Original line number Diff line number Diff line
@@ -40,7 +40,7 @@ from open_exposure_gateway.domain.models import (
    TrafficInfluenceState as RecordTrafficInfluenceState,
)
from open_exposure_gateway.domain.quality_on_demand import SRMOperationStatus
from open_exposure_gateway.domain.traffic_influence import Subject
from open_exposure_gateway.domain.traffic_influence import SRMTrafficInfluenceCapability, Subject
from open_exposure_gateway.ports.database.callbacks import CallbackRegistrationRepository
from open_exposure_gateway.ports.database.operations import OperationRepository
from open_exposure_gateway.ports.database.traffic_influences import TrafficInfluenceRepository
@@ -406,6 +406,27 @@ class TrafficInfluenceService:
        operation = await self._operation_repo.get_by_id(operation_id)
        return operation.metadata.get("request", {}) if operation is not None else {}

    def _build_response_best_effort(
        self,
        traffic_influence: TrafficInfluenceRecord,
        request_metadata: dict[str, Any],
        capability: SRMTrafficInfluenceCapability | None,
    ) -> TrafficInfluenceResponse:
        # A malformed/unexpected capability shape must not break a GET that local
        # bookkeeping alone can already answer -- same "best-effort" rationale as the SRM
        # call itself, just covering the translation step too.
        if capability is not None:
            try:
                return build_traffic_influence_response(
                    traffic_influence, request_metadata, capability
                )
            except ValidationError:
                logger.warning(
                    "srm_traffic_influence_translation_failed",
                    traffic_influence_id=str(traffic_influence.traffic_influence_id),
                )
        return build_traffic_influence_response(traffic_influence, request_metadata)

    async def get_traffic_influence(
        self,
        traffic_influence_id: str,
@@ -420,14 +441,15 @@ class TrafficInfluenceService:
        if traffic_influence is None:
            raise NotFoundException(message=f"Traffic influence {traffic_influence_id} not found")

        # Best-effort: SRM's capability_instance for this policy carries no field that
        # CAMARA's TrafficInfluence needs (see build_traffic_influence_response), and its
        # read endpoint is still unconfirmed (srm_client.py), so a failure here must never
        # block a GET that OEG's own bookkeeping can already answer -- including the
        # 'deletion in progress'/'deleted' terminal states, where SRM may have already
        # purged the capability_instance entirely.
        # Best-effort: its read endpoint is still unconfirmed (srm_client.py), so a
        # failure here must never block a GET that OEG's own bookkeeping can already
        # answer -- including the 'deletion in progress'/'deleted' terminal states, where
        # SRM may have already purged the capability_instance entirely. When it does
        # succeed, build_traffic_influence_response translates it back to CAMARA shape
        # the same way build_session_info does for QoD.
        capability = None
        try:
            await self.srm_client.get_traffic_influence(
            capability = await self.srm_client.get_traffic_influence(
                traffic_influence_id=traffic_influence_id,
                x_correlator=x_correlator,
            )
@@ -437,7 +459,7 @@ class TrafficInfluenceService:
            )

        request_metadata = await self._fetch_request_metadata(traffic_influence.operation_id)
        return build_traffic_influence_response(traffic_influence, request_metadata)
        return self._build_response_best_effort(traffic_influence, request_metadata, capability)

    async def list_traffic_influences(
        self,
@@ -449,18 +471,24 @@ class TrafficInfluenceService:

        records = await self._traffic_influence_repo.list_by_app_id(app_id)

        # Best-effort, same rationale as get_traffic_influence. OEG mints
        # traffic_influence_id and sends it to SRM as service_instance_id at activation
        # (build_activate_command), so that's the join key back to each local record.
        capabilities_by_id: dict[str, SRMTrafficInfluenceCapability] = {}
        try:
            await self.srm_client.get_traffic_influences(
            capabilities = await self.srm_client.get_traffic_influences(
                app_id=app_id,
                x_correlator=x_correlator,
            )
            capabilities_by_id = {c.service_instance_id: c for c in capabilities}
        except (NotFoundException, DownstreamServiceException, ValidationError):
            logger.warning("srm_traffic_influence_list_read_failed", app_id=str(app_id))

        responses = []
        for record in records:
            request_metadata = await self._fetch_request_metadata(record.operation_id)
            responses.append(build_traffic_influence_response(record, request_metadata))
            capability = capabilities_by_id.get(str(record.traffic_influence_id))
            responses.append(self._build_response_best_effort(record, request_metadata, capability))
        return responses

    async def delete_traffic_influence(
+92 −0
Original line number Diff line number Diff line
@@ -20,8 +20,12 @@ from open_exposure_gateway.domain.models import (
    TrafficInfluenceState,
)
from open_exposure_gateway.domain.traffic_influence import (
    NetworkCapabilityParametersSnapshot,
    NetworkCapabilityTarget,
    NetworkCapabilityTargetApplicationReference,
    SRMNetworkCapabilityActivateCommand,
    SRMNetworkCapabilityDeactivateCommand,
    SRMTrafficInfluenceCapability,
    Subject,
)
from tests.unit.fakes import (
@@ -258,6 +262,94 @@ class TestTrafficInfluenceGetFlow:
        response = api_client.get(f"{TI_BASE}/traffic-influences/not-a-uuid")
        assert response.status_code == 400

    def test_get_translates_live_srm_capability_into_response(
        self, api_client: TestClient, fake_srm: FakeSRMClient
    ) -> None:
        """When SRM has a live view, it's translated back to CAMARA shape (same idea as
        build_session_info for QoD) -- most importantly edgeCloudZoneId, which is the one
        field only SRM knows for sure once a request pinned by region (not zone) is
        actually realized."""
        body = {
            "apiConsumerId": TI_BODY["apiConsumerId"],
            "appId": TI_BODY["appId"],
            "edgeCloudRegion": "region-1",
            "sourceTrafficFilters": TI_BODY["sourceTrafficFilters"],
            "destinationTrafficFilters": TI_BODY["destinationTrafficFilters"],
        }
        traffic_influence_id = api_client.post(f"{TI_BASE}/traffic-influences", json=body).json()[
            "trafficInfluenceID"
        ]

        fake_srm.traffic_influences[traffic_influence_id] = SRMTrafficInfluenceCapability(
            service_instance_id=traffic_influence_id,
            capability_type="traffic_influence",
            state="active",
            zone_id="9f8d7c6b-5a4e-3d2c-1b0a-9f8d7c6b5a4e",
            app_provider_id="provider-1",
            external_ref=traffic_influence_id,
            parameters_snapshot=NetworkCapabilityParametersSnapshot(
                target=NetworkCapabilityTarget(
                    application_reference=NetworkCapabilityTargetApplicationReference(
                        external_app_id=TI_BODY["appId"]
                    )
                ),
                parameters={
                    "traffic_filters": {
                        "source": {"port": 45678},
                        "destination": {"port": 443, "protocol": "TCP"},
                    }
                },
            ),
        )

        response = api_client.get(f"{TI_BASE}/traffic-influences/{traffic_influence_id}")

        assert response.status_code == 200
        assert response.json()["edgeCloudZoneId"] == "9f8d7c6b-5a4e-3d2c-1b0a-9f8d7c6b5a4e"

    def test_get_falls_back_to_local_data_when_srm_read_fails(self, api_client: TestClient) -> None:
        """FakeSRMClient 404s get_traffic_influence when nothing was seeded -- GET must
        still succeed from local bookkeeping alone."""
        traffic_influence_id = api_client.post(
            f"{TI_BASE}/traffic-influences", json=TI_BODY
        ).json()["trafficInfluenceID"]

        response = api_client.get(f"{TI_BASE}/traffic-influences/{traffic_influence_id}")

        assert response.status_code == 200
        assert response.json()["edgeCloudZoneId"] == TI_BODY["edgeCloudZoneId"]

    def test_get_falls_back_to_local_data_when_srm_capability_is_malformed(
        self, api_client: TestClient, fake_srm: FakeSRMClient
    ) -> None:
        """A capability that fails to translate to CAMARA shape (e.g. a non-UUID zone_id)
        must not turn a working GET into a 500 -- fall back to local bookkeeping instead."""
        traffic_influence_id = api_client.post(
            f"{TI_BASE}/traffic-influences", json=TI_BODY
        ).json()["trafficInfluenceID"]

        fake_srm.traffic_influences[traffic_influence_id] = SRMTrafficInfluenceCapability(
            service_instance_id=traffic_influence_id,
            capability_type="traffic_influence",
            state="active",
            zone_id="not-a-uuid",
            app_provider_id="provider-1",
            external_ref=traffic_influence_id,
            parameters_snapshot=NetworkCapabilityParametersSnapshot(
                target=NetworkCapabilityTarget(
                    application_reference=NetworkCapabilityTargetApplicationReference(
                        external_app_id=TI_BODY["appId"]
                    )
                ),
                parameters={},
            ),
        )

        response = api_client.get(f"{TI_BASE}/traffic-influences/{traffic_influence_id}")

        assert response.status_code == 200
        assert response.json()["edgeCloudZoneId"] == TI_BODY["edgeCloudZoneId"]

    def test_list_traffic_influences_returns_empty_by_default(self, api_client: TestClient) -> None:
        response = api_client.get(f"{TI_BASE}/traffic-influences")
        assert response.status_code == 200