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

feat: implement onTrafficinfluencechange webhook delivery

parent 76b109eb
Loading
Loading
Loading
Loading
+19 −0
Original line number Diff line number Diff line
import httpx

from open_exposure_gateway.core.config import get_settings
from open_exposure_gateway.domain.traffic_influence import TrafficInfluenceChangedCloudEvent


class HttpTrafficInfluenceCallbackClient:
    def __init__(self) -> None:
        settings = get_settings()
        self.timeout = settings.callback_settings.timeout

    async def deliver(self, sink: str, event: TrafficInfluenceChangedCloudEvent) -> None:
        async with httpx.AsyncClient(timeout=self.timeout) as client:
            response = await client.post(
                sink,
                content=event.model_dump_json(),
                headers={"Content-Type": "application/cloudevents+json"},
            )
        response.raise_for_status()
+22 −1
Original line number Diff line number Diff line
from typing import Any
from uuid import UUID
from uuid import UUID, uuid4

from pydantic import ValidationError

@@ -19,6 +19,7 @@ from open_exposure_gateway.domain.traffic_influence import (
    SRMTrafficInfluenceCapability,
    TrafficFilter,
    TrafficFilters,
    TrafficInfluenceChangedCloudEvent,
)


@@ -133,3 +134,23 @@ def build_traffic_influence_response(
                }

    return TrafficInfluenceResponse.model_validate(response_data)


def build_traffic_influence_changed_event(
    traffic_influence: TrafficInfluenceRecord,
    request_metadata: dict[str, Any],
    source: str,
    occurred_at: str,
) -> TrafficInfluenceChangedCloudEvent:
    """`onTrafficInfluenceChanged` -- TI's equivalent of
    quality_on_demand_mapper.build_qos_status_changed_event. `data` is the same resource
    representation `build_traffic_influence_response` returns (the notification schema is
    `allOf: [PostTrafficInfluenceDevice]`, i.e. the whole resource, not a delta like QoD's).
    """
    response = build_traffic_influence_response(traffic_influence, request_metadata)
    return TrafficInfluenceChangedCloudEvent(
        id=uuid4(),
        source=source,
        time=occurred_at,
        data=response.model_dump(mode="json"),
    )
+99 −26
Original line number Diff line number Diff line
@@ -18,8 +18,10 @@ from open_exposure_gateway.api.camara.traffic_influence.vwip.schemas import (
from open_exposure_gateway.application.mappers.traffic_influence_mapper import (
    build_activate_command,
    build_deactivate_command,
    build_traffic_influence_changed_event,
    build_traffic_influence_response,
)
from open_exposure_gateway.core.config import get_settings
from open_exposure_gateway.core.exceptions import (
    BadRequestException,
    ConflictException,
@@ -28,6 +30,7 @@ from open_exposure_gateway.core.exceptions import (
)
from open_exposure_gateway.domain.edge_application_management import SRMOperationCompleted
from open_exposure_gateway.domain.models import (
    CallbackDelivery,
    CallbackRegistration,
    Operation,
    OperationStatus,
@@ -41,11 +44,17 @@ from open_exposure_gateway.domain.models import (
)
from open_exposure_gateway.domain.quality_on_demand import SRMOperationStatus
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.callbacks import (
    CallbackDeliveryRepository,
    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
from open_exposure_gateway.ports.traffic_influence_callback_port import (
    TrafficInfluenceCallbackDeliveryPort,
)

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

@@ -55,14 +64,6 @@ _OPERATION_COMPLETION_STATUS_MAP: dict[str, OperationStatus] = {
    "failed": OperationStatus.FAILED,
}

_TERMINAL_OPERATION_STATUSES = frozenset(
    {
        OperationStatus.COMPLETED,
        OperationStatus.PARTIALLY_COMPLETED,
        OperationStatus.FAILED,
    }
)

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

@@ -86,12 +87,16 @@ class TrafficInfluenceService:
        operation_repo: Optional[OperationRepository] = None,
        traffic_influence_repo: Optional[TrafficInfluenceRepository] = None,
        callback_registration_repo: Optional[CallbackRegistrationRepository] = None,
        callback_delivery_repo: Optional[CallbackDeliveryRepository] = None,
        callback_delivery_port: Optional[TrafficInfluenceCallbackDeliveryPort] = 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
        self._callback_delivery_repo = callback_delivery_repo
        self._callback_delivery_port = callback_delivery_port

    def _parse_traffic_influence_id(self, traffic_influence_id: str) -> UUID:
        try:
@@ -246,8 +251,8 @@ class TrafficInfluenceService:

        # 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.
        # the generic callback table instead. Delivery happens later, keyed on this
        # operation_id, from handle_completed/handle_status_changed/delete_traffic_influence.
        if request.subscriptionRequest is not None and self._callback_registration_repo is not None:
            await self._callback_registration_repo.save(
                CallbackRegistration(
@@ -285,6 +290,59 @@ class TrafficInfluenceService:
            subscriptionRequest=request.subscriptionRequest,
        )

    async def _deliver_traffic_influence_changed(
        self,
        registration_operation_id: UUID,
        traffic_influence: TrafficInfluenceRecord,
        occurred_at: str,
    ) -> None:
        if self._callback_registration_repo is None:
            return
        registration = await self._callback_registration_repo.get_by_operation_id(
            registration_operation_id
        )
        if registration is None or not registration.is_active:
            return
        if self._callback_delivery_port is None or self._callback_delivery_repo is None:
            raise RuntimeError("CallbackDeliveryPort/CallbackDeliveryRepository is not available")

        source = (
            f"{get_settings().public_base_url}/traffic-influence/vwip/traffic-influences/"
            f"{traffic_influence.traffic_influence_id}"
        )
        request_metadata = await self._fetch_request_metadata(registration_operation_id)
        cloud_event = build_traffic_influence_changed_event(
            traffic_influence=traffic_influence,
            request_metadata=request_metadata,
            source=source,
            occurred_at=occurred_at,
        )

        prior_attempts = await self._callback_delivery_repo.list_by_callback_registration_id(
            registration.id
        )
        last_error: Optional[str] = None
        try:
            await self._callback_delivery_port.deliver(registration.sink, cloud_event)
            state = "delivered"
        except Exception as exc:
            state = "failed"
            last_error = str(exc)
            logger.warning(
                "traffic_influence_callback_delivery_failed", sink=registration.sink, error=str(exc)
            )

        await self._callback_delivery_repo.save(
            CallbackDelivery(
                id=uuid4(),
                callback_registration_id=registration.id,
                operation_id=registration_operation_id,
                attempt=len(prior_attempts) + 1,
                state=state,
                last_error=last_error,
            )
        )

    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")
@@ -304,14 +362,6 @@ class TrafficInfluenceService:
        ):
            return

        if operation.status in _TERMINAL_OPERATION_STATUSES:
            logger.info(
                "redelivered_completion_ignored_for_terminal_operation",
                operation_id=event.operation_id,
                status=operation.status.value,
            )
            return

        status = _OPERATION_COMPLETION_STATUS_MAP[event.status]
        result = None
        if status != OperationStatus.FAILED:
@@ -343,7 +393,13 @@ class TrafficInfluenceService:
                if succeeded and instance is not None
                else {"state": RecordTrafficInfluenceState.ERROR}
            )
            await self._traffic_influence_repo.save(traffic_influence.model_copy(update=update))
            prior_state = traffic_influence.state
            saved = await self._traffic_influence_repo.save(
                traffic_influence.model_copy(update=update)
            )
            if prior_state == saved.state:
                return
            await self._deliver_traffic_influence_changed(operation_id, saved, event.completed_at)
            return

        # TRAFFIC_INFLUENCE_DEACTIVATE: the deactivate command mints its own
@@ -364,9 +420,19 @@ class TrafficInfluenceService:
            if succeeded
            else RecordTrafficInfluenceState.ERROR
        }
        await self._traffic_influence_repo.save(
        prior_state = traffic_influence.state
        saved = await self._traffic_influence_repo.save(
            traffic_influence.model_copy(update=deactivate_update)
        )
        if prior_state == saved.state:
            return
        # Registrations are keyed on the *create* operation_id, not this deactivate
        # operation's own id (same anchor QoD's delete_session notification uses) --
        # saved.operation_id is that original id, preserved unchanged through deactivation.
        # No QoD equivalent: QoD's qod_sessions.state has no terminal "deleted" value to
        # report (delete_session already notifies synchronously going UNAVAILABLE), but
        # TI's CAMARA spec requires 'deleted' itself to be an observable state.
        await self._deliver_traffic_influence_changed(saved.operation_id, saved, event.completed_at)

    async def handle_status_changed(self, event: SRMOperationStatus) -> None:
        if self._operation_repo is None or self._traffic_influence_repo is None:
@@ -399,12 +465,13 @@ class TrafficInfluenceService:
            if qos_status == "AVAILABLE"
            else RecordTrafficInfluenceState.ERROR
        )
        await self._traffic_influence_repo.save(
        prior_state = traffic_influence.state
        saved = await self._traffic_influence_repo.save(
            traffic_influence.model_copy(update={"state": state})
        )

        # Delivery of onTrafficInfluenceChanged is a separate, not-yet-built concern
        # (same scoping as the subscriptionRequest note in create_traffic_influence).
        if prior_state == saved.state:
            return
        await self._deliver_traffic_influence_changed(operation_id, saved, event.emitted_at)

    async def _fetch_request_metadata(self, operation_id: UUID) -> dict[str, Any]:
        if self._operation_repo is None:
@@ -541,8 +608,14 @@ class TrafficInfluenceService:
            "Failed to publish traffic influence deactivation",
        )

        await self._traffic_influence_repo.save(
        was_active = traffic_influence.state == RecordTrafficInfluenceState.ACTIVE
        saved = await self._traffic_influence_repo.save(
            traffic_influence.model_copy(
                update={"state": RecordTrafficInfluenceState.DELETION_IN_PROGRESS}
            )
        )
        # Same condition QoD's delete_session uses to skip notifying (via qod_session.state
        # == AVAILABLE): don't tell subscribers something's "going away" if it was never
        # reported as up in the first place (still ORDERED/CREATED when DELETE arrived).
        if was_active:
            await self._deliver_traffic_influence_changed(saved.operation_id, saved, requested_at)
+4 −0
Original line number Diff line number Diff line
@@ -5,6 +5,9 @@ from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession, async_sessionmaker
from open_exposure_gateway.ports.databus_port import DataBusPort
from open_exposure_gateway.ports.qod_callback_port import QodCallbackDeliveryPort
from open_exposure_gateway.ports.srm_port import SRMClientPort
from open_exposure_gateway.ports.traffic_influence_callback_port import (
    TrafficInfluenceCallbackDeliveryPort,
)


class AppState(Protocol):
@@ -13,3 +16,4 @@ class AppState(Protocol):
    db_engine: AsyncEngine
    session_maker: async_sessionmaker[AsyncSession]
    qod_callback_client: QodCallbackDeliveryPort
    traffic_influence_callback_client: TrafficInfluenceCallbackDeliveryPort
+20 −1
Original line number Diff line number Diff line
@@ -55,6 +55,9 @@ from open_exposure_gateway.ports.database.traffic_influences import TrafficInflu
from open_exposure_gateway.ports.databus_port import DataBusPort
from open_exposure_gateway.ports.qod_callback_port import QodCallbackDeliveryPort
from open_exposure_gateway.ports.srm_port import SRMClientPort
from open_exposure_gateway.ports.traffic_influence_callback_port import (
    TrafficInfluenceCallbackDeliveryPort,
)


@dataclass
@@ -151,6 +154,12 @@ def get_qod_callback_client(request: Request) -> QodCallbackDeliveryPort:
    return get_app_state(request=request).qod_callback_client


def get_traffic_influence_callback_client(
    request: Request,
) -> TrafficInfluenceCallbackDeliveryPort:
    return get_app_state(request=request).traffic_influence_callback_client


def get_edge_app_service(
    srm: SRMClientPort = Depends(get_client),
    publisher: DataBusPort = Depends(get_publisher),
@@ -211,7 +220,17 @@ def get_traffic_influence_service(
    callback_registration_repo: CallbackRegistrationRepository = Depends(
        get_callback_registration_repo
    ),
    callback_delivery_repo: CallbackDeliveryRepository = Depends(get_callback_delivery_repo),
    callback_delivery_port: TrafficInfluenceCallbackDeliveryPort = Depends(
        get_traffic_influence_callback_client
    ),
) -> TrafficInfluenceService:
    return TrafficInfluenceService(
        srm, publisher, operation_repo, traffic_influence_repo, callback_registration_repo
        srm,
        publisher,
        operation_repo,
        traffic_influence_repo,
        callback_registration_repo,
        callback_delivery_repo,
        callback_delivery_port,
    )
Loading