Commit 50d0a3a3 authored by Sergio Gimenez's avatar Sergio Gimenez
Browse files

feat(fm): establish outbound federation with a partner

parent 5a3e97a0
Loading
Loading
Loading
Loading
+15 −1
Original line number Diff line number Diff line
from uuid import UUID

from sqlalchemy import Select, select
from sqlalchemy import Select, insert, select
from sqlalchemy.ext.asyncio import AsyncSession

from federation_manager.adapters.database.tables import federation_contexts as contexts
@@ -34,6 +34,20 @@ class PostgresFederationContextRepo:
            )
        )

    async def add(self, context: FederationContext) -> None:
        await self._session.execute(
            insert(contexts).values(
                id=context.id,
                partner_op_id=context.partner_op_id,
                agreement_id=context.agreement_id,
                direction=context.direction,
                federation_context_id=context.federation_context_id,
                status_callback_url=context.status_callback_url,
                status=context.status,
            )
        )
        await self._session.commit()

    async def _one(self, stmt: Select[tuple[object, ...]]) -> FederationContext | None:
        row = (await self._session.execute(stmt)).one_or_none()
        if row is None:
+10 −2
Original line number Diff line number Diff line
from typing import Any
from uuid import UUID

from sqlalchemy import Select, select
@@ -19,10 +20,17 @@ class PostgresPartnerRepo:
    async def find_by_id(self, partner_id: UUID) -> PartnerOP | None:
        return await self._one(select(partner_ops).where(partner_ops.c.id == partner_id))

    async def list_active(self) -> list[PartnerOP]:
        stmt = select(partner_ops).where(partner_ops.c.status == "active")
        rows = (await self._session.execute(stmt)).all()
        return [self._to_partner(row) for row in rows]

    async def _one(self, stmt: Select[tuple[object, ...]]) -> PartnerOP | None:
        row = (await self._session.execute(stmt)).one_or_none()
        if row is None:
            return None
        return None if row is None else self._to_partner(row)

    @staticmethod
    def _to_partner(row: Any) -> PartnerOP:
        return PartnerOP(
            id=row.id,
            mcc_mnc=row.mcc_mnc,
+5 −1
Original line number Diff line number Diff line
@@ -37,7 +37,11 @@ class HttpxEwbiClient:
                body = response.json()
            except ValueError:
                raise PartnerRequestFailed(partner.id) from None
        return EwbiResponse(status_code=response.status_code, body=body)
        return EwbiResponse(
            status_code=response.status_code,
            body=body,
            location=response.headers.get("Location"),
        )

    @staticmethod
    def _url(partner: PartnerOP, path: str) -> str:
+24 −0
Original line number Diff line number Diff line
@@ -5,10 +5,12 @@ from federation_manager.domain.errors import (
    AgreementExpired,
    AgreementViolation,
    AuthenticationFailed,
    FederationContextMissing,
    NoRouteMatched,
    PartnerEndpointConfigurationError,
    PartnerNotActive,
    PartnerRequestFailed,
    PartnerResponseInvalid,
    PartnerTokenConfigurationError,
    PartnerTokenRequestFailed,
    PartnerUnknown,
@@ -115,6 +117,28 @@ def register_exception_handlers(app: FastAPI) -> None:
            request.url.path,
        )

    @app.exception_handler(PartnerResponseInvalid)
    async def _partner_response_invalid(
        request: Request, exc: PartnerResponseInvalid
    ) -> JSONResponse:
        return problem(
            502,
            "partner-response-invalid",
            "Invalid Partner Response",
            "The partner operator's response did not match the OPG.04 contract.",
            request.url.path,
        )

    @app.exception_handler(FederationContextMissing)
    async def _context_missing(request: Request, exc: FederationContextMissing) -> JSONResponse:
        return problem(
            409,
            "federation-not-established",
            "Federation Not Established",
            "No active federation context exists with the resolved partner operator.",
            request.url.path,
        )

    @app.exception_handler(PartnerEndpointConfigurationError)
    @app.exception_handler(PartnerTokenConfigurationError)
    async def _partner_misconfigured(request: Request, exc: Exception) -> JSONResponse:
+21 −11
Original line number Diff line number Diff line
@@ -11,19 +11,19 @@ from federation_manager.dependencies import get_outbound_federation_service

router = APIRouter(prefix="/internal/federation", tags=["internal-federation"])

# mirrors domain.ewbi.EWBI_SERVICE_PATHS; a test keeps them in sync
ApiType = Literal[
    "edge-cloud-deploy",
    "edge-cloud-scale",
    "edge-cloud-terminate",
    "network-capability",
    "device-location-retrieve",
    "device-status-retrieve",
]
# mirrors domain.ewbi.SERVICE_API_NAMES; a test keeps them in sync
ApiType = Literal["device-location-retrieve", "device-status-retrieve"]

_E164 = re_compile(r"^\+[1-9][0-9]{4,14}$")


class ServiceApiBody(BaseModel):
    model_config = ConfigDict(extra="forbid", populate_by_name=True)

    media_type: Literal["application/json"] = Field(default="application/json", alias="mediaType")
    api_content: dict[str, Any] = Field(alias="APIContent")


class OutboundFederationRequest(BaseModel):
    model_config = ConfigDict(extra="forbid")

@@ -31,7 +31,11 @@ class OutboundFederationRequest(BaseModel):
    identifier_type: Literal["msisdn", "ip"]
    identifier_value: str = Field(min_length=1, max_length=64)
    correlation_id: UUID
    payload: dict[str, Any] = Field(default_factory=dict)
    customer_id: UUID
    customer_info: str = Field(min_length=1, max_length=255)
    txn_identifier: str = Field(min_length=1, max_length=255)
    service_api_body: ServiceApiBody
    event_notification_dest: str | None = None

    @model_validator(mode="after")
    def _identifier_matches_type(self) -> "OutboundFederationRequest":
@@ -50,6 +54,7 @@ class OutboundFederationResponse(BaseModel):
    partner_op_id: UUID
    status_code: int
    body: Any = None
    location: str | None = None


@router.post("/outbound")
@@ -63,11 +68,16 @@ async def outbound(
            identifier_type=body.identifier_type,
            identifier_value=body.identifier_value,
            correlation_id=body.correlation_id,
            payload=body.payload,
            customer_id=body.customer_id,
            customer_info=body.customer_info,
            txn_identifier=body.txn_identifier,
            api_content=body.service_api_body.api_content,
            event_notification_dest=body.event_notification_dest,
        )
    )
    return OutboundFederationResponse(
        partner_op_id=result.partner_op_id,
        status_code=result.status_code,
        body=result.body,
        location=result.location,
    )
Loading