Connection Parameters¶
To maintain flexibility in how you specify the connection information required for your applications to properly connect to RabbitMQ, pika implements two classes for encapsulating the information, ConnectionParameters and URLParameters.
ConnectionParameters¶
The classic object for specifying all of the connection parameters required to connect to RabbitMQ, ConnectionParameters provides attributes for tweaking every possible connection option.
Example:
import pika
# Set the connection parameters to connect to rabbit-server1 on port 5672
# on the / virtual host using the username "guest" and password "guest"
credentials = pika.PlainCredentials('guest', 'guest')
parameters = pika.ConnectionParameters('rabbit-server1',
5672,
'/',
credentials)
ConnectionParameters
¶
Bases: Parameters
Connection parameters object that is passed into the connection adapter upon construction.
Source code in pika/connection.py
class ConnectionParameters(Parameters):
"""Connection parameters object that is passed into the connection adapter upon construction."""
# Protect against accidental assignment of an invalid attribute
__slots__ = ()
class _DEFAULT:
"""Designates default parameter value; internal use."""
... # noqa: PIE790, PYI013
T = TypeVar('T')
DefaultT = Union[T, Type[_DEFAULT]]
def __init__(self,
host: DefaultT[str] = _DEFAULT,
port: DefaultT[int] = _DEFAULT,
virtual_host: DefaultT[str] = _DEFAULT,
credentials: DefaultT[(
pika.credentials.PlainCredentials |
pika.credentials.ExternalCredentials)] = _DEFAULT,
channel_max: DefaultT[int] = _DEFAULT,
frame_max: DefaultT[int] = _DEFAULT,
heartbeat: DefaultT[int | Callable[[Connection, float], int] |
None] = _DEFAULT,
ssl_options: DefaultT[SSLOptions | None] = _DEFAULT,
connection_attempts: DefaultT[int] = _DEFAULT,
retry_delay: DefaultT[float] = _DEFAULT,
socket_timeout: DefaultT[float | None] = _DEFAULT,
stack_timeout: DefaultT[float | None] = _DEFAULT,
locale: DefaultT[str] = _DEFAULT,
blocked_connection_timeout: DefaultT[float | None] = _DEFAULT,
client_properties: DefaultT[dict[str, Any] | None] = _DEFAULT,
tcp_options: DefaultT[dict[str, Any] | None] = _DEFAULT,
**kwargs: Any) -> None:
"""
Create a new ConnectionParameters instance.
See `Parameters` for default values.
:param host: Hostname or IP Address to connect to
:param port: TCP port to connect to
:param virtual_host: RabbitMQ virtual host to use
:param credentials: auth credentials
:param channel_max: Maximum number of channels to allow
:param frame_max: The maximum byte size for an AMQP frame
:param heartbeat: Controls AMQP heartbeat timeout negotiation
during connection tuning. An integer value always overrides the value
proposed by broker. Use 0 to deactivate heartbeats and None to always
accept the broker's proposal. If a callable is given, it will be called
with the connection instance and the heartbeat timeout proposed by broker
as its arguments. The callback should return a non-negative integer that
will be used to override the broker's proposal.
:param ssl_options: None for plaintext or
`pika.SSLOptions` instance for SSL/TLS. Defaults to None.
:param connection_attempts: Maximum number of retry attempts
:param retry_delay: Time to wait in seconds, before the next
:param socket_timeout: Positive socket connect timeout in
seconds.
:param stack_timeout: Positive full protocol stack
(TCP/[SSL]/AMQP) bring-up timeout in seconds. It's recommended to
set this value higher than `socket_timeout`.
:param locale: Set the locale value
:param blocked_connection_timeout: If not None,
the value is a non-negative timeout, in seconds, for the
connection to remain blocked (triggered by Connection.Blocked from
broker); if the timeout expires before connection becomes unblocked,
the connection will be torn down, triggering the adapter-specific
mechanism for informing client app about the closed connection:
passing `ConnectionBlockedTimeout` exception to on_close_callback
in asynchronous adapters or raising it in `BlockingConnection`.
:param client_properties: None or dict of client properties used to
override the fields in the default client properties reported to
RabbitMQ via `Connection.StartOk` method.
:param tcp_options: None or a dict of TCP options to set for socket
"""
super().__init__()
if blocked_connection_timeout is not self._DEFAULT:
self.blocked_connection_timeout = cast(Optional[float],
blocked_connection_timeout)
if channel_max is not self._DEFAULT:
self.channel_max = cast(int, channel_max)
if client_properties is not self._DEFAULT:
self.client_properties = cast(Optional[Dict[str, Any]],
client_properties)
if connection_attempts is not self._DEFAULT:
self.connection_attempts = cast(int, connection_attempts)
if credentials is not self._DEFAULT:
self.credentials = cast(
Union[pika.credentials.PlainCredentials,
pika.credentials.ExternalCredentials], credentials)
if frame_max is not self._DEFAULT:
self.frame_max = cast(int, frame_max)
if heartbeat is not self._DEFAULT:
self.heartbeat = cast(
Optional[Union[int, Callable[[Connection, float], int]]],
heartbeat)
if host is not self._DEFAULT:
self.host = cast(str, host)
if locale is not self._DEFAULT:
self.locale = cast(str, locale)
if retry_delay is not self._DEFAULT:
self.retry_delay = cast(float, retry_delay)
if socket_timeout is not self._DEFAULT:
self.socket_timeout = cast(Optional[float], socket_timeout)
if stack_timeout is not self._DEFAULT:
self.stack_timeout = cast(Optional[float], stack_timeout)
if ssl_options is not self._DEFAULT:
self.ssl_options = cast(Optional[SSLOptions], ssl_options)
# Set port after SSL status is known
if port is not self._DEFAULT:
self.port = cast(int, port)
else:
self.port = self.DEFAULT_SSL_PORT if self.ssl_options else self.DEFAULT_PORT
if virtual_host is not self._DEFAULT:
self.virtual_host = cast(str, virtual_host)
if tcp_options is not self._DEFAULT:
self.tcp_options = cast(Optional[Dict[str, Any]], tcp_options)
if kwargs:
raise TypeError(f'unexpected kwargs: {kwargs!r}')
blocked_connection_timeout
instance-attribute
¶
blocked_connection_timeout = cast(
Optional[float], blocked_connection_timeout
)
client_properties
instance-attribute
¶
client_properties = cast(
Optional[Dict[str, Any]], client_properties
)
credentials
instance-attribute
¶
credentials = cast(
Union[
pika.credentials.PlainCredentials,
pika.credentials.ExternalCredentials,
],
credentials,
)
heartbeat
instance-attribute
¶
heartbeat = cast(
Optional[
Union[int, Callable[[Connection, float], int]]
],
heartbeat,
)
DEFAULT_BLOCKED_CONNECTION_TIMEOUT
class-attribute
instance-attribute
¶
DEFAULT_BLOCKED_CONNECTION_TIMEOUT = None
DEFAULT_CHANNEL_MAX
class-attribute
instance-attribute
¶
DEFAULT_CHANNEL_MAX = pika.channel.MAX_CHANNELS
DEFAULT_CREDENTIALS
class-attribute
instance-attribute
¶
DEFAULT_CREDENTIALS = pika.credentials.PlainCredentials(
DEFAULT_USERNAME, DEFAULT_PASSWORD
)
_credentials
instance-attribute
¶
_credentials: PlainCredentials | ExternalCredentials = (
self.DEFAULT_CREDENTIALS
)
_heartbeat
instance-attribute
¶
_heartbeat: (
int | Callable[[Connection, float], int] | None
) = None
_DEFAULT
¶
Designates default parameter value; internal use.
Source code in pika/connection.py
class _DEFAULT:
"""Designates default parameter value; internal use."""
... # noqa: PIE790, PYI013
__init__
¶
__init__(
host: DefaultT[str] = _DEFAULT,
port: DefaultT[int] = _DEFAULT,
virtual_host: DefaultT[str] = _DEFAULT,
credentials: DefaultT[
PlainCredentials | ExternalCredentials
] = _DEFAULT,
channel_max: DefaultT[int] = _DEFAULT,
frame_max: DefaultT[int] = _DEFAULT,
heartbeat: DefaultT[
int | Callable[[Connection, float], int] | None
] = _DEFAULT,
ssl_options: DefaultT[SSLOptions | None] = _DEFAULT,
connection_attempts: DefaultT[int] = _DEFAULT,
retry_delay: DefaultT[float] = _DEFAULT,
socket_timeout: DefaultT[float | None] = _DEFAULT,
stack_timeout: DefaultT[float | None] = _DEFAULT,
locale: DefaultT[str] = _DEFAULT,
blocked_connection_timeout: DefaultT[
float | None
] = _DEFAULT,
client_properties: DefaultT[
dict[str, Any] | None
] = _DEFAULT,
tcp_options: DefaultT[dict[str, Any] | None] = _DEFAULT,
**kwargs: Any,
) -> None
Create a new ConnectionParameters instance.
See Parameters for default values.
| PARAMETER | DESCRIPTION |
|---|---|
host
|
Hostname or IP Address to connect to
TYPE:
|
port
|
TCP port to connect to
TYPE:
|
virtual_host
|
RabbitMQ virtual host to use
TYPE:
|
credentials
|
auth credentials
TYPE:
|
channel_max
|
Maximum number of channels to allow
TYPE:
|
frame_max
|
The maximum byte size for an AMQP frame
TYPE:
|
heartbeat
|
Controls AMQP heartbeat timeout negotiation during connection tuning. An integer value always overrides the value proposed by broker. Use 0 to deactivate heartbeats and None to always accept the broker's proposal. If a callable is given, it will be called with the connection instance and the heartbeat timeout proposed by broker as its arguments. The callback should return a non-negative integer that will be used to override the broker's proposal.
TYPE:
|
ssl_options
|
None for plaintext or
TYPE:
|
connection_attempts
|
Maximum number of retry attempts
TYPE:
|
retry_delay
|
Time to wait in seconds, before the next
TYPE:
|
socket_timeout
|
Positive socket connect timeout in seconds.
TYPE:
|
stack_timeout
|
Positive full protocol stack
(TCP/[SSL]/AMQP) bring-up timeout in seconds. It's recommended to
set this value higher than
TYPE:
|
locale
|
Set the locale value
TYPE:
|
blocked_connection_timeout
|
If not None,
the value is a non-negative timeout, in seconds, for the
connection to remain blocked (triggered by Connection.Blocked from
broker); if the timeout expires before connection becomes unblocked,
the connection will be torn down, triggering the adapter-specific
mechanism for informing client app about the closed connection:
passing
TYPE:
|
client_properties
|
None or dict of client properties used to
override the fields in the default client properties reported to
RabbitMQ via
TYPE:
|
tcp_options
|
None or a dict of TCP options to set for socket
TYPE:
|
Source code in pika/connection.py
def __init__(self,
host: DefaultT[str] = _DEFAULT,
port: DefaultT[int] = _DEFAULT,
virtual_host: DefaultT[str] = _DEFAULT,
credentials: DefaultT[(
pika.credentials.PlainCredentials |
pika.credentials.ExternalCredentials)] = _DEFAULT,
channel_max: DefaultT[int] = _DEFAULT,
frame_max: DefaultT[int] = _DEFAULT,
heartbeat: DefaultT[int | Callable[[Connection, float], int] |
None] = _DEFAULT,
ssl_options: DefaultT[SSLOptions | None] = _DEFAULT,
connection_attempts: DefaultT[int] = _DEFAULT,
retry_delay: DefaultT[float] = _DEFAULT,
socket_timeout: DefaultT[float | None] = _DEFAULT,
stack_timeout: DefaultT[float | None] = _DEFAULT,
locale: DefaultT[str] = _DEFAULT,
blocked_connection_timeout: DefaultT[float | None] = _DEFAULT,
client_properties: DefaultT[dict[str, Any] | None] = _DEFAULT,
tcp_options: DefaultT[dict[str, Any] | None] = _DEFAULT,
**kwargs: Any) -> None:
"""
Create a new ConnectionParameters instance.
See `Parameters` for default values.
:param host: Hostname or IP Address to connect to
:param port: TCP port to connect to
:param virtual_host: RabbitMQ virtual host to use
:param credentials: auth credentials
:param channel_max: Maximum number of channels to allow
:param frame_max: The maximum byte size for an AMQP frame
:param heartbeat: Controls AMQP heartbeat timeout negotiation
during connection tuning. An integer value always overrides the value
proposed by broker. Use 0 to deactivate heartbeats and None to always
accept the broker's proposal. If a callable is given, it will be called
with the connection instance and the heartbeat timeout proposed by broker
as its arguments. The callback should return a non-negative integer that
will be used to override the broker's proposal.
:param ssl_options: None for plaintext or
`pika.SSLOptions` instance for SSL/TLS. Defaults to None.
:param connection_attempts: Maximum number of retry attempts
:param retry_delay: Time to wait in seconds, before the next
:param socket_timeout: Positive socket connect timeout in
seconds.
:param stack_timeout: Positive full protocol stack
(TCP/[SSL]/AMQP) bring-up timeout in seconds. It's recommended to
set this value higher than `socket_timeout`.
:param locale: Set the locale value
:param blocked_connection_timeout: If not None,
the value is a non-negative timeout, in seconds, for the
connection to remain blocked (triggered by Connection.Blocked from
broker); if the timeout expires before connection becomes unblocked,
the connection will be torn down, triggering the adapter-specific
mechanism for informing client app about the closed connection:
passing `ConnectionBlockedTimeout` exception to on_close_callback
in asynchronous adapters or raising it in `BlockingConnection`.
:param client_properties: None or dict of client properties used to
override the fields in the default client properties reported to
RabbitMQ via `Connection.StartOk` method.
:param tcp_options: None or a dict of TCP options to set for socket
"""
super().__init__()
if blocked_connection_timeout is not self._DEFAULT:
self.blocked_connection_timeout = cast(Optional[float],
blocked_connection_timeout)
if channel_max is not self._DEFAULT:
self.channel_max = cast(int, channel_max)
if client_properties is not self._DEFAULT:
self.client_properties = cast(Optional[Dict[str, Any]],
client_properties)
if connection_attempts is not self._DEFAULT:
self.connection_attempts = cast(int, connection_attempts)
if credentials is not self._DEFAULT:
self.credentials = cast(
Union[pika.credentials.PlainCredentials,
pika.credentials.ExternalCredentials], credentials)
if frame_max is not self._DEFAULT:
self.frame_max = cast(int, frame_max)
if heartbeat is not self._DEFAULT:
self.heartbeat = cast(
Optional[Union[int, Callable[[Connection, float], int]]],
heartbeat)
if host is not self._DEFAULT:
self.host = cast(str, host)
if locale is not self._DEFAULT:
self.locale = cast(str, locale)
if retry_delay is not self._DEFAULT:
self.retry_delay = cast(float, retry_delay)
if socket_timeout is not self._DEFAULT:
self.socket_timeout = cast(Optional[float], socket_timeout)
if stack_timeout is not self._DEFAULT:
self.stack_timeout = cast(Optional[float], stack_timeout)
if ssl_options is not self._DEFAULT:
self.ssl_options = cast(Optional[SSLOptions], ssl_options)
# Set port after SSL status is known
if port is not self._DEFAULT:
self.port = cast(int, port)
else:
self.port = self.DEFAULT_SSL_PORT if self.ssl_options else self.DEFAULT_PORT
if virtual_host is not self._DEFAULT:
self.virtual_host = cast(str, virtual_host)
if tcp_options is not self._DEFAULT:
self.tcp_options = cast(Optional[Dict[str, Any]], tcp_options)
if kwargs:
raise TypeError(f'unexpected kwargs: {kwargs!r}')
__repr__
¶
__repr__() -> str
Represent the info about the instance.
Source code in pika/connection.py
@override
def __repr__(self) -> str:
"""Represent the info about the instance."""
return (
f'<{self.__class__.__name__} host={self.host} port={self.port} virtual_host={self.virtual_host} ssl={bool(self.ssl_options)}>'
)
__eq__
¶
__eq__(other) -> bool
Source code in pika/connection.py
@override
def __eq__(self, other) -> bool:
if isinstance(other, Parameters):
return self._host == other._host and self._port == other._port
return NotImplemented
__ne__
¶
__ne__(other) -> bool
Source code in pika/connection.py
@override
def __ne__(self, other) -> bool:
result = self.__eq__(other)
if result is not NotImplemented:
return not result
return NotImplemented
URLParameters¶
The URLParameters class allows you to pass in an AMQP URL when creating the object and supports the host, port, virtual host, ssl, username and password in the base URL and other options are passed in via query parameters.
Example:
import pika
# Set the connection parameters to connect to rabbit-server1 on port 5672
# on the / virtual host using the username "guest" and password "guest"
parameters = pika.URLParameters('amqp://guest:guest@rabbit-server1:5672/%2F')
URLParameters
¶
Bases: Parameters
Connect to RabbitMQ via an AMQP URL in the format::
amqp://username:password@host:port/<virtual_host>[?query-string]
Ensure that the virtual host is URI encoded when specified. For example if
you are using the default "/" virtual host, the value should be %2f.
See Parameters for default values.
Valid query string values are:
- channel_max:
Override the default maximum channel count value
- client_properties:
dict of client properties used to override the fields in the default
client properties reported to RabbitMQ via `Connection.StartOk`
method
- connection_attempts:
Specify how many times pika should try and reconnect before it gives up
- frame_max:
Override the default maximum frame size for communication
- heartbeat:
Desired connection heartbeat timeout for negotiation. If not present
the broker's value is accepted. 0 turns heartbeat off.
- locale:
Override the default `en_US` locale value
- ssl_options:
None for plaintext; for SSL: dict of public ssl context-related
arguments that may be passed to :meth:`ssl.SSLSocket` as kwargs,
except `sock`, `server_side`,`do_handshake_on_connect`, `family`,
`type`, `proto`, `fileno`.
- retry_delay:
The number of seconds to sleep before attempting to connect on
connection failure.
- socket_timeout:
Socket connect timeout value in seconds (float or int)
- stack_timeout:
Positive full protocol stack (TCP/[SSL]/AMQP) bring-up timeout in
seconds. It's recommended to set this value higher than
`socket_timeout`.
- blocked_connection_timeout:
Set the timeout, in seconds, that the connection may remain blocked
(triggered by Connection.Blocked from broker); if the timeout
expires before connection becomes unblocked, the connection will be
torn down, triggering the connection's on_close_callback
- tcp_options:
Set the tcp options for the underlying socket.
| PARAMETER | DESCRIPTION |
|---|---|
url
|
The AMQP URL to connect to
TYPE:
|
Source code in pika/connection.py
class URLParameters(Parameters):
"""Connect to RabbitMQ via an AMQP URL in the format::
amqp://username:password@host:port/<virtual_host>[?query-string]
Ensure that the virtual host is URI encoded when specified. For example if
you are using the default "/" virtual host, the value should be `%2f`.
See `Parameters` for default values.
Valid query string values are:
- channel_max:
Override the default maximum channel count value
- client_properties:
dict of client properties used to override the fields in the default
client properties reported to RabbitMQ via `Connection.StartOk`
method
- connection_attempts:
Specify how many times pika should try and reconnect before it gives up
- frame_max:
Override the default maximum frame size for communication
- heartbeat:
Desired connection heartbeat timeout for negotiation. If not present
the broker's value is accepted. 0 turns heartbeat off.
- locale:
Override the default `en_US` locale value
- ssl_options:
None for plaintext; for SSL: dict of public ssl context-related
arguments that may be passed to :meth:`ssl.SSLSocket` as kwargs,
except `sock`, `server_side`,`do_handshake_on_connect`, `family`,
`type`, `proto`, `fileno`.
- retry_delay:
The number of seconds to sleep before attempting to connect on
connection failure.
- socket_timeout:
Socket connect timeout value in seconds (float or int)
- stack_timeout:
Positive full protocol stack (TCP/[SSL]/AMQP) bring-up timeout in
seconds. It's recommended to set this value higher than
`socket_timeout`.
- blocked_connection_timeout:
Set the timeout, in seconds, that the connection may remain blocked
(triggered by Connection.Blocked from broker); if the timeout
expires before connection becomes unblocked, the connection will be
torn down, triggering the connection's on_close_callback
- tcp_options:
Set the tcp options for the underlying socket.
:param url: The AMQP URL to connect to
"""
# Protect against accidental assignment of an invalid attribute
__slots__ = ('_all_url_query_values',)
# The name of the private function for parsing and setting a given URL query
# arg is constructed by catenating the query arg's name to this prefix
_SETTER_PREFIX = '_set_url_'
def __init__(self, url: str) -> None:
"""
Create a new URLParameters instance.
:param url: The URL value
"""
super().__init__()
self._all_url_query_values = None
# Handle the Protocol scheme
if url[0:4].lower() == 'amqp':
url = 'http' + url[4:]
parts = urlparse(url)
if parts.scheme == 'https':
# Create default context which will get overridden by the
# ssl_options URL arg, if any
self.ssl_options = pika.SSLOptions(
context=ssl.create_default_context())
elif parts.scheme == 'http':
self.ssl_options = None
elif parts.scheme:
raise ValueError(
f'Unexpected URL scheme {parts.scheme!r}; supported scheme '
'values: amqp, amqps')
if parts.hostname is not None:
self.host = parts.hostname
# Take care of port after SSL status is known
if parts.port is not None:
self.port = parts.port
else:
self.port = (self.DEFAULT_SSL_PORT
if self.ssl_options else self.DEFAULT_PORT)
if parts.username is not None:
assert parts.password is not None
self.credentials = pika.credentials.PlainCredentials(
url_unquote(parts.username), url_unquote(parts.password))
# Get the Virtual Host
if len(parts.path) > 1:
self.virtual_host = url_unquote(parts.path.split('/')[1])
# Handle query string values, validating and assigning them
self._all_url_query_values = url_parse_qs(parts.query)
for name, value in self._all_url_query_values.items():
try:
set_value = getattr(self, self._SETTER_PREFIX + name)
except AttributeError:
raise ValueError(f'Unknown URL parameter: {name!r}')
try:
(single_value,) = value
except ValueError:
raise ValueError(
f'Expected exactly one value for URL parameter '
f'{name}, but got {len(value)} values: {value}')
set_value(single_value)
def _set_url_blocked_connection_timeout(self, value: float) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
blocked_connection_timeout = float(value)
except ValueError as exc:
raise ValueError(
f'Invalid blocked_connection_timeout value {value!r}: {exc!r}')
self.blocked_connection_timeout = blocked_connection_timeout
def _set_url_channel_max(self, value: int) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
channel_max = int(value)
except ValueError as exc:
raise ValueError(f'Invalid channel_max value {value!r}: {exc!r}')
self.channel_max = channel_max
def _set_url_client_properties(self, value: str) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
self.client_properties = ast.literal_eval(value)
def _set_url_connection_attempts(self, value: int) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
connection_attempts = int(value)
except ValueError as exc:
raise ValueError(
f'Invalid connection_attempts value {value!r}: {exc!r}')
self.connection_attempts = connection_attempts
def _set_url_frame_max(self, value: int) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
frame_max = int(value)
except ValueError as exc:
raise ValueError(f'Invalid frame_max value {value!r}: {exc!r}')
self.frame_max = frame_max
def _set_url_heartbeat(self, value: int) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
heartbeat_timeout = int(value)
except ValueError as exc:
raise ValueError(f'Invalid heartbeat value {value!r}: {exc!r}')
self.heartbeat = heartbeat_timeout
def _set_url_locale(self, value: str) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
self.locale = value
def _set_url_retry_delay(self, value: float) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
retry_delay = float(value)
except ValueError as exc:
raise ValueError(f'Invalid retry_delay value {value!r}: {exc!r}')
self.retry_delay = retry_delay
def _set_url_socket_timeout(self, value: float) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
socket_timeout = float(value)
except ValueError as exc:
raise ValueError(f'Invalid socket_timeout value {value!r}: {exc!r}')
self.socket_timeout = socket_timeout
def _set_url_stack_timeout(self, value: float) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
stack_timeout = float(value)
except ValueError as exc:
raise ValueError(f'Invalid stack_timeout value {value!r}: {exc!r}')
self.stack_timeout = stack_timeout
def _set_url_ssl_options(self, value: str) -> None:
"""
Deserialize and apply the corresponding query string arg.
:param value: Raw query-string value to deserialize
"""
opts = ast.literal_eval(value)
if opts is None:
if self.ssl_options is not None:
raise ValueError(
'Specified ssl_options=None URI arg is inconsistent with '
'the specified amqps URI scheme.')
else:
# Older versions of Pika would take the opts dict and pass it
# directly as kwargs to the deprecated ssl.wrap_socket method.
# Here, we take the valid options and translate them into args
# for various SSLContext methods.
#
# https://docs.python.org/3/library/ssl.html#ssl.wrap_socket
#
# SSLContext.load_verify_locations(cafile=None, capath=None, cadata=None)
try:
opt_protocol = ssl.PROTOCOL_TLS_CLIENT
except AttributeError:
opt_protocol = ssl.PROTOCOL_TLSv1_2
if 'protocol' in opts:
opt_protocol = opts['protocol']
cxt = ssl.SSLContext(protocol=opt_protocol)
opt_cafile = opts.get('ca_certs') or opts.get('cafile')
opt_capath = opts.get('ca_path') or opts.get('capath')
opt_cadata = opts.get('ca_data') or opts.get('cadata')
cxt.load_verify_locations(opt_cafile, opt_capath, opt_cadata)
# SSLContext.load_cert_chain(certfile, keyfile=None, password=None)
if 'certfile' in opts:
opt_certfile = opts['certfile']
opt_keyfile = opts.get('keyfile')
opt_password = opts.get('password')
cxt.load_cert_chain(opt_certfile, opt_keyfile, opt_password)
if 'ciphers' in opts:
opt_ciphers = opts['ciphers']
cxt.set_ciphers(opt_ciphers)
server_hostname = opts.get('server_hostname')
self.ssl_options = pika.SSLOptions(context=cxt,
server_hostname=server_hostname)
def _set_url_tcp_options(self, value: str) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
self.tcp_options = ast.literal_eval(value)
ssl_options
instance-attribute
¶
ssl_options = pika.SSLOptions(
context=ssl.create_default_context()
)
credentials
instance-attribute
¶
credentials = pika.credentials.PlainCredentials(
url_unquote(parts.username), url_unquote(parts.password)
)
DEFAULT_BLOCKED_CONNECTION_TIMEOUT
class-attribute
instance-attribute
¶
DEFAULT_BLOCKED_CONNECTION_TIMEOUT = None
DEFAULT_CHANNEL_MAX
class-attribute
instance-attribute
¶
DEFAULT_CHANNEL_MAX = pika.channel.MAX_CHANNELS
DEFAULT_CREDENTIALS
class-attribute
instance-attribute
¶
DEFAULT_CREDENTIALS = pika.credentials.PlainCredentials(
DEFAULT_USERNAME, DEFAULT_PASSWORD
)
_credentials
instance-attribute
¶
_credentials: PlainCredentials | ExternalCredentials = (
self.DEFAULT_CREDENTIALS
)
_heartbeat
instance-attribute
¶
_heartbeat: (
int | Callable[[Connection, float], int] | None
) = None
blocked_connection_timeout
property
writable
¶
blocked_connection_timeout: float | None
| RETURNS | DESCRIPTION |
|---|---|
float | None
|
blocked connection timeout. Defaults to
|
channel_max
property
writable
¶
channel_max: int
| RETURNS | DESCRIPTION |
|---|---|
int
|
max preferred number of channels. Defaults to
|
client_properties
property
writable
¶
client_properties: dict[str, Any] | None
| RETURNS | DESCRIPTION |
|---|---|
dict[str, Any] | None
|
client properties used to override the fields in the default
client properties reported to RabbitMQ via |
connection_attempts
property
writable
¶
connection_attempts: int
| RETURNS | DESCRIPTION |
|---|---|
int
|
number of socket connection attempts. Defaults to
|
frame_max
property
writable
¶
frame_max: int
| RETURNS | DESCRIPTION |
|---|---|
int
|
desired maximum AMQP frame size to use. Defaults to
|
heartbeat
property
writable
¶
heartbeat: int | Callable[[Connection, float], int] | None
| RETURNS | DESCRIPTION |
|---|---|
int | Callable[[Connection, float], int] | None
|
AMQP connection heartbeat timeout value for negotiation during
connection tuning or callable which is invoked during connection tuning.
None to accept broker's value. 0 turns heartbeat off. Defaults to
|
locale
property
writable
¶
locale: str
| RETURNS | DESCRIPTION |
|---|---|
str
|
locale value to pass to broker; e.g., 'en_US'. Defaults to
|
retry_delay
property
writable
¶
retry_delay: float
| RETURNS | DESCRIPTION |
|---|---|
float
|
interval between socket connection attempts; see also
|
socket_timeout
property
writable
¶
socket_timeout: float | None
| RETURNS | DESCRIPTION |
|---|---|
float | None
|
socket connect timeout in seconds. Defaults to
|
stack_timeout
property
writable
¶
stack_timeout: float | None
| RETURNS | DESCRIPTION |
|---|---|
float | None
|
full protocol stack TCP/[SSL]/AMQP bring-up timeout in
seconds. Defaults to |
tcp_options
property
writable
¶
tcp_options: dict[str, int] | None
| RETURNS | DESCRIPTION |
|---|---|
dict[str, int] | None
|
None or a dict of options to pass to the underlying socket |
__init__
¶
__init__(url: str) -> None
Create a new URLParameters instance.
| PARAMETER | DESCRIPTION |
|---|---|
url
|
The URL value
TYPE:
|
Source code in pika/connection.py
def __init__(self, url: str) -> None:
"""
Create a new URLParameters instance.
:param url: The URL value
"""
super().__init__()
self._all_url_query_values = None
# Handle the Protocol scheme
if url[0:4].lower() == 'amqp':
url = 'http' + url[4:]
parts = urlparse(url)
if parts.scheme == 'https':
# Create default context which will get overridden by the
# ssl_options URL arg, if any
self.ssl_options = pika.SSLOptions(
context=ssl.create_default_context())
elif parts.scheme == 'http':
self.ssl_options = None
elif parts.scheme:
raise ValueError(
f'Unexpected URL scheme {parts.scheme!r}; supported scheme '
'values: amqp, amqps')
if parts.hostname is not None:
self.host = parts.hostname
# Take care of port after SSL status is known
if parts.port is not None:
self.port = parts.port
else:
self.port = (self.DEFAULT_SSL_PORT
if self.ssl_options else self.DEFAULT_PORT)
if parts.username is not None:
assert parts.password is not None
self.credentials = pika.credentials.PlainCredentials(
url_unquote(parts.username), url_unquote(parts.password))
# Get the Virtual Host
if len(parts.path) > 1:
self.virtual_host = url_unquote(parts.path.split('/')[1])
# Handle query string values, validating and assigning them
self._all_url_query_values = url_parse_qs(parts.query)
for name, value in self._all_url_query_values.items():
try:
set_value = getattr(self, self._SETTER_PREFIX + name)
except AttributeError:
raise ValueError(f'Unknown URL parameter: {name!r}')
try:
(single_value,) = value
except ValueError:
raise ValueError(
f'Expected exactly one value for URL parameter '
f'{name}, but got {len(value)} values: {value}')
set_value(single_value)
_set_url_blocked_connection_timeout
¶
_set_url_blocked_connection_timeout(value: float) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_blocked_connection_timeout(self, value: float) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
blocked_connection_timeout = float(value)
except ValueError as exc:
raise ValueError(
f'Invalid blocked_connection_timeout value {value!r}: {exc!r}')
self.blocked_connection_timeout = blocked_connection_timeout
_set_url_channel_max
¶
_set_url_channel_max(value: int) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_channel_max(self, value: int) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
channel_max = int(value)
except ValueError as exc:
raise ValueError(f'Invalid channel_max value {value!r}: {exc!r}')
self.channel_max = channel_max
_set_url_client_properties
¶
_set_url_client_properties(value: str) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_client_properties(self, value: str) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
self.client_properties = ast.literal_eval(value)
_set_url_connection_attempts
¶
_set_url_connection_attempts(value: int) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_connection_attempts(self, value: int) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
connection_attempts = int(value)
except ValueError as exc:
raise ValueError(
f'Invalid connection_attempts value {value!r}: {exc!r}')
self.connection_attempts = connection_attempts
_set_url_frame_max
¶
_set_url_frame_max(value: int) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_frame_max(self, value: int) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
frame_max = int(value)
except ValueError as exc:
raise ValueError(f'Invalid frame_max value {value!r}: {exc!r}')
self.frame_max = frame_max
_set_url_heartbeat
¶
_set_url_heartbeat(value: int) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_heartbeat(self, value: int) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
heartbeat_timeout = int(value)
except ValueError as exc:
raise ValueError(f'Invalid heartbeat value {value!r}: {exc!r}')
self.heartbeat = heartbeat_timeout
_set_url_locale
¶
_set_url_locale(value: str) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_locale(self, value: str) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
self.locale = value
_set_url_retry_delay
¶
_set_url_retry_delay(value: float) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_retry_delay(self, value: float) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
retry_delay = float(value)
except ValueError as exc:
raise ValueError(f'Invalid retry_delay value {value!r}: {exc!r}')
self.retry_delay = retry_delay
_set_url_socket_timeout
¶
_set_url_socket_timeout(value: float) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_socket_timeout(self, value: float) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
socket_timeout = float(value)
except ValueError as exc:
raise ValueError(f'Invalid socket_timeout value {value!r}: {exc!r}')
self.socket_timeout = socket_timeout
_set_url_stack_timeout
¶
_set_url_stack_timeout(value: float) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_stack_timeout(self, value: float) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
try:
stack_timeout = float(value)
except ValueError as exc:
raise ValueError(f'Invalid stack_timeout value {value!r}: {exc!r}')
self.stack_timeout = stack_timeout
_set_url_ssl_options
¶
_set_url_ssl_options(value: str) -> None
Deserialize and apply the corresponding query string arg.
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_ssl_options(self, value: str) -> None:
"""
Deserialize and apply the corresponding query string arg.
:param value: Raw query-string value to deserialize
"""
opts = ast.literal_eval(value)
if opts is None:
if self.ssl_options is not None:
raise ValueError(
'Specified ssl_options=None URI arg is inconsistent with '
'the specified amqps URI scheme.')
else:
# Older versions of Pika would take the opts dict and pass it
# directly as kwargs to the deprecated ssl.wrap_socket method.
# Here, we take the valid options and translate them into args
# for various SSLContext methods.
#
# https://docs.python.org/3/library/ssl.html#ssl.wrap_socket
#
# SSLContext.load_verify_locations(cafile=None, capath=None, cadata=None)
try:
opt_protocol = ssl.PROTOCOL_TLS_CLIENT
except AttributeError:
opt_protocol = ssl.PROTOCOL_TLSv1_2
if 'protocol' in opts:
opt_protocol = opts['protocol']
cxt = ssl.SSLContext(protocol=opt_protocol)
opt_cafile = opts.get('ca_certs') or opts.get('cafile')
opt_capath = opts.get('ca_path') or opts.get('capath')
opt_cadata = opts.get('ca_data') or opts.get('cadata')
cxt.load_verify_locations(opt_cafile, opt_capath, opt_cadata)
# SSLContext.load_cert_chain(certfile, keyfile=None, password=None)
if 'certfile' in opts:
opt_certfile = opts['certfile']
opt_keyfile = opts.get('keyfile')
opt_password = opts.get('password')
cxt.load_cert_chain(opt_certfile, opt_keyfile, opt_password)
if 'ciphers' in opts:
opt_ciphers = opts['ciphers']
cxt.set_ciphers(opt_ciphers)
server_hostname = opts.get('server_hostname')
self.ssl_options = pika.SSLOptions(context=cxt,
server_hostname=server_hostname)
_set_url_tcp_options
¶
_set_url_tcp_options(value: str) -> None
Deserialize and apply the corresponding query string arg
| PARAMETER | DESCRIPTION |
|---|---|
value
|
Raw query-string value to deserialize
TYPE:
|
Source code in pika/connection.py
def _set_url_tcp_options(self, value: str) -> None:
"""Deserialize and apply the corresponding query string arg
:param value: Raw query-string value to deserialize
"""
self.tcp_options = ast.literal_eval(value)
__repr__
¶
__repr__() -> str
Represent the info about the instance.
Source code in pika/connection.py
@override
def __repr__(self) -> str:
"""Represent the info about the instance."""
return (
f'<{self.__class__.__name__} host={self.host} port={self.port} virtual_host={self.virtual_host} ssl={bool(self.ssl_options)}>'
)
__eq__
¶
__eq__(other) -> bool
Source code in pika/connection.py
@override
def __eq__(self, other) -> bool:
if isinstance(other, Parameters):
return self._host == other._host and self._port == other._port
return NotImplemented
__ne__
¶
__ne__(other) -> bool
Source code in pika/connection.py
@override
def __ne__(self, other) -> bool:
result = self.__eq__(other)
if result is not NotImplemented:
return not result
return NotImplemented