Gevent Connection Adapter¶
[!WARNING]
GeventConnectionis deprecated and will be removed in Pika 2.0. UseConnectioninstead, 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:
|
on_open_callback
|
The method to call when the connection is open
TYPE:
|
on_open_error_callback
|
Called if the connection
can't be established or connection establishment is interrupted by
TYPE:
|
on_close_callback
|
Called when a previously fully
open connection is closed:
TYPE:
|
custom_ioloop
|
Use a custom Gevent ILoop
TYPE:
|
internal_connection_workflow
|
True for autonomous connection
establishment which is default; False for externally-managed
connection workflow via the
TYPE:
|
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:
|
on_done
|
Callback to report when connection workflow is done
TYPE:
|
custom_ioloop
|
Optional custom Gevent ILoop to use for the connection workflow; if None, a new _GeventSelectorIOLoop will be created
TYPE:
|
workflow
|
Optional connection workflow instance to use
TYPE:
|
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)
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',
}
_body_max_length
instance-attribute
¶
_body_max_length: int = (
spec.FRAME_MAX_SIZE
- spec.FRAME_HEADER_SIZE
- spec.FRAME_END_SIZE
)
_client_properties
property
¶
_client_properties: dict[str, Any]
Return the client properties dictionary.
_connection_workflow
instance-attribute
¶
_connection_workflow: (
AbstractAMQPConnectionWorkflow | None
) = None
_internal_connection_workflow
instance-attribute
¶
_internal_connection_workflow = internal_connection_workflow
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.
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, |
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.
publisher_confirms
property
¶
publisher_confirms: bool
Specifies if the active connection can use publisher confirmations.
__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:
|
on_open_callback
|
The method to call when the connection is open
TYPE:
|
on_open_error_callback
|
Called if the connection
can't be established or connection establishment is interrupted by
TYPE:
|
on_close_callback
|
Called when a previously fully
open connection is closed:
TYPE:
|
custom_ioloop
|
Use a custom Gevent ILoop
TYPE:
|
internal_connection_workflow
|
True for autonomous connection
establishment which is default; False for externally-managed
connection workflow via the
TYPE:
|
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:
|
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:
|
callback
|
Callback to call.
TYPE:
|
| 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:
|
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:
|
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:
|
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:
|
| 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:
|
reply_text
|
The text reason for why the channels are closing
TYPE:
|
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:
|
on_open_callback
|
The callback when the channel is opened. The callback will be
invoked with the
TYPE:
|
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:
|
error
|
The exception that caused the failure
TYPE:
|
| 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:
|
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:
|
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:
|
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:
|
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:
|
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:
|
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:
|
server_value
|
The server value
TYPE:
|
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:
|
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:
|
method_frame
|
method frame having
TYPE:
|
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:
|
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:
|
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:
|
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:
|
| 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:
|
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:
|
method_frame
|
method frame having
TYPE:
|
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
TYPE:
|
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:
|
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
TYPE:
|
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:
|
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:
|
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:
|
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
TYPE:
|
| 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:
|
| 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:
|
| 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 |
| 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:
|
method_classes
|
The method classes (derived from
TYPE:
|
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:
|
method
|
The method frame to call
TYPE:
|
callback
|
The callback for the RPC response
TYPE:
|
acceptable_replies
|
The replies this RPC call expects
TYPE:
|
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:
|
reply_text
|
The text reason for the close
TYPE:
|
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:
|
response
|
The encoded value to send
TYPE:
|
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:
|
| 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:
|
method_frame
|
The method frame to send
TYPE:
|
content
|
A content frame, which is tuple of properties and body.
TYPE:
|
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:
|
method
|
The method to send
TYPE:
|
content
|
If set, is a content frame, is tuple of properties and body.
TYPE:
|
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:
|
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:
|
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
TYPE:
|
connection_factory
|
A function that takes
TYPE:
|
nbio
|
TYPE:
|
workflow
|
Pass an instance of an implementation of the
TYPE:
|
on_done
|
as defined in
:py:meth:
TYPE:
|
| RETURNS | DESCRIPTION |
|---|---|
AbstractAMQPConnectionWorkflow
|
Connection workflow instance in use. The user should limit
their interaction with this object only to it's |
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:
|
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:
|
server_value
|
integral value of the heartbeat timeout proposed by broker; 0 (zero) to disable heartbeat.
TYPE:
|
| 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
TYPE:
|
shim_or_exc
|
TYPE:
|
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:
|
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
TYPE:
|
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
TYPE:
|
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:
|
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:
|
remove_default
|
Remove default exception raising callback
TYPE:
|
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:
|
on_open_callback
|
The callback when the channel is opened. The callback will be
invoked with the
TYPE:
|
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:
|
reply_text
|
The text reason for the close
TYPE:
|
| 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:
|
on_done
|
Callback to report when connection workflow is done
TYPE:
|
custom_ioloop
|
Optional custom Gevent ILoop to use for the connection workflow; if None, a new _GeventSelectorIOLoop will be created
TYPE:
|
workflow
|
Optional connection workflow instance to use
TYPE:
|
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:
|
| 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:
|
| 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:
|
| 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:
|
| 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:
|
| 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:
|
reason
|
The reason for the secret update
TYPE:
|
callback
|
Callback to call on
TYPE:
|
| 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])