From 8ae07567cfff19f26b8c2559a30e2d591622e966 Mon Sep 17 00:00:00 2001 From: Sergio Gimenez Date: Tue, 15 Sep 2026 16:01:52 +0200 Subject: [PATCH] Add topology hiding, readiness probe and structured logging --- pyproject.toml | 1 + .../adapters/databus/nats_adapter.py | 16 ++- src/federation_manager/api/platform/health.py | 30 ++++- src/federation_manager/application/events.py | 13 +- .../application/federation.py | 14 +- src/federation_manager/core/config.py | 1 + src/federation_manager/core/logging.py | 27 ++++ src/federation_manager/dependencies.py | 6 + src/federation_manager/domain/ports.py | 4 + src/federation_manager/domain/topology.py | 51 +++++++ src/federation_manager/main.py | 4 + src/federation_manager/middleware/__init__.py | 0 .../middleware/topology_hiding.py | 49 +++++++ tests/test_platform_health.py | 71 ++++++++++ tests/test_topology_hiding.py | 127 ++++++++++++++++++ uv.lock | 11 ++ 16 files changed, 418 insertions(+), 7 deletions(-) create mode 100644 src/federation_manager/core/logging.py create mode 100644 src/federation_manager/domain/topology.py create mode 100644 src/federation_manager/middleware/__init__.py create mode 100644 src/federation_manager/middleware/topology_hiding.py create mode 100644 tests/test_platform_health.py create mode 100644 tests/test_topology_hiding.py diff --git a/pyproject.toml b/pyproject.toml index 75ecd9b..68a8346 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -16,6 +16,7 @@ dependencies = [ "pydantic-settings>=2.4", "pyjwt[crypto]>=2.9", "sqlalchemy[asyncio]>=2.0", + "structlog>=25.1", ] [project.optional-dependencies] diff --git a/src/federation_manager/adapters/databus/nats_adapter.py b/src/federation_manager/adapters/databus/nats_adapter.py index f40ff8b..2565efe 100644 --- a/src/federation_manager/adapters/databus/nats_adapter.py +++ b/src/federation_manager/adapters/databus/nats_adapter.py @@ -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: @@ -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( @@ -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() diff --git a/src/federation_manager/api/platform/health.py b/src/federation_manager/api/platform/health.py index 26859e6..6ad9224 100644 --- a/src/federation_manager/api/platform/health.py +++ b/src/federation_manager/api/platform/health.py @@ -1,4 +1,11 @@ -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"]) @@ -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 diff --git a/src/federation_manager/application/events.py b/src/federation_manager/application/events.py index 370f943..a90cd19 100644 --- a/src/federation_manager/application/events.py +++ b/src/federation_manager/application/events.py @@ -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: @@ -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 diff --git a/src/federation_manager/application/federation.py b/src/federation_manager/application/federation.py index e206fb3..d6889bf 100644 --- a/src/federation_manager/application/federation.py +++ b/src/federation_manager/application/federation.py @@ -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, @@ -26,6 +27,8 @@ from federation_manager.domain.ports import ( PartnerRepositoryPort, ) +logger = get_logger(__name__) + OUTBOUND = "outbound" INBOUND = "inbound" AVAILABLE = "available" @@ -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 diff --git a/src/federation_manager/core/config.py b/src/federation_manager/core/config.py index 80cd56f..45990d9 100644 --- a/src/federation_manager/core/config.py +++ b/src/federation_manager/core/config.py @@ -8,6 +8,7 @@ class Settings(BaseSettings): postgres_url: str = "postgresql+asyncpg://fm:fm@localhost:5433/fm_db" postgres_echo: bool = False + log_level: str = "INFO" nats_url: str = "nats://localhost:4222" srm_internal_url: str = "http://localhost:8081" event_consumer_durable: str = "fm-event-worker" diff --git a/src/federation_manager/core/logging.py b/src/federation_manager/core/logging.py new file mode 100644 index 0000000..7d82062 --- /dev/null +++ b/src/federation_manager/core/logging.py @@ -0,0 +1,27 @@ +import logging + +import structlog + + +def configure_logging(level: str = "INFO") -> None: + logging.basicConfig(format="%(message)s", level=getattr(logging, level.upper(), logging.INFO)) + structlog.configure( + processors=[ + structlog.contextvars.merge_contextvars, + structlog.processors.add_log_level, + structlog.processors.TimeStamper(fmt="iso", utc=True), + structlog.processors.StackInfoRenderer(), + structlog.processors.format_exc_info, + structlog.processors.JSONRenderer(), + ], + wrapper_class=structlog.make_filtering_bound_logger( + getattr(logging, level.upper(), logging.INFO) + ), + cache_logger_on_first_use=True, + ) + structlog.contextvars.bind_contextvars(service="fm") + + +def get_logger(name: str) -> structlog.BoundLogger: + logger: structlog.BoundLogger = structlog.get_logger(name) + return logger diff --git a/src/federation_manager/dependencies.py b/src/federation_manager/dependencies.py index 0173b97..46c47b1 100644 --- a/src/federation_manager/dependencies.py +++ b/src/federation_manager/dependencies.py @@ -21,6 +21,7 @@ from federation_manager.application.queries import InboundQueryService from federation_manager.core.config import get_settings from federation_manager.domain.ports import ( AgreementRepositoryPort, + DataBusHealthPort, DataBusPublisherPort, EwbiClientPort, FederationContextRepositoryPort, @@ -126,6 +127,11 @@ def get_inbound_federation_service( ) +def get_databus_health(request: Request) -> DataBusHealthPort: + publisher: DataBusHealthPort = request.app.state.command_publisher + return publisher + + def get_command_publisher(request: Request) -> DataBusPublisherPort: publisher: DataBusPublisherPort = request.app.state.command_publisher return publisher diff --git a/src/federation_manager/domain/ports.py b/src/federation_manager/domain/ports.py index c6f9791..0a79619 100644 --- a/src/federation_manager/domain/ports.py +++ b/src/federation_manager/domain/ports.py @@ -94,5 +94,9 @@ class TransactionRepositoryPort(Protocol): ) -> None: ... +class DataBusHealthPort(Protocol): + def is_connected(self) -> bool: ... + + class DataBusPublisherPort(Protocol): async def publish(self, subject: str, payload: dict[str, object]) -> None: ... diff --git a/src/federation_manager/domain/topology.py b/src/federation_manager/domain/topology.py new file mode 100644 index 0000000..cd1710b --- /dev/null +++ b/src/federation_manager/domain/topology.py @@ -0,0 +1,51 @@ +import re +from typing import Any + +REDACTED = "[redacted]" + +_STRIP_KEYS = frozenset( + { + "cluster_name", + "host", + "hostname", + "internal_ip", + "internal_trace_id", + "namespace", + "node_name", + "pod_name", + "pod_uid", + } +) + +_STRIP_PATTERNS = ( + re.compile(r"\b\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}\b"), + # compressed IPv6, or four-plus groups: a clock time has at most two colons + re.compile(r"\b[0-9a-f]{0,4}::[0-9a-f:]{0,29}", re.IGNORECASE), + re.compile(r"\b(?:[0-9a-f]{1,4}:){4,7}[0-9a-f]{1,4}\b", re.IGNORECASE), + re.compile(r"\b[a-z0-9-]+(?:\.[a-z0-9-]+)*\.svc\.cluster\.local\b", re.IGNORECASE), + re.compile(r"\bnamespace[-/][a-z0-9-]+\b", re.IGNORECASE), +) + + +def sanitize(payload: Any) -> Any: + if isinstance(payload, dict): + return { + key: REDACTED if _is_sensitive_key(key) else sanitize(value) + for key, value in payload.items() + } + if isinstance(payload, list): + return [sanitize(item) for item in payload] + if isinstance(payload, str): + return _redact(payload) + return payload + + +def _is_sensitive_key(key: str) -> bool: + normalized = key.replace("-", "_").lower() + return normalized in _STRIP_KEYS or normalized.startswith("internal_") + + +def _redact(value: str) -> str: + for pattern in _STRIP_PATTERNS: + value = pattern.sub(REDACTED, value) + return value diff --git a/src/federation_manager/main.py b/src/federation_manager/main.py index eb45dde..129059e 100644 --- a/src/federation_manager/main.py +++ b/src/federation_manager/main.py @@ -41,12 +41,15 @@ from federation_manager.application.federation import ( ) from federation_manager.contracts.srm import SUBJECT_OPERATION_COMPLETED from federation_manager.core.config import Settings, get_settings +from federation_manager.core.logging import configure_logging from federation_manager.domain.ports import EwbiClientPort +from federation_manager.middleware.topology_hiding import TopologyHidingMiddleware @asynccontextmanager async def default_lifespan(app: FastAPI) -> AsyncIterator[None]: settings = get_settings() + configure_logging(settings.log_level) engine = build_engine(settings.postgres_url, echo=settings.postgres_echo) await create_schema(engine) app.state.session_maker = build_session_maker(engine) @@ -131,6 +134,7 @@ def create_app(lifespan: Lifespan[FastAPI] | None = None) -> FastAPI: version=__version__, lifespan=lifespan or default_lifespan, ) + app.add_middleware(TopologyHidingMiddleware) register_exception_handlers(app) app.openapi = _openapi_with_problem_schemas(app) # type: ignore[method-assign] app.include_router(health_router) diff --git a/src/federation_manager/middleware/__init__.py b/src/federation_manager/middleware/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/federation_manager/middleware/topology_hiding.py b/src/federation_manager/middleware/topology_hiding.py new file mode 100644 index 0000000..6cfbe4b --- /dev/null +++ b/src/federation_manager/middleware/topology_hiding.py @@ -0,0 +1,49 @@ +import json +from collections.abc import Awaitable, Callable + +from starlette.middleware.base import BaseHTTPMiddleware +from starlette.requests import Request +from starlette.responses import Response + +from federation_manager.domain.ewbi import EWBI_BASE_PATH +from federation_manager.domain.topology import sanitize + +_JSON_TYPES = ("application/json", "application/problem+json") + + +class TopologyHidingMiddleware(BaseHTTPMiddleware): + async def dispatch( + self, request: Request, call_next: Callable[[Request], Awaitable[Response]] + ) -> Response: + response = await call_next(request) + if not request.url.path.startswith(EWBI_BASE_PATH): + return response + media_type = response.headers.get("content-type", "") + if not any(media_type.startswith(json_type) for json_type in _JSON_TYPES): + return response + + body_iterator = getattr(response, "body_iterator", None) + if body_iterator is None: + return response + chunks: list[bytes] = [] + async for chunk in body_iterator: + chunks.append(chunk.encode() if isinstance(chunk, str) else bytes(chunk)) + body = b"".join(chunks) + try: + payload = json.loads(body) + except ValueError: + return Response( + content=body, + status_code=response.status_code, + headers=dict(response.headers), + media_type=media_type, + ) + sanitized = json.dumps(sanitize(payload)).encode() + headers = dict(response.headers) + headers.pop("content-length", None) + return Response( + content=sanitized, + status_code=response.status_code, + headers=headers, + media_type=media_type, + ) diff --git a/tests/test_platform_health.py b/tests/test_platform_health.py new file mode 100644 index 0000000..b23bc94 --- /dev/null +++ b/tests/test_platform_health.py @@ -0,0 +1,71 @@ +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager +from typing import Any + +from fastapi import FastAPI +from fastapi.testclient import TestClient + +from federation_manager.dependencies import get_databus_health, get_session +from federation_manager.main import create_app + + +@asynccontextmanager +async def _no_infra(app: FastAPI) -> AsyncIterator[None]: + yield + + +class FakeSession: + def __init__(self, healthy: bool = True) -> None: + self.healthy = healthy + + async def execute(self, statement: Any) -> None: + if not self.healthy: + raise ConnectionError("postgres is down") + + +class FakeDataBus: + def __init__(self, connected: bool = True) -> None: + self.connected = connected + + def is_connected(self) -> bool: + return self.connected + + +def _client(postgres: bool = True, databus: bool = True) -> TestClient: + app = create_app(lifespan=_no_infra) + app.dependency_overrides[get_session] = lambda: FakeSession(postgres) + app.dependency_overrides[get_databus_health] = lambda: FakeDataBus(databus) + return TestClient(app) + + +def test_healthz_reports_liveness_without_touching_dependencies() -> None: + app = create_app(lifespan=_no_infra) + + response = TestClient(app).get("/healthz") + + assert response.status_code == 200 + assert response.json() == {"status": "ok"} + + +def test_readyz_is_ready_when_postgres_and_the_databus_answer() -> None: + response = _client().get("/readyz") + + assert response.status_code == 200 + assert response.json() == { + "status": "ready", + "checks": {"postgres": True, "databus": True}, + } + + +def test_readyz_is_503_when_postgres_is_unreachable() -> None: + response = _client(postgres=False).get("/readyz") + + assert response.status_code == 503 + assert response.json()["checks"] == {"postgres": False, "databus": True} + + +def test_readyz_is_503_when_the_databus_is_disconnected() -> None: + response = _client(databus=False).get("/readyz") + + assert response.status_code == 503 + assert response.json()["checks"] == {"postgres": True, "databus": False} diff --git a/tests/test_topology_hiding.py b/tests/test_topology_hiding.py new file mode 100644 index 0000000..b5e2249 --- /dev/null +++ b/tests/test_topology_hiding.py @@ -0,0 +1,127 @@ +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager +from typing import Any + +import pytest +from fastapi import FastAPI +from fastapi.testclient import TestClient + +from federation_manager.dependencies import ( + get_federation_context_repo, + get_jwt_validator, + get_partner_repo, +) +from federation_manager.domain.topology import REDACTED, sanitize +from federation_manager.main import create_app +from tests.fakes import FakeJwtValidator, InMemoryFederationContextRepo, InMemoryPartnerRepo + +LEAKY = { + "appInstanceInfo": { + "accesspointInfo": [{"internal_ip": "10.42.3.17", "host": "node-7.oop-prod"}], + "message": "scheduled on 10.42.3.17 in namespace/oop-prod", + }, + "endpoint": "video-es.oop-prod.svc.cluster.local", + "pod_uid": "0f2a-77", + "internal_trace_id": "abc123", + "zoneId": "f47ac10b-58cc-4372-a567-0e02b2c3d479", +} + + +@pytest.mark.parametrize( + ("payload", "expected"), + [ + ({"internal_ip": "10.0.0.1"}, {"internal_ip": REDACTED}), + ({"cluster-name": "oop-prod"}, {"cluster-name": REDACTED}), + ({"internal_anything": "x"}, {"internal_anything": REDACTED}), + ({"note": "runs at 10.42.3.17"}, {"note": f"runs at {REDACTED}"}), + ( + {"note": "video.oop-prod.svc.cluster.local"}, + {"note": REDACTED}, + ), + ({"note": "in namespace/oop-prod now"}, {"note": f"in {REDACTED} now"}), + ({"note": "fe80::1 link local"}, {"note": f"{REDACTED} link local"}), + ], +) +def test_sensitive_values_are_redacted(payload: dict[str, Any], expected: dict[str, Any]) -> None: + assert sanitize(payload) == expected + + +def test_legitimate_data_survives() -> None: + payload = { + "zoneId": "f47ac10b-58cc-4372-a567-0e02b2c3d479", + "appInstIdentifier": "82f97ea55f15469281df467c01c774c8", + "federationContextId": "fed-ctx-1", + "appVersion": "1.2.0", + "count": 3, + "ready": True, + "missing": None, + } + + assert sanitize(payload) == payload + + +def test_a_uuid_is_not_mistaken_for_an_ipv6_address() -> None: + payload = {"id": "f47ac10b-58cc-4372-a567-0e02b2c3d479"} + + assert sanitize(payload) == payload + + +def test_nested_structures_are_walked() -> None: + cleaned = sanitize(LEAKY) + + assert cleaned["pod_uid"] == REDACTED + assert cleaned["internal_trace_id"] == REDACTED + assert cleaned["endpoint"] == REDACTED + accesspoint = cleaned["appInstanceInfo"]["accesspointInfo"][0] + assert accesspoint == {"internal_ip": REDACTED, "host": REDACTED} + assert "10.42.3.17" not in cleaned["appInstanceInfo"]["message"] + assert "oop-prod" not in cleaned["appInstanceInfo"]["message"] + assert cleaned["zoneId"] == LEAKY["zoneId"] + + +@asynccontextmanager +async def _no_infra(app: FastAPI) -> AsyncIterator[None]: + yield + + +def _app_leaking_on(path: str) -> FastAPI: + app = create_app(lifespan=_no_infra) + + @app.get(path) + def leak() -> dict[str, Any]: + return LEAKY + + return app + + +def test_ewbi_responses_are_sanitised_by_the_middleware() -> None: + app = _app_leaking_on("/operatorplatform/federation/v1/leaky") + + body = TestClient(app).get("/operatorplatform/federation/v1/leaky").json() + + assert body["pod_uid"] == REDACTED + assert body["endpoint"] == REDACTED + assert "10.42.3.17" not in str(body) + assert body["zoneId"] == LEAKY["zoneId"] + + +def test_internal_responses_are_left_alone() -> None: + app = _app_leaking_on("/internal/leaky") + + body = TestClient(app).get("/internal/leaky").json() + + assert body == LEAKY + + +def test_problem_responses_keep_their_media_type_and_headers() -> None: + app = create_app(lifespan=_no_infra) + app.dependency_overrides[get_partner_repo] = lambda: InMemoryPartnerRepo([]) + app.dependency_overrides[get_jwt_validator] = lambda: FakeJwtValidator({}) + app.dependency_overrides[get_federation_context_repo] = InMemoryFederationContextRepo + + response = TestClient(app).get("/operatorplatform/federation/v1/fed-ctx-1/health") + + assert response.status_code == 401 + assert response.headers["content-type"].startswith("application/problem+json") + assert response.headers["WWW-Authenticate"] == 'Bearer scope="fed-mgmt"' + assert response.json()["type"] == "urn:oop:ewbi:error:authentication-failed" diff --git a/uv.lock b/uv.lock index caee66c..2f3ec50 100644 --- a/uv.lock +++ b/uv.lock @@ -630,6 +630,7 @@ dependencies = [ { name = "pydantic-settings" }, { name = "pyjwt", extra = ["crypto"] }, { name = "sqlalchemy", extra = ["asyncio"] }, + { name = "structlog" }, ] [package.optional-dependencies] @@ -659,6 +660,7 @@ requires-dist = [ { name = "pyyaml", marker = "extra == 'dev'", specifier = ">=6.0" }, { name = "ruff", marker = "extra == 'dev'", specifier = ">=0.6" }, { name = "sqlalchemy", extras = ["asyncio"], specifier = ">=2.0" }, + { name = "structlog", specifier = ">=25.1" }, ] provides-extras = ["dev"] @@ -1701,6 +1703,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/c8/cb/6a6a47d5b464bd08695d254f3da6e7986cc70c9fa5d778eda57538edfe56/starlette-1.6.0-py3-none-any.whl", hash = "sha256:a86dd39d14bb45f85a3d18525215a9ef0cfd1f192ac793220e72598c90335f0c", size = 75969, upload-time = "2026-08-08T18:27:56.196Z" }, ] +[[package]] +name = "structlog" +version = "26.1.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/5e/89/b4a0bcfdf4f71a3dea31379f095929613d7e4528a0996bca6aa964cd0dca/structlog-26.1.0.tar.gz", hash = "sha256:f63a716cbd1b1291cf7661de7794b455acfa4c43c5bcf1630e6ad5ddc1adb3b7", size = 1459881, upload-time = "2026-06-06T07:33:39.348Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/a9/18/489c97b834dfff9cf2fc2507cede4bcd4b11e67f84bc462acd1992496f86/structlog-26.1.0-py3-none-any.whl", hash = "sha256:e081a26d6c373e6d201eca24eede26d8ffab07f88f477822e679183428d3d91e", size = 73764, upload-time = "2026-06-06T07:33:38.046Z" }, +] + [[package]] name = "typer" version = "0.27.2" -- GitLab