Commit b037ae1c authored by Sergio Gimenez's avatar Sergio Gimenez
Browse files

add federation contexts

parent b299e266
Loading
Loading
Loading
Loading
+49 −0
Original line number Diff line number Diff line
from uuid import UUID

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

from federation_manager.adapters.database.tables import federation_contexts as contexts
from federation_manager.domain.models import FederationContext


class PostgresFederationContextRepo:
    def __init__(self, session: AsyncSession) -> None:
        self._session = session

    async def find_active_outbound(self, partner_id: UUID) -> FederationContext | None:
        return await self._one(
            select(contexts)
            .where(
                contexts.c.partner_op_id == partner_id,
                contexts.c.direction == "outbound",
                contexts.c.status == "available",
            )
            .order_by(contexts.c.created_at.desc())
            .limit(1)
        )

    async def find_inbound(
        self, partner_id: UUID, federation_context_id: str
    ) -> FederationContext | None:
        return await self._one(
            select(contexts).where(
                contexts.c.partner_op_id == partner_id,
                contexts.c.direction == "inbound",
                contexts.c.federation_context_id == federation_context_id,
            )
        )

    async def _one(self, stmt: Select[tuple[object, ...]]) -> FederationContext | None:
        row = (await self._session.execute(stmt)).one_or_none()
        if row is None:
            return None
        return FederationContext(
            id=row.id,
            partner_op_id=row.partner_op_id,
            direction=row.direction,
            federation_context_id=row.federation_context_id,
            status=row.status,
            agreement_id=row.agreement_id,
            status_callback_url=row.status_callback_url,
        )
+24 −2
Original line number Diff line number Diff line
@@ -54,6 +54,26 @@ federation_agreements = Table(
    Index("idx_federation_agreements_status_validity", "status", "valid_from", "valid_until"),
)

federation_contexts = Table(
    "federation_contexts",
    metadata,
    Column("id", PGUUID(as_uuid=True), primary_key=True, server_default=func.gen_random_uuid()),
    Column("partner_op_id", PGUUID(as_uuid=True), ForeignKey("partner_ops.id"), nullable=False),
    Column("agreement_id", PGUUID(as_uuid=True), ForeignKey("federation_agreements.id")),
    Column("direction", String(10), nullable=False),
    Column("federation_context_id", String(255), nullable=False),
    Column("status_callback_url", Text),
    Column("status", String(20), nullable=False),
    Column("created_at", DateTime(timezone=True), nullable=False, server_default=func.now()),
    Column("updated_at", DateTime(timezone=True), nullable=False, server_default=func.now()),
    UniqueConstraint(
        "partner_op_id",
        "direction",
        "federation_context_id",
        name="uq_federation_contexts_partner_direction_id",
    ),
)

routing_rules = Table(
    "routing_rules",
    metadata,
@@ -77,11 +97,12 @@ federation_transactions = Table(
    Column("id", PGUUID(as_uuid=True), primary_key=True, server_default=func.gen_random_uuid()),
    Column("partner_op_id", PGUUID(as_uuid=True), ForeignKey("partner_ops.id"), nullable=False),
    Column("agreement_id", PGUUID(as_uuid=True), ForeignKey("federation_agreements.id")),
    Column("federation_context_row_id", PGUUID(as_uuid=True), ForeignKey("federation_contexts.id")),
    Column("direction", String(10), nullable=False),
    Column("federation_operation_id", PGUUID(as_uuid=True)),
    Column("operation_id", PGUUID(as_uuid=True)),
    Column("correlation_id", PGUUID(as_uuid=True)),
    Column("federation_correlation_id", String(255)),
    Column("external_txn_id", String(255)),
    Column("api_type", String(100), nullable=False),
    Column("status", String(20), nullable=False, server_default="pending"),
    Column("request_summary", JSONB, nullable=False),
@@ -97,8 +118,9 @@ federation_transactions = Table(
        unique=True,
        postgresql_where=text("federation_operation_id IS NOT NULL"),
    ),
    Index("idx_fed_tx_context", "federation_context_row_id"),
    Index("idx_fed_tx_operation", "operation_id"),
    Index("idx_fed_tx_status", "status"),
    Index("idx_fed_tx_started", "started_at"),
    Index("idx_fed_tx_fed_correlation", "federation_correlation_id"),
    Index("idx_fed_tx_external_txn", "external_txn_id"),
)
+2 −1
Original line number Diff line number Diff line
@@ -19,11 +19,12 @@ class PostgresTransactionRepo:
                id=transaction.id,
                partner_op_id=transaction.partner_op_id,
                agreement_id=transaction.agreement_id,
                federation_context_row_id=transaction.federation_context_row_id,
                direction=transaction.direction,
                federation_operation_id=transaction.federation_operation_id,
                operation_id=transaction.operation_id,
                correlation_id=transaction.correlation_id,
                federation_correlation_id=transaction.federation_correlation_id,
                external_txn_id=transaction.external_txn_id,
                api_type=transaction.api_type,
                status=transaction.status,
                request_summary=transaction.request_summary,
+10 −0
Original line number Diff line number Diff line
@@ -5,6 +5,9 @@ from fastapi import Depends, Request
from sqlalchemy.ext.asyncio import AsyncSession

from federation_manager.adapters.database.agreement_repo import PostgresAgreementRepo
from federation_manager.adapters.database.federation_context_repo import (
    PostgresFederationContextRepo,
)
from federation_manager.adapters.database.partner_repo import PostgresPartnerRepo
from federation_manager.adapters.database.routing_repo import PostgresRoutingRuleRepo
from federation_manager.adapters.database.transaction_repo import PostgresTransactionRepo
@@ -14,6 +17,7 @@ from federation_manager.application.outbound import OutboundFederationService
from federation_manager.domain.ports import (
    AgreementRepositoryPort,
    EwbiClientPort,
    FederationContextRepositoryPort,
    JwtValidatorPort,
    PartnerRepositoryPort,
    PartnerTokenProviderPort,
@@ -40,6 +44,12 @@ def get_agreement_repo(
    return PostgresAgreementRepo(session)


def get_federation_context_repo(
    session: Annotated[AsyncSession, Depends(get_session)],
) -> FederationContextRepositoryPort:
    return PostgresFederationContextRepo(session)


def get_routing_rule_repo(
    session: Annotated[AsyncSession, Depends(get_session)],
) -> RoutingRuleRepositoryPort:
+17 −1
Original line number Diff line number Diff line
@@ -75,6 +75,21 @@ class RoutingRule:
    is_active: bool = True


@dataclass(frozen=True)
class FederationContext:
    id: UUID
    partner_op_id: UUID
    direction: str  # inbound: we issued the context id | outbound: the partner issued it
    federation_context_id: str  # opaque OPG.04 FederationContextId, not necessarily a UUID
    # OPG.04 `Status` enum, lowercase, plus `terminated` once DeleteFederationDetails has run
    status: str  # available | locked | not_available | temporary_failure | failed | terminated
    agreement_id: UUID | None = None
    status_callback_url: str | None = None

    def is_active(self) -> bool:
        return self.status == "available"


@dataclass
class FederationTransaction:
    id: UUID
@@ -85,7 +100,8 @@ class FederationTransaction:
    request_summary: dict[str, object]
    started_at: datetime
    agreement_id: UUID | None = None
    federation_correlation_id: str | None = None  # unused: no correlation header on EWBI
    federation_context_row_id: UUID | None = None
    external_txn_id: str | None = None  # OPG.04 txnIdentifier / apiTxnId when the operation has one
    federation_operation_id: UUID | None = None
    operation_id: UUID | None = None  # never returned to partners
    correlation_id: UUID | None = None  # never returned to partners
Loading