Commit 1c2d6fc7 authored by George Papathanail's avatar George Papathanail
Browse files

feat: wire federation-lifecycle routes to FmClient

parent 89a0e9cc
Loading
Loading
Loading
Loading
Loading
+112 −42
Original line number Diff line number Diff line
"""FastAPI routes for the GSMA Federation Manager API (EWBI OPG v1.2.0).

Endpoints defined by ``API_definitions/federation-manager.yaml``.

This module currently exposes the ``FederationManagement`` tag plus
``zone_subscribe``, ``get_zone_data`` and ``zone_unsubscribe`` from
``AvailabilityZoneInfoSynchronization``, and every handler is a surface stub that
raises :class:`NotImplementedException` (HTTP 501). The request/response contract
(path, schema, error codes) is live and visible in the OpenAPI document; the
fulfilment logic is added in a later step.
Endpoints defined by ``API_definitions/federation-manager.yaml``. Every handler is a
stateless synchronous relay to FM (ADR-0046, ADR-0048): the raw request body is
forwarded to FM's matching internal endpoint unaltered, and FM's response body and
status are relayed back unaltered. OEG does no request/response parsing or
reshaping on this path -- the Pydantic schemas in ``schemas.py`` are referenced
here only to document the response contract in the OpenAPI document (via
``response_model``); they are not used to build, validate or serialize the actual
request/response bodies at runtime, which would risk silently altering a payload
OEG's hand-maintained models don't fully cover.

This module exposes the ``FederationManagement`` tag plus ``zone_subscribe``,
``get_zone_data`` and ``zone_unsubscribe`` from
``AvailabilityZoneInfoSynchronization``. ``update_federation``
(``PATCH /{federationContextId}/partner``) is out of scope.
"""

from typing import Any
from typing import Annotated, Any

from fastapi import APIRouter
from fastapi import APIRouter, Depends, Request, Response

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_2_0.schemas import (
    FederationContextId,
    FederationContextIdResponse,
    FederationDetails,
    FederationRequestData,
    FederationResponseData,
    ProblemDetails,
    ZoneIdentifier,
    ZoneRegisteredData,
    ZoneRegistrationRequestData,
    ZoneRegistrationResponseData,
)
from open_exposure_gateway.core.exceptions import NotImplementedException
from open_exposure_gateway.dependencies import get_fm_client

# GSMA OPG serves the Federation Management API at {apiRoot}/operatorplatform/federation/v1.
BASE_PATH = "/operatorplatform/federation/v1"
@@ -35,6 +41,8 @@ BASE_PATH = "/operatorplatform/federation/v1"
# federation-manager.yaml tag.
router = APIRouter(prefix=BASE_PATH)

FmClientDep = Annotated[FmClient, Depends(get_fm_client)]

# Every Federation Manager error response is an RFC 7807 ProblemDetails (federation-manager.yaml).
_ERROR_RESPONSES: dict[int | str, dict[str, Any]] = {
    400: {"model": ProblemDetails, "description": "Bad request"},
@@ -43,7 +51,6 @@ _ERROR_RESPONSES: dict[int | str, dict[str, Any]] = {
    409: {"model": ProblemDetails, "description": "Conflict"},
    422: {"model": ProblemDetails, "description": "Unprocessable entity"},
    500: {"model": ProblemDetails, "description": "Internal server error"},
    501: {"model": ProblemDetails, "description": "Not implemented"},
    503: {"model": ProblemDetails, "description": "Service unavailable"},
    520: {"model": ProblemDetails, "description": "Unknown error"},
}
@@ -53,7 +60,40 @@ def _responses(*codes: int) -> dict[int | str, dict[str, Any]]:
    return {code: _ERROR_RESPONSES[code] for code in codes}


_NOT_IMPLEMENTED = "Federation Manager is not implemented in this release"
async def _relay(
    fm_client: FmClient,
    method: str,
    path: str,
    request: Request,
    x_correlator: str | 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. ``/partner``),
    which ``FmClient`` maps onto FM's internal endpoint.
    """
    body = await request.body()
    try:
        fm_response = await fm_client.relay(
            method,
            path,
            content=body or None,
            content_type=request.headers.get("content-type"),
            x_correlator=x_correlator,
        )
    except FmUnavailableError as exc:
        problem = ProblemDetails(title="Federation Manager unavailable", detail=str(exc))
        return Response(
            content=problem.model_dump_json(exclude_none=True),
            status_code=503,
            media_type="application/problem+json",
        )

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


@router.post(
@@ -62,11 +102,12 @@ _NOT_IMPLEMENTED = "Federation Manager is not implemented in this release"
    summary="Create a one-direction federation with a partner operator platform",
    operation_id="create_federation",
    response_model=FederationResponseData,
    response_model_exclude_none=True,
    responses=_responses(400, 401, 404, 409, 422, 500, 501, 503, 520),
    responses=_responses(400, 401, 404, 409, 422, 500, 503, 520),
)
async def create_federation(request: FederationRequestData) -> Any:
    raise NotImplementedException(message=_NOT_IMPLEMENTED)
async def create_federation(
    request: Request, fm_client: FmClientDep, x_correlator: XCorrelatorHeader = None
) -> Response:
    return await _relay(fm_client, "POST", "/partner", request, x_correlator)


@router.get(
@@ -75,11 +116,15 @@ async def create_federation(request: FederationRequestData) -> Any:
    summary="Retrieve details about the federation context with the partner OP",
    operation_id="get_federation_details",
    response_model=FederationDetails,
    response_model_exclude_none=True,
    responses=_responses(400, 401, 404, 409, 422, 500, 501, 503, 520),
    responses=_responses(400, 401, 404, 409, 422, 500, 503, 520),
)
async def get_federation_details(federationContextId: FederationContextId) -> Any:
    raise NotImplementedException(message=_NOT_IMPLEMENTED)
async def get_federation_details(
    federationContextId: FederationContextId,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(fm_client, "GET", f"/{federationContextId}/partner", request, x_correlator)


@router.delete(
@@ -88,10 +133,17 @@ async def get_federation_details(federationContextId: FederationContextId) -> An
    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, 501, 503, 520),
    responses=_responses(400, 401, 404, 409, 422, 500, 503, 520),
)
async def delete_federation_details(
    federationContextId: FederationContextId,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client, "DELETE", f"/{federationContextId}/partner", request, x_correlator
    )
async def delete_federation_details(federationContextId: FederationContextId) -> Any:
    raise NotImplementedException(message=_NOT_IMPLEMENTED)


@router.get(
@@ -100,11 +152,12 @@ async def delete_federation_details(federationContextId: FederationContextId) ->
    summary="Retrieve the existing federationContextId with the partner operator platform",
    operation_id="get_federation_context_id",
    response_model=FederationContextIdResponse,
    response_model_exclude_none=True,
    responses=_responses(400, 401, 404, 409, 422, 500, 501, 503, 520),
    responses=_responses(400, 401, 404, 409, 422, 500, 503, 520),
)
async def get_federation_context_id() -> Any:
    raise NotImplementedException(message=_NOT_IMPLEMENTED)
async def get_federation_context_id(
    request: Request, fm_client: FmClientDep, x_correlator: XCorrelatorHeader = None
) -> Response:
    return await _relay(fm_client, "GET", "/fed-context-id", request, x_correlator)


@router.post(
@@ -113,13 +166,15 @@ async def get_federation_context_id() -> Any:
    summary="Subscribe to partner OP availability zones and reserve zone resources",
    operation_id="zone_subscribe",
    response_model=ZoneRegistrationResponseData,
    response_model_exclude_none=True,
    responses=_responses(400, 401, 404, 409, 422, 500, 501, 503, 520),
    responses=_responses(400, 401, 404, 409, 422, 500, 503, 520),
)
async def zone_subscribe(
    federationContextId: FederationContextId, request: ZoneRegistrationRequestData
) -> Any:
    raise NotImplementedException(message=_NOT_IMPLEMENTED)
    federationContextId: FederationContextId,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(fm_client, "POST", f"/{federationContextId}/zones", request, x_correlator)


@router.get(
@@ -131,11 +186,18 @@ async def zone_subscribe(
    ),
    operation_id="get_zone_data",
    response_model=ZoneRegisteredData,
    response_model_exclude_none=True,
    responses=_responses(400, 401, 404, 409, 422, 500, 501, 503, 520),
    responses=_responses(400, 401, 404, 409, 422, 500, 503, 520),
)
async def get_zone_data(
    federationContextId: FederationContextId,
    zoneId: ZoneIdentifier,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client, "GET", f"/{federationContextId}/zones/{zoneId}", request, x_correlator
    )
async def get_zone_data(federationContextId: FederationContextId, zoneId: ZoneIdentifier) -> Any:
    raise NotImplementedException(message=_NOT_IMPLEMENTED)


@router.delete(
@@ -147,7 +209,15 @@ async def get_zone_data(federationContextId: FederationContextId, zoneId: ZoneId
    ),
    operation_id="zone_unsubscribe",
    status_code=200,
    responses=_responses(400, 401, 404, 409, 422, 500, 501, 503, 520),
    responses=_responses(400, 401, 404, 409, 422, 500, 503, 520),
)
async def zone_unsubscribe(
    federationContextId: FederationContextId,
    zoneId: ZoneIdentifier,
    request: Request,
    fm_client: FmClientDep,
    x_correlator: XCorrelatorHeader = None,
) -> Response:
    return await _relay(
        fm_client, "DELETE", f"/{federationContextId}/zones/{zoneId}", request, x_correlator
    )
async def zone_unsubscribe(federationContextId: FederationContextId, zoneId: ZoneIdentifier) -> Any:
    raise NotImplementedException(message=_NOT_IMPLEMENTED)
+2 −0
Original line number Diff line number Diff line
@@ -2,6 +2,7 @@ from typing import Protocol

from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession, async_sessionmaker

from open_exposure_gateway.adapters.http.fm_client import FmClient
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
@@ -13,3 +14,4 @@ class AppState(Protocol):
    db_engine: AsyncEngine
    session_maker: async_sessionmaker[AsyncSession]
    qod_callback_client: QodCallbackDeliveryPort
    fm_client: FmClient
+5 −0
Original line number Diff line number Diff line
@@ -24,6 +24,7 @@ from open_exposure_gateway.adapters.database.repos.operations import (
from open_exposure_gateway.adapters.database.repos.qod_sessions import (
    SqlQodSessionRepository,
)
from open_exposure_gateway.adapters.http.fm_client import FmClient
from open_exposure_gateway.api.camara.common import XCorrelatorHeader
from open_exposure_gateway.api.error_handlers import x_correlator_header
from open_exposure_gateway.application.services.edge_application_management_service import (
@@ -140,6 +141,10 @@ def get_qod_callback_client(request: Request) -> QodCallbackDeliveryPort:
    return get_app_state(request=request).qod_callback_client


def get_fm_client(request: Request) -> FmClient:
    return get_app_state(request=request).fm_client


def get_edge_app_service(
    srm: SRMClientPort = Depends(get_client),
    publisher: DataBusPort = Depends(get_publisher),
+3 −0
Original line number Diff line number Diff line
@@ -29,6 +29,7 @@ from open_exposure_gateway.adapters.databus.nats_adapter import (
    NatsOperationStatusConsumer,
)
from open_exposure_gateway.adapters.http.callback_client import HttpCallbackClient
from open_exposure_gateway.adapters.http.fm_client import FmClient
from open_exposure_gateway.adapters.http.qod_callback_client import HttpQodCallbackClient
from open_exposure_gateway.adapters.http.srm_client import SRMClient
from open_exposure_gateway.api.camara.edge_application_management.vwip.router import (
@@ -239,6 +240,7 @@ async def default_lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
        raise

    qod_callback_client = HttpQodCallbackClient()
    fm_client = FmClient()

    try:
        publisher = NatsMessagePublisher(settings.nats_settings)
@@ -278,6 +280,7 @@ async def default_lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
    app.state.db_engine = db_engine
    app.state.session_maker = session_maker
    app.state.qod_callback_client = qod_callback_client
    app.state.fm_client = fm_client

    yield

+113 −55
Original line number Diff line number Diff line
"""GSMA Federation Manager -- surface stub.
"""GSMA Federation Manager -- stateless relay to FM (ADR-0046, ADR-0048).

The FederationManagement routes are wired into the app and visible in the OpenAPI
document, but every handler returns HTTP 501 until fulfilment logic is added.
Every route forwards the raw request body to FM's matching internal endpoint and
relays FM's response back unaltered. These tests exercise that relay contract
through a fake FmClient (no real FM, no real network) rather than FmClient's own
transport behaviour, which test_fm_client.py already covers.
"""

from datetime import datetime, timezone
from collections.abc import Generator
from typing import Any

import httpx
import pytest
from fastapi.testclient import TestClient
from pydantic import ValidationError

from open_exposure_gateway.adapters.http.fm_client import FmUnavailableError
from open_exposure_gateway.api.gsma.federation_manager.v1_2_0.router import BASE_PATH
from open_exposure_gateway.api.gsma.federation_manager.v1_2_0.schemas import (
    FederationRequestData,
    ZoneRegistrationRequestData,
)
from open_exposure_gateway.dependencies import get_fm_client
from open_exposure_gateway.main import app

client = TestClient(app, raise_server_exceptions=False)

_CTX = "fed-ctx-1"
_ZONE = "zone-a"
_VALID_CREATE = {
    "origOPFederationId": "orig-op-1",
    "initialDate": "2026-09-04T12:00:00Z",
    "partnerStatusLink": "https://orig-op.example.com/status",
}
_VALID_ZONE_SUBSCRIBE = {
    "acceptedAvailabilityZones": ["zone-a"],
    "availZoneNotifLink": "https://orig-op.example.com/zone-notif",


class _FakeFmClient:
    def __init__(self) -> None:
        self.calls: list[dict[str, Any]] = []
        self.response = httpx.Response(
            200,
            content=b'{"federationContextId": "fed-ctx-1"}',
            headers={"content-type": "application/json"},
        )
        self.raise_unavailable = False

    async def relay(
        self,
        method: str,
        path: str,
        content: bytes | None,
        content_type: str | None,
        x_correlator: str | None,
    ) -> httpx.Response:
        self.calls.append(
            {
                "method": method,
                "path": path,
                "content": content,
                "content_type": content_type,
                "x_correlator": x_correlator,
            }
        )
        if self.raise_unavailable:
            raise FmUnavailableError("FM request timed out")
        return self.response


@pytest.fixture
def fake_fm_client() -> Generator[_FakeFmClient, None, None]:
    fake = _FakeFmClient()
    app.dependency_overrides[get_fm_client] = lambda: fake
    yield fake
    app.dependency_overrides.clear()


@pytest.fixture
def client() -> TestClient:
    return TestClient(app)


_ENDPOINTS = [
    ("post", f"{BASE_PATH}/partner", _VALID_CREATE),
    ("get", f"{BASE_PATH}/{_CTX}/partner", None),
    ("delete", f"{BASE_PATH}/{_CTX}/partner", None),
    ("get", f"{BASE_PATH}/fed-context-id", None),
    ("post", f"{BASE_PATH}/{_CTX}/zones", _VALID_ZONE_SUBSCRIBE),
    ("get", f"{BASE_PATH}/{_CTX}/zones/{_ZONE}", None),
    ("delete", f"{BASE_PATH}/{_CTX}/zones/{_ZONE}", None),
    ("post", f"{BASE_PATH}/partner", "POST", "/partner", b'{"a": 1}'),
    ("get", f"{BASE_PATH}/{_CTX}/partner", "GET", f"/{_CTX}/partner", None),
    ("delete", f"{BASE_PATH}/{_CTX}/partner", "DELETE", f"/{_CTX}/partner", None),
    ("get", f"{BASE_PATH}/fed-context-id", "GET", "/fed-context-id", None),
    ("post", f"{BASE_PATH}/{_CTX}/zones", "POST", f"/{_CTX}/zones", b'{"a": 1}'),
    ("get", f"{BASE_PATH}/{_CTX}/zones/{_ZONE}", "GET", f"/{_CTX}/zones/{_ZONE}", None),
    ("delete", f"{BASE_PATH}/{_CTX}/zones/{_ZONE}", "DELETE", f"/{_CTX}/zones/{_ZONE}", None),
]

_OPERATION_IDS = {
@@ -53,10 +88,57 @@ _OPERATION_IDS = {
}


@pytest.mark.parametrize(("method", "path", "body"), _ENDPOINTS)
def test_endpoint_returns_501(method: str, path: str, body: dict[str, Any] | None) -> None:
    response = client.request(method, path, json=body)
    assert response.status_code == 501
@pytest.mark.parametrize(("http_method", "url", "fm_method", "fm_path", "body"), _ENDPOINTS)
def test_endpoint_relays_to_fm_unaltered(
    client: TestClient,
    fake_fm_client: _FakeFmClient,
    http_method: str,
    url: str,
    fm_method: str,
    fm_path: str,
    body: bytes | None,
) -> None:
    headers = {"x-correlator": "corr-1"}
    if body:
        headers["content-type"] = "application/json"
    response = client.request(http_method, url, content=body, headers=headers)

    assert response.status_code == 200
    assert response.content == b'{"federationContextId": "fed-ctx-1"}'
    assert len(fake_fm_client.calls) == 1
    call = fake_fm_client.calls[0]
    assert call["method"] == fm_method
    assert call["path"] == fm_path
    assert call["content"] == body
    assert call["x_correlator"] == "corr-1"


def test_fm_error_response_is_relayed_unaltered(
    client: TestClient, fake_fm_client: _FakeFmClient
) -> None:
    fake_fm_client.response = httpx.Response(
        409,
        content=b'{"detail": "federation already exists"}',
        headers={"content-type": "application/problem+json"},
    )

    response = client.get(f"{BASE_PATH}/{_CTX}/partner")

    assert response.status_code == 409
    assert response.headers["content-type"] == "application/problem+json"
    assert response.content == b'{"detail": "federation already exists"}'


def test_fm_unavailable_returns_problem_details_503(
    client: TestClient, fake_fm_client: _FakeFmClient
) -> None:
    fake_fm_client.raise_unavailable = True

    response = client.get(f"{BASE_PATH}/fed-context-id")

    assert response.status_code == 503
    assert response.headers["content-type"] == "application/problem+json"
    assert response.json()["title"] == "Federation Manager unavailable"


def test_openapi_exposes_federation_management_operations() -> None:
@@ -68,27 +150,3 @@ def test_openapi_exposes_federation_management_operations() -> None:
        if isinstance(op, dict) and "operationId" in op
    }
    assert _OPERATION_IDS <= seen


def test_federation_request_data_round_trips() -> None:
    model = FederationRequestData.model_validate(_VALID_CREATE)
    assert model.origOPFederationId == "orig-op-1"
    assert model.initialDate == datetime(2026, 9, 4, 12, 0, tzinfo=timezone.utc)


def test_federation_request_data_requires_partner_status_link() -> None:
    payload = {k: v for k, v in _VALID_CREATE.items() if k != "partnerStatusLink"}
    with pytest.raises(ValidationError):
        FederationRequestData.model_validate(payload)


def test_zone_registration_request_data_round_trips() -> None:
    model = ZoneRegistrationRequestData.model_validate(_VALID_ZONE_SUBSCRIBE)
    assert model.acceptedAvailabilityZones == ["zone-a"]


def test_zone_registration_request_data_rejects_empty_zone_list() -> None:
    with pytest.raises(ValidationError):
        ZoneRegistrationRequestData.model_validate(
            {"acceptedAvailabilityZones": [], "availZoneNotifLink": "https://x.example.com"}
        )