Skip to content

Gevent Connection Adapter

[!WARNING] GeventConnection is deprecated and will be removed in Pika 2.0. Use Connection instead, which runs its own IOLoop on a background thread and works with any framework, including Gevent.

gevent_connection

Use pika with the Gevent IOLoop.

GeventConnection

Bases: BaseConnection

Implementation of pika's BaseConnection.

An async selector-based connection which integrates with Gevent.

Source code in pika/adapters/gevent_connection.py
class GeventConnection(BaseConnection):
    """
    Implementation of pika's ``BaseConnection``.

    An async selector-based connection which integrates with Gevent.
    """

    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: gevent._interfaces.ILoop | AbstractIOServices |
        None = None,
            internal_connection_workflow: bool = True) -> None:
        """
        Create a new GeventConnection instance and connect to RabbitMQ on Gevent's event-loop.

        :param parameters: The connection
            parameters
        :param on_open_callback: The method to call when the
            connection is open
        :param on_open_error_callback: 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: 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: Use a custom Gevent ILoop
        :param internal_connection_workflow: True for autonomous connection
            establishment which is default; False for externally-managed
            connection workflow via the `create_connection()` factory
        """
        warnings.warn(
            'GeventConnection is deprecated and will be removed in Pika 2.0. '
            'Use Connection instead, which works with any framework '
            'including Gevent. See '
            'https://pika.github.io/pika/latest/modules/adapters/thread_safe/',
            DeprecationWarning,
            stacklevel=2,
        )
        if pika._utils.ON_WINDOWS:
            raise RuntimeError('GeventConnection is not supported on Windows.')

        custom_ioloop = (custom_ioloop or
                         _GeventSelectorIOLoop(gevent.get_hub()))

        if isinstance(custom_ioloop, AbstractIOServices):
            nbio = custom_ioloop
        else:
            nbio = _GeventSelectorIOServicesAdapter(custom_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: gevent._interfaces.ILoop | 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 Gevent ILoop to use for the
            connection workflow; if None, a new _GeventSelectorIOLoop will be
            created
        :param workflow: Optional connection workflow instance to use
        """
        custom_ioloop = (custom_ioloop or
                         _GeventSelectorIOLoop(gevent.get_hub()))

        nbio = _GeventSelectorIOServicesAdapter(custom_ioloop)

        def connection_factory(params) -> GeventConnection:
            """
            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)

__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: ILoop | AbstractIOServices | None = None,
    internal_connection_workflow: bool = True,
) -> None

Create a new GeventConnection instance and connect to RabbitMQ on Gevent's event-loop.

PARAMETER DESCRIPTION
parameters

The connection parameters

TYPE: Parameters | None DEFAULT: None

on_open_callback

The method to call when the connection is open

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

on_open_error_callback

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

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

Use a custom Gevent ILoop

TYPE: ILoop | AbstractIOServices | 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

Source code in pika/adapters/gevent_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: gevent._interfaces.ILoop | AbstractIOServices |
    None = None,
        internal_connection_workflow: bool = True) -> None:
    """
    Create a new GeventConnection instance and connect to RabbitMQ on Gevent's event-loop.

    :param parameters: The connection
        parameters
    :param on_open_callback: The method to call when the
        connection is open
    :param on_open_error_callback: 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: 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: Use a custom Gevent ILoop
    :param internal_connection_workflow: True for autonomous connection
        establishment which is default; False for externally-managed
        connection workflow via the `create_connection()` factory
    """
    warnings.warn(
        'GeventConnection is deprecated and will be removed in Pika 2.0. '
        'Use Connection instead, which works with any framework '
        'including Gevent. See '
        'https://pika.github.io/pika/latest/modules/adapters/thread_safe/',
        DeprecationWarning,
        stacklevel=2,
    )
    if pika._utils.ON_WINDOWS:
        raise RuntimeError('GeventConnection is not supported on Windows.')

    custom_ioloop = (custom_ioloop or
                     _GeventSelectorIOLoop(gevent.get_hub()))

    if isinstance(custom_ioloop, AbstractIOServices):
        nbio = custom_ioloop
    else:
        nbio = _GeventSelectorIOServicesAdapter(custom_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: ILoop | 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 Gevent ILoop to use for the connection workflow; if None, a new _GeventSelectorIOLoop will be created

TYPE: ILoop | None DEFAULT: None

workflow

Optional connection workflow instance to use

TYPE: AbstractAMQPConnectionWorkflow | None DEFAULT: None

Source code in pika/adapters/gevent_connection.py
@classmethod
@override
def create_connection(
    cls,
    connection_configs: Sequence[connection.Parameters],
    on_done: Callable[[(connection.Connection |
                        connection_workflow.AMQPConnectorException)], None],
    custom_ioloop: gevent._interfaces.ILoop | 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 Gevent ILoop to use for the
        connection workflow; if None, a new _GeventSelectorIOLoop will be
        created
    :param workflow: Optional connection workflow instance to use
    """
    custom_ioloop = (custom_ioloop or
                     _GeventSelectorIOLoop(gevent.get_hub()))

    nbio = _GeventSelectorIOServicesAdapter(custom_ioloop)

    def connection_factory(params) -> GeventConnection:
        """
        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)

Be sure to check out the asynchronous examples.

There is no specific Gevent example yet.

Class Reference

GeventConnection

Bases: BaseConnection

Implementation of pika's BaseConnection.

An async selector-based connection which integrates with Gevent.

Source code in pika/adapters/gevent_connection.py
class GeventConnection(BaseConnection):
    """
    Implementation of pika's ``BaseConnection``.

    An async selector-based connection which integrates with Gevent.
    """

    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: gevent._interfaces.ILoop | AbstractIOServices |
        None = None,
            internal_connection_workflow: bool = True) -> None:
        """
        Create a new GeventConnection instance and connect to RabbitMQ on Gevent's event-loop.

        :param parameters: The connection
            parameters
        :param on_open_callback: The method to call when the
            connection is open
        :param on_open_error_callback: 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: 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: Use a custom Gevent ILoop
        :param internal_connection_workflow: True for autonomous connection
            establishment which is default; False for externally-managed
            connection workflow via the `create_connection()` factory
        """
        warnings.warn(
            'GeventConnection is deprecated and will be removed in Pika 2.0. '
            'Use Connection instead, which works with any framework '
            'including Gevent. See '
            'https://pika.github.io/pika/latest/modules/adapters/thread_safe/',
            DeprecationWarning,
            stacklevel=2,
        )
        if pika._utils.ON_WINDOWS:
            raise RuntimeError('GeventConnection is not supported on Windows.')

        custom_ioloop = (custom_ioloop or
                         _GeventSelectorIOLoop(gevent.get_hub()))

        if isinstance(custom_ioloop, AbstractIOServices):
            nbio = custom_ioloop
        else:
            nbio = _GeventSelectorIOServicesAdapter(custom_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: gevent._interfaces.ILoop | 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 Gevent ILoop to use for the
            connection workflow; if None, a new _GeventSelectorIOLoop will be
            created
        :param workflow: Optional connection workflow instance to use
        """
        custom_ioloop = (custom_ioloop or
                         _GeventSelectorIOLoop(gevent.get_hub()))

        nbio = _GeventSelectorIOServicesAdapter(custom_ioloop)

        def connection_factory(params) -> GeventConnection:
            """
            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)

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: ILoop | AbstractIOServices | None = None,
    internal_connection_workflow: bool = True,
) -> None

Create a new GeventConnection instance and connect to RabbitMQ on Gevent's event-loop.

PARAMETER DESCRIPTION
parameters

The connection parameters

TYPE: Parameters | None DEFAULT: None

on_open_callback

The method to call when the connection is open

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

on_open_error_callback

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

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

Use a custom Gevent ILoop

TYPE: ILoop | AbstractIOServices | 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

Source code in pika/adapters/gevent_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: gevent._interfaces.ILoop | AbstractIOServices |
    None = None,
        internal_connection_workflow: bool = True) -> None:
    """
    Create a new GeventConnection instance and connect to RabbitMQ on Gevent's event-loop.

    :param parameters: The connection
        parameters
    :param on_open_callback: The method to call when the
        connection is open
    :param on_open_error_callback: 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: 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: Use a custom Gevent ILoop
    :param internal_connection_workflow: True for autonomous connection
        establishment which is default; False for externally-managed
        connection workflow via the `create_connection()` factory
    """
    warnings.warn(
        'GeventConnection is deprecated and will be removed in Pika 2.0. '
        'Use Connection instead, which works with any framework '
        'including Gevent. See '
        'https://pika.github.io/pika/latest/modules/adapters/thread_safe/',
        DeprecationWarning,
        stacklevel=2,
    )
    if pika._utils.ON_WINDOWS:
        raise RuntimeError('GeventConnection is not supported on Windows.')

    custom_ioloop = (custom_ioloop or
                     _GeventSelectorIOLoop(gevent.get_hub()))

    if isinstance(custom_ioloop, AbstractIOServices):
        nbio = custom_ioloop
    else:
        nbio = _GeventSelectorIOServicesAdapter(custom_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

_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: ILoop | 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 Gevent ILoop to use for the connection workflow; if None, a new _GeventSelectorIOLoop will be created

TYPE: ILoop | None DEFAULT: None

workflow

Optional connection workflow instance to use

TYPE: AbstractAMQPConnectionWorkflow | None DEFAULT: None

Source code in pika/adapters/gevent_connection.py
@classmethod
@override
def create_connection(
    cls,
    connection_configs: Sequence[connection.Parameters],
    on_done: Callable[[(connection.Connection |
                        connection_workflow.AMQPConnectorException)], None],
    custom_ioloop: gevent._interfaces.ILoop | 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 Gevent ILoop to use for the
        connection workflow; if None, a new _GeventSelectorIOLoop will be
        created
    :param workflow: Optional connection workflow instance to use
    """
    custom_ioloop = (custom_ioloop or
                     _GeventSelectorIOLoop(gevent.get_hub()))

    nbio = _GeventSelectorIOServicesAdapter(custom_ioloop)

    def connection_factory(params) -> GeventConnection:
        """
        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])