Commit 540d7ec2 authored by Sergio Gimenez's avatar Sergio Gimenez
Browse files

test(fm): two-stack federation over real OAuth2

Two FMs as separate uvicorn processes, each with its own database,
plus the real Keycloak. FM-A bootstraps with
FM_BOOTSTRAP_FEDERATION=true, fetches a client-credentials token as
originating-op-1 and calls FM-B's POST /partner over HTTP. FM-B
validates the JWT against Keycloak's JWKS, matches azp to its
partner row and mints a context id. Asserts both databases hold the
same federationContextId and that FM-B's health endpoint answers
AVAILABLE with a real token; a token from the wrong client cannot
read the context.

Separate processes because get_settings() is an lru_cache singleton,
so two apps cannot hold different config in one process. Requires
FM_ALLOW_INSECURE_PARTNER_ENDPOINTS since the local stack has no
TLS. Requires the compose stack (Keycloak + Postgres) running.
parent 896d6b5f
Loading
Loading
Loading
Loading
+257 −0
Original line number Diff line number Diff line
import asyncio
import os
import socket
import subprocess
import sys
import time
from collections.abc import Iterator
from dataclasses import dataclass
from pathlib import Path
from uuid import UUID, uuid4

import httpx
import pytest
from sqlalchemy import delete, insert
from sqlalchemy.ext.asyncio import create_async_engine

from federation_manager.adapters.database.core import (
    build_engine,
    build_session_maker,
    create_schema,
)
from federation_manager.adapters.database.federation_context_repo import (
    PostgresFederationContextRepo,
)
from federation_manager.adapters.database.tables import federation_contexts, partner_ops
from federation_manager.domain.models import FederationContext

pytestmark = pytest.mark.integration

PG_ROOT = os.getenv("FM_POSTGRES_ROOT", "postgresql+asyncpg://fm:fm@localhost:5433")
DB_A, DB_B = "fm_db", "fm_db_b"
ISSUER = os.getenv("FM_KEYCLOAK_ISSUER", "http://localhost:8090/realms/federation")
TOKEN_ENDPOINT = f"{ISSUER}/protocol/openid-connect/token"
CLIENT_A, SECRET_A = "originating-op-1", "dd7vNwFqjNpYwaghlEwMbw10g0klWDHb"
CLIENT_B, SECRET_B = "originating-op-2", "2mhznERfWclLDuVojY77Lp4Qd2r4e8Ms"


@dataclass(frozen=True)
class Stacks:
    port_a: int
    port_b: int
    partner_a: UUID
    partner_b: UUID


def _free_port() -> int:
    with socket.socket() as s:
        s.bind(("127.0.0.1", 0))
        port: int = s.getsockname()[1]
        return port


def _require(url: str, what: str) -> None:
    try:
        httpx.get(url, timeout=2.0)
    except httpx.HTTPError:
        pytest.skip(
            f"{what} not reachable at {url}; run docker compose -f docker-compose.dev.yaml up -d"
        )


async def _create_database(name: str) -> None:
    engine = create_async_engine(f"{PG_ROOT}/postgres", isolation_level="AUTOCOMMIT")
    async with engine.connect() as conn:
        exists = await conn.exec_driver_sql(f"SELECT 1 FROM pg_database WHERE datname = '{name}'")
        if exists.first() is None:
            await conn.exec_driver_sql(f'CREATE DATABASE "{name}"')
    await engine.dispose()


def _start_fm(port: int, env: dict[str, str]) -> subprocess.Popen[bytes]:
    process = subprocess.Popen(
        [
            sys.executable,
            "-m",
            "uvicorn",
            "federation_manager.main:app",
            "--host",
            "127.0.0.1",
            "--port",
            str(port),
            "--log-level",
            "warning",
        ],
        env={**os.environ, **env},
        stdout=subprocess.PIPE,
        stderr=subprocess.STDOUT,
    )
    deadline = time.monotonic() + 40
    while time.monotonic() < deadline:
        if process.poll() is not None:
            output = process.stdout.read().decode() if process.stdout else ""
            raise AssertionError(f"FM on port {port} exited early:\n{output}")
        try:
            if httpx.get(f"http://127.0.0.1:{port}/healthz", timeout=1.0).status_code == 200:
                return process
        except httpx.HTTPError:
            time.sleep(0.3)
    process.kill()
    raise AssertionError(f"FM on port {port} did not become ready")


def _fm_env(database: str, port: int, federation_id: str, **extra: str) -> dict[str, str]:
    return {
        "FM_POSTGRES_URL": f"{PG_ROOT}/{database}",
        "FM_KEYCLOAK_ISSUER": ISSUER,
        "FM_FEDERATION_ID": federation_id,
        "FM_COUNTRY_CODE": "ES",
        "FM_MCC": "214",
        "FM_MNCS": '["07"]',
        "FM_PARTNER_STATUS_LINK": f"http://127.0.0.1:{port}/operatorplatform/federation/v1/partner-status",
        "FM_ALLOW_INSECURE_PARTNER_ENDPOINTS": "true",
        **extra,
    }


def _access_token(client_id: str, secret: str) -> str:
    response = httpx.post(
        TOKEN_ENDPOINT,
        data={
            "grant_type": "client_credentials",
            "client_id": client_id,
            "client_secret": secret,
            "scope": "fed-mgmt",
        },
        timeout=10.0,
    )
    response.raise_for_status()
    token: str = response.json()["access_token"]
    return token


@pytest.fixture
def stacks(tmp_path: Path) -> Iterator[Stacks]:
    _require(f"{ISSUER}/.well-known/openid-configuration", "Keycloak")

    asyncio.run(_create_database(DB_B))

    port_a, port_b = _free_port(), _free_port()
    secret_file = tmp_path / "partner-b-secret"
    secret_file.write_text(SECRET_A, encoding="utf-8")

    partner_b_id, partner_a_id = uuid4(), uuid4()
    mcc_mnc_b, mcc_mnc_a = uuid4().hex[:10], uuid4().hex[:10]

    async def seed() -> None:
        for database, values in (
            (
                DB_A,
                {
                    "id": partner_b_id,
                    "mcc_mnc": mcc_mnc_b,
                    "oauth2_client_id": CLIENT_B,
                    "base_url": f"http://127.0.0.1:{port_b}",
                    "our_client_id": CLIENT_A,
                    "our_client_secret_ref": str(secret_file),
                    "token_endpoint": TOKEN_ENDPOINT,
                    "status": "active",
                },
            ),
            (
                DB_B,
                {
                    "id": partner_a_id,
                    "mcc_mnc": mcc_mnc_a,
                    "oauth2_client_id": CLIENT_A,
                    "base_url": f"http://127.0.0.1:{port_a}",
                    "status": "active",
                },
            ),
        ):
            engine = build_engine(f"{PG_ROOT}/{database}")
            await create_schema(engine)
            async with build_session_maker(engine)() as session:
                await session.execute(insert(partner_ops).values(**values))
                await session.commit()
            await engine.dispose()

    asyncio.run(seed())

    processes = []
    try:
        processes.append(_start_fm(port_b, _fm_env(DB_B, port_b, "op-b")))
        processes.append(
            _start_fm(
                port_a,
                _fm_env(DB_A, port_a, "op-a", FM_BOOTSTRAP_FEDERATION="true"),
            )
        )
        yield Stacks(port_a=port_a, port_b=port_b, partner_a=partner_a_id, partner_b=partner_b_id)
    finally:
        for process in processes:
            process.kill()
            process.wait(timeout=10)

        async def cleanup() -> None:
            for database, partner_id in ((DB_A, partner_b_id), (DB_B, partner_a_id)):
                engine = build_engine(f"{PG_ROOT}/{database}")
                async with build_session_maker(engine)() as session:
                    await session.execute(
                        delete(federation_contexts).where(
                            federation_contexts.c.partner_op_id == partner_id
                        )
                    )
                    await session.execute(delete(partner_ops).where(partner_ops.c.id == partner_id))
                    await session.commit()
                await engine.dispose()

        asyncio.run(cleanup())


async def _context(database: str, partner_id: UUID, direction: str) -> FederationContext | None:
    engine = build_engine(f"{PG_ROOT}/{database}")
    async with build_session_maker(engine)() as session:
        repo = PostgresFederationContextRepo(session)
        context = await (
            repo.find_active_outbound(partner_id)
            if direction == "outbound"
            else repo.find_active_inbound(partner_id)
        )
    await engine.dispose()
    return context


def test_two_stacks_federate_over_real_oauth2(stacks: Stacks) -> None:
    outbound = asyncio.run(_context(DB_A, stacks.partner_b, "outbound"))
    inbound = asyncio.run(_context(DB_B, stacks.partner_a, "inbound"))

    assert outbound is not None, "FM-A stored no outbound context; bootstrap federation failed"
    assert inbound is not None, "FM-B stored no inbound context"
    assert outbound.federation_context_id == inbound.federation_context_id
    assert outbound.status == "available"
    assert inbound.status == "available"

    context_id = outbound.federation_context_id
    health = httpx.get(
        f"http://127.0.0.1:{stacks.port_b}/operatorplatform/federation/v1/{context_id}/health",
        headers={"Authorization": f"Bearer {_access_token(CLIENT_A, SECRET_A)}"},
        timeout=10.0,
    )

    assert health.status_code == 200
    assert health.json()["federationHealthStatus"]["federationStatus"] == "AVAILABLE"


def test_partner_b_rejects_an_unknown_client(stacks: Stacks) -> None:
    outbound = asyncio.run(_context(DB_A, stacks.partner_b, "outbound"))
    assert outbound is not None
    context_id = outbound.federation_context_id
    health = httpx.get(
        f"http://127.0.0.1:{stacks.port_b}/operatorplatform/federation/v1/{context_id}/health",
        headers={"Authorization": f"Bearer {_access_token(CLIENT_B, SECRET_B)}"},
        timeout=10.0,
    )

    assert health.status_code in (401, 404)
    assert "federationHealthStatus" not in health.text