Commit ff67d322 authored by George Papathanail's avatar George Papathanail
Browse files

feat(traffic-influence): add TrafficInfluenceService

create_traffic_influence: resolves edgeCloudZoneId/edgeCloudRegion to a
resource_zone_id (region has no SRM field of its own -- resolved locally
via the same zones lookup EAM uses), uses service_specification_id =
appId (ADR-0011: this holds for an already-onboarded app, unlike QoD's
synthetic id), persists Operation + TrafficInfluence rows, persists
subscriptionRequest into the generic CallbackRegistration table, then
publishes the activate command and returns an ORDERED response.

handle_completed: branches on TRAFFIC_INFLUENCE vs
TRAFFIC_INFLUENCE_DEACTIVATE. The deactivate completion can't be looked
up by operation_id (the deactivate command mints its own, distinct from
the row's stored operation_id) -- recovered instead from
traffic_influence_id stashed in the deactivate operation's own metadata,
so both activate and deactivate outcomes update the row's state
(active/error/deleted), unlike QoD's handle_completed which only ever
resolves the activate side.

get/list_traffic_influence proxy straight to the SRM client (open
contract gap, see next commits). delete_traffic_influence mirrors QoD's
404-before-external_ref guard and additionally marks the local row
"deletion in progress" per the CAMARA state enum.
parent d9302372
Loading
Loading
Loading
Loading
+337 −0
Original line number Diff line number Diff line
from datetime import datetime, timezone
from typing import Optional
from uuid import UUID, uuid4

import structlog
from pydantic import BaseModel

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,
)
from open_exposure_gateway.api.camara.traffic_influence.vwip.schemas import (
    TrafficInfluenceState as ApiTrafficInfluenceState,
)
from open_exposure_gateway.application.mappers.traffic_influence_mapper import (
    build_activate_command,
    build_deactivate_command,
)
from open_exposure_gateway.core.exceptions import DownstreamServiceException, NotFoundException
from open_exposure_gateway.domain.models import (
    CallbackRegistration,
    Operation,
    OperationStatus,
    OperationType,
)
from open_exposure_gateway.domain.models import (
    TrafficInfluence as TrafficInfluenceRecord,
)
from open_exposure_gateway.domain.models import (
    TrafficInfluenceState as RecordTrafficInfluenceState,
)
from open_exposure_gateway.domain.srm_events import SRMOperationCompleted
from open_exposure_gateway.domain.traffic_influence import 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
from open_exposure_gateway.ports.databus_port import DataBusPort
from open_exposure_gateway.ports.srm_port import SRMClientPort

logger: structlog.BoundLogger = structlog.get_logger(__name__)

_OPERATION_COMPLETION_STATUS_MAP: dict[str, OperationStatus] = {
    "completed": OperationStatus.COMPLETED,
    "partially_completed": OperationStatus.PARTIALLY_COMPLETED,
    "failed": OperationStatus.FAILED,
}

# CallbackRegistration.api_family discriminator for Traffic Influence's subscriptionRequest.
_API_FAMILY = "traffic_influence"


class TrafficInfluenceService:
    def __init__(
        self,
        srm_client: SRMClientPort,
        publisher: Optional[DataBusPort] = None,
        operation_repo: Optional[OperationRepository] = None,
        traffic_influence_repo: Optional[TrafficInfluenceRepository] = None,
        callback_registration_repo: Optional[CallbackRegistrationRepository] = None,
    ) -> None:
        self.srm_client = srm_client
        self._publisher = publisher
        self._operation_repo = operation_repo
        self._traffic_influence_repo = traffic_influence_repo
        self._callback_registration_repo = callback_registration_repo

    def _new_operation_metadata(self, x_correlator: Optional[str]) -> tuple[UUID, str, str]:
        operation_id = uuid4()
        correlation_id = x_correlator or str(uuid4())
        requested_at = datetime.now(timezone.utc).isoformat()
        return operation_id, correlation_id, requested_at

    async def _publish(self, subject: Subject, command: BaseModel, error_msg: str) -> None:
        if self._publisher is None:
            raise RuntimeError("DataBus publisher is not available")
        try:
            await self._publisher.publish(subject, command.model_dump(mode="json"))
        except Exception as exc:
            raise DownstreamServiceException(error_msg) from exc

    async def _resolve_resource_zone_id(
        self,
        request: PostTrafficInfluence,
        x_correlator: Optional[str],
    ) -> str | None:
        if request.edgeCloudZoneId is not None:
            return str(request.edgeCloudZoneId)
        if request.edgeCloudRegion is not None:
            zones = await self.srm_client.get_resource_zones(
                region=request.edgeCloudRegion,
                status=None,
                x_correlator=x_correlator,
            )
            if zones:
                return zones[0].resource_zone_id
        return None

    async def create_traffic_influence(
        self,
        request: PostTrafficInfluence,
        tenant_id: str,
        app_provider_id: str,
        x_correlator: Optional[str] = None,
    ) -> TrafficInfluenceResponse:
        if self._operation_repo is None or self._traffic_influence_repo is None:
            raise RuntimeError("Operation/TrafficInfluence repositories are not available")

        operation_id, correlation_id, requested_at = self._new_operation_metadata(x_correlator)
        traffic_influence_id = uuid4()
        # edgeCloudRegion has no SRM field of its own (srm/canonical-parameters-schema.md
        # has no "region" concept); it's resolved locally to a zone before the command
        # is built, same zones lookup EAM already uses.
        resource_zone_id = await self._resolve_resource_zone_id(request, x_correlator)

        command = build_activate_command(
            request=request,
            operation_id=operation_id,
            correlation_id=correlation_id,
            requested_at=requested_at,
            # service_specification_id == appId for an already-onboarded app
            # (srm/interface-contract.md §E.4, ADR-0011) -- unlike QoD, which mints a
            # throwaway id because its sessions aren't tied to an onboarded catalog entry.
            service_specification_id=request.appId,
            app_provider_id=app_provider_id,
            resource_zone_id=resource_zone_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.TRAFFIC_INFLUENCE,
                status=OperationStatus.PENDING,
                subject=Subject.TASK_ACTIVATE,
                metadata={"app_id": str(request.appId)},
            )
        )

        source_filter = request.sourceTrafficFilters
        destination_filter = request.destinationTrafficFilters
        await self._traffic_influence_repo.save(
            TrafficInfluenceRecord(
                traffic_influence_id=traffic_influence_id,
                operation_id=operation_id,
                app_id=request.appId,
                edge_cloud_zone_id=request.edgeCloudZoneId,
                edge_cloud_region=request.edgeCloudRegion,
                source_port=source_filter.sourcePort if source_filter else None,
                destination_port=destination_filter.destinationPort if destination_filter else None,
                destination_protocol=(
                    destination_filter.destinationProtocol if destination_filter else None
                ),
                state=RecordTrafficInfluenceState.ORDERED,
            )
        )

        # subscriptionRequest is never sent to SRM (srm/canonical-parameters-schema.md:
        # "NOT stored -- OEG owns subscription/callback state"); it's persisted here into
        # the generic callback table instead. Delivery of onTrafficInfluenceChanged is a
        # separate, not-yet-built concern -- this only records the registration.
        if request.subscriptionRequest is not None and self._callback_registration_repo is not None:
            await self._callback_registration_repo.save(
                CallbackRegistration(
                    id=uuid4(),
                    operation_id=operation_id,
                    tenant_id=tenant_id,
                    api_family=_API_FAMILY,
                    sink=request.subscriptionRequest.sink,
                    event_types=list(request.subscriptionRequest.types),
                    expires_at=request.subscriptionRequest.config.subscriptionExpireTime,
                )
            )

        await self._publish(
            Subject.TASK_ACTIVATE,
            command,
            "Failed to publish traffic influence activation",
        )

        return TrafficInfluenceResponse(
            trafficInfluenceID=traffic_influence_id,
            state=ApiTrafficInfluenceState.ORDERED,
            apiConsumerId=request.apiConsumerId,
            appId=request.appId,
            appInstanceId=request.appInstanceId,
            edgeCloudRegion=request.edgeCloudRegion,
            edgeCloudZoneId=request.edgeCloudZoneId,
            sourceTrafficFilters=request.sourceTrafficFilters,
            destinationTrafficFilters=request.destinationTrafficFilters,
            subscriptionRequest=request.subscriptionRequest,
        )

    async def handle_completed(self, event: SRMOperationCompleted) -> None:
        if self._operation_repo is None or self._traffic_influence_repo is None:
            raise RuntimeError("Operation/TrafficInfluence repositories are not available")

        operation_id = UUID(event.operation_id)
        operation = await self._operation_repo.get_by_id(operation_id)
        if operation is None:
            logger.warning(
                "operation_completed_for_unknown_operation", operation_id=event.operation_id
            )
            return

        # Not ours: event.srm.operation.completed is shared across every domain.
        if operation.operation_type not in (
            OperationType.TRAFFIC_INFLUENCE,
            OperationType.TRAFFIC_INFLUENCE_DEACTIVATE,
        ):
            return

        status = _OPERATION_COMPLETION_STATUS_MAP[event.status]
        result = None
        if status != OperationStatus.FAILED:
            result = {"instances": [i.model_dump(mode="json") for i in event.instances]}
        await self._operation_repo.save(
            operation.model_copy(
                update={
                    "status": status,
                    "result": result,
                    "error": event.error,
                    "completed_at": datetime.fromisoformat(event.completed_at),
                }
            )
        )

        instance = event.instances[0] if event.instances else None
        succeeded = instance is not None and instance.status == "completed"

        if operation.operation_type == OperationType.TRAFFIC_INFLUENCE:
            traffic_influence = await self._traffic_influence_repo.get_by_operation_id(operation_id)
            if traffic_influence is None:
                logger.warning(
                    "traffic_influence_completed_for_unknown_operation",
                    operation_id=event.operation_id,
                )
                return
            update = (
                {"state": RecordTrafficInfluenceState.ACTIVE, "external_ref": instance.external_ref}
                if succeeded and instance is not None
                else {"state": RecordTrafficInfluenceState.ERROR}
            )
            await self._traffic_influence_repo.save(traffic_influence.model_copy(update=update))
            return

        # TRAFFIC_INFLUENCE_DEACTIVATE: the deactivate command mints its own
        # operation_id (distinct from the activate one the row is keyed by), so the
        # deactivated row is recovered from the operation's own metadata instead.
        raw_id = operation.metadata.get("traffic_influence_id")
        if raw_id is None:
            logger.warning(
                "deactivate_completed_without_traffic_influence_id",
                operation_id=event.operation_id,
            )
            return
        traffic_influence = await self._traffic_influence_repo.get_by_id(UUID(raw_id))
        if traffic_influence is None:
            return
        deactivate_update = {
            "state": RecordTrafficInfluenceState.DELETED
            if succeeded
            else RecordTrafficInfluenceState.ERROR
        }
        await self._traffic_influence_repo.save(
            traffic_influence.model_copy(update=deactivate_update)
        )

    async def get_traffic_influence(
        self,
        traffic_influence_id: str,
        x_correlator: Optional[str] = None,
    ) -> TrafficInfluenceResponse:
        return await self.srm_client.get_traffic_influence(
            traffic_influence_id=traffic_influence_id,
            x_correlator=x_correlator,
        )

    async def list_traffic_influences(
        self,
        app_id: Optional[UUID] = None,
        x_correlator: Optional[str] = None,
    ) -> list[TrafficInfluenceResponse]:
        return await self.srm_client.get_traffic_influences(
            app_id=app_id,
            x_correlator=x_correlator,
        )

    async def delete_traffic_influence(
        self,
        traffic_influence_id: str,
        tenant_id: str,
        app_provider_id: str,
        x_correlator: Optional[str] = None,
    ) -> None:
        if self._operation_repo is None or self._traffic_influence_repo is None:
            raise RuntimeError("Operation/TrafficInfluence repositories are not available")

        traffic_influence = await self._traffic_influence_repo.get_by_id(UUID(traffic_influence_id))
        if traffic_influence is None or traffic_influence.external_ref is None:
            raise NotFoundException(message=f"Traffic influence {traffic_influence_id} not found")

        operation_id, correlation_id, requested_at = self._new_operation_metadata(x_correlator)
        command = build_deactivate_command(
            external_ref=traffic_influence.external_ref,
            operation_id=operation_id,
            correlation_id=correlation_id,
            requested_at=requested_at,
            app_provider_id=app_provider_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.TRAFFIC_INFLUENCE_DEACTIVATE,
                status=OperationStatus.PENDING,
                subject=Subject.TASK_DEACTIVATE,
                metadata={"traffic_influence_id": str(traffic_influence_id)},
            )
        )
        await self._traffic_influence_repo.save(
            traffic_influence.model_copy(
                update={"state": RecordTrafficInfluenceState.DELETION_IN_PROGRESS}
            )
        )

        await self._publish(
            Subject.TASK_DEACTIVATE,
            command,
            "Failed to publish traffic influence deactivation",
        )