Loading pyproject.toml +1 −0 Original line number Diff line number Diff line Loading @@ -16,6 +16,7 @@ dependencies = [ "pydantic-settings>=2.4", "pyjwt[crypto]>=2.9", "sqlalchemy[asyncio]>=2.0", "structlog>=25.1", ] [project.optional-dependencies] Loading src/federation_manager/adapters/databus/nats_adapter.py +14 −2 Original line number Diff line number Diff line Loading @@ -13,6 +13,9 @@ from federation_manager.contracts.srm import ( TASK_STREAM, TASK_STREAM_MAX_AGE_SECONDS, ) from federation_manager.core.logging import get_logger logger = get_logger(__name__) class NatsCommandPublisher: Loading @@ -31,6 +34,9 @@ class NatsCommandPublisher: self._nc = None self._js = None def is_connected(self) -> bool: return self._nc is not None and self._nc.is_connected async def ensure_task_stream(self) -> None: js = self._require_js() config = StreamConfig( Loading Loading @@ -93,8 +99,14 @@ class NatsEventConsumer: async def on_message(message: Msg) -> None: try: await handler(json.loads(message.data)) except Exception: # Leave it unacked so JetStream redelivers up to max_deliver. except Exception as error: # left unacked so JetStream redelivers up to max_deliver logger.error( "event_handling_failed", subject=message.subject, error=type(error).__name__, exc_info=True, ) return await message.ack() Loading src/federation_manager/api/platform/health.py +29 −1 Original line number Diff line number Diff line from fastapi import APIRouter from typing import Annotated from fastapi import APIRouter, Depends, Response from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncSession from federation_manager.dependencies import get_databus_health, get_session from federation_manager.domain.ports import DataBusHealthPort router = APIRouter(tags=["platform"]) Loading @@ -6,3 +13,24 @@ router = APIRouter(tags=["platform"]) @router.get("/healthz") def healthz() -> dict[str, str]: return {"status": "ok"} @router.get("/readyz") async def readyz( response: Response, session: Annotated[AsyncSession, Depends(get_session)], databus: Annotated[DataBusHealthPort, Depends(get_databus_health)], ) -> dict[str, object]: checks = {"postgres": await _postgres_ready(session), "databus": databus.is_connected()} ready = all(checks.values()) if not ready: response.status_code = 503 return {"status": "ready" if ready else "not ready", "checks": checks} async def _postgres_ready(session: AsyncSession) -> bool: try: await session.execute(text("SELECT 1")) except Exception: return False return True src/federation_manager/application/events.py +12 −1 Original line number Diff line number Diff line Loading @@ -8,12 +8,16 @@ from federation_manager.contracts.ewbi import ( InstanceStatusCallback, ) from federation_manager.contracts.srm import CompletedInstanceV1, SrmOperationCompletedV1 from federation_manager.core.logging import get_logger from federation_manager.domain.models import FederationTransaction from federation_manager.domain.ports import ( CallbackClientPort, PartnerRepositoryPort, TransactionRepositoryPort, ) from federation_manager.domain.topology import sanitize logger = get_logger(__name__) def _utcnow() -> datetime: Loading Loading @@ -69,11 +73,18 @@ class OperationCompletedConsumer: delivered = True for body in self._callback_bodies(transaction, event): payload = body.model_dump(mode="json", by_alias=True, exclude_none=True) payload = sanitize(body.model_dump(mode="json", by_alias=True, exclude_none=True)) delivered &= await self._callbacks.deliver(partner, transaction.callback_url, payload) await self._transactions.record_callback( transaction.id, status="delivered" if delivered else "failed" ) if not delivered: logger.warning( "instance_status_callback_failed", partner_op_id=str(transaction.partner_op_id), api_type=transaction.api_type, callback_url=transaction.callback_url, ) def _callback_bodies( self, transaction: FederationTransaction, event: SrmOperationCompletedV1 Loading src/federation_manager/application/federation.py +11 −3 Original line number Diff line number Diff line Loading @@ -10,6 +10,7 @@ from federation_manager.contracts.ewbi import ( FederationResponseData, MobileNetworkIds, ) from federation_manager.core.logging import get_logger from federation_manager.domain.errors import ( FederationAlreadyExists, FederationContextUnknown, Loading @@ -26,6 +27,8 @@ from federation_manager.domain.ports import ( PartnerRepositoryPort, ) logger = get_logger(__name__) OUTBOUND = "outbound" INBOUND = "inbound" AVAILABLE = "available" Loading Loading @@ -111,9 +114,14 @@ class FederationEstablishmentService: continue try: established.append(await self.establish(partner)) except FederationError: # One unreachable or half-configured partner must not stop FM from starting. continue except FederationError as error: # one unreachable or half-configured partner must not stop FM from starting logger.warning( "bootstrap_federation_failed", partner_op_id=str(partner.id), partner_mcc_mnc=partner.mcc_mnc, error=type(error).__name__, ) return established Loading Loading
pyproject.toml +1 −0 Original line number Diff line number Diff line Loading @@ -16,6 +16,7 @@ dependencies = [ "pydantic-settings>=2.4", "pyjwt[crypto]>=2.9", "sqlalchemy[asyncio]>=2.0", "structlog>=25.1", ] [project.optional-dependencies] Loading
src/federation_manager/adapters/databus/nats_adapter.py +14 −2 Original line number Diff line number Diff line Loading @@ -13,6 +13,9 @@ from federation_manager.contracts.srm import ( TASK_STREAM, TASK_STREAM_MAX_AGE_SECONDS, ) from federation_manager.core.logging import get_logger logger = get_logger(__name__) class NatsCommandPublisher: Loading @@ -31,6 +34,9 @@ class NatsCommandPublisher: self._nc = None self._js = None def is_connected(self) -> bool: return self._nc is not None and self._nc.is_connected async def ensure_task_stream(self) -> None: js = self._require_js() config = StreamConfig( Loading Loading @@ -93,8 +99,14 @@ class NatsEventConsumer: async def on_message(message: Msg) -> None: try: await handler(json.loads(message.data)) except Exception: # Leave it unacked so JetStream redelivers up to max_deliver. except Exception as error: # left unacked so JetStream redelivers up to max_deliver logger.error( "event_handling_failed", subject=message.subject, error=type(error).__name__, exc_info=True, ) return await message.ack() Loading
src/federation_manager/api/platform/health.py +29 −1 Original line number Diff line number Diff line from fastapi import APIRouter from typing import Annotated from fastapi import APIRouter, Depends, Response from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncSession from federation_manager.dependencies import get_databus_health, get_session from federation_manager.domain.ports import DataBusHealthPort router = APIRouter(tags=["platform"]) Loading @@ -6,3 +13,24 @@ router = APIRouter(tags=["platform"]) @router.get("/healthz") def healthz() -> dict[str, str]: return {"status": "ok"} @router.get("/readyz") async def readyz( response: Response, session: Annotated[AsyncSession, Depends(get_session)], databus: Annotated[DataBusHealthPort, Depends(get_databus_health)], ) -> dict[str, object]: checks = {"postgres": await _postgres_ready(session), "databus": databus.is_connected()} ready = all(checks.values()) if not ready: response.status_code = 503 return {"status": "ready" if ready else "not ready", "checks": checks} async def _postgres_ready(session: AsyncSession) -> bool: try: await session.execute(text("SELECT 1")) except Exception: return False return True
src/federation_manager/application/events.py +12 −1 Original line number Diff line number Diff line Loading @@ -8,12 +8,16 @@ from federation_manager.contracts.ewbi import ( InstanceStatusCallback, ) from federation_manager.contracts.srm import CompletedInstanceV1, SrmOperationCompletedV1 from federation_manager.core.logging import get_logger from federation_manager.domain.models import FederationTransaction from federation_manager.domain.ports import ( CallbackClientPort, PartnerRepositoryPort, TransactionRepositoryPort, ) from federation_manager.domain.topology import sanitize logger = get_logger(__name__) def _utcnow() -> datetime: Loading Loading @@ -69,11 +73,18 @@ class OperationCompletedConsumer: delivered = True for body in self._callback_bodies(transaction, event): payload = body.model_dump(mode="json", by_alias=True, exclude_none=True) payload = sanitize(body.model_dump(mode="json", by_alias=True, exclude_none=True)) delivered &= await self._callbacks.deliver(partner, transaction.callback_url, payload) await self._transactions.record_callback( transaction.id, status="delivered" if delivered else "failed" ) if not delivered: logger.warning( "instance_status_callback_failed", partner_op_id=str(transaction.partner_op_id), api_type=transaction.api_type, callback_url=transaction.callback_url, ) def _callback_bodies( self, transaction: FederationTransaction, event: SrmOperationCompletedV1 Loading
src/federation_manager/application/federation.py +11 −3 Original line number Diff line number Diff line Loading @@ -10,6 +10,7 @@ from federation_manager.contracts.ewbi import ( FederationResponseData, MobileNetworkIds, ) from federation_manager.core.logging import get_logger from federation_manager.domain.errors import ( FederationAlreadyExists, FederationContextUnknown, Loading @@ -26,6 +27,8 @@ from federation_manager.domain.ports import ( PartnerRepositoryPort, ) logger = get_logger(__name__) OUTBOUND = "outbound" INBOUND = "inbound" AVAILABLE = "available" Loading Loading @@ -111,9 +114,14 @@ class FederationEstablishmentService: continue try: established.append(await self.establish(partner)) except FederationError: # One unreachable or half-configured partner must not stop FM from starting. continue except FederationError as error: # one unreachable or half-configured partner must not stop FM from starting logger.warning( "bootstrap_federation_failed", partner_op_id=str(partner.id), partner_mcc_mnc=partner.mcc_mnc, error=type(error).__name__, ) return established Loading