Skip to content

Thread-Safe Connection Adapter

Connection wraps SelectConnection and runs its IOLoop in a dedicated background thread. All channel operations are routed through add_callback_threadsafe so that socket I/O stays confined to a single thread. External threads can call any Channel method concurrently without coordination.

This eliminates the need to manually schedule callbacks via add_callback_threadsafe when publishing or acknowledging from multiple threads.

When to use Connection

Use Connection when your application needs to publish or consume messages from multiple threads. It provides a blocking, synchronous API similar to BlockingConnection but with full thread safety:

  • Multiple threads can call basic_publish simultaneously
  • Consumer callbacks run on a dedicated worker thread, not the IOLoop thread, so slow processing does not stall heartbeats
  • Channel methods (basic_ack, queue_declare, etc.) are safe to call from within consumer callbacks

Basic usage

import threading
import pika
from pika.adapters.thread_safe_connection import Connection

conn = Connection(pika.ConnectionParameters('localhost'))
ch = conn.channel()

ch.queue_declare('work')

# Publish from multiple threads simultaneously
for i in range(5):
    threading.Thread(
        target=ch.basic_publish,
        kwargs=dict(exchange='', routing_key='work', body=f'msg-{i}'.encode()),
    ).start()

conn.close()

Consumer threading model

The on_message_callback registered with basic_consume is dispatched on a per-channel worker thread. This means:

  • Blocking operations (database writes, HTTP calls, even queue_declare or basic_qos) are safe inside delivery callbacks
  • The IOLoop thread is never starved by slow consumer processing, so heartbeats are always sent on time
  • Messages are delivered to the callback in order (single worker thread per channel)
  • All Channel methods (basic_ack, basic_nack, basic_reject, basic_publish) are safe to call from within the callback

Work queue bounds and back-pressure

Each worker thread consumes from a bounded queue holding at most work_queue_maxsize pending dispatches (default 1000, matching the RabbitMQ Java client). When a slow callback lets the queue fill up, further enqueues block the IOLoop thread, which stops reading from the socket and applies back-pressure to the broker through TCP flow control instead of buffering events in memory without limit.

Heartbeats are not sent while the IOLoop is blocked, so the broker's heartbeat timeout bounds how long a wedged consumer can stall the connection before the broker closes it.

Both knobs are constructor parameters on Connection:

conn = Connection(
    pika.ConnectionParameters('localhost'),
    work_queue_maxsize=1000,     # 0 means unbounded
    work_queue_put_timeout=30.0, # seconds; must be positive, default 30 s
)

work_queue_put_timeout defaults to 30 seconds and must be a positive number; None (infinite) is rejected with ValueError because it lets a full queue stall or deadlock the IOLoop thread. Set an override comfortably below the negotiated heartbeat-death window (twice the heartbeat interval) so a stall surfaces as WorkQueueFullError rather than an opaque broker-side close. When the queue stays full for that long, pika.exceptions.WorkQueueFullError is raised on the IOLoop thread, tearing the connection down rather than silently dropping the event (which would lose an auto-acked delivery or force redelivery of a manually-acked one).

Comparison with other adapters

BlockingConnection SelectConnection + add_callback_threadsafe Connection
Thread safety Not thread-safe Manual callback scheduling required Fully thread-safe
Consumer threading Same thread as IOLoop Same thread as IOLoop Dedicated worker thread
Heartbeat risk Slow consumers can stall heartbeats Slow consumers can stall heartbeats IOLoop runs independently
API style Synchronous, blocking Asynchronous, callback-based Synchronous, blocking

Class Reference

Connection

Pika connection that is safe to use from multiple threads.

.. note:: Each instance starts a background thread named pika-ioloop-N (where N is a per-process sequence number) so that multiple connections are distinguishable in stack traces and thread listings.

Internally wraps :class:~pika.adapters.SelectConnection and runs its IOLoop in a single dedicated background thread (the IOLoop thread). All channel operations submitted from external threads are routed through :meth:add_callback_threadsafe so that _tx_buffers has exactly one owner — the IOLoop thread.

Usage::

conn = Connection(pika.ConnectionParameters('localhost'))
ch = conn.channel()

# safe to call from any number of threads simultaneously
threading.Thread(target=ch.basic_publish,
                 kwargs=dict(exchange='', routing_key='q',
                             body='hello')).start()

conn.close()

Can also be used as a context manager::

with Connection(pika.ConnectionParameters('localhost')) as conn:
    ch = conn.channel()
    ch.basic_publish(exchange='', routing_key='q', body='hello')
PARAMETER DESCRIPTION
parameters

Connection parameters.

TYPE: Parameters

on_open_error_callback

Called in the IOLoop thread if the connection cannot be established. Signature: on_open_error_callback(connection, exception)

DEFAULT: None

on_close_callback

Called in the IOLoop thread when the connection is closed. Signature: on_close_callback(connection, reason)

DEFAULT: None

timeout

Seconds to wait for the AMQP connection to be established. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely (relying on the socket timeout in parameters). On timeout the IOLoop is stopped and :class:TimeoutError is raised.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

work_queue_maxsize

Maximum number of user-callback dispatches (deliveries, publisher confirms, returned messages, cancel and blocked/unblocked notifications) that may be pending on each worker's queue. Defaults to :data:DEFAULT_WORK_QUEUE_MAXSIZE (1000, matching the RabbitMQ Java client). While a queue is full the IOLoop blocks, applying back-pressure to the broker via TCP. Heartbeats are not sent while the IOLoop is blocked, so the broker's heartbeat timeout bounds how long a wedged consumer can stall the connection. Pass 0 for an unbounded queue.

TYPE: int DEFAULT: DEFAULT_WORK_QUEUE_MAXSIZE

work_queue_put_timeout

Seconds to wait for space on a full work queue. Defaults to :data:DEFAULT_WORK_QUEUE_PUT_TIMEOUT (30 s); must be a positive number, and None (infinite) is rejected because it lets a full queue stall or deadlock the IOLoop thread. Set an override comfortably below the negotiated heartbeat-death window (twice the heartbeat interval) so a stall surfaces here rather than as an opaque broker-side close. On timeout :class:pika.exceptions.WorkQueueFullError is raised on the IOLoop thread, closing the connection rather than silently dropping the event (which would lose an auto-acked delivery or force redelivery of a manually-acked one).

TYPE: float DEFAULT: DEFAULT_WORK_QUEUE_PUT_TIMEOUT

RAISES DESCRIPTION
ValueError

if work_queue_put_timeout is None or is not a positive number.

Exception

if the connection cannot be established.

TimeoutError

if timeout expires before the connection opens.

Source code in pika/adapters/thread_safe_connection.py
class Connection:
    """Pika connection that is safe to use from multiple threads.

    .. note:: Each instance starts a background thread named
        ``pika-ioloop-N`` (where *N* is a per-process sequence number)
        so that multiple connections are distinguishable in stack traces
        and thread listings.

    Internally wraps :class:`~pika.adapters.SelectConnection` and runs its
    IOLoop in a single dedicated background thread (the *IOLoop thread*).
    All channel operations submitted from external threads are routed
    through :meth:`add_callback_threadsafe` so that ``_tx_buffers`` has
    exactly one owner — the IOLoop thread.

    Usage::

        conn = Connection(pika.ConnectionParameters('localhost'))
        ch = conn.channel()

        # safe to call from any number of threads simultaneously
        threading.Thread(target=ch.basic_publish,
                         kwargs=dict(exchange='', routing_key='q',
                                     body='hello')).start()

        conn.close()

    Can also be used as a context manager::

        with Connection(pika.ConnectionParameters('localhost')) as conn:
            ch = conn.channel()
            ch.basic_publish(exchange='', routing_key='q', body='hello')

    :param parameters: Connection parameters.
    :type parameters: pika.connection.Parameters
    :param on_open_error_callback:
        Called in the IOLoop thread if the connection cannot be established.
        Signature: ``on_open_error_callback(connection, exception)``
    :param on_close_callback:
        Called in the IOLoop thread when the connection is closed.
        Signature: ``on_close_callback(connection, reason)``
    :param timeout: Seconds to wait for the AMQP connection
        to be established.  Defaults to :data:`DEFAULT_RPC_TIMEOUT` (10 s).
        Pass ``None`` to wait indefinitely (relying on the socket timeout
        in *parameters*).  On timeout the IOLoop is stopped and
        :class:`TimeoutError` is raised.
    :param work_queue_maxsize: Maximum number of user-callback dispatches
        (deliveries, publisher confirms, returned messages, cancel and
        blocked/unblocked notifications) that may be pending on each
        worker's queue.  Defaults to :data:`DEFAULT_WORK_QUEUE_MAXSIZE`
        (1000, matching the RabbitMQ Java client).  While a queue is full
        the IOLoop blocks, applying back-pressure to the broker via TCP.
        Heartbeats are not sent while the IOLoop is blocked, so the
        broker's heartbeat timeout bounds how long a wedged consumer can
        stall the connection.  Pass ``0`` for an unbounded queue.
    :param work_queue_put_timeout: Seconds to wait for space on a full
        work queue.  Defaults to :data:`DEFAULT_WORK_QUEUE_PUT_TIMEOUT`
        (30 s); must be a positive number, and ``None`` (infinite) is
        rejected because it lets a full queue stall or deadlock the IOLoop
        thread.  Set an override comfortably below the negotiated
        heartbeat-death window (twice the heartbeat interval) so a stall
        surfaces here rather than as an opaque broker-side close.  On
        timeout :class:`pika.exceptions.WorkQueueFullError` is raised on
        the IOLoop thread, closing the connection rather than silently
        dropping the event (which would lose an auto-acked delivery or
        force redelivery of a manually-acked one).
    :raises ValueError: if *work_queue_put_timeout* is ``None`` or is not a
        positive number.
    :raises Exception: if the connection cannot be established.
    :raises TimeoutError: if *timeout* expires before the connection opens.
    """

    _instance_counter = itertools.count(1)

    def __init__(
            self,
            parameters,
            on_open_error_callback=None,
            on_close_callback=None,
            timeout: float | None = DEFAULT_RPC_TIMEOUT,
            work_queue_maxsize: int = DEFAULT_WORK_QUEUE_MAXSIZE,
            work_queue_put_timeout: float = DEFAULT_WORK_QUEUE_PUT_TIMEOUT
    ) -> None:
        self._user_on_open_error_callback = on_open_error_callback
        self._user_on_close_callback = on_close_callback
        self._work_queue_maxsize = work_queue_maxsize
        self._work_queue_put_timeout = _validate_put_timeout(
            work_queue_put_timeout)

        self._connect_error = None
        self._connected_event = threading.Event()

        self._instance_id = next(self._instance_counter)

        self._channel_waiters_lock = threading.Lock()
        self._closed_reason: BaseException | None = None
        # The traceback _closed_reason carried when it was recorded.  Raising
        # the shared instance appends frames to it permanently, so every raise
        # site restores this snapshot; see _with_close_traceback.
        self._closed_reason_tb: TracebackType | None = None
        self._blocking_waiters: list[tuple[threading.Event,
                                           list[BaseException | None]]] = []
        self._channels: list[Channel] = []

        # Single-worker pool for connection-level event callbacks
        # (Connection.Blocked / Unblocked).  Keeps user code off the
        # IOLoop thread so a slow listener cannot stall heartbeats.
        self._connection_work_pool = _BoundedWorkPool(
            maxsize=work_queue_maxsize,
            put_timeout=work_queue_put_timeout,
            thread_name=f'pika-conn-{self._instance_id}',
        )
        self._connection_pool_shutdown = False

        try:
            self._connection = SelectConnection(
                parameters=parameters,
                on_open_callback=self._on_connection_open,
                on_open_error_callback=self._on_connection_open_error,
                on_close_callback=self._on_connection_closed,
            )
        except Exception:
            # SelectConnection construction failed before the IOLoop thread
            # started, so _run_ioloop's cleanup tail will never run.  Shut
            # the connection-event pool down here to avoid leaking its
            # worker thread.
            self._shutdown_connection_pool()
            raise

        def _run_ioloop() -> None:
            try:
                self._connection.ioloop.start()
            except Exception as exc:
                # An unhandled exception in a callback killed the IOLoop.
                # Wake every blocked caller so they do not hang forever.
                LOGGER.exception('IOLoop thread crashed')
                with self._channel_waiters_lock:
                    if self._closed_reason is None:
                        self._record_closed_reason(exc)
                    for evt, err in self._blocking_waiters:
                        if err[0] is None:
                            err[0] = self._closed_reason
                        evt.set()
                    self._blocking_waiters.clear()
                # Wake __init__ if it crashed before the connection opened.
                if not self._connected_event.is_set():
                    self._connect_error = exc
                    self._connected_event.set()
            # IOLoop has exited - safe to block waiting for pool workers.
            self._shutdown_all_consumer_pools()
            self._shutdown_connection_pool()

        self._ioloop_thread = threading.Thread(
            target=_run_ioloop,
            name=f'pika-ioloop-{self._instance_id}',
            daemon=True,
        )
        self._ioloop_thread.start()

        # Block the calling thread until the connection is open or fails.
        if not self._connected_event.wait(timeout=timeout):
            self._connect_error = TimeoutError(
                f'connection attempt timed out after {timeout} seconds')
            self._connection.ioloop.add_callback_threadsafe(
                self._connection.ioloop.stop)
            # Wait for the IOLoop thread to exit so it can shut down the
            # consumer and connection pools before we raise.  Otherwise
            # those pools' worker threads outlive __init__.
            self._ioloop_thread.join(timeout=timeout)
            raise self._connect_error
        if self._connect_error is not None:
            # IOLoop thread already exited (open-error path stops the
            # IOLoop) but its cleanup tail may still be running; wait
            # briefly so pool shutdown completes before we raise.
            self._ioloop_thread.join(timeout=timeout)
            raise self._connect_error

    # ------------------------------------------------------------------
    # IOLoop-thread callbacks
    # ------------------------------------------------------------------

    def _on_connection_open(self, _connection) -> None:
        self._connected_event.set()

    def _on_connection_open_error(self, _connection, error) -> None:
        self._connect_error = error
        # Stop the IOLoop so the background thread can exit cleanly.
        self._connection.ioloop.stop()
        self._connected_event.set()
        if self._user_on_open_error_callback:
            self._user_on_open_error_callback(_connection, error)

    def _on_connection_closed(self, _connection, reason) -> None:
        # Connection is gone - stop the IOLoop so the thread exits.
        # Pool shutdown happens after ioloop.start() returns in _run_ioloop
        # (calling shutdown(wait=True) here would deadlock the IOLoop thread
        # if a pool worker is mid-callback).
        self._connection.ioloop.stop()
        with self._channel_waiters_lock:
            self._record_closed_reason(reason)
            for evt, err in self._blocking_waiters:
                err[0] = reason
                evt.set()
            self._blocking_waiters.clear()
        if self._user_on_close_callback:
            self._user_on_close_callback(_connection, reason)

    def _shutdown_all_consumer_pools(self,
                                     timeout: float | None = None) -> None:
        """
        Shut down all tracked channel consumer pools.

        :param timeout: Total seconds to wait across all channel pools. ``None`` waits indefinitely
            (the clean-shutdown default). A finite budget is split across channels so the whole
            sweep stays bounded.
        """
        with self._channel_waiters_lock:
            channels = list(self._channels)
        if not channels:
            return
        if timeout is None:
            for ch in channels:
                ch._shutdown_pool()
            return
        # Signal every pool to start draining before joining any of them, so
        # all workers drain concurrently.  Joining one at a time without this
        # would let an early wedged channel burn the whole budget while later
        # channels' workers had not even been told to stop, so those healthy
        # workers would be abandoned at timeout=0 with a spurious warning.
        for ch in channels:
            ch._signal_pool_shutdown()
        # Share the remaining budget across channels using a common deadline so
        # a slow early channel cannot consume the entire allowance and starve
        # later ones of their chance to drain.
        deadline = time.monotonic() + timeout
        for ch in channels:
            ch._shutdown_pool(timeout=max(0.0, deadline - time.monotonic()))

    def _is_pool_worker_thread(self) -> bool:
        """
        True if the calling thread is one of this connection's pool worker threads.

        Covers every channel's consumer worker plus the connection's own event-callback worker
        (:attr:`_connection_work_pool`).  :meth:`close` uses this to detect that it is running from
        inside a delivery/event callback so it can avoid joining :attr:`_ioloop_thread`: that thread's
        post-``ioloop.start()`` cleanup tail (:meth:`_shutdown_all_consumer_pools`,
        :meth:`_shutdown_connection_pool`) joins every pool worker in turn, including this one, so
        joining it here would deadlock each thread waiting on the other (issue #1686).

        Reading each pool's ``_thread`` without its lock is safe here: a worker only executes this
        check from within its own run loop, by which point ``_thread`` was already assigned (under
        the pool's lock, before ``start()``) by the thread that created it - that assignment
        happens-before anything the worker itself observes.
        """
        current = threading.current_thread()
        if current is self._connection_work_pool._thread:
            return True
        with self._channel_waiters_lock:
            channels = list(self._channels)
        return any(current is ch._consumer_work_pool._thread for ch in channels)

    def _shutdown_connection_pool(self, timeout: float | None = None) -> None:
        """
        Shut down the connection-level event-callback pool.

        :param timeout: Seconds to wait for the worker to drain and exit. ``None`` waits
            indefinitely (the clean-shutdown default). A wedged worker outliving a finite *timeout*
            is left running (daemon thread) rather than stalling the caller; a warning is logged.
        """
        # Atomically claim the shutdown so the IOLoop tail and a concurrent
        # close() do not both run the join.  The join stays outside the lock so
        # it cannot stall other channel operations for the whole timeout.
        with self._channel_waiters_lock:
            if self._connection_pool_shutdown:
                return
            self._connection_pool_shutdown = True
        if not self._connection_work_pool.shutdown(wait=True, timeout=timeout):
            LOGGER.warning(
                'Connection %s event work pool did not drain within '
                '%s seconds; abandoning its worker thread', self._instance_id,
                timeout)

    # ------------------------------------------------------------------
    # Public API
    # ------------------------------------------------------------------

    def channel(self, timeout: float | None = DEFAULT_RPC_TIMEOUT) -> Channel:
        """
        Open a new channel and return a :class:`Channel`.

        Blocks the calling thread until the channel is open.  The returned channel's methods are
        safe to call from any thread.

        :param timeout: Seconds to wait for Channel.OpenOk. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :raises Exception: if the connection is closed before the channel opens.
        :raises TimeoutError: if *timeout* expires before the channel opens.
        """
        ready = threading.Event()
        result = [None]
        error: list[BaseException | None] = [None]

        with self._channel_waiters_lock:
            reason = self._closed_reason
            if reason is not None:
                raise _with_close_traceback(reason, self._closed_reason_tb)
            self._blocking_waiters.append((ready, error))

        def _open() -> None:

            def _on_open(ch) -> None:
                result[0] = ch
                ready.set()

            try:
                self._connection.channel(on_open_callback=_on_open)
            except Exception as exc:
                error[0] = exc
                ready.set()

        self._connection.ioloop.add_callback_threadsafe(_open)

        try:
            if not ready.wait(timeout=timeout):
                raise TimeoutError(
                    f'channel open timed out after {timeout} seconds')
        finally:
            with self._channel_waiters_lock:
                try:
                    self._blocking_waiters.remove((ready, error))
                except ValueError:
                    pass

        if error[0] is not None:
            raise _reraisable(error[0], self._closed_reason,
                              self._closed_reason_tb)

        # Race guard: the connection may have closed on the IOLoop thread
        # after _on_open fired but before we woke up.  Without this check,
        # _shutdown_all_consumer_pools() would have already run and the
        # newly-constructed channel's consumer pool would never be reached.
        with self._channel_waiters_lock:
            reason = self._closed_reason
            if reason is not None:
                raise _with_close_traceback(reason, self._closed_reason_tb)
            ch = Channel(result[0],
                         self,
                         work_queue_maxsize=self._work_queue_maxsize,
                         work_queue_put_timeout=self._work_queue_put_timeout)
            self._channels.append(ch)
        # Register the broker-close hook on the IOLoop thread, where the raw
        # channel's callback stack is safe to mutate.  A broker close arriving
        # in the brief gap before this runs simply falls back to the
        # connection-shutdown sweep, which still drops the channel from
        # tracking.
        self._schedule_unchecked(
            lambda: ch._channel.add_on_close_callback(ch._on_broker_close))
        return ch

    def close(self, timeout: float | None = 10) -> None:
        """
        Close the connection and block until the IOLoop thread exits.

        Schedules a clean Connection.Close handshake and waits up to *timeout* seconds for the
        IOLoop thread to finish.  If the thread is still alive after the timeout (e.g. the broker
        never sends Connection.CloseOk), the IOLoop is force-stopped via ``add_callback_threadsafe``
        and joined once more under a second *timeout* budget.  A worker still draining when that
        budget expires is left to finish on its own (it is a daemon thread), so a wedged delivery
        callback cannot hang ``close`` indefinitely.

        Safe to call from any thread including the IOLoop thread itself (e.g. from within a channel
        callback).  When called from the IOLoop thread the close is initiated synchronously and the
        method returns immediately without joining.

        Also safe to call from a delivery or connection-event callback running on one of this
        connection's pool worker threads (e.g. a consumer's ``on_message_callback``).  The IOLoop
        thread's own post-close cleanup joins every pool worker, including whichever one is running
        the callback that called ``close()``, so joining :attr:`_ioloop_thread` from there would
        deadlock the two threads against each other. In that case the close is scheduled and this
        method returns immediately without joining; the worker exits normally once the callback
        returns.  Because this path does not join the IOLoop thread, *timeout* is not applied on it:
        the close is best-effort and an unresponsive broker is not force-stopped from here, so a
        healthy handshake is assumed (see #1706).

        Calling ``close()`` on an already-closed connection is a no-op.

        :param timeout: Seconds to wait for a clean close before force-stopping the IOLoop. Defaults
            to 10 seconds, which is sufficient for a healthy broker on any reasonable network. The
            pathological force path may wait up to twice this (once for the clean join, once for the
            forced join). Pass ``None`` to wait indefinitely. Ignored when ``close()`` is called
            from one of this connection's pool worker threads (see above).
        """
        if threading.current_thread() is self._ioloop_thread:
            try:
                self._connection.close()
            except Exception:
                LOGGER.debug('connection.close() raised from IOLoop thread',
                             exc_info=True)
            return
        with self._channel_waiters_lock:
            if self._closed_reason is not None:
                return

        def _safe_close() -> None:
            try:
                self._connection.close()
            except Exception:
                # Already closing or closed — _on_connection_closed will wake waiters.
                LOGGER.debug(
                    'connection.close() raised (already closing or closed)',
                    exc_info=True)

        self._connection.ioloop.add_callback_threadsafe(_safe_close)

        if self._is_pool_worker_thread():
            # Called from within a delivery or connection-event callback running on one of this
            # connection's pool worker threads. The IOLoop thread's post-close cleanup tail
            # (_shutdown_all_consumer_pools / _shutdown_connection_pool) joins every pool worker,
            # including this one, once ioloop.start() returns - so joining self._ioloop_thread here
            # would deadlock: this thread waiting on the IOLoop thread, which is waiting on this
            # thread to return from the very callback that called close(). The close is already
            # scheduled above, so return without joining. This worker then returns from the
            # callback and, when the cleanup tail calls shutdown() on its pool and joins it, exits
            # via the drained-and-shutdown check at the top of _run_worker. (The shutdown() self-
            # join guard is not what saves us here: that fires only when a worker shuts down its own
            # pool, whereas here the IOLoop thread, not this worker, calls shutdown().)
            return

        self._ioloop_thread.join(timeout=timeout)
        if self._ioloop_thread.is_alive():
            self._connection.ioloop.add_callback_threadsafe(
                self._connection.ioloop.stop)
            # The IOLoop thread runs its own pool-drain tail after
            # ioloop.start() returns; a wedged consumer callback can keep it
            # from exiting.  Bound the join (and the pool shutdowns below) by a
            # fresh budget so a stuck worker cannot hang close() forever - the
            # worker is a daemon thread and is left to finish on its own.
            force_deadline = (None if timeout is None else time.monotonic() +
                              timeout)

            def _remaining() -> float | None:
                if force_deadline is None:
                    return None
                return max(0.0, force_deadline - time.monotonic())

            self._ioloop_thread.join(timeout=_remaining())
            if self._ioloop_thread.is_alive():
                LOGGER.warning(
                    'Connection %s IOLoop thread did not exit within the close '
                    'timeout; abandoning it', self._instance_id)
            # _on_connection_closed may not have fired if the IOLoop was
            # force-stopped before the broker sent Connection.CloseOk.
            # Wake any threads still blocked so they do not hang forever.
            forced = Exception(
                'connection force-closed: broker did not respond to close')
            with self._channel_waiters_lock:
                if self._closed_reason is None:
                    self._record_closed_reason(forced)
                    for evt, err in self._blocking_waiters:
                        if err[0] is None:
                            err[0] = self._closed_reason
                        evt.set()
                    self._blocking_waiters.clear()
            # Normally the IOLoop tail already claimed these pools, making the
            # calls no-ops; bound them in case this thread reaches them first.
            self._shutdown_all_consumer_pools(timeout=_remaining())
            self._shutdown_connection_pool(timeout=_remaining())

    def abort(self, timeout: float | None = 10) -> None:
        """
        Close the connection, swallowing any errors.

        Equivalent to :meth:`close` but never raises.  Currently :meth:`close` itself does not raise
        (it logs a warning on timeout), but ``abort`` is provided for API symmetry with other AMQP
        0-9-1 clients and to make error-recovery intent explicit at the call site.

        Safe to call from any thread.

        :param timeout: Seconds to wait for a clean close before force-stopping the IOLoop. Defaults
            to 10 seconds. Pass ``None`` to wait indefinitely.
        """
        try:
            self.close(timeout=timeout)
        except Exception:
            LOGGER.debug('connection abort() suppressed error', exc_info=True)

    def __enter__(self) -> Self:
        return self

    def __exit__(self, exc_type, exc_val, exc_tb) -> Literal[False]:
        self.close()
        return False

    def add_callback_threadsafe(self, callback) -> None:
        """
        Schedule *callback* to run in the IOLoop thread.

        Safe to call from any thread.  Exposed so callers can schedule arbitrary work without going
        through the wrapper methods.

        Raises rather than accepting work that can never run: once the connection is closed the
        IOLoop thread has exited, so a queued callback would be silently dropped.  The caller gets
        the recorded close reason, the same exception every other method on the connection raises
        once closed, or a :class:`~pika.exceptions.ConnectionWrongStateError` when no reason was
        recorded yet.

        Calls from the IOLoop thread itself are exempt and never raise, mirroring :meth:`close`.
        A close callback runs on that thread while the connection is already closed, so raising
        would abort the rest of teardown rather than reach the application.  The callback is
        accepted and logged at debug level, but a stopping IOLoop will most likely never run it.

        :param callback: Zero-argument callable.
        :raises BaseException: the recorded close reason, if the connection is closed and this is
            not the IOLoop thread.
        :raises pika.exceptions.ConnectionWrongStateError: if the connection is closed but no reason
            was recorded yet, and this is not the IOLoop thread.
        """
        try:
            self._check_not_closed()
        except Exception:
            # Teardown runs on the IOLoop thread: _on_connection_closed records
            # the close reason and then invokes the user's on_close_callback,
            # whose natural idiom is to schedule follow-up work.  Raising there
            # would propagate out through CallbackManager.process and abort the
            # rest of teardown, so accept the callback from that thread instead.
            # The IOLoop is already stopping, so the callback most likely will
            # not run; that beats leaving the connection half torn down.
            if threading.current_thread() is not self._ioloop_thread:
                raise
            LOGGER.debug(
                'add_callback_threadsafe() called from the IOLoop thread on a '
                'closed connection; the callback may not run',
                exc_info=True)
        self._schedule_unchecked(callback)

    def _schedule_unchecked(self, callback) -> None:
        """
        Hand *callback* to the IOLoop without the closed-connection check.

        For internal callers that must not be rejected mid-teardown.  Several channel methods
        register a waiter before scheduling and unregister it in a ``finally``; raising here would
        skip that cleanup and leak the waiter, so they check whatever state they care about
        themselves and then schedule through this method.

        :param callback: Zero-argument callable.
        """
        self._connection.ioloop.add_callback_threadsafe(callback)

    def _record_closed_reason(self, reason: BaseException) -> None:
        """
        Record the connection's close reason and the traceback it arrived with.

        Callers must already hold ``_channel_waiters_lock``: the reason and its traceback are read
        together and must not be seen half-updated.

        :param reason: The exception every blocked caller will be given. A value that is not an
            exception carries no traceback to snapshot; pika always passes an exception, but
            ``_on_connection_closed`` is driven by the underlying connection and does not enforce
            that.
        """
        self._closed_reason = reason
        self._closed_reason_tb = getattr(reason, '__traceback__', None)

    def _check_not_closed(self) -> None:
        """
        Raise if the connection is known to be closed.

        Best-effort: the connection may still close between this check and the caller's use of the
        IOLoop, which no amount of locking can prevent.  It catches the case that matters in
        practice - work submitted to a connection the caller already closed.

        Raises the recorded ``_closed_reason`` itself, the same exception every other method on the
        class raises once closed, with its close-time traceback restored so re-raising the shared
        instance does not grow it (see :func:`_with_close_traceback`).  Only when the connection is
        closed but no reason was recorded yet does it raise a fresh
        :class:`~pika.exceptions.ConnectionWrongStateError`.

        :raises BaseException: ``_closed_reason`` if one was recorded.
        :raises pika.exceptions.ConnectionWrongStateError: if the underlying connection is closed
            but no reason was recorded yet, which is the brief window between the connection
            reaching the closed state and ``_on_connection_closed`` running.
        """
        with self._channel_waiters_lock:
            reason = self._closed_reason
            if reason is not None:
                raise _with_close_traceback(reason, self._closed_reason_tb)
        if self._connection.is_closed:
            raise ConnectionWrongStateError(
                'Connection.add_callback_threadsafe() called on '
                'closed connection.')

    def add_on_connection_blocked_callback(self, callback) -> None:
        """
        Register a callback for ``Connection.Blocked`` notifications.

        RabbitMQ sends ``Connection.Blocked`` when the broker is running
        low on memory or disk.  In this state RabbitMQ stops processing
        incoming data, so a publisher receiving this notification should
        suspend publishing until the connection is unblocked.

        Dispatched on a connection-level worker thread (one per
        connection), so the callback may safely call any
        :class:`Channel` method without stalling heartbeats.
        The RabbitMQ Java and .NET clients run this listener inline on
        the I/O thread; pika's wrapper deliberately diverges so a slow
        listener cannot stall heartbeats.

        Safe to call from any thread.

        :param callback:
            ``callback(connection, method_frame)`` where *connection* is
            this :class:`Connection` and *method_frame* contains
            a :class:`pika.spec.Connection.Blocked`.
        :raises Exception: if the connection is already closed.
        """
        self._register_connection_event_callback(
            'add_on_connection_blocked_callback', callback)

    def add_on_connection_unblocked_callback(self, callback) -> None:
        """
        Register a callback for ``Connection.Unblocked`` notifications.

        Sent by RabbitMQ once a previously blocked connection is no
        longer resource-constrained, letting publishers resume.

        Dispatched on the connection-level worker thread (same as
        ``add_on_connection_blocked_callback``).  The RabbitMQ Java and
        .NET clients run this listener inline on the I/O thread; pika's
        wrapper deliberately diverges so a slow listener cannot stall
        heartbeats.

        Safe to call from any thread.

        :param callback:
            ``callback(connection, method_frame)`` where *connection* is
            this :class:`Connection` and *method_frame* contains
            a :class:`pika.spec.Connection.Unblocked`.
        :raises Exception: if the connection is already closed.
        """
        self._register_connection_event_callback(
            'add_on_connection_unblocked_callback', callback)

    def _register_connection_event_callback(self, raw_method_name: str,
                                            callback) -> None:
        """
        Register a connection-level event callback via the IOLoop.

        Wraps the user callback so it dispatches on the connection work pool, then schedules
        registration with the underlying :class:`pika.connection.Connection` on the IOLoop thread.
        :param raw_method_name: Unqualified AMQP method name (e.g. ``"Basic.Publish"``)
        :param callback: User callback to dispatch on the connection work pool."
        """
        with self._channel_waiters_lock:
            reason = self._closed_reason
            if reason is not None:
                raise _with_close_traceback(reason, self._closed_reason_tb)

        def _wrapped(_raw_conn, method_frame) -> None:
            _submit_or_terminate(
                self._connection_work_pool, self._connection,
                'Connection event dropped: work pool shut down',
                Channel._safe_dispatch, raw_method_name, callback, self,
                method_frame)

        def _register() -> None:
            try:
                getattr(self._connection, raw_method_name)(_wrapped)
            except Exception:
                LOGGER.warning('%s failed', raw_method_name, exc_info=True)

        self._connection.ioloop.add_callback_threadsafe(_register)

    @property
    def is_open(self) -> bool:
        return self._connection.is_open

    @property
    def is_closed(self) -> bool:
        return self._connection.is_closed

_blocking_waiters instance-attribute

_blocking_waiters: list[
    tuple[Event, list[BaseException | None]]
] = []

_channel_waiters_lock instance-attribute

_channel_waiters_lock = threading.Lock()

_channels instance-attribute

_channels: list[Channel] = []

_closed_reason instance-attribute

_closed_reason: BaseException | None = None

_closed_reason_tb instance-attribute

_closed_reason_tb: TracebackType | None = None

_connect_error instance-attribute

_connect_error = None

_connected_event instance-attribute

_connected_event = threading.Event()

_connection instance-attribute

_connection = SelectConnection(
    parameters=parameters,
    on_open_callback=self._on_connection_open,
    on_open_error_callback=self._on_connection_open_error,
    on_close_callback=self._on_connection_closed,
)

_connection_pool_shutdown instance-attribute

_connection_pool_shutdown = False

_connection_work_pool instance-attribute

_connection_work_pool = _BoundedWorkPool(
    maxsize=work_queue_maxsize,
    put_timeout=work_queue_put_timeout,
    thread_name=f'pika-conn-{self._instance_id}',
)

_instance_counter class-attribute instance-attribute

_instance_counter = itertools.count(1)

_instance_id instance-attribute

_instance_id = next(self._instance_counter)

_ioloop_thread instance-attribute

_ioloop_thread = threading.Thread(
    target=_run_ioloop,
    name=f'pika-ioloop-{self._instance_id}',
    daemon=True,
)

_user_on_close_callback instance-attribute

_user_on_close_callback = on_close_callback

_user_on_open_error_callback instance-attribute

_user_on_open_error_callback = on_open_error_callback

_work_queue_maxsize instance-attribute

_work_queue_maxsize = work_queue_maxsize

_work_queue_put_timeout instance-attribute

_work_queue_put_timeout = _validate_put_timeout(
    work_queue_put_timeout
)

is_closed property

is_closed: bool

is_open property

is_open: bool

__enter__

__enter__() -> Self
Source code in pika/adapters/thread_safe_connection.py
def __enter__(self) -> Self:
    return self

__exit__

__exit__(exc_type, exc_val, exc_tb) -> Literal[False]
Source code in pika/adapters/thread_safe_connection.py
def __exit__(self, exc_type, exc_val, exc_tb) -> Literal[False]:
    self.close()
    return False

__init__

__init__(
    parameters,
    on_open_error_callback=None,
    on_close_callback=None,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
    work_queue_maxsize: int = DEFAULT_WORK_QUEUE_MAXSIZE,
    work_queue_put_timeout: float = DEFAULT_WORK_QUEUE_PUT_TIMEOUT,
) -> None
Source code in pika/adapters/thread_safe_connection.py
def __init__(
        self,
        parameters,
        on_open_error_callback=None,
        on_close_callback=None,
        timeout: float | None = DEFAULT_RPC_TIMEOUT,
        work_queue_maxsize: int = DEFAULT_WORK_QUEUE_MAXSIZE,
        work_queue_put_timeout: float = DEFAULT_WORK_QUEUE_PUT_TIMEOUT
) -> None:
    self._user_on_open_error_callback = on_open_error_callback
    self._user_on_close_callback = on_close_callback
    self._work_queue_maxsize = work_queue_maxsize
    self._work_queue_put_timeout = _validate_put_timeout(
        work_queue_put_timeout)

    self._connect_error = None
    self._connected_event = threading.Event()

    self._instance_id = next(self._instance_counter)

    self._channel_waiters_lock = threading.Lock()
    self._closed_reason: BaseException | None = None
    # The traceback _closed_reason carried when it was recorded.  Raising
    # the shared instance appends frames to it permanently, so every raise
    # site restores this snapshot; see _with_close_traceback.
    self._closed_reason_tb: TracebackType | None = None
    self._blocking_waiters: list[tuple[threading.Event,
                                       list[BaseException | None]]] = []
    self._channels: list[Channel] = []

    # Single-worker pool for connection-level event callbacks
    # (Connection.Blocked / Unblocked).  Keeps user code off the
    # IOLoop thread so a slow listener cannot stall heartbeats.
    self._connection_work_pool = _BoundedWorkPool(
        maxsize=work_queue_maxsize,
        put_timeout=work_queue_put_timeout,
        thread_name=f'pika-conn-{self._instance_id}',
    )
    self._connection_pool_shutdown = False

    try:
        self._connection = SelectConnection(
            parameters=parameters,
            on_open_callback=self._on_connection_open,
            on_open_error_callback=self._on_connection_open_error,
            on_close_callback=self._on_connection_closed,
        )
    except Exception:
        # SelectConnection construction failed before the IOLoop thread
        # started, so _run_ioloop's cleanup tail will never run.  Shut
        # the connection-event pool down here to avoid leaking its
        # worker thread.
        self._shutdown_connection_pool()
        raise

    def _run_ioloop() -> None:
        try:
            self._connection.ioloop.start()
        except Exception as exc:
            # An unhandled exception in a callback killed the IOLoop.
            # Wake every blocked caller so they do not hang forever.
            LOGGER.exception('IOLoop thread crashed')
            with self._channel_waiters_lock:
                if self._closed_reason is None:
                    self._record_closed_reason(exc)
                for evt, err in self._blocking_waiters:
                    if err[0] is None:
                        err[0] = self._closed_reason
                    evt.set()
                self._blocking_waiters.clear()
            # Wake __init__ if it crashed before the connection opened.
            if not self._connected_event.is_set():
                self._connect_error = exc
                self._connected_event.set()
        # IOLoop has exited - safe to block waiting for pool workers.
        self._shutdown_all_consumer_pools()
        self._shutdown_connection_pool()

    self._ioloop_thread = threading.Thread(
        target=_run_ioloop,
        name=f'pika-ioloop-{self._instance_id}',
        daemon=True,
    )
    self._ioloop_thread.start()

    # Block the calling thread until the connection is open or fails.
    if not self._connected_event.wait(timeout=timeout):
        self._connect_error = TimeoutError(
            f'connection attempt timed out after {timeout} seconds')
        self._connection.ioloop.add_callback_threadsafe(
            self._connection.ioloop.stop)
        # Wait for the IOLoop thread to exit so it can shut down the
        # consumer and connection pools before we raise.  Otherwise
        # those pools' worker threads outlive __init__.
        self._ioloop_thread.join(timeout=timeout)
        raise self._connect_error
    if self._connect_error is not None:
        # IOLoop thread already exited (open-error path stops the
        # IOLoop) but its cleanup tail may still be running; wait
        # briefly so pool shutdown completes before we raise.
        self._ioloop_thread.join(timeout=timeout)
        raise self._connect_error

_check_not_closed

_check_not_closed() -> None

Raise if the connection is known to be closed.

Best-effort: the connection may still close between this check and the caller's use of the IOLoop, which no amount of locking can prevent. It catches the case that matters in practice - work submitted to a connection the caller already closed.

Raises the recorded _closed_reason itself, the same exception every other method on the class raises once closed, with its close-time traceback restored so re-raising the shared instance does not grow it (see :func:_with_close_traceback). Only when the connection is closed but no reason was recorded yet does it raise a fresh :class:~pika.exceptions.ConnectionWrongStateError.

RAISES DESCRIPTION
BaseException

_closed_reason if one was recorded.

pika.exceptions.ConnectionWrongStateError

if the underlying connection is closed but no reason was recorded yet, which is the brief window between the connection reaching the closed state and _on_connection_closed running.

Source code in pika/adapters/thread_safe_connection.py
def _check_not_closed(self) -> None:
    """
    Raise if the connection is known to be closed.

    Best-effort: the connection may still close between this check and the caller's use of the
    IOLoop, which no amount of locking can prevent.  It catches the case that matters in
    practice - work submitted to a connection the caller already closed.

    Raises the recorded ``_closed_reason`` itself, the same exception every other method on the
    class raises once closed, with its close-time traceback restored so re-raising the shared
    instance does not grow it (see :func:`_with_close_traceback`).  Only when the connection is
    closed but no reason was recorded yet does it raise a fresh
    :class:`~pika.exceptions.ConnectionWrongStateError`.

    :raises BaseException: ``_closed_reason`` if one was recorded.
    :raises pika.exceptions.ConnectionWrongStateError: if the underlying connection is closed
        but no reason was recorded yet, which is the brief window between the connection
        reaching the closed state and ``_on_connection_closed`` running.
    """
    with self._channel_waiters_lock:
        reason = self._closed_reason
        if reason is not None:
            raise _with_close_traceback(reason, self._closed_reason_tb)
    if self._connection.is_closed:
        raise ConnectionWrongStateError(
            'Connection.add_callback_threadsafe() called on '
            'closed connection.')

_is_pool_worker_thread

_is_pool_worker_thread() -> bool

True if the calling thread is one of this connection's pool worker threads.

Covers every channel's consumer worker plus the connection's own event-callback worker (:attr:_connection_work_pool). :meth:close uses this to detect that it is running from inside a delivery/event callback so it can avoid joining :attr:_ioloop_thread: that thread's post-ioloop.start() cleanup tail (:meth:_shutdown_all_consumer_pools, :meth:_shutdown_connection_pool) joins every pool worker in turn, including this one, so joining it here would deadlock each thread waiting on the other (issue #1686).

Reading each pool's _thread without its lock is safe here: a worker only executes this check from within its own run loop, by which point _thread was already assigned (under the pool's lock, before start()) by the thread that created it - that assignment happens-before anything the worker itself observes.

Source code in pika/adapters/thread_safe_connection.py
def _is_pool_worker_thread(self) -> bool:
    """
    True if the calling thread is one of this connection's pool worker threads.

    Covers every channel's consumer worker plus the connection's own event-callback worker
    (:attr:`_connection_work_pool`).  :meth:`close` uses this to detect that it is running from
    inside a delivery/event callback so it can avoid joining :attr:`_ioloop_thread`: that thread's
    post-``ioloop.start()`` cleanup tail (:meth:`_shutdown_all_consumer_pools`,
    :meth:`_shutdown_connection_pool`) joins every pool worker in turn, including this one, so
    joining it here would deadlock each thread waiting on the other (issue #1686).

    Reading each pool's ``_thread`` without its lock is safe here: a worker only executes this
    check from within its own run loop, by which point ``_thread`` was already assigned (under
    the pool's lock, before ``start()``) by the thread that created it - that assignment
    happens-before anything the worker itself observes.
    """
    current = threading.current_thread()
    if current is self._connection_work_pool._thread:
        return True
    with self._channel_waiters_lock:
        channels = list(self._channels)
    return any(current is ch._consumer_work_pool._thread for ch in channels)

_on_connection_closed

_on_connection_closed(_connection, reason) -> None
Source code in pika/adapters/thread_safe_connection.py
def _on_connection_closed(self, _connection, reason) -> None:
    # Connection is gone - stop the IOLoop so the thread exits.
    # Pool shutdown happens after ioloop.start() returns in _run_ioloop
    # (calling shutdown(wait=True) here would deadlock the IOLoop thread
    # if a pool worker is mid-callback).
    self._connection.ioloop.stop()
    with self._channel_waiters_lock:
        self._record_closed_reason(reason)
        for evt, err in self._blocking_waiters:
            err[0] = reason
            evt.set()
        self._blocking_waiters.clear()
    if self._user_on_close_callback:
        self._user_on_close_callback(_connection, reason)

_on_connection_open

_on_connection_open(_connection) -> None
Source code in pika/adapters/thread_safe_connection.py
def _on_connection_open(self, _connection) -> None:
    self._connected_event.set()

_on_connection_open_error

_on_connection_open_error(_connection, error) -> None
Source code in pika/adapters/thread_safe_connection.py
def _on_connection_open_error(self, _connection, error) -> None:
    self._connect_error = error
    # Stop the IOLoop so the background thread can exit cleanly.
    self._connection.ioloop.stop()
    self._connected_event.set()
    if self._user_on_open_error_callback:
        self._user_on_open_error_callback(_connection, error)

_record_closed_reason

_record_closed_reason(reason: BaseException) -> None

Record the connection's close reason and the traceback it arrived with.

Callers must already hold _channel_waiters_lock: the reason and its traceback are read together and must not be seen half-updated.

PARAMETER DESCRIPTION
reason

The exception every blocked caller will be given. A value that is not an exception carries no traceback to snapshot; pika always passes an exception, but _on_connection_closed is driven by the underlying connection and does not enforce that.

TYPE: BaseException

Source code in pika/adapters/thread_safe_connection.py
def _record_closed_reason(self, reason: BaseException) -> None:
    """
    Record the connection's close reason and the traceback it arrived with.

    Callers must already hold ``_channel_waiters_lock``: the reason and its traceback are read
    together and must not be seen half-updated.

    :param reason: The exception every blocked caller will be given. A value that is not an
        exception carries no traceback to snapshot; pika always passes an exception, but
        ``_on_connection_closed`` is driven by the underlying connection and does not enforce
        that.
    """
    self._closed_reason = reason
    self._closed_reason_tb = getattr(reason, '__traceback__', None)

_register_connection_event_callback

_register_connection_event_callback(
    raw_method_name: str, callback
) -> None

Register a connection-level event callback via the IOLoop.

Wraps the user callback so it dispatches on the connection work pool, then schedules registration with the underlying :class:pika.connection.Connection on the IOLoop thread.

PARAMETER DESCRIPTION
raw_method_name

Unqualified AMQP method name (e.g. "Basic.Publish")

TYPE: str

callback

User callback to dispatch on the connection work pool."

Source code in pika/adapters/thread_safe_connection.py
def _register_connection_event_callback(self, raw_method_name: str,
                                        callback) -> None:
    """
    Register a connection-level event callback via the IOLoop.

    Wraps the user callback so it dispatches on the connection work pool, then schedules
    registration with the underlying :class:`pika.connection.Connection` on the IOLoop thread.
    :param raw_method_name: Unqualified AMQP method name (e.g. ``"Basic.Publish"``)
    :param callback: User callback to dispatch on the connection work pool."
    """
    with self._channel_waiters_lock:
        reason = self._closed_reason
        if reason is not None:
            raise _with_close_traceback(reason, self._closed_reason_tb)

    def _wrapped(_raw_conn, method_frame) -> None:
        _submit_or_terminate(
            self._connection_work_pool, self._connection,
            'Connection event dropped: work pool shut down',
            Channel._safe_dispatch, raw_method_name, callback, self,
            method_frame)

    def _register() -> None:
        try:
            getattr(self._connection, raw_method_name)(_wrapped)
        except Exception:
            LOGGER.warning('%s failed', raw_method_name, exc_info=True)

    self._connection.ioloop.add_callback_threadsafe(_register)

_schedule_unchecked

_schedule_unchecked(callback) -> None

Hand callback to the IOLoop without the closed-connection check.

For internal callers that must not be rejected mid-teardown. Several channel methods register a waiter before scheduling and unregister it in a finally; raising here would skip that cleanup and leak the waiter, so they check whatever state they care about themselves and then schedule through this method.

PARAMETER DESCRIPTION
callback

Zero-argument callable.

Source code in pika/adapters/thread_safe_connection.py
def _schedule_unchecked(self, callback) -> None:
    """
    Hand *callback* to the IOLoop without the closed-connection check.

    For internal callers that must not be rejected mid-teardown.  Several channel methods
    register a waiter before scheduling and unregister it in a ``finally``; raising here would
    skip that cleanup and leak the waiter, so they check whatever state they care about
    themselves and then schedule through this method.

    :param callback: Zero-argument callable.
    """
    self._connection.ioloop.add_callback_threadsafe(callback)

_shutdown_all_consumer_pools

_shutdown_all_consumer_pools(
    timeout: float | None = None,
) -> None

Shut down all tracked channel consumer pools.

PARAMETER DESCRIPTION
timeout

Total seconds to wait across all channel pools. None waits indefinitely (the clean-shutdown default). A finite budget is split across channels so the whole sweep stays bounded.

TYPE: float | None DEFAULT: None

Source code in pika/adapters/thread_safe_connection.py
def _shutdown_all_consumer_pools(self,
                                 timeout: float | None = None) -> None:
    """
    Shut down all tracked channel consumer pools.

    :param timeout: Total seconds to wait across all channel pools. ``None`` waits indefinitely
        (the clean-shutdown default). A finite budget is split across channels so the whole
        sweep stays bounded.
    """
    with self._channel_waiters_lock:
        channels = list(self._channels)
    if not channels:
        return
    if timeout is None:
        for ch in channels:
            ch._shutdown_pool()
        return
    # Signal every pool to start draining before joining any of them, so
    # all workers drain concurrently.  Joining one at a time without this
    # would let an early wedged channel burn the whole budget while later
    # channels' workers had not even been told to stop, so those healthy
    # workers would be abandoned at timeout=0 with a spurious warning.
    for ch in channels:
        ch._signal_pool_shutdown()
    # Share the remaining budget across channels using a common deadline so
    # a slow early channel cannot consume the entire allowance and starve
    # later ones of their chance to drain.
    deadline = time.monotonic() + timeout
    for ch in channels:
        ch._shutdown_pool(timeout=max(0.0, deadline - time.monotonic()))

_shutdown_connection_pool

_shutdown_connection_pool(
    timeout: float | None = None,
) -> None

Shut down the connection-level event-callback pool.

PARAMETER DESCRIPTION
timeout

Seconds to wait for the worker to drain and exit. None waits indefinitely (the clean-shutdown default). A wedged worker outliving a finite timeout is left running (daemon thread) rather than stalling the caller; a warning is logged.

TYPE: float | None DEFAULT: None

Source code in pika/adapters/thread_safe_connection.py
def _shutdown_connection_pool(self, timeout: float | None = None) -> None:
    """
    Shut down the connection-level event-callback pool.

    :param timeout: Seconds to wait for the worker to drain and exit. ``None`` waits
        indefinitely (the clean-shutdown default). A wedged worker outliving a finite *timeout*
        is left running (daemon thread) rather than stalling the caller; a warning is logged.
    """
    # Atomically claim the shutdown so the IOLoop tail and a concurrent
    # close() do not both run the join.  The join stays outside the lock so
    # it cannot stall other channel operations for the whole timeout.
    with self._channel_waiters_lock:
        if self._connection_pool_shutdown:
            return
        self._connection_pool_shutdown = True
    if not self._connection_work_pool.shutdown(wait=True, timeout=timeout):
        LOGGER.warning(
            'Connection %s event work pool did not drain within '
            '%s seconds; abandoning its worker thread', self._instance_id,
            timeout)

abort

abort(timeout: float | None = 10) -> None

Close the connection, swallowing any errors.

Equivalent to :meth:close but never raises. Currently :meth:close itself does not raise (it logs a warning on timeout), but abort is provided for API symmetry with other AMQP 0-9-1 clients and to make error-recovery intent explicit at the call site.

Safe to call from any thread.

PARAMETER DESCRIPTION
timeout

Seconds to wait for a clean close before force-stopping the IOLoop. Defaults to 10 seconds. Pass None to wait indefinitely.

TYPE: float | None DEFAULT: 10

Source code in pika/adapters/thread_safe_connection.py
def abort(self, timeout: float | None = 10) -> None:
    """
    Close the connection, swallowing any errors.

    Equivalent to :meth:`close` but never raises.  Currently :meth:`close` itself does not raise
    (it logs a warning on timeout), but ``abort`` is provided for API symmetry with other AMQP
    0-9-1 clients and to make error-recovery intent explicit at the call site.

    Safe to call from any thread.

    :param timeout: Seconds to wait for a clean close before force-stopping the IOLoop. Defaults
        to 10 seconds. Pass ``None`` to wait indefinitely.
    """
    try:
        self.close(timeout=timeout)
    except Exception:
        LOGGER.debug('connection abort() suppressed error', exc_info=True)

add_callback_threadsafe

add_callback_threadsafe(callback) -> None

Schedule callback to run in the IOLoop thread.

Safe to call from any thread. Exposed so callers can schedule arbitrary work without going through the wrapper methods.

Raises rather than accepting work that can never run: once the connection is closed the IOLoop thread has exited, so a queued callback would be silently dropped. The caller gets the recorded close reason, the same exception every other method on the connection raises once closed, or a :class:~pika.exceptions.ConnectionWrongStateError when no reason was recorded yet.

Calls from the IOLoop thread itself are exempt and never raise, mirroring :meth:close. A close callback runs on that thread while the connection is already closed, so raising would abort the rest of teardown rather than reach the application. The callback is accepted and logged at debug level, but a stopping IOLoop will most likely never run it.

PARAMETER DESCRIPTION
callback

Zero-argument callable.

RAISES DESCRIPTION
BaseException

the recorded close reason, if the connection is closed and this is not the IOLoop thread.

pika.exceptions.ConnectionWrongStateError

if the connection is closed but no reason was recorded yet, and this is not the IOLoop thread.

Source code in pika/adapters/thread_safe_connection.py
def add_callback_threadsafe(self, callback) -> None:
    """
    Schedule *callback* to run in the IOLoop thread.

    Safe to call from any thread.  Exposed so callers can schedule arbitrary work without going
    through the wrapper methods.

    Raises rather than accepting work that can never run: once the connection is closed the
    IOLoop thread has exited, so a queued callback would be silently dropped.  The caller gets
    the recorded close reason, the same exception every other method on the connection raises
    once closed, or a :class:`~pika.exceptions.ConnectionWrongStateError` when no reason was
    recorded yet.

    Calls from the IOLoop thread itself are exempt and never raise, mirroring :meth:`close`.
    A close callback runs on that thread while the connection is already closed, so raising
    would abort the rest of teardown rather than reach the application.  The callback is
    accepted and logged at debug level, but a stopping IOLoop will most likely never run it.

    :param callback: Zero-argument callable.
    :raises BaseException: the recorded close reason, if the connection is closed and this is
        not the IOLoop thread.
    :raises pika.exceptions.ConnectionWrongStateError: if the connection is closed but no reason
        was recorded yet, and this is not the IOLoop thread.
    """
    try:
        self._check_not_closed()
    except Exception:
        # Teardown runs on the IOLoop thread: _on_connection_closed records
        # the close reason and then invokes the user's on_close_callback,
        # whose natural idiom is to schedule follow-up work.  Raising there
        # would propagate out through CallbackManager.process and abort the
        # rest of teardown, so accept the callback from that thread instead.
        # The IOLoop is already stopping, so the callback most likely will
        # not run; that beats leaving the connection half torn down.
        if threading.current_thread() is not self._ioloop_thread:
            raise
        LOGGER.debug(
            'add_callback_threadsafe() called from the IOLoop thread on a '
            'closed connection; the callback may not run',
            exc_info=True)
    self._schedule_unchecked(callback)

add_on_connection_blocked_callback

add_on_connection_blocked_callback(callback) -> None

Register a callback for Connection.Blocked notifications.

RabbitMQ sends Connection.Blocked when the broker is running low on memory or disk. In this state RabbitMQ stops processing incoming data, so a publisher receiving this notification should suspend publishing until the connection is unblocked.

Dispatched on a connection-level worker thread (one per connection), so the callback may safely call any :class:Channel method without stalling heartbeats. The RabbitMQ Java and .NET clients run this listener inline on the I/O thread; pika's wrapper deliberately diverges so a slow listener cannot stall heartbeats.

Safe to call from any thread.

PARAMETER DESCRIPTION
callback

callback(connection, method_frame) where connection is this :class:Connection and method_frame contains a :class:pika.spec.Connection.Blocked.

RAISES DESCRIPTION
Exception

if the connection is already closed.

Source code in pika/adapters/thread_safe_connection.py
def add_on_connection_blocked_callback(self, callback) -> None:
    """
    Register a callback for ``Connection.Blocked`` notifications.

    RabbitMQ sends ``Connection.Blocked`` when the broker is running
    low on memory or disk.  In this state RabbitMQ stops processing
    incoming data, so a publisher receiving this notification should
    suspend publishing until the connection is unblocked.

    Dispatched on a connection-level worker thread (one per
    connection), so the callback may safely call any
    :class:`Channel` method without stalling heartbeats.
    The RabbitMQ Java and .NET clients run this listener inline on
    the I/O thread; pika's wrapper deliberately diverges so a slow
    listener cannot stall heartbeats.

    Safe to call from any thread.

    :param callback:
        ``callback(connection, method_frame)`` where *connection* is
        this :class:`Connection` and *method_frame* contains
        a :class:`pika.spec.Connection.Blocked`.
    :raises Exception: if the connection is already closed.
    """
    self._register_connection_event_callback(
        'add_on_connection_blocked_callback', callback)

add_on_connection_unblocked_callback

add_on_connection_unblocked_callback(callback) -> None

Register a callback for Connection.Unblocked notifications.

Sent by RabbitMQ once a previously blocked connection is no longer resource-constrained, letting publishers resume.

Dispatched on the connection-level worker thread (same as add_on_connection_blocked_callback). The RabbitMQ Java and .NET clients run this listener inline on the I/O thread; pika's wrapper deliberately diverges so a slow listener cannot stall heartbeats.

Safe to call from any thread.

PARAMETER DESCRIPTION
callback

callback(connection, method_frame) where connection is this :class:Connection and method_frame contains a :class:pika.spec.Connection.Unblocked.

RAISES DESCRIPTION
Exception

if the connection is already closed.

Source code in pika/adapters/thread_safe_connection.py
def add_on_connection_unblocked_callback(self, callback) -> None:
    """
    Register a callback for ``Connection.Unblocked`` notifications.

    Sent by RabbitMQ once a previously blocked connection is no
    longer resource-constrained, letting publishers resume.

    Dispatched on the connection-level worker thread (same as
    ``add_on_connection_blocked_callback``).  The RabbitMQ Java and
    .NET clients run this listener inline on the I/O thread; pika's
    wrapper deliberately diverges so a slow listener cannot stall
    heartbeats.

    Safe to call from any thread.

    :param callback:
        ``callback(connection, method_frame)`` where *connection* is
        this :class:`Connection` and *method_frame* contains
        a :class:`pika.spec.Connection.Unblocked`.
    :raises Exception: if the connection is already closed.
    """
    self._register_connection_event_callback(
        'add_on_connection_unblocked_callback', callback)

channel

channel(
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> Channel

Open a new channel and return a :class:Channel.

Blocks the calling thread until the channel is open. The returned channel's methods are safe to call from any thread.

PARAMETER DESCRIPTION
timeout

Seconds to wait for Channel.OpenOk. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RAISES DESCRIPTION
Exception

if the connection is closed before the channel opens.

TimeoutError

if timeout expires before the channel opens.

Source code in pika/adapters/thread_safe_connection.py
def channel(self, timeout: float | None = DEFAULT_RPC_TIMEOUT) -> Channel:
    """
    Open a new channel and return a :class:`Channel`.

    Blocks the calling thread until the channel is open.  The returned channel's methods are
    safe to call from any thread.

    :param timeout: Seconds to wait for Channel.OpenOk. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :raises Exception: if the connection is closed before the channel opens.
    :raises TimeoutError: if *timeout* expires before the channel opens.
    """
    ready = threading.Event()
    result = [None]
    error: list[BaseException | None] = [None]

    with self._channel_waiters_lock:
        reason = self._closed_reason
        if reason is not None:
            raise _with_close_traceback(reason, self._closed_reason_tb)
        self._blocking_waiters.append((ready, error))

    def _open() -> None:

        def _on_open(ch) -> None:
            result[0] = ch
            ready.set()

        try:
            self._connection.channel(on_open_callback=_on_open)
        except Exception as exc:
            error[0] = exc
            ready.set()

    self._connection.ioloop.add_callback_threadsafe(_open)

    try:
        if not ready.wait(timeout=timeout):
            raise TimeoutError(
                f'channel open timed out after {timeout} seconds')
    finally:
        with self._channel_waiters_lock:
            try:
                self._blocking_waiters.remove((ready, error))
            except ValueError:
                pass

    if error[0] is not None:
        raise _reraisable(error[0], self._closed_reason,
                          self._closed_reason_tb)

    # Race guard: the connection may have closed on the IOLoop thread
    # after _on_open fired but before we woke up.  Without this check,
    # _shutdown_all_consumer_pools() would have already run and the
    # newly-constructed channel's consumer pool would never be reached.
    with self._channel_waiters_lock:
        reason = self._closed_reason
        if reason is not None:
            raise _with_close_traceback(reason, self._closed_reason_tb)
        ch = Channel(result[0],
                     self,
                     work_queue_maxsize=self._work_queue_maxsize,
                     work_queue_put_timeout=self._work_queue_put_timeout)
        self._channels.append(ch)
    # Register the broker-close hook on the IOLoop thread, where the raw
    # channel's callback stack is safe to mutate.  A broker close arriving
    # in the brief gap before this runs simply falls back to the
    # connection-shutdown sweep, which still drops the channel from
    # tracking.
    self._schedule_unchecked(
        lambda: ch._channel.add_on_close_callback(ch._on_broker_close))
    return ch

close

close(timeout: float | None = 10) -> None

Close the connection and block until the IOLoop thread exits.

Schedules a clean Connection.Close handshake and waits up to timeout seconds for the IOLoop thread to finish. If the thread is still alive after the timeout (e.g. the broker never sends Connection.CloseOk), the IOLoop is force-stopped via add_callback_threadsafe and joined once more under a second timeout budget. A worker still draining when that budget expires is left to finish on its own (it is a daemon thread), so a wedged delivery callback cannot hang close indefinitely.

Safe to call from any thread including the IOLoop thread itself (e.g. from within a channel callback). When called from the IOLoop thread the close is initiated synchronously and the method returns immediately without joining.

Also safe to call from a delivery or connection-event callback running on one of this connection's pool worker threads (e.g. a consumer's on_message_callback). The IOLoop thread's own post-close cleanup joins every pool worker, including whichever one is running the callback that called close(), so joining :attr:_ioloop_thread from there would deadlock the two threads against each other. In that case the close is scheduled and this method returns immediately without joining; the worker exits normally once the callback returns. Because this path does not join the IOLoop thread, timeout is not applied on it: the close is best-effort and an unresponsive broker is not force-stopped from here, so a healthy handshake is assumed (see #1706).

Calling close() on an already-closed connection is a no-op.

PARAMETER DESCRIPTION
timeout

Seconds to wait for a clean close before force-stopping the IOLoop. Defaults to 10 seconds, which is sufficient for a healthy broker on any reasonable network. The pathological force path may wait up to twice this (once for the clean join, once for the forced join). Pass None to wait indefinitely. Ignored when close() is called from one of this connection's pool worker threads (see above).

TYPE: float | None DEFAULT: 10

Source code in pika/adapters/thread_safe_connection.py
def close(self, timeout: float | None = 10) -> None:
    """
    Close the connection and block until the IOLoop thread exits.

    Schedules a clean Connection.Close handshake and waits up to *timeout* seconds for the
    IOLoop thread to finish.  If the thread is still alive after the timeout (e.g. the broker
    never sends Connection.CloseOk), the IOLoop is force-stopped via ``add_callback_threadsafe``
    and joined once more under a second *timeout* budget.  A worker still draining when that
    budget expires is left to finish on its own (it is a daemon thread), so a wedged delivery
    callback cannot hang ``close`` indefinitely.

    Safe to call from any thread including the IOLoop thread itself (e.g. from within a channel
    callback).  When called from the IOLoop thread the close is initiated synchronously and the
    method returns immediately without joining.

    Also safe to call from a delivery or connection-event callback running on one of this
    connection's pool worker threads (e.g. a consumer's ``on_message_callback``).  The IOLoop
    thread's own post-close cleanup joins every pool worker, including whichever one is running
    the callback that called ``close()``, so joining :attr:`_ioloop_thread` from there would
    deadlock the two threads against each other. In that case the close is scheduled and this
    method returns immediately without joining; the worker exits normally once the callback
    returns.  Because this path does not join the IOLoop thread, *timeout* is not applied on it:
    the close is best-effort and an unresponsive broker is not force-stopped from here, so a
    healthy handshake is assumed (see #1706).

    Calling ``close()`` on an already-closed connection is a no-op.

    :param timeout: Seconds to wait for a clean close before force-stopping the IOLoop. Defaults
        to 10 seconds, which is sufficient for a healthy broker on any reasonable network. The
        pathological force path may wait up to twice this (once for the clean join, once for the
        forced join). Pass ``None`` to wait indefinitely. Ignored when ``close()`` is called
        from one of this connection's pool worker threads (see above).
    """
    if threading.current_thread() is self._ioloop_thread:
        try:
            self._connection.close()
        except Exception:
            LOGGER.debug('connection.close() raised from IOLoop thread',
                         exc_info=True)
        return
    with self._channel_waiters_lock:
        if self._closed_reason is not None:
            return

    def _safe_close() -> None:
        try:
            self._connection.close()
        except Exception:
            # Already closing or closed — _on_connection_closed will wake waiters.
            LOGGER.debug(
                'connection.close() raised (already closing or closed)',
                exc_info=True)

    self._connection.ioloop.add_callback_threadsafe(_safe_close)

    if self._is_pool_worker_thread():
        # Called from within a delivery or connection-event callback running on one of this
        # connection's pool worker threads. The IOLoop thread's post-close cleanup tail
        # (_shutdown_all_consumer_pools / _shutdown_connection_pool) joins every pool worker,
        # including this one, once ioloop.start() returns - so joining self._ioloop_thread here
        # would deadlock: this thread waiting on the IOLoop thread, which is waiting on this
        # thread to return from the very callback that called close(). The close is already
        # scheduled above, so return without joining. This worker then returns from the
        # callback and, when the cleanup tail calls shutdown() on its pool and joins it, exits
        # via the drained-and-shutdown check at the top of _run_worker. (The shutdown() self-
        # join guard is not what saves us here: that fires only when a worker shuts down its own
        # pool, whereas here the IOLoop thread, not this worker, calls shutdown().)
        return

    self._ioloop_thread.join(timeout=timeout)
    if self._ioloop_thread.is_alive():
        self._connection.ioloop.add_callback_threadsafe(
            self._connection.ioloop.stop)
        # The IOLoop thread runs its own pool-drain tail after
        # ioloop.start() returns; a wedged consumer callback can keep it
        # from exiting.  Bound the join (and the pool shutdowns below) by a
        # fresh budget so a stuck worker cannot hang close() forever - the
        # worker is a daemon thread and is left to finish on its own.
        force_deadline = (None if timeout is None else time.monotonic() +
                          timeout)

        def _remaining() -> float | None:
            if force_deadline is None:
                return None
            return max(0.0, force_deadline - time.monotonic())

        self._ioloop_thread.join(timeout=_remaining())
        if self._ioloop_thread.is_alive():
            LOGGER.warning(
                'Connection %s IOLoop thread did not exit within the close '
                'timeout; abandoning it', self._instance_id)
        # _on_connection_closed may not have fired if the IOLoop was
        # force-stopped before the broker sent Connection.CloseOk.
        # Wake any threads still blocked so they do not hang forever.
        forced = Exception(
            'connection force-closed: broker did not respond to close')
        with self._channel_waiters_lock:
            if self._closed_reason is None:
                self._record_closed_reason(forced)
                for evt, err in self._blocking_waiters:
                    if err[0] is None:
                        err[0] = self._closed_reason
                    evt.set()
                self._blocking_waiters.clear()
        # Normally the IOLoop tail already claimed these pools, making the
        # calls no-ops; bound them in case this thread reaches them first.
        self._shutdown_all_consumer_pools(timeout=_remaining())
        self._shutdown_connection_pool(timeout=_remaining())

Channel

Thread-safe wrapper around :class:pika.channel.Channel.

Every write operation is routed through the parent connection's add_callback_threadsafe so that _tx_buffers is only ever touched from the IOLoop thread.

.. rubric:: Consumer callback threading model

The on_message_callback registered with :meth:basic_consume is dispatched on a per-channel worker thread, not the IOLoop thread. This means:

  • Blocking operations (database writes, HTTP calls, even :meth:queue_declare / :meth:basic_qos / :meth:basic_cancel) are safe inside delivery callbacks.
  • As long as the callback keeps up with deliveries, the IOLoop thread is not starved by consumer processing, so heartbeats are sent on time.
  • Messages are delivered to the callback in order (a single worker thread per channel).
  • All :class:Channel methods (:meth:basic_ack, :meth:basic_nack, :meth:basic_reject, :meth:basic_publish, :meth:queue_declare, etc.) are safe to call from within the callback.

.. rubric:: Sustained-backlog caveat

"Safe" assumes the callback keeps pace with incoming deliveries. The worker consumes from a bounded queue (see work_queue_maxsize). If a slow callback lets that queue fill, the IOLoop thread blocks while enqueuing the next delivery - sending no heartbeats during the stall - and, if the queue stays full for work_queue_put_timeout seconds, raises :class:pika.exceptions.WorkQueueFullError, which tears the connection down rather than silently dropping the event. A persistently slow consumer therefore does not stall heartbeats forever; it fails fast with an explicit error. Keep callbacks faster than the sustained delivery rate, or raise work_queue_maxsize to absorb bursts.

.. rubric:: IOLoop-thread callbacks

Callables passed directly to :meth:~Connection.add_callback_threadsafe still run on the IOLoop thread. These must return quickly and must not call blocking channel methods (which would deadlock).

Source code in pika/adapters/thread_safe_connection.py
class Channel:
    """
    Thread-safe wrapper around :class:`pika.channel.Channel`.

    Every write operation is routed through the parent connection's
    ``add_callback_threadsafe`` so that ``_tx_buffers`` is only ever
    touched from the IOLoop thread.

    .. rubric:: Consumer callback threading model

    The *on_message_callback* registered with :meth:`basic_consume` is
    dispatched on a per-channel worker thread, **not** the IOLoop
    thread.  This means:

    - Blocking operations (database writes, HTTP calls, even
      :meth:`queue_declare` / :meth:`basic_qos` / :meth:`basic_cancel`)
      are safe inside delivery callbacks.
    - As long as the callback keeps up with deliveries, the IOLoop
      thread is not starved by consumer processing, so heartbeats are
      sent on time.
    - Messages are delivered to the callback **in order** (a single
      worker thread per channel).
    - All :class:`Channel` methods (:meth:`basic_ack`,
      :meth:`basic_nack`, :meth:`basic_reject`, :meth:`basic_publish`,
      :meth:`queue_declare`, etc.) are safe to call from within the
      callback.

    .. rubric:: Sustained-backlog caveat

    "Safe" assumes the callback keeps pace with incoming deliveries.
    The worker consumes from a bounded queue (see *work_queue_maxsize*).
    If a slow callback lets that queue fill, the IOLoop thread blocks
    while enqueuing the next delivery - sending no heartbeats during the
    stall - and, if the queue stays full for *work_queue_put_timeout*
    seconds, raises :class:`pika.exceptions.WorkQueueFullError`, which
    tears the connection down rather than silently dropping the event.
    A persistently slow consumer therefore does not stall heartbeats
    forever; it fails fast with an explicit error.  Keep callbacks
    faster than the sustained delivery rate, or raise *work_queue_maxsize*
    to absorb bursts.

    .. rubric:: IOLoop-thread callbacks

    Callables passed directly to
    :meth:`~Connection.add_callback_threadsafe` still run
    on the IOLoop thread.  These must return quickly and must not call
    blocking channel methods (which would deadlock).
    """

    def __init__(
            self,
            channel,
            wrapper,
            work_queue_maxsize: int = DEFAULT_WORK_QUEUE_MAXSIZE,
            work_queue_put_timeout: float = DEFAULT_WORK_QUEUE_PUT_TIMEOUT
    ) -> None:
        work_queue_put_timeout = _validate_put_timeout(work_queue_put_timeout)
        self._channel = channel
        self._wrapper = wrapper
        self._consumer_work_pool = _BoundedWorkPool(
            maxsize=work_queue_maxsize,
            put_timeout=work_queue_put_timeout,
            thread_name='pika-consumer',
        )
        self._pool_shutdown = False
        self._next_publish_seq_no: int | None = None
        self._confirm_select_ok = None
        # Serializes confirm_delivery so two threads cannot both pass the
        # _confirm_select_ok guard while the first is still waiting for
        # Confirm.SelectOk.  Per-channel rather than the connection-wide
        # _channel_waiters_lock: the call blocks for a broker round trip, which
        # must not stall operations on other channels.
        self._confirm_lock = threading.Lock()

    def _check_not_closed(self) -> None:
        """
        Raise if the connection is known to be closed.

        Called from fire-and-forget methods to prevent silently dropping work when the connection is
        already gone.
        """
        with self._wrapper._channel_waiters_lock:
            reason = self._wrapper._closed_reason
            if reason is not None:
                raise _with_close_traceback(reason,
                                            self._wrapper._closed_reason_tb)

    @staticmethod
    def _safe_dispatch(label, callback, *args) -> None:
        """
        Run *callback*, logging any exception instead of letting it propagate.

        Wraps a user callback so one failing dispatch neither escapes into the caller nor prevents
        later dispatches from running.  Most callers submit the wrapped callback to a
        :class:`_BoundedWorkPool`, so it runs on a pool worker; :meth:`basic_publish` calls it
        inline on the IOLoop thread. Either way the exception is contained here.

        :param label: Human-readable name of the callback for log lines.
        :param callback: The user callback.
        :param args: Positional arguments forwarded to *callback*.
        """
        try:
            callback(*args)
        except Exception:
            LOGGER.exception('Unhandled exception in %s', label)

    def _signal_pool_shutdown(self) -> None:
        """
        Tell the consumer pool to begin draining without joining its worker.

        Sets the pool's shutdown flag and wakes an idle worker, but returns immediately.  Used to
        signal every channel's pool up front so their workers drain concurrently before
        :meth:`_shutdown_pool` joins them one at a time under a shared budget.
        """
        self._consumer_work_pool.shutdown(wait=False)

    def _claim_pool_shutdown(self) -> bool:
        """
        Claim this channel's pool shutdown and drop it from connection tracking.

        Returns ``True`` if this call is the first to claim the shutdown - so the caller owns any
        follow-up work, such as joining the worker - and ``False`` if a prior close already claimed
        it.  Being idempotent lets a user :meth:`close`, the broker-close hook
        (:meth:`_on_broker_close`) and the connection-shutdown sweep race while each channel is
        dropped from ``_channels`` exactly once.

        The tracking list exists only so :meth:`~Connection._shutdown_all_consumer_pools` can reach
        every live pool, and a channel whose pool is shutting down has nothing left to drain; left
        in place the entries accumulate for the life of the connection, retaining each closed
        channel's pool, raw channel and per-RPC callbacks.

        :returns: ``True`` if this call claimed the shutdown, ``False`` if a prior close already did.
        """
        # Safe to mutate the list while a sweep is in flight:
        # _shutdown_all_consumer_pools snapshots it under this same lock before
        # iterating.
        with self._wrapper._channel_waiters_lock:
            if self._pool_shutdown:
                return False
            self._pool_shutdown = True
            try:
                self._wrapper._channels.remove(self)
            except ValueError:
                pass
        return True

    def _shutdown_pool(self, timeout: float | None = None) -> None:
        """
        Shut down the consumer work pool, joining its worker.

        Claims the shutdown and drops the channel from tracking via :meth:`_claim_pool_shutdown`,
        then joins the worker.  Used by the paths that can afford to wait off the IOLoop thread: a
        user :meth:`close` (on the caller's thread) and the connection-shutdown sweep (after the
        IOLoop has stopped).

        :param timeout: Seconds to wait for the worker to drain and exit. ``None`` waits
            indefinitely. A wedged worker that outlives a finite *timeout* is left running (it is a
            daemon thread) rather than stalling the caller; a warning is logged.
        """
        # The join stays outside the lock: it can block for the whole timeout,
        # and holding the connection-wide lock that long would stall every other
        # channel operation.
        if not self._claim_pool_shutdown():
            return
        if not self._consumer_work_pool.shutdown(wait=True, timeout=timeout):
            LOGGER.warning(
                'Channel %s consumer work pool did not drain within '
                '%s seconds; abandoning its worker thread',
                self._channel.channel_number, timeout)

    def _on_broker_close(self, _channel, reason) -> None:
        """
        Drop this channel from tracking when the broker closes it on its own.

        Registered once per channel, on the IOLoop thread, right after the channel opens.  A client
        close - a user :meth:`close`, or the per-channel closes a graceful :meth:`Connection.close`
        issues - reports :class:`~pika.exceptions.ChannelClosedByClient` and is left to
        :meth:`_shutdown_pool`, so this handles only broker- or error-initiated closes.  Without it
        those channels stay in ``_channels`` (with their pool, worker thread, raw channel and
        per-RPC callbacks) until the connection itself is torn down - the leak #1688 describes,
        reached through a close path :meth:`_shutdown_pool` never sees.

        Runs on the IOLoop thread, so it must not join the worker: joining a worker that is mid
        delivery-callback and waiting on this same IOLoop would deadlock, which is why pool joins
        are otherwise deferred off this thread (see :meth:`~Connection._on_connection_closed`).  It
        only signals the pool to drain; the daemon worker exits on its own once the channel is
        closed and no further deliveries arrive.  The connection-shutdown sweep, running after the
        IOLoop stops, joins any channel this has not already dropped.

        :param _channel: The raw channel reporting the close (unused).
        :param reason: The exception describing why the channel closed.
        """
        from pika.exceptions import ChannelClosedByClient
        if isinstance(reason, ChannelClosedByClient):
            return
        if self._claim_pool_shutdown():
            self._consumer_work_pool.shutdown(wait=False)

    def _register_waiter(self) -> tuple[Event, list[BaseException | None]]:
        """
        Create and register a blocking waiter.

        :returns: (ready, error) tuple for use with ``ready.wait()``
        :raises Exception: if the connection is already closed.
        """
        ready = threading.Event()
        error: list[BaseException | None] = [None]
        with self._wrapper._channel_waiters_lock:
            reason = self._wrapper._closed_reason
            if reason is not None:
                raise _with_close_traceback(reason,
                                            self._wrapper._closed_reason_tb)
            self._wrapper._blocking_waiters.append((ready, error))
        return ready, error

    def _unregister_waiter(self, ready: Event,
                           error: list[BaseException | None]) -> None:
        """
        Remove a waiter from the blocking list.

        :param ready: Threading event signalling that the RPC response has
            arrived
        :param error: list containing an exception if the RPC response was an
            error, or None if it was successful
        """
        with self._wrapper._channel_waiters_lock:
            try:
                self._wrapper._blocking_waiters.remove((ready, error))
            except ValueError:
                pass

    def _release_close_callback(self, on_chan_close, *extra) -> None:
        """
        Unregister the per-RPC channel callbacks now that the call has finished.

        :meth:`~pika.channel.Channel.add_on_close_callback` registers with ``one_shot=False`` and
        each RPC passes a distinct closure, so nothing ever collapses them: left in place they
        accumulate for the life of the channel, and ``CallbackManager.add`` rescans that growing
        list on every later registration.  A completed RPC has no use for its close callback, so
        drop it.

        Removal is scheduled rather than done inline because the channel's callback stack belongs to
        the IOLoop thread.  Ordering is safe: the ``_invoke`` that registers is queued before this,
        and the IOLoop drains its callback queue in FIFO order, so the removal never runs first.

        :param on_chan_close: The close callback registered for this RPC.
        :param extra: Zero-argument callables removing any other per-RPC callbacks.
        """

        def _remove() -> None:
            self._channel.remove_on_close_callback(on_chan_close)
            for remove_one in extra:
                remove_one()

        self._wrapper._schedule_unchecked(_remove)

    def _blocking_rpc(self,
                      method_name: str,
                      channel_method,
                      timeout: float | None,
                      *args,
                      on_sent: Callable[[], None] | None = None,
                      **kwargs) -> Any:
        """
        Execute a channel RPC and block until the broker responds.

        Handles the waiter lifecycle: registers the calling thread's event
        in ``_blocking_waiters``, schedules the RPC on the IOLoop thread,
        waits for the response (or timeout/error), unregisters, and returns
        the result or raises.

        *channel_method* is called on the IOLoop thread with (*args,
        **kwargs, callback=<success_cb>).  The success callback receives
        the broker's response frame and must be accepted as a keyword
        argument named ``callback``.

        :param method_name: Human-readable name for timeout messages.
        :param channel_method: Bound method on the raw channel.
        :param timeout: Seconds to wait.
        :param on_sent: Optional zero-argument callable run on the IOLoop thread
            immediately after *channel_method* writes its frame, ahead of any
            callback already queued behind it.  For state that must be ordered
            against the frame going out rather than against the broker's
            response coming back.
        :returns: The broker response frame.
        :raises TimeoutError: if *timeout* expires.
        :raises Exception: if the connection or channel closes first.
        """
        ready, error = self._register_waiter()
        result = [None]

        def _on_ok(method_frame) -> None:
            result[0] = method_frame
            ready.set()

        def _on_chan_close(ch, reason) -> None:
            if not ready.is_set():
                if error[0] is None:
                    error[0] = reason
                ready.set()

        def _invoke() -> None:
            try:
                self._channel.add_on_close_callback(_on_chan_close)
                channel_method(*args, **kwargs, callback=_on_ok)
            except Exception as exc:
                error[0] = exc
                ready.set()
                return
            # The frame is on the wire.  Run the post-send hook separately: if
            # channel_method delivered its response synchronously (result and
            # ready already set) a raising hook must not clobber that result,
            # since _blocking_rpc checks error ahead of result.  Surface the
            # hook's failure only while the response is still outstanding.
            if on_sent is not None:
                try:
                    on_sent()
                except Exception as exc:
                    if not ready.is_set():
                        error[0] = exc
                        ready.set()
                    else:
                        LOGGER.exception(
                            'on_sent hook raised after %s completed',
                            method_name)

        self._wrapper._schedule_unchecked(_invoke)

        try:
            if not ready.wait(timeout=timeout):
                raise TimeoutError(
                    f'{method_name} timed out after {timeout} seconds')
        finally:
            self._unregister_waiter(ready, error)
            self._release_close_callback(_on_chan_close)

        if error[0] is not None:
            raise _reraisable(error[0], self._wrapper._closed_reason,
                              self._wrapper._closed_reason_tb)
        return result[0]

    def basic_publish(self,
                      exchange,
                      routing_key,
                      body,
                      properties=None,
                      mandatory: bool = False,
                      on_publish=None) -> None:
        """
        Schedule a publish in the IOLoop thread (fire-and-forget).

        Safe to call from any thread simultaneously.

        :param exchange: The exchange to publish to.
        :param routing_key: The routing key to publish with.
        :param body: The message body to publish.
        :param properties: Properties for the message.
        :param mandatory: If True, return unroutable messages to the publisher
        :param on_publish: Optional callback invoked on the **IOLoop thread** immediately after the
            publish frame is written successfully, with the delivery tag (int) as its sole argument.
            Only meaningful when publisher confirms are enabled via :meth:`confirm_delivery`;
            ignored otherwise. Must return quickly (same contract as any
            :meth:`~Connection.add_callback_threadsafe` callback). An exception raised by this
            callback is logged and suppressed rather than propagated, so it does not tear down the
            connection.
        :raises Exception: if the connection is already closed.
        """
        self._check_not_closed()

        def _publish() -> None:
            try:
                self._channel.basic_publish(
                    exchange=exchange,
                    routing_key=routing_key,
                    body=body,
                    properties=properties,
                    mandatory=mandatory,
                )
            except Exception:
                LOGGER.warning('basic_publish failed (channel may have closed)',
                               exc_info=True)
                return
            if self._next_publish_seq_no is not None:
                self._next_publish_seq_no += 1
                if on_publish is not None:
                    # Guard on_publish with _safe_dispatch so an exception
                    # from user code is logged and contained instead of
                    # propagating out of the IOLoop thread and killing the
                    # whole connection (#1687).  Unlike the consumer, confirm,
                    # return and cancel callbacks, which _submit_or_terminate
                    # hands to the work pool, on_publish runs inline here on
                    # the IOLoop thread as the basic_publish docstring
                    # documents; _safe_dispatch adds only the exception guard,
                    # not the pool dispatch.
                    self._safe_dispatch('publish callback', on_publish,
                                        self._next_publish_seq_no)

        self._wrapper._schedule_unchecked(_publish)

    def basic_ack(self, delivery_tag: int = 0, multiple: bool = False) -> None:
        """
        Schedule an acknowledgement in the IOLoop thread (fire-and-forget).

        Safe to call from any thread simultaneously.

        :param delivery_tag: Server-assigned delivery tag
        :param multiple: If True, apply to all messages up to and including this delivery tag
        :raises Exception: if the connection is already closed.
        """
        self._check_not_closed()

        def _ack() -> None:
            try:
                self._channel.basic_ack(delivery_tag=delivery_tag,
                                        multiple=multiple)
            except Exception:
                LOGGER.warning('basic_ack failed (channel may have closed)',
                               exc_info=True)

        self._wrapper._schedule_unchecked(_ack)

    def basic_nack(self,
                   delivery_tag: int = 0,
                   multiple: bool = False,
                   requeue: bool = True) -> None:
        """
        Schedule a negative acknowledgement in the IOLoop thread (fire-and-forget).

        Safe to call from any thread simultaneously.

        :param delivery_tag: Server-assigned delivery tag
        :param multiple: If True, apply to all messages up to and including this delivery tag
        :param requeue: If True, requeue the message on the broker
        :raises Exception: if the connection is already closed.
        """
        self._check_not_closed()

        def _nack() -> None:
            try:
                self._channel.basic_nack(delivery_tag=delivery_tag,
                                         multiple=multiple,
                                         requeue=requeue)
            except Exception:
                LOGGER.warning('basic_nack failed (channel may have closed)',
                               exc_info=True)

        self._wrapper._schedule_unchecked(_nack)

    def basic_reject(self, delivery_tag: int = 0, requeue: bool = True) -> None:
        """
        Schedule a rejection in the IOLoop thread (fire-and-forget).

        Safe to call from any thread simultaneously.

        :param delivery_tag: Server-assigned delivery tag
        :param requeue: If True, requeue the message on the broker
        :raises Exception: if the connection is already closed.
        """
        self._check_not_closed()

        def _reject() -> None:
            try:
                self._channel.basic_reject(delivery_tag=delivery_tag,
                                           requeue=requeue)
            except Exception:
                LOGGER.warning('basic_reject failed (channel may have closed)',
                               exc_info=True)

        self._wrapper._schedule_unchecked(_reject)

    def basic_qos(self,
                  prefetch_size: int = 0,
                  prefetch_count: int = 0,
                  global_qos: bool = False,
                  timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """Set channel QoS and block until Basic.QosOk arrives.

        Safe to call from any thread.

        :param prefetch_size: Prefetch window in octets (0 = no limit).
        :param prefetch_count: Prefetch window in whole messages (0 = no limit).
        :param global_qos: Apply QoS to all consumers on the channel.
        :param timeout: Seconds to wait for the response.
            Defaults to :data:`DEFAULT_RPC_TIMEOUT` (10 s).
            Pass ``None`` to wait indefinitely.
        :returns: The Basic.QosOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'basic_qos',
            self._channel.basic_qos,
            timeout,
            prefetch_size=prefetch_size,
            prefetch_count=prefetch_count,
            global_qos=global_qos,
        )

    def basic_get(
            self,
            queue,
            auto_ack: bool = False,
            timeout: float | None = DEFAULT_RPC_TIMEOUT) -> tuple[None, ...]:
        """
        Get a single message from the broker and block until it arrives.

        Returns a ``(method, properties, body)`` tuple if a message is available, or ``(None, None,
        None)`` if the queue is empty.

        Safe to call from any thread.

        :param queue: The queue to get a message from.
        :param auto_ack: Do not require acknowledgement.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns:``(method, properties, body)`` or ``(None, None, None)``.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        ready, error = self._register_waiter()
        result = [None, None, None]

        def _on_get_ok(ch, method, properties, body) -> None:
            result[0] = method
            result[1] = properties
            result[2] = body
            ready.set()

        def _on_get_empty(method_frame) -> None:
            ready.set()

        def _on_chan_close(ch, reason) -> None:
            if not ready.is_set():
                if error[0] is None:
                    error[0] = reason
                ready.set()

        def _get() -> None:
            try:
                self._channel.add_on_close_callback(_on_chan_close)
                self._channel.add_callback(
                    _on_get_empty,
                    replies=[spec.Basic.GetEmpty],
                    one_shot=True,
                )
                self._channel.basic_get(
                    queue=queue,
                    callback=_on_get_ok,
                    auto_ack=auto_ack,
                )
            except Exception as exc:
                error[0] = exc
                ready.set()

        self._wrapper._schedule_unchecked(_get)

        try:
            if not ready.wait(timeout=timeout):
                raise TimeoutError(
                    f'basic_get timed out after {timeout} seconds')
        finally:
            self._unregister_waiter(ready, error)
            # A one-shot is only consumed if it fires, so the Basic.GetEmpty
            # callback outlives every get that actually returned a message.
            self._release_close_callback(
                _on_chan_close, lambda: self._channel.remove_callback(
                    _on_get_empty, [spec.Basic.GetEmpty]))

        if error[0] is not None:
            raise _reraisable(error[0], self._wrapper._closed_reason,
                              self._wrapper._closed_reason_tb)
        return tuple(result)

    def add_on_cancel_callback(self, callback) -> None:
        """
        Register a callback for server-initiated consumer cancellation.

        The broker sends ``Basic.Cancel`` when a consumer is cancelled by the server (for example,
        when the queue the consumer is bound to is deleted).  Without registering this callback, a
        consumer can silently stop receiving messages.

        Dispatched on the per-channel worker thread (same as delivery callbacks), so the callback
        may safely call any :class:`Channel` method.  The RabbitMQ Java and .NET clients run this
        listener inline on the I/O thread; pika's wrapper deliberately diverges so a slow listener
        cannot stall heartbeats.

        Safe to call from any thread.

        :param callback:``callback(method_frame)`` where *method_frame* contains a
            :class:`pika.spec.Basic.Cancel`.
        :raises Exception: if the connection is already closed.
        """
        self._check_not_closed()

        def _wrapped(method_frame) -> None:
            _submit_or_terminate(
                self._consumer_work_pool, self._wrapper._connection,
                'Server-initiated cancel dropped: work pool shut down',
                self._safe_dispatch, 'cancel listener', callback, method_frame)

        def _register() -> None:
            try:
                self._channel.add_on_cancel_callback(_wrapped)
            except Exception:
                LOGGER.warning('add_on_cancel_callback failed', exc_info=True)

        self._wrapper._schedule_unchecked(_register)

    def add_on_return_callback(self, callback) -> None:
        """
        Register a callback for messages returned by the broker.

        When a message is published with ``mandatory=True`` and cannot be routed to any queue, the
        broker returns it via a :class:`pika.spec.Basic.Return`.  The *callback* receives the
        returned message.

        Dispatched on the per-channel worker thread (same as delivery callbacks), so the callback
        may safely call any :class:`Channel` method.  The RabbitMQ Java and .NET clients run this
        listener inline on the I/O thread; pika's wrapper deliberately diverges so a slow listener
        cannot stall heartbeats.

        Safe to call from any thread.

        :param callback:``callback(channel, method, properties, body)`` where *channel* is this
            :class:`Channel`, *method* is a :class:`pika.spec.Basic.Return`, *properties* is a
            :class:`pika.spec.BasicProperties`, and *body* is :class:`bytes`.
        :raises Exception: if the connection is already closed.
        """
        self._check_not_closed()

        def _wrapped(_raw_ch, method, properties, body) -> None:
            _submit_or_terminate(
                self._consumer_work_pool, self._wrapper._connection,
                'Returned message dropped: work pool shut down',
                self._safe_dispatch, 'return listener', callback, self, method,
                properties, body)

        def _register() -> None:
            try:
                self._channel.add_on_return_callback(_wrapped)
            except Exception:
                LOGGER.warning('add_on_return_callback failed', exc_info=True)

        self._wrapper._schedule_unchecked(_register)

    def confirm_delivery(self,
                         ack_nack_callback,
                         timeout: float | None = DEFAULT_RPC_TIMEOUT):
        """
        Enable publisher confirms and block until Confirm.SelectOk arrives.

        Idempotent: calling this method more than once on the same
        channel - concurrently from several threads included - returns
        the original Confirm.SelectOk without sending a frame to the
        broker or resetting the delivery-tag counter, matching the
        behavior of the RabbitMQ Java and .NET clients.

        Concurrent callers are serialized on a per-channel lock, so a
        second thread waits for the first call's round trip to finish
        before it is handed the cached frame.  Its own *timeout* bounds
        only the wait for the broker's response, not that queueing, so
        it can block for the first caller's timeout on top of its own.

        The *ack_nack_callback* is dispatched on the channel's worker
        thread (same as delivery callbacks), not the IOLoop thread.
        The RabbitMQ Java and .NET clients run this listener inline on
        the I/O thread; pika's wrapper deliberately diverges so a slow
        listener cannot stall heartbeats.

        Safe to call from any thread.

        A :class:`TimeoutError` does not mean confirms are off.  The
        Confirm.Select frame is written before the wait begins and may
        well have reached the broker, which then numbers publishes from
        1, so the delivery-tag counter is armed regardless of whether
        Confirm.SelectOk arrived in time.  A subsequent publish therefore
        still fires its ``on_publish`` and advances
        :attr:`next_publish_seq_no`.  Retry ``confirm_delivery`` rather
        than assuming confirms were left disabled.

        :param ack_nack_callback:
            ``callback(method_frame)`` called for each Basic.Ack or
            Basic.Nack received from the broker.
        :param timeout: Seconds to wait for the response.
            Defaults to :data:`DEFAULT_RPC_TIMEOUT` (10 s).
            Pass ``None`` to wait indefinitely.
        :returns: The Confirm.SelectOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        with self._confirm_lock:
            if self._confirm_select_ok is not None:
                return self._confirm_select_ok

            def _wrapped_ack_nack(method_frame) -> None:
                _submit_or_terminate(
                    self._consumer_work_pool, self._wrapper._connection,
                    'Publisher confirm dropped: work pool shut down',
                    self._safe_dispatch, 'publisher confirm callback',
                    ack_nack_callback, method_frame)

            def _arm_seq_no() -> None:
                # Runs on the IOLoop thread the moment Confirm.Select is
                # written, so it is ordered ahead of every publish the IOLoop
                # has queued behind it.  Arming from the calling thread once the
                # RPC returns instead would miss each publish issued during the
                # round trip: the broker numbers those from 1, leaving the
                # counter permanently behind the broker's delivery tags.
                #
                # Arm only once.  A confirm_delivery that timed out may still
                # have reached the broker, and the retry's second Confirm.Select
                # does not reset the broker's counter, so neither may this.
                if self._next_publish_seq_no is None:
                    self._next_publish_seq_no = 0

            result = self._blocking_rpc(
                'confirm_delivery',
                self._channel.confirm_delivery,
                timeout,
                on_sent=_arm_seq_no,
                ack_nack_callback=_wrapped_ack_nack,
            )
            self._confirm_select_ok = result
            return result

    def basic_consume(self,
                      queue,
                      on_message_callback,
                      auto_ack: bool = False,
                      exclusive: bool = False,
                      consumer_tag=None,
                      arguments=None,
                      timeout: float | None = DEFAULT_RPC_TIMEOUT):
        """
        Register a consumer and block until Basic.ConsumeOk arrives.

        The *on_message_callback* is dispatched on the channel's worker thread, not the IOLoop
        thread.  All :class:`Channel` methods are safe to call from within the callback.

        Because the channel uses a single worker thread, deliveries are processed serially.  A
        callback that blocks (e.g. on a database write or a call to :meth:`queue_declare`) delays
        subsequent deliveries on the same channel until it returns.

        Safe to call from any thread.

        :param queue: Queue to consume from.
        :param on_message_callback:``callback(channel, method, properties, body)``
        :param auto_ack: Disable manual acknowledgement.
        :param exclusive: Request exclusive consumer access.
        :param consumer_tag: Client-provided tag; generated if omitted.
        :param arguments: Additional AMQP arguments.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The consumer tag assigned by the broker.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """

        def _wrapped_callback(ch, method, properties, body) -> None:
            _submit_or_terminate(
                self._consumer_work_pool, self._wrapper._connection,
                'Consumer delivery dropped: work pool shut down',
                self._safe_dispatch, 'consumer callback', on_message_callback,
                self, method, properties, body)

        frame = self._blocking_rpc(
            'basic_consume',
            self._channel.basic_consume,
            timeout,
            queue=queue,
            on_message_callback=_wrapped_callback,
            auto_ack=auto_ack,
            exclusive=exclusive,
            consumer_tag=consumer_tag,
            arguments=arguments,
        )
        return frame.method.consumer_tag

    def basic_cancel(self,
                     consumer_tag,
                     timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Cancel a consumer and block until Basic.CancelOk arrives.

        Safe to call from any thread.

        :param consumer_tag: Tag returned by :meth:`basic_consume`.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Basic.CancelOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'basic_cancel',
            self._channel.basic_cancel,
            timeout,
            consumer_tag=consumer_tag,
        )

    def queue_declare(self,
                      queue,
                      passive: bool = False,
                      durable: bool = False,
                      exclusive: bool = False,
                      auto_delete: bool = False,
                      arguments=None,
                      timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Declare a queue and block the calling thread until Queue.DeclareOk arrives.

        Safe to call from any thread.

        :param queue: The queue name. If empty, the broker will generate a unique name.
        :param passive: If True, only check whether the queue or exchange exists
        :param durable: If True, the queue survives broker restart
        :param exclusive: If True, restrict access to the current connection
        :param auto_delete: If True, delete the queue or exchange when no longer in use
        :param arguments: Custom arguments for the queue declaration
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Queue.DeclareOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'queue_declare',
            self._channel.queue_declare,
            timeout,
            queue=queue,
            passive=passive,
            durable=durable,
            exclusive=exclusive,
            auto_delete=auto_delete,
            arguments=arguments,
        )

    def exchange_declare(self,
                         exchange,
                         exchange_type: str = 'direct',
                         passive: bool = False,
                         durable: bool = False,
                         auto_delete: bool = False,
                         internal: bool = False,
                         arguments=None,
                         timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Declare an exchange and block until Exchange.DeclareOk arrives.

        Safe to call from any thread.

        :param exchange: The exchange name.
        :param exchange_type: The exchange type (direct, fanout, topic, headers).
        :param passive: Only check if the exchange exists.
        :param durable: Survive broker restart.
        :param auto_delete: Delete when no queues are bound.
        :param internal: Can only be published to by other exchanges.
        :param arguments: Custom arguments for the exchange.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Exchange.DeclareOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'exchange_declare',
            self._channel.exchange_declare,
            timeout,
            exchange=exchange,
            exchange_type=exchange_type,
            passive=passive,
            durable=durable,
            auto_delete=auto_delete,
            internal=internal,
            arguments=arguments,
        )

    def queue_bind(self,
                   queue,
                   exchange,
                   routing_key=None,
                   arguments=None,
                   timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Bind a queue to an exchange and block until Queue.BindOk arrives.

        Safe to call from any thread.

        :param queue: The queue to bind.
        :param exchange: The exchange to bind to.
        :param routing_key: The routing key to bind on. Defaults to the queue name.
        :param arguments: Custom arguments for the binding.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Queue.BindOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'queue_bind',
            self._channel.queue_bind,
            timeout,
            queue=queue,
            exchange=exchange,
            routing_key=routing_key,
            arguments=arguments,
        )

    def queue_unbind(self,
                     queue,
                     exchange=None,
                     routing_key=None,
                     arguments=None,
                     timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Unbind a queue from an exchange and block until Queue.UnbindOk arrives.

        Safe to call from any thread.

        :param queue: The queue to unbind.
        :param exchange: The exchange to unbind from.
        :param routing_key: The routing key to unbind. Defaults to the queue name.
        :param arguments: Custom arguments for the unbinding.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Queue.UnbindOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'queue_unbind',
            self._channel.queue_unbind,
            timeout,
            queue=queue,
            exchange=exchange,
            routing_key=routing_key,
            arguments=arguments,
        )

    def queue_delete(self,
                     queue,
                     if_unused: bool = False,
                     if_empty: bool = False,
                     timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Delete a queue and block until Queue.DeleteOk arrives.

        Safe to call from any thread.

        :param queue: The queue to delete.
        :param if_unused: Only delete if the queue has no consumers.
        :param if_empty: Only delete if the queue is empty.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Queue.DeleteOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'queue_delete',
            self._channel.queue_delete,
            timeout,
            queue=queue,
            if_unused=if_unused,
            if_empty=if_empty,
        )

    def queue_purge(self,
                    queue,
                    timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Purge all messages from a queue and block until Queue.PurgeOk arrives.

        Safe to call from any thread.

        :param queue: The queue to purge.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Queue.PurgeOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'queue_purge',
            self._channel.queue_purge,
            timeout,
            queue=queue,
        )

    def exchange_bind(self,
                      destination,
                      source,
                      routing_key: str = '',
                      arguments=None,
                      timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Bind an exchange to another exchange and block until Exchange.BindOk.

        Safe to call from any thread.

        :param destination: The destination exchange to bind.
        :param source: The source exchange to bind to.
        :param routing_key: The routing key to bind on.
        :param arguments: Custom arguments for the binding.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Exchange.BindOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'exchange_bind',
            self._channel.exchange_bind,
            timeout,
            destination=destination,
            source=source,
            routing_key=routing_key,
            arguments=arguments,
        )

    def exchange_unbind(self,
                        destination,
                        source,
                        routing_key: str = '',
                        arguments=None,
                        timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Unbind an exchange from another exchange and block until Exchange.UnbindOk.

        Safe to call from any thread.

        :param destination: The destination exchange to unbind.
        :param source: The source exchange to unbind from.
        :param routing_key: The routing key to unbind.
        :param arguments: Custom arguments for the unbinding.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Exchange.UnbindOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'exchange_unbind',
            self._channel.exchange_unbind,
            timeout,
            destination=destination,
            source=source,
            routing_key=routing_key,
            arguments=arguments,
        )

    def exchange_delete(self,
                        exchange=None,
                        if_unused: bool = False,
                        timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
        """
        Delete an exchange and block until Exchange.DeleteOk arrives.

        Safe to call from any thread.

        :param exchange: The exchange name.
        :param if_unused: Only delete if the exchange has no bindings.
        :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
            (10 s). Pass ``None`` to wait indefinitely.
        :returns: The Exchange.DeleteOk method frame.
        :raises Exception: if the connection is closed before the response arrives.
        :raises TimeoutError: if *timeout* expires before the response arrives.
        """
        return self._blocking_rpc(
            'exchange_delete',
            self._channel.exchange_delete,
            timeout,
            exchange=exchange,
            if_unused=if_unused,
        )

    def close(self,
              reply_code: int = 0,
              reply_text: str = 'Normal shutdown',
              timeout: float | None = 10) -> None:
        """
        Close the channel and block until the Channel.CloseOk arrives.

        Shuts down the consumer work pool after the channel is closed, allowing any in-flight
        delivery callbacks to complete.

        If the channel is already closed or closing, returns immediately. Safe to call from any
        thread.

        *timeout* bounds the whole call: it caps the wait for Channel.CloseOk and then the leftover
        budget caps the consumer-pool drain, so a slow or wedged delivery callback cannot stall
        ``close`` past *timeout*.  A worker still running when the budget is exhausted is left to
        finish on its own (it is a daemon thread).

        :param reply_code: Close reason code to send to the broker.
        :param reply_text: Close reason text to send to the broker.
        :param timeout: Seconds to wait for Channel.CloseOk before treating the channel as closed
            regardless. Defaults to 10 seconds. Pass ``None`` to wait indefinitely.
        :raises Exception: if the connection is closed or the channel is closed by the broker rather
            than by this client.
        """
        ready, error = self._register_waiter()
        deadline = None if timeout is None else time.monotonic() + timeout

        def _close() -> None:
            if self._channel.is_closed or self._channel.is_closing:
                ready.set()
                return

            def _on_channel_close(channel, reason) -> None:
                from pika.exceptions import ChannelClosedByClient
                if not isinstance(reason, ChannelClosedByClient):
                    error[0] = reason
                ready.set()

            try:
                self._channel.add_on_close_callback(_on_channel_close)
                self._channel.close(reply_code=reply_code,
                                    reply_text=reply_text)
            except Exception as exc:
                error[0] = exc
                ready.set()

        self._wrapper._schedule_unchecked(_close)

        try:
            if not ready.wait(timeout=timeout):
                LOGGER.warning('Channel %s close timed out after %s seconds',
                               self._channel.channel_number, timeout)
            # Drain the pool within whatever remains of the caller's budget so
            # the whole close honors *timeout* rather than up to twice it.
            if deadline is None:
                pool_timeout: float | None = None
            else:
                pool_timeout = max(0.0, deadline - time.monotonic())
            self._shutdown_pool(timeout=pool_timeout)
        finally:
            self._unregister_waiter(ready, error)

        if error[0] is not None:
            raise _reraisable(error[0], self._wrapper._closed_reason,
                              self._wrapper._closed_reason_tb)

    def abort(self,
              reply_code: int = 0,
              reply_text: str = 'Normal shutdown',
              timeout: float | None = 10) -> None:
        """
        Close the channel, swallowing any errors.

        Equivalent to :meth:`close` but never raises.  Useful in error-recovery paths where the
        channel may already be in a bad state and the caller only wants a best-effort shutdown.

        Safe to call from any thread.

        :param reply_code: Close reason code to send to the broker.
        :param reply_text: Close reason text to send to the broker.
        :param timeout: Seconds to wait for Channel.CloseOk before treating the channel as closed
            regardless. Defaults to 10 seconds. Pass ``None`` to wait indefinitely.
        """
        try:
            self.close(reply_code=reply_code,
                       reply_text=reply_text,
                       timeout=timeout)
        except Exception:
            LOGGER.debug('channel abort() suppressed error', exc_info=True)

    @property
    def next_publish_seq_no(self) -> int | None:
        """
        The delivery tag that will be assigned to the next published message.

        Returns ``None`` if publisher confirms have not been enabled via :meth:`confirm_delivery`.
        Once confirms are enabled, returns an integer starting at 1 that increments after each
        successful publish, matching the behavior of the RabbitMQ Java and .NET clients.

        This property is safe to read from any thread, but the value is only stable if no other
        thread is concurrently publishing (the counter advances on the IOLoop thread inside the
        scheduled publish callback).
        """
        if self._next_publish_seq_no is None:
            return None
        return self._next_publish_seq_no + 1

    @property
    def channel_number(self):
        return self._channel.channel_number

    @property
    def is_open(self):
        return self._channel.is_open

    @property
    def is_closed(self):
        return self._channel.is_closed

_channel instance-attribute

_channel = channel

_confirm_lock instance-attribute

_confirm_lock = threading.Lock()

_confirm_select_ok instance-attribute

_confirm_select_ok = None

_consumer_work_pool instance-attribute

_consumer_work_pool = _BoundedWorkPool(
    maxsize=work_queue_maxsize,
    put_timeout=work_queue_put_timeout,
    thread_name='pika-consumer',
)

_next_publish_seq_no instance-attribute

_next_publish_seq_no: int | None = None

_pool_shutdown instance-attribute

_pool_shutdown = False

_wrapper instance-attribute

_wrapper = wrapper

channel_number property

channel_number

is_closed property

is_closed

is_open property

is_open

next_publish_seq_no property

next_publish_seq_no: int | None

The delivery tag that will be assigned to the next published message.

Returns None if publisher confirms have not been enabled via :meth:confirm_delivery. Once confirms are enabled, returns an integer starting at 1 that increments after each successful publish, matching the behavior of the RabbitMQ Java and .NET clients.

This property is safe to read from any thread, but the value is only stable if no other thread is concurrently publishing (the counter advances on the IOLoop thread inside the scheduled publish callback).

__init__

__init__(
    channel,
    wrapper,
    work_queue_maxsize: int = DEFAULT_WORK_QUEUE_MAXSIZE,
    work_queue_put_timeout: float = DEFAULT_WORK_QUEUE_PUT_TIMEOUT,
) -> None
Source code in pika/adapters/thread_safe_connection.py
def __init__(
        self,
        channel,
        wrapper,
        work_queue_maxsize: int = DEFAULT_WORK_QUEUE_MAXSIZE,
        work_queue_put_timeout: float = DEFAULT_WORK_QUEUE_PUT_TIMEOUT
) -> None:
    work_queue_put_timeout = _validate_put_timeout(work_queue_put_timeout)
    self._channel = channel
    self._wrapper = wrapper
    self._consumer_work_pool = _BoundedWorkPool(
        maxsize=work_queue_maxsize,
        put_timeout=work_queue_put_timeout,
        thread_name='pika-consumer',
    )
    self._pool_shutdown = False
    self._next_publish_seq_no: int | None = None
    self._confirm_select_ok = None
    # Serializes confirm_delivery so two threads cannot both pass the
    # _confirm_select_ok guard while the first is still waiting for
    # Confirm.SelectOk.  Per-channel rather than the connection-wide
    # _channel_waiters_lock: the call blocks for a broker round trip, which
    # must not stall operations on other channels.
    self._confirm_lock = threading.Lock()

_blocking_rpc

_blocking_rpc(
    method_name: str,
    channel_method,
    timeout: float | None,
    *args,
    on_sent: Callable[[], None] | None = None,
    **kwargs,
) -> Any

Execute a channel RPC and block until the broker responds.

Handles the waiter lifecycle: registers the calling thread's event in _blocking_waiters, schedules the RPC on the IOLoop thread, waits for the response (or timeout/error), unregisters, and returns the result or raises.

channel_method is called on the IOLoop thread with (args, *kwargs, callback=). The success callback receives the broker's response frame and must be accepted as a keyword argument named callback.

PARAMETER DESCRIPTION
method_name

Human-readable name for timeout messages.

TYPE: str

channel_method

Bound method on the raw channel.

timeout

Seconds to wait.

TYPE: float | None

on_sent

Optional zero-argument callable run on the IOLoop thread immediately after channel_method writes its frame, ahead of any callback already queued behind it. For state that must be ordered against the frame going out rather than against the broker's response coming back.

TYPE: Callable[[], None] | None DEFAULT: None

RETURNS DESCRIPTION
Any

The broker response frame.

RAISES DESCRIPTION
TimeoutError

if timeout expires.

Exception

if the connection or channel closes first.

Source code in pika/adapters/thread_safe_connection.py
def _blocking_rpc(self,
                  method_name: str,
                  channel_method,
                  timeout: float | None,
                  *args,
                  on_sent: Callable[[], None] | None = None,
                  **kwargs) -> Any:
    """
    Execute a channel RPC and block until the broker responds.

    Handles the waiter lifecycle: registers the calling thread's event
    in ``_blocking_waiters``, schedules the RPC on the IOLoop thread,
    waits for the response (or timeout/error), unregisters, and returns
    the result or raises.

    *channel_method* is called on the IOLoop thread with (*args,
    **kwargs, callback=<success_cb>).  The success callback receives
    the broker's response frame and must be accepted as a keyword
    argument named ``callback``.

    :param method_name: Human-readable name for timeout messages.
    :param channel_method: Bound method on the raw channel.
    :param timeout: Seconds to wait.
    :param on_sent: Optional zero-argument callable run on the IOLoop thread
        immediately after *channel_method* writes its frame, ahead of any
        callback already queued behind it.  For state that must be ordered
        against the frame going out rather than against the broker's
        response coming back.
    :returns: The broker response frame.
    :raises TimeoutError: if *timeout* expires.
    :raises Exception: if the connection or channel closes first.
    """
    ready, error = self._register_waiter()
    result = [None]

    def _on_ok(method_frame) -> None:
        result[0] = method_frame
        ready.set()

    def _on_chan_close(ch, reason) -> None:
        if not ready.is_set():
            if error[0] is None:
                error[0] = reason
            ready.set()

    def _invoke() -> None:
        try:
            self._channel.add_on_close_callback(_on_chan_close)
            channel_method(*args, **kwargs, callback=_on_ok)
        except Exception as exc:
            error[0] = exc
            ready.set()
            return
        # The frame is on the wire.  Run the post-send hook separately: if
        # channel_method delivered its response synchronously (result and
        # ready already set) a raising hook must not clobber that result,
        # since _blocking_rpc checks error ahead of result.  Surface the
        # hook's failure only while the response is still outstanding.
        if on_sent is not None:
            try:
                on_sent()
            except Exception as exc:
                if not ready.is_set():
                    error[0] = exc
                    ready.set()
                else:
                    LOGGER.exception(
                        'on_sent hook raised after %s completed',
                        method_name)

    self._wrapper._schedule_unchecked(_invoke)

    try:
        if not ready.wait(timeout=timeout):
            raise TimeoutError(
                f'{method_name} timed out after {timeout} seconds')
    finally:
        self._unregister_waiter(ready, error)
        self._release_close_callback(_on_chan_close)

    if error[0] is not None:
        raise _reraisable(error[0], self._wrapper._closed_reason,
                          self._wrapper._closed_reason_tb)
    return result[0]

_check_not_closed

_check_not_closed() -> None

Raise if the connection is known to be closed.

Called from fire-and-forget methods to prevent silently dropping work when the connection is already gone.

Source code in pika/adapters/thread_safe_connection.py
def _check_not_closed(self) -> None:
    """
    Raise if the connection is known to be closed.

    Called from fire-and-forget methods to prevent silently dropping work when the connection is
    already gone.
    """
    with self._wrapper._channel_waiters_lock:
        reason = self._wrapper._closed_reason
        if reason is not None:
            raise _with_close_traceback(reason,
                                        self._wrapper._closed_reason_tb)

_claim_pool_shutdown

_claim_pool_shutdown() -> bool

Claim this channel's pool shutdown and drop it from connection tracking.

Returns True if this call is the first to claim the shutdown - so the caller owns any follow-up work, such as joining the worker - and False if a prior close already claimed it. Being idempotent lets a user :meth:close, the broker-close hook (:meth:_on_broker_close) and the connection-shutdown sweep race while each channel is dropped from _channels exactly once.

The tracking list exists only so :meth:~Connection._shutdown_all_consumer_pools can reach every live pool, and a channel whose pool is shutting down has nothing left to drain; left in place the entries accumulate for the life of the connection, retaining each closed channel's pool, raw channel and per-RPC callbacks.

RETURNS DESCRIPTION
bool

True if this call claimed the shutdown, False if a prior close already did.

Source code in pika/adapters/thread_safe_connection.py
def _claim_pool_shutdown(self) -> bool:
    """
    Claim this channel's pool shutdown and drop it from connection tracking.

    Returns ``True`` if this call is the first to claim the shutdown - so the caller owns any
    follow-up work, such as joining the worker - and ``False`` if a prior close already claimed
    it.  Being idempotent lets a user :meth:`close`, the broker-close hook
    (:meth:`_on_broker_close`) and the connection-shutdown sweep race while each channel is
    dropped from ``_channels`` exactly once.

    The tracking list exists only so :meth:`~Connection._shutdown_all_consumer_pools` can reach
    every live pool, and a channel whose pool is shutting down has nothing left to drain; left
    in place the entries accumulate for the life of the connection, retaining each closed
    channel's pool, raw channel and per-RPC callbacks.

    :returns: ``True`` if this call claimed the shutdown, ``False`` if a prior close already did.
    """
    # Safe to mutate the list while a sweep is in flight:
    # _shutdown_all_consumer_pools snapshots it under this same lock before
    # iterating.
    with self._wrapper._channel_waiters_lock:
        if self._pool_shutdown:
            return False
        self._pool_shutdown = True
        try:
            self._wrapper._channels.remove(self)
        except ValueError:
            pass
    return True

_on_broker_close

_on_broker_close(_channel, reason) -> None

Drop this channel from tracking when the broker closes it on its own.

Registered once per channel, on the IOLoop thread, right after the channel opens. A client close - a user :meth:close, or the per-channel closes a graceful :meth:Connection.close issues - reports :class:~pika.exceptions.ChannelClosedByClient and is left to :meth:_shutdown_pool, so this handles only broker- or error-initiated closes. Without it those channels stay in _channels (with their pool, worker thread, raw channel and per-RPC callbacks) until the connection itself is torn down - the leak #1688 describes, reached through a close path :meth:_shutdown_pool never sees.

Runs on the IOLoop thread, so it must not join the worker: joining a worker that is mid delivery-callback and waiting on this same IOLoop would deadlock, which is why pool joins are otherwise deferred off this thread (see :meth:~Connection._on_connection_closed). It only signals the pool to drain; the daemon worker exits on its own once the channel is closed and no further deliveries arrive. The connection-shutdown sweep, running after the IOLoop stops, joins any channel this has not already dropped.

PARAMETER DESCRIPTION
_channel

The raw channel reporting the close (unused).

reason

The exception describing why the channel closed.

Source code in pika/adapters/thread_safe_connection.py
def _on_broker_close(self, _channel, reason) -> None:
    """
    Drop this channel from tracking when the broker closes it on its own.

    Registered once per channel, on the IOLoop thread, right after the channel opens.  A client
    close - a user :meth:`close`, or the per-channel closes a graceful :meth:`Connection.close`
    issues - reports :class:`~pika.exceptions.ChannelClosedByClient` and is left to
    :meth:`_shutdown_pool`, so this handles only broker- or error-initiated closes.  Without it
    those channels stay in ``_channels`` (with their pool, worker thread, raw channel and
    per-RPC callbacks) until the connection itself is torn down - the leak #1688 describes,
    reached through a close path :meth:`_shutdown_pool` never sees.

    Runs on the IOLoop thread, so it must not join the worker: joining a worker that is mid
    delivery-callback and waiting on this same IOLoop would deadlock, which is why pool joins
    are otherwise deferred off this thread (see :meth:`~Connection._on_connection_closed`).  It
    only signals the pool to drain; the daemon worker exits on its own once the channel is
    closed and no further deliveries arrive.  The connection-shutdown sweep, running after the
    IOLoop stops, joins any channel this has not already dropped.

    :param _channel: The raw channel reporting the close (unused).
    :param reason: The exception describing why the channel closed.
    """
    from pika.exceptions import ChannelClosedByClient
    if isinstance(reason, ChannelClosedByClient):
        return
    if self._claim_pool_shutdown():
        self._consumer_work_pool.shutdown(wait=False)

_register_waiter

_register_waiter() -> tuple[
    Event, list[BaseException | None]
]

Create and register a blocking waiter.

RETURNS DESCRIPTION
tuple[Event, list[BaseException | None]]

(ready, error) tuple for use with ready.wait()

RAISES DESCRIPTION
Exception

if the connection is already closed.

Source code in pika/adapters/thread_safe_connection.py
def _register_waiter(self) -> tuple[Event, list[BaseException | None]]:
    """
    Create and register a blocking waiter.

    :returns: (ready, error) tuple for use with ``ready.wait()``
    :raises Exception: if the connection is already closed.
    """
    ready = threading.Event()
    error: list[BaseException | None] = [None]
    with self._wrapper._channel_waiters_lock:
        reason = self._wrapper._closed_reason
        if reason is not None:
            raise _with_close_traceback(reason,
                                        self._wrapper._closed_reason_tb)
        self._wrapper._blocking_waiters.append((ready, error))
    return ready, error

_release_close_callback

_release_close_callback(on_chan_close, *extra) -> None

Unregister the per-RPC channel callbacks now that the call has finished.

:meth:~pika.channel.Channel.add_on_close_callback registers with one_shot=False and each RPC passes a distinct closure, so nothing ever collapses them: left in place they accumulate for the life of the channel, and CallbackManager.add rescans that growing list on every later registration. A completed RPC has no use for its close callback, so drop it.

Removal is scheduled rather than done inline because the channel's callback stack belongs to the IOLoop thread. Ordering is safe: the _invoke that registers is queued before this, and the IOLoop drains its callback queue in FIFO order, so the removal never runs first.

PARAMETER DESCRIPTION
on_chan_close

The close callback registered for this RPC.

extra

Zero-argument callables removing any other per-RPC callbacks.

DEFAULT: ()

Source code in pika/adapters/thread_safe_connection.py
def _release_close_callback(self, on_chan_close, *extra) -> None:
    """
    Unregister the per-RPC channel callbacks now that the call has finished.

    :meth:`~pika.channel.Channel.add_on_close_callback` registers with ``one_shot=False`` and
    each RPC passes a distinct closure, so nothing ever collapses them: left in place they
    accumulate for the life of the channel, and ``CallbackManager.add`` rescans that growing
    list on every later registration.  A completed RPC has no use for its close callback, so
    drop it.

    Removal is scheduled rather than done inline because the channel's callback stack belongs to
    the IOLoop thread.  Ordering is safe: the ``_invoke`` that registers is queued before this,
    and the IOLoop drains its callback queue in FIFO order, so the removal never runs first.

    :param on_chan_close: The close callback registered for this RPC.
    :param extra: Zero-argument callables removing any other per-RPC callbacks.
    """

    def _remove() -> None:
        self._channel.remove_on_close_callback(on_chan_close)
        for remove_one in extra:
            remove_one()

    self._wrapper._schedule_unchecked(_remove)

_safe_dispatch staticmethod

_safe_dispatch(label, callback, *args) -> None

Run callback, logging any exception instead of letting it propagate.

Wraps a user callback so one failing dispatch neither escapes into the caller nor prevents later dispatches from running. Most callers submit the wrapped callback to a :class:_BoundedWorkPool, so it runs on a pool worker; :meth:basic_publish calls it inline on the IOLoop thread. Either way the exception is contained here.

PARAMETER DESCRIPTION
label

Human-readable name of the callback for log lines.

callback

The user callback.

args

Positional arguments forwarded to callback.

DEFAULT: ()

Source code in pika/adapters/thread_safe_connection.py
@staticmethod
def _safe_dispatch(label, callback, *args) -> None:
    """
    Run *callback*, logging any exception instead of letting it propagate.

    Wraps a user callback so one failing dispatch neither escapes into the caller nor prevents
    later dispatches from running.  Most callers submit the wrapped callback to a
    :class:`_BoundedWorkPool`, so it runs on a pool worker; :meth:`basic_publish` calls it
    inline on the IOLoop thread. Either way the exception is contained here.

    :param label: Human-readable name of the callback for log lines.
    :param callback: The user callback.
    :param args: Positional arguments forwarded to *callback*.
    """
    try:
        callback(*args)
    except Exception:
        LOGGER.exception('Unhandled exception in %s', label)

_shutdown_pool

_shutdown_pool(timeout: float | None = None) -> None

Shut down the consumer work pool, joining its worker.

Claims the shutdown and drops the channel from tracking via :meth:_claim_pool_shutdown, then joins the worker. Used by the paths that can afford to wait off the IOLoop thread: a user :meth:close (on the caller's thread) and the connection-shutdown sweep (after the IOLoop has stopped).

PARAMETER DESCRIPTION
timeout

Seconds to wait for the worker to drain and exit. None waits indefinitely. A wedged worker that outlives a finite timeout is left running (it is a daemon thread) rather than stalling the caller; a warning is logged.

TYPE: float | None DEFAULT: None

Source code in pika/adapters/thread_safe_connection.py
def _shutdown_pool(self, timeout: float | None = None) -> None:
    """
    Shut down the consumer work pool, joining its worker.

    Claims the shutdown and drops the channel from tracking via :meth:`_claim_pool_shutdown`,
    then joins the worker.  Used by the paths that can afford to wait off the IOLoop thread: a
    user :meth:`close` (on the caller's thread) and the connection-shutdown sweep (after the
    IOLoop has stopped).

    :param timeout: Seconds to wait for the worker to drain and exit. ``None`` waits
        indefinitely. A wedged worker that outlives a finite *timeout* is left running (it is a
        daemon thread) rather than stalling the caller; a warning is logged.
    """
    # The join stays outside the lock: it can block for the whole timeout,
    # and holding the connection-wide lock that long would stall every other
    # channel operation.
    if not self._claim_pool_shutdown():
        return
    if not self._consumer_work_pool.shutdown(wait=True, timeout=timeout):
        LOGGER.warning(
            'Channel %s consumer work pool did not drain within '
            '%s seconds; abandoning its worker thread',
            self._channel.channel_number, timeout)

_signal_pool_shutdown

_signal_pool_shutdown() -> None

Tell the consumer pool to begin draining without joining its worker.

Sets the pool's shutdown flag and wakes an idle worker, but returns immediately. Used to signal every channel's pool up front so their workers drain concurrently before :meth:_shutdown_pool joins them one at a time under a shared budget.

Source code in pika/adapters/thread_safe_connection.py
def _signal_pool_shutdown(self) -> None:
    """
    Tell the consumer pool to begin draining without joining its worker.

    Sets the pool's shutdown flag and wakes an idle worker, but returns immediately.  Used to
    signal every channel's pool up front so their workers drain concurrently before
    :meth:`_shutdown_pool` joins them one at a time under a shared budget.
    """
    self._consumer_work_pool.shutdown(wait=False)

_unregister_waiter

_unregister_waiter(
    ready: Event, error: list[BaseException | None]
) -> None

Remove a waiter from the blocking list.

PARAMETER DESCRIPTION
ready

Threading event signalling that the RPC response has arrived

TYPE: Event

error

list containing an exception if the RPC response was an error, or None if it was successful

TYPE: list[BaseException | None]

Source code in pika/adapters/thread_safe_connection.py
def _unregister_waiter(self, ready: Event,
                       error: list[BaseException | None]) -> None:
    """
    Remove a waiter from the blocking list.

    :param ready: Threading event signalling that the RPC response has
        arrived
    :param error: list containing an exception if the RPC response was an
        error, or None if it was successful
    """
    with self._wrapper._channel_waiters_lock:
        try:
            self._wrapper._blocking_waiters.remove((ready, error))
        except ValueError:
            pass

abort

abort(
    reply_code: int = 0,
    reply_text: str = 'Normal shutdown',
    timeout: float | None = 10,
) -> None

Close the channel, swallowing any errors.

Equivalent to :meth:close but never raises. Useful in error-recovery paths where the channel may already be in a bad state and the caller only wants a best-effort shutdown.

Safe to call from any thread.

PARAMETER DESCRIPTION
reply_code

Close reason code to send to the broker.

TYPE: int DEFAULT: 0

reply_text

Close reason text to send to the broker.

TYPE: str DEFAULT: 'Normal shutdown'

timeout

Seconds to wait for Channel.CloseOk before treating the channel as closed regardless. Defaults to 10 seconds. Pass None to wait indefinitely.

TYPE: float | None DEFAULT: 10

Source code in pika/adapters/thread_safe_connection.py
def abort(self,
          reply_code: int = 0,
          reply_text: str = 'Normal shutdown',
          timeout: float | None = 10) -> None:
    """
    Close the channel, swallowing any errors.

    Equivalent to :meth:`close` but never raises.  Useful in error-recovery paths where the
    channel may already be in a bad state and the caller only wants a best-effort shutdown.

    Safe to call from any thread.

    :param reply_code: Close reason code to send to the broker.
    :param reply_text: Close reason text to send to the broker.
    :param timeout: Seconds to wait for Channel.CloseOk before treating the channel as closed
        regardless. Defaults to 10 seconds. Pass ``None`` to wait indefinitely.
    """
    try:
        self.close(reply_code=reply_code,
                   reply_text=reply_text,
                   timeout=timeout)
    except Exception:
        LOGGER.debug('channel abort() suppressed error', exc_info=True)

add_on_cancel_callback

add_on_cancel_callback(callback) -> None

Register a callback for server-initiated consumer cancellation.

The broker sends Basic.Cancel when a consumer is cancelled by the server (for example, when the queue the consumer is bound to is deleted). Without registering this callback, a consumer can silently stop receiving messages.

Dispatched on the per-channel worker thread (same as delivery callbacks), so the callback may safely call any :class:Channel method. The RabbitMQ Java and .NET clients run this listener inline on the I/O thread; pika's wrapper deliberately diverges so a slow listener cannot stall heartbeats.

Safe to call from any thread.

PARAMETER DESCRIPTION
callback

callback(method_frame) where method_frame contains a :class:pika.spec.Basic.Cancel.

RAISES DESCRIPTION
Exception

if the connection is already closed.

Source code in pika/adapters/thread_safe_connection.py
def add_on_cancel_callback(self, callback) -> None:
    """
    Register a callback for server-initiated consumer cancellation.

    The broker sends ``Basic.Cancel`` when a consumer is cancelled by the server (for example,
    when the queue the consumer is bound to is deleted).  Without registering this callback, a
    consumer can silently stop receiving messages.

    Dispatched on the per-channel worker thread (same as delivery callbacks), so the callback
    may safely call any :class:`Channel` method.  The RabbitMQ Java and .NET clients run this
    listener inline on the I/O thread; pika's wrapper deliberately diverges so a slow listener
    cannot stall heartbeats.

    Safe to call from any thread.

    :param callback:``callback(method_frame)`` where *method_frame* contains a
        :class:`pika.spec.Basic.Cancel`.
    :raises Exception: if the connection is already closed.
    """
    self._check_not_closed()

    def _wrapped(method_frame) -> None:
        _submit_or_terminate(
            self._consumer_work_pool, self._wrapper._connection,
            'Server-initiated cancel dropped: work pool shut down',
            self._safe_dispatch, 'cancel listener', callback, method_frame)

    def _register() -> None:
        try:
            self._channel.add_on_cancel_callback(_wrapped)
        except Exception:
            LOGGER.warning('add_on_cancel_callback failed', exc_info=True)

    self._wrapper._schedule_unchecked(_register)

add_on_return_callback

add_on_return_callback(callback) -> None

Register a callback for messages returned by the broker.

When a message is published with mandatory=True and cannot be routed to any queue, the broker returns it via a :class:pika.spec.Basic.Return. The callback receives the returned message.

Dispatched on the per-channel worker thread (same as delivery callbacks), so the callback may safely call any :class:Channel method. The RabbitMQ Java and .NET clients run this listener inline on the I/O thread; pika's wrapper deliberately diverges so a slow listener cannot stall heartbeats.

Safe to call from any thread.

PARAMETER DESCRIPTION
callback

callback(channel, method, properties, body) where channel is this :class:Channel, method is a :class:pika.spec.Basic.Return, properties is a :class:pika.spec.BasicProperties, and body is :class:bytes.

RAISES DESCRIPTION
Exception

if the connection is already closed.

Source code in pika/adapters/thread_safe_connection.py
def add_on_return_callback(self, callback) -> None:
    """
    Register a callback for messages returned by the broker.

    When a message is published with ``mandatory=True`` and cannot be routed to any queue, the
    broker returns it via a :class:`pika.spec.Basic.Return`.  The *callback* receives the
    returned message.

    Dispatched on the per-channel worker thread (same as delivery callbacks), so the callback
    may safely call any :class:`Channel` method.  The RabbitMQ Java and .NET clients run this
    listener inline on the I/O thread; pika's wrapper deliberately diverges so a slow listener
    cannot stall heartbeats.

    Safe to call from any thread.

    :param callback:``callback(channel, method, properties, body)`` where *channel* is this
        :class:`Channel`, *method* is a :class:`pika.spec.Basic.Return`, *properties* is a
        :class:`pika.spec.BasicProperties`, and *body* is :class:`bytes`.
    :raises Exception: if the connection is already closed.
    """
    self._check_not_closed()

    def _wrapped(_raw_ch, method, properties, body) -> None:
        _submit_or_terminate(
            self._consumer_work_pool, self._wrapper._connection,
            'Returned message dropped: work pool shut down',
            self._safe_dispatch, 'return listener', callback, self, method,
            properties, body)

    def _register() -> None:
        try:
            self._channel.add_on_return_callback(_wrapped)
        except Exception:
            LOGGER.warning('add_on_return_callback failed', exc_info=True)

    self._wrapper._schedule_unchecked(_register)

basic_ack

basic_ack(
    delivery_tag: int = 0, multiple: bool = False
) -> None

Schedule an acknowledgement in the IOLoop thread (fire-and-forget).

Safe to call from any thread simultaneously.

PARAMETER DESCRIPTION
delivery_tag

Server-assigned delivery tag

TYPE: int DEFAULT: 0

multiple

If True, apply to all messages up to and including this delivery tag

TYPE: bool DEFAULT: False

RAISES DESCRIPTION
Exception

if the connection is already closed.

Source code in pika/adapters/thread_safe_connection.py
def basic_ack(self, delivery_tag: int = 0, multiple: bool = False) -> None:
    """
    Schedule an acknowledgement in the IOLoop thread (fire-and-forget).

    Safe to call from any thread simultaneously.

    :param delivery_tag: Server-assigned delivery tag
    :param multiple: If True, apply to all messages up to and including this delivery tag
    :raises Exception: if the connection is already closed.
    """
    self._check_not_closed()

    def _ack() -> None:
        try:
            self._channel.basic_ack(delivery_tag=delivery_tag,
                                    multiple=multiple)
        except Exception:
            LOGGER.warning('basic_ack failed (channel may have closed)',
                           exc_info=True)

    self._wrapper._schedule_unchecked(_ack)

basic_cancel

basic_cancel(
    consumer_tag,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Cancel a consumer and block until Basic.CancelOk arrives.

Safe to call from any thread.

PARAMETER DESCRIPTION
consumer_tag

Tag returned by :meth:basic_consume.

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Basic.CancelOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def basic_cancel(self,
                 consumer_tag,
                 timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Cancel a consumer and block until Basic.CancelOk arrives.

    Safe to call from any thread.

    :param consumer_tag: Tag returned by :meth:`basic_consume`.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Basic.CancelOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'basic_cancel',
        self._channel.basic_cancel,
        timeout,
        consumer_tag=consumer_tag,
    )

basic_consume

basic_consume(
    queue,
    on_message_callback,
    auto_ack: bool = False,
    exclusive: bool = False,
    consumer_tag=None,
    arguments=None,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
)

Register a consumer and block until Basic.ConsumeOk arrives.

The on_message_callback is dispatched on the channel's worker thread, not the IOLoop thread. All :class:Channel methods are safe to call from within the callback.

Because the channel uses a single worker thread, deliveries are processed serially. A callback that blocks (e.g. on a database write or a call to :meth:queue_declare) delays subsequent deliveries on the same channel until it returns.

Safe to call from any thread.

PARAMETER DESCRIPTION
queue

Queue to consume from.

on_message_callback

callback(channel, method, properties, body)

auto_ack

Disable manual acknowledgement.

TYPE: bool DEFAULT: False

exclusive

Request exclusive consumer access.

TYPE: bool DEFAULT: False

consumer_tag

Client-provided tag; generated if omitted.

DEFAULT: None

arguments

Additional AMQP arguments.

DEFAULT: None

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION

The consumer tag assigned by the broker.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def basic_consume(self,
                  queue,
                  on_message_callback,
                  auto_ack: bool = False,
                  exclusive: bool = False,
                  consumer_tag=None,
                  arguments=None,
                  timeout: float | None = DEFAULT_RPC_TIMEOUT):
    """
    Register a consumer and block until Basic.ConsumeOk arrives.

    The *on_message_callback* is dispatched on the channel's worker thread, not the IOLoop
    thread.  All :class:`Channel` methods are safe to call from within the callback.

    Because the channel uses a single worker thread, deliveries are processed serially.  A
    callback that blocks (e.g. on a database write or a call to :meth:`queue_declare`) delays
    subsequent deliveries on the same channel until it returns.

    Safe to call from any thread.

    :param queue: Queue to consume from.
    :param on_message_callback:``callback(channel, method, properties, body)``
    :param auto_ack: Disable manual acknowledgement.
    :param exclusive: Request exclusive consumer access.
    :param consumer_tag: Client-provided tag; generated if omitted.
    :param arguments: Additional AMQP arguments.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The consumer tag assigned by the broker.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """

    def _wrapped_callback(ch, method, properties, body) -> None:
        _submit_or_terminate(
            self._consumer_work_pool, self._wrapper._connection,
            'Consumer delivery dropped: work pool shut down',
            self._safe_dispatch, 'consumer callback', on_message_callback,
            self, method, properties, body)

    frame = self._blocking_rpc(
        'basic_consume',
        self._channel.basic_consume,
        timeout,
        queue=queue,
        on_message_callback=_wrapped_callback,
        auto_ack=auto_ack,
        exclusive=exclusive,
        consumer_tag=consumer_tag,
        arguments=arguments,
    )
    return frame.method.consumer_tag

basic_get

basic_get(
    queue,
    auto_ack: bool = False,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> tuple[None, ...]

Get a single message from the broker and block until it arrives.

Returns a (method, properties, body) tuple if a message is available, or (None, None, None) if the queue is empty.

Safe to call from any thread.

PARAMETER DESCRIPTION
queue

The queue to get a message from.

auto_ack

Do not require acknowledgement.

TYPE: bool DEFAULT: False

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
tuple[None, ...]

(method, properties, body) or (None, None, None).

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def basic_get(
        self,
        queue,
        auto_ack: bool = False,
        timeout: float | None = DEFAULT_RPC_TIMEOUT) -> tuple[None, ...]:
    """
    Get a single message from the broker and block until it arrives.

    Returns a ``(method, properties, body)`` tuple if a message is available, or ``(None, None,
    None)`` if the queue is empty.

    Safe to call from any thread.

    :param queue: The queue to get a message from.
    :param auto_ack: Do not require acknowledgement.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns:``(method, properties, body)`` or ``(None, None, None)``.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    ready, error = self._register_waiter()
    result = [None, None, None]

    def _on_get_ok(ch, method, properties, body) -> None:
        result[0] = method
        result[1] = properties
        result[2] = body
        ready.set()

    def _on_get_empty(method_frame) -> None:
        ready.set()

    def _on_chan_close(ch, reason) -> None:
        if not ready.is_set():
            if error[0] is None:
                error[0] = reason
            ready.set()

    def _get() -> None:
        try:
            self._channel.add_on_close_callback(_on_chan_close)
            self._channel.add_callback(
                _on_get_empty,
                replies=[spec.Basic.GetEmpty],
                one_shot=True,
            )
            self._channel.basic_get(
                queue=queue,
                callback=_on_get_ok,
                auto_ack=auto_ack,
            )
        except Exception as exc:
            error[0] = exc
            ready.set()

    self._wrapper._schedule_unchecked(_get)

    try:
        if not ready.wait(timeout=timeout):
            raise TimeoutError(
                f'basic_get timed out after {timeout} seconds')
    finally:
        self._unregister_waiter(ready, error)
        # A one-shot is only consumed if it fires, so the Basic.GetEmpty
        # callback outlives every get that actually returned a message.
        self._release_close_callback(
            _on_chan_close, lambda: self._channel.remove_callback(
                _on_get_empty, [spec.Basic.GetEmpty]))

    if error[0] is not None:
        raise _reraisable(error[0], self._wrapper._closed_reason,
                          self._wrapper._closed_reason_tb)
    return tuple(result)

basic_nack

basic_nack(
    delivery_tag: int = 0,
    multiple: bool = False,
    requeue: bool = True,
) -> None

Schedule a negative acknowledgement in the IOLoop thread (fire-and-forget).

Safe to call from any thread simultaneously.

PARAMETER DESCRIPTION
delivery_tag

Server-assigned delivery tag

TYPE: int DEFAULT: 0

multiple

If True, apply to all messages up to and including this delivery tag

TYPE: bool DEFAULT: False

requeue

If True, requeue the message on the broker

TYPE: bool DEFAULT: True

RAISES DESCRIPTION
Exception

if the connection is already closed.

Source code in pika/adapters/thread_safe_connection.py
def basic_nack(self,
               delivery_tag: int = 0,
               multiple: bool = False,
               requeue: bool = True) -> None:
    """
    Schedule a negative acknowledgement in the IOLoop thread (fire-and-forget).

    Safe to call from any thread simultaneously.

    :param delivery_tag: Server-assigned delivery tag
    :param multiple: If True, apply to all messages up to and including this delivery tag
    :param requeue: If True, requeue the message on the broker
    :raises Exception: if the connection is already closed.
    """
    self._check_not_closed()

    def _nack() -> None:
        try:
            self._channel.basic_nack(delivery_tag=delivery_tag,
                                     multiple=multiple,
                                     requeue=requeue)
        except Exception:
            LOGGER.warning('basic_nack failed (channel may have closed)',
                           exc_info=True)

    self._wrapper._schedule_unchecked(_nack)

basic_publish

basic_publish(
    exchange,
    routing_key,
    body,
    properties=None,
    mandatory: bool = False,
    on_publish=None,
) -> None

Schedule a publish in the IOLoop thread (fire-and-forget).

Safe to call from any thread simultaneously.

PARAMETER DESCRIPTION
exchange

The exchange to publish to.

routing_key

The routing key to publish with.

body

The message body to publish.

properties

Properties for the message.

DEFAULT: None

mandatory

If True, return unroutable messages to the publisher

TYPE: bool DEFAULT: False

on_publish

Optional callback invoked on the IOLoop thread immediately after the publish frame is written successfully, with the delivery tag (int) as its sole argument. Only meaningful when publisher confirms are enabled via :meth:confirm_delivery; ignored otherwise. Must return quickly (same contract as any :meth:~Connection.add_callback_threadsafe callback). An exception raised by this callback is logged and suppressed rather than propagated, so it does not tear down the connection.

DEFAULT: None

RAISES DESCRIPTION
Exception

if the connection is already closed.

Source code in pika/adapters/thread_safe_connection.py
def basic_publish(self,
                  exchange,
                  routing_key,
                  body,
                  properties=None,
                  mandatory: bool = False,
                  on_publish=None) -> None:
    """
    Schedule a publish in the IOLoop thread (fire-and-forget).

    Safe to call from any thread simultaneously.

    :param exchange: The exchange to publish to.
    :param routing_key: The routing key to publish with.
    :param body: The message body to publish.
    :param properties: Properties for the message.
    :param mandatory: If True, return unroutable messages to the publisher
    :param on_publish: Optional callback invoked on the **IOLoop thread** immediately after the
        publish frame is written successfully, with the delivery tag (int) as its sole argument.
        Only meaningful when publisher confirms are enabled via :meth:`confirm_delivery`;
        ignored otherwise. Must return quickly (same contract as any
        :meth:`~Connection.add_callback_threadsafe` callback). An exception raised by this
        callback is logged and suppressed rather than propagated, so it does not tear down the
        connection.
    :raises Exception: if the connection is already closed.
    """
    self._check_not_closed()

    def _publish() -> None:
        try:
            self._channel.basic_publish(
                exchange=exchange,
                routing_key=routing_key,
                body=body,
                properties=properties,
                mandatory=mandatory,
            )
        except Exception:
            LOGGER.warning('basic_publish failed (channel may have closed)',
                           exc_info=True)
            return
        if self._next_publish_seq_no is not None:
            self._next_publish_seq_no += 1
            if on_publish is not None:
                # Guard on_publish with _safe_dispatch so an exception
                # from user code is logged and contained instead of
                # propagating out of the IOLoop thread and killing the
                # whole connection (#1687).  Unlike the consumer, confirm,
                # return and cancel callbacks, which _submit_or_terminate
                # hands to the work pool, on_publish runs inline here on
                # the IOLoop thread as the basic_publish docstring
                # documents; _safe_dispatch adds only the exception guard,
                # not the pool dispatch.
                self._safe_dispatch('publish callback', on_publish,
                                    self._next_publish_seq_no)

    self._wrapper._schedule_unchecked(_publish)

basic_qos

basic_qos(
    prefetch_size: int = 0,
    prefetch_count: int = 0,
    global_qos: bool = False,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Set channel QoS and block until Basic.QosOk arrives.

Safe to call from any thread.

PARAMETER DESCRIPTION
prefetch_size

Prefetch window in octets (0 = no limit).

TYPE: int DEFAULT: 0

prefetch_count

Prefetch window in whole messages (0 = no limit).

TYPE: int DEFAULT: 0

global_qos

Apply QoS to all consumers on the channel.

TYPE: bool DEFAULT: False

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Basic.QosOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def basic_qos(self,
              prefetch_size: int = 0,
              prefetch_count: int = 0,
              global_qos: bool = False,
              timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """Set channel QoS and block until Basic.QosOk arrives.

    Safe to call from any thread.

    :param prefetch_size: Prefetch window in octets (0 = no limit).
    :param prefetch_count: Prefetch window in whole messages (0 = no limit).
    :param global_qos: Apply QoS to all consumers on the channel.
    :param timeout: Seconds to wait for the response.
        Defaults to :data:`DEFAULT_RPC_TIMEOUT` (10 s).
        Pass ``None`` to wait indefinitely.
    :returns: The Basic.QosOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'basic_qos',
        self._channel.basic_qos,
        timeout,
        prefetch_size=prefetch_size,
        prefetch_count=prefetch_count,
        global_qos=global_qos,
    )

basic_reject

basic_reject(
    delivery_tag: int = 0, requeue: bool = True
) -> None

Schedule a rejection in the IOLoop thread (fire-and-forget).

Safe to call from any thread simultaneously.

PARAMETER DESCRIPTION
delivery_tag

Server-assigned delivery tag

TYPE: int DEFAULT: 0

requeue

If True, requeue the message on the broker

TYPE: bool DEFAULT: True

RAISES DESCRIPTION
Exception

if the connection is already closed.

Source code in pika/adapters/thread_safe_connection.py
def basic_reject(self, delivery_tag: int = 0, requeue: bool = True) -> None:
    """
    Schedule a rejection in the IOLoop thread (fire-and-forget).

    Safe to call from any thread simultaneously.

    :param delivery_tag: Server-assigned delivery tag
    :param requeue: If True, requeue the message on the broker
    :raises Exception: if the connection is already closed.
    """
    self._check_not_closed()

    def _reject() -> None:
        try:
            self._channel.basic_reject(delivery_tag=delivery_tag,
                                       requeue=requeue)
        except Exception:
            LOGGER.warning('basic_reject failed (channel may have closed)',
                           exc_info=True)

    self._wrapper._schedule_unchecked(_reject)

close

close(
    reply_code: int = 0,
    reply_text: str = 'Normal shutdown',
    timeout: float | None = 10,
) -> None

Close the channel and block until the Channel.CloseOk arrives.

Shuts down the consumer work pool after the channel is closed, allowing any in-flight delivery callbacks to complete.

If the channel is already closed or closing, returns immediately. Safe to call from any thread.

timeout bounds the whole call: it caps the wait for Channel.CloseOk and then the leftover budget caps the consumer-pool drain, so a slow or wedged delivery callback cannot stall close past timeout. A worker still running when the budget is exhausted is left to finish on its own (it is a daemon thread).

PARAMETER DESCRIPTION
reply_code

Close reason code to send to the broker.

TYPE: int DEFAULT: 0

reply_text

Close reason text to send to the broker.

TYPE: str DEFAULT: 'Normal shutdown'

timeout

Seconds to wait for Channel.CloseOk before treating the channel as closed regardless. Defaults to 10 seconds. Pass None to wait indefinitely.

TYPE: float | None DEFAULT: 10

RAISES DESCRIPTION
Exception

if the connection is closed or the channel is closed by the broker rather than by this client.

Source code in pika/adapters/thread_safe_connection.py
def close(self,
          reply_code: int = 0,
          reply_text: str = 'Normal shutdown',
          timeout: float | None = 10) -> None:
    """
    Close the channel and block until the Channel.CloseOk arrives.

    Shuts down the consumer work pool after the channel is closed, allowing any in-flight
    delivery callbacks to complete.

    If the channel is already closed or closing, returns immediately. Safe to call from any
    thread.

    *timeout* bounds the whole call: it caps the wait for Channel.CloseOk and then the leftover
    budget caps the consumer-pool drain, so a slow or wedged delivery callback cannot stall
    ``close`` past *timeout*.  A worker still running when the budget is exhausted is left to
    finish on its own (it is a daemon thread).

    :param reply_code: Close reason code to send to the broker.
    :param reply_text: Close reason text to send to the broker.
    :param timeout: Seconds to wait for Channel.CloseOk before treating the channel as closed
        regardless. Defaults to 10 seconds. Pass ``None`` to wait indefinitely.
    :raises Exception: if the connection is closed or the channel is closed by the broker rather
        than by this client.
    """
    ready, error = self._register_waiter()
    deadline = None if timeout is None else time.monotonic() + timeout

    def _close() -> None:
        if self._channel.is_closed or self._channel.is_closing:
            ready.set()
            return

        def _on_channel_close(channel, reason) -> None:
            from pika.exceptions import ChannelClosedByClient
            if not isinstance(reason, ChannelClosedByClient):
                error[0] = reason
            ready.set()

        try:
            self._channel.add_on_close_callback(_on_channel_close)
            self._channel.close(reply_code=reply_code,
                                reply_text=reply_text)
        except Exception as exc:
            error[0] = exc
            ready.set()

    self._wrapper._schedule_unchecked(_close)

    try:
        if not ready.wait(timeout=timeout):
            LOGGER.warning('Channel %s close timed out after %s seconds',
                           self._channel.channel_number, timeout)
        # Drain the pool within whatever remains of the caller's budget so
        # the whole close honors *timeout* rather than up to twice it.
        if deadline is None:
            pool_timeout: float | None = None
        else:
            pool_timeout = max(0.0, deadline - time.monotonic())
        self._shutdown_pool(timeout=pool_timeout)
    finally:
        self._unregister_waiter(ready, error)

    if error[0] is not None:
        raise _reraisable(error[0], self._wrapper._closed_reason,
                          self._wrapper._closed_reason_tb)

confirm_delivery

confirm_delivery(
    ack_nack_callback,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
)

Enable publisher confirms and block until Confirm.SelectOk arrives.

Idempotent: calling this method more than once on the same channel - concurrently from several threads included - returns the original Confirm.SelectOk without sending a frame to the broker or resetting the delivery-tag counter, matching the behavior of the RabbitMQ Java and .NET clients.

Concurrent callers are serialized on a per-channel lock, so a second thread waits for the first call's round trip to finish before it is handed the cached frame. Its own timeout bounds only the wait for the broker's response, not that queueing, so it can block for the first caller's timeout on top of its own.

The ack_nack_callback is dispatched on the channel's worker thread (same as delivery callbacks), not the IOLoop thread. The RabbitMQ Java and .NET clients run this listener inline on the I/O thread; pika's wrapper deliberately diverges so a slow listener cannot stall heartbeats.

Safe to call from any thread.

A :class:TimeoutError does not mean confirms are off. The Confirm.Select frame is written before the wait begins and may well have reached the broker, which then numbers publishes from 1, so the delivery-tag counter is armed regardless of whether Confirm.SelectOk arrived in time. A subsequent publish therefore still fires its on_publish and advances :attr:next_publish_seq_no. Retry confirm_delivery rather than assuming confirms were left disabled.

PARAMETER DESCRIPTION
ack_nack_callback

callback(method_frame) called for each Basic.Ack or Basic.Nack received from the broker.

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION

The Confirm.SelectOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def confirm_delivery(self,
                     ack_nack_callback,
                     timeout: float | None = DEFAULT_RPC_TIMEOUT):
    """
    Enable publisher confirms and block until Confirm.SelectOk arrives.

    Idempotent: calling this method more than once on the same
    channel - concurrently from several threads included - returns
    the original Confirm.SelectOk without sending a frame to the
    broker or resetting the delivery-tag counter, matching the
    behavior of the RabbitMQ Java and .NET clients.

    Concurrent callers are serialized on a per-channel lock, so a
    second thread waits for the first call's round trip to finish
    before it is handed the cached frame.  Its own *timeout* bounds
    only the wait for the broker's response, not that queueing, so
    it can block for the first caller's timeout on top of its own.

    The *ack_nack_callback* is dispatched on the channel's worker
    thread (same as delivery callbacks), not the IOLoop thread.
    The RabbitMQ Java and .NET clients run this listener inline on
    the I/O thread; pika's wrapper deliberately diverges so a slow
    listener cannot stall heartbeats.

    Safe to call from any thread.

    A :class:`TimeoutError` does not mean confirms are off.  The
    Confirm.Select frame is written before the wait begins and may
    well have reached the broker, which then numbers publishes from
    1, so the delivery-tag counter is armed regardless of whether
    Confirm.SelectOk arrived in time.  A subsequent publish therefore
    still fires its ``on_publish`` and advances
    :attr:`next_publish_seq_no`.  Retry ``confirm_delivery`` rather
    than assuming confirms were left disabled.

    :param ack_nack_callback:
        ``callback(method_frame)`` called for each Basic.Ack or
        Basic.Nack received from the broker.
    :param timeout: Seconds to wait for the response.
        Defaults to :data:`DEFAULT_RPC_TIMEOUT` (10 s).
        Pass ``None`` to wait indefinitely.
    :returns: The Confirm.SelectOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    with self._confirm_lock:
        if self._confirm_select_ok is not None:
            return self._confirm_select_ok

        def _wrapped_ack_nack(method_frame) -> None:
            _submit_or_terminate(
                self._consumer_work_pool, self._wrapper._connection,
                'Publisher confirm dropped: work pool shut down',
                self._safe_dispatch, 'publisher confirm callback',
                ack_nack_callback, method_frame)

        def _arm_seq_no() -> None:
            # Runs on the IOLoop thread the moment Confirm.Select is
            # written, so it is ordered ahead of every publish the IOLoop
            # has queued behind it.  Arming from the calling thread once the
            # RPC returns instead would miss each publish issued during the
            # round trip: the broker numbers those from 1, leaving the
            # counter permanently behind the broker's delivery tags.
            #
            # Arm only once.  A confirm_delivery that timed out may still
            # have reached the broker, and the retry's second Confirm.Select
            # does not reset the broker's counter, so neither may this.
            if self._next_publish_seq_no is None:
                self._next_publish_seq_no = 0

        result = self._blocking_rpc(
            'confirm_delivery',
            self._channel.confirm_delivery,
            timeout,
            on_sent=_arm_seq_no,
            ack_nack_callback=_wrapped_ack_nack,
        )
        self._confirm_select_ok = result
        return result

exchange_bind

exchange_bind(
    destination,
    source,
    routing_key: str = '',
    arguments=None,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Bind an exchange to another exchange and block until Exchange.BindOk.

Safe to call from any thread.

PARAMETER DESCRIPTION
destination

The destination exchange to bind.

source

The source exchange to bind to.

routing_key

The routing key to bind on.

TYPE: str DEFAULT: ''

arguments

Custom arguments for the binding.

DEFAULT: None

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Exchange.BindOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def exchange_bind(self,
                  destination,
                  source,
                  routing_key: str = '',
                  arguments=None,
                  timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Bind an exchange to another exchange and block until Exchange.BindOk.

    Safe to call from any thread.

    :param destination: The destination exchange to bind.
    :param source: The source exchange to bind to.
    :param routing_key: The routing key to bind on.
    :param arguments: Custom arguments for the binding.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Exchange.BindOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'exchange_bind',
        self._channel.exchange_bind,
        timeout,
        destination=destination,
        source=source,
        routing_key=routing_key,
        arguments=arguments,
    )

exchange_declare

exchange_declare(
    exchange,
    exchange_type: str = 'direct',
    passive: bool = False,
    durable: bool = False,
    auto_delete: bool = False,
    internal: bool = False,
    arguments=None,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Declare an exchange and block until Exchange.DeclareOk arrives.

Safe to call from any thread.

PARAMETER DESCRIPTION
exchange

The exchange name.

exchange_type

The exchange type (direct, fanout, topic, headers).

TYPE: str DEFAULT: 'direct'

passive

Only check if the exchange exists.

TYPE: bool DEFAULT: False

durable

Survive broker restart.

TYPE: bool DEFAULT: False

auto_delete

Delete when no queues are bound.

TYPE: bool DEFAULT: False

internal

Can only be published to by other exchanges.

TYPE: bool DEFAULT: False

arguments

Custom arguments for the exchange.

DEFAULT: None

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Exchange.DeclareOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def exchange_declare(self,
                     exchange,
                     exchange_type: str = 'direct',
                     passive: bool = False,
                     durable: bool = False,
                     auto_delete: bool = False,
                     internal: bool = False,
                     arguments=None,
                     timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Declare an exchange and block until Exchange.DeclareOk arrives.

    Safe to call from any thread.

    :param exchange: The exchange name.
    :param exchange_type: The exchange type (direct, fanout, topic, headers).
    :param passive: Only check if the exchange exists.
    :param durable: Survive broker restart.
    :param auto_delete: Delete when no queues are bound.
    :param internal: Can only be published to by other exchanges.
    :param arguments: Custom arguments for the exchange.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Exchange.DeclareOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'exchange_declare',
        self._channel.exchange_declare,
        timeout,
        exchange=exchange,
        exchange_type=exchange_type,
        passive=passive,
        durable=durable,
        auto_delete=auto_delete,
        internal=internal,
        arguments=arguments,
    )

exchange_delete

exchange_delete(
    exchange=None,
    if_unused: bool = False,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Delete an exchange and block until Exchange.DeleteOk arrives.

Safe to call from any thread.

PARAMETER DESCRIPTION
exchange

The exchange name.

DEFAULT: None

if_unused

Only delete if the exchange has no bindings.

TYPE: bool DEFAULT: False

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Exchange.DeleteOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def exchange_delete(self,
                    exchange=None,
                    if_unused: bool = False,
                    timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Delete an exchange and block until Exchange.DeleteOk arrives.

    Safe to call from any thread.

    :param exchange: The exchange name.
    :param if_unused: Only delete if the exchange has no bindings.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Exchange.DeleteOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'exchange_delete',
        self._channel.exchange_delete,
        timeout,
        exchange=exchange,
        if_unused=if_unused,
    )

exchange_unbind

exchange_unbind(
    destination,
    source,
    routing_key: str = '',
    arguments=None,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Unbind an exchange from another exchange and block until Exchange.UnbindOk.

Safe to call from any thread.

PARAMETER DESCRIPTION
destination

The destination exchange to unbind.

source

The source exchange to unbind from.

routing_key

The routing key to unbind.

TYPE: str DEFAULT: ''

arguments

Custom arguments for the unbinding.

DEFAULT: None

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Exchange.UnbindOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def exchange_unbind(self,
                    destination,
                    source,
                    routing_key: str = '',
                    arguments=None,
                    timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Unbind an exchange from another exchange and block until Exchange.UnbindOk.

    Safe to call from any thread.

    :param destination: The destination exchange to unbind.
    :param source: The source exchange to unbind from.
    :param routing_key: The routing key to unbind.
    :param arguments: Custom arguments for the unbinding.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Exchange.UnbindOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'exchange_unbind',
        self._channel.exchange_unbind,
        timeout,
        destination=destination,
        source=source,
        routing_key=routing_key,
        arguments=arguments,
    )

queue_bind

queue_bind(
    queue,
    exchange,
    routing_key=None,
    arguments=None,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Bind a queue to an exchange and block until Queue.BindOk arrives.

Safe to call from any thread.

PARAMETER DESCRIPTION
queue

The queue to bind.

exchange

The exchange to bind to.

routing_key

The routing key to bind on. Defaults to the queue name.

DEFAULT: None

arguments

Custom arguments for the binding.

DEFAULT: None

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Queue.BindOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def queue_bind(self,
               queue,
               exchange,
               routing_key=None,
               arguments=None,
               timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Bind a queue to an exchange and block until Queue.BindOk arrives.

    Safe to call from any thread.

    :param queue: The queue to bind.
    :param exchange: The exchange to bind to.
    :param routing_key: The routing key to bind on. Defaults to the queue name.
    :param arguments: Custom arguments for the binding.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Queue.BindOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'queue_bind',
        self._channel.queue_bind,
        timeout,
        queue=queue,
        exchange=exchange,
        routing_key=routing_key,
        arguments=arguments,
    )

queue_declare

queue_declare(
    queue,
    passive: bool = False,
    durable: bool = False,
    exclusive: bool = False,
    auto_delete: bool = False,
    arguments=None,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Declare a queue and block the calling thread until Queue.DeclareOk arrives.

Safe to call from any thread.

PARAMETER DESCRIPTION
queue

The queue name. If empty, the broker will generate a unique name.

passive

If True, only check whether the queue or exchange exists

TYPE: bool DEFAULT: False

durable

If True, the queue survives broker restart

TYPE: bool DEFAULT: False

exclusive

If True, restrict access to the current connection

TYPE: bool DEFAULT: False

auto_delete

If True, delete the queue or exchange when no longer in use

TYPE: bool DEFAULT: False

arguments

Custom arguments for the queue declaration

DEFAULT: None

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Queue.DeclareOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def queue_declare(self,
                  queue,
                  passive: bool = False,
                  durable: bool = False,
                  exclusive: bool = False,
                  auto_delete: bool = False,
                  arguments=None,
                  timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Declare a queue and block the calling thread until Queue.DeclareOk arrives.

    Safe to call from any thread.

    :param queue: The queue name. If empty, the broker will generate a unique name.
    :param passive: If True, only check whether the queue or exchange exists
    :param durable: If True, the queue survives broker restart
    :param exclusive: If True, restrict access to the current connection
    :param auto_delete: If True, delete the queue or exchange when no longer in use
    :param arguments: Custom arguments for the queue declaration
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Queue.DeclareOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'queue_declare',
        self._channel.queue_declare,
        timeout,
        queue=queue,
        passive=passive,
        durable=durable,
        exclusive=exclusive,
        auto_delete=auto_delete,
        arguments=arguments,
    )

queue_delete

queue_delete(
    queue,
    if_unused: bool = False,
    if_empty: bool = False,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Delete a queue and block until Queue.DeleteOk arrives.

Safe to call from any thread.

PARAMETER DESCRIPTION
queue

The queue to delete.

if_unused

Only delete if the queue has no consumers.

TYPE: bool DEFAULT: False

if_empty

Only delete if the queue is empty.

TYPE: bool DEFAULT: False

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Queue.DeleteOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def queue_delete(self,
                 queue,
                 if_unused: bool = False,
                 if_empty: bool = False,
                 timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Delete a queue and block until Queue.DeleteOk arrives.

    Safe to call from any thread.

    :param queue: The queue to delete.
    :param if_unused: Only delete if the queue has no consumers.
    :param if_empty: Only delete if the queue is empty.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Queue.DeleteOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'queue_delete',
        self._channel.queue_delete,
        timeout,
        queue=queue,
        if_unused=if_unused,
        if_empty=if_empty,
    )

queue_purge

queue_purge(
    queue, timeout: float | None = DEFAULT_RPC_TIMEOUT
) -> None

Purge all messages from a queue and block until Queue.PurgeOk arrives.

Safe to call from any thread.

PARAMETER DESCRIPTION
queue

The queue to purge.

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Queue.PurgeOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def queue_purge(self,
                queue,
                timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Purge all messages from a queue and block until Queue.PurgeOk arrives.

    Safe to call from any thread.

    :param queue: The queue to purge.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Queue.PurgeOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'queue_purge',
        self._channel.queue_purge,
        timeout,
        queue=queue,
    )

queue_unbind

queue_unbind(
    queue,
    exchange=None,
    routing_key=None,
    arguments=None,
    timeout: float | None = DEFAULT_RPC_TIMEOUT,
) -> None

Unbind a queue from an exchange and block until Queue.UnbindOk arrives.

Safe to call from any thread.

PARAMETER DESCRIPTION
queue

The queue to unbind.

exchange

The exchange to unbind from.

DEFAULT: None

routing_key

The routing key to unbind. Defaults to the queue name.

DEFAULT: None

arguments

Custom arguments for the unbinding.

DEFAULT: None

timeout

Seconds to wait for the response. Defaults to :data:DEFAULT_RPC_TIMEOUT (10 s). Pass None to wait indefinitely.

TYPE: float | None DEFAULT: DEFAULT_RPC_TIMEOUT

RETURNS DESCRIPTION
None

The Queue.UnbindOk method frame.

RAISES DESCRIPTION
Exception

if the connection is closed before the response arrives.

TimeoutError

if timeout expires before the response arrives.

Source code in pika/adapters/thread_safe_connection.py
def queue_unbind(self,
                 queue,
                 exchange=None,
                 routing_key=None,
                 arguments=None,
                 timeout: float | None = DEFAULT_RPC_TIMEOUT) -> None:
    """
    Unbind a queue from an exchange and block until Queue.UnbindOk arrives.

    Safe to call from any thread.

    :param queue: The queue to unbind.
    :param exchange: The exchange to unbind from.
    :param routing_key: The routing key to unbind. Defaults to the queue name.
    :param arguments: Custom arguments for the unbinding.
    :param timeout: Seconds to wait for the response. Defaults to :data:`DEFAULT_RPC_TIMEOUT`
        (10 s). Pass ``None`` to wait indefinitely.
    :returns: The Queue.UnbindOk method frame.
    :raises Exception: if the connection is closed before the response arrives.
    :raises TimeoutError: if *timeout* expires before the response arrives.
    """
    return self._blocking_rpc(
        'queue_unbind',
        self._channel.queue_unbind,
        timeout,
        queue=queue,
        exchange=exchange,
        routing_key=routing_key,
        arguments=arguments,
    )