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

Merge branch 'refactor/fm' into 'develop'

Refactor/fm

See merge request !27
parents b28c0dd4 c72628ff
Loading
Loading
Loading
Loading
Loading
+65 −0
Original line number Diff line number Diff line
import httpx
import structlog

from open_exposure_gateway.core.config import get_settings

logger = structlog.get_logger(__name__)

# OEG's platform-namespace federation routes (/platform/v1/federation/...) mirror
# FM's own /internal/... routes 1:1, so this is a plain prefix swap -- nothing else.
_INTERNAL_PREFIX = "/internal"


class FmUnavailableError(Exception):
    """FM did not answer the internal call (connection failure or timeout)."""


class FmClient:
    def __init__(self) -> None:
        settings = get_settings()
        self.base_url = str(settings.fm_settings.base_url).rstrip("/")
        self.timeout = settings.fm_settings.timeout

    async def relay(
        self,
        method: str,
        path: str,
        content: bytes | None,
        content_type: str | None,
        x_correlator: str | None,
        params: tuple[tuple[str, str], ...] | None = None,
    ) -> httpx.Response:
        """Forward one federation-lifecycle operation to FM, unaltered.

        `path` is the operator-facing suffix after ``/platform/v1/federation``
        (e.g. ``/partners/{partnerOpId}/federations``,
        ``/partners/{partnerOpId}/federations/{federationContextId}/zones/{zoneId}``).
        The request body and query string are forwarded as the caller sent them --
        no parsing, no reshaping -- and the caller relays the returned response
        body and status the same way.
        """
        url = f"{self.base_url}{_INTERNAL_PREFIX}{path}"
        headers = {"X-Correlator": x_correlator} if x_correlator else None
        if content_type:
            headers = {**(headers or {}), "Content-Type": content_type}

        log = logger.bind(method=method, url=url, x_correlator=x_correlator)

        try:
            async with httpx.AsyncClient(timeout=self.timeout) as client:
                return await client.request(
                    method=method,
                    url=url,
                    content=content,
                    headers=headers,
                    params=params,
                )
        except httpx.TimeoutException as exc:
            log.exception("FM request timed out")
            raise FmUnavailableError("FM request timed out") from exc
        except httpx.ConnectError as exc:
            log.exception("FM connection failed")
            raise FmUnavailableError("Could not connect to FM") from exc
        except httpx.RequestError as exc:
            log.exception("FM request error")
            raise FmUnavailableError("FM request failed") from exc
+7843 −0

File added.

Preview size limit exceeded, changes collapsed.

+359 −0
Original line number Diff line number Diff line
"""FastAPI routes for the platform-namespace federation-lifecycle relay (ADR-0046, ADR-0048).

Bodies are raw GSMA (OPG.04 v6.0 / artifact v1.4.0), forwarded unaltered, but paths
are OEG's own platform namespace rather than the GSMA EWBI paths: GSMA addresses the
callee operator through the URL (operator A calls B's FM directly), which only works
because there is no intermediate hop. Here the call chain is Portal -> OEG -> FM-A,
with FM-A only then making the real EWBI call to FM-B, so the partner has to be named
explicitly -- as ``partnerOpId``, our own registered id for the partner -- on every
route. FM uses it to look up B's token URL, client id and secret from B's
registration, and to reject a ``federationContextId`` that belongs to a different
partner. Every handler is otherwise a stateless synchronous relay to FM: the raw
request body (and query string) is forwarded to FM's matching ``/internal`` endpoint
unaltered, and FM's response body and status are relayed back unaltered. The two POST
bodies (``FederationRequestData``, ``ZoneRegistrationRequestData``) are validated for
shape before being forwarded, but the *validated* object is only used as a gate; the
original raw bytes are what's actually sent to FM, so a field OEG's model doesn't
happen to cover is still forwarded unchanged rather than dropped. The remaining
schemas in ``schemas.py`` are referenced here only to document the response contract
in the OpenAPI document (via ``response_model``); FM's response is never parsed
against them.

This module exposes the ``FederationManagement`` tag plus ``zone_subscribe``,
``get_zone_data``, ``get_zone_details`` and ``zone_unsubscribe`` from
``AvailabilityZoneInfoSynchronization``. ``update_federation``
(``UpdateFederation``) is out of scope.
"""

import json
from typing import Annotated, Any

from fastapi import APIRouter, Depends, Query, Request, Response
from pydantic import BaseModel, ValidationError

from open_exposure_gateway.adapters.http.fm_client import FmClient, FmUnavailableError
from open_exposure_gateway.api.camara.common import XCorrelatorHeader
from open_exposure_gateway.api.gsma.federation_manager.v1_4_0.schemas import (
    FederationContextId,
    FederationContextIdResponse,
    FederationDetails,
    FederationRequestData,
    FederationResponseData,
    PartnerOpId,
    ProblemDetails,
    ZoneIdentifier,
    ZoneRegisteredData,
    ZoneRegistrationRequestData,
    ZoneRegistrationResponseData,
)
from open_exposure_gateway.dependencies import get_fm_client

# OEG's own platform namespace, mirrored 1:1 at FM's `/internal` (FmClient does a
# prefix swap, nothing else -- see fm_client.py).
BASE_PATH = "/platform/v1/federation"

PARTNER_PREFIX = "/partners/{partnerOpId}/federations"

# Tags are set per route (not on the router) so each endpoint carries only its own
# GSMA operation's tag.
router = APIRouter(prefix=BASE_PATH)

FmClientDep = Annotated[FmClient, Depends(get_fm_client)]

# Every Federation Manager error response is an RFC 7807 ProblemDetails.
_ERROR_RESPONSES: dict[int | str, dict[str, Any]] = {
    400: {"model": ProblemDetails, "description": "Bad request"},
    401: {"model": ProblemDetails, "description": "Unauthorized"},
    404: {"model": ProblemDetails, "description": "Not found"},
    409: {"model": ProblemDetails, "description": "Conflict"},
    422: {"model": ProblemDetails, "description": "Unprocessable entity"},
    500: {"model": ProblemDetails, "description": "Internal server error"},
    502: {"model": ProblemDetails, "description": "Bad gateway"},
    503: {"model": ProblemDetails, "description": "Service unavailable"},
    520: {"model": ProblemDetails, "description": "Unknown error"},
}


def _responses(*codes: int) -> dict[int | str, dict[str, Any]]:
    return {code: _ERROR_RESPONSES[code] for code in codes}


def _problem_response(status_code: int, title: str, detail: str) -> Response:
    problem = ProblemDetails(title=title, detail=detail)
    return Response(
        content=problem.model_dump_json(exclude_none=True),
        status_code=status_code,
        media_type="application/problem+json",
    )


def _is_valid_json(content: bytes) -> bool:
    try:
        json.loads(content)
    except ValueError:
        return False
    return True


async def _relay(
    fm_client: FmClient,
    method: str,
    path: str,
    request: Request,
    x_correlator: str | None,
    request_schema: type[BaseModel] | None = None,
) -> Response:
    """Forward one federation-lifecycle operation to FM and relay its answer unaltered.

    ``path`` is the operator-facing suffix after ``BASE_PATH`` (e.g.
    ``/partners/{partnerOpId}/federations``), which ``FmClient`` maps onto FM's
    matching internal endpoint. The request's query string is forwarded as-is
    (e.g. ``GetZoneData``'s ``zoneId`` filter). A body that isn't even
    syntactically JSON can't reach FM meaningfully, so that's rejected here
    (ADR-0047: OEG-authored 400); the mirror case on the way back -- FM answering
    with a body it labelled JSON that isn't actually parseable -- is ADR-0047's 502.

    When ``request_schema`` is given, the parsed body is also validated against it
    as a shape check before being forwarded. The validated object itself is
    discarded afterwards -- what's sent to FM is always the original raw bytes, so
    this is a gate, not a reshape.
    """
    body = await request.body()
    if body:
        try:
            parsed = json.loads(body)
        except ValueError:
            return _problem_response(
                400, "Malformed request body", "Request body is not valid JSON"
            )
        if request_schema is not None:
            try:
                request_schema.model_validate(parsed)
            except ValidationError as exc:
                return _problem_response(
                    400,
                    "Request body does not match the expected schema",
                    str(exc.errors())[:512],
                )

    try:
        fm_response = await fm_client.relay(
            method,
            path,
            content=body or None,
            content_type=request.headers.get("content-type"),
            x_correlator=x_correlator,
            params=tuple(request.query_params.multi_items()),
        )
    except FmUnavailableError as exc:
        return _problem_response(503, "Federation Manager unavailable", str(exc))

    if fm_response.content and not _is_valid_json(fm_response.content):
        return _problem_response(
            502,
            "Federation Manager returned an unusable response",
            "Response body is not valid JSON",
        )

    return Response(
        content=fm_response.content,
        status_code=fm_response.status_code,
        media_type=fm_response.headers.get("content-type"),
    )


@router.post(
    PARTNER_PREFIX,
    tags=["Federation Manager"],
    summary="Create a one-direction federation with a partner operator platform",
    operation_id="create_federation",
    response_model=FederationResponseData,
    responses=_responses(400, 401, 404, 409, 422, 500, 502, 503, 520),
)
async def create_federation(
    partnerOpId: PartnerOpId,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client,
        "POST",
        f"/partners/{partnerOpId}/federations",
        request,
        x_correlator,
        request_schema=FederationRequestData,
    )


@router.get(
    PARTNER_PREFIX,
    tags=["Federation Manager"],
    summary="Retrieve the existing federationContextId with the partner operator platform",
    operation_id="get_federation_context_id",
    response_model=FederationContextIdResponse,
    responses=_responses(400, 401, 404, 409, 422, 500, 502, 503, 520),
)
async def get_federation_context_id(
    partnerOpId: PartnerOpId,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client, "GET", f"/partners/{partnerOpId}/federations", request, x_correlator
    )


@router.get(
    f"{PARTNER_PREFIX}/{{federationContextId}}",
    tags=["Federation Manager"],
    summary="Retrieve details about the federation context with the partner OP",
    operation_id="get_federation_details",
    response_model=FederationDetails,
    responses=_responses(400, 401, 404, 409, 422, 500, 502, 503, 520),
)
async def get_federation_details(
    partnerOpId: PartnerOpId,
    federationContextId: FederationContextId,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client,
        "GET",
        f"/partners/{partnerOpId}/federations/{federationContextId}",
        request,
        x_correlator,
    )


@router.delete(
    f"{PARTNER_PREFIX}/{{federationContextId}}",
    tags=["Federation Manager"],
    summary="Remove an existing federation with the partner OP",
    operation_id="delete_federation_details",
    status_code=200,
    responses=_responses(400, 401, 404, 409, 422, 500, 502, 503, 520),
)
async def delete_federation_details(
    partnerOpId: PartnerOpId,
    federationContextId: FederationContextId,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client,
        "DELETE",
        f"/partners/{partnerOpId}/federations/{federationContextId}",
        request,
        x_correlator,
    )


@router.get(
    f"{PARTNER_PREFIX}/{{federationContextId}}/zones",
    tags=["Availability Zone Info Synchronization"],
    summary="List the zones the partner OP offers, optionally filtered to one zone",
    operation_id="get_zone_data",
    response_model=ZoneRegisteredData,
    responses=_responses(400, 401, 404, 409, 422, 500, 502, 503, 520),
)
async def get_zone_data(
    partnerOpId: PartnerOpId,
    federationContextId: FederationContextId,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
    zoneId: Annotated[ZoneIdentifier | None, Query()] = None,
) -> Response:
    return await _relay(
        fm_client,
        "GET",
        f"/partners/{partnerOpId}/federations/{federationContextId}/zones",
        request,
        x_correlator,
    )


@router.post(
    f"{PARTNER_PREFIX}/{{federationContextId}}/zones",
    tags=["Availability Zone Info Synchronization"],
    summary="Subscribe to partner OP availability zones and reserve zone resources",
    operation_id="zone_subscribe",
    response_model=ZoneRegistrationResponseData,
    responses=_responses(400, 401, 404, 409, 422, 500, 502, 503, 520),
)
async def zone_subscribe(
    partnerOpId: PartnerOpId,
    federationContextId: FederationContextId,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client,
        "POST",
        f"/partners/{partnerOpId}/federations/{federationContextId}/zones",
        request,
        x_correlator,
        request_schema=ZoneRegistrationRequestData,
    )


@router.get(
    f"{PARTNER_PREFIX}/{{federationContextId}}/zones/{{zoneId}}",
    tags=["Availability Zone Info Synchronization"],
    summary=(
        "Retrieves details about the computation and network resources that partner "
        "OP has reserved for this zone"
    ),
    operation_id="get_zone_details",
    response_model=ZoneRegisteredData,
    responses=_responses(400, 401, 404, 409, 422, 500, 502, 503, 520),
)
async def get_zone_details(
    partnerOpId: PartnerOpId,
    federationContextId: FederationContextId,
    zoneId: ZoneIdentifier,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client,
        "GET",
        f"/partners/{partnerOpId}/federations/{federationContextId}/zones/{zoneId}",
        request,
        x_correlator,
    )


@router.delete(
    f"{PARTNER_PREFIX}/{{federationContextId}}/zones/{{zoneId}}",
    tags=["Availability Zone Info Synchronization"],
    summary=(
        "Assert usage of a partner OP zone. Originating OP informs partner OP that "
        "it will no longer access the specified zone"
    ),
    operation_id="zone_unsubscribe",
    status_code=200,
    responses=_responses(400, 401, 404, 409, 422, 500, 502, 503, 520),
)
async def zone_unsubscribe(
    partnerOpId: PartnerOpId,
    federationContextId: FederationContextId,
    zoneId: ZoneIdentifier,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client,
        "DELETE",
        f"/partners/{partnerOpId}/federations/{federationContextId}/zones/{zoneId}",
        request,
        x_correlator,
    )
+363 −0

File added.

Preview size limit exceeded, changes collapsed.

+6 −0
Original line number Diff line number Diff line
@@ -69,6 +69,11 @@ class LocationRetrievalSettings(BaseModel):
        return self.service_specification_id == DEFAULT_LOCATION_RETRIEVAL_SERVICE_SPECIFICATION_ID


class FmSettings(BaseModel):
    base_url: HttpUrl = HttpUrl("http://localhost:8082")
    timeout: float = 10.0


class Settings(BaseSettings):
    model_config = SettingsConfigDict(
        env_file=".env",
@@ -91,6 +96,7 @@ class Settings(BaseSettings):
    callback_settings: CallbackSettings = CallbackSettings()
    qod_settings: QodSettings = QodSettings()
    location_retrieval_settings: LocationRetrievalSettings = LocationRetrievalSettings()
    fm_settings: FmSettings = FmSettings()


@lru_cache
Loading