Skip to content

Select Connection Adapter

Asynchronous adapter without third-party dependencies.

Interacting with Pika from another thread

pika.SelectConnection's I/O loop provides add_callback_threadsafe() to allow interacting with Pika from another thread.

Class Reference

select_connection

A connection adapter that tries to use the best polling method for the platform pika is running on.

EPollPoller

Bases: PollPoller

EPoll works on Linux and can have better performance than Poll in certain scenarios.

Both are faster than select.

Source code in pika/adapters/select_connection.py
class EPollPoller(PollPoller):
    """
    EPoll works on Linux and can have better performance than Poll in certain scenarios.

    Both are faster than select.
    """

    POLL_TIMEOUT_MULT = 1

    @staticmethod
    @override
    def _create_poller() -> Any:
        # By name: `select.epoll` is Linux-only.
        return getattr(select, 'epoll')()  # noqa: B009

IOLoop

Bases: AbstractSelectorIOLoop

I/O loop implementation that picks a suitable poller (select, poll, epoll, kqueue) to use based on platform.

Implements the pika.adapters.utils.selector_ioloop_adapter.AbstractSelectorIOLoop interface.

Source code in pika/adapters/select_connection.py
class IOLoop(AbstractSelectorIOLoop):
    """
    I/O loop implementation that picks a suitable poller (`select`, `poll`, `epoll`, `kqueue`) to
    use based on platform.

    Implements the
    `pika.adapters.utils.selector_ioloop_adapter.AbstractSelectorIOLoop`
    interface.
    """

    # READ/WRITE/ERROR per `AbstractSelectorIOLoop` requirements
    READ = PollEvents.READ
    WRITE = PollEvents.WRITE
    ERROR = PollEvents.ERROR

    def __init__(self) -> None:
        self._timer = _Timer()

        # Callbacks requested via `add_callback`
        self._callbacks: collections.deque[Callable[
            [], None]] | list[Any] = collections.deque()

        self._poller = self._get_poller(self._get_remaining_interval,
                                        self.process_timeouts)

    @override
    def close(self) -> None:
        """
        Release IOLoop's resources.

        `IOLoop.close` is intended to be called by the application or test code only after
        `IOLoop.start()` returns. After calling `close()`, no other interaction with the closed
        instance of `IOLoop` should be performed.
        """
        if self._callbacks is not None:
            self._poller.close()
            self._timer.close()
            # Set _callbacks to empty list rather than None so that race from
            # another thread calling add_callback_threadsafe() won't result in
            # AttributeError
            self._callbacks = []

    @staticmethod
    def _get_poller(get_wait_seconds: Callable[[], float | None],
                    process_timeouts: Callable[[], None]) -> _PollerBase:
        """
        Determine the best poller to use for this environment and instantiate it.

        :param get_wait_seconds: Function for getting the maximum number of
                                 seconds to wait for IO for use by the poller
        :param process_timeouts: Function for processing timeouts for use by the
                                 poller

        :returns: The instantiated poller instance supporting `_PollerBase` API
        """
        poller: _PollerBase | None = None

        kwargs: POLLER_PARAMS = {
            'get_wait_seconds': get_wait_seconds,
            'process_timeouts': process_timeouts
        }

        def requested(name: str) -> bool:
            """Whether `SELECT_TYPE` permits the named poller."""
            return not SELECT_TYPE or SELECT_TYPE == name

        if hasattr(select, 'epoll') and requested('epoll'):
            LOGGER.debug('Using EPollPoller')
            poller = EPollPoller(**kwargs)

        if not poller and hasattr(select, 'kqueue') and requested('kqueue'):
            LOGGER.debug('Using KQueuePoller')
            poller = KQueuePoller(**kwargs)

        if (not poller and hasattr(select, 'poll') and
                hasattr(select.poll(), 'modify') and requested('poll')):
            LOGGER.debug('Using PollPoller')
            poller = PollPoller(**kwargs)

        if not poller:
            LOGGER.debug('Using SelectPoller')
            poller = SelectPoller(**kwargs)

        return poller

    @override
    def call_later(self, delay: float, callback: Callable[[],
                                                          None]) -> _Timeout:
        """
        Add the callback to the IOLoop timer to be called after delay seconds from the time of call
        on best-effort basis. Returns a handle to the timeout.

        :param delay: The number of seconds to wait to call callback
        :param callback: The callback method
        :returns: handle to the created timeout that may be passed to `remove_timeout()`
        """
        return self._timer.call_later(delay, callback)

    @override
    def remove_timeout(self, timeout_handle: _Timeout) -> None:
        """
        Remove a timeout.

        :param timeout_handle: Handle of timeout to remove
        """
        self._timer.remove_timeout(timeout_handle)

    def add_callback_threadsafe(self, callback: Callable[[], None]) -> None:
        """
        Requests a call to the given function as soon as possible in the context of this IOLoop's
        thread.

        NOTE: This is the only thread-safe method in IOLoop. All other
        manipulations of IOLoop must be performed from the IOLoop's thread.

        For example, a thread may request a call to the `stop` method of an
        ioloop that is running in a different thread via
        `ioloop.add_callback_threadsafe(ioloop.stop)`

        :param callback: The callback method
        """
        if not callable(callback):
            raise TypeError(
                f'callback must be a callable, but got {callback!r}')

        # NOTE: `deque.append` is atomic
        self._callbacks.append(callback)

        # Wake up the IOLoop which may be running in another thread
        self._poller.wake_threadsafe()

        LOGGER.debug('add_callback_threadsafe: added callback=%r', callback)

    # To satisfy `AbstractSelectorIOLoop` requirement
    add_callback = add_callback_threadsafe

    def process_timeouts(self) -> None:
        """
        [Extension] Process pending callbacks and timeouts, invoking those whose time has come.

        Internal use only.
        """
        # Avoid I/O starvation by postponing new callbacks to the next iteration
        assert isinstance(self._callbacks, collections.deque)
        for _ in range(len(self._callbacks)):
            callback = self._callbacks.popleft()
            LOGGER.debug('process_timeouts: invoking callback=%r', callback)
            callback()

        self._timer.process_timeouts()

    def _get_remaining_interval(self) -> float | None:
        """
        Get the remaining interval to the next callback or timeout expiration.

        :returns: non-negative number of seconds until next callback or timer expiration; None if
            there are no callbacks and timers
        """
        if self._callbacks:
            return 0

        return self._timer.get_remaining_interval()

    @override
    def add_handler(self, fd: int, handler: Callable[[int, int], None],
                    events: int) -> None:
        """
        Start watching the given file descriptor for events.

        :param fd: The file descriptor
        :param handler: When requested event(s) occur, `handler(fd, events)` will be called.
        :param events: The event mask using READ, WRITE, ERROR.
        """
        self._poller.add_handler(fd, handler, events)

    @override
    def update_handler(self, fd: int, events: int) -> None:
        """
        Changes the events we watch for.

        :param fd: The file descriptor
        :param events: The event mask using READ, WRITE, ERROR
        """
        self._poller.update_handler(fd, events)

    @override
    def remove_handler(self, fd: int) -> None:
        """
        Stop watching the given file descriptor for events.

        :param fd: The file descriptor
        """
        self._poller.remove_handler(fd)

    @override
    def start(self) -> None:
        """
        [API] Start the main poller loop.

        It will loop until requested to exit. See `IOLoop.stop`.
        """
        self._poller.start()

    @override
    def stop(self) -> None:
        """
        [API] Request exit from the ioloop.

        The loop is NOT guaranteed to stop before this method returns.

        To invoke `stop()` safely from a thread other than this IOLoop's thread, call it via
        `add_callback_threadsafe`; e.g.,

        `ioloop.add_callback_threadsafe(ioloop.stop)`
        """
        self._poller.stop()

    def activate_poller(self) -> None:
        """[Extension] Activate the poller."""
        self._poller.activate_poller()

    def deactivate_poller(self) -> None:
        """[Extension] Deactivate the poller."""
        self._poller.deactivate_poller()

    def poll(self) -> None:
        """[Extension] Wait for events of interest on registered file
        descriptors until an event of interest occurs or next timer deadline or
        `_PollerBase._MAX_POLL_TIMEOUT`, whichever is sooner, and dispatch the
        corresponding event handlers.

        """
        self._poller.poll()

activate_poller

activate_poller() -> None

[Extension] Activate the poller.

Source code in pika/adapters/select_connection.py
def activate_poller(self) -> None:
    """[Extension] Activate the poller."""
    self._poller.activate_poller()

add_callback_threadsafe

add_callback_threadsafe(
    callback: Callable[[], None],
) -> None

Requests a call to the given function as soon as possible in the context of this IOLoop's thread.

NOTE: This is the only thread-safe method in IOLoop. All other manipulations of IOLoop must be performed from the IOLoop's thread.

For example, a thread may request a call to the stop method of an ioloop that is running in a different thread via ioloop.add_callback_threadsafe(ioloop.stop)

PARAMETER DESCRIPTION
callback

The callback method

TYPE: Callable[[], None]

Source code in pika/adapters/select_connection.py
def add_callback_threadsafe(self, callback: Callable[[], None]) -> None:
    """
    Requests a call to the given function as soon as possible in the context of this IOLoop's
    thread.

    NOTE: This is the only thread-safe method in IOLoop. All other
    manipulations of IOLoop must be performed from the IOLoop's thread.

    For example, a thread may request a call to the `stop` method of an
    ioloop that is running in a different thread via
    `ioloop.add_callback_threadsafe(ioloop.stop)`

    :param callback: The callback method
    """
    if not callable(callback):
        raise TypeError(
            f'callback must be a callable, but got {callback!r}')

    # NOTE: `deque.append` is atomic
    self._callbacks.append(callback)

    # Wake up the IOLoop which may be running in another thread
    self._poller.wake_threadsafe()

    LOGGER.debug('add_callback_threadsafe: added callback=%r', callback)

add_handler

add_handler(
    fd: int,
    handler: Callable[[int, int], None],
    events: int,
) -> None

Start watching the given file descriptor for events.

PARAMETER DESCRIPTION
fd

The file descriptor

TYPE: int

handler

When requested event(s) occur, handler(fd, events) will be called.

TYPE: Callable[[int, int], None]

events

The event mask using READ, WRITE, ERROR.

TYPE: int

Source code in pika/adapters/select_connection.py
@override
def add_handler(self, fd: int, handler: Callable[[int, int], None],
                events: int) -> None:
    """
    Start watching the given file descriptor for events.

    :param fd: The file descriptor
    :param handler: When requested event(s) occur, `handler(fd, events)` will be called.
    :param events: The event mask using READ, WRITE, ERROR.
    """
    self._poller.add_handler(fd, handler, events)

call_later

call_later(
    delay: float, callback: Callable[[], None]
) -> _Timeout

Add the callback to the IOLoop timer to be called after delay seconds from the time of call on best-effort basis. Returns a handle to the timeout.

PARAMETER DESCRIPTION
delay

The number of seconds to wait to call callback

TYPE: float

callback

The callback method

TYPE: Callable[[], None]

RETURNS DESCRIPTION
_Timeout

handle to the created timeout that may be passed to remove_timeout()

Source code in pika/adapters/select_connection.py
@override
def call_later(self, delay: float, callback: Callable[[],
                                                      None]) -> _Timeout:
    """
    Add the callback to the IOLoop timer to be called after delay seconds from the time of call
    on best-effort basis. Returns a handle to the timeout.

    :param delay: The number of seconds to wait to call callback
    :param callback: The callback method
    :returns: handle to the created timeout that may be passed to `remove_timeout()`
    """
    return self._timer.call_later(delay, callback)

close

close() -> None

Release IOLoop's resources.

IOLoop.close is intended to be called by the application or test code only after IOLoop.start() returns. After calling close(), no other interaction with the closed instance of IOLoop should be performed.

Source code in pika/adapters/select_connection.py
@override
def close(self) -> None:
    """
    Release IOLoop's resources.

    `IOLoop.close` is intended to be called by the application or test code only after
    `IOLoop.start()` returns. After calling `close()`, no other interaction with the closed
    instance of `IOLoop` should be performed.
    """
    if self._callbacks is not None:
        self._poller.close()
        self._timer.close()
        # Set _callbacks to empty list rather than None so that race from
        # another thread calling add_callback_threadsafe() won't result in
        # AttributeError
        self._callbacks = []

deactivate_poller

deactivate_poller() -> None

[Extension] Deactivate the poller.

Source code in pika/adapters/select_connection.py
def deactivate_poller(self) -> None:
    """[Extension] Deactivate the poller."""
    self._poller.deactivate_poller()

poll

poll() -> None

[Extension] Wait for events of interest on registered file descriptors until an event of interest occurs or next timer deadline or _PollerBase._MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the corresponding event handlers.

Source code in pika/adapters/select_connection.py
def poll(self) -> None:
    """[Extension] Wait for events of interest on registered file
    descriptors until an event of interest occurs or next timer deadline or
    `_PollerBase._MAX_POLL_TIMEOUT`, whichever is sooner, and dispatch the
    corresponding event handlers.

    """
    self._poller.poll()

process_timeouts

process_timeouts() -> None

[Extension] Process pending callbacks and timeouts, invoking those whose time has come.

Internal use only.

Source code in pika/adapters/select_connection.py
def process_timeouts(self) -> None:
    """
    [Extension] Process pending callbacks and timeouts, invoking those whose time has come.

    Internal use only.
    """
    # Avoid I/O starvation by postponing new callbacks to the next iteration
    assert isinstance(self._callbacks, collections.deque)
    for _ in range(len(self._callbacks)):
        callback = self._callbacks.popleft()
        LOGGER.debug('process_timeouts: invoking callback=%r', callback)
        callback()

    self._timer.process_timeouts()

remove_handler

remove_handler(fd: int) -> None

Stop watching the given file descriptor for events.

PARAMETER DESCRIPTION
fd

The file descriptor

TYPE: int

Source code in pika/adapters/select_connection.py
@override
def remove_handler(self, fd: int) -> None:
    """
    Stop watching the given file descriptor for events.

    :param fd: The file descriptor
    """
    self._poller.remove_handler(fd)

remove_timeout

remove_timeout(timeout_handle: _Timeout) -> None

Remove a timeout.

PARAMETER DESCRIPTION
timeout_handle

Handle of timeout to remove

TYPE: _Timeout

Source code in pika/adapters/select_connection.py
@override
def remove_timeout(self, timeout_handle: _Timeout) -> None:
    """
    Remove a timeout.

    :param timeout_handle: Handle of timeout to remove
    """
    self._timer.remove_timeout(timeout_handle)

start

start() -> None

[API] Start the main poller loop.

It will loop until requested to exit. See IOLoop.stop.

Source code in pika/adapters/select_connection.py
@override
def start(self) -> None:
    """
    [API] Start the main poller loop.

    It will loop until requested to exit. See `IOLoop.stop`.
    """
    self._poller.start()

stop

stop() -> None

[API] Request exit from the ioloop.

The loop is NOT guaranteed to stop before this method returns.

To invoke stop() safely from a thread other than this IOLoop's thread, call it via add_callback_threadsafe; e.g.,

ioloop.add_callback_threadsafe(ioloop.stop)

Source code in pika/adapters/select_connection.py
@override
def stop(self) -> None:
    """
    [API] Request exit from the ioloop.

    The loop is NOT guaranteed to stop before this method returns.

    To invoke `stop()` safely from a thread other than this IOLoop's thread, call it via
    `add_callback_threadsafe`; e.g.,

    `ioloop.add_callback_threadsafe(ioloop.stop)`
    """
    self._poller.stop()

update_handler

update_handler(fd: int, events: int) -> None

Changes the events we watch for.

PARAMETER DESCRIPTION
fd

The file descriptor

TYPE: int

events

The event mask using READ, WRITE, ERROR

TYPE: int

Source code in pika/adapters/select_connection.py
@override
def update_handler(self, fd: int, events: int) -> None:
    """
    Changes the events we watch for.

    :param fd: The file descriptor
    :param events: The event mask using READ, WRITE, ERROR
    """
    self._poller.update_handler(fd, events)

KQueuePoller

Bases: _PollerBase

KQueuePoller works on BSD based systems and is faster than select.

select.kqueue, select.kevent, and the KQ_* constants exist only on BSD platforms, and typeshed declares them behind a sys.platform guard to match. Every method that touches them therefore opens with the same guard, which narrows the platform for the type checker and never fires at runtime: _get_poller() only constructs this class when select.kqueue exists.

Source code in pika/adapters/select_connection.py
class KQueuePoller(_PollerBase):
    """
    KQueuePoller works on BSD based systems and is faster than select.

    `select.kqueue`, `select.kevent`, and the `KQ_*` constants exist only on BSD
    platforms, and typeshed declares them behind a `sys.platform` guard to match.
    Every method that touches them therefore opens with the same guard, which
    narrows the platform for the type checker and never fires at runtime:
    `_get_poller()` only constructs this class when `select.kqueue` exists.
    """

    def __init__(self, get_wait_seconds: Callable[[], float | None],
                 process_timeouts: Callable[[], None]) -> None:
        """Create an instance of the KQueuePoller."""
        # `Any` because the value is either `None` or a `select.kqueue` object,
        # a type the checker cannot name on platforms without that attribute.
        self._kqueue: Any = None
        super().__init__(get_wait_seconds, process_timeouts)

    @staticmethod
    def _map_event(kevent: Any) -> int:
        """
        Return the event type associated with a kevent object.

        :param kevent: a kevent object as returned by kqueue.control()
        """
        if sys.platform == 'linux' or sys.platform == 'win32':
            raise NotImplementedError('kqueue is available on BSD only.')

        mask = 0
        kq_filter_read = select.KQ_FILTER_READ
        kq_filter_write = select.KQ_FILTER_WRITE
        kq_ev_eof = select.KQ_EV_EOF
        kq_ev_error = select.KQ_EV_ERROR

        if kevent.filter == kq_filter_read:
            mask = PollEvents.READ
        elif kevent.filter == kq_filter_write:
            mask = PollEvents.WRITE
            if kevent.flags & kq_ev_eof:
                # May be set when the peer reader disconnects. We don't check
                # KQ_EV_EOF for KQ_FILTER_READ because in that case it may be
                # set before the remaining data is consumed from sockbuf.
                mask |= PollEvents.ERROR
        elif kevent.flags & kq_ev_error:
            mask = PollEvents.ERROR
        else:
            LOGGER.critical('Unexpected kevent: %s', kevent)

        return mask

    @override
    def poll(self) -> None:
        """Wait for events of interest on registered file descriptors until an event of interest
        occurs or next timer deadline or _MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the
        corresponding event handlers.
        """
        while True:
            try:
                assert self._kqueue is not None
                kevents = self._kqueue.control(None, 1000, self._get_max_wait())
                break
            except _SELECT_ERRORS as error:
                if _is_resumable(error):
                    continue
                raise

        fd_event_map: dict[int, int] = collections.defaultdict(int)
        for event in kevents:
            fd_event_map[event.ident] |= self._map_event(event)

        self._dispatch_fd_events(fd_event_map)

    @override
    def _init_poller(self) -> None:
        """Notify the implementation to allocate the poller resource."""
        assert self._kqueue is None

        if sys.platform == 'linux' or sys.platform == 'win32':
            raise NotImplementedError('kqueue is available on BSD only.')

        self._kqueue = select.kqueue()

    @override
    def _uninit_poller(self) -> None:
        """Notify the implementation to release the poller resource."""
        if self._kqueue is not None:
            self._kqueue.close()
            self._kqueue = None

    @override
    def _register_fd(self, fileno: int, events: int) -> None:
        """
        The base class invokes this method to notify the implementation to register the file
        descriptor with the polling object. The request must be ignored if the poller is not
        activated.

        :param fileno: The file descriptor
        :param events: The event mask using READ, WRITE, ERROR
        """
        self._modify_fd_events(fileno,
                               events=events,
                               events_to_clear=0,
                               events_to_set=events)

    @override
    def _modify_fd_events(self, fileno: int, events: int, events_to_clear: int,
                          events_to_set: int) -> None:
        """
        The base class invokes this method to notify the implementation to modify an already
        registered file descriptor. The request must be ignored if the poller is not activated.

        :param fileno: The file descriptor
        :param events: absolute events (READ, WRITE, ERROR)
        :param events_to_clear: The events to clear (READ, WRITE, ERROR)
        :param events_to_set: The events to set (READ, WRITE, ERROR)
        """
        if self._kqueue is None:
            return

        if sys.platform == 'linux' or sys.platform == 'win32':
            raise NotImplementedError('kqueue is available on BSD only.')

        kevents = []

        if events_to_clear & PollEvents.READ:
            kevents.append(
                select.kevent(fileno,
                              filter=select.KQ_FILTER_READ,
                              flags=select.KQ_EV_DELETE))
        if events_to_set & PollEvents.READ:
            kevents.append(
                select.kevent(fileno,
                              filter=select.KQ_FILTER_READ,
                              flags=select.KQ_EV_ADD))
        if events_to_clear & PollEvents.WRITE:
            kevents.append(
                select.kevent(fileno,
                              filter=select.KQ_FILTER_WRITE,
                              flags=select.KQ_EV_DELETE))
        if events_to_set & PollEvents.WRITE:
            kevents.append(
                select.kevent(fileno,
                              filter=select.KQ_FILTER_WRITE,
                              flags=select.KQ_EV_ADD))

        self._kqueue.control(kevents, 0)

    @override
    def _unregister_fd(self, fileno: int, events_to_clear: int) -> None:
        """
        The base class invokes this method to notify the implementation to unregister the file
        descriptor being tracked by the polling object. The request must be ignored if the poller is
        not activated.

        :param fileno: The file descriptor
        :param events_to_clear: The events to clear (READ, WRITE, ERROR)
        """
        self._modify_fd_events(fileno,
                               events=0,
                               events_to_clear=events_to_clear,
                               events_to_set=0)

__init__

__init__(
    get_wait_seconds: Callable[[], float | None],
    process_timeouts: Callable[[], None],
) -> None

Create an instance of the KQueuePoller.

Source code in pika/adapters/select_connection.py
def __init__(self, get_wait_seconds: Callable[[], float | None],
             process_timeouts: Callable[[], None]) -> None:
    """Create an instance of the KQueuePoller."""
    # `Any` because the value is either `None` or a `select.kqueue` object,
    # a type the checker cannot name on platforms without that attribute.
    self._kqueue: Any = None
    super().__init__(get_wait_seconds, process_timeouts)

poll

poll() -> None

Wait for events of interest on registered file descriptors until an event of interest occurs or next timer deadline or _MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the corresponding event handlers.

Source code in pika/adapters/select_connection.py
@override
def poll(self) -> None:
    """Wait for events of interest on registered file descriptors until an event of interest
    occurs or next timer deadline or _MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the
    corresponding event handlers.
    """
    while True:
        try:
            assert self._kqueue is not None
            kevents = self._kqueue.control(None, 1000, self._get_max_wait())
            break
        except _SELECT_ERRORS as error:
            if _is_resumable(error):
                continue
            raise

    fd_event_map: dict[int, int] = collections.defaultdict(int)
    for event in kevents:
        fd_event_map[event.ident] |= self._map_event(event)

    self._dispatch_fd_events(fd_event_map)

PollEvents

Event flags for I/O.

Source code in pika/adapters/select_connection.py
class PollEvents:
    """Event flags for I/O."""

    # Use epoll's constants to keep life easy
    READ = getattr(select, 'POLLIN', 0x01)  # available for read
    WRITE = getattr(select, 'POLLOUT', 0x04)  # available for write
    ERROR = getattr(select, 'POLLERR', 0x08)  # error on associated fd
    HANGUP = getattr(select, 'POLLHUP', 0x10)  # hangup on associated fd

PollPoller

Bases: _PollerBase

Poll works on Linux and can have better performance than EPoll in certain scenarios.

Both are faster than select.

Source code in pika/adapters/select_connection.py
class PollPoller(_PollerBase):
    """
    Poll works on Linux and can have better performance than EPoll in certain scenarios.

    Both are faster than select.
    """

    POLL_TIMEOUT_MULT = 1000

    def __init__(self, get_wait_seconds: Callable[[], float | None],
                 process_timeouts: Callable[[], None]) -> None:
        """Create an instance of the PollPoller."""
        self._poll: Any = None
        super().__init__(get_wait_seconds, process_timeouts)

    @staticmethod
    def _create_poller() -> Any:
        # By name: `select.poll` is absent on Windows.
        return getattr(select, 'poll')()  # noqa: B009

    @override
    def poll(self) -> None:
        """Wait for events of interest on registered file descriptors until an event of interest
        occurs or next timer deadline or _MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the
        corresponding event handlers.
        """
        while True:
            try:
                events = self._poll.poll(self._get_max_wait())
                break
            except _SELECT_ERRORS as error:
                if _is_resumable(error):
                    continue
                raise

        fd_event_map: dict[int, int] = collections.defaultdict(int)
        for fileno, event in events:
            # NOTE: On OS X, when poll() sets POLLHUP, it's mutually-exclusive with
            # POLLOUT and it doesn't seem to set POLLERR along with POLLHUP when
            # socket connection fails, for example. So, we need to at least add
            # POLLERR when we see POLLHUP
            if (event & PollEvents.HANGUP) and pika._utils.ON_OSX:
                event |= PollEvents.ERROR

            fd_event_map[fileno] |= event

        self._dispatch_fd_events(fd_event_map)

    @override
    def _init_poller(self) -> None:
        """Notify the implementation to allocate the poller resource."""
        assert self._poll is None

        self._poll = self._create_poller()

    @override
    def _uninit_poller(self) -> None:
        """Notify the implementation to release the poller resource."""
        if self._poll is not None:
            if hasattr(self._poll, 'close'):
                self._poll.close(
                )  # pyright: ignore[reportAttributeAccessIssue]

            self._poll = None

    @override
    def _register_fd(self, fileno: int, events: int) -> None:
        """
        The base class invokes this method to notify the implementation to register the file
        descriptor with the polling object. The request must be ignored if the poller is not
        activated.

        :param fileno: The file descriptor
        :param events: The event mask using READ, WRITE, ERROR
        """
        if self._poll is not None:
            self._poll.register(fileno, events)

    @override
    def _modify_fd_events(self, fileno: int, events: int, events_to_clear: int,
                          events_to_set: int) -> None:
        """
        The base class invokes this method to notify the implementation to modify an already
        registered file descriptor. The request must be ignored if the poller is not activated.

        :param fileno: The file descriptor
        :param events: absolute events (READ, WRITE, ERROR)
        :param events_to_clear: The events to clear (READ, WRITE, ERROR)
        :param events_to_set: The events to set (READ, WRITE, ERROR)
        """
        if self._poll is not None:
            self._poll.modify(fileno, events)

    @override
    def _unregister_fd(self, fileno: int, events_to_clear: int) -> None:
        """
        The base class invokes this method to notify the implementation to unregister the file
        descriptor being tracked by the polling object. The request must be ignored if the poller is
        not activated.

        :param fileno: The file descriptor
        :param events_to_clear: The events to clear (READ, WRITE, ERROR)
        """
        if self._poll is not None:
            self._poll.unregister(fileno)

__init__

__init__(
    get_wait_seconds: Callable[[], float | None],
    process_timeouts: Callable[[], None],
) -> None

Create an instance of the PollPoller.

Source code in pika/adapters/select_connection.py
def __init__(self, get_wait_seconds: Callable[[], float | None],
             process_timeouts: Callable[[], None]) -> None:
    """Create an instance of the PollPoller."""
    self._poll: Any = None
    super().__init__(get_wait_seconds, process_timeouts)

poll

poll() -> None

Wait for events of interest on registered file descriptors until an event of interest occurs or next timer deadline or _MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the corresponding event handlers.

Source code in pika/adapters/select_connection.py
@override
def poll(self) -> None:
    """Wait for events of interest on registered file descriptors until an event of interest
    occurs or next timer deadline or _MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the
    corresponding event handlers.
    """
    while True:
        try:
            events = self._poll.poll(self._get_max_wait())
            break
        except _SELECT_ERRORS as error:
            if _is_resumable(error):
                continue
            raise

    fd_event_map: dict[int, int] = collections.defaultdict(int)
    for fileno, event in events:
        # NOTE: On OS X, when poll() sets POLLHUP, it's mutually-exclusive with
        # POLLOUT and it doesn't seem to set POLLERR along with POLLHUP when
        # socket connection fails, for example. So, we need to at least add
        # POLLERR when we see POLLHUP
        if (event & PollEvents.HANGUP) and pika._utils.ON_OSX:
            event |= PollEvents.ERROR

        fd_event_map[fileno] |= event

    self._dispatch_fd_events(fd_event_map)

SelectConnection

Bases: BaseConnection

An asynchronous connection adapter that attempts to use the fastest event loop adapter for the given platform.

Source code in pika/adapters/select_connection.py
class SelectConnection(BaseConnection):
    """An asynchronous connection adapter that attempts to use the fastest event loop adapter for
    the given platform.
    """

    def __init__(
            self,
            parameters: connection.Parameters | None = None,
            on_open_callback: Callable[[connection.Connection], None] |
        None = None,
            on_open_error_callback: Callable[
                [connection.Connection, BaseException], None] | None = None,
            on_close_callback: Callable[[connection.Connection, BaseException],
                                        None] | None = None,
            custom_ioloop: nbio_interface.AbstractIOServices | IOLoop |
        None = None,
            internal_connection_workflow: bool = True) -> None:
        """
        Create a new instance of the Connection object.

        :param parameters: Connection parameters
        :param on_open_callback: Method to call on connection open
        :param on_open_error_callback: Callback (or None) with signature
            ``(Connection, BaseException) -> None``; called if the connection
            can't be established or connection establishment is interrupted by
            `Connection.close()`: on_open_error_callback(Connection, exception).
        :param on_close_callback: Callback (or None) with signature
            ``(Connection, BaseException) -> None``; called when a previously fully
            open connection is closed:
            `on_close_callback(Connection, exception)`, where `exception` is
            either an instance of `exceptions.ConnectionClosed` if closed by
            user or broker or exception of another type that describes the cause
            of connection failure.
        :param custom_ioloop: Provide a custom I/O Loop object.
        :param internal_connection_workflow: True for autonomous connection
            establishment which is default; False for externally-managed
            connection workflow via the `create_connection()` factory.
        :raises RuntimeError:
        """
        if isinstance(custom_ioloop, nbio_interface.AbstractIOServices):
            nbio = custom_ioloop
        else:
            nbio = SelectorIOServicesAdapter(custom_ioloop or IOLoop())

        super().__init__(
            parameters,
            on_open_callback,
            on_open_error_callback,
            on_close_callback,
            nbio,
            internal_connection_workflow=internal_connection_workflow)

    @classmethod
    @override
    def create_connection(
        cls,
        connection_configs: Sequence[connection.Parameters],
        on_done: Callable[[(connection.Connection |
                            connection_workflow.AMQPConnectorException)], None],
        custom_ioloop: Any | None = None,
        workflow: connection_workflow.AbstractAMQPConnectionWorkflow |
        None = None
    ) -> connection_workflow.AbstractAMQPConnectionWorkflow:
        """
        :param connection_configs: One or more connection parameter objects
        :param on_done: Callback to report when connection workflow is done
        :param custom_ioloop: Optional custom IOLoop or nbio interface to use
            for the connection workflow
        :param workflow: Optional connection workflow instance to use; if None,
            a default workflow will be created
        """
        nbio = SelectorIOServicesAdapter(custom_ioloop or IOLoop())

        def connection_factory(params) -> SelectConnection:
            """
            Connection factory.

            :param params: Connection parameters
            """
            if params is None:
                raise ValueError('Expected pika.connection.Parameters '
                                 'instance, but got None in params arg.')
            return cls(parameters=params,
                       custom_ioloop=nbio,
                       internal_connection_workflow=False)

        return cls._start_connection_workflow(
            connection_configs=connection_configs,
            connection_factory=connection_factory,
            nbio=nbio,
            workflow=workflow,
            on_done=on_done)

    def _get_write_buffer_size(self) -> int:
        """
        :returns: Current size of output data buffered by the transport
        """
        assert self._transport is not None
        return self._transport.get_write_buffer_size()

__init__

__init__(
    parameters: Parameters | None = None,
    on_open_callback: Callable[[Connection], None]
    | None = None,
    on_open_error_callback: Callable[
        [Connection, BaseException], None
    ]
    | None = None,
    on_close_callback: Callable[
        [Connection, BaseException], None
    ]
    | None = None,
    custom_ioloop: AbstractIOServices
    | IOLoop
    | None = None,
    internal_connection_workflow: bool = True,
) -> None

Create a new instance of the Connection object.

PARAMETER DESCRIPTION
parameters

Connection parameters

TYPE: Parameters | None DEFAULT: None

on_open_callback

Method to call on connection open

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

on_open_error_callback

Callback (or None) with signature (Connection, BaseException) -> None; called if the connection can't be established or connection establishment is interrupted by Connection.close(): on_open_error_callback(Connection, exception).

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

on_close_callback

Callback (or None) with signature (Connection, BaseException) -> None; called when a previously fully open connection is closed: on_close_callback(Connection, exception), where exception is either an instance of exceptions.ConnectionClosed if closed by user or broker or exception of another type that describes the cause of connection failure.

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

custom_ioloop

Provide a custom I/O Loop object.

TYPE: AbstractIOServices | IOLoop | None DEFAULT: None

internal_connection_workflow

True for autonomous connection establishment which is default; False for externally-managed connection workflow via the create_connection() factory.

TYPE: bool DEFAULT: True

RAISES DESCRIPTION
RuntimeError
Source code in pika/adapters/select_connection.py
def __init__(
        self,
        parameters: connection.Parameters | None = None,
        on_open_callback: Callable[[connection.Connection], None] |
    None = None,
        on_open_error_callback: Callable[
            [connection.Connection, BaseException], None] | None = None,
        on_close_callback: Callable[[connection.Connection, BaseException],
                                    None] | None = None,
        custom_ioloop: nbio_interface.AbstractIOServices | IOLoop |
    None = None,
        internal_connection_workflow: bool = True) -> None:
    """
    Create a new instance of the Connection object.

    :param parameters: Connection parameters
    :param on_open_callback: Method to call on connection open
    :param on_open_error_callback: Callback (or None) with signature
        ``(Connection, BaseException) -> None``; called if the connection
        can't be established or connection establishment is interrupted by
        `Connection.close()`: on_open_error_callback(Connection, exception).
    :param on_close_callback: Callback (or None) with signature
        ``(Connection, BaseException) -> None``; called when a previously fully
        open connection is closed:
        `on_close_callback(Connection, exception)`, where `exception` is
        either an instance of `exceptions.ConnectionClosed` if closed by
        user or broker or exception of another type that describes the cause
        of connection failure.
    :param custom_ioloop: Provide a custom I/O Loop object.
    :param internal_connection_workflow: True for autonomous connection
        establishment which is default; False for externally-managed
        connection workflow via the `create_connection()` factory.
    :raises RuntimeError:
    """
    if isinstance(custom_ioloop, nbio_interface.AbstractIOServices):
        nbio = custom_ioloop
    else:
        nbio = SelectorIOServicesAdapter(custom_ioloop or IOLoop())

    super().__init__(
        parameters,
        on_open_callback,
        on_open_error_callback,
        on_close_callback,
        nbio,
        internal_connection_workflow=internal_connection_workflow)

create_connection classmethod

create_connection(
    connection_configs: Sequence[Parameters],
    on_done: Callable[
        [Connection | AMQPConnectorException], None
    ],
    custom_ioloop: Any | None = None,
    workflow: AbstractAMQPConnectionWorkflow | None = None,
) -> AbstractAMQPConnectionWorkflow
PARAMETER DESCRIPTION
connection_configs

One or more connection parameter objects

TYPE: Sequence[Parameters]

on_done

Callback to report when connection workflow is done

TYPE: Callable[[Connection | AMQPConnectorException], None]

custom_ioloop

Optional custom IOLoop or nbio interface to use for the connection workflow

TYPE: Any | None DEFAULT: None

workflow

Optional connection workflow instance to use; if None, a default workflow will be created

TYPE: AbstractAMQPConnectionWorkflow | None DEFAULT: None

Source code in pika/adapters/select_connection.py
@classmethod
@override
def create_connection(
    cls,
    connection_configs: Sequence[connection.Parameters],
    on_done: Callable[[(connection.Connection |
                        connection_workflow.AMQPConnectorException)], None],
    custom_ioloop: Any | None = None,
    workflow: connection_workflow.AbstractAMQPConnectionWorkflow |
    None = None
) -> connection_workflow.AbstractAMQPConnectionWorkflow:
    """
    :param connection_configs: One or more connection parameter objects
    :param on_done: Callback to report when connection workflow is done
    :param custom_ioloop: Optional custom IOLoop or nbio interface to use
        for the connection workflow
    :param workflow: Optional connection workflow instance to use; if None,
        a default workflow will be created
    """
    nbio = SelectorIOServicesAdapter(custom_ioloop or IOLoop())

    def connection_factory(params) -> SelectConnection:
        """
        Connection factory.

        :param params: Connection parameters
        """
        if params is None:
            raise ValueError('Expected pika.connection.Parameters '
                             'instance, but got None in params arg.')
        return cls(parameters=params,
                   custom_ioloop=nbio,
                   internal_connection_workflow=False)

    return cls._start_connection_workflow(
        connection_configs=connection_configs,
        connection_factory=connection_factory,
        nbio=nbio,
        workflow=workflow,
        on_done=on_done)

SelectPoller

Bases: _PollerBase

Default behavior is to use Select since it's the widest supported and has all of the methods we need for child classes as well.

One should only need to override the update_handler and start methods for additional types.

Source code in pika/adapters/select_connection.py
class SelectPoller(_PollerBase):
    """
    Default behavior is to use Select since it's the widest supported and has all of the methods we
    need for child classes as well.

    One should only need to override the update_handler and start methods for additional types.
    """

    # if the poller uses MS specify 1000
    POLL_TIMEOUT_MULT = 1

    @override
    def poll(self) -> None:
        """Wait for events of interest on registered file descriptors until an event of interest
        occurs or next timer deadline or _MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the
        corresponding event handlers.
        """
        assert self._fd_events is not None

        while True:
            try:
                if (self._fd_events[PollEvents.READ] or
                        self._fd_events[PollEvents.WRITE] or
                        self._fd_events[PollEvents.ERROR]):
                    read, write, error = select.select(
                        self._fd_events[PollEvents.READ],
                        self._fd_events[PollEvents.WRITE],
                        self._fd_events[PollEvents.ERROR], self._get_max_wait())
                else:
                    # NOTE When called without any FDs, select fails on
                    # Windows with error 10022, 'An invalid argument was
                    # supplied'.
                    time.sleep(self._get_max_wait())
                    read, write, error = [], [], []
                break
            except _SELECT_ERRORS as error:
                if _is_resumable(error):
                    continue
                raise

        # Build an event bit mask for each fileno we've received an event for
        fd_event_map: dict[int, int] = collections.defaultdict(int)
        for fd_set, evt in zip(
            (read, write, error),
            (PollEvents.READ, PollEvents.WRITE, PollEvents.ERROR)):
            for fileno in fd_set:
                fd_event_map[fileno] |= evt

        self._dispatch_fd_events(fd_event_map)

    @override
    def _init_poller(self) -> None:
        """Notify the implementation to allocate the poller resource."""

        # It's a no op in SelectPoller

    @override
    def _uninit_poller(self) -> None:
        """Notify the implementation to release the poller resource."""

        # It's a no op in SelectPoller

    @override
    def _register_fd(self, fileno: int, events: int) -> None:
        """
        The base class invokes this method to notify the implementation to register the file
        descriptor with the polling object. The request must be ignored if the poller is not
        activated.

        :param fileno: The file descriptor
        :param events: The event mask using READ, WRITE, ERROR
        """

        # It's a no op in SelectPoller

    @override
    def _modify_fd_events(self, fileno: int, events: int, events_to_clear: int,
                          events_to_set: int) -> None:
        """
        The base class invokes this method to notify the implementation to modify an already
        registered file descriptor. The request must be ignored if the poller is not activated.

        :param fileno: The file descriptor
        :param events: absolute events (READ, WRITE, ERROR)
        :param events_to_clear: The events to clear (READ, WRITE, ERROR)
        :param events_to_set: The events to set (READ, WRITE, ERROR)
        """

        # It's a no op in SelectPoller

    @override
    def _unregister_fd(self, fileno: int, events_to_clear: int) -> None:
        """
        The base class invokes this method to notify the implementation to unregister the file
        descriptor being tracked by the polling object. The request must be ignored if the poller is
        not activated.

        This is a no-op in SelectPoller.

        :param fileno: The file descriptor
        :param events_to_clear: The events to clear (READ, WRITE, ERROR)
        """

poll

poll() -> None

Wait for events of interest on registered file descriptors until an event of interest occurs or next timer deadline or _MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the corresponding event handlers.

Source code in pika/adapters/select_connection.py
@override
def poll(self) -> None:
    """Wait for events of interest on registered file descriptors until an event of interest
    occurs or next timer deadline or _MAX_POLL_TIMEOUT, whichever is sooner, and dispatch the
    corresponding event handlers.
    """
    assert self._fd_events is not None

    while True:
        try:
            if (self._fd_events[PollEvents.READ] or
                    self._fd_events[PollEvents.WRITE] or
                    self._fd_events[PollEvents.ERROR]):
                read, write, error = select.select(
                    self._fd_events[PollEvents.READ],
                    self._fd_events[PollEvents.WRITE],
                    self._fd_events[PollEvents.ERROR], self._get_max_wait())
            else:
                # NOTE When called without any FDs, select fails on
                # Windows with error 10022, 'An invalid argument was
                # supplied'.
                time.sleep(self._get_max_wait())
                read, write, error = [], [], []
            break
        except _SELECT_ERRORS as error:
            if _is_resumable(error):
                continue
            raise

    # Build an event bit mask for each fileno we've received an event for
    fd_event_map: dict[int, int] = collections.defaultdict(int)
    for fd_set, evt in zip(
        (read, write, error),
        (PollEvents.READ, PollEvents.WRITE, PollEvents.ERROR)):
        for fileno in fd_set:
            fd_event_map[fileno] |= evt

    self._dispatch_fd_events(fd_event_map)

SelectConnection

Bases: BaseConnection

An asynchronous connection adapter that attempts to use the fastest event loop adapter for the given platform.

Source code in pika/adapters/select_connection.py
class SelectConnection(BaseConnection):
    """An asynchronous connection adapter that attempts to use the fastest event loop adapter for
    the given platform.
    """

    def __init__(
            self,
            parameters: connection.Parameters | None = None,
            on_open_callback: Callable[[connection.Connection], None] |
        None = None,
            on_open_error_callback: Callable[
                [connection.Connection, BaseException], None] | None = None,
            on_close_callback: Callable[[connection.Connection, BaseException],
                                        None] | None = None,
            custom_ioloop: nbio_interface.AbstractIOServices | IOLoop |
        None = None,
            internal_connection_workflow: bool = True) -> None:
        """
        Create a new instance of the Connection object.

        :param parameters: Connection parameters
        :param on_open_callback: Method to call on connection open
        :param on_open_error_callback: Callback (or None) with signature
            ``(Connection, BaseException) -> None``; called if the connection
            can't be established or connection establishment is interrupted by
            `Connection.close()`: on_open_error_callback(Connection, exception).
        :param on_close_callback: Callback (or None) with signature
            ``(Connection, BaseException) -> None``; called when a previously fully
            open connection is closed:
            `on_close_callback(Connection, exception)`, where `exception` is
            either an instance of `exceptions.ConnectionClosed` if closed by
            user or broker or exception of another type that describes the cause
            of connection failure.
        :param custom_ioloop: Provide a custom I/O Loop object.
        :param internal_connection_workflow: True for autonomous connection
            establishment which is default; False for externally-managed
            connection workflow via the `create_connection()` factory.
        :raises RuntimeError:
        """
        if isinstance(custom_ioloop, nbio_interface.AbstractIOServices):
            nbio = custom_ioloop
        else:
            nbio = SelectorIOServicesAdapter(custom_ioloop or IOLoop())

        super().__init__(
            parameters,
            on_open_callback,
            on_open_error_callback,
            on_close_callback,
            nbio,
            internal_connection_workflow=internal_connection_workflow)

    @classmethod
    @override
    def create_connection(
        cls,
        connection_configs: Sequence[connection.Parameters],
        on_done: Callable[[(connection.Connection |
                            connection_workflow.AMQPConnectorException)], None],
        custom_ioloop: Any | None = None,
        workflow: connection_workflow.AbstractAMQPConnectionWorkflow |
        None = None
    ) -> connection_workflow.AbstractAMQPConnectionWorkflow:
        """
        :param connection_configs: One or more connection parameter objects
        :param on_done: Callback to report when connection workflow is done
        :param custom_ioloop: Optional custom IOLoop or nbio interface to use
            for the connection workflow
        :param workflow: Optional connection workflow instance to use; if None,
            a default workflow will be created
        """
        nbio = SelectorIOServicesAdapter(custom_ioloop or IOLoop())

        def connection_factory(params) -> SelectConnection:
            """
            Connection factory.

            :param params: Connection parameters
            """
            if params is None:
                raise ValueError('Expected pika.connection.Parameters '
                                 'instance, but got None in params arg.')
            return cls(parameters=params,
                       custom_ioloop=nbio,
                       internal_connection_workflow=False)

        return cls._start_connection_workflow(
            connection_configs=connection_configs,
            connection_factory=connection_factory,
            nbio=nbio,
            workflow=workflow,
            on_done=on_done)

    def _get_write_buffer_size(self) -> int:
        """
        :returns: Current size of output data buffered by the transport
        """
        assert self._transport is not None
        return self._transport.get_write_buffer_size()

CONNECTION_CLOSED class-attribute instance-attribute

CONNECTION_CLOSED = 0

CONNECTION_CLOSING class-attribute instance-attribute

CONNECTION_CLOSING = 6

CONNECTION_INIT class-attribute instance-attribute

CONNECTION_INIT = 1

CONNECTION_OPEN class-attribute instance-attribute

CONNECTION_OPEN = 5

CONNECTION_PROTOCOL class-attribute instance-attribute

CONNECTION_PROTOCOL = 2

CONNECTION_START class-attribute instance-attribute

CONNECTION_START = 3

CONNECTION_TUNE class-attribute instance-attribute

CONNECTION_TUNE = 4

ON_CONNECTION_CLOSED class-attribute instance-attribute

ON_CONNECTION_CLOSED = '_on_connection_closed'

ON_CONNECTION_ERROR class-attribute instance-attribute

ON_CONNECTION_ERROR = '_on_connection_error'

ON_CONNECTION_OPEN_OK class-attribute instance-attribute

ON_CONNECTION_OPEN_OK = '_on_connection_open_ok'

_STATE_NAMES class-attribute

_STATE_NAMES: dict[int, str] = {
    CONNECTION_CLOSED: 'CLOSED',
    CONNECTION_INIT: 'INIT',
    CONNECTION_PROTOCOL: 'PROTOCOL',
    CONNECTION_START: 'START',
    CONNECTION_TUNE: 'TUNE',
    CONNECTION_OPEN: 'OPEN',
    CONNECTION_CLOSING: 'CLOSING',
}

_blocked_conn_timer instance-attribute

_blocked_conn_timer: object = None

_body_max_length instance-attribute

_body_max_length: int = (
    spec.FRAME_MAX_SIZE
    - spec.FRAME_HEADER_SIZE
    - spec.FRAME_END_SIZE
)

_channels instance-attribute

_channels: dict[int, Channel] = {}

_client_properties property

_client_properties: dict[str, Any]

Return the client properties dictionary.

_connection_workflow instance-attribute

_connection_workflow: (
    AbstractAMQPConnectionWorkflow | None
) = None

_error instance-attribute

_error: Exception | None = None

_frame_buffer instance-attribute

_frame_buffer: bytearray = bytearray()

_got_eof instance-attribute

_got_eof: bool = False

_heartbeat_checker instance-attribute

_heartbeat_checker: HeartbeatChecker | None = None

_internal_connection_workflow instance-attribute

_internal_connection_workflow = internal_connection_workflow

_nbio instance-attribute

_nbio: AbstractIOServices = nbio

_opened instance-attribute

_opened: bool = False

_processing_frame_buffer instance-attribute

_processing_frame_buffer: bool = False

_transport instance-attribute

_transport: AbstractStreamTransport | None = None

accept_consumer_cancel_ok property

accept_consumer_cancel_ok: bool

Specifies if the server accepts a client-sent basic.cancel-ok in reply to a broker-sent basic.cancel on the active connection.

basic_nack property

basic_nack: bool

Specifies if the server supports basic.nack on the active connection.

callbacks instance-attribute

callbacks = pika.callback.CallbackManager()

connection_state instance-attribute

connection_state: int = self.CONNECTION_CLOSED

consumer_cancel_notify property

consumer_cancel_notify: bool

Specifies if the server supports consumer cancel notification on the active connection.

exchange_exchange_bindings property

exchange_exchange_bindings: bool

Specifies if the active connection supports exchange to exchange bindings.

ioloop property

ioloop: Any
RETURNS DESCRIPTION
Any

the native I/O loop instance underlying async services selected by user or the default selected by the specialized connection adapter (e.g., Twisted reactor, asyncio.SelectorEventLoop, select_connection.IOLoop, etc.)

is_closed property

is_closed: bool

Returns a boolean reporting the current connection state.

is_closing property

is_closing: bool

Returns True if connection is in the process of closing due to client-initiated close request, but closing is not yet complete.

is_open property

is_open: bool

Returns a boolean reporting the current connection state.

known_hosts instance-attribute

known_hosts: str | bytes | None = None

params instance-attribute

params = copy.deepcopy(parameters)

publisher_confirms property

publisher_confirms: bool

Specifies if the active connection can use publisher confirmations.

server_capabilities instance-attribute

server_capabilities: dict[str, bool] | None = None

server_properties instance-attribute

server_properties: dict[str, Any] | None = None

__init__

__init__(
    parameters: Parameters | None = None,
    on_open_callback: Callable[[Connection], None]
    | None = None,
    on_open_error_callback: Callable[
        [Connection, BaseException], None
    ]
    | None = None,
    on_close_callback: Callable[
        [Connection, BaseException], None
    ]
    | None = None,
    custom_ioloop: AbstractIOServices
    | IOLoop
    | None = None,
    internal_connection_workflow: bool = True,
) -> None

Create a new instance of the Connection object.

PARAMETER DESCRIPTION
parameters

Connection parameters

TYPE: Parameters | None DEFAULT: None

on_open_callback

Method to call on connection open

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

on_open_error_callback

Callback (or None) with signature (Connection, BaseException) -> None; called if the connection can't be established or connection establishment is interrupted by Connection.close(): on_open_error_callback(Connection, exception).

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

on_close_callback

Callback (or None) with signature (Connection, BaseException) -> None; called when a previously fully open connection is closed: on_close_callback(Connection, exception), where exception is either an instance of exceptions.ConnectionClosed if closed by user or broker or exception of another type that describes the cause of connection failure.

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

custom_ioloop

Provide a custom I/O Loop object.

TYPE: AbstractIOServices | IOLoop | None DEFAULT: None

internal_connection_workflow

True for autonomous connection establishment which is default; False for externally-managed connection workflow via the create_connection() factory.

TYPE: bool DEFAULT: True

RAISES DESCRIPTION
RuntimeError
Source code in pika/adapters/select_connection.py
def __init__(
        self,
        parameters: connection.Parameters | None = None,
        on_open_callback: Callable[[connection.Connection], None] |
    None = None,
        on_open_error_callback: Callable[
            [connection.Connection, BaseException], None] | None = None,
        on_close_callback: Callable[[connection.Connection, BaseException],
                                    None] | None = None,
        custom_ioloop: nbio_interface.AbstractIOServices | IOLoop |
    None = None,
        internal_connection_workflow: bool = True) -> None:
    """
    Create a new instance of the Connection object.

    :param parameters: Connection parameters
    :param on_open_callback: Method to call on connection open
    :param on_open_error_callback: Callback (or None) with signature
        ``(Connection, BaseException) -> None``; called if the connection
        can't be established or connection establishment is interrupted by
        `Connection.close()`: on_open_error_callback(Connection, exception).
    :param on_close_callback: Callback (or None) with signature
        ``(Connection, BaseException) -> None``; called when a previously fully
        open connection is closed:
        `on_close_callback(Connection, exception)`, where `exception` is
        either an instance of `exceptions.ConnectionClosed` if closed by
        user or broker or exception of another type that describes the cause
        of connection failure.
    :param custom_ioloop: Provide a custom I/O Loop object.
    :param internal_connection_workflow: True for autonomous connection
        establishment which is default; False for externally-managed
        connection workflow via the `create_connection()` factory.
    :raises RuntimeError:
    """
    if isinstance(custom_ioloop, nbio_interface.AbstractIOServices):
        nbio = custom_ioloop
    else:
        nbio = SelectorIOServicesAdapter(custom_ioloop or IOLoop())

    super().__init__(
        parameters,
        on_open_callback,
        on_open_error_callback,
        on_close_callback,
        nbio,
        internal_connection_workflow=internal_connection_workflow)

__repr__

__repr__() -> str
Source code in pika/adapters/base_connection.py
@override
def __repr__(self) -> str:
    return (
        f'<{self.__class__.__name__} {self._STATE_NAMES[self.connection_state]} transport={self._transport} params={self.params}>'
    )

_abort_connection_workflow

_abort_connection_workflow() -> None

Asynchronously abort connection workflow.

Upon completion, Connection._on_stream_terminated() will be called with None as the error argument.

Assumption: may be called only while connection is opening.

Source code in pika/adapters/base_connection.py
def _abort_connection_workflow(self) -> None:
    """
    Asynchronously abort connection workflow.

    Upon completion, `Connection._on_stream_terminated()` will be called with None
    as the error argument.

    Assumption: may be called only while connection is opening.
    """
    assert not self._opened, (
        '_abort_connection_workflow() may be called only when '
        'connection is opening.')

    if self._transport is None:
        # NOTE: this is possible only when user calls Connection.close() to
        # interrupt internally-initiated connection establishment.
        # self._connection_workflow.abort() would not call
        # Connection.close() before pairing of connection with transport.
        assert self._internal_connection_workflow, (
            'Unexpected _abort_connection_workflow() call with '
            'no transport in external connection workflow mode.')

        # This will result in call to _on_connection_workflow_done() upon
        # completion
        assert self._connection_workflow is not None
        self._connection_workflow.abort()
    else:
        # NOTE: we can't use self._connection_workflow.abort() in this case,
        # because it would result in infinite recursion as we're called
        # from Connection.close() and _connection_workflow.abort() calls
        # Connection.close() to abort a connection that's already been
        # paired with a transport. During internally-initiated connection
        # establishment, AMQPConnectionWorkflow will discover that user
        # aborted the connection when it receives
        # pika.exceptions.ConnectionOpenAborted.

        # This completes asynchronously, culminating in call to our method
        # `connection_lost()`
        self._transport.abort()

_adapter_add_callback_threadsafe

_adapter_add_callback_threadsafe(
    callback: Callable[..., Any],
) -> None

Add a callback to be called from the I/O loop thread.

PARAMETER DESCRIPTION
callback

Callback to call.

TYPE: Callable[..., Any]

Source code in pika/adapters/base_connection.py
@override
def _adapter_add_callback_threadsafe(self, callback: Callable[...,
                                                              Any]) -> None:
    """
    Add a callback to be called from the I/O loop thread.

    :param callback: Callback to call.
    """
    if not callable(callback):
        raise TypeError(
            f'callback must be a callable, but got {callback!r}')

    self._nbio.add_callback_threadsafe(callback)

_adapter_call_later

_adapter_call_later(
    delay: float, callback: Callable[[], None]
) -> object

Schedule a callback to be called after a delay.

PARAMETER DESCRIPTION
delay

Delay in seconds.

TYPE: float

callback

Callback to call.

TYPE: Callable[[], None]

RETURNS DESCRIPTION
object

Timeout handle that can be used to cancel.

Source code in pika/adapters/base_connection.py
@override
def _adapter_call_later(self, delay: float,
                        callback: Callable[[], None]) -> object:
    """
    Schedule a callback to be called after a delay.

    :param delay: Delay in seconds.
    :param callback: Callback to call.
    :returns: Timeout handle that can be used to cancel.
    """
    return self._nbio.call_later(delay, callback)

_adapter_connect_stream

_adapter_connect_stream() -> None

Initiate full-stack connection establishment asynchronously for internally-initiated connection bring-up.

Upon failed completion, we will invoke Connection._on_stream_terminated(). NOTE: On success, the stack will be up already, so there is no corresponding callback.

Source code in pika/adapters/base_connection.py
@override
def _adapter_connect_stream(self) -> None:
    """
    Initiate full-stack connection establishment asynchronously for internally-initiated
    connection bring-up.

    Upon failed completion, we will invoke `Connection._on_stream_terminated()`. NOTE: On
    success, the stack will be up already, so there is no corresponding callback.
    """
    self._connection_workflow = connection_workflow.AMQPConnectionWorkflow(
        _until_first_amqp_attempt=True)

    self._connection_workflow.set_io_services(self._nbio)

    def create_connector() -> connection_workflow.AMQPConnector:
        """`AMQPConnector` factory."""
        return connection_workflow.AMQPConnector(
            lambda _params: _StreamingProtocolShim(self), self._nbio)

    self._connection_workflow.start(
        [self.params],
        connector_factory=create_connector,
        native_loop=self._nbio.get_native_ioloop(),
        on_done=functools.partial(
            self._unshim_connection_workflow_callback,
            cast(Any, self._on_connection_workflow_done)))

_adapter_disconnect_stream

_adapter_disconnect_stream() -> None

Asynchronously bring down the streaming transport layer and invoke Connection._on_stream_terminated() asynchronously when complete.

Source code in pika/adapters/base_connection.py
@override
def _adapter_disconnect_stream(self) -> None:
    """Asynchronously bring down the streaming transport layer and invoke
    `Connection._on_stream_terminated()` asynchronously when complete.

    """
    if not self._opened:
        self._abort_connection_workflow()
    else:
        # This completes asynchronously, culminating in call to our method
        # `connection_lost()`
        assert self._transport is not None
        self._transport.abort()

_adapter_emit_data

_adapter_emit_data(data: bytes) -> None

Take ownership of data and send it to AMQP server as soon as possible.

PARAMETER DESCRIPTION
data

TYPE: bytes

Source code in pika/adapters/base_connection.py
@override
def _adapter_emit_data(self, data: bytes) -> None:
    """
    Take ownership of data and send it to AMQP server as soon as possible.

    :param data:
    """
    assert self._transport is not None
    self._transport.write(data)

_adapter_remove_timeout

_adapter_remove_timeout(timeout_id: object) -> None

Remove a scheduled timeout.

PARAMETER DESCRIPTION
timeout_id

Timeout handle to cancel.

TYPE: object

Source code in pika/adapters/base_connection.py
@override
def _adapter_remove_timeout(self, timeout_id: object) -> None:
    """
    Remove a scheduled timeout.

    :param timeout_id: Timeout handle to cancel.
    """
    cast(nbio_interface.AbstractTimerReference, timeout_id).cancel()

_add_channel_callbacks

_add_channel_callbacks(channel_number: int) -> None

Add the appropriate callbacks for the specified channel number.

PARAMETER DESCRIPTION
channel_number

The channel number for the callbacks

TYPE: int

Source code in pika/connection.py
def _add_channel_callbacks(self, channel_number: int) -> None:
    """
    Add the appropriate callbacks for the specified channel number.

    :param channel_number: The channel number for the callbacks
    """
    # This permits us to garbage-collect our reference to the channel
    # regardless of whether it was closed by client or broker, and do so
    # after all channel-close callbacks.
    self._channels[channel_number]._add_on_cleanup_callback(
        self._on_channel_cleanup)

_add_connection_start_callback

_add_connection_start_callback() -> None

Add a callback for when a Connection.Start frame is received from the broker.

Source code in pika/connection.py
def _add_connection_start_callback(self) -> None:
    """Add a callback for when a Connection.Start frame is received from the broker."""
    self.callbacks.add(0, spec.Connection.Start, self._on_connection_start)

_add_connection_tune_callback

_add_connection_tune_callback() -> None

Add a callback for when a Connection.Tune frame is received.

Source code in pika/connection.py
def _add_connection_tune_callback(self) -> None:
    """Add a callback for when a Connection.Tune frame is received."""
    self.callbacks.add(0, spec.Connection.Tune, self._on_connection_tune)

_check_for_protocol_mismatch

_check_for_protocol_mismatch(value: Method[Start]) -> None

Invoked when starting a connection to make sure it's a supported protocol.

PARAMETER DESCRIPTION
value

The frame to check

TYPE: Method[Start]

RAISES DESCRIPTION
ProtocolVersionMismatch
Source code in pika/connection.py
def _check_for_protocol_mismatch(
        self, value: frame.Method[spec.Connection.Start]) -> None:
    """
    Invoked when starting a connection to make sure it's a supported protocol.

    :param value: The frame to check
    :raises ProtocolVersionMismatch:
    """
    if ((value.method.version_major, value.method.version_minor)
            != spec.PROTOCOL_VERSION[0:2]):
        raise exceptions.ProtocolVersionMismatch(frame.ProtocolHeader(),
                                                 value)

_close_channels

_close_channels(reply_code: int, reply_text: str) -> None

Initiate graceful closing of channels that are in OPEN or OPENING states, passing reply_code and reply_text.

PARAMETER DESCRIPTION
reply_code

The code for why the channels are being closed

TYPE: int

reply_text

The text reason for why the channels are closing

TYPE: str

Source code in pika/connection.py
def _close_channels(self, reply_code: int, reply_text: str) -> None:
    """
    Initiate graceful closing of channels that are in OPEN or OPENING states, passing reply_code
    and reply_text.

    :param reply_code: The code for why the channels are being closed
    :param reply_text: The text reason for why the channels are closing
    """
    assert self.is_open, str(self)

    for channel_number in list(self._channels.keys()):
        chan = self._channels[channel_number]
        if not (chan.is_closing or chan.is_closed):
            chan.close(reply_code, reply_text)

_create_channel

_create_channel(
    channel_number: int,
    on_open_callback: Callable[[Channel], Any] | None,
) -> Channel

Create a new channel using the specified channel number and calling back the method specified by on_open_callback.

PARAMETER DESCRIPTION
channel_number

The channel number to use

TYPE: int

on_open_callback

The callback when the channel is opened. The callback will be invoked with the Channel instance as its only argument.

TYPE: Callable[[Channel], Any] | None

Source code in pika/connection.py
def _create_channel(
        self, channel_number: int,
        on_open_callback: Callable[[Channel], Any] | None) -> Channel:
    """
    Create a new channel using the specified channel number and calling back the method
    specified by on_open_callback.

    :param channel_number: The channel number to use
    :param on_open_callback: The callback when the channel is opened. The callback will be
        invoked with the `Channel` instance as its only argument.
    """
    LOGGER.debug('Creating channel %s', channel_number)
    return pika.channel.Channel(self, channel_number, on_open_callback)

_create_heartbeat_checker

_create_heartbeat_checker() -> HeartbeatChecker | None

Create a heartbeat checker instance if there is a heartbeat interval set.

Source code in pika/connection.py
def _create_heartbeat_checker(
        self) -> pika.heartbeat.HeartbeatChecker | None:
    """Create a heartbeat checker instance if there is a heartbeat interval set."""
    if isinstance(self.params.heartbeat, int) and self.params.heartbeat > 0:
        LOGGER.debug('Creating a HeartbeatChecker: %r',
                     self.params.heartbeat)
        return pika.heartbeat.HeartbeatChecker(self, self.params.heartbeat)

    return None

_default_on_connection_error

_default_on_connection_error(
    _connection_unused: Connection, error: Exception
) -> None

Default behavior when the connecting connection cannot connect and user didn't supply own on_connection_error callback.

PARAMETER DESCRIPTION
_connection_unused

The connection instance (unused)

TYPE: Connection

error

The exception that caused the failure

TYPE: Exception

RAISES DESCRIPTION
Exception

the given error

Source code in pika/connection.py
def _default_on_connection_error(self, _connection_unused: Connection,
                                 error: Exception) -> None:
    """
    Default behavior when the connecting connection cannot connect and user didn't supply own
    `on_connection_error` callback.

    :param _connection_unused: The connection instance (unused)
    :param error: The exception that caused the failure
    :raises Exception: the given error
    """
    raise error

_deliver_frame_to_channel

_deliver_frame_to_channel(value: Frame) -> None

Deliver the frame to the channel specified in the frame.

PARAMETER DESCRIPTION
value

The frame to deliver

TYPE: Frame

Source code in pika/connection.py
def _deliver_frame_to_channel(self, value: frame.Frame) -> None:
    """
    Deliver the frame to the channel specified in the frame.

    :param value: The frame to deliver
    """
    if value.channel_number not in self._channels:
        # This should never happen and would constitute breach of the
        # protocol
        LOGGER.critical(
            'Received %s frame for unregistered channel %i on %s',
            value.NAME, value.channel_number, self)
        return

    self._channels[value.channel_number]._handle_content_frame(value)

_ensure_closed

_ensure_closed() -> None

If the connection is not closed, close it.

Source code in pika/connection.py
def _ensure_closed(self) -> None:
    """If the connection is not closed, close it."""
    if self.is_open:
        self.close()

_get_body_frame_max_length

_get_body_frame_max_length() -> int

Calculate the maximum amount of bytes that can be in a body frame.

Source code in pika/connection.py
def _get_body_frame_max_length(self) -> int:
    """Calculate the maximum amount of bytes that can be in a body frame."""
    return (self.params.frame_max - spec.FRAME_HEADER_SIZE -
            spec.FRAME_END_SIZE)

_get_credentials

_get_credentials(
    method_frame: Method[Start],
) -> tuple[str, bytes | None]

Get credentials for authentication.

PARAMETER DESCRIPTION
method_frame

The Connection.Start frame

TYPE: Method[Start]

Source code in pika/connection.py
def _get_credentials(
    self, method_frame: frame.Method[spec.Connection.Start]
) -> tuple[str, bytes | None]:
    """
    Get credentials for authentication.

    :param method_frame: The Connection.Start frame
    """
    (auth_type,
     response) = self.params.credentials.response_for(method_frame.method)
    if not auth_type:
        raise exceptions.AuthenticationError(self.params.credentials.TYPE,
                                             host=self.params.host,
                                             port=self.params.port)
    self.params.credentials.erase_credentials()
    return auth_type, response

_get_write_buffer_size

_get_write_buffer_size() -> int
RETURNS DESCRIPTION
int

Current size of output data buffered by the transport

Source code in pika/adapters/select_connection.py
def _get_write_buffer_size(self) -> int:
    """
    :returns: Current size of output data buffered by the transport
    """
    assert self._transport is not None
    return self._transport.get_write_buffer_size()

_handle_connection_workflow_failure

_handle_connection_workflow_failure(
    error: Exception | None,
) -> None

Handle failure of self-initiated stack bring-up and call Connection._on_stream_terminated() if connection is not in closed state yet. Called by adapter layer when the full-stack connection workflow fails.

PARAMETER DESCRIPTION
error

Exception (or None) describing the reason for failure or None if the connection workflow was aborted.

TYPE: Exception | None

Source code in pika/adapters/base_connection.py
def _handle_connection_workflow_failure(self,
                                        error: Exception | None) -> None:
    """
    Handle failure of self-initiated stack bring-up and call
    `Connection._on_stream_terminated()` if connection is not in closed state yet. Called by
    adapter layer when the full-stack connection workflow fails.

    :param error: Exception (or None) describing the reason for failure or None if the
        connection workflow was aborted.
    """
    if error is None:
        LOGGER.info('Self-initiated stack bring-up aborted.')
    else:
        LOGGER.error('Self-initiated stack bring-up failed: %r', error)

    if not self.is_closed:
        self._on_stream_terminated(error)
    else:
        # This may happen when AMQP layer bring up was started but did not
        # complete
        LOGGER.debug('_handle_connection_workflow_failure(): '
                     'suppressing - connection already closed.')

_has_pending_callbacks

_has_pending_callbacks(value: Method[Method]) -> int | None

Return true if there are any callbacks pending for the specified frame.

PARAMETER DESCRIPTION
value

The frame to check

TYPE: Method[Method]

Source code in pika/connection.py
def _has_pending_callbacks(
        self, value: frame.Method[amqp_object.Method]) -> int | None:
    """
    Return true if there are any callbacks pending for the specified frame.

    :param value: The frame to check
    """
    return self.callbacks.pending(value.channel_number, value.method)

_init_connection_state

_init_connection_state() -> None

Initialize or reset all of our internal state variables for a given connection.

If we disconnect and reconnect, all of our state needs to be wiped.

Source code in pika/adapters/base_connection.py
@override
def _init_connection_state(self) -> None:
    """
    Initialize or reset all of our internal state variables for a given connection.

    If we disconnect and reconnect, all of our state needs to be wiped.
    """
    super()._init_connection_state()

    self._connection_workflow = None
    self._transport = None
    self._got_eof = False

_is_method_frame

_is_method_frame(value: Frame) -> bool

Returns true if the frame is a method frame.

PARAMETER DESCRIPTION
value

The frame to evaluate

TYPE: Frame

Source code in pika/connection.py
def _is_method_frame(self, value: frame.Frame) -> bool:
    """
    Returns true if the frame is a method frame.

    :param value: The frame to evaluate
    """
    return value.frame_type == spec.FRAME_METHOD

_is_protocol_header_frame

_is_protocol_header_frame(value: Frame) -> bool

Returns True if it's a protocol header frame.

PARAMETER DESCRIPTION
value

Frame to inspect

TYPE: Frame

Source code in pika/connection.py
def _is_protocol_header_frame(self, value: frame.Frame) -> bool:
    """
    Returns True if it's a protocol header frame.

    :param value: Frame to inspect
    """
    return isinstance(value, frame.ProtocolHeader)

_negotiate_integer_value staticmethod

_negotiate_integer_value(
    client_value: int | None, server_value: int | None
) -> int

Negotiates two values.

If either of them is 0 or None, returns the other one. If both are positive integers, returns the smallest one.

PARAMETER DESCRIPTION
client_value

The client value

TYPE: int | None

server_value

The server value

TYPE: int | None

Source code in pika/connection.py
@staticmethod
def _negotiate_integer_value(client_value: int | None,
                             server_value: int | None) -> int:
    """
    Negotiates two values.

    If either of them is 0 or None, returns the other one. If both are positive integers,
    returns the smallest one.

    :param client_value: The client value
    :param server_value: The server value
    """
    if client_value is None:
        client_value = 0
    if server_value is None:
        server_value = 0

    # this is consistent with how Java client and Bunny
    # perform negotiation, see pika/pika#874
    if client_value == 0 or server_value == 0:
        val = max(client_value, server_value)
    else:
        val = min(client_value, server_value)

    return val

_next_channel_number

_next_channel_number() -> int

Return the next available channel number or raise an exception.

Source code in pika/connection.py
def _next_channel_number(self) -> int:
    """Return the next available channel number or raise an exception."""
    limit = self.params.channel_max or pika.channel.MAX_CHANNELS
    if len(self._channels) >= limit:
        raise exceptions.NoFreeChannels(host=self.params.host,
                                        port=self.params.port)

    for num in range(1, len(self._channels) + 1):
        if num not in self._channels:
            return num
    return len(self._channels) + 1

_on_blocked_connection_timeout

_on_blocked_connection_timeout() -> None

Called when the "connection blocked timeout" expires.

When this happens, we tear down the connection

Source code in pika/connection.py
def _on_blocked_connection_timeout(self) -> None:
    """
    Called when the "connection blocked timeout" expires.

    When this happens, we tear down the connection
    """
    self._blocked_conn_timer = None
    self._terminate_stream(
        exceptions.ConnectionBlockedTimeout(
            'Blocked connection timeout expired.',
            host=self.params.host,
            port=self.params.port))

_on_channel_cleanup

_on_channel_cleanup(channel: Channel) -> None

Remove the channel from the dict of channels when Channel.CloseOk is sent.

If connection is closing and no more channels remain, proceed to _on_close_ready.

PARAMETER DESCRIPTION
channel

channel instance

TYPE: Channel

Source code in pika/connection.py
def _on_channel_cleanup(self, channel: Channel) -> None:
    """
    Remove the channel from the dict of channels when Channel.CloseOk is sent.

    If connection is closing and no more channels remain, proceed to `_on_close_ready`.

    :param channel: channel instance
    """
    try:
        del self._channels[channel.channel_number]
        LOGGER.debug('Removed channel %s', channel.channel_number)
    except KeyError:
        LOGGER.error('Channel %r not in channels', channel.channel_number)
    if self.is_closing:
        if not self._channels:
            # Initiate graceful closing of the connection
            self._on_close_ready()
        else:
            # Once Connection enters CLOSING state, all remaining channels
            # should also be in CLOSING state. Deviation from this would
            # prevent Connection from completing its closing procedure.
            channels_not_in_closing_state = [
                chan for chan in self._channels.values()
                if not chan.is_closing
            ]
            if channels_not_in_closing_state:
                LOGGER.critical(
                    'Connection in CLOSING state has non-CLOSING '
                    'channels: %r', channels_not_in_closing_state)

_on_close_ready

_on_close_ready() -> None

Called when the Connection is in a state that it can close after a close has been requested by client.

This happens after all of the channels are closed that were open when the close request was made.

Source code in pika/connection.py
def _on_close_ready(self) -> None:
    """
    Called when the Connection is in a state that it can close after a close has been requested
    by client.

    This happens after all of the channels are closed that were open when the close request was
    made.
    """
    if self.is_closed:
        LOGGER.warning('_on_close_ready invoked when already closed')
        return

    assert isinstance(self._error, exceptions.ConnectionClosed)
    self._send_connection_close(self._error.reply_code,
                                self._error.reply_text)

_on_connection_blocked

_on_connection_blocked(
    _connection: Connection, method_frame: Method[Blocked]
) -> None

Handle Connection.Blocked notification from RabbitMQ broker.

PARAMETER DESCRIPTION
_connection

The connection instance (unused)

TYPE: Connection

method_frame

method frame having method member of type pika.spec.Connection.Blocked

TYPE: Method[Blocked]

Source code in pika/connection.py
def _on_connection_blocked(
        self, _connection: Connection,
        method_frame: frame.Method[spec.Connection.Blocked]) -> None:
    """
    Handle Connection.Blocked notification from RabbitMQ broker.

    :param _connection: The connection instance (unused)
    :param method_frame: method frame having `method` member of type
        `pika.spec.Connection.Blocked`
    """
    LOGGER.warning('Received %s from broker', method_frame)

    if self._blocked_conn_timer is not None:
        # RabbitMQ is not supposed to repeat Connection.Blocked, but it
        # doesn't hurt to be careful
        LOGGER.warning(
            '_blocked_conn_timer %s already set when '
            '_on_connection_blocked is called', self._blocked_conn_timer)
    else:
        assert self.params.blocked_connection_timeout is not None
        self._blocked_conn_timer = self._adapter_call_later(
            self.params.blocked_connection_timeout,
            self._on_blocked_connection_timeout)

_on_connection_close_from_broker

_on_connection_close_from_broker(
    method_frame: Method[Close],
) -> None

Called when the connection is closed remotely via Connection.Close frame from broker.

PARAMETER DESCRIPTION
method_frame

The Connection.Close frame

TYPE: Method[Close]

Source code in pika/connection.py
def _on_connection_close_from_broker(
        self, method_frame: frame.Method[spec.Connection.Close]) -> None:
    """
    Called when the connection is closed remotely via Connection.Close frame from broker.

    :param method_frame: The Connection.Close frame
    """
    LOGGER.debug('_on_connection_close_from_broker: frame=%s', method_frame)

    # The generated spec types reply_code as optional, but a frame
    # decoded off the wire always carries it.
    reply_code = method_frame.method.reply_code
    assert reply_code is not None

    self._terminate_stream(
        exceptions.ConnectionClosedByBroker(reply_code,
                                            method_frame.method.reply_text,
                                            host=self.params.host,
                                            port=self.params.port))

_on_connection_close_ok

_on_connection_close_ok(
    method_frame: Method[CloseOk],
) -> None

Called when Connection.CloseOk is received from remote.

PARAMETER DESCRIPTION
method_frame

The Connection.CloseOk frame

TYPE: Method[CloseOk]

Source code in pika/connection.py
def _on_connection_close_ok(
        self, method_frame: frame.Method[spec.Connection.CloseOk]) -> None:
    """
    Called when Connection.CloseOk is received from remote.

    :param method_frame: The Connection.CloseOk frame
    """
    LOGGER.debug('_on_connection_close_ok: frame=%s', method_frame)

    self._terminate_stream(None)

_on_connection_open_ok

_on_connection_open_ok(
    method_frame: Method[OpenOk],
) -> None

This is called once we have tuned the connection with the server and called the Connection.Open on the server and it has replied with Connection.Ok.

PARAMETER DESCRIPTION
method_frame

Server response frame

TYPE: Method[OpenOk]

Source code in pika/connection.py
def _on_connection_open_ok(
        self, method_frame: frame.Method[spec.Connection.OpenOk]) -> None:
    """
    This is called once we have tuned the connection with the server and called the
    Connection.Open on the server and it has replied with Connection.Ok.

    :param method_frame: Server response frame
    """
    self._opened = True

    self.known_hosts = method_frame.method.known_hosts

    # We're now connected at the AMQP level
    self._set_connection_state(self.CONNECTION_OPEN)

    # Call our initial callback that we're open
    self.callbacks.process(0, self.ON_CONNECTION_OPEN_OK, self, self)

_on_connection_start

_on_connection_start(method_frame: Method[Start]) -> None

This is called as a callback once we have received a Connection.Start from the server.

PARAMETER DESCRIPTION
method_frame

The frame received

TYPE: Method[Start]

RAISES DESCRIPTION
UnexpectedFrameError
Source code in pika/connection.py
def _on_connection_start(
        self, method_frame: frame.Method[spec.Connection.Start]) -> None:
    """
    This is called as a callback once we have received a Connection.Start from the server.

    :param method_frame: The frame received
    :raises UnexpectedFrameError:
    """
    self._set_connection_state(self.CONNECTION_START)

    try:
        if self._is_protocol_header_frame(method_frame):
            raise exceptions.UnexpectedFrameError(method_frame)
        self._check_for_protocol_mismatch(method_frame)
        self._set_server_information(method_frame)
        self._add_connection_tune_callback()
        self._send_connection_start_ok(*self._get_credentials(method_frame))
    except Exception as error:
        LOGGER.exception('Error processing Connection.Start.')
        self._terminate_stream(error)

_on_connection_tune

_on_connection_tune(method_frame: Method[Tune]) -> None

Once the Broker sends back a Connection.Tune, we will set our tuning variables that have been returned to us and kick off the Heartbeat monitor if required, send our TuneOk and then the Connection. Open rpc call on channel 0.

PARAMETER DESCRIPTION
method_frame

The frame received

TYPE: Method[Tune]

Source code in pika/connection.py
def _on_connection_tune(
        self, method_frame: frame.Method[spec.Connection.Tune]) -> None:
    """
    Once the Broker sends back a Connection.Tune, we will set our tuning variables that have
    been returned to us and kick off the Heartbeat monitor if required, send our TuneOk and then
    the Connection. Open rpc call on channel 0.

    :param method_frame: The frame received
    """
    self._set_connection_state(self.CONNECTION_TUNE)

    # Get our max channels, frames and heartbeat interval
    self.params.channel_max = Connection._negotiate_integer_value(
        self.params.channel_max, method_frame.method.channel_max)
    self.params.frame_max = Connection._negotiate_integer_value(
        self.params.frame_max, method_frame.method.frame_max)

    if callable(self.params.heartbeat):
        ret_heartbeat = self.params.heartbeat(self,
                                              method_frame.method.heartbeat)
        if ret_heartbeat is None or callable(ret_heartbeat):
            # Enforce callback-specific restrictions on callback's return value
            raise TypeError('heartbeat callback must not return None '
                            f'or callable, but got {ret_heartbeat!r}')

        # Leave it to heartbeat setter deal with the rest of the validation
        self.params.heartbeat = ret_heartbeat

    # Negotiate heartbeat timeout
    self.params.heartbeat = self._tune_heartbeat_timeout(
        client_value=self.params.heartbeat,
        server_value=method_frame.method.heartbeat)

    # Calculate the maximum pieces for body frames
    self._body_max_length = self._get_body_frame_max_length()

    # Create a new heartbeat checker if needed
    self._heartbeat_checker = self._create_heartbeat_checker()

    # Send the TuneOk response with what we've agreed upon
    self._send_connection_tune_ok()

    # Send the Connection.Open RPC call for the vhost
    self._send_connection_open()

_on_connection_unblocked

_on_connection_unblocked(
    _connection: Connection, method_frame: Method[Unblocked]
) -> None

Handle Connection.Unblocked notification from RabbitMQ broker.

PARAMETER DESCRIPTION
_connection

The connection instance (unused)

TYPE: Connection

method_frame

method frame having method member of type pika.spec.Connection.Blocked

TYPE: Method[Unblocked]

Source code in pika/connection.py
def _on_connection_unblocked(
        self, _connection: Connection,
        method_frame: frame.Method[spec.Connection.Unblocked]) -> None:
    """
    Handle Connection.Unblocked notification from RabbitMQ broker.

    :param _connection: The connection instance (unused)
    :param method_frame: method frame having `method` member of type
        `pika.spec.Connection.Blocked`
    """
    LOGGER.info('Received %s from broker', method_frame)

    if self._blocked_conn_timer is None:
        # RabbitMQ is supposed to pair Connection.Blocked/Unblocked, but it
        # doesn't hurt to be careful
        LOGGER.warning('_blocked_conn_timer was not active when '
                       '_on_connection_unblocked called')
    else:
        self._adapter_remove_timeout(self._blocked_conn_timer)
        self._blocked_conn_timer = None

_on_connection_workflow_done

_on_connection_workflow_done(
    conn_or_exc: BaseConnection | Exception,
) -> None

AMQPConnectionWorkflow completion callback.

PARAMETER DESCRIPTION
conn_or_exc

Our own connection instance on success; exception on failure. See AbstractAMQPConnectionWorkflow.start() for details. Type: BaseConnection | Exception.

TYPE: BaseConnection | Exception

Source code in pika/adapters/base_connection.py
def _on_connection_workflow_done(
        self, conn_or_exc: BaseConnection | Exception) -> None:
    """
    `AMQPConnectionWorkflow` completion callback.

    :param conn_or_exc: Our own connection
        instance on success; exception on failure. See
        `AbstractAMQPConnectionWorkflow.start()` for details.
        Type: `BaseConnection | Exception`.
    """
    LOGGER.debug('Full-stack connection workflow completed: %r',
                 conn_or_exc)

    self._connection_workflow = None

    # Notify protocol of failure
    if isinstance(conn_or_exc, Exception):
        self._transport = None
        error: Exception | None
        if isinstance(conn_or_exc,
                      connection_workflow.AMQPConnectionWorkflowAborted):
            LOGGER.info('Full-stack connection workflow aborted: %r',
                        conn_or_exc)
            error = None
        else:
            LOGGER.error('Full-stack connection workflow failed: %r',
                         conn_or_exc)
            if (isinstance(conn_or_exc,
                           connection_workflow.AMQPConnectionWorkflowFailed)
                    and isinstance(
                        conn_or_exc.exceptions[-1],
                        connection_workflow.AMQPConnectorSocketConnectError)
               ):
                last_exc = conn_or_exc.exceptions[-1]
                error = pika.exceptions.AMQPConnectionError(
                    conn_or_exc, host=last_exc.host, port=last_exc.port)
            else:
                error = conn_or_exc

        self._handle_connection_workflow_failure(error)
    else:
        # NOTE: On success, the stack will be up already, so there is no
        #       corresponding callback.
        assert conn_or_exc is self, (
            f'Expected self conn={self!r} from workflow, but got {conn_or_exc!r}.'
        )

_on_data_available

_on_data_available(data_in: bytes) -> None

This is called by our Adapter, passing in the data from the socket.

As long as we have buffer try and map out frame data.

PARAMETER DESCRIPTION
data_in

The data that is available to read

TYPE: bytes

Source code in pika/connection.py
def _on_data_available(self, data_in: bytes) -> None:
    """
    This is called by our Adapter, passing in the data from the socket.

    As long as we have buffer try and map out frame data.

    :param data_in: The data that is available to read
    """
    buffer = self._frame_buffer
    buffer += data_in

    if self._processing_frame_buffer:
        # Re-entrant call: a frame callback caused more data to be read.
        # The data was appended above; the outer invocation's loop will
        # decode it, since trimming mid-loop would corrupt its offset.
        return

    # Decode frames by walking an offset and trim the buffer once per data
    # event; trimming (copying) the buffer after every frame is quadratic
    # when one read contains many frames.
    self._processing_frame_buffer = True
    offset = 0
    try:
        while offset < len(buffer):
            consumed_count, frame_value = frame.decode_frame(buffer, offset)
            if not frame_value:
                break
            offset += consumed_count
            self.bytes_received += consumed_count
            self._process_frame(frame_value)
            if self._frame_buffer is not buffer:
                # A frame callback reset the connection state (e.g. stream
                # terminated); leftover data is no longer relevant.
                return
    finally:
        self._processing_frame_buffer = False
        if offset and self._frame_buffer is buffer:
            del buffer[:offset]

_on_stream_connected

_on_stream_connected() -> None

Invoked when the socket is connected and it's time to start speaking AMQP with the broker.

Source code in pika/connection.py
def _on_stream_connected(self) -> None:
    """Invoked when the socket is connected and it's time to start speaking AMQP with the
    broker.
    """
    self._set_connection_state(self.CONNECTION_PROTOCOL)

    # Start the communication with the RabbitMQ Broker
    self._send_frame(frame.ProtocolHeader())

_on_stream_terminated

_on_stream_terminated(error: Exception | None) -> None

Handle termination of stack (including TCP layer) or failure to establish the stack.

Notify registered ON_CONNECTION_ERROR or ON_CONNECTION_CLOSED callbacks, depending on whether the connection was opening or open.

PARAMETER DESCRIPTION
error

Exception (or None). None means that the transport was aborted internally and exception in self._error represents the cause. Otherwise it's an exception object that describes the unexpected loss of connection.

TYPE: Exception | None

Source code in pika/connection.py
def _on_stream_terminated(self, error: Exception | None) -> None:
    """
    Handle termination of stack (including TCP layer) or failure to establish the stack.

    Notify registered ON_CONNECTION_ERROR or ON_CONNECTION_CLOSED callbacks, depending on
    whether the connection was opening or open.

    :param error: Exception (or None). None means that the transport was aborted internally and
        exception in `self._error` represents the cause. Otherwise it's an exception object that
        describes the unexpected loss of connection.
    """
    LOGGER.info(
        'AMQP stack terminated, failed to connect, or aborted: '
        'opened=%r, error-arg=%r; pending-error=%r', self._opened, error,
        self._error)

    if error is not None:
        if self._error is not None:
            LOGGER.debug(
                '_on_stream_terminated(): overriding '
                'pending-error=%r with %r', self._error, error)
        self._error = error
    else:
        assert self._error is not None, (
            '_on_stream_terminated() expected self._error to be populated '
            'with reason for terminating stack.')

    # Stop the heartbeat checker if it exists
    self._remove_heartbeat()

    # Remove connection management callbacks
    self._remove_callbacks(0,
                           [spec.Connection.Close, spec.Connection.Start])

    if self.params.blocked_connection_timeout is not None:
        self._remove_callbacks(
            0, [spec.Connection.Blocked, spec.Connection.Unblocked])

    if not self._opened and isinstance(
            self._error,
        (exceptions.StreamLostError, exceptions.ConnectionClosedByBroker)):
        # Heuristically deduce error based on connection state
        original_error = self._error
        deduced_error: Exception | None = None
        if self.connection_state == self.CONNECTION_PROTOCOL:
            LOGGER.error('Probably incompatible Protocol Versions')
            deduced_error = exceptions.IncompatibleProtocolError(
                repr(self._error),
                host=self.params.host,
                port=self.params.port)
        elif self.connection_state == self.CONNECTION_START:
            LOGGER.error(
                'Connection closed while authenticating indicating a '
                'probable authentication error')
            deduced_error = exceptions.ProbableAuthenticationError(
                repr(self._error),
                host=self.params.host,
                port=self.params.port)
        elif self.connection_state == self.CONNECTION_TUNE:
            LOGGER.error('Connection closed while tuning the connection '
                         'indicating a probable permission error when '
                         'accessing a virtual host')
            deduced_error = exceptions.ProbableAccessDeniedError(
                repr(self._error),
                host=self.params.host,
                port=self.params.port)
        elif self.connection_state not in [
                self.CONNECTION_OPEN, self.CONNECTION_CLOSED,
                self.CONNECTION_CLOSING
        ]:
            LOGGER.warning('Unexpected connection state on disconnect: %i',
                           self.connection_state)
        if deduced_error is not None:
            # pika/pika#1390: keep the original error (and the __cause__
            # traceback it carries) reachable behind the heuristic one.
            deduced_error.__cause__ = original_error
            self._error = deduced_error

    # Transition to closed state
    self._set_connection_state(self.CONNECTION_CLOSED)

    # Inform our channel proxies, if any are still around
    for channel in list(self._channels.keys()):
        if channel not in self._channels:
            continue

        self._channels[channel]._on_close_meta(self._error)

    # Inform interested parties
    if not self._opened:
        LOGGER.info('Connection setup terminated due to %r', self._error)
        self.callbacks.process(0, self.ON_CONNECTION_ERROR, self, self,
                               self._error)
    else:
        LOGGER.info('Stack terminated due to %r', self._error)
        self.callbacks.process(0, self.ON_CONNECTION_CLOSED, self, self,
                               self._error)

    # Reset connection properties
    self._init_connection_state()

_output_marshaled_frames

_output_marshaled_frames(
    marshaled_frames: Sequence[bytes],
) -> None

Output list of marshaled frames to buffer and update stats.

PARAMETER DESCRIPTION
marshaled_frames

A list of frames marshaled to bytes

TYPE: Sequence[bytes]

Source code in pika/connection.py
def _output_marshaled_frames(self,
                             marshaled_frames: Sequence[bytes]) -> None:
    """
    Output list of marshaled frames to buffer and update stats.

    :param marshaled_frames: A list of frames marshaled to bytes
    """
    for marshaled_frame in marshaled_frames:
        self.bytes_sent += len(marshaled_frame)
        self.frames_sent += 1
        self._adapter_emit_data(marshaled_frame)

_process_callbacks

_process_callbacks(frame_value: Frame) -> bool

Process the callbacks for the frame if the frame is a method frame and if it has any callbacks pending.

PARAMETER DESCRIPTION
frame_value

The frame to process

TYPE: Frame

Source code in pika/connection.py
def _process_callbacks(self, frame_value: frame.Frame) -> bool:
    """
    Process the callbacks for the frame if the frame is a method frame and if it has any
    callbacks pending.

    :param frame_value: The frame to process
    """
    if (self._is_method_frame(frame_value) and
            self._has_pending_callbacks(cast(frame.Method, frame_value))):
        self.callbacks.process(
            cast(frame.Method, frame_value).channel_number,  # Prefix
            cast(frame.Method, frame_value).method,  # Key
            self,  # Caller
            cast(frame.Method, frame_value))  # Args
        return True
    return False

_process_frame

_process_frame(frame_value: Frame | ProtocolHeader) -> None

Process an inbound frame from the socket.

PARAMETER DESCRIPTION
frame_value

The frame to process

TYPE: Frame | ProtocolHeader

Source code in pika/connection.py
def _process_frame(
    self, frame_value: (frame.Frame | frame.ProtocolHeader)) -> None:
    """
    Process an inbound frame from the socket.

    :param frame_value: The frame to process
    """
    # Will receive a frame type of -1 if protocol version mismatch
    frame_type = frame_value.frame_type
    if frame_type < 0:
        return

    # ProtocolHeader (frame_type == -1) is handled above; narrow
    # the type for mypy.
    assert not isinstance(frame_value, frame.ProtocolHeader)

    # Keep track of how many frames have been read
    self.frames_received += 1

    # Process any callbacks, if True, exit method
    if self._process_callbacks(frame_value):
        return

    # If a heartbeat is received, update the checker
    if frame_type == spec.FRAME_HEARTBEAT:
        if self._heartbeat_checker:
            self._heartbeat_checker.received()
        else:
            LOGGER.warning('Received heartbeat frame without a heartbeat '
                           'checker')

    # If the frame has a channel number beyond the base channel, deliver it
    elif frame_value.channel_number > 0:
        self._deliver_frame_to_channel(frame_value)

_proto_connection_lost

_proto_connection_lost(error: BaseException | None) -> None

Called upon loss or closing of TCP connection.

:py:class:.utils.nbio_interface.AbstractStreamProtocol implementation.

NOTE: connection_made() and connection_lost() are each called just once and in that order. All other callbacks are called between them.

PARAMETER DESCRIPTION
error

An exception (check for BaseException) indicates connection failure. None indicates that connection was closed on this side, such as when it's aborted or when AbstractStreamProtocol.eof_received() returns a falsy result. Type: BaseException | None.

TYPE: BaseException | None

RAISES DESCRIPTION
Exception

Exception-based exception on error

Source code in pika/adapters/base_connection.py
def _proto_connection_lost(self, error: BaseException | None) -> None:
    """
    Called upon loss or closing of TCP connection.

    :py:class:`.utils.nbio_interface.AbstractStreamProtocol` implementation.

    NOTE: `connection_made()` and `connection_lost()` are each called just
    once and in that order. All other callbacks are called between them.

    :param error: An exception (check for
        `BaseException`) indicates connection failure. None indicates that
        connection was closed on this side, such as when it's aborted or
        when `AbstractStreamProtocol.eof_received()` returns a falsy result.
        Type: `BaseException | None`.
    :raises Exception: Exception-based exception on error
    """
    self._transport = None

    if error is None:
        # Either result of `eof_received()` or abort
        reason = None
        if self._got_eof:
            reason = pika.exceptions.StreamLostError(
                'Transport indicated EOF',
                host=self.params.host,
                port=self.params.port)
    else:
        reason = pika.exceptions.StreamLostError(
            f'Stream connection lost: {error!r}',
            host=self.params.host,
            port=self.params.port)
        # pika/pika#1390: keep the original traceback reachable. This pins
        # the error's traceback (and the frame locals it captured) for as
        # long as this reason is held as self._error, an accepted cost for
        # diagnosability.
        reason.__cause__ = error

    LOGGER.log(logging.DEBUG if reason is None else logging.ERROR,
               'connection_lost: %r', reason)

    self._on_stream_terminated(reason)

_proto_connection_made

_proto_connection_made(
    transport: AbstractStreamTransport,
) -> None

Introduces transport to protocol after transport is connected.

:py:class:.utils.nbio_interface.AbstractStreamProtocol implementation.

PARAMETER DESCRIPTION
transport

TYPE: AbstractStreamTransport

RAISES DESCRIPTION
Exception

Exception-based exception on error

Source code in pika/adapters/base_connection.py
def _proto_connection_made(
        self, transport: nbio_interface.AbstractStreamTransport) -> None:
    """
    Introduces transport to protocol after transport is connected.

    :py:class:`.utils.nbio_interface.AbstractStreamProtocol` implementation.

    :param transport:
    :raises Exception: Exception-based exception on error
    """
    self._transport = transport

    # Let connection know that stream is available
    self._on_stream_connected()

_proto_data_received

_proto_data_received(data: bytes) -> None

Called to deliver incoming data from the server to the protocol.

:py:class:.utils.nbio_interface.AbstractStreamProtocol implementation.

PARAMETER DESCRIPTION
data

Non-empty data bytes.

TYPE: bytes

RAISES DESCRIPTION
Exception

Exception-based exception on error

Source code in pika/adapters/base_connection.py
def _proto_data_received(self, data: bytes) -> None:
    """
    Called to deliver incoming data from the server to the protocol.

    :py:class:`.utils.nbio_interface.AbstractStreamProtocol` implementation.

    :param data: Non-empty data bytes.
    :raises Exception: Exception-based exception on error
    """
    self._on_data_available(data)

_proto_eof_received

_proto_eof_received() -> bool

Called after the remote peer shuts its write end of the connection. :py:class:.utils.nbio_interface.AbstractStreamProtocol implementation.

RETURNS DESCRIPTION
bool

A falsy value (including None) will cause the transport to close itself, resulting in an eventual connection_lost() call from the transport. If a truthy value is returned, it will be the protocol's responsibility to close/abort the transport.

RAISES DESCRIPTION
Exception

Exception-based exception on error

Source code in pika/adapters/base_connection.py
def _proto_eof_received(self) -> bool:
    """
    Called after the remote peer shuts its write end of the connection.
    :py:class:`.utils.nbio_interface.AbstractStreamProtocol` implementation.

    :returns: A falsy value (including None) will cause the transport to close itself, resulting
        in an eventual `connection_lost()` call from the transport. If a truthy value is
        returned, it will be the protocol's responsibility to close/abort the transport.
    :raises Exception: Exception-based exception on error
    """
    LOGGER.error('Transport indicated EOF.')

    self._got_eof = True

    # This is how a reset connection will typically present itself
    # when we have nothing to send to the server over plaintext stream.
    #
    # Have transport tear down the connection and invoke our
    # `connection_lost` method
    return False

_remove_callbacks

_remove_callbacks(
    channel_number: int,
    method_classes: Sequence[type[Method]],
) -> None

Remove the callbacks for the specified channel number and list of method frames.

PARAMETER DESCRIPTION
channel_number

The channel number to remove the callback on

TYPE: int

method_classes

The method classes (derived from pika.amqp_object.Method) for the callbacks

TYPE: Sequence[type[Method]]

Source code in pika/connection.py
def _remove_callbacks(
        self, channel_number: int,
        method_classes: Sequence[type[amqp_object.Method]]) -> None:
    """
    Remove the callbacks for the specified channel number and list of method frames.

    :param channel_number: The channel number to remove the callback on
    :param method_classes: The method classes (derived from `pika.amqp_object.Method`) for the
        callbacks
    """
    for method_cls in method_classes:
        self.callbacks.remove(str(channel_number), method_cls)

_remove_heartbeat

_remove_heartbeat() -> None

Stop the heartbeat checker if it exists.

Source code in pika/connection.py
def _remove_heartbeat(self) -> None:
    """Stop the heartbeat checker if it exists."""
    if self._heartbeat_checker:
        self._heartbeat_checker.stop()
        self._heartbeat_checker = None

_rpc

_rpc(
    channel_number: int,
    method: Method,
    callback: Callable[..., Any] | None = None,
    acceptable_replies: Sequence[type[Method]]
    | None = None,
) -> None

Make an RPC call for the given callback, channel number and method.

acceptable_replies lists out what responses we'll process from the server with the specified callback.

PARAMETER DESCRIPTION
channel_number

The channel number for the RPC call

TYPE: int

method

The method frame to call

TYPE: Method

callback

The callback for the RPC response

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

acceptable_replies

The replies this RPC call expects

TYPE: Sequence[type[Method]] | None DEFAULT: None

Source code in pika/connection.py
def _rpc(
    self,
    channel_number: int,
    method: amqp_object.Method,
    callback: Callable[..., Any] | None = None,
    acceptable_replies: Sequence[type[amqp_object.Method]] | None = None
) -> None:
    """
    Make an RPC call for the given callback, channel number and method.

    acceptable_replies lists out what responses we'll process from the server with the specified
    callback.

    :param channel_number: The channel number for the RPC call
    :param method: The method frame to call
    :param callback: The callback for the RPC response
    :param acceptable_replies: The replies this RPC call expects
    """
    # Validate that acceptable_replies is a list or None
    if acceptable_replies and not isinstance(acceptable_replies, list):
        raise TypeError('acceptable_replies should be list or None')

    # Validate the callback is callable
    if callback is not None:
        validators.require_callback(callback)
        assert acceptable_replies is not None
        for reply in acceptable_replies:
            self.callbacks.add(channel_number, reply, callback)

    # Send the rpc call to RabbitMQ
    self._send_method(channel_number, method)

_send_connection_close

_send_connection_close(
    reply_code: int, reply_text: str
) -> None

Send a Connection.Close method frame.

PARAMETER DESCRIPTION
reply_code

The reason for the close

TYPE: int

reply_text

The text reason for the close

TYPE: str

Source code in pika/connection.py
def _send_connection_close(self, reply_code: int, reply_text: str) -> None:
    """
    Send a Connection.Close method frame.

    :param reply_code: The reason for the close
    :param reply_text: The text reason for the close
    """
    self._rpc(0, spec.Connection.Close(reply_code, reply_text, 0, 0),
              self._on_connection_close_ok, [spec.Connection.CloseOk])

_send_connection_open

_send_connection_open() -> None

Send a Connection.Open frame.

Source code in pika/connection.py
def _send_connection_open(self) -> None:
    """Send a Connection.Open frame."""
    self._rpc(0, spec.Connection.Open(self.params.virtual_host,
                                      insist=True),
              self._on_connection_open_ok, [spec.Connection.OpenOk])

_send_connection_start_ok

_send_connection_start_ok(
    authentication_type: str, response: bytes | None
) -> None

Send a Connection.StartOk frame.

PARAMETER DESCRIPTION
authentication_type

The auth type value

TYPE: str

response

The encoded value to send

TYPE: bytes | None

Source code in pika/connection.py
def _send_connection_start_ok(self, authentication_type: str,
                              response: bytes | None) -> None:
    """
    Send a Connection.StartOk frame.

    :param authentication_type: The auth type value
    :param response: The encoded value to send
    """
    self._send_method(
        0,
        spec.Connection.StartOk(self._client_properties,
                                authentication_type, response,
                                self.params.locale))

_send_connection_tune_ok

_send_connection_tune_ok() -> None

Send a Connection.TuneOk frame.

Source code in pika/connection.py
def _send_connection_tune_ok(self) -> None:
    """Send a Connection.TuneOk frame."""
    # Tuning resolves any heartbeat callable and negotiates an integral
    # timeout before this is sent.
    heartbeat = self.params.heartbeat
    assert heartbeat is not None and not callable(heartbeat)

    self._send_method(
        0,
        spec.Connection.TuneOk(self.params.channel_max,
                               self.params.frame_max, heartbeat))

_send_frame

_send_frame(frame_value: Frame | ProtocolHeader) -> None

This appends the fully generated frame to send to the broker to the output buffer which will be then sent via the connection adapter.

PARAMETER DESCRIPTION
frame_value

The frame to write

TYPE: Frame | ProtocolHeader

RAISES DESCRIPTION
exceptions.ConnectionClosed
Source code in pika/connection.py
def _send_frame(
    self, frame_value: (frame.Frame | frame.ProtocolHeader)) -> None:
    """
    This appends the fully generated frame to send to the broker to the output buffer which will
    be then sent via the connection adapter.

    :param frame_value: The frame to write
    :raises exceptions.ConnectionClosed:
    """
    if self.is_closed:
        LOGGER.error('Attempted to send frame when closed')
        raise exceptions.ConnectionWrongStateError(
            'Attempted to send a frame on closed connection.')

    marshaled_frame = frame_value.marshal()
    self._output_marshaled_frames([marshaled_frame])

_send_message

_send_message(
    channel_number: int,
    method_frame: Method,
    content: tuple[BasicProperties, bytes],
) -> None

Publish a message.

PARAMETER DESCRIPTION
channel_number

The channel number for the frame

TYPE: int

method_frame

The method frame to send

TYPE: Method

content

A content frame, which is tuple of properties and body.

TYPE: tuple[BasicProperties, bytes]

Source code in pika/connection.py
def _send_message(self, channel_number: int,
                  method_frame: amqp_object.Method,
                  content: tuple[spec.BasicProperties, bytes]) -> None:
    """
    Publish a message.

    :param channel_number: The channel number for the frame
    :param method_frame: The method frame to send
    :param content: A content frame, which is tuple of properties and body.
    """
    length = len(content[1])
    marshaled_body_frames = []

    # Note: we construct the Method, Header and Content objects, marshal them
    # *then* output in case the marshaling operation throws an exception
    frame_method = frame.Method(channel_number, method_frame)
    frame_header = frame.Header(channel_number, length, content[0])
    marshaled_body_frames.append(frame_method.marshal())
    marshaled_body_frames.append(frame_header.marshal())

    if content[1]:
        chunks = math.ceil(float(length) / self._body_max_length)
        for chunk in range(chunks):
            start = chunk * self._body_max_length
            end = start + self._body_max_length
            end = min(end, length)
            frame_body = frame.Body(channel_number, content[1][start:end])
            marshaled_body_frames.append(frame_body.marshal())

    self._output_marshaled_frames(marshaled_body_frames)

_send_method

_send_method(
    channel_number: int,
    method: Method,
    content: tuple[BasicProperties, bytes] | None = None,
) -> None

Constructs a RPC method frame and then sends it to the broker.

PARAMETER DESCRIPTION
channel_number

The channel number for the frame

TYPE: int

method

The method to send

TYPE: Method

content

If set, is a content frame, is tuple of properties and body.

TYPE: tuple[BasicProperties, bytes] | None DEFAULT: None

Source code in pika/connection.py
def _send_method(
        self,
        channel_number: int,
        method: amqp_object.Method,
        content: tuple[spec.BasicProperties, bytes] | None = None) -> None:
    """
    Constructs a RPC method frame and then sends it to the broker.

    :param channel_number: The channel number for the frame
    :param method: The method to send
    :param content: If set, is a content frame, is tuple of properties and body.
    """
    if content:
        self._send_message(channel_number, method, content)
    else:
        self._send_frame(frame.Method(channel_number, method))

_set_connection_state

_set_connection_state(connection_state: int) -> None

Set the connection state.

PARAMETER DESCRIPTION
connection_state

The connection state to set

TYPE: int

Source code in pika/connection.py
def _set_connection_state(self, connection_state: int) -> None:
    """
    Set the connection state.

    :param connection_state: The connection state to set
    """
    LOGGER.debug('New Connection state: %s (prev=%s)',
                 self._STATE_NAMES[connection_state],
                 self._STATE_NAMES[self.connection_state])

    self.connection_state = connection_state

_set_server_information

_set_server_information(
    method_frame: Method[Start],
) -> None

Set the server properties and capabilities.

PARAMETER DESCRIPTION
method_frame

The Connection.Start frame

TYPE: Method[Start]

Source code in pika/connection.py
def _set_server_information(
        self, method_frame: frame.Method[spec.Connection.Start]) -> None:
    """
    Set the server properties and capabilities.

    :param method_frame: The Connection.Start frame
    """
    self.server_properties = method_frame.method.server_properties  # pyright: ignore[reportAttributeAccessIssue]
    assert self.server_properties is not None
    # Note: 'capabilities' is deliberately kept in server_properties.
    # Code that removed it never worked (it used hasattr on a dict, dead
    # since 2011), so callers have always seen the key there; removing it
    # now would break them.
    self.server_capabilities = self.server_properties.get(
        'capabilities', {})

_start_connection_workflow classmethod

_start_connection_workflow(
    connection_configs: Sequence[Parameters],
    connection_factory: Callable[
        [Parameters], BaseConnection
    ],
    nbio: AbstractIOServices,
    workflow: AbstractAMQPConnectionWorkflow | None,
    on_done: Callable[
        [Connection | AMQPConnectorException], None
    ],
) -> AbstractAMQPConnectionWorkflow

Helper function for custom implementations of create_connection().

PARAMETER DESCRIPTION
connection_configs

A sequence of one or more pika.connection.Parameters-based objects.

TYPE: Sequence[Parameters]

connection_factory

A function that takes pika.connection.Parameters as its only arg and returns a brand new pika.connection.Connection-based adapter instance each time it is called. The factory must instantiate the connection with internal_connection_workflow=False.

TYPE: Callable[[Parameters], BaseConnection]

nbio

TYPE: AbstractIOServices

workflow

Pass an instance of an implementation of the connection_workflow.AbstractAMQPConnectionWorkflow interface; defaults to a connection_workflow.AMQPConnectionWorkflow instance with default values for optional args. Type: connection_workflow.AbstractAMQPConnectionWorkflow | None.

TYPE: AbstractAMQPConnectionWorkflow | None

on_done

as defined in :py:meth:connection_workflow.AbstractAMQPConnectionWorkflow.start().

TYPE: Callable[[Connection | AMQPConnectorException], None]

RETURNS DESCRIPTION
AbstractAMQPConnectionWorkflow

Connection workflow instance in use. The user should limit their interaction with this object only to it's abort() method.

Source code in pika/adapters/base_connection.py
@classmethod
def _start_connection_workflow(
    cls, connection_configs: Sequence[connection.Parameters],
    connection_factory: Callable[[connection.Parameters], BaseConnection],
    nbio: nbio_interface.AbstractIOServices,
    workflow: connection_workflow.AbstractAMQPConnectionWorkflow | None,
    on_done: Callable[[(connection.Connection |
                        connection_workflow.AMQPConnectorException)], None]
) -> connection_workflow.AbstractAMQPConnectionWorkflow:
    """
    Helper function for custom implementations of `create_connection()`.

    :param connection_configs: A sequence of one or more
        `pika.connection.Parameters`-based objects.
    :param connection_factory: A function that takes
        `pika.connection.Parameters` as its only arg and returns a brand new
        `pika.connection.Connection`-based adapter instance each time it is
        called. The factory must instantiate the connection with
        `internal_connection_workflow=False`.
    :param nbio:
    :param workflow:
        Pass an instance of an implementation of the
        `connection_workflow.AbstractAMQPConnectionWorkflow` interface;
        defaults to a `connection_workflow.AMQPConnectionWorkflow` instance
        with default values for optional args. Type:
        `connection_workflow.AbstractAMQPConnectionWorkflow | None`.
    :param on_done: as defined in
        :py:meth:`connection_workflow.AbstractAMQPConnectionWorkflow.start()`.
    :returns: Connection workflow instance in use. The user should limit
        their interaction with this object only to it's `abort()` method.
    """
    if workflow is None:
        workflow = connection_workflow.AMQPConnectionWorkflow()
        LOGGER.debug('Created default connection workflow %r', workflow)

    if isinstance(workflow, connection_workflow.AMQPConnectionWorkflow):
        workflow.set_io_services(nbio)

    def create_connector() -> connection_workflow.AMQPConnector:
        """`AMQPConnector` factory."""
        return connection_workflow.AMQPConnector(
            lambda params: _StreamingProtocolShim(connection_factory(params)
                                                 ), nbio)

    workflow.start(connection_configs=connection_configs,
                   connector_factory=create_connector,
                   native_loop=nbio.get_native_ioloop(),
                   on_done=functools.partial(
                       cls._unshim_connection_workflow_callback,
                       cast(Any, on_done)))

    return workflow

_terminate_stream

_terminate_stream(error: Exception | None) -> None

Deactivate heartbeat instance if activated already, and initiate termination of the stream (TCP) connection asynchronously.

When connection terminates, the appropriate user callback will be invoked with the given error: "on open error" or "on connection closed".

PARAMETER DESCRIPTION
error

Exception (or None) describing the reason for termination; None for normal closing, such as upon receipt of Connection.CloseOk.

TYPE: Exception | None

Source code in pika/connection.py
def _terminate_stream(self, error: Exception | None) -> None:
    """
    Deactivate heartbeat instance if activated already, and initiate termination of the stream
    (TCP) connection asynchronously.

    When connection terminates, the appropriate user callback will be
    invoked with the given error: "on open error" or "on connection closed".

    :param error: Exception (or None) describing the reason
        for termination; None for normal closing, such as upon receipt of
        Connection.CloseOk.
    """
    assert isinstance(error, (type(None), Exception)), (
        f'error arg is neither None nor instance of Exception: {error!r}.')

    if error is not None:
        # Save the exception for user callback once the stream closes
        self._error = error
    else:
        assert self._error is not None, (
            '_terminate_stream() expected self._error to be set when '
            'passed None error arg.')

    # So it won't mess with the stack
    self._remove_heartbeat()

    # Begin disconnection of stream or termination of connection workflow
    self._adapter_disconnect_stream()

_tune_heartbeat_timeout staticmethod

_tune_heartbeat_timeout(
    client_value: int | None, server_value: int
) -> int

Determine heartbeat timeout per AMQP 0-9-1 rules.

Per https://www.rabbitmq.com/resources/specs/amqp0-9-1.pdf,

Both peers negotiate the limits to the lowest agreed value as follows: - The server MUST tell the client what limits it proposes. - The client responds and MAY reduce those limits for its connection

If the client specifies a value, it always takes precedence.

PARAMETER DESCRIPTION
client_value

None to accept server_value; otherwise, an integral number in seconds; 0 (zero) to disable heartbeat.

TYPE: int | None

server_value

integral value of the heartbeat timeout proposed by broker; 0 (zero) to disable heartbeat.

TYPE: int

RETURNS DESCRIPTION
int

the value of the heartbeat timeout to use and return to broker

Source code in pika/connection.py
@staticmethod
def _tune_heartbeat_timeout(client_value: int | None,
                            server_value: int) -> int:
    """
    Determine heartbeat timeout per AMQP 0-9-1 rules.

    Per https://www.rabbitmq.com/resources/specs/amqp0-9-1.pdf,

    > Both peers negotiate the limits to the lowest agreed value as follows:
    > - The server MUST tell the client what limits it proposes.
    > - The client responds and **MAY reduce those limits** for its
        connection

    If the client specifies a value, it always takes precedence.

    :param client_value: None to accept server_value; otherwise, an integral
        number in seconds; 0 (zero) to disable heartbeat.
    :param server_value: integral value of the heartbeat timeout proposed by
        broker; 0 (zero) to disable heartbeat.

    :returns: the value of the heartbeat timeout to use and return to broker
    """
    if client_value is None:
        # Accept server's limit
        timeout = server_value
    else:
        timeout = client_value

    return timeout

_unshim_connection_workflow_callback staticmethod

_unshim_connection_workflow_callback(
    user_on_done: Callable[
        [BaseConnection | BaseException], None
    ],
    shim_or_exc: _StreamingProtocolShim | Exception,
) -> None
PARAMETER DESCRIPTION
user_on_done

user's on_done callback as defined in :py:meth:connection_workflow.AbstractAMQPConnectionWorkflow.start().

TYPE: Callable[[BaseConnection | BaseException], None]

shim_or_exc

_StreamingProtocolShim | Exception

TYPE: _StreamingProtocolShim | Exception

Source code in pika/adapters/base_connection.py
@staticmethod
def _unshim_connection_workflow_callback(
        user_on_done: Callable[[BaseConnection | BaseException], None],
        shim_or_exc: _StreamingProtocolShim | Exception) -> None:
    """

    :param user_on_done: user's `on_done` callback as defined in
        :py:meth:`connection_workflow.AbstractAMQPConnectionWorkflow.start()`.
    :param shim_or_exc: `_StreamingProtocolShim | Exception`
    """
    result: BaseConnection | BaseException
    if isinstance(shim_or_exc, _StreamingProtocolShim):
        result = shim_or_exc.conn
    else:
        result = shim_or_exc

    user_on_done(result)

add_on_close_callback

add_on_close_callback(
    callback: Callable[[Connection, Exception], Any],
) -> None

Add a callback notification when the connection has closed.

The callback will be passed the connection and an exception instance. The exception will either be an instance of exceptions.ConnectionClosed if a fully-open connection was closed by user or broker or exception of another type that describes the cause of connection closure/failure. If the stream died because some other exception was raised - for example from one of your own callbacks - most adapters chain that exception onto the reason as its __cause__, so inspecting reason.__cause__ reveals where it came from.

PARAMETER DESCRIPTION
callback

Callback to call on close, having the signature: callback(pika.connection.Connection, exception)

TYPE: Callable[[Connection, Exception], Any]

Source code in pika/connection.py
def add_on_close_callback(
        self, callback: Callable[[Connection, Exception], Any]) -> None:
    """
    Add a callback notification when the connection has closed.

    The callback will be passed the connection and an exception instance. The exception will
    either be an instance of `exceptions.ConnectionClosed` if a fully-open connection was closed
    by user or broker or exception of another type that describes the cause of connection
    closure/failure. If the stream died because some other exception was raised - for example
    from one of your own callbacks - most adapters chain that exception onto the reason as its
    `__cause__`, so inspecting `reason.__cause__` reveals where it came from.

    :param callback: Callback to call on close, having the signature:
        callback(pika.connection.Connection, exception)
    """
    validators.require_callback(callback)
    self.callbacks.add(0, self.ON_CONNECTION_CLOSED, callback, False)

add_on_connection_blocked_callback

add_on_connection_blocked_callback(
    callback: Callable[[Connection, Method[Blocked]], Any],
) -> None

RabbitMQ AMQP extension - Add a callback to be notified when the connection gets blocked (Connection.Blocked received from RabbitMQ) due to the broker running low on resources (memory or disk). In this state RabbitMQ suspends processing incoming data until the connection is unblocked, so it's a good idea for publishers receiving this notification to suspend publishing until the connection becomes unblocked.

See also Connection.add_on_connection_unblocked_callback()

See also ConnectionParameters.blocked_connection_timeout.

PARAMETER DESCRIPTION
callback

Callback to call on Connection.Blocked, having the signature callback(connection, pika.frame.Method), where the method frame's method member is of type pika.spec.Connection.Blocked

TYPE: Callable[[Connection, Method[Blocked]], Any]

Source code in pika/connection.py
def add_on_connection_blocked_callback(
    self,
    callback: Callable[[Connection, frame.Method[spec.Connection.Blocked]],
                       Any]
) -> None:
    """RabbitMQ AMQP extension - Add a callback to be notified when the
    connection gets blocked (`Connection.Blocked` received from RabbitMQ)
    due to the broker running low on resources (memory or disk). In this
    state RabbitMQ suspends processing incoming data until the connection
    is unblocked, so it's a good idea for publishers receiving this
    notification to suspend publishing until the connection becomes
    unblocked.

    See also `Connection.add_on_connection_unblocked_callback()`

    See also `ConnectionParameters.blocked_connection_timeout`.

    :param callback: Callback to call on `Connection.Blocked`,
        having the signature `callback(connection, pika.frame.Method)`,
        where the method frame's `method` member is of type
        `pika.spec.Connection.Blocked`

    """
    validators.require_callback(callback)
    self.callbacks.add(0,
                       spec.Connection.Blocked,
                       functools.partial(callback, self),
                       one_shot=False)

add_on_connection_unblocked_callback

add_on_connection_unblocked_callback(
    callback: Callable[
        [Connection, Method[Unblocked]], Any
    ],
) -> None

RabbitMQ AMQP extension - Add a callback to be notified when the connection gets unblocked (Connection.Unblocked frame is received from RabbitMQ) letting publishers know it's ok to start publishing again.

PARAMETER DESCRIPTION
callback

Callback to call on Connection.Unblocked, having the signature callback(connection, pika.frame.Method), where the method frame's method member is of type pika.spec.Connection.Unblocked

TYPE: Callable[[Connection, Method[Unblocked]], Any]

Source code in pika/connection.py
def add_on_connection_unblocked_callback(
    self, callback: Callable[
        [Connection, frame.Method[spec.Connection.Unblocked]], Any]
) -> None:
    """RabbitMQ AMQP extension - Add a callback to be notified when the
    connection gets unblocked (`Connection.Unblocked` frame is received from
    RabbitMQ) letting publishers know it's ok to start publishing again.

    :param callback: Callback to call on
        `Connection.Unblocked`, having the signature
        `callback(connection, pika.frame.Method)`, where the method frame's
        `method` member is of type `pika.spec.Connection.Unblocked`

    """
    validators.require_callback(callback)
    self.callbacks.add(0,
                       spec.Connection.Unblocked,
                       functools.partial(callback, self),
                       one_shot=False)

add_on_open_callback

add_on_open_callback(
    callback: Callable[[Connection], Any],
) -> None

Add a callback notification when the connection has opened.

The callback will be passed the connection instance as its only arg.

PARAMETER DESCRIPTION
callback

Callback to call when open

TYPE: Callable[[Connection], Any]

Source code in pika/connection.py
def add_on_open_callback(self, callback: Callable[[Connection],
                                                  Any]) -> None:
    """
    Add a callback notification when the connection has opened.

    The callback will be passed the connection instance as its only arg.

    :param callback: Callback to call when open
    """
    validators.require_callback(callback)
    self.callbacks.add(0, self.ON_CONNECTION_OPEN_OK, callback, False)

add_on_open_error_callback

add_on_open_error_callback(
    callback: Callable[
        [Connection, BaseException | None], Any
    ],
    remove_default: bool = True,
) -> None

Add a callback notification when the connection can not be opened.

The callback method should accept the connection instance that could not connect, and either a string or an exception as its second arg.

PARAMETER DESCRIPTION
callback

Callback to call when can't connect, having the signature _(Connection, Exception)

TYPE: Callable[[Connection, BaseException | None], Any]

remove_default

Remove default exception raising callback

TYPE: bool DEFAULT: True

Source code in pika/connection.py
def add_on_open_error_callback(self,
                               callback: Callable[
                                   [Connection, BaseException | None], Any],
                               remove_default: bool = True) -> None:
    """
    Add a callback notification when the connection can not be opened.

    The callback method should accept the connection instance that could not connect, and either
    a string or an exception as its second arg.

    :param callback: Callback to call when can't connect, having the signature _(Connection,
        Exception)
    :param remove_default: Remove default exception raising callback
    """
    validators.require_callback(callback)
    if remove_default:
        self.callbacks.remove(0, self.ON_CONNECTION_ERROR,
                              self._default_on_connection_error)
    self.callbacks.add(0, self.ON_CONNECTION_ERROR, callback, False)

channel

channel(
    channel_number: int | None = None,
    on_open_callback: Callable[[Channel], Any]
    | None = None,
) -> Channel

Create a new channel with the next available channel number or pass in a channel number to use. Must be non-zero if you would like to specify but it is recommended that you let Pika manage the channel numbers.

PARAMETER DESCRIPTION
channel_number

The channel number to use, defaults to the next available.

TYPE: int | None DEFAULT: None

on_open_callback

The callback when the channel is opened. The callback will be invoked with the Channel instance as its only argument.

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

Source code in pika/connection.py
def channel(
        self,
        channel_number: int | None = None,
        on_open_callback: Callable[[Channel], Any] | None = None
) -> Channel:
    """
    Create a new channel with the next available channel number or pass in a channel number to
    use. Must be non-zero if you would like to specify but it is recommended that you let Pika
    manage the channel numbers.

    :param channel_number: The channel number to use, defaults to the next available.
    :param on_open_callback: The callback when the channel is opened. The callback will be
        invoked with the `Channel` instance as its only argument.
    """
    if not self.is_open:
        raise exceptions.ConnectionWrongStateError(
            f'Channel allocation requires an open connection: {self}')

    validators.rpc_completion_callback(on_open_callback)

    if not channel_number:
        channel_number = self._next_channel_number()

    self._channels[channel_number] = self._create_channel(
        channel_number, on_open_callback)
    self._add_channel_callbacks(channel_number)
    self._channels[channel_number].open()
    return self._channels[channel_number]

close

close(
    reply_code: int = 200,
    reply_text: str = 'Normal shutdown',
) -> None

Disconnect from RabbitMQ.

If there are any open channels, it will attempt to close them prior to fully disconnecting. Channels which have active consumers will attempt to send a Basic.Cancel to RabbitMQ to cleanly stop the delivery of messages prior to closing the channel.

PARAMETER DESCRIPTION
reply_code

The code number for the close

TYPE: int DEFAULT: 200

reply_text

The text reason for the close

TYPE: str DEFAULT: 'Normal shutdown'

RAISES DESCRIPTION
pika.exceptions.ConnectionWrongStateError

if connection is closed or closing.

Source code in pika/connection.py
def close(self,
          reply_code: int = 200,
          reply_text: str = 'Normal shutdown') -> None:
    """
    Disconnect from RabbitMQ.

    If there are any open channels, it will attempt to close them prior to fully disconnecting.
    Channels which have active consumers will attempt to send a Basic.Cancel to RabbitMQ to
    cleanly stop the delivery of messages prior to closing the channel.

    :param reply_code: The code number for the close
    :param reply_text: The text reason for the close
    :raises pika.exceptions.ConnectionWrongStateError: if connection is closed or closing.
    """
    if self.is_closing or self.is_closed:
        msg = (
            f'Illegal close({reply_code}, {reply_text!r}) request on {self} because it '
            f'was called while connection state={self._STATE_NAMES[self.connection_state]}.'
        )
        LOGGER.error(msg)
        raise exceptions.ConnectionWrongStateError(msg)

    # NOTE The connection is either in opening or open state

    # Initiate graceful closing of channels that are OPEN or OPENING
    if self._channels:
        self._close_channels(reply_code, reply_text)

    prev_state = self.connection_state

    # Transition to closing
    self._set_connection_state(self.CONNECTION_CLOSING)
    LOGGER.info('Closing connection (%s): %r', reply_code, reply_text)

    if not self._opened:
        # It was opening, but not fully open yet, so we won't attempt
        # graceful AMQP Connection.Close.
        LOGGER.info('Connection.close() is terminating stream and '
                    'bypassing graceful AMQP close, since AMQP is still '
                    'opening.')

        error = exceptions.ConnectionOpenAborted(
            'Connection.close() called before connection '
            f'finished opening: prev_state={self._STATE_NAMES[prev_state]} ({reply_code}): {reply_text!r}',
            host=self.params.host,
            port=self.params.port)
        self._terminate_stream(error)

    else:
        self._error = exceptions.ConnectionClosedByClient(
            reply_code,
            reply_text,
            host=self.params.host,
            port=self.params.port)

        # If there are channels that haven't finished closing yet, then
        # _on_close_ready will finally be called from _on_channel_cleanup once
        # all channels have been closed
        if not self._channels:
            # We can initiate graceful closing of the connection right away,
            # since no more channels remain
            self._on_close_ready()
        else:
            LOGGER.info(
                'Connection.close is waiting for %d channels to close: %s',
                len(self._channels), self)

create_connection classmethod

create_connection(
    connection_configs: Sequence[Parameters],
    on_done: Callable[
        [Connection | AMQPConnectorException], None
    ],
    custom_ioloop: Any | None = None,
    workflow: AbstractAMQPConnectionWorkflow | None = None,
) -> AbstractAMQPConnectionWorkflow
PARAMETER DESCRIPTION
connection_configs

One or more connection parameter objects

TYPE: Sequence[Parameters]

on_done

Callback to report when connection workflow is done

TYPE: Callable[[Connection | AMQPConnectorException], None]

custom_ioloop

Optional custom IOLoop or nbio interface to use for the connection workflow

TYPE: Any | None DEFAULT: None

workflow

Optional connection workflow instance to use; if None, a default workflow will be created

TYPE: AbstractAMQPConnectionWorkflow | None DEFAULT: None

Source code in pika/adapters/select_connection.py
@classmethod
@override
def create_connection(
    cls,
    connection_configs: Sequence[connection.Parameters],
    on_done: Callable[[(connection.Connection |
                        connection_workflow.AMQPConnectorException)], None],
    custom_ioloop: Any | None = None,
    workflow: connection_workflow.AbstractAMQPConnectionWorkflow |
    None = None
) -> connection_workflow.AbstractAMQPConnectionWorkflow:
    """
    :param connection_configs: One or more connection parameter objects
    :param on_done: Callback to report when connection workflow is done
    :param custom_ioloop: Optional custom IOLoop or nbio interface to use
        for the connection workflow
    :param workflow: Optional connection workflow instance to use; if None,
        a default workflow will be created
    """
    nbio = SelectorIOServicesAdapter(custom_ioloop or IOLoop())

    def connection_factory(params) -> SelectConnection:
        """
        Connection factory.

        :param params: Connection parameters
        """
        if params is None:
            raise ValueError('Expected pika.connection.Parameters '
                             'instance, but got None in params arg.')
        return cls(parameters=params,
                   custom_ioloop=nbio,
                   internal_connection_workflow=False)

    return cls._start_connection_workflow(
        connection_configs=connection_configs,
        connection_factory=connection_factory,
        nbio=nbio,
        workflow=workflow,
        on_done=on_done)

remove_on_close_callback

remove_on_close_callback(
    callback: Callable[[Connection, Exception], Any],
) -> bool

Remove a callback registered via :meth:add_on_close_callback.

PARAMETER DESCRIPTION
callback

The callback to remove

TYPE: Callable[[Connection, Exception], Any]

RETURNS DESCRIPTION
bool

True if a callback was removed, otherwise False

Source code in pika/connection.py
def remove_on_close_callback(
        self, callback: Callable[[Connection, Exception], Any]) -> bool:
    """
    Remove a callback registered via :meth:`add_on_close_callback`.

    :param callback: The callback to remove
    :returns: True if a callback was removed, otherwise False
    """
    return self.callbacks.remove_matching(0, self.ON_CONNECTION_CLOSED,
                                          lambda cb: cb == callback)

remove_on_connection_blocked_callback

remove_on_connection_blocked_callback(
    callback: Callable[[Connection, Method[Blocked]], Any],
) -> bool

Remove a callback registered via :meth:add_on_connection_blocked_callback.

The callback was wrapped in a functools.partial at registration time, so it is matched here by the original callable rather than by equality.

PARAMETER DESCRIPTION
callback

The callback to remove

TYPE: Callable[[Connection, Method[Blocked]], Any]

RETURNS DESCRIPTION
bool

True if a callback was removed, otherwise False

Source code in pika/connection.py
def remove_on_connection_blocked_callback(
    self,
    callback: Callable[[Connection, frame.Method[spec.Connection.Blocked]],
                       Any]
) -> bool:
    """Remove a callback registered via
    :meth:`add_on_connection_blocked_callback`.

    The callback was wrapped in a ``functools.partial`` at registration time,
    so it is matched here by the original callable rather than by equality.

    :param callback: The callback to remove
    :returns: True if a callback was removed, otherwise False
    """

    def matches(stored: Callable[..., Any]) -> bool:
        return (isinstance(stored, functools.partial) and
                stored.func == callback and stored.args == (self,))

    return self.callbacks.remove_matching(0, spec.Connection.Blocked,
                                          matches)

remove_on_connection_unblocked_callback

remove_on_connection_unblocked_callback(
    callback: Callable[
        [Connection, Method[Unblocked]], Any
    ],
) -> bool

Remove a callback registered via :meth:add_on_connection_unblocked_callback.

The callback was wrapped in a functools.partial at registration time, so it is matched here by the original callable rather than by equality.

PARAMETER DESCRIPTION
callback

The callback to remove

TYPE: Callable[[Connection, Method[Unblocked]], Any]

RETURNS DESCRIPTION
bool

True if a callback was removed, otherwise False

Source code in pika/connection.py
def remove_on_connection_unblocked_callback(
    self, callback: Callable[
        [Connection, frame.Method[spec.Connection.Unblocked]], Any]
) -> bool:
    """Remove a callback registered via
    :meth:`add_on_connection_unblocked_callback`.

    The callback was wrapped in a ``functools.partial`` at registration time,
    so it is matched here by the original callable rather than by equality.

    :param callback: The callback to remove
    :returns: True if a callback was removed, otherwise False
    """

    def matches(stored: Callable[..., Any]) -> bool:
        return (isinstance(stored, functools.partial) and
                stored.func == callback and stored.args == (self,))

    return self.callbacks.remove_matching(0, spec.Connection.Unblocked,
                                          matches)

remove_on_open_callback

remove_on_open_callback(
    callback: Callable[[Connection], Any],
) -> bool

Remove a callback registered via :meth:add_on_open_callback.

PARAMETER DESCRIPTION
callback

The callback to remove

TYPE: Callable[[Connection], Any]

RETURNS DESCRIPTION
bool

True if a callback was removed, otherwise False

Source code in pika/connection.py
def remove_on_open_callback(self, callback: Callable[[Connection],
                                                     Any]) -> bool:
    """
    Remove a callback registered via :meth:`add_on_open_callback`.

    :param callback: The callback to remove
    :returns: True if a callback was removed, otherwise False
    """
    return self.callbacks.remove_matching(0, self.ON_CONNECTION_OPEN_OK,
                                          lambda cb: cb == callback)

remove_on_open_error_callback

remove_on_open_error_callback(
    callback: Callable[
        [Connection, BaseException | None], Any
    ],
) -> bool

Remove a callback registered via :meth:add_on_open_error_callback.

PARAMETER DESCRIPTION
callback

The callback to remove

TYPE: Callable[[Connection, BaseException | None], Any]

RETURNS DESCRIPTION
bool

True if a callback was removed, otherwise False

Source code in pika/connection.py
def remove_on_open_error_callback(
        self, callback: Callable[[Connection, BaseException | None],
                                 Any]) -> bool:
    """
    Remove a callback registered via :meth:`add_on_open_error_callback`.

    :param callback: The callback to remove
    :returns: True if a callback was removed, otherwise False
    """
    return self.callbacks.remove_matching(0, self.ON_CONNECTION_ERROR,
                                          lambda cb: cb == callback)

update_secret

update_secret(
    new_secret: str,
    reason: str,
    callback: Callable[[Method[UpdateSecretOk]], Any]
    | None = None,
) -> None

RabbitMQ AMQP extension - This method updates the secret used to authenticate this connection. It is used when secrets have an expiration date and need to be renewed, like OAuth 2 tokens. Pass a callback to be notified of the response from the server.

PARAMETER DESCRIPTION
new_secret

The new secret

TYPE: str

reason

The reason for the secret update

TYPE: str

callback

Callback to call on Connection.UpdateSecretOk, having the signature callback(pika.frame.Method), where the method frame's method member is of type pika.spec.Connection.UpdateSecretOk

TYPE: Callable[[Method[UpdateSecretOk]], Any] | None DEFAULT: None

RAISES DESCRIPTION
pika.exceptions.ConnectionWrongStateError

if connection is not open.

Source code in pika/connection.py
def update_secret(
    self,
    new_secret: str,
    reason: str,
    callback: Callable[[frame.Method[spec.Connection.UpdateSecretOk]],
                       Any] | None = None
) -> None:
    """RabbitMQ AMQP extension - This method updates the secret used to authenticate this connection.
    It is used when secrets have an expiration date and need to be renewed, like OAuth 2 tokens.
    Pass a callback to be notified of the response from the server.

    :param new_secret: The new secret
    :param reason: The reason for the secret update
    :param callback: Callback to call on
        `Connection.UpdateSecretOk`, having the signature
        `callback(pika.frame.Method)`, where the method frame's
        `method` member is of type `pika.spec.Connection.UpdateSecretOk`

    :raises pika.exceptions.ConnectionWrongStateError: if connection is
        not open.
    """
    if not self.is_open:
        raise exceptions.ConnectionWrongStateError(
            f'Secret update requires an open connection: {self}')

    validators.rpc_completion_callback(callback)
    self._rpc(0, spec.Connection.UpdateSecret(new_secret, reason), callback,
              [spec.Connection.UpdateSecretOk])