Skip to content

Routing

RabbitRouter

RabbitRouter

Modular router — groups related routes with shared defaults.

Use routers to organize handlers by domain:

orders_router = RabbitRouter(prefix="orders", exchange="orders-exchange")

@orders_router.subscriber(queue="orders-queue", routing_key="created")
def handle_order(order: Order) -> None:
    ...

Include in broker/app:

broker.include_router(orders_router)

Source code in src/rabbitkit/core/router.py
class RabbitRouter:
    """Modular router — groups related routes with shared defaults.

    Use routers to organize handlers by domain:
    ```python
    orders_router = RabbitRouter(prefix="orders", exchange="orders-exchange")

    @orders_router.subscriber(queue="orders-queue", routing_key="created")
    def handle_order(order: Order) -> None:
        ...
    ```

    Include in broker/app:
    ```python
    broker.include_router(orders_router)
    ```
    """

    def __init__(
        self,
        prefix: str = "",
        exchange: RabbitExchange | str | None = None,
        middlewares: list[BaseMiddleware] | None = None,
        serializer: Serializer[Any] | None = None,
        tags: frozenset[str] | set[str] | None = None,
    ) -> None:
        self._prefix = prefix
        self._default_exchange = self._normalize_exchange(exchange)
        self._middlewares = middlewares or []
        self._serializer = serializer
        self._tags = frozenset(tags) if tags else frozenset()
        self._registry = SubscriberRegistry()

    @property
    def prefix(self) -> str:
        return self._prefix

    @property
    def routes(self) -> list[RouteDefinition]:
        return self._registry.routes

    def subscriber(
        self,
        queue: RabbitQueue | str,
        exchange: RabbitExchange | str | None = None,
        routing_key: str = "",
        ack_policy: AckPolicy = AckPolicy.AUTO,
        middlewares: list[BaseMiddleware] | None = None,
        serializer: Serializer[Any] | None = None,
        retry: RetryConfig | RetryDisabled | None = None,
        tags: frozenset[str] | set[str] | None = None,
        description: str = "",
        name: str | None = None,
        prefetch_count: int | None = None,
        filter_fn: Callable[[RabbitMessage], bool] | None = None,
        reject_without_dlx: str | None = None,
    ) -> Callable[[Callable[..., Any]], Callable[..., Any]]:
        """Register a subscriber on this router.

        Applies router-level defaults:
        - prefix is prepended to routing_key
        - exchange falls back to router's default exchange
        - middlewares are merged (router-level + route-level)
        - serializer falls back to router's serializer
        - tags are merged (router-level + route-level)
        """
        # Apply prefix to routing key
        if self._prefix and routing_key:
            effective_rk = f"{self._prefix}.{routing_key}"
        else:
            effective_rk = routing_key or self._prefix

        # Fall back to router's default exchange
        effective_exchange = exchange if exchange is not None else self._default_exchange

        # Merge middlewares (router-level first, then route-level)
        route_mw = middlewares or []
        effective_mw = self._middlewares + route_mw

        # Fall back to router's serializer
        effective_serializer = serializer if serializer is not None else self._serializer

        # Merge tags
        route_tags = frozenset(tags) if tags else frozenset()
        effective_tags = self._tags | route_tags

        return self._registry.subscriber(
            queue=queue,
            exchange=effective_exchange,
            routing_key=effective_rk,
            ack_policy=ack_policy,
            middlewares=effective_mw,
            serializer=effective_serializer,
            retry=retry,
            tags=effective_tags,
            description=description,
            name=name,
            prefetch_count=prefetch_count,
            filter_fn=filter_fn,
            reject_without_dlx=reject_without_dlx,
        )

    def publisher(
        self,
        exchange: RabbitExchange | str | None = None,
        routing_key: str = "",
    ) -> Callable[[Callable[..., Any]], Callable[..., Any]]:
        """Register a result publisher on this router."""
        return self._registry.publisher(exchange=exchange, routing_key=routing_key)

    def _normalize_exchange(self, exchange: RabbitExchange | str | None) -> RabbitExchange | None:
        """Normalize exchange to RabbitExchange or None."""
        if isinstance(exchange, str):
            return RabbitExchange(name=exchange)
        return exchange

Methods:

subscriber(queue: RabbitQueue | str, exchange: RabbitExchange | str | None = None, routing_key: str = '', ack_policy: AckPolicy = AckPolicy.AUTO, middlewares: list[BaseMiddleware] | None = None, serializer: Serializer[Any] | None = None, retry: RetryConfig | RetryDisabled | None = None, tags: frozenset[str] | set[str] | None = None, description: str = '', name: str | None = None, prefetch_count: int | None = None, filter_fn: Callable[[RabbitMessage], bool] | None = None, reject_without_dlx: str | None = None) -> Callable[[Callable[..., Any]], Callable[..., Any]]

Register a subscriber on this router.

Applies router-level defaults: - prefix is prepended to routing_key - exchange falls back to router's default exchange - middlewares are merged (router-level + route-level) - serializer falls back to router's serializer - tags are merged (router-level + route-level)

Source code in src/rabbitkit/core/router.py
def subscriber(
    self,
    queue: RabbitQueue | str,
    exchange: RabbitExchange | str | None = None,
    routing_key: str = "",
    ack_policy: AckPolicy = AckPolicy.AUTO,
    middlewares: list[BaseMiddleware] | None = None,
    serializer: Serializer[Any] | None = None,
    retry: RetryConfig | RetryDisabled | None = None,
    tags: frozenset[str] | set[str] | None = None,
    description: str = "",
    name: str | None = None,
    prefetch_count: int | None = None,
    filter_fn: Callable[[RabbitMessage], bool] | None = None,
    reject_without_dlx: str | None = None,
) -> Callable[[Callable[..., Any]], Callable[..., Any]]:
    """Register a subscriber on this router.

    Applies router-level defaults:
    - prefix is prepended to routing_key
    - exchange falls back to router's default exchange
    - middlewares are merged (router-level + route-level)
    - serializer falls back to router's serializer
    - tags are merged (router-level + route-level)
    """
    # Apply prefix to routing key
    if self._prefix and routing_key:
        effective_rk = f"{self._prefix}.{routing_key}"
    else:
        effective_rk = routing_key or self._prefix

    # Fall back to router's default exchange
    effective_exchange = exchange if exchange is not None else self._default_exchange

    # Merge middlewares (router-level first, then route-level)
    route_mw = middlewares or []
    effective_mw = self._middlewares + route_mw

    # Fall back to router's serializer
    effective_serializer = serializer if serializer is not None else self._serializer

    # Merge tags
    route_tags = frozenset(tags) if tags else frozenset()
    effective_tags = self._tags | route_tags

    return self._registry.subscriber(
        queue=queue,
        exchange=effective_exchange,
        routing_key=effective_rk,
        ack_policy=ack_policy,
        middlewares=effective_mw,
        serializer=effective_serializer,
        retry=retry,
        tags=effective_tags,
        description=description,
        name=name,
        prefetch_count=prefetch_count,
        filter_fn=filter_fn,
        reject_without_dlx=reject_without_dlx,
    )

publisher(exchange: RabbitExchange | str | None = None, routing_key: str = '') -> Callable[[Callable[..., Any]], Callable[..., Any]]

Register a result publisher on this router.

Source code in src/rabbitkit/core/router.py
def publisher(
    self,
    exchange: RabbitExchange | str | None = None,
    routing_key: str = "",
) -> Callable[[Callable[..., Any]], Callable[..., Any]]:
    """Register a result publisher on this router."""
    return self._registry.publisher(exchange=exchange, routing_key=routing_key)

RabbitExchange

RabbitExchange dataclass

Exchange declaration model.

Source code in src/rabbitkit/core/topology.py
@dataclass(frozen=True, slots=True)
class RabbitExchange:
    """Exchange declaration model."""

    name: str
    type: ExchangeType = ExchangeType.DIRECT
    durable: bool = True
    auto_delete: bool = False
    passive: bool = False
    internal: bool = False
    arguments: dict[str, Any] = field(default_factory=dict)
    bind_to: str | None = None
    bind_arguments: dict[str, Any] = field(default_factory=dict)
    routing_key: str = ""

    def __post_init__(self) -> None:
        self.validate()

    def validate(self) -> None:
        """Validate exchange configuration."""
        if not self.name and self.type != ExchangeType.DIRECT:
            msg = "Non-default exchanges must have a name"
            raise TopologyValidationError(msg)
        if self.internal and self.auto_delete:
            msg = "Internal exchanges cannot be auto_delete (they are never published to directly)."
            raise TopologyValidationError(msg)
        validate_amqp_shortstr("Exchange name", self.name)
        validate_amqp_shortstr("Exchange routing_key", self.routing_key)

    def to_declare_kwargs(self) -> dict[str, Any]:
        """Build exchange_declare kwargs for pika/aio-pika."""
        return {
            "exchange": self.name,
            "exchange_type": self.type.value,
            "durable": self.durable,
            "auto_delete": self.auto_delete,
            "passive": self.passive,
            "internal": self.internal,
            "arguments": self.arguments or None,
        }

    def to_bind_kwargs(self) -> dict[str, Any] | None:
        """Build exchange_bind kwargs. Returns None if no binding."""
        if self.bind_to is None:
            return None
        return {
            "destination": self.name,
            "source": self.bind_to,
            "routing_key": self.routing_key,
            "arguments": self.bind_arguments or None,
        }

Methods:

validate() -> None

Validate exchange configuration.

Source code in src/rabbitkit/core/topology.py
def validate(self) -> None:
    """Validate exchange configuration."""
    if not self.name and self.type != ExchangeType.DIRECT:
        msg = "Non-default exchanges must have a name"
        raise TopologyValidationError(msg)
    if self.internal and self.auto_delete:
        msg = "Internal exchanges cannot be auto_delete (they are never published to directly)."
        raise TopologyValidationError(msg)
    validate_amqp_shortstr("Exchange name", self.name)
    validate_amqp_shortstr("Exchange routing_key", self.routing_key)

to_declare_kwargs() -> dict[str, Any]

Build exchange_declare kwargs for pika/aio-pika.

Source code in src/rabbitkit/core/topology.py
def to_declare_kwargs(self) -> dict[str, Any]:
    """Build exchange_declare kwargs for pika/aio-pika."""
    return {
        "exchange": self.name,
        "exchange_type": self.type.value,
        "durable": self.durable,
        "auto_delete": self.auto_delete,
        "passive": self.passive,
        "internal": self.internal,
        "arguments": self.arguments or None,
    }

to_bind_kwargs() -> dict[str, Any] | None

Build exchange_bind kwargs. Returns None if no binding.

Source code in src/rabbitkit/core/topology.py
def to_bind_kwargs(self) -> dict[str, Any] | None:
    """Build exchange_bind kwargs. Returns None if no binding."""
    if self.bind_to is None:
        return None
    return {
        "destination": self.name,
        "source": self.bind_to,
        "routing_key": self.routing_key,
        "arguments": self.bind_arguments or None,
    }

RabbitQueue

RabbitQueue dataclass

Queue declaration model with type-specific validation.

Source code in src/rabbitkit/core/topology.py
@dataclass(frozen=True, slots=True)
class RabbitQueue:
    """Queue declaration model with type-specific validation."""

    name: str
    durable: bool = True
    exclusive: bool = False
    passive: bool = False
    auto_delete: bool = False
    routing_key: str = ""
    bind_arguments: dict[str, Any] = field(default_factory=dict)
    queue_type: QueueType = QueueType.CLASSIC

    # DLQ
    dead_letter_exchange: str | None = None
    dead_letter_routing_key: str | None = None

    # Limits
    message_ttl: int | None = None  # ms
    max_length: int | None = None
    max_length_bytes: int | None = None
    # x-consumer-timeout (ms) — per-queue override of the server's
    # consumer ack timeout (server default: 30 minutes). If a delivered
    # message stays unacked past it, RabbitMQ force-closes the consumer's
    # channel. The server does not advertise its limit to clients, so if a
    # handler can legitimately hold a message longer than 30 minutes, raise
    # it HERE at declaration time (RabbitMQ >= 3.12; classic/quorum only).
    consumer_timeout: int | None = None

    # Classic-only
    lazy: bool = False  # x-queue-mode: lazy (classic only)
    max_priority: int | None = None  # classic only (0-255)

    # Quorum-specific
    delivery_limit: int | None = None  # x-delivery-limit (quorum only)
    # x-single-active-consumer — NOT quorum-specific: valid on classic
    # (RabbitMQ 3.8+) and quorum queues alike. The standard active/standby
    # pattern: N replicas subscribe, the broker delivers to exactly one and
    # fails over automatically when it dies.
    single_active_consumer: bool = False

    # Overflow
    overflow: str | None = None  # "drop-head" | "reject-publish" | "reject-publish-dlx"

    # Expiry
    expires: int | None = None  # ms — auto-delete after idle

    # Extra arguments (escape hatch)
    arguments: dict[str, Any] = field(default_factory=dict)

    def __post_init__(self) -> None:
        self.validate()

    def validate(self) -> None:
        """Enforce queue-type-specific constraints.

        Raises TopologyValidationError (a ValueError subclass) for
        invalid combinations.
        Uses warnings.warn() for unusual but legal combinations.
        """
        if not self.name:
            msg = "Queue name is required"
            raise TopologyValidationError(msg)
        validate_amqp_shortstr("Queue name", self.name)
        validate_amqp_shortstr("Queue routing_key", self.routing_key)

        # Quorum constraints
        if self.queue_type == QueueType.QUORUM:
            if not self.durable:
                msg = "Quorum queues must be durable"
                raise TopologyValidationError(msg)
            if self.exclusive:
                msg = "Quorum queues cannot be exclusive"
                raise TopologyValidationError(msg)
            if self.lazy:
                msg = "Quorum queues do not support lazy mode (x-queue-mode)"
                raise TopologyValidationError(msg)
            if self.max_priority is not None:
                msg = "Quorum queues do not support priorities"
                raise TopologyValidationError(msg)

        # Stream constraints
        if self.queue_type == QueueType.STREAM:
            if not self.durable:
                msg = "Stream queues must be durable"
                raise TopologyValidationError(msg)
            if self.exclusive:
                msg = "Stream queues cannot be exclusive"
                raise TopologyValidationError(msg)
            if self.lazy:
                msg = "Stream queues do not support lazy mode"
                raise TopologyValidationError(msg)
            if self.max_priority is not None:
                msg = "Stream queues do not support priorities"
                raise TopologyValidationError(msg)
            if self.message_ttl is not None:
                msg = "Stream queues do not support message TTL"
                raise TopologyValidationError(msg)
            if self.consumer_timeout is not None:
                msg = "Stream queues do not support consumer_timeout (per-message ack timeouts do not apply to streams)"
                raise TopologyValidationError(msg)

        if self.consumer_timeout is not None and self.consumer_timeout <= 0:
            msg = f"Queue '{self.name}': consumer_timeout must be a positive number of milliseconds"
            raise TopologyValidationError(msg)

        # Classic constraints
        if self.queue_type == QueueType.CLASSIC:
            if self.delivery_limit is not None:
                msg = "Classic queues do not support delivery_limit (quorum only)"
                raise TopologyValidationError(msg)

        # Warnings for unusual combos
        if self.lazy:
            warnings.warn(
                f"Queue '{self.name}': lazy=True sets the deprecated x-queue-mode=lazy "
                "argument. RabbitMQ >=3.12 defaults classic queues to CQv2, which already "
                "keeps message bodies out of memory in a lazy-like manner -- x-queue-mode "
                "is a silent no-op there. On RabbitMQ <3.12 (or a classic queue explicitly "
                "downgraded to v1) it still has effect. If you're targeting >=3.12, drop "
                "lazy=True; the default queue behavior already covers this.",
                UserWarning,
                stacklevel=2,
            )

        if self.auto_delete and self.durable:
            warnings.warn(
                f"Queue '{self.name}': auto_delete=True with durable=True is unusual — "
                "the queue will be deleted when the last consumer disconnects, "
                "despite being durable",
                UserWarning,
                stacklevel=2,
            )

        if self.passive and any(
            [
                self.lazy,
                self.max_priority is not None,
                self.delivery_limit is not None,
                self.message_ttl is not None,
                self.max_length is not None,
                self.consumer_timeout is not None,
            ]
        ):
            warnings.warn(
                f"Queue '{self.name}': passive=True with creation-only options set — "
                "these options are ignored for passive declarations",
                UserWarning,
                stacklevel=2,
            )

    def to_declare_kwargs(self) -> dict[str, Any]:
        """Build queue_declare kwargs with merged x-arguments."""
        args: dict[str, Any] = {}

        # Queue type
        args["x-queue-type"] = self.queue_type.value

        # DLQ
        if self.dead_letter_exchange is not None:
            args["x-dead-letter-exchange"] = self.dead_letter_exchange
        if self.dead_letter_routing_key is not None:
            args["x-dead-letter-routing-key"] = self.dead_letter_routing_key

        # Limits
        if self.message_ttl is not None:
            args["x-message-ttl"] = self.message_ttl
        if self.max_length is not None:
            args["x-max-length"] = self.max_length
        if self.max_length_bytes is not None:
            args["x-max-length-bytes"] = self.max_length_bytes
        if self.consumer_timeout is not None:
            args["x-consumer-timeout"] = self.consumer_timeout

        # Classic-only
        if self.lazy:
            args["x-queue-mode"] = "lazy"
        if self.max_priority is not None:
            args["x-max-priority"] = self.max_priority

        # Quorum-specific
        if self.delivery_limit is not None:
            args["x-delivery-limit"] = self.delivery_limit
        if self.single_active_consumer:
            args["x-single-active-consumer"] = True

        # Overflow
        if self.overflow is not None:
            args["x-overflow"] = self.overflow

        # Expiry
        if self.expires is not None:
            args["x-expires"] = self.expires

        # Merge user-provided arguments (escape hatch takes precedence)
        args.update(self.arguments)

        return {
            "queue": self.name,
            "durable": self.durable,
            "exclusive": self.exclusive,
            "auto_delete": self.auto_delete,
            "passive": self.passive,
            "arguments": args,
        }

    def to_bind_kwargs(self, exchange: str) -> dict[str, Any]:
        """Build queue_bind kwargs."""
        return {
            "queue": self.name,
            "exchange": exchange,
            "routing_key": self.routing_key,
            "arguments": self.bind_arguments or None,
        }

Methods:

validate() -> None

Enforce queue-type-specific constraints.

Raises TopologyValidationError (a ValueError subclass) for invalid combinations. Uses warnings.warn() for unusual but legal combinations.

Source code in src/rabbitkit/core/topology.py
def validate(self) -> None:
    """Enforce queue-type-specific constraints.

    Raises TopologyValidationError (a ValueError subclass) for
    invalid combinations.
    Uses warnings.warn() for unusual but legal combinations.
    """
    if not self.name:
        msg = "Queue name is required"
        raise TopologyValidationError(msg)
    validate_amqp_shortstr("Queue name", self.name)
    validate_amqp_shortstr("Queue routing_key", self.routing_key)

    # Quorum constraints
    if self.queue_type == QueueType.QUORUM:
        if not self.durable:
            msg = "Quorum queues must be durable"
            raise TopologyValidationError(msg)
        if self.exclusive:
            msg = "Quorum queues cannot be exclusive"
            raise TopologyValidationError(msg)
        if self.lazy:
            msg = "Quorum queues do not support lazy mode (x-queue-mode)"
            raise TopologyValidationError(msg)
        if self.max_priority is not None:
            msg = "Quorum queues do not support priorities"
            raise TopologyValidationError(msg)

    # Stream constraints
    if self.queue_type == QueueType.STREAM:
        if not self.durable:
            msg = "Stream queues must be durable"
            raise TopologyValidationError(msg)
        if self.exclusive:
            msg = "Stream queues cannot be exclusive"
            raise TopologyValidationError(msg)
        if self.lazy:
            msg = "Stream queues do not support lazy mode"
            raise TopologyValidationError(msg)
        if self.max_priority is not None:
            msg = "Stream queues do not support priorities"
            raise TopologyValidationError(msg)
        if self.message_ttl is not None:
            msg = "Stream queues do not support message TTL"
            raise TopologyValidationError(msg)
        if self.consumer_timeout is not None:
            msg = "Stream queues do not support consumer_timeout (per-message ack timeouts do not apply to streams)"
            raise TopologyValidationError(msg)

    if self.consumer_timeout is not None and self.consumer_timeout <= 0:
        msg = f"Queue '{self.name}': consumer_timeout must be a positive number of milliseconds"
        raise TopologyValidationError(msg)

    # Classic constraints
    if self.queue_type == QueueType.CLASSIC:
        if self.delivery_limit is not None:
            msg = "Classic queues do not support delivery_limit (quorum only)"
            raise TopologyValidationError(msg)

    # Warnings for unusual combos
    if self.lazy:
        warnings.warn(
            f"Queue '{self.name}': lazy=True sets the deprecated x-queue-mode=lazy "
            "argument. RabbitMQ >=3.12 defaults classic queues to CQv2, which already "
            "keeps message bodies out of memory in a lazy-like manner -- x-queue-mode "
            "is a silent no-op there. On RabbitMQ <3.12 (or a classic queue explicitly "
            "downgraded to v1) it still has effect. If you're targeting >=3.12, drop "
            "lazy=True; the default queue behavior already covers this.",
            UserWarning,
            stacklevel=2,
        )

    if self.auto_delete and self.durable:
        warnings.warn(
            f"Queue '{self.name}': auto_delete=True with durable=True is unusual — "
            "the queue will be deleted when the last consumer disconnects, "
            "despite being durable",
            UserWarning,
            stacklevel=2,
        )

    if self.passive and any(
        [
            self.lazy,
            self.max_priority is not None,
            self.delivery_limit is not None,
            self.message_ttl is not None,
            self.max_length is not None,
            self.consumer_timeout is not None,
        ]
    ):
        warnings.warn(
            f"Queue '{self.name}': passive=True with creation-only options set — "
            "these options are ignored for passive declarations",
            UserWarning,
            stacklevel=2,
        )

to_declare_kwargs() -> dict[str, Any]

Build queue_declare kwargs with merged x-arguments.

Source code in src/rabbitkit/core/topology.py
def to_declare_kwargs(self) -> dict[str, Any]:
    """Build queue_declare kwargs with merged x-arguments."""
    args: dict[str, Any] = {}

    # Queue type
    args["x-queue-type"] = self.queue_type.value

    # DLQ
    if self.dead_letter_exchange is not None:
        args["x-dead-letter-exchange"] = self.dead_letter_exchange
    if self.dead_letter_routing_key is not None:
        args["x-dead-letter-routing-key"] = self.dead_letter_routing_key

    # Limits
    if self.message_ttl is not None:
        args["x-message-ttl"] = self.message_ttl
    if self.max_length is not None:
        args["x-max-length"] = self.max_length
    if self.max_length_bytes is not None:
        args["x-max-length-bytes"] = self.max_length_bytes
    if self.consumer_timeout is not None:
        args["x-consumer-timeout"] = self.consumer_timeout

    # Classic-only
    if self.lazy:
        args["x-queue-mode"] = "lazy"
    if self.max_priority is not None:
        args["x-max-priority"] = self.max_priority

    # Quorum-specific
    if self.delivery_limit is not None:
        args["x-delivery-limit"] = self.delivery_limit
    if self.single_active_consumer:
        args["x-single-active-consumer"] = True

    # Overflow
    if self.overflow is not None:
        args["x-overflow"] = self.overflow

    # Expiry
    if self.expires is not None:
        args["x-expires"] = self.expires

    # Merge user-provided arguments (escape hatch takes precedence)
    args.update(self.arguments)

    return {
        "queue": self.name,
        "durable": self.durable,
        "exclusive": self.exclusive,
        "auto_delete": self.auto_delete,
        "passive": self.passive,
        "arguments": args,
    }

to_bind_kwargs(exchange: str) -> dict[str, Any]

Build queue_bind kwargs.

Source code in src/rabbitkit/core/topology.py
def to_bind_kwargs(self, exchange: str) -> dict[str, Any]:
    """Build queue_bind kwargs."""
    return {
        "queue": self.name,
        "exchange": exchange,
        "routing_key": self.routing_key,
        "arguments": self.bind_arguments or None,
    }

TopologyMode

TopologyMode

Bases: str, Enum

Topology declaration modes.

See Contract 6 in the plan for precedence rules: - AUTO_DECLARE: declare exchanges/queues/bindings on startup - PASSIVE_ONLY: all declarations use passive=True - MANUAL: skip all topology operations

Source code in src/rabbitkit/core/types.py
class TopologyMode(str, Enum):
    """Topology declaration modes.

    See Contract 6 in the plan for precedence rules:
    - AUTO_DECLARE: declare exchanges/queues/bindings on startup
    - PASSIVE_ONLY: all declarations use passive=True
    - MANUAL: skip all topology operations
    """

    AUTO_DECLARE = "auto_declare"
    PASSIVE_ONLY = "passive_only"
    MANUAL = "manual"

TopologyDispatcher

TopologyDispatcher

Resolves TopologyMode into a concrete TopoAction per entity.

The transport computes to_declare_kwargs() and performs the actual (sync or async) channel call based on the returned TopoAction, keeping all transport-specific I/O out of this class.

Source code in src/rabbitkit/core/topology_dispatch.py
class TopologyDispatcher:
    """Resolves ``TopologyMode`` into a concrete ``TopoAction`` per entity.

    The transport computes ``to_declare_kwargs()`` and performs the actual
    (sync or async) channel call based on the returned ``TopoAction``,
    keeping all transport-specific I/O out of this class.
    """

    def __init__(self, mode: TopologyMode) -> None:
        self._mode = mode

    @property
    def mode(self) -> TopologyMode:
        """The topology mode this dispatcher was configured with."""
        return self._mode

    def exchange_action(self, exchange: RabbitExchange) -> TopoAction:
        """Action to take for ``declare_exchange(exchange)``."""
        if self._mode == TopologyMode.MANUAL:
            return TopoAction.SKIP
        if self._mode == TopologyMode.PASSIVE_ONLY or exchange.passive:
            return TopoAction.PASSIVE
        return TopoAction.DECLARE

    def queue_action(self, queue: RabbitQueue) -> TopoAction:
        """Action to take for ``declare_queue(queue)``."""
        if self._mode == TopologyMode.MANUAL:
            return TopoAction.SKIP
        if self._mode == TopologyMode.PASSIVE_ONLY or queue.passive:
            return TopoAction.PASSIVE
        return TopoAction.DECLARE

    def binding_action(self) -> TopoAction:
        """Action to take for ``bind_queue`` / ``bind_exchange``.

        Bindings have no passive variant — they are skipped only under
        ``MANUAL`` and performed otherwise.
        """
        if self._mode == TopologyMode.MANUAL:
            return TopoAction.SKIP
        return TopoAction.DECLARE

Attributes

mode: TopologyMode property

The topology mode this dispatcher was configured with.

Methods:

exchange_action(exchange: RabbitExchange) -> TopoAction

Action to take for declare_exchange(exchange).

Source code in src/rabbitkit/core/topology_dispatch.py
def exchange_action(self, exchange: RabbitExchange) -> TopoAction:
    """Action to take for ``declare_exchange(exchange)``."""
    if self._mode == TopologyMode.MANUAL:
        return TopoAction.SKIP
    if self._mode == TopologyMode.PASSIVE_ONLY or exchange.passive:
        return TopoAction.PASSIVE
    return TopoAction.DECLARE

queue_action(queue: RabbitQueue) -> TopoAction

Action to take for declare_queue(queue).

Source code in src/rabbitkit/core/topology_dispatch.py
def queue_action(self, queue: RabbitQueue) -> TopoAction:
    """Action to take for ``declare_queue(queue)``."""
    if self._mode == TopologyMode.MANUAL:
        return TopoAction.SKIP
    if self._mode == TopologyMode.PASSIVE_ONLY or queue.passive:
        return TopoAction.PASSIVE
    return TopoAction.DECLARE

binding_action() -> TopoAction

Action to take for bind_queue / bind_exchange.

Bindings have no passive variant — they are skipped only under MANUAL and performed otherwise.

Source code in src/rabbitkit/core/topology_dispatch.py
def binding_action(self) -> TopoAction:
    """Action to take for ``bind_queue`` / ``bind_exchange``.

    Bindings have no passive variant — they are skipped only under
    ``MANUAL`` and performed otherwise.
    """
    if self._mode == TopologyMode.MANUAL:
        return TopoAction.SKIP
    return TopoAction.DECLARE

TopoAction

Bases: Enum

What the transport should do for a topology entity.

Source code in src/rabbitkit/core/topology_dispatch.py
class TopoAction(Enum):
    """What the transport should do for a topology entity."""

    SKIP = auto()  # TopologyMode.MANUAL — do nothing
    PASSIVE = auto()  # PASSIVE_ONLY or entity.passive — passive existence check
    DECLARE = auto()  # active declaration