Commit 1375ead9 authored by Sergio Gimenez's avatar Sergio Gimenez
Browse files

feat(fm): accept InstallApp and publish command.srm.service.deploy

POST /operatorplatform/federation/v1/{federationContextId}/
application/lcm with an Idempotency-Key: authenticate the partner,
resolve their inbound context, check the agreement permits
install-app in that zone, resolve the app triple to a service
specification, persist a pending transaction, publish
command.srm.service.deploy, flip to in_progress and answer 202 with
the zone and a new instance identifier.

Replaying the same key returns the same instance without
republishing; the same key with a different body is a 409
(idempotency-key-reused). Nothing is published when authorisation
fails. The instance identifier is a UUID's hex form because GSMA's
InstanceIdentifier pattern forbids dashes; zone ids parse as UUIDs
per ADR-0017, so a non-UUID zone is an agreement violation.

The SRM command envelope is extra="forbid" (validated against the
SRM pydantic model on feat/location-retrieval-api), so fields must
never be added to it. FM now requires NATS at startup.
parent 2e4c9f5a
Loading
Loading
Loading
Loading
+11 −0
Original line number Diff line number Diff line
@@ -12,6 +12,7 @@ from federation_manager.domain.errors import (
    FederationAlreadyExists,
    FederationContextMissing,
    FederationContextUnknown,
    IdempotencyKeyReused,
    NoRouteMatched,
    PartnerEndpointConfigurationError,
    PartnerNotActive,
@@ -181,6 +182,16 @@ def register_exception_handlers(app: FastAPI) -> None:
            request.url.path,
        )

    @app.exception_handler(IdempotencyKeyReused)
    async def _idempotency_reused(request: Request, exc: IdempotencyKeyReused) -> JSONResponse:
        return problem(
            409,
            "idempotency-key-reused",
            "Idempotency Key Reused",
            "This Idempotency-Key was already used for a different request.",
            request.url.path,
        )

    @app.exception_handler(FederationContextUnknown)
    async def _context_unknown(request: Request, exc: FederationContextUnknown) -> JSONResponse:
        return problem(
+33 −0
Original line number Diff line number Diff line
from typing import Annotated

from fastapi import APIRouter, Depends, Header, Path

from federation_manager.api.errors import EWBI_ERROR_RESPONSES
from federation_manager.api.security import get_bearer_token
from federation_manager.application.deployment import InboundDeploymentService
from federation_manager.contracts.ewbi import InstallAppRequest, InstallAppResponse
from federation_manager.dependencies import get_inbound_deployment_service
from federation_manager.domain.ewbi import EWBI_BASE_PATH

router = APIRouter(prefix=EWBI_BASE_PATH, tags=["ApplicationDeploymentManagement"])

FederationContextIdPath = Annotated[str, Path(pattern=r"^[A-Za-z0-9][A-Za-z0-9-]*$")]


@router.post(
    "/{federationContextId}/application/lcm",
    operation_id="InstallApp",
    status_code=202,
    responses=EWBI_ERROR_RESPONSES,
)
async def install_app(
    federationContextId: FederationContextIdPath,  # noqa: N803 - GSMA path template name
    body: InstallAppRequest,
    service: Annotated[InboundDeploymentService, Depends(get_inbound_deployment_service)],
    token: Annotated[str, Depends(get_bearer_token)],
    idempotency_key: Annotated[str, Header(alias="Idempotency-Key", min_length=1)],
) -> InstallAppResponse:
    accepted = await service.install(token, federationContextId, idempotency_key, body)
    return InstallAppResponse(
        zone_id=accepted.zone_id, app_inst_identifier=accepted.app_instance_identifier
    )
+153 −0
Original line number Diff line number Diff line
import json
from collections.abc import Callable
from dataclasses import dataclass
from datetime import datetime, timezone
from hashlib import sha256
from uuid import UUID, uuid4

from federation_manager.application.authentication import PartnerAuthenticator
from federation_manager.application.authorization import AgreementChecker
from federation_manager.contracts.ewbi import InstallAppRequest
from federation_manager.contracts.srm import (
    SUBJECT_DEPLOY,
    DeployPayloadV1,
    DeployTargetV1,
    SrmServiceDeployV1,
)
from federation_manager.domain.errors import (
    AgreementViolation,
    FederationContextUnknown,
    IdempotencyKeyReused,
)
from federation_manager.domain.models import FederationTransaction
from federation_manager.domain.ports import (
    DataBusPublisherPort,
    FederationContextRepositoryPort,
    TransactionRepositoryPort,
)

API_TYPE = "install-app"
INBOUND = "inbound"


def _utcnow() -> datetime:
    return datetime.now(timezone.utc)


@dataclass(frozen=True)
class DeploymentAccepted:
    zone_id: str
    app_instance_identifier: str


class InboundDeploymentService:
    def __init__(
        self,
        authenticator: PartnerAuthenticator,
        contexts: FederationContextRepositoryPort,
        agreements: AgreementChecker,
        transactions: TransactionRepositoryPort,
        publisher: DataBusPublisherPort,
        *,
        clock: Callable[[], datetime] = _utcnow,
        id_factory: Callable[[], UUID] = uuid4,
    ) -> None:
        self._authenticator = authenticator
        self._contexts = contexts
        self._agreements = agreements
        self._transactions = transactions
        self._publisher = publisher
        self._clock = clock
        self._new_id = id_factory

    async def install(
        self,
        token: str,
        federation_context_id: str,
        idempotency_key: str,
        request: InstallAppRequest,
    ) -> DeploymentAccepted:
        partner = await self._authenticator.authenticate(token)
        context = await self._contexts.find_inbound(partner.id, federation_context_id)
        if context is None or context.is_terminated():
            raise FederationContextUnknown(partner.id)

        fingerprint = _fingerprint(request)
        replay = await self._transactions.find_by_idempotency_key(
            partner.id, API_TYPE, idempotency_key
        )
        if replay is not None:
            if replay.request_fingerprint != fingerprint:
                raise IdempotencyKeyReused(idempotency_key)
            return DeploymentAccepted(
                zone_id=str(replay.request_summary["zone_id"]),
                app_instance_identifier=str(replay.external_resource_id),
            )

        zone_id = _zone_uuid(request.zone_info.zone_id)
        agreement = await self._agreements.require(partner, API_TYPE, zone_id)
        specification_id = agreement.resolve_app_spec(
            request.app_id, request.app_version, request.zone_info.flavour_id
        )
        if specification_id is None:
            raise AgreementViolation

        app_instance_id = self._new_id()
        operation_id = self._new_id()
        correlation_id = self._new_id()
        transaction = FederationTransaction(
            id=self._new_id(),
            partner_op_id=partner.id,
            agreement_id=agreement.id,
            federation_context_row_id=context.id,
            direction=INBOUND,
            operation_id=operation_id,
            correlation_id=correlation_id,
            api_type=API_TYPE,
            status="pending",
            idempotency_key=idempotency_key,
            request_fingerprint=fingerprint,
            external_resource_id=app_instance_id.hex,
            callback_url=request.app_inst_callback_link,
            callback_status="pending",
            request_summary={
                "app_id": request.app_id,
                "app_version": request.app_version,
                "flavour_id": request.zone_info.flavour_id,
                "zone_id": request.zone_info.zone_id,
                "service_specification_id": str(specification_id),
            },
            started_at=self._clock(),
        )
        await self._transactions.add(transaction)

        command = SrmServiceDeployV1(
            operation_id=operation_id,
            correlation_id=str(correlation_id),
            requested_at=self._clock(),
            app_provider_id=str(partner.id),
            federation_partner_ref=partner.mcc_mnc,
            source="federation",
            service_specification_id=specification_id,
            targets=[DeployTargetV1(app_instance_id=app_instance_id, zone_id=zone_id)],
            deploy=DeployPayloadV1(),
        )
        await self._publisher.publish(SUBJECT_DEPLOY, command.model_dump(mode="json"))
        await self._transactions.mark_in_progress(transaction.id)

        return DeploymentAccepted(
            zone_id=request.zone_info.zone_id, app_instance_identifier=app_instance_id.hex
        )


def _fingerprint(request: InstallAppRequest) -> str:
    canonical = json.dumps(request.model_dump(mode="json", by_alias=True), sort_keys=True)
    return sha256(canonical.encode()).hexdigest()


def _zone_uuid(zone_id: str) -> UUID:
    # ADR-0017: the EWBI zone identifier is SRM's zone UUID, passed through unchanged.
    try:
        return UUID(zone_id)
    except ValueError:
        raise AgreementViolation from None
+35 −0
Original line number Diff line number Diff line
@@ -195,3 +195,38 @@ class ProblemDetails(BaseModel):
    invalid_params: list[InvalidParam] | None = Field(
        default=None, validation_alias="invalidParams", serialization_alias="invalidParams"
    )


class ZoneInfo(BaseModel):
    model_config = ConfigDict(extra="ignore", populate_by_name=True)

    zone_id: str = Field(validation_alias="zoneId", serialization_alias="zoneId")
    flavour_id: str = Field(validation_alias="flavourId", serialization_alias="flavourId")
    resource_consumption: str | None = Field(
        default=None,
        validation_alias="resourceConsumption",
        serialization_alias="resourceConsumption",
    )


class InstallAppRequest(BaseModel):
    model_config = ConfigDict(extra="ignore", populate_by_name=True)

    app_id: str = Field(validation_alias="appId", serialization_alias="appId")
    app_version: str = Field(validation_alias="appVersion", serialization_alias="appVersion")
    app_provider_id: str = Field(
        validation_alias="appProviderId", serialization_alias="appProviderId"
    )
    zone_info: ZoneInfo = Field(validation_alias="zoneInfo", serialization_alias="zoneInfo")
    app_inst_callback_link: str = Field(
        validation_alias="appInstCallbackLink", serialization_alias="appInstCallbackLink"
    )


class InstallAppResponse(BaseModel):
    model_config = ConfigDict(populate_by_name=True)

    zone_id: str = Field(validation_alias="zoneId", serialization_alias="zoneId")
    app_inst_identifier: str = Field(
        validation_alias="appInstIdentifier", serialization_alias="appInstIdentifier"
    )
+23 −0
Original line number Diff line number Diff line
@@ -13,11 +13,13 @@ from federation_manager.adapters.database.routing_repo import PostgresRoutingRul
from federation_manager.adapters.database.transaction_repo import PostgresTransactionRepo
from federation_manager.application.authentication import PartnerAuthenticator
from federation_manager.application.authorization import AgreementChecker
from federation_manager.application.deployment import InboundDeploymentService
from federation_manager.application.federation import InboundFederationService, LocalOperator
from federation_manager.application.outbound import OutboundFederationService
from federation_manager.core.config import get_settings
from federation_manager.domain.ports import (
    AgreementRepositoryPort,
    DataBusPublisherPort,
    EwbiClientPort,
    FederationContextRepositoryPort,
    JwtValidatorPort,
@@ -119,3 +121,24 @@ def get_inbound_federation_service(
            platform_caps=settings.platform_caps,
        ),
    )


def get_command_publisher(request: Request) -> DataBusPublisherPort:
    publisher: DataBusPublisherPort = request.app.state.command_publisher
    return publisher


def get_inbound_deployment_service(
    authenticator: Annotated[PartnerAuthenticator, Depends(get_partner_authenticator)],
    context_repo: Annotated[FederationContextRepositoryPort, Depends(get_federation_context_repo)],
    agreement_repo: Annotated[AgreementRepositoryPort, Depends(get_agreement_repo)],
    transaction_repo: Annotated[TransactionRepositoryPort, Depends(get_transaction_repo)],
    publisher: Annotated[DataBusPublisherPort, Depends(get_command_publisher)],
) -> InboundDeploymentService:
    return InboundDeploymentService(
        authenticator,
        context_repo,
        AgreementChecker(agreement_repo),
        transaction_repo,
        publisher,
    )
Loading