Skip to content

Twisted Connection Adapter

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

twisted_connection

Using Pika with a Twisted reactor.

The interfaces in this module are Deferred-based when possible. This means that the connection.channel() method and most of the channel methods return Deferreds instead of taking a callback argument and that basic_consume() returns a Twisted DeferredQueue where messages from the server will be stored. Refer to the docstrings for TwistedProtocolConnection.channel() and the TwistedChannel class for details.

ClosableDeferredQueue

Bases: DeferredQueue

Like the normal Twisted DeferredQueue, but after close() is called with an exception instance all pending Deferreds are errbacked and further attempts to call get() or put() return a Failure wrapping that exception.

Source code in pika/adapters/twisted_connection.py
class ClosableDeferredQueue(defer.DeferredQueue):
    """Like the normal Twisted DeferredQueue, but after close() is called with an exception instance
    all pending Deferreds are errbacked and further attempts to call get() or put() return a Failure
    wrapping that exception.
    """

    def __init__(self,
                 size: int | None = None,
                 backlog: int | None = None) -> None:
        self.closed: Exception | None = None
        super().__init__(size, backlog)

    @override
    def put(  # type: ignore[override]
            self, obj: Any) -> defer.Deferred[Any] | None:
        """
        Like the original :meth:`DeferredQueue.put` method, but returns an errback if the queue is
        closed.

        :param obj: Object to put into the queue
        :returns: Nothing on success, or a failed Deferred if the queue is closed
        """
        if self.closed:
            LOGGER.error('Impossible to put to the queue, it is closed.')
            return defer.fail(self.closed)
        # `DeferredQueue.put()` returns nothing, so the success path yields
        # `None`. The return type stays optional for the closed path above.
        defer.DeferredQueue.put(self, obj)
        return None

    @override
    def get(self) -> defer.Deferred[Any]:
        """
        Returns a Deferred that will fire with the next item in the queue, when it's available.

        The Deferred will errback if the queue is closed.

        :returns: Deferred that fires with the next item.
        """
        if self.closed:
            LOGGER.error('Impossible to get from the queue, it is closed.')
            return defer.fail(self.closed)
        return defer.DeferredQueue.get(self)

    def close(self, reason: Exception | None) -> None:
        """
        Closes the queue.

        Errback the pending calls to :meth:`get()`.

        :param reason: The reason for closing the queue
        """
        if self.closed:
            LOGGER.warning('Queue was already closed with reason: %s.',
                           self.closed)
        self.closed = reason
        while self.waiting:
            self.waiting.pop().errback(reason)
        self.pending = []

close

close(reason: Exception | None) -> None

Closes the queue.

Errback the pending calls to :meth:get().

PARAMETER DESCRIPTION
reason

The reason for closing the queue

TYPE: Exception | None

Source code in pika/adapters/twisted_connection.py
def close(self, reason: Exception | None) -> None:
    """
    Closes the queue.

    Errback the pending calls to :meth:`get()`.

    :param reason: The reason for closing the queue
    """
    if self.closed:
        LOGGER.warning('Queue was already closed with reason: %s.',
                       self.closed)
    self.closed = reason
    while self.waiting:
        self.waiting.pop().errback(reason)
    self.pending = []

get

get() -> Deferred[Any]

Returns a Deferred that will fire with the next item in the queue, when it's available.

The Deferred will errback if the queue is closed.

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires with the next item.

Source code in pika/adapters/twisted_connection.py
@override
def get(self) -> defer.Deferred[Any]:
    """
    Returns a Deferred that will fire with the next item in the queue, when it's available.

    The Deferred will errback if the queue is closed.

    :returns: Deferred that fires with the next item.
    """
    if self.closed:
        LOGGER.error('Impossible to get from the queue, it is closed.')
        return defer.fail(self.closed)
    return defer.DeferredQueue.get(self)

put

put(obj: Any) -> Deferred[Any] | None

Like the original :meth:DeferredQueue.put method, but returns an errback if the queue is closed.

PARAMETER DESCRIPTION
obj

Object to put into the queue

TYPE: Any

RETURNS DESCRIPTION
Deferred[Any] | None

Nothing on success, or a failed Deferred if the queue is closed

Source code in pika/adapters/twisted_connection.py
@override
def put(  # type: ignore[override]
        self, obj: Any) -> defer.Deferred[Any] | None:
    """
    Like the original :meth:`DeferredQueue.put` method, but returns an errback if the queue is
    closed.

    :param obj: Object to put into the queue
    :returns: Nothing on success, or a failed Deferred if the queue is closed
    """
    if self.closed:
        LOGGER.error('Impossible to put to the queue, it is closed.')
        return defer.fail(self.closed)
    # `DeferredQueue.put()` returns nothing, so the success path yields
    # `None`. The return type stays optional for the closed path above.
    defer.DeferredQueue.put(self, obj)
    return None

TwistedChannel

A wrapper around Pika's Channel.

Channel methods that normally take a callback argument are wrapped to return a Deferred that fires with whatever would be passed to the callback. If the channel gets closed, all pending Deferreds are errbacked with a ChannelClosed exception. The returned Deferreds fire with whatever arguments the callback to the original method would receive.

Some methods like basic_consume and basic_get are wrapped in a special way, see their docstrings for details.

Source code in pika/adapters/twisted_connection.py
class TwistedChannel:
    """
    A wrapper around Pika's Channel.

    Channel methods that normally take a callback argument are wrapped to return a Deferred that
    fires with whatever would be passed to the callback. If the channel gets closed, all pending
    Deferreds are errbacked with a ChannelClosed exception. The returned Deferreds fire with
    whatever arguments the callback to the original method would receive.

    Some methods like basic_consume and basic_get are wrapped in a special way, see their docstrings
    for details.
    """

    def __init__(self, channel: channel.Channel) -> None:
        self._channel = channel
        self._closed: Exception | None = None
        self._calls: set[defer.Deferred[Any]] = set()
        self._consumers: dict[str, ClosableDeferredQueue] = {}
        # Store Basic.Get calls so we can handle GetEmpty replies
        self._basic_get_deferred: defer.Deferred[Any] | None = None
        self._channel.add_callback(self._on_getempty, [spec.Basic.GetEmpty],
                                   False)
        # We need this mapping to close the ClosableDeferredQueue when a queue
        # is deleted.
        self._queue_name_to_consumer_tags: dict[str, set[str]] = {}
        # Whether RabbitMQ delivery confirmation has been enabled
        self._delivery_confirmation = False
        self._delivery_message_id: int | None = None
        self._deliveries: dict[int, defer.Deferred[Any]] = {}
        # Holds a ReceivedMessage object representing a message received via
        # Basic.Return in publisher-acknowledgments mode.
        self._puback_return: ReceivedMessage | None = None

        self.on_closed: defer.Deferred[Any] = defer.Deferred()
        self._channel.add_on_close_callback(self._on_channel_closed)
        self._channel.add_on_cancel_callback(
            self._on_consumer_cancelled_by_broker)

    @override
    def __repr__(self) -> str:
        return f'<{self.__class__.__name__} channel={self._channel!r}>'

    def _on_channel_closed(self, _channel: channel.Channel,
                           reason: Exception | None) -> None:
        # enter the closed state
        self._closed = reason
        # errback all pending calls
        for d in self._calls:
            d.errback(self._closed)
        # errback all pending deliveries
        for d in self._deliveries.values():
            d.errback(self._closed)
        # close all open queues
        for consumer in self._consumers.values():
            consumer.close(self._closed)
        # release references to stored objects
        self._calls = set()
        self._deliveries = {}
        self._consumers = {}
        self.on_closed.callback(self._closed)

    def _on_consumer_cancelled_by_broker(
        self, method_frame: pika.frame.Method[pika.spec.Basic.Cancel]
    ) -> (pika.frame.Method[pika.spec.Basic.Cancel] |
          pika.frame.Method[pika.spec.Basic.CancelOk]):
        """
        Called by impl when broker cancels consumer via Basic.Cancel.

        This is a RabbitMQ-specific feature. The circumstances include deletion of queue being
        consumed as well as failure of a HA node responsible for the queue being consumed.

        :param method_frame: method frame with the `spec.Basic.Cancel` method
        """
        return self._on_consumer_cancelled(method_frame)

    def _on_consumer_cancelled(
        self, frame: (pika.frame.Method[pika.spec.Basic.Cancel] |
                      pika.frame.Method[pika.spec.Basic.CancelOk])
    ) -> (pika.frame.Method[pika.spec.Basic.Cancel] |
          pika.frame.Method[pika.spec.Basic.CancelOk]):
        """
        Called when the broker cancels a consumer via Basic.Cancel or when the broker responds to a
        Basic.Cancel request by Basic.CancelOk.

        :param frame: method frame with the `spec.Basic.Cancel` or `spec.Basic.CancelOk` method
        """
        consumer_tag = frame.method.consumer_tag
        if consumer_tag not in self._consumers:
            # Could be cancelled by user or broker earlier
            LOGGER.warning('basic_cancel - consumer not found: %s',
                           consumer_tag)
            return frame
        self._consumers[consumer_tag].close(exceptions.ConsumerCancelled())
        del self._consumers[consumer_tag]
        # Remove from the queue-to-ctags index:
        for ctags in self._queue_name_to_consumer_tags.values():
            try:
                ctags.remove(consumer_tag)
            except KeyError:
                continue
        return frame

    def _on_getempty(
            self,
            _method_frame: pika.frame.Method[pika.spec.Basic.Get]) -> None:
        """
        Callback the Basic.Get deferred with None.

        :param _method_frame: Method frame from Basic.Get response (unused)
        """
        if self._basic_get_deferred is None:
            LOGGER.warning('Got Basic.GetEmpty but no Basic.Get calls '
                           'were pending.')
            return
        self._basic_get_deferred.callback(None)

    def _wrap_channel_method(self,
                             name: str) -> Callable[..., defer.Deferred[Any]]:
        """
        Wrap Pika's Channel method to make it return a Deferred that fires when the method completes
        and errbacks if the channel gets closed. If the original method's callback would receive
        more than one argument, the Deferred fires with a tuple of argument values.

        :param name: Attribute name to look up on the underlying channel
        """
        method = getattr(self._channel, name)

        @functools.wraps(method)
        def wrapped(*args, **kwargs):
            if self._closed:
                return defer.fail(self._closed)

            d: defer.Deferred = defer.Deferred()
            self._calls.add(d)
            d.addCallback(self._clear_call, d)

            def single_argument(*args) -> None:
                """
                Make sure that the deferred is called with a single argument.

                In case the original callback fires with more than one, convert to a tuple.
                """
                if len(args) > 1:
                    d.callback(tuple(args))
                else:
                    d.callback(*args)

            kwargs['callback'] = single_argument

            try:
                method(*args, **kwargs)
            except Exception:
                return defer.fail()
            return d

        return wrapped

    def _clear_call(self, ret: Any, d: defer.Deferred[Any]) -> Any:
        self._calls.discard(d)
        return ret

    # Public Channel attributes

    @property
    def channel_number(self) -> int:
        return self._channel.channel_number

    @property
    def connection(self) -> pika.connection.Connection:
        return self._channel.connection

    @property
    def is_closed(self) -> bool:
        """Returns True if the channel is closed."""
        return self._channel.is_closed

    @property
    def is_closing(self) -> bool:
        """Returns True if client-initiated closing of the channel is in progress."""
        return self._channel.is_closing

    @property
    def is_open(self) -> bool:
        """Returns True if the channel is open."""
        return self._channel.is_open

    @property
    def flow_active(self) -> bool:
        return self._channel.flow_active

    @property
    def consumer_tags(self) -> list[str]:
        return self._channel.consumer_tags

    # Deferred-equivalents of public Channel methods

    def callback_deferred(self, deferred: defer.Deferred[Any],
                          replies: Sequence[type[amqp_object.Method]]) -> None:
        """
        Pass in a Deferred and a list replies from the RabbitMQ broker which you'd like the Deferred
        to be callbacked on with the frame as callback value.

        :param deferred: The Deferred to callback
        :param replies: The replies to callback on
        """
        self._channel.add_callback(deferred.callback, replies)

    # Public Channel methods

    def add_on_return_callback(
            self, callback: Callable[[ReceivedMessage], None]) -> None:
        """
        Pass a callback function that will be called when a published message is rejected and
        returned by the server via `Basic.Return`.

        :param callback: The method to call on callback with the
            message as only argument. The message is a named tuple with
            the following attributes
            - channel: this TwistedChannel
            - method: pika.spec.Basic.Return
            - properties: pika.spec.BasicProperties
            - body: bytes
        """
        self._channel.add_on_return_callback(
            lambda _channel, method, properties, body: callback(
                ReceivedMessage(
                    channel=self,
                    method=method,
                    properties=properties,
                    body=body,
                )))

    def basic_ack(self, delivery_tag: int = 0, multiple: bool = False) -> None:
        """
        Acknowledge one or more messages.

        When sent by the client, this method acknowledges one or more messages delivered via the
        Deliver or Get-Ok methods. When sent by server, this method acknowledges one or more
        messages published with the Publish method on a channel in confirm mode. The acknowledgement
        can be for a single message or a set of messages up to and including a specific message.

        :param delivery_tag: int/long The server-assigned delivery tag
        :param multiple: If set to True, the delivery tag is treated as "up to and including", so
            that multiple messages can be acknowledged with a single method. If set to False, the
            delivery tag refers to a single message. If the multiple field is 1, and the delivery
            tag is zero, this indicates acknowledgement of all outstanding messages.
        """
        return self._channel.basic_ack(delivery_tag=delivery_tag,
                                       multiple=multiple)

    def basic_cancel(
        self,
        consumer_tag: str = ''
    ) -> defer.Deferred[(pika.frame.Method[pika.spec.Basic.CancelOk] |
                         pika.frame.Method[pika.spec.Basic.Cancel])]:
        """
        This method cancels a consumer.

        This does not affect already delivered messages, but it does mean the server will not send
        any more messages for that consumer. The client may receive an arbitrary number of messages
        in between sending the cancel method and receiving the cancel-ok reply. It may also be sent
        from the server to the client in the event of the consumer being unexpectedly cancelled
        (i.e. cancelled for any reason other than the server receiving the corresponding
        basic.cancel from the client). This allows clients to be notified of the loss of consumers
        due to events such as queue deletion.

        This method wraps :meth:`Channel.basic_cancel <pika.channel.Channel.basic_cancel>` and
        closes any deferred queue associated with that consumer.

        :param consumer_tag: Identifier for the consumer
        :returns: Deferred that fires on the Basic.CancelOk response
        :raises ValueError:
        """
        wrapped = self._wrap_channel_method('basic_cancel')
        d = wrapped(consumer_tag=consumer_tag)
        return d.addCallback(self._on_consumer_cancelled)

    def basic_consume(
            self,
            queue: str,
            auto_ack: bool = False,
            exclusive: bool = False,
            consumer_tag: str | None = None,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Consume from a server queue.

        Sends the AMQP 0-9-1 command Basic.Consume to the broker and binds
        messages for the consumer_tag to a
        :class:`ClosableDeferredQueue`. If you do not pass in a
        consumer_tag, one will be automatically generated for you.

        For more information on basic_consume, see:
        Tutorial 2 at https://www.rabbitmq.com/getstarted.html
        https://www.rabbitmq.com/confirms.html
        https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.consume

        :param queue: The queue to consume from. Use the empty string to
            specify the most recent server-named queue for this channel.
        :param auto_ack: if set to True, automatic acknowledgement mode
            will be used (see https://www.rabbitmq.com/confirms.html). This
            corresponds with the 'no_ack' parameter in the basic.consume AMQP
            0.9.1 method
        :param exclusive: Don't allow other consumers on the queue
        :param consumer_tag: Specify your own consumer tag
        :param arguments: Custom key/value pair arguments for the consumer
        :returns: Deferred that fires with a tuple
            ``(queue_object, consumer_tag)``. The Deferred will errback with an
            instance of :class:`exceptions.ChannelClosed` if the call fails.
            The queue object is an instance of :class:`ClosableDeferredQueue`,
            where data received from the queue will be stored. Clients should
            use its :meth:`get() <ClosableDeferredQueue.get>` method to fetch
            an individual message, which will return a Deferred firing with a
            namedtuple whose attributes are:
            - channel: this TwistedChannel
            - method: pika.spec.Basic.Deliver
            - properties: pika.spec.BasicProperties
            - body: bytes
        """
        if self._closed:
            return defer.fail(self._closed)

        queue_obj = ClosableDeferredQueue()
        d: defer.Deferred[Any] = defer.Deferred()
        self._calls.add(d)

        def on_consume_ok(frame) -> None:
            consumer_tag = frame.method.consumer_tag
            self._queue_name_to_consumer_tags.setdefault(
                queue, set()).add(consumer_tag)
            self._consumers[consumer_tag] = queue_obj
            self._calls.discard(d)
            d.callback((queue_obj, consumer_tag))

        def on_message_callback(_channel, method, properties, body) -> None:
            """Add the ReceivedMessage to the queue, while replacing the channel implementation."""
            queue_obj.put(
                ReceivedMessage(
                    channel=self,
                    method=method,
                    properties=properties,
                    body=body,
                ))

        try:
            self._channel.basic_consume(
                queue=queue,
                on_message_callback=on_message_callback,
                auto_ack=auto_ack,
                exclusive=exclusive,
                consumer_tag=consumer_tag,
                arguments=arguments,
                callback=on_consume_ok,
            )
        except Exception:
            return defer.fail()

        return d

    def basic_get(self,
                  queue: str,
                  auto_ack: bool = False) -> defer.Deferred[Any]:
        """
        Get a single message from the AMQP broker.

        Will return If the queue is empty, it will return None.
        If you want to
        be notified of Basic.GetEmpty, use the Channel.add_callback method
        adding your Basic.GetEmpty callback which should expect only one
        parameter, frame. Due to implementation details, this cannot be called
        a second time until the callback is executed.  For more information on
        basic_get and its parameters, see:

        https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.get

        This method wraps :meth:`Channel.basic_get
        <pika.channel.Channel.basic_get>`.

        :param queue: The queue from which to get a message. Use the empty
                      string to specify the most recent server-named queue
                      for this channel.
        :param auto_ack: Tell the broker to not expect a reply
        :returns: Deferred that fires with a namedtuple whose attributes are:
             - channel: this TwistedChannel
             - method: pika.spec.Basic.GetOk
             - properties: pika.spec.BasicProperties
             - body: bytes
            If the queue is empty, None will be returned.
        :raises pika.exceptions.DuplicateGetOkCallback:
        """
        if self._basic_get_deferred is not None:
            raise exceptions.DuplicateGetOkCallback()

        def create_namedtuple(result) -> ReceivedMessage | None:
            if result is None:
                return None
            _channel, method, properties, body = result
            return ReceivedMessage(
                channel=self,
                method=method,
                properties=properties,
                body=body,
            )

        def cleanup_attribute(result):
            self._basic_get_deferred = None
            return result

        d = self._wrap_channel_method('basic_get')(queue=queue,
                                                   auto_ack=auto_ack)
        d.addCallback(create_namedtuple)
        d.addBoth(cleanup_attribute)
        self._basic_get_deferred = d
        return d

    def basic_nack(self,
                   delivery_tag: int = 0,
                   multiple: bool = False,
                   requeue: bool = True) -> None:
        """
        This method allows a client to reject one or more incoming messages.

        It can be used to interrupt and cancel large incoming messages, or return untreatable
        messages to their original queue.

        :param delivery_tag: int/long The server-assigned delivery tag
        :param multiple: If set to True, the delivery tag is treated as "up to and including", so
            that multiple messages can be acknowledged with a single method. If set to False, the
            delivery tag refers to a single message. If the multiple field is 1, and the delivery
            tag is zero, this indicates acknowledgement of all outstanding messages.
        :param requeue: If requeue is true, the server will attempt to requeue the message. If
            requeue is false or the requeue attempt fails the messages are discarded or dead-
            lettered.
        """
        return self._channel.basic_nack(
            delivery_tag=delivery_tag,
            multiple=multiple,
            requeue=requeue,
        )

    def basic_publish(self,
                      exchange: str,
                      routing_key: str,
                      body: bytes,
                      properties: spec.BasicProperties | None = None,
                      mandatory: bool = False) -> defer.Deferred[Any]:
        """
        Publish to the channel with the given exchange, routing key and body.

        This method wraps :meth:`Channel.basic_publish <pika.channel.Channel.basic_publish>`, but
        makes sure the channel is not closed before publishing.

        For more information on basic_publish and what the parameters do, see:

        https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.publish

        :param exchange: The exchange to publish to
        :param routing_key: The routing key to bind on
        :param body: The message body
        :param properties: Basic.properties
        :param mandatory: The mandatory flag
        :returns: A Deferred that fires with the result of the channel's basic_publish.
        :raises UnroutableError: raised when a message published in publisher-acknowledgments mode
            (see `BlockingChannel.confirm_delivery`) is returned via `Basic.Return` followed by
            `Basic.Ack`.
        :raises NackError: raised when a message published in publisher-acknowledgements mode is
            Nack'ed by the broker. See `BlockingChannel.confirm_delivery`.
        """
        if self._closed:
            return defer.fail(self._closed)
        self._channel.basic_publish(exchange=exchange,
                                    routing_key=routing_key,
                                    body=body,
                                    properties=properties,
                                    mandatory=mandatory)
        if not self._delivery_confirmation:
            # `Channel.basic_publish()` returns nothing, so without delivery
            # confirmation there is no result to report beyond completion.
            return defer.succeed(None)
        # See https://www.rabbitmq.com/confirms.html#publisher-confirms
        assert self._delivery_message_id is not None
        self._delivery_message_id += 1
        self._deliveries[self._delivery_message_id] = defer.Deferred()
        return self._deliveries[self._delivery_message_id]

    def basic_qos(self,
                  prefetch_size: int = 0,
                  prefetch_count: int = 0,
                  global_qos: bool = False) -> defer.Deferred[Any]:
        """
        Specify quality of service.

        This method requests a specific quality of service. The QoS can be specified for the current
        channel or for all channels on the connection. The client can request that messages be sent
        in advance so that when the client finishes processing a message, the following message is
        already held locally, rather than needing to be sent down the channel. Prefetching gives a
        performance improvement.

        :param prefetch_size: This field specifies the prefetch window size. The server will send a
            message in advance if it is equal to or smaller in size than the available prefetch size
            (and also falls into other prefetch limits). May be set to zero, meaning "no specific
            limit", although other prefetch limits may still apply. The prefetch-size is ignored by
            consumers who have enabled the no-ack option.
        :param prefetch_count: Specifies a prefetch window in terms of whole messages. This field
            may be used in combination with the prefetch-size field; a message will only be sent in
            advance if both prefetch windows (and those at the channel and connection level) allow
            it. The prefetch-count is ignored by consumers who have enabled the no-ack option.
        :param global_qos: Should the QoS apply to all channels on the connection.
        :returns: Deferred that fires on the Basic.QosOk response
        """
        return self._wrap_channel_method('basic_qos')(
            prefetch_size=prefetch_size,
            prefetch_count=prefetch_count,
            global_qos=global_qos,
        )

    def basic_reject(self, delivery_tag: int, requeue: bool = True) -> None:
        """
        Reject an incoming message.

        This method allows a client to reject a message. It can be used to interrupt and cancel
        large incoming messages, or return untreatable messages to their original queue.

        :param delivery_tag: int/long The server-assigned delivery tag
        :param requeue: If requeue is true, the server will attempt to requeue the message. If
            requeue is false or the requeue attempt fails the messages are discarded or dead-
            lettered.
        :raises TypeError:
        """
        return self._channel.basic_reject(delivery_tag=delivery_tag,
                                          requeue=requeue)

    def basic_recover(self, requeue: bool = False) -> defer.Deferred[Any]:
        """
        This method asks the server to redeliver all unacknowledged messages on a specified channel.
        Zero or more messages may be redelivered. This method replaces the asynchronous Recover.

        :param requeue: If False, the message will be redelivered to the original recipient. If
            True, the server will attempt to requeue the message, potentially then delivering it to
            an alternative subscriber.
        :returns: Deferred that fires on the Basic.RecoverOk response
        """
        return self._wrap_channel_method('basic_recover')(requeue=requeue)

    def close(self,
              reply_code: int = 0,
              reply_text: str = 'Normal shutdown') -> None:
        """
        Invoke a graceful shutdown of the channel with the AMQP Broker.

        If channel is OPENING, transition to CLOSING and suppress the incoming Channel.OpenOk, if
        any.

        :param reply_code: The reason code to send to broker
        :param reply_text: The reason text to send to broker
        :raises ChannelWrongStateError: if channel is closed or closing
        """
        return self._channel.close(reply_code=reply_code, reply_text=reply_text)

    def confirm_delivery(self) -> defer.Deferred[Any]:
        """
        Turn on Confirm mode in the channel.

        Pass in a callback to be notified by the Broker when a message has been confirmed as received or
        rejected (Basic.Ack, Basic.Nack) from the broker to the publisher.

        For more information see:
            https://www.rabbitmq.com/confirms.html#publisher-confirms

        :returns: Deferred that fires on the Confirm.SelectOk response
        """
        if self._delivery_confirmation:
            LOGGER.error('confirm_delivery: confirmation was already enabled.')
            return defer.succeed(None)
        wrapped = self._wrap_channel_method('confirm_delivery')
        d = wrapped(ack_nack_callback=self._on_delivery_confirmation)

        def set_delivery_confirmation(result):
            self._delivery_confirmation = True
            self._delivery_message_id = 0
            LOGGER.debug('Delivery confirmation enabled.')
            return result

        d.addCallback(set_delivery_confirmation)
        # Unroutable messages returned after this point will be in the context
        # of publisher acknowledgments
        self._channel.add_on_return_callback(self._on_puback_message_returned)
        return d

    def _on_delivery_confirmation(
        self, method_frame: (pika.frame.Method[pika.spec.Basic.Ack] |
                             pika.frame.Method[pika.spec.Basic.Nack])
    ) -> None:
        """
        Invoked by pika when RabbitMQ responds to a Basic.Publish RPC command, passing in either a
        Basic.Ack or Basic.Nack frame with the delivery tag of the message that was published. The
        delivery tag is an integer counter indicating the message number that was sent on the
        channel via Basic.Publish. Here we're just doing house keeping to keep track of stats and
        remove message numbers that we expect a delivery confirmation of from the list used to keep
        track of messages that are pending confirmation.

        :param method_frame: Basic.Ack or Basic.Nack frame
        """
        delivery_tag = method_frame.method.delivery_tag
        if delivery_tag not in self._deliveries:
            LOGGER.error('Delivery tag %s not found in the pending deliveries',
                         delivery_tag)
            return
        if method_frame.method.multiple:
            tags = [tag for tag in self._deliveries if tag <= delivery_tag]
            tags.sort()
        else:
            tags = [delivery_tag]
        for tag in tags:
            d = self._deliveries[tag]
            del self._deliveries[tag]
            if isinstance(method_frame.method, pika.spec.Basic.Nack):
                # Broker was unable to process the message due to internal
                # error
                LOGGER.warning(
                    "Message was Nack'ed by broker: nack=%r; channel=%s;",
                    method_frame.method, self.channel_number)
                if self._puback_return is not None:
                    returned_messages: list[ReceivedMessage] = [
                        self._puback_return
                    ]
                    self._puback_return = None
                else:
                    returned_messages = []
                d.errback(exceptions.NackError(returned_messages))
            else:
                assert isinstance(method_frame.method, pika.spec.Basic.Ack)
                if self._puback_return is not None:
                    # Unroutable message was returned
                    returned_messages = [self._puback_return]
                    self._puback_return = None
                    d.errback(exceptions.UnroutableError(returned_messages))
                else:
                    d.callback(method_frame.method)

    def _on_puback_message_returned(self, channel: channel.Channel,
                                    method: spec.Basic.Return,
                                    properties: spec.BasicProperties,
                                    body: bytes) -> None:
        """
        Called as the result of Basic.Return from broker in publisher-acknowledgements mode.

        Saves the info as a ReturnedMessage instance in self._puback_return.

        :param channel: our self._impl channel
        :param method:
        :param properties: message properties
        :param body: returned message body; empty string if no body
        """
        assert isinstance(method, spec.Basic.Return), method
        assert isinstance(properties, spec.BasicProperties), properties

        LOGGER.warning(
            'Published message was returned: _delivery_confirmation=%s; '
            'channel=%s; method=%r; properties=%r; body_size=%d; '
            'body_prefix=%.255r', self._delivery_confirmation,
            channel.channel_number, method, properties,
            len(body) if body is not None else None, body)

        self._puback_return = ReceivedMessage(channel=self,
                                              method=method,
                                              properties=properties,
                                              body=body)

    def exchange_bind(
            self,
            destination: str,
            source: str,
            routing_key: str = '',
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Bind an exchange to another exchange.

        :param destination: The destination exchange to bind
        :param source: The source exchange to bind to
        :param routing_key: The routing key to bind on
        :param arguments: Custom key/value pair arguments for the binding
        :raises ValueError:
        :returns: Deferred that fires on the Exchange.BindOk response
        """
        return self._wrap_channel_method('exchange_bind')(
            destination=destination,
            source=source,
            routing_key=routing_key,
            arguments=arguments,
        )

    def exchange_declare(
            self,
            exchange: str,
            exchange_type: str | ExchangeType = ExchangeType.direct,
            passive: bool = False,
            durable: bool = False,
            auto_delete: bool = False,
            internal: bool = False,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        This method creates an exchange if it does not already exist, and if the exchange exists,
        verifies that it is of the correct and expected class.

        If passive set, the server will reply with Declare-Ok if the exchange
        already exists with the same name, and raise an error if not and if the
        exchange does not already exist, the server MUST raise a channel
        exception with reply code 404 (not found).

        :param exchange: The exchange name consists of a non-empty sequence
            of these characters: letters, digits, hyphen, underscore, period,
            or colon
        :param exchange_type: The exchange type to use
        :param passive: Perform a declare or just check to see if it
            exists
        :param durable: Survive a reboot of RabbitMQ
        :param auto_delete: Remove when no more queues are bound to it
        :param internal: Can only be published to by other exchanges
        :param arguments: Custom key/value pair arguments for the exchange
        :returns: Deferred that fires on the Exchange.DeclareOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('exchange_declare')(
            exchange=exchange,
            exchange_type=exchange_type,
            passive=passive,
            durable=durable,
            auto_delete=auto_delete,
            internal=internal,
            arguments=arguments,
        )

    def exchange_delete(self,
                        exchange: str | None = None,
                        if_unused: bool = False) -> defer.Deferred[Any]:
        """
        Delete the exchange.

        :param exchange: The exchange name
        :param if_unused: only delete if the exchange is unused
        :returns: Deferred that fires on the Exchange.DeleteOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('exchange_delete')(
            exchange=exchange,
            if_unused=if_unused,
        )

    def exchange_unbind(
            self,
            destination: str,
            source: str,
            routing_key: str = '',
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Unbind an exchange from another exchange.

        :param destination: The destination exchange to unbind
        :param source: The source exchange to unbind from
        :param routing_key: The routing key to unbind
        :param arguments: Custom key/value pair arguments for the binding
        :returns: Deferred that fires on the Exchange.UnbindOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('exchange_unbind')(
            destination=destination,
            source=source,
            routing_key=routing_key,
            arguments=arguments,
        )

    def flow(self, active: bool = True) -> defer.Deferred[Any]:
        """
        Turn Channel flow control off and on.

        Returns a Deferred that will fire with a bool indicating the channel flow state. For more
        information, please reference:

        https://www.rabbitmq.com/amqp-0-9-1-reference.html#channel.flow

        :param active: Turn flow on or off
        :returns: Deferred that fires with the channel flow state
        :raises ValueError:
        """
        return self._wrap_channel_method('flow')(active=active)

    def open(self) -> None:
        """Open the channel."""
        return self._channel.open()

    def queue_bind(
            self,
            queue: str,
            exchange: str,
            routing_key: str | None = None,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Bind the queue to the specified exchange.

        :param queue: The queue to bind to the exchange
        :param exchange: The source exchange to bind to
        :param routing_key: The routing key to bind on
        :param arguments: Custom key/value pair arguments for the binding
        :returns: Deferred that fires on the Queue.BindOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('queue_bind')(
            queue=queue,
            exchange=exchange,
            routing_key=routing_key,
            arguments=arguments,
        )

    def queue_declare(
            self,
            queue: str,
            passive: bool = False,
            durable: bool = False,
            exclusive: bool = False,
            auto_delete: bool = False,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Declare queue, create if needed.

        This method creates or checks a queue. When creating a new queue the client can specify various
        properties that control the durability of the queue and its contents,
        and the level of sharing for the queue.

        Use an empty string as the queue name for the broker to auto-generate
        one

        :param queue: The queue name; if empty string, the broker will
            create a unique queue name
        :param passive: Only check to see if the queue exists
        :param durable: Survive reboots of the broker
        :param exclusive: Only allow access by the current connection
        :param auto_delete: Delete after consumer cancels or disconnects
        :param arguments: Custom key/value arguments for the queue
        :returns: Deferred that fires on the Queue.DeclareOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('queue_declare')(
            queue=queue,
            passive=passive,
            durable=durable,
            exclusive=exclusive,
            auto_delete=auto_delete,
            arguments=arguments,
        )

    def queue_delete(self,
                     queue: str,
                     if_unused: bool = False,
                     if_empty: bool = False) -> defer.Deferred[Any]:
        """
        Delete a queue from the broker.

        This method wraps :meth:`Channel.queue_delete <pika.channel.Channel.queue_delete>`, and
        removes the reference to the queue object after it gets deleted on the server.

        :param queue: The queue to delete
        :param if_unused: only delete if it's unused
        :param if_empty: only delete if the queue is empty
        :returns: Deferred that fires on the Queue.DeleteOk response
        :raises ValueError:
        """
        wrapped = self._wrap_channel_method('queue_delete')
        d = wrapped(queue=queue, if_unused=if_unused, if_empty=if_empty)

        def _clear_consumer(ret, queue_name):
            for consumer_tag in list(
                    self._queue_name_to_consumer_tags.get(queue_name, set())):
                self._consumers[consumer_tag].close(
                    exceptions.ConsumerCancelled(
                        f'Queue {queue_name} was deleted.'))
                del self._consumers[consumer_tag]
                self._queue_name_to_consumer_tags[queue_name].remove(
                    consumer_tag)
            return ret

        return d.addCallback(_clear_consumer, queue)

    def queue_purge(self, queue: str) -> defer.Deferred[Any]:
        """
        Purge all of the messages from the specified queue.

        :param queue: The queue to purge
        :returns: Deferred that fires on the Queue.PurgeOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('queue_purge')(queue=queue)

    def queue_unbind(
            self,
            queue: str,
            exchange: str | None,
            routing_key: str | None = None,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Unbind a queue from an exchange.

        :param queue: The queue to unbind from the exchange
        :param exchange: The source exchange to bind from
        :param routing_key: The routing key to unbind
        :param arguments: Custom key/value pair arguments for the binding
        :returns: Deferred that fires on the Queue.UnbindOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('queue_unbind')(
            queue=queue,
            exchange=exchange,
            routing_key=routing_key,
            arguments=arguments,
        )

    def tx_commit(self) -> defer.Deferred[Any]:
        """
        Commit a transaction.

        :returns: Deferred that fires on the Tx.CommitOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('tx_commit')()

    def tx_rollback(self) -> defer.Deferred[Any]:
        """
        Rollback a transaction.

        :returns: Deferred that fires on the Tx.RollbackOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('tx_rollback')()

    def tx_select(self) -> defer.Deferred[Any]:
        """
        Select standard transaction mode.

        This method sets the channel to use standard transactions. The client must use this method
        at least once on a channel before using the Commit or Rollback methods.

        :returns: Deferred that fires on the Tx.SelectOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('tx_select')()

is_closed property

is_closed: bool

Returns True if the channel is closed.

is_closing property

is_closing: bool

Returns True if client-initiated closing of the channel is in progress.

is_open property

is_open: bool

Returns True if the channel is open.

add_on_return_callback

add_on_return_callback(
    callback: Callable[[ReceivedMessage], None],
) -> None

Pass a callback function that will be called when a published message is rejected and returned by the server via Basic.Return.

PARAMETER DESCRIPTION
callback

The method to call on callback with the message as only argument. The message is a named tuple with the following attributes - channel: this TwistedChannel - method: pika.spec.Basic.Return - properties: pika.spec.BasicProperties - body: bytes

TYPE: Callable[[ReceivedMessage], None]

Source code in pika/adapters/twisted_connection.py
def add_on_return_callback(
        self, callback: Callable[[ReceivedMessage], None]) -> None:
    """
    Pass a callback function that will be called when a published message is rejected and
    returned by the server via `Basic.Return`.

    :param callback: The method to call on callback with the
        message as only argument. The message is a named tuple with
        the following attributes
        - channel: this TwistedChannel
        - method: pika.spec.Basic.Return
        - properties: pika.spec.BasicProperties
        - body: bytes
    """
    self._channel.add_on_return_callback(
        lambda _channel, method, properties, body: callback(
            ReceivedMessage(
                channel=self,
                method=method,
                properties=properties,
                body=body,
            )))

basic_ack

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

Acknowledge one or more messages.

When sent by the client, this method acknowledges one or more messages delivered via the Deliver or Get-Ok methods. When sent by server, this method acknowledges one or more messages published with the Publish method on a channel in confirm mode. The acknowledgement can be for a single message or a set of messages up to and including a specific message.

PARAMETER DESCRIPTION
delivery_tag

int/long The server-assigned delivery tag

TYPE: int DEFAULT: 0

multiple

If set to True, the delivery tag is treated as "up to and including", so that multiple messages can be acknowledged with a single method. If set to False, the delivery tag refers to a single message. If the multiple field is 1, and the delivery tag is zero, this indicates acknowledgement of all outstanding messages.

TYPE: bool DEFAULT: False

Source code in pika/adapters/twisted_connection.py
def basic_ack(self, delivery_tag: int = 0, multiple: bool = False) -> None:
    """
    Acknowledge one or more messages.

    When sent by the client, this method acknowledges one or more messages delivered via the
    Deliver or Get-Ok methods. When sent by server, this method acknowledges one or more
    messages published with the Publish method on a channel in confirm mode. The acknowledgement
    can be for a single message or a set of messages up to and including a specific message.

    :param delivery_tag: int/long The server-assigned delivery tag
    :param multiple: If set to True, the delivery tag is treated as "up to and including", so
        that multiple messages can be acknowledged with a single method. If set to False, the
        delivery tag refers to a single message. If the multiple field is 1, and the delivery
        tag is zero, this indicates acknowledgement of all outstanding messages.
    """
    return self._channel.basic_ack(delivery_tag=delivery_tag,
                                   multiple=multiple)

basic_cancel

basic_cancel(
    consumer_tag: str = '',
) -> Deferred[Method[CancelOk] | Method[Cancel]]

This method cancels a consumer.

This does not affect already delivered messages, but it does mean the server will not send any more messages for that consumer. The client may receive an arbitrary number of messages in between sending the cancel method and receiving the cancel-ok reply. It may also be sent from the server to the client in the event of the consumer being unexpectedly cancelled (i.e. cancelled for any reason other than the server receiving the corresponding basic.cancel from the client). This allows clients to be notified of the loss of consumers due to events such as queue deletion.

This method wraps :meth:Channel.basic_cancel <pika.channel.Channel.basic_cancel> and closes any deferred queue associated with that consumer.

PARAMETER DESCRIPTION
consumer_tag

Identifier for the consumer

TYPE: str DEFAULT: ''

RETURNS DESCRIPTION
Deferred[Method[CancelOk] | Method[Cancel]]

Deferred that fires on the Basic.CancelOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def basic_cancel(
    self,
    consumer_tag: str = ''
) -> defer.Deferred[(pika.frame.Method[pika.spec.Basic.CancelOk] |
                     pika.frame.Method[pika.spec.Basic.Cancel])]:
    """
    This method cancels a consumer.

    This does not affect already delivered messages, but it does mean the server will not send
    any more messages for that consumer. The client may receive an arbitrary number of messages
    in between sending the cancel method and receiving the cancel-ok reply. It may also be sent
    from the server to the client in the event of the consumer being unexpectedly cancelled
    (i.e. cancelled for any reason other than the server receiving the corresponding
    basic.cancel from the client). This allows clients to be notified of the loss of consumers
    due to events such as queue deletion.

    This method wraps :meth:`Channel.basic_cancel <pika.channel.Channel.basic_cancel>` and
    closes any deferred queue associated with that consumer.

    :param consumer_tag: Identifier for the consumer
    :returns: Deferred that fires on the Basic.CancelOk response
    :raises ValueError:
    """
    wrapped = self._wrap_channel_method('basic_cancel')
    d = wrapped(consumer_tag=consumer_tag)
    return d.addCallback(self._on_consumer_cancelled)

basic_consume

basic_consume(
    queue: str,
    auto_ack: bool = False,
    exclusive: bool = False,
    consumer_tag: str | None = None,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Consume from a server queue.

Sends the AMQP 0-9-1 command Basic.Consume to the broker and binds messages for the consumer_tag to a :class:ClosableDeferredQueue. If you do not pass in a consumer_tag, one will be automatically generated for you.

For more information on basic_consume, see: Tutorial 2 at https://www.rabbitmq.com/getstarted.html https://www.rabbitmq.com/confirms.html https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.consume

PARAMETER DESCRIPTION
queue

The queue to consume from. Use the empty string to specify the most recent server-named queue for this channel.

TYPE: str

auto_ack

if set to True, automatic acknowledgement mode will be used (see https://www.rabbitmq.com/confirms.html). This corresponds with the 'no_ack' parameter in the basic.consume AMQP 0.9.1 method

TYPE: bool DEFAULT: False

exclusive

Don't allow other consumers on the queue

TYPE: bool DEFAULT: False

consumer_tag

Specify your own consumer tag

TYPE: str | None DEFAULT: None

arguments

Custom key/value pair arguments for the consumer

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires with a tuple (queue_object, consumer_tag). The Deferred will errback with an instance of :class:exceptions.ChannelClosed if the call fails. The queue object is an instance of :class:ClosableDeferredQueue, where data received from the queue will be stored. Clients should use its :meth:get() <ClosableDeferredQueue.get> method to fetch an individual message, which will return a Deferred firing with a namedtuple whose attributes are: - channel: this TwistedChannel - method: pika.spec.Basic.Deliver - properties: pika.spec.BasicProperties - body: bytes

Source code in pika/adapters/twisted_connection.py
def basic_consume(
        self,
        queue: str,
        auto_ack: bool = False,
        exclusive: bool = False,
        consumer_tag: str | None = None,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Consume from a server queue.

    Sends the AMQP 0-9-1 command Basic.Consume to the broker and binds
    messages for the consumer_tag to a
    :class:`ClosableDeferredQueue`. If you do not pass in a
    consumer_tag, one will be automatically generated for you.

    For more information on basic_consume, see:
    Tutorial 2 at https://www.rabbitmq.com/getstarted.html
    https://www.rabbitmq.com/confirms.html
    https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.consume

    :param queue: The queue to consume from. Use the empty string to
        specify the most recent server-named queue for this channel.
    :param auto_ack: if set to True, automatic acknowledgement mode
        will be used (see https://www.rabbitmq.com/confirms.html). This
        corresponds with the 'no_ack' parameter in the basic.consume AMQP
        0.9.1 method
    :param exclusive: Don't allow other consumers on the queue
    :param consumer_tag: Specify your own consumer tag
    :param arguments: Custom key/value pair arguments for the consumer
    :returns: Deferred that fires with a tuple
        ``(queue_object, consumer_tag)``. The Deferred will errback with an
        instance of :class:`exceptions.ChannelClosed` if the call fails.
        The queue object is an instance of :class:`ClosableDeferredQueue`,
        where data received from the queue will be stored. Clients should
        use its :meth:`get() <ClosableDeferredQueue.get>` method to fetch
        an individual message, which will return a Deferred firing with a
        namedtuple whose attributes are:
        - channel: this TwistedChannel
        - method: pika.spec.Basic.Deliver
        - properties: pika.spec.BasicProperties
        - body: bytes
    """
    if self._closed:
        return defer.fail(self._closed)

    queue_obj = ClosableDeferredQueue()
    d: defer.Deferred[Any] = defer.Deferred()
    self._calls.add(d)

    def on_consume_ok(frame) -> None:
        consumer_tag = frame.method.consumer_tag
        self._queue_name_to_consumer_tags.setdefault(
            queue, set()).add(consumer_tag)
        self._consumers[consumer_tag] = queue_obj
        self._calls.discard(d)
        d.callback((queue_obj, consumer_tag))

    def on_message_callback(_channel, method, properties, body) -> None:
        """Add the ReceivedMessage to the queue, while replacing the channel implementation."""
        queue_obj.put(
            ReceivedMessage(
                channel=self,
                method=method,
                properties=properties,
                body=body,
            ))

    try:
        self._channel.basic_consume(
            queue=queue,
            on_message_callback=on_message_callback,
            auto_ack=auto_ack,
            exclusive=exclusive,
            consumer_tag=consumer_tag,
            arguments=arguments,
            callback=on_consume_ok,
        )
    except Exception:
        return defer.fail()

    return d

basic_get

basic_get(
    queue: str, auto_ack: bool = False
) -> Deferred[Any]

Get a single message from the AMQP broker.

Will return If the queue is empty, it will return None. If you want to be notified of Basic.GetEmpty, use the Channel.add_callback method adding your Basic.GetEmpty callback which should expect only one parameter, frame. Due to implementation details, this cannot be called a second time until the callback is executed. For more information on basic_get and its parameters, see:

https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.get

This method wraps :meth:Channel.basic_get <pika.channel.Channel.basic_get>.

PARAMETER DESCRIPTION
queue

The queue from which to get a message. Use the empty string to specify the most recent server-named queue for this channel.

TYPE: str

auto_ack

Tell the broker to not expect a reply

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires with a namedtuple whose attributes are: - channel: this TwistedChannel - method: pika.spec.Basic.GetOk - properties: pika.spec.BasicProperties - body: bytes If the queue is empty, None will be returned.

RAISES DESCRIPTION
pika.exceptions.DuplicateGetOkCallback
Source code in pika/adapters/twisted_connection.py
def basic_get(self,
              queue: str,
              auto_ack: bool = False) -> defer.Deferred[Any]:
    """
    Get a single message from the AMQP broker.

    Will return If the queue is empty, it will return None.
    If you want to
    be notified of Basic.GetEmpty, use the Channel.add_callback method
    adding your Basic.GetEmpty callback which should expect only one
    parameter, frame. Due to implementation details, this cannot be called
    a second time until the callback is executed.  For more information on
    basic_get and its parameters, see:

    https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.get

    This method wraps :meth:`Channel.basic_get
    <pika.channel.Channel.basic_get>`.

    :param queue: The queue from which to get a message. Use the empty
                  string to specify the most recent server-named queue
                  for this channel.
    :param auto_ack: Tell the broker to not expect a reply
    :returns: Deferred that fires with a namedtuple whose attributes are:
         - channel: this TwistedChannel
         - method: pika.spec.Basic.GetOk
         - properties: pika.spec.BasicProperties
         - body: bytes
        If the queue is empty, None will be returned.
    :raises pika.exceptions.DuplicateGetOkCallback:
    """
    if self._basic_get_deferred is not None:
        raise exceptions.DuplicateGetOkCallback()

    def create_namedtuple(result) -> ReceivedMessage | None:
        if result is None:
            return None
        _channel, method, properties, body = result
        return ReceivedMessage(
            channel=self,
            method=method,
            properties=properties,
            body=body,
        )

    def cleanup_attribute(result):
        self._basic_get_deferred = None
        return result

    d = self._wrap_channel_method('basic_get')(queue=queue,
                                               auto_ack=auto_ack)
    d.addCallback(create_namedtuple)
    d.addBoth(cleanup_attribute)
    self._basic_get_deferred = d
    return d

basic_nack

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

This method allows a client to reject one or more incoming messages.

It can be used to interrupt and cancel large incoming messages, or return untreatable messages to their original queue.

PARAMETER DESCRIPTION
delivery_tag

int/long The server-assigned delivery tag

TYPE: int DEFAULT: 0

multiple

If set to True, the delivery tag is treated as "up to and including", so that multiple messages can be acknowledged with a single method. If set to False, the delivery tag refers to a single message. If the multiple field is 1, and the delivery tag is zero, this indicates acknowledgement of all outstanding messages.

TYPE: bool DEFAULT: False

requeue

If requeue is true, the server will attempt to requeue the message. If requeue is false or the requeue attempt fails the messages are discarded or dead- lettered.

TYPE: bool DEFAULT: True

Source code in pika/adapters/twisted_connection.py
def basic_nack(self,
               delivery_tag: int = 0,
               multiple: bool = False,
               requeue: bool = True) -> None:
    """
    This method allows a client to reject one or more incoming messages.

    It can be used to interrupt and cancel large incoming messages, or return untreatable
    messages to their original queue.

    :param delivery_tag: int/long The server-assigned delivery tag
    :param multiple: If set to True, the delivery tag is treated as "up to and including", so
        that multiple messages can be acknowledged with a single method. If set to False, the
        delivery tag refers to a single message. If the multiple field is 1, and the delivery
        tag is zero, this indicates acknowledgement of all outstanding messages.
    :param requeue: If requeue is true, the server will attempt to requeue the message. If
        requeue is false or the requeue attempt fails the messages are discarded or dead-
        lettered.
    """
    return self._channel.basic_nack(
        delivery_tag=delivery_tag,
        multiple=multiple,
        requeue=requeue,
    )

basic_publish

basic_publish(
    exchange: str,
    routing_key: str,
    body: bytes,
    properties: BasicProperties | None = None,
    mandatory: bool = False,
) -> Deferred[Any]

Publish to the channel with the given exchange, routing key and body.

This method wraps :meth:Channel.basic_publish <pika.channel.Channel.basic_publish>, but makes sure the channel is not closed before publishing.

For more information on basic_publish and what the parameters do, see:

https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.publish

PARAMETER DESCRIPTION
exchange

The exchange to publish to

TYPE: str

routing_key

The routing key to bind on

TYPE: str

body

The message body

TYPE: bytes

properties

Basic.properties

TYPE: BasicProperties | None DEFAULT: None

mandatory

The mandatory flag

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

A Deferred that fires with the result of the channel's basic_publish.

RAISES DESCRIPTION
UnroutableError

raised when a message published in publisher-acknowledgments mode (see BlockingChannel.confirm_delivery) is returned via Basic.Return followed by Basic.Ack.

NackError

raised when a message published in publisher-acknowledgements mode is Nack'ed by the broker. See BlockingChannel.confirm_delivery.

Source code in pika/adapters/twisted_connection.py
def basic_publish(self,
                  exchange: str,
                  routing_key: str,
                  body: bytes,
                  properties: spec.BasicProperties | None = None,
                  mandatory: bool = False) -> defer.Deferred[Any]:
    """
    Publish to the channel with the given exchange, routing key and body.

    This method wraps :meth:`Channel.basic_publish <pika.channel.Channel.basic_publish>`, but
    makes sure the channel is not closed before publishing.

    For more information on basic_publish and what the parameters do, see:

    https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.publish

    :param exchange: The exchange to publish to
    :param routing_key: The routing key to bind on
    :param body: The message body
    :param properties: Basic.properties
    :param mandatory: The mandatory flag
    :returns: A Deferred that fires with the result of the channel's basic_publish.
    :raises UnroutableError: raised when a message published in publisher-acknowledgments mode
        (see `BlockingChannel.confirm_delivery`) is returned via `Basic.Return` followed by
        `Basic.Ack`.
    :raises NackError: raised when a message published in publisher-acknowledgements mode is
        Nack'ed by the broker. See `BlockingChannel.confirm_delivery`.
    """
    if self._closed:
        return defer.fail(self._closed)
    self._channel.basic_publish(exchange=exchange,
                                routing_key=routing_key,
                                body=body,
                                properties=properties,
                                mandatory=mandatory)
    if not self._delivery_confirmation:
        # `Channel.basic_publish()` returns nothing, so without delivery
        # confirmation there is no result to report beyond completion.
        return defer.succeed(None)
    # See https://www.rabbitmq.com/confirms.html#publisher-confirms
    assert self._delivery_message_id is not None
    self._delivery_message_id += 1
    self._deliveries[self._delivery_message_id] = defer.Deferred()
    return self._deliveries[self._delivery_message_id]

basic_qos

basic_qos(
    prefetch_size: int = 0,
    prefetch_count: int = 0,
    global_qos: bool = False,
) -> Deferred[Any]

Specify quality of service.

This method requests a specific quality of service. The QoS can be specified for the current channel or for all channels on the connection. The client can request that messages be sent in advance so that when the client finishes processing a message, the following message is already held locally, rather than needing to be sent down the channel. Prefetching gives a performance improvement.

PARAMETER DESCRIPTION
prefetch_size

This field specifies the prefetch window size. The server will send a message in advance if it is equal to or smaller in size than the available prefetch size (and also falls into other prefetch limits). May be set to zero, meaning "no specific limit", although other prefetch limits may still apply. The prefetch-size is ignored by consumers who have enabled the no-ack option.

TYPE: int DEFAULT: 0

prefetch_count

Specifies a prefetch window in terms of whole messages. This field may be used in combination with the prefetch-size field; a message will only be sent in advance if both prefetch windows (and those at the channel and connection level) allow it. The prefetch-count is ignored by consumers who have enabled the no-ack option.

TYPE: int DEFAULT: 0

global_qos

Should the QoS apply to all channels on the connection.

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Basic.QosOk response

Source code in pika/adapters/twisted_connection.py
def basic_qos(self,
              prefetch_size: int = 0,
              prefetch_count: int = 0,
              global_qos: bool = False) -> defer.Deferred[Any]:
    """
    Specify quality of service.

    This method requests a specific quality of service. The QoS can be specified for the current
    channel or for all channels on the connection. The client can request that messages be sent
    in advance so that when the client finishes processing a message, the following message is
    already held locally, rather than needing to be sent down the channel. Prefetching gives a
    performance improvement.

    :param prefetch_size: This field specifies the prefetch window size. The server will send a
        message in advance if it is equal to or smaller in size than the available prefetch size
        (and also falls into other prefetch limits). May be set to zero, meaning "no specific
        limit", although other prefetch limits may still apply. The prefetch-size is ignored by
        consumers who have enabled the no-ack option.
    :param prefetch_count: Specifies a prefetch window in terms of whole messages. This field
        may be used in combination with the prefetch-size field; a message will only be sent in
        advance if both prefetch windows (and those at the channel and connection level) allow
        it. The prefetch-count is ignored by consumers who have enabled the no-ack option.
    :param global_qos: Should the QoS apply to all channels on the connection.
    :returns: Deferred that fires on the Basic.QosOk response
    """
    return self._wrap_channel_method('basic_qos')(
        prefetch_size=prefetch_size,
        prefetch_count=prefetch_count,
        global_qos=global_qos,
    )

basic_recover

basic_recover(requeue: bool = False) -> Deferred[Any]

This method asks the server to redeliver all unacknowledged messages on a specified channel. Zero or more messages may be redelivered. This method replaces the asynchronous Recover.

PARAMETER DESCRIPTION
requeue

If False, the message will be redelivered to the original recipient. If True, the server will attempt to requeue the message, potentially then delivering it to an alternative subscriber.

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Basic.RecoverOk response

Source code in pika/adapters/twisted_connection.py
def basic_recover(self, requeue: bool = False) -> defer.Deferred[Any]:
    """
    This method asks the server to redeliver all unacknowledged messages on a specified channel.
    Zero or more messages may be redelivered. This method replaces the asynchronous Recover.

    :param requeue: If False, the message will be redelivered to the original recipient. If
        True, the server will attempt to requeue the message, potentially then delivering it to
        an alternative subscriber.
    :returns: Deferred that fires on the Basic.RecoverOk response
    """
    return self._wrap_channel_method('basic_recover')(requeue=requeue)

basic_reject

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

Reject an incoming message.

This method allows a client to reject a message. It can be used to interrupt and cancel large incoming messages, or return untreatable messages to their original queue.

PARAMETER DESCRIPTION
delivery_tag

int/long The server-assigned delivery tag

TYPE: int

requeue

If requeue is true, the server will attempt to requeue the message. If requeue is false or the requeue attempt fails the messages are discarded or dead- lettered.

TYPE: bool DEFAULT: True

RAISES DESCRIPTION
TypeError
Source code in pika/adapters/twisted_connection.py
def basic_reject(self, delivery_tag: int, requeue: bool = True) -> None:
    """
    Reject an incoming message.

    This method allows a client to reject a message. It can be used to interrupt and cancel
    large incoming messages, or return untreatable messages to their original queue.

    :param delivery_tag: int/long The server-assigned delivery tag
    :param requeue: If requeue is true, the server will attempt to requeue the message. If
        requeue is false or the requeue attempt fails the messages are discarded or dead-
        lettered.
    :raises TypeError:
    """
    return self._channel.basic_reject(delivery_tag=delivery_tag,
                                      requeue=requeue)

callback_deferred

callback_deferred(
    deferred: Deferred[Any], replies: Sequence[type[Method]]
) -> None

Pass in a Deferred and a list replies from the RabbitMQ broker which you'd like the Deferred to be callbacked on with the frame as callback value.

PARAMETER DESCRIPTION
deferred

The Deferred to callback

TYPE: Deferred[Any]

replies

The replies to callback on

TYPE: Sequence[type[Method]]

Source code in pika/adapters/twisted_connection.py
def callback_deferred(self, deferred: defer.Deferred[Any],
                      replies: Sequence[type[amqp_object.Method]]) -> None:
    """
    Pass in a Deferred and a list replies from the RabbitMQ broker which you'd like the Deferred
    to be callbacked on with the frame as callback value.

    :param deferred: The Deferred to callback
    :param replies: The replies to callback on
    """
    self._channel.add_callback(deferred.callback, replies)

close

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

Invoke a graceful shutdown of the channel with the AMQP Broker.

If channel is OPENING, transition to CLOSING and suppress the incoming Channel.OpenOk, if any.

PARAMETER DESCRIPTION
reply_code

The reason code to send to broker

TYPE: int DEFAULT: 0

reply_text

The reason text to send to broker

TYPE: str DEFAULT: 'Normal shutdown'

RAISES DESCRIPTION
ChannelWrongStateError

if channel is closed or closing

Source code in pika/adapters/twisted_connection.py
def close(self,
          reply_code: int = 0,
          reply_text: str = 'Normal shutdown') -> None:
    """
    Invoke a graceful shutdown of the channel with the AMQP Broker.

    If channel is OPENING, transition to CLOSING and suppress the incoming Channel.OpenOk, if
    any.

    :param reply_code: The reason code to send to broker
    :param reply_text: The reason text to send to broker
    :raises ChannelWrongStateError: if channel is closed or closing
    """
    return self._channel.close(reply_code=reply_code, reply_text=reply_text)

confirm_delivery

confirm_delivery() -> Deferred[Any]

Turn on Confirm mode in the channel.

Pass in a callback to be notified by the Broker when a message has been confirmed as received or rejected (Basic.Ack, Basic.Nack) from the broker to the publisher.

For more information see: https://www.rabbitmq.com/confirms.html#publisher-confirms

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Confirm.SelectOk response

Source code in pika/adapters/twisted_connection.py
def confirm_delivery(self) -> defer.Deferred[Any]:
    """
    Turn on Confirm mode in the channel.

    Pass in a callback to be notified by the Broker when a message has been confirmed as received or
    rejected (Basic.Ack, Basic.Nack) from the broker to the publisher.

    For more information see:
        https://www.rabbitmq.com/confirms.html#publisher-confirms

    :returns: Deferred that fires on the Confirm.SelectOk response
    """
    if self._delivery_confirmation:
        LOGGER.error('confirm_delivery: confirmation was already enabled.')
        return defer.succeed(None)
    wrapped = self._wrap_channel_method('confirm_delivery')
    d = wrapped(ack_nack_callback=self._on_delivery_confirmation)

    def set_delivery_confirmation(result):
        self._delivery_confirmation = True
        self._delivery_message_id = 0
        LOGGER.debug('Delivery confirmation enabled.')
        return result

    d.addCallback(set_delivery_confirmation)
    # Unroutable messages returned after this point will be in the context
    # of publisher acknowledgments
    self._channel.add_on_return_callback(self._on_puback_message_returned)
    return d

exchange_bind

exchange_bind(
    destination: str,
    source: str,
    routing_key: str = '',
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Bind an exchange to another exchange.

PARAMETER DESCRIPTION
destination

The destination exchange to bind

TYPE: str

source

The source exchange to bind to

TYPE: str

routing_key

The routing key to bind on

TYPE: str DEFAULT: ''

arguments

Custom key/value pair arguments for the binding

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Exchange.BindOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def exchange_bind(
        self,
        destination: str,
        source: str,
        routing_key: str = '',
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Bind an exchange to another exchange.

    :param destination: The destination exchange to bind
    :param source: The source exchange to bind to
    :param routing_key: The routing key to bind on
    :param arguments: Custom key/value pair arguments for the binding
    :raises ValueError:
    :returns: Deferred that fires on the Exchange.BindOk response
    """
    return self._wrap_channel_method('exchange_bind')(
        destination=destination,
        source=source,
        routing_key=routing_key,
        arguments=arguments,
    )

exchange_declare

exchange_declare(
    exchange: str,
    exchange_type: str | ExchangeType = direct,
    passive: bool = False,
    durable: bool = False,
    auto_delete: bool = False,
    internal: bool = False,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

This method creates an exchange if it does not already exist, and if the exchange exists, verifies that it is of the correct and expected class.

If passive set, the server will reply with Declare-Ok if the exchange already exists with the same name, and raise an error if not and if the exchange does not already exist, the server MUST raise a channel exception with reply code 404 (not found).

PARAMETER DESCRIPTION
exchange

The exchange name consists of a non-empty sequence of these characters: letters, digits, hyphen, underscore, period, or colon

TYPE: str

exchange_type

The exchange type to use

TYPE: str | ExchangeType DEFAULT: direct

passive

Perform a declare or just check to see if it exists

TYPE: bool DEFAULT: False

durable

Survive a reboot of RabbitMQ

TYPE: bool DEFAULT: False

auto_delete

Remove when no more queues are bound to it

TYPE: bool DEFAULT: False

internal

Can only be published to by other exchanges

TYPE: bool DEFAULT: False

arguments

Custom key/value pair arguments for the exchange

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Exchange.DeclareOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def exchange_declare(
        self,
        exchange: str,
        exchange_type: str | ExchangeType = ExchangeType.direct,
        passive: bool = False,
        durable: bool = False,
        auto_delete: bool = False,
        internal: bool = False,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    This method creates an exchange if it does not already exist, and if the exchange exists,
    verifies that it is of the correct and expected class.

    If passive set, the server will reply with Declare-Ok if the exchange
    already exists with the same name, and raise an error if not and if the
    exchange does not already exist, the server MUST raise a channel
    exception with reply code 404 (not found).

    :param exchange: The exchange name consists of a non-empty sequence
        of these characters: letters, digits, hyphen, underscore, period,
        or colon
    :param exchange_type: The exchange type to use
    :param passive: Perform a declare or just check to see if it
        exists
    :param durable: Survive a reboot of RabbitMQ
    :param auto_delete: Remove when no more queues are bound to it
    :param internal: Can only be published to by other exchanges
    :param arguments: Custom key/value pair arguments for the exchange
    :returns: Deferred that fires on the Exchange.DeclareOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('exchange_declare')(
        exchange=exchange,
        exchange_type=exchange_type,
        passive=passive,
        durable=durable,
        auto_delete=auto_delete,
        internal=internal,
        arguments=arguments,
    )

exchange_delete

exchange_delete(
    exchange: str | None = None, if_unused: bool = False
) -> Deferred[Any]

Delete the exchange.

PARAMETER DESCRIPTION
exchange

The exchange name

TYPE: str | None DEFAULT: None

if_unused

only delete if the exchange is unused

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Exchange.DeleteOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def exchange_delete(self,
                    exchange: str | None = None,
                    if_unused: bool = False) -> defer.Deferred[Any]:
    """
    Delete the exchange.

    :param exchange: The exchange name
    :param if_unused: only delete if the exchange is unused
    :returns: Deferred that fires on the Exchange.DeleteOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('exchange_delete')(
        exchange=exchange,
        if_unused=if_unused,
    )

exchange_unbind

exchange_unbind(
    destination: str,
    source: str,
    routing_key: str = '',
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Unbind an exchange from another exchange.

PARAMETER DESCRIPTION
destination

The destination exchange to unbind

TYPE: str

source

The source exchange to unbind from

TYPE: str

routing_key

The routing key to unbind

TYPE: str DEFAULT: ''

arguments

Custom key/value pair arguments for the binding

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Exchange.UnbindOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def exchange_unbind(
        self,
        destination: str,
        source: str,
        routing_key: str = '',
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Unbind an exchange from another exchange.

    :param destination: The destination exchange to unbind
    :param source: The source exchange to unbind from
    :param routing_key: The routing key to unbind
    :param arguments: Custom key/value pair arguments for the binding
    :returns: Deferred that fires on the Exchange.UnbindOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('exchange_unbind')(
        destination=destination,
        source=source,
        routing_key=routing_key,
        arguments=arguments,
    )

flow

flow(active: bool = True) -> Deferred[Any]

Turn Channel flow control off and on.

Returns a Deferred that will fire with a bool indicating the channel flow state. For more information, please reference:

https://www.rabbitmq.com/amqp-0-9-1-reference.html#channel.flow

PARAMETER DESCRIPTION
active

Turn flow on or off

TYPE: bool DEFAULT: True

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires with the channel flow state

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def flow(self, active: bool = True) -> defer.Deferred[Any]:
    """
    Turn Channel flow control off and on.

    Returns a Deferred that will fire with a bool indicating the channel flow state. For more
    information, please reference:

    https://www.rabbitmq.com/amqp-0-9-1-reference.html#channel.flow

    :param active: Turn flow on or off
    :returns: Deferred that fires with the channel flow state
    :raises ValueError:
    """
    return self._wrap_channel_method('flow')(active=active)

open

open() -> None

Open the channel.

Source code in pika/adapters/twisted_connection.py
def open(self) -> None:
    """Open the channel."""
    return self._channel.open()

queue_bind

queue_bind(
    queue: str,
    exchange: str,
    routing_key: str | None = None,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Bind the queue to the specified exchange.

PARAMETER DESCRIPTION
queue

The queue to bind to the exchange

TYPE: str

exchange

The source exchange to bind to

TYPE: str

routing_key

The routing key to bind on

TYPE: str | None DEFAULT: None

arguments

Custom key/value pair arguments for the binding

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.BindOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_bind(
        self,
        queue: str,
        exchange: str,
        routing_key: str | None = None,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Bind the queue to the specified exchange.

    :param queue: The queue to bind to the exchange
    :param exchange: The source exchange to bind to
    :param routing_key: The routing key to bind on
    :param arguments: Custom key/value pair arguments for the binding
    :returns: Deferred that fires on the Queue.BindOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('queue_bind')(
        queue=queue,
        exchange=exchange,
        routing_key=routing_key,
        arguments=arguments,
    )

queue_declare

queue_declare(
    queue: str,
    passive: bool = False,
    durable: bool = False,
    exclusive: bool = False,
    auto_delete: bool = False,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Declare queue, create if needed.

This method creates or checks a queue. When creating a new queue the client can specify various properties that control the durability of the queue and its contents, and the level of sharing for the queue.

Use an empty string as the queue name for the broker to auto-generate one

PARAMETER DESCRIPTION
queue

The queue name; if empty string, the broker will create a unique queue name

TYPE: str

passive

Only check to see if the queue exists

TYPE: bool DEFAULT: False

durable

Survive reboots of the broker

TYPE: bool DEFAULT: False

exclusive

Only allow access by the current connection

TYPE: bool DEFAULT: False

auto_delete

Delete after consumer cancels or disconnects

TYPE: bool DEFAULT: False

arguments

Custom key/value arguments for the queue

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.DeclareOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_declare(
        self,
        queue: str,
        passive: bool = False,
        durable: bool = False,
        exclusive: bool = False,
        auto_delete: bool = False,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Declare queue, create if needed.

    This method creates or checks a queue. When creating a new queue the client can specify various
    properties that control the durability of the queue and its contents,
    and the level of sharing for the queue.

    Use an empty string as the queue name for the broker to auto-generate
    one

    :param queue: The queue name; if empty string, the broker will
        create a unique queue name
    :param passive: Only check to see if the queue exists
    :param durable: Survive reboots of the broker
    :param exclusive: Only allow access by the current connection
    :param auto_delete: Delete after consumer cancels or disconnects
    :param arguments: Custom key/value arguments for the queue
    :returns: Deferred that fires on the Queue.DeclareOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('queue_declare')(
        queue=queue,
        passive=passive,
        durable=durable,
        exclusive=exclusive,
        auto_delete=auto_delete,
        arguments=arguments,
    )

queue_delete

queue_delete(
    queue: str,
    if_unused: bool = False,
    if_empty: bool = False,
) -> Deferred[Any]

Delete a queue from the broker.

This method wraps :meth:Channel.queue_delete <pika.channel.Channel.queue_delete>, and removes the reference to the queue object after it gets deleted on the server.

PARAMETER DESCRIPTION
queue

The queue to delete

TYPE: str

if_unused

only delete if it's unused

TYPE: bool DEFAULT: False

if_empty

only delete if the queue is empty

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.DeleteOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_delete(self,
                 queue: str,
                 if_unused: bool = False,
                 if_empty: bool = False) -> defer.Deferred[Any]:
    """
    Delete a queue from the broker.

    This method wraps :meth:`Channel.queue_delete <pika.channel.Channel.queue_delete>`, and
    removes the reference to the queue object after it gets deleted on the server.

    :param queue: The queue to delete
    :param if_unused: only delete if it's unused
    :param if_empty: only delete if the queue is empty
    :returns: Deferred that fires on the Queue.DeleteOk response
    :raises ValueError:
    """
    wrapped = self._wrap_channel_method('queue_delete')
    d = wrapped(queue=queue, if_unused=if_unused, if_empty=if_empty)

    def _clear_consumer(ret, queue_name):
        for consumer_tag in list(
                self._queue_name_to_consumer_tags.get(queue_name, set())):
            self._consumers[consumer_tag].close(
                exceptions.ConsumerCancelled(
                    f'Queue {queue_name} was deleted.'))
            del self._consumers[consumer_tag]
            self._queue_name_to_consumer_tags[queue_name].remove(
                consumer_tag)
        return ret

    return d.addCallback(_clear_consumer, queue)

queue_purge

queue_purge(queue: str) -> Deferred[Any]

Purge all of the messages from the specified queue.

PARAMETER DESCRIPTION
queue

The queue to purge

TYPE: str

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.PurgeOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_purge(self, queue: str) -> defer.Deferred[Any]:
    """
    Purge all of the messages from the specified queue.

    :param queue: The queue to purge
    :returns: Deferred that fires on the Queue.PurgeOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('queue_purge')(queue=queue)

queue_unbind

queue_unbind(
    queue: str,
    exchange: str | None,
    routing_key: str | None = None,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Unbind a queue from an exchange.

PARAMETER DESCRIPTION
queue

The queue to unbind from the exchange

TYPE: str

exchange

The source exchange to bind from

TYPE: str | None

routing_key

The routing key to unbind

TYPE: str | None DEFAULT: None

arguments

Custom key/value pair arguments for the binding

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.UnbindOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_unbind(
        self,
        queue: str,
        exchange: str | None,
        routing_key: str | None = None,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Unbind a queue from an exchange.

    :param queue: The queue to unbind from the exchange
    :param exchange: The source exchange to bind from
    :param routing_key: The routing key to unbind
    :param arguments: Custom key/value pair arguments for the binding
    :returns: Deferred that fires on the Queue.UnbindOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('queue_unbind')(
        queue=queue,
        exchange=exchange,
        routing_key=routing_key,
        arguments=arguments,
    )

tx_commit

tx_commit() -> Deferred[Any]

Commit a transaction.

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Tx.CommitOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def tx_commit(self) -> defer.Deferred[Any]:
    """
    Commit a transaction.

    :returns: Deferred that fires on the Tx.CommitOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('tx_commit')()

tx_rollback

tx_rollback() -> Deferred[Any]

Rollback a transaction.

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Tx.RollbackOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def tx_rollback(self) -> defer.Deferred[Any]:
    """
    Rollback a transaction.

    :returns: Deferred that fires on the Tx.RollbackOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('tx_rollback')()

tx_select

tx_select() -> Deferred[Any]

Select standard transaction mode.

This method sets the channel to use standard transactions. The client must use this method at least once on a channel before using the Commit or Rollback methods.

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Tx.SelectOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def tx_select(self) -> defer.Deferred[Any]:
    """
    Select standard transaction mode.

    This method sets the channel to use standard transactions. The client must use this method
    at least once on a channel before using the Commit or Rollback methods.

    :returns: Deferred that fires on the Tx.SelectOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('tx_select')()

TwistedProtocolConnection

Bases: Protocol

A Pika-specific implementation of a Twisted Protocol.

Allows using Twisted's non-blocking connectTCP/connectSSL methods for connecting to the server.

TwistedProtocolConnection objects have a ready instance variable that's a Deferred which fires when the connection is ready to be used (the initial AMQP handshaking has been done). You have to wait for this Deferred to fire before requesting a channel.

Once the connection is ready, you will be able to use the closed instance variable: a Deferred which fires when the connection is closed.

Since it's Twisted handling connection establishing it does not accept connect callbacks, you have to implement that within Twisted. Also remember that the host, port and ssl values of the connection parameters are ignored because, yet again, it's Twisted who manages the connection.

Source code in pika/adapters/twisted_connection.py
class TwistedProtocolConnection(protocol.Protocol):
    """
    A Pika-specific implementation of a Twisted Protocol.

    Allows using Twisted's non-blocking connectTCP/connectSSL methods for connecting to the
    server.

    TwistedProtocolConnection objects have a `ready` instance variable that's a
    Deferred which fires when the connection is ready to be used (the initial
    AMQP handshaking has been done). You *have* to wait for this Deferred to
    fire before requesting a channel.

    Once the connection is ready, you will be able to use the `closed` instance
    variable: a Deferred which fires when the connection is closed.

    Since it's Twisted handling connection establishing it does not accept
    connect callbacks, you have to implement that within Twisted. Also remember
    that the host, port and ssl values of the connection parameters are ignored
    because, yet again, it's Twisted who manages the connection.
    """

    def __init__(self,
                 parameters: pika.connection.ConnectionParameters | None = None,
                 custom_reactor: Any = None) -> None:
        warnings.warn(
            'TwistedProtocolConnection is deprecated and will be removed in '
            'Pika 2.0. Use Connection instead, which works with any '
            'framework including Twisted. See '
            'https://pika.github.io/pika/latest/modules/adapters/thread_safe/',
            DeprecationWarning,
            stacklevel=2,
        )
        self.ready: defer.Deferred[Any] | None = defer.Deferred()
        self.ready.addCallback(lambda _: self.connectionReady())
        self.closed: defer.Deferred[Any] | None = None
        self._impl = _TwistedConnectionAdapter(
            parameters=parameters,
            on_open_callback=self._on_connection_ready,
            on_open_error_callback=self._on_connection_failed,
            on_close_callback=self._on_connection_closed,
            custom_reactor=custom_reactor,
        )
        self._calls: set[defer.Deferred[Any]] = set()

    def channel(self,
                channel_number: int | None = None) -> Deferred[TwistedChannel]:
        """
        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.
        :returns: a Deferred that fires with an instance of a wrapper around the Pika Channel class.
        """
        d: defer.Deferred[Any] = defer.Deferred()
        self._impl.channel(channel_number, d.callback)
        self._calls.add(d)
        d.addCallback(self._clear_call, d)
        return d.addCallback(TwistedChannel)

    @property
    def is_open(self) -> bool:
        # For compatibility with previous releases.
        return self._impl.is_open

    @property
    def is_closed(self) -> bool:
        # For compatibility with previous releases.
        return self._impl.is_closed

    def close(
            self,
            reply_code: int = 200,
            reply_text: str = 'Normal shutdown') -> defer.Deferred[Any] | None:
        if not self._impl.is_closed:
            self._impl.close(reply_code, reply_text)
        return self.closed

    # IProtocol methods

    @override
    def dataReceived(self, data: bytes) -> None:
        # Pass the bytes to Pika for parsing
        self._impl.data_received(data)

    @override
    def connectionLost(
        self,
        reason: twisted.python.failure.Failure = protocol.connectionDone
    ) -> None:
        self._impl.connection_lost(reason)
        # Let the caller know there's been an error
        d, self.ready = self.ready, None
        if d:
            d.errback(reason)

    @override
    def makeConnection(
            self, transport: twisted.internet.interfaces.ITransport) -> None:
        self._impl.connection_made(transport)
        protocol.Protocol.makeConnection(self, transport)

    # Our own methods

    def connectionReady(
        self
    ) -> (TwistedProtocolConnection | defer.Deferred[TwistedProtocolConnection]
         ):
        """This method will be called when the underlying connection is ready."""
        return self

    def _on_connection_ready(self,
                             _connection: pika.connection.Connection) -> None:
        d, self.ready = self.ready, None
        if d:
            self.closed = defer.Deferred()
            d.callback(None)

    def _on_connection_failed(self,
                              _connection: pika.connection.Connection,
                              _error_message: Exception | None = None) -> None:
        d, self.ready = self.ready, None
        if d:
            attempts = self._impl.params.connection_attempts
            exc = exceptions.AMQPConnectionError(attempts,
                                                 host=self._impl.params.host,
                                                 port=self._impl.params.port)
            d.errback(exc)

    def _on_connection_closed(self, _connection: pika.connection.Connection,
                              exception: Exception | Failure) -> None:
        # `Failure` is not an `Exception` subclass, so it has to be named in the
        # annotation. `Connection` itself only ever passes an `Exception`, but
        # this stays callable with a `Failure` for the branch below.
        #
        # `d` is a separate name from the loop variable because `self.closed` is
        # optional while the pending calls are not.
        for call in self._calls:
            call.errback(exception)
        self._calls = set()

        d, self.closed = self.closed, None
        if d:
            if isinstance(exception, Failure):
                # Calling `callback` with a Failure instance would trigger the
                # errback path, so unwrap it first.
                d.callback(exception.value)
            else:
                d.callback(exception)

    def _clear_call(self, ret, d):
        self._calls.discard(d)
        return ret

channel

channel(
    channel_number: int | None = None,
) -> Deferred[TwistedChannel]

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

PARAMETER DESCRIPTION
channel_number

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

TYPE: int | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[TwistedChannel]

a Deferred that fires with an instance of a wrapper around the Pika Channel class.

Source code in pika/adapters/twisted_connection.py
def channel(self,
            channel_number: int | None = None) -> Deferred[TwistedChannel]:
    """
    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.
    :returns: a Deferred that fires with an instance of a wrapper around the Pika Channel class.
    """
    d: defer.Deferred[Any] = defer.Deferred()
    self._impl.channel(channel_number, d.callback)
    self._calls.add(d)
    d.addCallback(self._clear_call, d)
    return d.addCallback(TwistedChannel)

connectionReady

connectionReady() -> (
    TwistedProtocolConnection
    | Deferred[TwistedProtocolConnection]
)

This method will be called when the underlying connection is ready.

Source code in pika/adapters/twisted_connection.py
def connectionReady(
    self
) -> (TwistedProtocolConnection | defer.Deferred[TwistedProtocolConnection]
     ):
    """This method will be called when the underlying connection is ready."""
    return self

Class Reference

TwistedProtocolConnection

Bases: Protocol

A Pika-specific implementation of a Twisted Protocol.

Allows using Twisted's non-blocking connectTCP/connectSSL methods for connecting to the server.

TwistedProtocolConnection objects have a ready instance variable that's a Deferred which fires when the connection is ready to be used (the initial AMQP handshaking has been done). You have to wait for this Deferred to fire before requesting a channel.

Once the connection is ready, you will be able to use the closed instance variable: a Deferred which fires when the connection is closed.

Since it's Twisted handling connection establishing it does not accept connect callbacks, you have to implement that within Twisted. Also remember that the host, port and ssl values of the connection parameters are ignored because, yet again, it's Twisted who manages the connection.

Source code in pika/adapters/twisted_connection.py
class TwistedProtocolConnection(protocol.Protocol):
    """
    A Pika-specific implementation of a Twisted Protocol.

    Allows using Twisted's non-blocking connectTCP/connectSSL methods for connecting to the
    server.

    TwistedProtocolConnection objects have a `ready` instance variable that's a
    Deferred which fires when the connection is ready to be used (the initial
    AMQP handshaking has been done). You *have* to wait for this Deferred to
    fire before requesting a channel.

    Once the connection is ready, you will be able to use the `closed` instance
    variable: a Deferred which fires when the connection is closed.

    Since it's Twisted handling connection establishing it does not accept
    connect callbacks, you have to implement that within Twisted. Also remember
    that the host, port and ssl values of the connection parameters are ignored
    because, yet again, it's Twisted who manages the connection.
    """

    def __init__(self,
                 parameters: pika.connection.ConnectionParameters | None = None,
                 custom_reactor: Any = None) -> None:
        warnings.warn(
            'TwistedProtocolConnection is deprecated and will be removed in '
            'Pika 2.0. Use Connection instead, which works with any '
            'framework including Twisted. See '
            'https://pika.github.io/pika/latest/modules/adapters/thread_safe/',
            DeprecationWarning,
            stacklevel=2,
        )
        self.ready: defer.Deferred[Any] | None = defer.Deferred()
        self.ready.addCallback(lambda _: self.connectionReady())
        self.closed: defer.Deferred[Any] | None = None
        self._impl = _TwistedConnectionAdapter(
            parameters=parameters,
            on_open_callback=self._on_connection_ready,
            on_open_error_callback=self._on_connection_failed,
            on_close_callback=self._on_connection_closed,
            custom_reactor=custom_reactor,
        )
        self._calls: set[defer.Deferred[Any]] = set()

    def channel(self,
                channel_number: int | None = None) -> Deferred[TwistedChannel]:
        """
        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.
        :returns: a Deferred that fires with an instance of a wrapper around the Pika Channel class.
        """
        d: defer.Deferred[Any] = defer.Deferred()
        self._impl.channel(channel_number, d.callback)
        self._calls.add(d)
        d.addCallback(self._clear_call, d)
        return d.addCallback(TwistedChannel)

    @property
    def is_open(self) -> bool:
        # For compatibility with previous releases.
        return self._impl.is_open

    @property
    def is_closed(self) -> bool:
        # For compatibility with previous releases.
        return self._impl.is_closed

    def close(
            self,
            reply_code: int = 200,
            reply_text: str = 'Normal shutdown') -> defer.Deferred[Any] | None:
        if not self._impl.is_closed:
            self._impl.close(reply_code, reply_text)
        return self.closed

    # IProtocol methods

    @override
    def dataReceived(self, data: bytes) -> None:
        # Pass the bytes to Pika for parsing
        self._impl.data_received(data)

    @override
    def connectionLost(
        self,
        reason: twisted.python.failure.Failure = protocol.connectionDone
    ) -> None:
        self._impl.connection_lost(reason)
        # Let the caller know there's been an error
        d, self.ready = self.ready, None
        if d:
            d.errback(reason)

    @override
    def makeConnection(
            self, transport: twisted.internet.interfaces.ITransport) -> None:
        self._impl.connection_made(transport)
        protocol.Protocol.makeConnection(self, transport)

    # Our own methods

    def connectionReady(
        self
    ) -> (TwistedProtocolConnection | defer.Deferred[TwistedProtocolConnection]
         ):
        """This method will be called when the underlying connection is ready."""
        return self

    def _on_connection_ready(self,
                             _connection: pika.connection.Connection) -> None:
        d, self.ready = self.ready, None
        if d:
            self.closed = defer.Deferred()
            d.callback(None)

    def _on_connection_failed(self,
                              _connection: pika.connection.Connection,
                              _error_message: Exception | None = None) -> None:
        d, self.ready = self.ready, None
        if d:
            attempts = self._impl.params.connection_attempts
            exc = exceptions.AMQPConnectionError(attempts,
                                                 host=self._impl.params.host,
                                                 port=self._impl.params.port)
            d.errback(exc)

    def _on_connection_closed(self, _connection: pika.connection.Connection,
                              exception: Exception | Failure) -> None:
        # `Failure` is not an `Exception` subclass, so it has to be named in the
        # annotation. `Connection` itself only ever passes an `Exception`, but
        # this stays callable with a `Failure` for the branch below.
        #
        # `d` is a separate name from the loop variable because `self.closed` is
        # optional while the pending calls are not.
        for call in self._calls:
            call.errback(exception)
        self._calls = set()

        d, self.closed = self.closed, None
        if d:
            if isinstance(exception, Failure):
                # Calling `callback` with a Failure instance would trigger the
                # errback path, so unwrap it first.
                d.callback(exception.value)
            else:
                d.callback(exception)

    def _clear_call(self, ret, d):
        self._calls.discard(d)
        return ret

_calls instance-attribute

_calls: set[Deferred[Any]] = set()

_impl instance-attribute

_impl = _TwistedConnectionAdapter(
    parameters=parameters,
    on_open_callback=self._on_connection_ready,
    on_open_error_callback=self._on_connection_failed,
    on_close_callback=self._on_connection_closed,
    custom_reactor=custom_reactor,
)

closed instance-attribute

closed: Deferred[Any] | None = None

is_closed property

is_closed: bool

is_open property

is_open: bool

ready instance-attribute

ready: Deferred[Any] | None = defer.Deferred()

__init__

__init__(
    parameters: ConnectionParameters | None = None,
    custom_reactor: Any = None,
) -> None
Source code in pika/adapters/twisted_connection.py
def __init__(self,
             parameters: pika.connection.ConnectionParameters | None = None,
             custom_reactor: Any = None) -> None:
    warnings.warn(
        'TwistedProtocolConnection is deprecated and will be removed in '
        'Pika 2.0. Use Connection instead, which works with any '
        'framework including Twisted. See '
        'https://pika.github.io/pika/latest/modules/adapters/thread_safe/',
        DeprecationWarning,
        stacklevel=2,
    )
    self.ready: defer.Deferred[Any] | None = defer.Deferred()
    self.ready.addCallback(lambda _: self.connectionReady())
    self.closed: defer.Deferred[Any] | None = None
    self._impl = _TwistedConnectionAdapter(
        parameters=parameters,
        on_open_callback=self._on_connection_ready,
        on_open_error_callback=self._on_connection_failed,
        on_close_callback=self._on_connection_closed,
        custom_reactor=custom_reactor,
    )
    self._calls: set[defer.Deferred[Any]] = set()

_clear_call

_clear_call(ret, d)
Source code in pika/adapters/twisted_connection.py
def _clear_call(self, ret, d):
    self._calls.discard(d)
    return ret

_on_connection_closed

_on_connection_closed(
    _connection: Connection, exception: Exception | Failure
) -> None
Source code in pika/adapters/twisted_connection.py
def _on_connection_closed(self, _connection: pika.connection.Connection,
                          exception: Exception | Failure) -> None:
    # `Failure` is not an `Exception` subclass, so it has to be named in the
    # annotation. `Connection` itself only ever passes an `Exception`, but
    # this stays callable with a `Failure` for the branch below.
    #
    # `d` is a separate name from the loop variable because `self.closed` is
    # optional while the pending calls are not.
    for call in self._calls:
        call.errback(exception)
    self._calls = set()

    d, self.closed = self.closed, None
    if d:
        if isinstance(exception, Failure):
            # Calling `callback` with a Failure instance would trigger the
            # errback path, so unwrap it first.
            d.callback(exception.value)
        else:
            d.callback(exception)

_on_connection_failed

_on_connection_failed(
    _connection: Connection,
    _error_message: Exception | None = None,
) -> None
Source code in pika/adapters/twisted_connection.py
def _on_connection_failed(self,
                          _connection: pika.connection.Connection,
                          _error_message: Exception | None = None) -> None:
    d, self.ready = self.ready, None
    if d:
        attempts = self._impl.params.connection_attempts
        exc = exceptions.AMQPConnectionError(attempts,
                                             host=self._impl.params.host,
                                             port=self._impl.params.port)
        d.errback(exc)

_on_connection_ready

_on_connection_ready(_connection: Connection) -> None
Source code in pika/adapters/twisted_connection.py
def _on_connection_ready(self,
                         _connection: pika.connection.Connection) -> None:
    d, self.ready = self.ready, None
    if d:
        self.closed = defer.Deferred()
        d.callback(None)

channel

channel(
    channel_number: int | None = None,
) -> Deferred[TwistedChannel]

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

PARAMETER DESCRIPTION
channel_number

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

TYPE: int | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[TwistedChannel]

a Deferred that fires with an instance of a wrapper around the Pika Channel class.

Source code in pika/adapters/twisted_connection.py
def channel(self,
            channel_number: int | None = None) -> Deferred[TwistedChannel]:
    """
    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.
    :returns: a Deferred that fires with an instance of a wrapper around the Pika Channel class.
    """
    d: defer.Deferred[Any] = defer.Deferred()
    self._impl.channel(channel_number, d.callback)
    self._calls.add(d)
    d.addCallback(self._clear_call, d)
    return d.addCallback(TwistedChannel)

close

close(
    reply_code: int = 200,
    reply_text: str = 'Normal shutdown',
) -> Deferred[Any] | None
Source code in pika/adapters/twisted_connection.py
def close(
        self,
        reply_code: int = 200,
        reply_text: str = 'Normal shutdown') -> defer.Deferred[Any] | None:
    if not self._impl.is_closed:
        self._impl.close(reply_code, reply_text)
    return self.closed

connectionLost

connectionLost(reason: Failure = connectionDone) -> None
Source code in pika/adapters/twisted_connection.py
@override
def connectionLost(
    self,
    reason: twisted.python.failure.Failure = protocol.connectionDone
) -> None:
    self._impl.connection_lost(reason)
    # Let the caller know there's been an error
    d, self.ready = self.ready, None
    if d:
        d.errback(reason)

connectionReady

connectionReady() -> (
    TwistedProtocolConnection
    | Deferred[TwistedProtocolConnection]
)

This method will be called when the underlying connection is ready.

Source code in pika/adapters/twisted_connection.py
def connectionReady(
    self
) -> (TwistedProtocolConnection | defer.Deferred[TwistedProtocolConnection]
     ):
    """This method will be called when the underlying connection is ready."""
    return self

dataReceived

dataReceived(data: bytes) -> None
Source code in pika/adapters/twisted_connection.py
@override
def dataReceived(self, data: bytes) -> None:
    # Pass the bytes to Pika for parsing
    self._impl.data_received(data)

makeConnection

makeConnection(transport: ITransport) -> None
Source code in pika/adapters/twisted_connection.py
@override
def makeConnection(
        self, transport: twisted.internet.interfaces.ITransport) -> None:
    self._impl.connection_made(transport)
    protocol.Protocol.makeConnection(self, transport)

TwistedChannel

A wrapper around Pika's Channel.

Channel methods that normally take a callback argument are wrapped to return a Deferred that fires with whatever would be passed to the callback. If the channel gets closed, all pending Deferreds are errbacked with a ChannelClosed exception. The returned Deferreds fire with whatever arguments the callback to the original method would receive.

Some methods like basic_consume and basic_get are wrapped in a special way, see their docstrings for details.

Source code in pika/adapters/twisted_connection.py
class TwistedChannel:
    """
    A wrapper around Pika's Channel.

    Channel methods that normally take a callback argument are wrapped to return a Deferred that
    fires with whatever would be passed to the callback. If the channel gets closed, all pending
    Deferreds are errbacked with a ChannelClosed exception. The returned Deferreds fire with
    whatever arguments the callback to the original method would receive.

    Some methods like basic_consume and basic_get are wrapped in a special way, see their docstrings
    for details.
    """

    def __init__(self, channel: channel.Channel) -> None:
        self._channel = channel
        self._closed: Exception | None = None
        self._calls: set[defer.Deferred[Any]] = set()
        self._consumers: dict[str, ClosableDeferredQueue] = {}
        # Store Basic.Get calls so we can handle GetEmpty replies
        self._basic_get_deferred: defer.Deferred[Any] | None = None
        self._channel.add_callback(self._on_getempty, [spec.Basic.GetEmpty],
                                   False)
        # We need this mapping to close the ClosableDeferredQueue when a queue
        # is deleted.
        self._queue_name_to_consumer_tags: dict[str, set[str]] = {}
        # Whether RabbitMQ delivery confirmation has been enabled
        self._delivery_confirmation = False
        self._delivery_message_id: int | None = None
        self._deliveries: dict[int, defer.Deferred[Any]] = {}
        # Holds a ReceivedMessage object representing a message received via
        # Basic.Return in publisher-acknowledgments mode.
        self._puback_return: ReceivedMessage | None = None

        self.on_closed: defer.Deferred[Any] = defer.Deferred()
        self._channel.add_on_close_callback(self._on_channel_closed)
        self._channel.add_on_cancel_callback(
            self._on_consumer_cancelled_by_broker)

    @override
    def __repr__(self) -> str:
        return f'<{self.__class__.__name__} channel={self._channel!r}>'

    def _on_channel_closed(self, _channel: channel.Channel,
                           reason: Exception | None) -> None:
        # enter the closed state
        self._closed = reason
        # errback all pending calls
        for d in self._calls:
            d.errback(self._closed)
        # errback all pending deliveries
        for d in self._deliveries.values():
            d.errback(self._closed)
        # close all open queues
        for consumer in self._consumers.values():
            consumer.close(self._closed)
        # release references to stored objects
        self._calls = set()
        self._deliveries = {}
        self._consumers = {}
        self.on_closed.callback(self._closed)

    def _on_consumer_cancelled_by_broker(
        self, method_frame: pika.frame.Method[pika.spec.Basic.Cancel]
    ) -> (pika.frame.Method[pika.spec.Basic.Cancel] |
          pika.frame.Method[pika.spec.Basic.CancelOk]):
        """
        Called by impl when broker cancels consumer via Basic.Cancel.

        This is a RabbitMQ-specific feature. The circumstances include deletion of queue being
        consumed as well as failure of a HA node responsible for the queue being consumed.

        :param method_frame: method frame with the `spec.Basic.Cancel` method
        """
        return self._on_consumer_cancelled(method_frame)

    def _on_consumer_cancelled(
        self, frame: (pika.frame.Method[pika.spec.Basic.Cancel] |
                      pika.frame.Method[pika.spec.Basic.CancelOk])
    ) -> (pika.frame.Method[pika.spec.Basic.Cancel] |
          pika.frame.Method[pika.spec.Basic.CancelOk]):
        """
        Called when the broker cancels a consumer via Basic.Cancel or when the broker responds to a
        Basic.Cancel request by Basic.CancelOk.

        :param frame: method frame with the `spec.Basic.Cancel` or `spec.Basic.CancelOk` method
        """
        consumer_tag = frame.method.consumer_tag
        if consumer_tag not in self._consumers:
            # Could be cancelled by user or broker earlier
            LOGGER.warning('basic_cancel - consumer not found: %s',
                           consumer_tag)
            return frame
        self._consumers[consumer_tag].close(exceptions.ConsumerCancelled())
        del self._consumers[consumer_tag]
        # Remove from the queue-to-ctags index:
        for ctags in self._queue_name_to_consumer_tags.values():
            try:
                ctags.remove(consumer_tag)
            except KeyError:
                continue
        return frame

    def _on_getempty(
            self,
            _method_frame: pika.frame.Method[pika.spec.Basic.Get]) -> None:
        """
        Callback the Basic.Get deferred with None.

        :param _method_frame: Method frame from Basic.Get response (unused)
        """
        if self._basic_get_deferred is None:
            LOGGER.warning('Got Basic.GetEmpty but no Basic.Get calls '
                           'were pending.')
            return
        self._basic_get_deferred.callback(None)

    def _wrap_channel_method(self,
                             name: str) -> Callable[..., defer.Deferred[Any]]:
        """
        Wrap Pika's Channel method to make it return a Deferred that fires when the method completes
        and errbacks if the channel gets closed. If the original method's callback would receive
        more than one argument, the Deferred fires with a tuple of argument values.

        :param name: Attribute name to look up on the underlying channel
        """
        method = getattr(self._channel, name)

        @functools.wraps(method)
        def wrapped(*args, **kwargs):
            if self._closed:
                return defer.fail(self._closed)

            d: defer.Deferred = defer.Deferred()
            self._calls.add(d)
            d.addCallback(self._clear_call, d)

            def single_argument(*args) -> None:
                """
                Make sure that the deferred is called with a single argument.

                In case the original callback fires with more than one, convert to a tuple.
                """
                if len(args) > 1:
                    d.callback(tuple(args))
                else:
                    d.callback(*args)

            kwargs['callback'] = single_argument

            try:
                method(*args, **kwargs)
            except Exception:
                return defer.fail()
            return d

        return wrapped

    def _clear_call(self, ret: Any, d: defer.Deferred[Any]) -> Any:
        self._calls.discard(d)
        return ret

    # Public Channel attributes

    @property
    def channel_number(self) -> int:
        return self._channel.channel_number

    @property
    def connection(self) -> pika.connection.Connection:
        return self._channel.connection

    @property
    def is_closed(self) -> bool:
        """Returns True if the channel is closed."""
        return self._channel.is_closed

    @property
    def is_closing(self) -> bool:
        """Returns True if client-initiated closing of the channel is in progress."""
        return self._channel.is_closing

    @property
    def is_open(self) -> bool:
        """Returns True if the channel is open."""
        return self._channel.is_open

    @property
    def flow_active(self) -> bool:
        return self._channel.flow_active

    @property
    def consumer_tags(self) -> list[str]:
        return self._channel.consumer_tags

    # Deferred-equivalents of public Channel methods

    def callback_deferred(self, deferred: defer.Deferred[Any],
                          replies: Sequence[type[amqp_object.Method]]) -> None:
        """
        Pass in a Deferred and a list replies from the RabbitMQ broker which you'd like the Deferred
        to be callbacked on with the frame as callback value.

        :param deferred: The Deferred to callback
        :param replies: The replies to callback on
        """
        self._channel.add_callback(deferred.callback, replies)

    # Public Channel methods

    def add_on_return_callback(
            self, callback: Callable[[ReceivedMessage], None]) -> None:
        """
        Pass a callback function that will be called when a published message is rejected and
        returned by the server via `Basic.Return`.

        :param callback: The method to call on callback with the
            message as only argument. The message is a named tuple with
            the following attributes
            - channel: this TwistedChannel
            - method: pika.spec.Basic.Return
            - properties: pika.spec.BasicProperties
            - body: bytes
        """
        self._channel.add_on_return_callback(
            lambda _channel, method, properties, body: callback(
                ReceivedMessage(
                    channel=self,
                    method=method,
                    properties=properties,
                    body=body,
                )))

    def basic_ack(self, delivery_tag: int = 0, multiple: bool = False) -> None:
        """
        Acknowledge one or more messages.

        When sent by the client, this method acknowledges one or more messages delivered via the
        Deliver or Get-Ok methods. When sent by server, this method acknowledges one or more
        messages published with the Publish method on a channel in confirm mode. The acknowledgement
        can be for a single message or a set of messages up to and including a specific message.

        :param delivery_tag: int/long The server-assigned delivery tag
        :param multiple: If set to True, the delivery tag is treated as "up to and including", so
            that multiple messages can be acknowledged with a single method. If set to False, the
            delivery tag refers to a single message. If the multiple field is 1, and the delivery
            tag is zero, this indicates acknowledgement of all outstanding messages.
        """
        return self._channel.basic_ack(delivery_tag=delivery_tag,
                                       multiple=multiple)

    def basic_cancel(
        self,
        consumer_tag: str = ''
    ) -> defer.Deferred[(pika.frame.Method[pika.spec.Basic.CancelOk] |
                         pika.frame.Method[pika.spec.Basic.Cancel])]:
        """
        This method cancels a consumer.

        This does not affect already delivered messages, but it does mean the server will not send
        any more messages for that consumer. The client may receive an arbitrary number of messages
        in between sending the cancel method and receiving the cancel-ok reply. It may also be sent
        from the server to the client in the event of the consumer being unexpectedly cancelled
        (i.e. cancelled for any reason other than the server receiving the corresponding
        basic.cancel from the client). This allows clients to be notified of the loss of consumers
        due to events such as queue deletion.

        This method wraps :meth:`Channel.basic_cancel <pika.channel.Channel.basic_cancel>` and
        closes any deferred queue associated with that consumer.

        :param consumer_tag: Identifier for the consumer
        :returns: Deferred that fires on the Basic.CancelOk response
        :raises ValueError:
        """
        wrapped = self._wrap_channel_method('basic_cancel')
        d = wrapped(consumer_tag=consumer_tag)
        return d.addCallback(self._on_consumer_cancelled)

    def basic_consume(
            self,
            queue: str,
            auto_ack: bool = False,
            exclusive: bool = False,
            consumer_tag: str | None = None,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Consume from a server queue.

        Sends the AMQP 0-9-1 command Basic.Consume to the broker and binds
        messages for the consumer_tag to a
        :class:`ClosableDeferredQueue`. If you do not pass in a
        consumer_tag, one will be automatically generated for you.

        For more information on basic_consume, see:
        Tutorial 2 at https://www.rabbitmq.com/getstarted.html
        https://www.rabbitmq.com/confirms.html
        https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.consume

        :param queue: The queue to consume from. Use the empty string to
            specify the most recent server-named queue for this channel.
        :param auto_ack: if set to True, automatic acknowledgement mode
            will be used (see https://www.rabbitmq.com/confirms.html). This
            corresponds with the 'no_ack' parameter in the basic.consume AMQP
            0.9.1 method
        :param exclusive: Don't allow other consumers on the queue
        :param consumer_tag: Specify your own consumer tag
        :param arguments: Custom key/value pair arguments for the consumer
        :returns: Deferred that fires with a tuple
            ``(queue_object, consumer_tag)``. The Deferred will errback with an
            instance of :class:`exceptions.ChannelClosed` if the call fails.
            The queue object is an instance of :class:`ClosableDeferredQueue`,
            where data received from the queue will be stored. Clients should
            use its :meth:`get() <ClosableDeferredQueue.get>` method to fetch
            an individual message, which will return a Deferred firing with a
            namedtuple whose attributes are:
            - channel: this TwistedChannel
            - method: pika.spec.Basic.Deliver
            - properties: pika.spec.BasicProperties
            - body: bytes
        """
        if self._closed:
            return defer.fail(self._closed)

        queue_obj = ClosableDeferredQueue()
        d: defer.Deferred[Any] = defer.Deferred()
        self._calls.add(d)

        def on_consume_ok(frame) -> None:
            consumer_tag = frame.method.consumer_tag
            self._queue_name_to_consumer_tags.setdefault(
                queue, set()).add(consumer_tag)
            self._consumers[consumer_tag] = queue_obj
            self._calls.discard(d)
            d.callback((queue_obj, consumer_tag))

        def on_message_callback(_channel, method, properties, body) -> None:
            """Add the ReceivedMessage to the queue, while replacing the channel implementation."""
            queue_obj.put(
                ReceivedMessage(
                    channel=self,
                    method=method,
                    properties=properties,
                    body=body,
                ))

        try:
            self._channel.basic_consume(
                queue=queue,
                on_message_callback=on_message_callback,
                auto_ack=auto_ack,
                exclusive=exclusive,
                consumer_tag=consumer_tag,
                arguments=arguments,
                callback=on_consume_ok,
            )
        except Exception:
            return defer.fail()

        return d

    def basic_get(self,
                  queue: str,
                  auto_ack: bool = False) -> defer.Deferred[Any]:
        """
        Get a single message from the AMQP broker.

        Will return If the queue is empty, it will return None.
        If you want to
        be notified of Basic.GetEmpty, use the Channel.add_callback method
        adding your Basic.GetEmpty callback which should expect only one
        parameter, frame. Due to implementation details, this cannot be called
        a second time until the callback is executed.  For more information on
        basic_get and its parameters, see:

        https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.get

        This method wraps :meth:`Channel.basic_get
        <pika.channel.Channel.basic_get>`.

        :param queue: The queue from which to get a message. Use the empty
                      string to specify the most recent server-named queue
                      for this channel.
        :param auto_ack: Tell the broker to not expect a reply
        :returns: Deferred that fires with a namedtuple whose attributes are:
             - channel: this TwistedChannel
             - method: pika.spec.Basic.GetOk
             - properties: pika.spec.BasicProperties
             - body: bytes
            If the queue is empty, None will be returned.
        :raises pika.exceptions.DuplicateGetOkCallback:
        """
        if self._basic_get_deferred is not None:
            raise exceptions.DuplicateGetOkCallback()

        def create_namedtuple(result) -> ReceivedMessage | None:
            if result is None:
                return None
            _channel, method, properties, body = result
            return ReceivedMessage(
                channel=self,
                method=method,
                properties=properties,
                body=body,
            )

        def cleanup_attribute(result):
            self._basic_get_deferred = None
            return result

        d = self._wrap_channel_method('basic_get')(queue=queue,
                                                   auto_ack=auto_ack)
        d.addCallback(create_namedtuple)
        d.addBoth(cleanup_attribute)
        self._basic_get_deferred = d
        return d

    def basic_nack(self,
                   delivery_tag: int = 0,
                   multiple: bool = False,
                   requeue: bool = True) -> None:
        """
        This method allows a client to reject one or more incoming messages.

        It can be used to interrupt and cancel large incoming messages, or return untreatable
        messages to their original queue.

        :param delivery_tag: int/long The server-assigned delivery tag
        :param multiple: If set to True, the delivery tag is treated as "up to and including", so
            that multiple messages can be acknowledged with a single method. If set to False, the
            delivery tag refers to a single message. If the multiple field is 1, and the delivery
            tag is zero, this indicates acknowledgement of all outstanding messages.
        :param requeue: If requeue is true, the server will attempt to requeue the message. If
            requeue is false or the requeue attempt fails the messages are discarded or dead-
            lettered.
        """
        return self._channel.basic_nack(
            delivery_tag=delivery_tag,
            multiple=multiple,
            requeue=requeue,
        )

    def basic_publish(self,
                      exchange: str,
                      routing_key: str,
                      body: bytes,
                      properties: spec.BasicProperties | None = None,
                      mandatory: bool = False) -> defer.Deferred[Any]:
        """
        Publish to the channel with the given exchange, routing key and body.

        This method wraps :meth:`Channel.basic_publish <pika.channel.Channel.basic_publish>`, but
        makes sure the channel is not closed before publishing.

        For more information on basic_publish and what the parameters do, see:

        https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.publish

        :param exchange: The exchange to publish to
        :param routing_key: The routing key to bind on
        :param body: The message body
        :param properties: Basic.properties
        :param mandatory: The mandatory flag
        :returns: A Deferred that fires with the result of the channel's basic_publish.
        :raises UnroutableError: raised when a message published in publisher-acknowledgments mode
            (see `BlockingChannel.confirm_delivery`) is returned via `Basic.Return` followed by
            `Basic.Ack`.
        :raises NackError: raised when a message published in publisher-acknowledgements mode is
            Nack'ed by the broker. See `BlockingChannel.confirm_delivery`.
        """
        if self._closed:
            return defer.fail(self._closed)
        self._channel.basic_publish(exchange=exchange,
                                    routing_key=routing_key,
                                    body=body,
                                    properties=properties,
                                    mandatory=mandatory)
        if not self._delivery_confirmation:
            # `Channel.basic_publish()` returns nothing, so without delivery
            # confirmation there is no result to report beyond completion.
            return defer.succeed(None)
        # See https://www.rabbitmq.com/confirms.html#publisher-confirms
        assert self._delivery_message_id is not None
        self._delivery_message_id += 1
        self._deliveries[self._delivery_message_id] = defer.Deferred()
        return self._deliveries[self._delivery_message_id]

    def basic_qos(self,
                  prefetch_size: int = 0,
                  prefetch_count: int = 0,
                  global_qos: bool = False) -> defer.Deferred[Any]:
        """
        Specify quality of service.

        This method requests a specific quality of service. The QoS can be specified for the current
        channel or for all channels on the connection. The client can request that messages be sent
        in advance so that when the client finishes processing a message, the following message is
        already held locally, rather than needing to be sent down the channel. Prefetching gives a
        performance improvement.

        :param prefetch_size: This field specifies the prefetch window size. The server will send a
            message in advance if it is equal to or smaller in size than the available prefetch size
            (and also falls into other prefetch limits). May be set to zero, meaning "no specific
            limit", although other prefetch limits may still apply. The prefetch-size is ignored by
            consumers who have enabled the no-ack option.
        :param prefetch_count: Specifies a prefetch window in terms of whole messages. This field
            may be used in combination with the prefetch-size field; a message will only be sent in
            advance if both prefetch windows (and those at the channel and connection level) allow
            it. The prefetch-count is ignored by consumers who have enabled the no-ack option.
        :param global_qos: Should the QoS apply to all channels on the connection.
        :returns: Deferred that fires on the Basic.QosOk response
        """
        return self._wrap_channel_method('basic_qos')(
            prefetch_size=prefetch_size,
            prefetch_count=prefetch_count,
            global_qos=global_qos,
        )

    def basic_reject(self, delivery_tag: int, requeue: bool = True) -> None:
        """
        Reject an incoming message.

        This method allows a client to reject a message. It can be used to interrupt and cancel
        large incoming messages, or return untreatable messages to their original queue.

        :param delivery_tag: int/long The server-assigned delivery tag
        :param requeue: If requeue is true, the server will attempt to requeue the message. If
            requeue is false or the requeue attempt fails the messages are discarded or dead-
            lettered.
        :raises TypeError:
        """
        return self._channel.basic_reject(delivery_tag=delivery_tag,
                                          requeue=requeue)

    def basic_recover(self, requeue: bool = False) -> defer.Deferred[Any]:
        """
        This method asks the server to redeliver all unacknowledged messages on a specified channel.
        Zero or more messages may be redelivered. This method replaces the asynchronous Recover.

        :param requeue: If False, the message will be redelivered to the original recipient. If
            True, the server will attempt to requeue the message, potentially then delivering it to
            an alternative subscriber.
        :returns: Deferred that fires on the Basic.RecoverOk response
        """
        return self._wrap_channel_method('basic_recover')(requeue=requeue)

    def close(self,
              reply_code: int = 0,
              reply_text: str = 'Normal shutdown') -> None:
        """
        Invoke a graceful shutdown of the channel with the AMQP Broker.

        If channel is OPENING, transition to CLOSING and suppress the incoming Channel.OpenOk, if
        any.

        :param reply_code: The reason code to send to broker
        :param reply_text: The reason text to send to broker
        :raises ChannelWrongStateError: if channel is closed or closing
        """
        return self._channel.close(reply_code=reply_code, reply_text=reply_text)

    def confirm_delivery(self) -> defer.Deferred[Any]:
        """
        Turn on Confirm mode in the channel.

        Pass in a callback to be notified by the Broker when a message has been confirmed as received or
        rejected (Basic.Ack, Basic.Nack) from the broker to the publisher.

        For more information see:
            https://www.rabbitmq.com/confirms.html#publisher-confirms

        :returns: Deferred that fires on the Confirm.SelectOk response
        """
        if self._delivery_confirmation:
            LOGGER.error('confirm_delivery: confirmation was already enabled.')
            return defer.succeed(None)
        wrapped = self._wrap_channel_method('confirm_delivery')
        d = wrapped(ack_nack_callback=self._on_delivery_confirmation)

        def set_delivery_confirmation(result):
            self._delivery_confirmation = True
            self._delivery_message_id = 0
            LOGGER.debug('Delivery confirmation enabled.')
            return result

        d.addCallback(set_delivery_confirmation)
        # Unroutable messages returned after this point will be in the context
        # of publisher acknowledgments
        self._channel.add_on_return_callback(self._on_puback_message_returned)
        return d

    def _on_delivery_confirmation(
        self, method_frame: (pika.frame.Method[pika.spec.Basic.Ack] |
                             pika.frame.Method[pika.spec.Basic.Nack])
    ) -> None:
        """
        Invoked by pika when RabbitMQ responds to a Basic.Publish RPC command, passing in either a
        Basic.Ack or Basic.Nack frame with the delivery tag of the message that was published. The
        delivery tag is an integer counter indicating the message number that was sent on the
        channel via Basic.Publish. Here we're just doing house keeping to keep track of stats and
        remove message numbers that we expect a delivery confirmation of from the list used to keep
        track of messages that are pending confirmation.

        :param method_frame: Basic.Ack or Basic.Nack frame
        """
        delivery_tag = method_frame.method.delivery_tag
        if delivery_tag not in self._deliveries:
            LOGGER.error('Delivery tag %s not found in the pending deliveries',
                         delivery_tag)
            return
        if method_frame.method.multiple:
            tags = [tag for tag in self._deliveries if tag <= delivery_tag]
            tags.sort()
        else:
            tags = [delivery_tag]
        for tag in tags:
            d = self._deliveries[tag]
            del self._deliveries[tag]
            if isinstance(method_frame.method, pika.spec.Basic.Nack):
                # Broker was unable to process the message due to internal
                # error
                LOGGER.warning(
                    "Message was Nack'ed by broker: nack=%r; channel=%s;",
                    method_frame.method, self.channel_number)
                if self._puback_return is not None:
                    returned_messages: list[ReceivedMessage] = [
                        self._puback_return
                    ]
                    self._puback_return = None
                else:
                    returned_messages = []
                d.errback(exceptions.NackError(returned_messages))
            else:
                assert isinstance(method_frame.method, pika.spec.Basic.Ack)
                if self._puback_return is not None:
                    # Unroutable message was returned
                    returned_messages = [self._puback_return]
                    self._puback_return = None
                    d.errback(exceptions.UnroutableError(returned_messages))
                else:
                    d.callback(method_frame.method)

    def _on_puback_message_returned(self, channel: channel.Channel,
                                    method: spec.Basic.Return,
                                    properties: spec.BasicProperties,
                                    body: bytes) -> None:
        """
        Called as the result of Basic.Return from broker in publisher-acknowledgements mode.

        Saves the info as a ReturnedMessage instance in self._puback_return.

        :param channel: our self._impl channel
        :param method:
        :param properties: message properties
        :param body: returned message body; empty string if no body
        """
        assert isinstance(method, spec.Basic.Return), method
        assert isinstance(properties, spec.BasicProperties), properties

        LOGGER.warning(
            'Published message was returned: _delivery_confirmation=%s; '
            'channel=%s; method=%r; properties=%r; body_size=%d; '
            'body_prefix=%.255r', self._delivery_confirmation,
            channel.channel_number, method, properties,
            len(body) if body is not None else None, body)

        self._puback_return = ReceivedMessage(channel=self,
                                              method=method,
                                              properties=properties,
                                              body=body)

    def exchange_bind(
            self,
            destination: str,
            source: str,
            routing_key: str = '',
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Bind an exchange to another exchange.

        :param destination: The destination exchange to bind
        :param source: The source exchange to bind to
        :param routing_key: The routing key to bind on
        :param arguments: Custom key/value pair arguments for the binding
        :raises ValueError:
        :returns: Deferred that fires on the Exchange.BindOk response
        """
        return self._wrap_channel_method('exchange_bind')(
            destination=destination,
            source=source,
            routing_key=routing_key,
            arguments=arguments,
        )

    def exchange_declare(
            self,
            exchange: str,
            exchange_type: str | ExchangeType = ExchangeType.direct,
            passive: bool = False,
            durable: bool = False,
            auto_delete: bool = False,
            internal: bool = False,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        This method creates an exchange if it does not already exist, and if the exchange exists,
        verifies that it is of the correct and expected class.

        If passive set, the server will reply with Declare-Ok if the exchange
        already exists with the same name, and raise an error if not and if the
        exchange does not already exist, the server MUST raise a channel
        exception with reply code 404 (not found).

        :param exchange: The exchange name consists of a non-empty sequence
            of these characters: letters, digits, hyphen, underscore, period,
            or colon
        :param exchange_type: The exchange type to use
        :param passive: Perform a declare or just check to see if it
            exists
        :param durable: Survive a reboot of RabbitMQ
        :param auto_delete: Remove when no more queues are bound to it
        :param internal: Can only be published to by other exchanges
        :param arguments: Custom key/value pair arguments for the exchange
        :returns: Deferred that fires on the Exchange.DeclareOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('exchange_declare')(
            exchange=exchange,
            exchange_type=exchange_type,
            passive=passive,
            durable=durable,
            auto_delete=auto_delete,
            internal=internal,
            arguments=arguments,
        )

    def exchange_delete(self,
                        exchange: str | None = None,
                        if_unused: bool = False) -> defer.Deferred[Any]:
        """
        Delete the exchange.

        :param exchange: The exchange name
        :param if_unused: only delete if the exchange is unused
        :returns: Deferred that fires on the Exchange.DeleteOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('exchange_delete')(
            exchange=exchange,
            if_unused=if_unused,
        )

    def exchange_unbind(
            self,
            destination: str,
            source: str,
            routing_key: str = '',
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Unbind an exchange from another exchange.

        :param destination: The destination exchange to unbind
        :param source: The source exchange to unbind from
        :param routing_key: The routing key to unbind
        :param arguments: Custom key/value pair arguments for the binding
        :returns: Deferred that fires on the Exchange.UnbindOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('exchange_unbind')(
            destination=destination,
            source=source,
            routing_key=routing_key,
            arguments=arguments,
        )

    def flow(self, active: bool = True) -> defer.Deferred[Any]:
        """
        Turn Channel flow control off and on.

        Returns a Deferred that will fire with a bool indicating the channel flow state. For more
        information, please reference:

        https://www.rabbitmq.com/amqp-0-9-1-reference.html#channel.flow

        :param active: Turn flow on or off
        :returns: Deferred that fires with the channel flow state
        :raises ValueError:
        """
        return self._wrap_channel_method('flow')(active=active)

    def open(self) -> None:
        """Open the channel."""
        return self._channel.open()

    def queue_bind(
            self,
            queue: str,
            exchange: str,
            routing_key: str | None = None,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Bind the queue to the specified exchange.

        :param queue: The queue to bind to the exchange
        :param exchange: The source exchange to bind to
        :param routing_key: The routing key to bind on
        :param arguments: Custom key/value pair arguments for the binding
        :returns: Deferred that fires on the Queue.BindOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('queue_bind')(
            queue=queue,
            exchange=exchange,
            routing_key=routing_key,
            arguments=arguments,
        )

    def queue_declare(
            self,
            queue: str,
            passive: bool = False,
            durable: bool = False,
            exclusive: bool = False,
            auto_delete: bool = False,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Declare queue, create if needed.

        This method creates or checks a queue. When creating a new queue the client can specify various
        properties that control the durability of the queue and its contents,
        and the level of sharing for the queue.

        Use an empty string as the queue name for the broker to auto-generate
        one

        :param queue: The queue name; if empty string, the broker will
            create a unique queue name
        :param passive: Only check to see if the queue exists
        :param durable: Survive reboots of the broker
        :param exclusive: Only allow access by the current connection
        :param auto_delete: Delete after consumer cancels or disconnects
        :param arguments: Custom key/value arguments for the queue
        :returns: Deferred that fires on the Queue.DeclareOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('queue_declare')(
            queue=queue,
            passive=passive,
            durable=durable,
            exclusive=exclusive,
            auto_delete=auto_delete,
            arguments=arguments,
        )

    def queue_delete(self,
                     queue: str,
                     if_unused: bool = False,
                     if_empty: bool = False) -> defer.Deferred[Any]:
        """
        Delete a queue from the broker.

        This method wraps :meth:`Channel.queue_delete <pika.channel.Channel.queue_delete>`, and
        removes the reference to the queue object after it gets deleted on the server.

        :param queue: The queue to delete
        :param if_unused: only delete if it's unused
        :param if_empty: only delete if the queue is empty
        :returns: Deferred that fires on the Queue.DeleteOk response
        :raises ValueError:
        """
        wrapped = self._wrap_channel_method('queue_delete')
        d = wrapped(queue=queue, if_unused=if_unused, if_empty=if_empty)

        def _clear_consumer(ret, queue_name):
            for consumer_tag in list(
                    self._queue_name_to_consumer_tags.get(queue_name, set())):
                self._consumers[consumer_tag].close(
                    exceptions.ConsumerCancelled(
                        f'Queue {queue_name} was deleted.'))
                del self._consumers[consumer_tag]
                self._queue_name_to_consumer_tags[queue_name].remove(
                    consumer_tag)
            return ret

        return d.addCallback(_clear_consumer, queue)

    def queue_purge(self, queue: str) -> defer.Deferred[Any]:
        """
        Purge all of the messages from the specified queue.

        :param queue: The queue to purge
        :returns: Deferred that fires on the Queue.PurgeOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('queue_purge')(queue=queue)

    def queue_unbind(
            self,
            queue: str,
            exchange: str | None,
            routing_key: str | None = None,
            arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
        """
        Unbind a queue from an exchange.

        :param queue: The queue to unbind from the exchange
        :param exchange: The source exchange to bind from
        :param routing_key: The routing key to unbind
        :param arguments: Custom key/value pair arguments for the binding
        :returns: Deferred that fires on the Queue.UnbindOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('queue_unbind')(
            queue=queue,
            exchange=exchange,
            routing_key=routing_key,
            arguments=arguments,
        )

    def tx_commit(self) -> defer.Deferred[Any]:
        """
        Commit a transaction.

        :returns: Deferred that fires on the Tx.CommitOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('tx_commit')()

    def tx_rollback(self) -> defer.Deferred[Any]:
        """
        Rollback a transaction.

        :returns: Deferred that fires on the Tx.RollbackOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('tx_rollback')()

    def tx_select(self) -> defer.Deferred[Any]:
        """
        Select standard transaction mode.

        This method sets the channel to use standard transactions. The client must use this method
        at least once on a channel before using the Commit or Rollback methods.

        :returns: Deferred that fires on the Tx.SelectOk response
        :raises ValueError:
        """
        return self._wrap_channel_method('tx_select')()

_basic_get_deferred instance-attribute

_basic_get_deferred: Deferred[Any] | None = None

_calls instance-attribute

_calls: set[Deferred[Any]] = set()

_channel instance-attribute

_channel = channel

_closed instance-attribute

_closed: Exception | None = None

_consumers instance-attribute

_consumers: dict[str, ClosableDeferredQueue] = {}

_deliveries instance-attribute

_deliveries: dict[int, Deferred[Any]] = {}

_delivery_confirmation instance-attribute

_delivery_confirmation = False

_delivery_message_id instance-attribute

_delivery_message_id: int | None = None

_puback_return instance-attribute

_puback_return: ReceivedMessage | None = None

_queue_name_to_consumer_tags instance-attribute

_queue_name_to_consumer_tags: dict[str, set[str]] = {}

channel_number property

channel_number: int

connection property

connection: Connection

consumer_tags property

consumer_tags: list[str]

flow_active property

flow_active: bool

is_closed property

is_closed: bool

Returns True if the channel is closed.

is_closing property

is_closing: bool

Returns True if client-initiated closing of the channel is in progress.

is_open property

is_open: bool

Returns True if the channel is open.

on_closed instance-attribute

on_closed: Deferred[Any] = defer.Deferred()

__init__

__init__(channel: Channel) -> None
Source code in pika/adapters/twisted_connection.py
def __init__(self, channel: channel.Channel) -> None:
    self._channel = channel
    self._closed: Exception | None = None
    self._calls: set[defer.Deferred[Any]] = set()
    self._consumers: dict[str, ClosableDeferredQueue] = {}
    # Store Basic.Get calls so we can handle GetEmpty replies
    self._basic_get_deferred: defer.Deferred[Any] | None = None
    self._channel.add_callback(self._on_getempty, [spec.Basic.GetEmpty],
                               False)
    # We need this mapping to close the ClosableDeferredQueue when a queue
    # is deleted.
    self._queue_name_to_consumer_tags: dict[str, set[str]] = {}
    # Whether RabbitMQ delivery confirmation has been enabled
    self._delivery_confirmation = False
    self._delivery_message_id: int | None = None
    self._deliveries: dict[int, defer.Deferred[Any]] = {}
    # Holds a ReceivedMessage object representing a message received via
    # Basic.Return in publisher-acknowledgments mode.
    self._puback_return: ReceivedMessage | None = None

    self.on_closed: defer.Deferred[Any] = defer.Deferred()
    self._channel.add_on_close_callback(self._on_channel_closed)
    self._channel.add_on_cancel_callback(
        self._on_consumer_cancelled_by_broker)

__repr__

__repr__() -> str
Source code in pika/adapters/twisted_connection.py
@override
def __repr__(self) -> str:
    return f'<{self.__class__.__name__} channel={self._channel!r}>'

_clear_call

_clear_call(ret: Any, d: Deferred[Any]) -> Any
Source code in pika/adapters/twisted_connection.py
def _clear_call(self, ret: Any, d: defer.Deferred[Any]) -> Any:
    self._calls.discard(d)
    return ret

_on_channel_closed

_on_channel_closed(
    _channel: Channel, reason: Exception | None
) -> None
Source code in pika/adapters/twisted_connection.py
def _on_channel_closed(self, _channel: channel.Channel,
                       reason: Exception | None) -> None:
    # enter the closed state
    self._closed = reason
    # errback all pending calls
    for d in self._calls:
        d.errback(self._closed)
    # errback all pending deliveries
    for d in self._deliveries.values():
        d.errback(self._closed)
    # close all open queues
    for consumer in self._consumers.values():
        consumer.close(self._closed)
    # release references to stored objects
    self._calls = set()
    self._deliveries = {}
    self._consumers = {}
    self.on_closed.callback(self._closed)

_on_consumer_cancelled

_on_consumer_cancelled(
    frame: Method[Cancel] | Method[CancelOk],
) -> Method[Cancel] | Method[CancelOk]

Called when the broker cancels a consumer via Basic.Cancel or when the broker responds to a Basic.Cancel request by Basic.CancelOk.

PARAMETER DESCRIPTION
frame

method frame with the spec.Basic.Cancel or spec.Basic.CancelOk method

TYPE: Method[Cancel] | Method[CancelOk]

Source code in pika/adapters/twisted_connection.py
def _on_consumer_cancelled(
    self, frame: (pika.frame.Method[pika.spec.Basic.Cancel] |
                  pika.frame.Method[pika.spec.Basic.CancelOk])
) -> (pika.frame.Method[pika.spec.Basic.Cancel] |
      pika.frame.Method[pika.spec.Basic.CancelOk]):
    """
    Called when the broker cancels a consumer via Basic.Cancel or when the broker responds to a
    Basic.Cancel request by Basic.CancelOk.

    :param frame: method frame with the `spec.Basic.Cancel` or `spec.Basic.CancelOk` method
    """
    consumer_tag = frame.method.consumer_tag
    if consumer_tag not in self._consumers:
        # Could be cancelled by user or broker earlier
        LOGGER.warning('basic_cancel - consumer not found: %s',
                       consumer_tag)
        return frame
    self._consumers[consumer_tag].close(exceptions.ConsumerCancelled())
    del self._consumers[consumer_tag]
    # Remove from the queue-to-ctags index:
    for ctags in self._queue_name_to_consumer_tags.values():
        try:
            ctags.remove(consumer_tag)
        except KeyError:
            continue
    return frame

_on_consumer_cancelled_by_broker

_on_consumer_cancelled_by_broker(
    method_frame: Method[Cancel],
) -> Method[Cancel] | Method[CancelOk]

Called by impl when broker cancels consumer via Basic.Cancel.

This is a RabbitMQ-specific feature. The circumstances include deletion of queue being consumed as well as failure of a HA node responsible for the queue being consumed.

PARAMETER DESCRIPTION
method_frame

method frame with the spec.Basic.Cancel method

TYPE: Method[Cancel]

Source code in pika/adapters/twisted_connection.py
def _on_consumer_cancelled_by_broker(
    self, method_frame: pika.frame.Method[pika.spec.Basic.Cancel]
) -> (pika.frame.Method[pika.spec.Basic.Cancel] |
      pika.frame.Method[pika.spec.Basic.CancelOk]):
    """
    Called by impl when broker cancels consumer via Basic.Cancel.

    This is a RabbitMQ-specific feature. The circumstances include deletion of queue being
    consumed as well as failure of a HA node responsible for the queue being consumed.

    :param method_frame: method frame with the `spec.Basic.Cancel` method
    """
    return self._on_consumer_cancelled(method_frame)

_on_delivery_confirmation

_on_delivery_confirmation(
    method_frame: Method[Ack] | Method[Nack],
) -> None

Invoked by pika when RabbitMQ responds to a Basic.Publish RPC command, passing in either a Basic.Ack or Basic.Nack frame with the delivery tag of the message that was published. The delivery tag is an integer counter indicating the message number that was sent on the channel via Basic.Publish. Here we're just doing house keeping to keep track of stats and remove message numbers that we expect a delivery confirmation of from the list used to keep track of messages that are pending confirmation.

PARAMETER DESCRIPTION
method_frame

Basic.Ack or Basic.Nack frame

TYPE: Method[Ack] | Method[Nack]

Source code in pika/adapters/twisted_connection.py
def _on_delivery_confirmation(
    self, method_frame: (pika.frame.Method[pika.spec.Basic.Ack] |
                         pika.frame.Method[pika.spec.Basic.Nack])
) -> None:
    """
    Invoked by pika when RabbitMQ responds to a Basic.Publish RPC command, passing in either a
    Basic.Ack or Basic.Nack frame with the delivery tag of the message that was published. The
    delivery tag is an integer counter indicating the message number that was sent on the
    channel via Basic.Publish. Here we're just doing house keeping to keep track of stats and
    remove message numbers that we expect a delivery confirmation of from the list used to keep
    track of messages that are pending confirmation.

    :param method_frame: Basic.Ack or Basic.Nack frame
    """
    delivery_tag = method_frame.method.delivery_tag
    if delivery_tag not in self._deliveries:
        LOGGER.error('Delivery tag %s not found in the pending deliveries',
                     delivery_tag)
        return
    if method_frame.method.multiple:
        tags = [tag for tag in self._deliveries if tag <= delivery_tag]
        tags.sort()
    else:
        tags = [delivery_tag]
    for tag in tags:
        d = self._deliveries[tag]
        del self._deliveries[tag]
        if isinstance(method_frame.method, pika.spec.Basic.Nack):
            # Broker was unable to process the message due to internal
            # error
            LOGGER.warning(
                "Message was Nack'ed by broker: nack=%r; channel=%s;",
                method_frame.method, self.channel_number)
            if self._puback_return is not None:
                returned_messages: list[ReceivedMessage] = [
                    self._puback_return
                ]
                self._puback_return = None
            else:
                returned_messages = []
            d.errback(exceptions.NackError(returned_messages))
        else:
            assert isinstance(method_frame.method, pika.spec.Basic.Ack)
            if self._puback_return is not None:
                # Unroutable message was returned
                returned_messages = [self._puback_return]
                self._puback_return = None
                d.errback(exceptions.UnroutableError(returned_messages))
            else:
                d.callback(method_frame.method)

_on_getempty

_on_getempty(_method_frame: Method[Get]) -> None

Callback the Basic.Get deferred with None.

PARAMETER DESCRIPTION
_method_frame

Method frame from Basic.Get response (unused)

TYPE: Method[Get]

Source code in pika/adapters/twisted_connection.py
def _on_getempty(
        self,
        _method_frame: pika.frame.Method[pika.spec.Basic.Get]) -> None:
    """
    Callback the Basic.Get deferred with None.

    :param _method_frame: Method frame from Basic.Get response (unused)
    """
    if self._basic_get_deferred is None:
        LOGGER.warning('Got Basic.GetEmpty but no Basic.Get calls '
                       'were pending.')
        return
    self._basic_get_deferred.callback(None)

_on_puback_message_returned

_on_puback_message_returned(
    channel: Channel,
    method: Return,
    properties: BasicProperties,
    body: bytes,
) -> None

Called as the result of Basic.Return from broker in publisher-acknowledgements mode.

Saves the info as a ReturnedMessage instance in self._puback_return.

PARAMETER DESCRIPTION
channel

our self._impl channel

TYPE: Channel

method

TYPE: Return

properties

message properties

TYPE: BasicProperties

body

returned message body; empty string if no body

TYPE: bytes

Source code in pika/adapters/twisted_connection.py
def _on_puback_message_returned(self, channel: channel.Channel,
                                method: spec.Basic.Return,
                                properties: spec.BasicProperties,
                                body: bytes) -> None:
    """
    Called as the result of Basic.Return from broker in publisher-acknowledgements mode.

    Saves the info as a ReturnedMessage instance in self._puback_return.

    :param channel: our self._impl channel
    :param method:
    :param properties: message properties
    :param body: returned message body; empty string if no body
    """
    assert isinstance(method, spec.Basic.Return), method
    assert isinstance(properties, spec.BasicProperties), properties

    LOGGER.warning(
        'Published message was returned: _delivery_confirmation=%s; '
        'channel=%s; method=%r; properties=%r; body_size=%d; '
        'body_prefix=%.255r', self._delivery_confirmation,
        channel.channel_number, method, properties,
        len(body) if body is not None else None, body)

    self._puback_return = ReceivedMessage(channel=self,
                                          method=method,
                                          properties=properties,
                                          body=body)

_wrap_channel_method

_wrap_channel_method(
    name: str,
) -> Callable[..., Deferred[Any]]

Wrap Pika's Channel method to make it return a Deferred that fires when the method completes and errbacks if the channel gets closed. If the original method's callback would receive more than one argument, the Deferred fires with a tuple of argument values.

PARAMETER DESCRIPTION
name

Attribute name to look up on the underlying channel

TYPE: str

Source code in pika/adapters/twisted_connection.py
def _wrap_channel_method(self,
                         name: str) -> Callable[..., defer.Deferred[Any]]:
    """
    Wrap Pika's Channel method to make it return a Deferred that fires when the method completes
    and errbacks if the channel gets closed. If the original method's callback would receive
    more than one argument, the Deferred fires with a tuple of argument values.

    :param name: Attribute name to look up on the underlying channel
    """
    method = getattr(self._channel, name)

    @functools.wraps(method)
    def wrapped(*args, **kwargs):
        if self._closed:
            return defer.fail(self._closed)

        d: defer.Deferred = defer.Deferred()
        self._calls.add(d)
        d.addCallback(self._clear_call, d)

        def single_argument(*args) -> None:
            """
            Make sure that the deferred is called with a single argument.

            In case the original callback fires with more than one, convert to a tuple.
            """
            if len(args) > 1:
                d.callback(tuple(args))
            else:
                d.callback(*args)

        kwargs['callback'] = single_argument

        try:
            method(*args, **kwargs)
        except Exception:
            return defer.fail()
        return d

    return wrapped

add_on_return_callback

add_on_return_callback(
    callback: Callable[[ReceivedMessage], None],
) -> None

Pass a callback function that will be called when a published message is rejected and returned by the server via Basic.Return.

PARAMETER DESCRIPTION
callback

The method to call on callback with the message as only argument. The message is a named tuple with the following attributes - channel: this TwistedChannel - method: pika.spec.Basic.Return - properties: pika.spec.BasicProperties - body: bytes

TYPE: Callable[[ReceivedMessage], None]

Source code in pika/adapters/twisted_connection.py
def add_on_return_callback(
        self, callback: Callable[[ReceivedMessage], None]) -> None:
    """
    Pass a callback function that will be called when a published message is rejected and
    returned by the server via `Basic.Return`.

    :param callback: The method to call on callback with the
        message as only argument. The message is a named tuple with
        the following attributes
        - channel: this TwistedChannel
        - method: pika.spec.Basic.Return
        - properties: pika.spec.BasicProperties
        - body: bytes
    """
    self._channel.add_on_return_callback(
        lambda _channel, method, properties, body: callback(
            ReceivedMessage(
                channel=self,
                method=method,
                properties=properties,
                body=body,
            )))

basic_ack

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

Acknowledge one or more messages.

When sent by the client, this method acknowledges one or more messages delivered via the Deliver or Get-Ok methods. When sent by server, this method acknowledges one or more messages published with the Publish method on a channel in confirm mode. The acknowledgement can be for a single message or a set of messages up to and including a specific message.

PARAMETER DESCRIPTION
delivery_tag

int/long The server-assigned delivery tag

TYPE: int DEFAULT: 0

multiple

If set to True, the delivery tag is treated as "up to and including", so that multiple messages can be acknowledged with a single method. If set to False, the delivery tag refers to a single message. If the multiple field is 1, and the delivery tag is zero, this indicates acknowledgement of all outstanding messages.

TYPE: bool DEFAULT: False

Source code in pika/adapters/twisted_connection.py
def basic_ack(self, delivery_tag: int = 0, multiple: bool = False) -> None:
    """
    Acknowledge one or more messages.

    When sent by the client, this method acknowledges one or more messages delivered via the
    Deliver or Get-Ok methods. When sent by server, this method acknowledges one or more
    messages published with the Publish method on a channel in confirm mode. The acknowledgement
    can be for a single message or a set of messages up to and including a specific message.

    :param delivery_tag: int/long The server-assigned delivery tag
    :param multiple: If set to True, the delivery tag is treated as "up to and including", so
        that multiple messages can be acknowledged with a single method. If set to False, the
        delivery tag refers to a single message. If the multiple field is 1, and the delivery
        tag is zero, this indicates acknowledgement of all outstanding messages.
    """
    return self._channel.basic_ack(delivery_tag=delivery_tag,
                                   multiple=multiple)

basic_cancel

basic_cancel(
    consumer_tag: str = '',
) -> Deferred[Method[CancelOk] | Method[Cancel]]

This method cancels a consumer.

This does not affect already delivered messages, but it does mean the server will not send any more messages for that consumer. The client may receive an arbitrary number of messages in between sending the cancel method and receiving the cancel-ok reply. It may also be sent from the server to the client in the event of the consumer being unexpectedly cancelled (i.e. cancelled for any reason other than the server receiving the corresponding basic.cancel from the client). This allows clients to be notified of the loss of consumers due to events such as queue deletion.

This method wraps :meth:Channel.basic_cancel <pika.channel.Channel.basic_cancel> and closes any deferred queue associated with that consumer.

PARAMETER DESCRIPTION
consumer_tag

Identifier for the consumer

TYPE: str DEFAULT: ''

RETURNS DESCRIPTION
Deferred[Method[CancelOk] | Method[Cancel]]

Deferred that fires on the Basic.CancelOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def basic_cancel(
    self,
    consumer_tag: str = ''
) -> defer.Deferred[(pika.frame.Method[pika.spec.Basic.CancelOk] |
                     pika.frame.Method[pika.spec.Basic.Cancel])]:
    """
    This method cancels a consumer.

    This does not affect already delivered messages, but it does mean the server will not send
    any more messages for that consumer. The client may receive an arbitrary number of messages
    in between sending the cancel method and receiving the cancel-ok reply. It may also be sent
    from the server to the client in the event of the consumer being unexpectedly cancelled
    (i.e. cancelled for any reason other than the server receiving the corresponding
    basic.cancel from the client). This allows clients to be notified of the loss of consumers
    due to events such as queue deletion.

    This method wraps :meth:`Channel.basic_cancel <pika.channel.Channel.basic_cancel>` and
    closes any deferred queue associated with that consumer.

    :param consumer_tag: Identifier for the consumer
    :returns: Deferred that fires on the Basic.CancelOk response
    :raises ValueError:
    """
    wrapped = self._wrap_channel_method('basic_cancel')
    d = wrapped(consumer_tag=consumer_tag)
    return d.addCallback(self._on_consumer_cancelled)

basic_consume

basic_consume(
    queue: str,
    auto_ack: bool = False,
    exclusive: bool = False,
    consumer_tag: str | None = None,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Consume from a server queue.

Sends the AMQP 0-9-1 command Basic.Consume to the broker and binds messages for the consumer_tag to a :class:ClosableDeferredQueue. If you do not pass in a consumer_tag, one will be automatically generated for you.

For more information on basic_consume, see: Tutorial 2 at https://www.rabbitmq.com/getstarted.html https://www.rabbitmq.com/confirms.html https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.consume

PARAMETER DESCRIPTION
queue

The queue to consume from. Use the empty string to specify the most recent server-named queue for this channel.

TYPE: str

auto_ack

if set to True, automatic acknowledgement mode will be used (see https://www.rabbitmq.com/confirms.html). This corresponds with the 'no_ack' parameter in the basic.consume AMQP 0.9.1 method

TYPE: bool DEFAULT: False

exclusive

Don't allow other consumers on the queue

TYPE: bool DEFAULT: False

consumer_tag

Specify your own consumer tag

TYPE: str | None DEFAULT: None

arguments

Custom key/value pair arguments for the consumer

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires with a tuple (queue_object, consumer_tag). The Deferred will errback with an instance of :class:exceptions.ChannelClosed if the call fails. The queue object is an instance of :class:ClosableDeferredQueue, where data received from the queue will be stored. Clients should use its :meth:get() <ClosableDeferredQueue.get> method to fetch an individual message, which will return a Deferred firing with a namedtuple whose attributes are: - channel: this TwistedChannel - method: pika.spec.Basic.Deliver - properties: pika.spec.BasicProperties - body: bytes

Source code in pika/adapters/twisted_connection.py
def basic_consume(
        self,
        queue: str,
        auto_ack: bool = False,
        exclusive: bool = False,
        consumer_tag: str | None = None,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Consume from a server queue.

    Sends the AMQP 0-9-1 command Basic.Consume to the broker and binds
    messages for the consumer_tag to a
    :class:`ClosableDeferredQueue`. If you do not pass in a
    consumer_tag, one will be automatically generated for you.

    For more information on basic_consume, see:
    Tutorial 2 at https://www.rabbitmq.com/getstarted.html
    https://www.rabbitmq.com/confirms.html
    https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.consume

    :param queue: The queue to consume from. Use the empty string to
        specify the most recent server-named queue for this channel.
    :param auto_ack: if set to True, automatic acknowledgement mode
        will be used (see https://www.rabbitmq.com/confirms.html). This
        corresponds with the 'no_ack' parameter in the basic.consume AMQP
        0.9.1 method
    :param exclusive: Don't allow other consumers on the queue
    :param consumer_tag: Specify your own consumer tag
    :param arguments: Custom key/value pair arguments for the consumer
    :returns: Deferred that fires with a tuple
        ``(queue_object, consumer_tag)``. The Deferred will errback with an
        instance of :class:`exceptions.ChannelClosed` if the call fails.
        The queue object is an instance of :class:`ClosableDeferredQueue`,
        where data received from the queue will be stored. Clients should
        use its :meth:`get() <ClosableDeferredQueue.get>` method to fetch
        an individual message, which will return a Deferred firing with a
        namedtuple whose attributes are:
        - channel: this TwistedChannel
        - method: pika.spec.Basic.Deliver
        - properties: pika.spec.BasicProperties
        - body: bytes
    """
    if self._closed:
        return defer.fail(self._closed)

    queue_obj = ClosableDeferredQueue()
    d: defer.Deferred[Any] = defer.Deferred()
    self._calls.add(d)

    def on_consume_ok(frame) -> None:
        consumer_tag = frame.method.consumer_tag
        self._queue_name_to_consumer_tags.setdefault(
            queue, set()).add(consumer_tag)
        self._consumers[consumer_tag] = queue_obj
        self._calls.discard(d)
        d.callback((queue_obj, consumer_tag))

    def on_message_callback(_channel, method, properties, body) -> None:
        """Add the ReceivedMessage to the queue, while replacing the channel implementation."""
        queue_obj.put(
            ReceivedMessage(
                channel=self,
                method=method,
                properties=properties,
                body=body,
            ))

    try:
        self._channel.basic_consume(
            queue=queue,
            on_message_callback=on_message_callback,
            auto_ack=auto_ack,
            exclusive=exclusive,
            consumer_tag=consumer_tag,
            arguments=arguments,
            callback=on_consume_ok,
        )
    except Exception:
        return defer.fail()

    return d

basic_get

basic_get(
    queue: str, auto_ack: bool = False
) -> Deferred[Any]

Get a single message from the AMQP broker.

Will return If the queue is empty, it will return None. If you want to be notified of Basic.GetEmpty, use the Channel.add_callback method adding your Basic.GetEmpty callback which should expect only one parameter, frame. Due to implementation details, this cannot be called a second time until the callback is executed. For more information on basic_get and its parameters, see:

https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.get

This method wraps :meth:Channel.basic_get <pika.channel.Channel.basic_get>.

PARAMETER DESCRIPTION
queue

The queue from which to get a message. Use the empty string to specify the most recent server-named queue for this channel.

TYPE: str

auto_ack

Tell the broker to not expect a reply

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires with a namedtuple whose attributes are: - channel: this TwistedChannel - method: pika.spec.Basic.GetOk - properties: pika.spec.BasicProperties - body: bytes If the queue is empty, None will be returned.

RAISES DESCRIPTION
pika.exceptions.DuplicateGetOkCallback
Source code in pika/adapters/twisted_connection.py
def basic_get(self,
              queue: str,
              auto_ack: bool = False) -> defer.Deferred[Any]:
    """
    Get a single message from the AMQP broker.

    Will return If the queue is empty, it will return None.
    If you want to
    be notified of Basic.GetEmpty, use the Channel.add_callback method
    adding your Basic.GetEmpty callback which should expect only one
    parameter, frame. Due to implementation details, this cannot be called
    a second time until the callback is executed.  For more information on
    basic_get and its parameters, see:

    https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.get

    This method wraps :meth:`Channel.basic_get
    <pika.channel.Channel.basic_get>`.

    :param queue: The queue from which to get a message. Use the empty
                  string to specify the most recent server-named queue
                  for this channel.
    :param auto_ack: Tell the broker to not expect a reply
    :returns: Deferred that fires with a namedtuple whose attributes are:
         - channel: this TwistedChannel
         - method: pika.spec.Basic.GetOk
         - properties: pika.spec.BasicProperties
         - body: bytes
        If the queue is empty, None will be returned.
    :raises pika.exceptions.DuplicateGetOkCallback:
    """
    if self._basic_get_deferred is not None:
        raise exceptions.DuplicateGetOkCallback()

    def create_namedtuple(result) -> ReceivedMessage | None:
        if result is None:
            return None
        _channel, method, properties, body = result
        return ReceivedMessage(
            channel=self,
            method=method,
            properties=properties,
            body=body,
        )

    def cleanup_attribute(result):
        self._basic_get_deferred = None
        return result

    d = self._wrap_channel_method('basic_get')(queue=queue,
                                               auto_ack=auto_ack)
    d.addCallback(create_namedtuple)
    d.addBoth(cleanup_attribute)
    self._basic_get_deferred = d
    return d

basic_nack

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

This method allows a client to reject one or more incoming messages.

It can be used to interrupt and cancel large incoming messages, or return untreatable messages to their original queue.

PARAMETER DESCRIPTION
delivery_tag

int/long The server-assigned delivery tag

TYPE: int DEFAULT: 0

multiple

If set to True, the delivery tag is treated as "up to and including", so that multiple messages can be acknowledged with a single method. If set to False, the delivery tag refers to a single message. If the multiple field is 1, and the delivery tag is zero, this indicates acknowledgement of all outstanding messages.

TYPE: bool DEFAULT: False

requeue

If requeue is true, the server will attempt to requeue the message. If requeue is false or the requeue attempt fails the messages are discarded or dead- lettered.

TYPE: bool DEFAULT: True

Source code in pika/adapters/twisted_connection.py
def basic_nack(self,
               delivery_tag: int = 0,
               multiple: bool = False,
               requeue: bool = True) -> None:
    """
    This method allows a client to reject one or more incoming messages.

    It can be used to interrupt and cancel large incoming messages, or return untreatable
    messages to their original queue.

    :param delivery_tag: int/long The server-assigned delivery tag
    :param multiple: If set to True, the delivery tag is treated as "up to and including", so
        that multiple messages can be acknowledged with a single method. If set to False, the
        delivery tag refers to a single message. If the multiple field is 1, and the delivery
        tag is zero, this indicates acknowledgement of all outstanding messages.
    :param requeue: If requeue is true, the server will attempt to requeue the message. If
        requeue is false or the requeue attempt fails the messages are discarded or dead-
        lettered.
    """
    return self._channel.basic_nack(
        delivery_tag=delivery_tag,
        multiple=multiple,
        requeue=requeue,
    )

basic_publish

basic_publish(
    exchange: str,
    routing_key: str,
    body: bytes,
    properties: BasicProperties | None = None,
    mandatory: bool = False,
) -> Deferred[Any]

Publish to the channel with the given exchange, routing key and body.

This method wraps :meth:Channel.basic_publish <pika.channel.Channel.basic_publish>, but makes sure the channel is not closed before publishing.

For more information on basic_publish and what the parameters do, see:

https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.publish

PARAMETER DESCRIPTION
exchange

The exchange to publish to

TYPE: str

routing_key

The routing key to bind on

TYPE: str

body

The message body

TYPE: bytes

properties

Basic.properties

TYPE: BasicProperties | None DEFAULT: None

mandatory

The mandatory flag

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

A Deferred that fires with the result of the channel's basic_publish.

RAISES DESCRIPTION
UnroutableError

raised when a message published in publisher-acknowledgments mode (see BlockingChannel.confirm_delivery) is returned via Basic.Return followed by Basic.Ack.

NackError

raised when a message published in publisher-acknowledgements mode is Nack'ed by the broker. See BlockingChannel.confirm_delivery.

Source code in pika/adapters/twisted_connection.py
def basic_publish(self,
                  exchange: str,
                  routing_key: str,
                  body: bytes,
                  properties: spec.BasicProperties | None = None,
                  mandatory: bool = False) -> defer.Deferred[Any]:
    """
    Publish to the channel with the given exchange, routing key and body.

    This method wraps :meth:`Channel.basic_publish <pika.channel.Channel.basic_publish>`, but
    makes sure the channel is not closed before publishing.

    For more information on basic_publish and what the parameters do, see:

    https://www.rabbitmq.com/amqp-0-9-1-reference.html#basic.publish

    :param exchange: The exchange to publish to
    :param routing_key: The routing key to bind on
    :param body: The message body
    :param properties: Basic.properties
    :param mandatory: The mandatory flag
    :returns: A Deferred that fires with the result of the channel's basic_publish.
    :raises UnroutableError: raised when a message published in publisher-acknowledgments mode
        (see `BlockingChannel.confirm_delivery`) is returned via `Basic.Return` followed by
        `Basic.Ack`.
    :raises NackError: raised when a message published in publisher-acknowledgements mode is
        Nack'ed by the broker. See `BlockingChannel.confirm_delivery`.
    """
    if self._closed:
        return defer.fail(self._closed)
    self._channel.basic_publish(exchange=exchange,
                                routing_key=routing_key,
                                body=body,
                                properties=properties,
                                mandatory=mandatory)
    if not self._delivery_confirmation:
        # `Channel.basic_publish()` returns nothing, so without delivery
        # confirmation there is no result to report beyond completion.
        return defer.succeed(None)
    # See https://www.rabbitmq.com/confirms.html#publisher-confirms
    assert self._delivery_message_id is not None
    self._delivery_message_id += 1
    self._deliveries[self._delivery_message_id] = defer.Deferred()
    return self._deliveries[self._delivery_message_id]

basic_qos

basic_qos(
    prefetch_size: int = 0,
    prefetch_count: int = 0,
    global_qos: bool = False,
) -> Deferred[Any]

Specify quality of service.

This method requests a specific quality of service. The QoS can be specified for the current channel or for all channels on the connection. The client can request that messages be sent in advance so that when the client finishes processing a message, the following message is already held locally, rather than needing to be sent down the channel. Prefetching gives a performance improvement.

PARAMETER DESCRIPTION
prefetch_size

This field specifies the prefetch window size. The server will send a message in advance if it is equal to or smaller in size than the available prefetch size (and also falls into other prefetch limits). May be set to zero, meaning "no specific limit", although other prefetch limits may still apply. The prefetch-size is ignored by consumers who have enabled the no-ack option.

TYPE: int DEFAULT: 0

prefetch_count

Specifies a prefetch window in terms of whole messages. This field may be used in combination with the prefetch-size field; a message will only be sent in advance if both prefetch windows (and those at the channel and connection level) allow it. The prefetch-count is ignored by consumers who have enabled the no-ack option.

TYPE: int DEFAULT: 0

global_qos

Should the QoS apply to all channels on the connection.

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Basic.QosOk response

Source code in pika/adapters/twisted_connection.py
def basic_qos(self,
              prefetch_size: int = 0,
              prefetch_count: int = 0,
              global_qos: bool = False) -> defer.Deferred[Any]:
    """
    Specify quality of service.

    This method requests a specific quality of service. The QoS can be specified for the current
    channel or for all channels on the connection. The client can request that messages be sent
    in advance so that when the client finishes processing a message, the following message is
    already held locally, rather than needing to be sent down the channel. Prefetching gives a
    performance improvement.

    :param prefetch_size: This field specifies the prefetch window size. The server will send a
        message in advance if it is equal to or smaller in size than the available prefetch size
        (and also falls into other prefetch limits). May be set to zero, meaning "no specific
        limit", although other prefetch limits may still apply. The prefetch-size is ignored by
        consumers who have enabled the no-ack option.
    :param prefetch_count: Specifies a prefetch window in terms of whole messages. This field
        may be used in combination with the prefetch-size field; a message will only be sent in
        advance if both prefetch windows (and those at the channel and connection level) allow
        it. The prefetch-count is ignored by consumers who have enabled the no-ack option.
    :param global_qos: Should the QoS apply to all channels on the connection.
    :returns: Deferred that fires on the Basic.QosOk response
    """
    return self._wrap_channel_method('basic_qos')(
        prefetch_size=prefetch_size,
        prefetch_count=prefetch_count,
        global_qos=global_qos,
    )

basic_recover

basic_recover(requeue: bool = False) -> Deferred[Any]

This method asks the server to redeliver all unacknowledged messages on a specified channel. Zero or more messages may be redelivered. This method replaces the asynchronous Recover.

PARAMETER DESCRIPTION
requeue

If False, the message will be redelivered to the original recipient. If True, the server will attempt to requeue the message, potentially then delivering it to an alternative subscriber.

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Basic.RecoverOk response

Source code in pika/adapters/twisted_connection.py
def basic_recover(self, requeue: bool = False) -> defer.Deferred[Any]:
    """
    This method asks the server to redeliver all unacknowledged messages on a specified channel.
    Zero or more messages may be redelivered. This method replaces the asynchronous Recover.

    :param requeue: If False, the message will be redelivered to the original recipient. If
        True, the server will attempt to requeue the message, potentially then delivering it to
        an alternative subscriber.
    :returns: Deferred that fires on the Basic.RecoverOk response
    """
    return self._wrap_channel_method('basic_recover')(requeue=requeue)

basic_reject

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

Reject an incoming message.

This method allows a client to reject a message. It can be used to interrupt and cancel large incoming messages, or return untreatable messages to their original queue.

PARAMETER DESCRIPTION
delivery_tag

int/long The server-assigned delivery tag

TYPE: int

requeue

If requeue is true, the server will attempt to requeue the message. If requeue is false or the requeue attempt fails the messages are discarded or dead- lettered.

TYPE: bool DEFAULT: True

RAISES DESCRIPTION
TypeError
Source code in pika/adapters/twisted_connection.py
def basic_reject(self, delivery_tag: int, requeue: bool = True) -> None:
    """
    Reject an incoming message.

    This method allows a client to reject a message. It can be used to interrupt and cancel
    large incoming messages, or return untreatable messages to their original queue.

    :param delivery_tag: int/long The server-assigned delivery tag
    :param requeue: If requeue is true, the server will attempt to requeue the message. If
        requeue is false or the requeue attempt fails the messages are discarded or dead-
        lettered.
    :raises TypeError:
    """
    return self._channel.basic_reject(delivery_tag=delivery_tag,
                                      requeue=requeue)

callback_deferred

callback_deferred(
    deferred: Deferred[Any], replies: Sequence[type[Method]]
) -> None

Pass in a Deferred and a list replies from the RabbitMQ broker which you'd like the Deferred to be callbacked on with the frame as callback value.

PARAMETER DESCRIPTION
deferred

The Deferred to callback

TYPE: Deferred[Any]

replies

The replies to callback on

TYPE: Sequence[type[Method]]

Source code in pika/adapters/twisted_connection.py
def callback_deferred(self, deferred: defer.Deferred[Any],
                      replies: Sequence[type[amqp_object.Method]]) -> None:
    """
    Pass in a Deferred and a list replies from the RabbitMQ broker which you'd like the Deferred
    to be callbacked on with the frame as callback value.

    :param deferred: The Deferred to callback
    :param replies: The replies to callback on
    """
    self._channel.add_callback(deferred.callback, replies)

close

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

Invoke a graceful shutdown of the channel with the AMQP Broker.

If channel is OPENING, transition to CLOSING and suppress the incoming Channel.OpenOk, if any.

PARAMETER DESCRIPTION
reply_code

The reason code to send to broker

TYPE: int DEFAULT: 0

reply_text

The reason text to send to broker

TYPE: str DEFAULT: 'Normal shutdown'

RAISES DESCRIPTION
ChannelWrongStateError

if channel is closed or closing

Source code in pika/adapters/twisted_connection.py
def close(self,
          reply_code: int = 0,
          reply_text: str = 'Normal shutdown') -> None:
    """
    Invoke a graceful shutdown of the channel with the AMQP Broker.

    If channel is OPENING, transition to CLOSING and suppress the incoming Channel.OpenOk, if
    any.

    :param reply_code: The reason code to send to broker
    :param reply_text: The reason text to send to broker
    :raises ChannelWrongStateError: if channel is closed or closing
    """
    return self._channel.close(reply_code=reply_code, reply_text=reply_text)

confirm_delivery

confirm_delivery() -> Deferred[Any]

Turn on Confirm mode in the channel.

Pass in a callback to be notified by the Broker when a message has been confirmed as received or rejected (Basic.Ack, Basic.Nack) from the broker to the publisher.

For more information see: https://www.rabbitmq.com/confirms.html#publisher-confirms

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Confirm.SelectOk response

Source code in pika/adapters/twisted_connection.py
def confirm_delivery(self) -> defer.Deferred[Any]:
    """
    Turn on Confirm mode in the channel.

    Pass in a callback to be notified by the Broker when a message has been confirmed as received or
    rejected (Basic.Ack, Basic.Nack) from the broker to the publisher.

    For more information see:
        https://www.rabbitmq.com/confirms.html#publisher-confirms

    :returns: Deferred that fires on the Confirm.SelectOk response
    """
    if self._delivery_confirmation:
        LOGGER.error('confirm_delivery: confirmation was already enabled.')
        return defer.succeed(None)
    wrapped = self._wrap_channel_method('confirm_delivery')
    d = wrapped(ack_nack_callback=self._on_delivery_confirmation)

    def set_delivery_confirmation(result):
        self._delivery_confirmation = True
        self._delivery_message_id = 0
        LOGGER.debug('Delivery confirmation enabled.')
        return result

    d.addCallback(set_delivery_confirmation)
    # Unroutable messages returned after this point will be in the context
    # of publisher acknowledgments
    self._channel.add_on_return_callback(self._on_puback_message_returned)
    return d

exchange_bind

exchange_bind(
    destination: str,
    source: str,
    routing_key: str = '',
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Bind an exchange to another exchange.

PARAMETER DESCRIPTION
destination

The destination exchange to bind

TYPE: str

source

The source exchange to bind to

TYPE: str

routing_key

The routing key to bind on

TYPE: str DEFAULT: ''

arguments

Custom key/value pair arguments for the binding

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Exchange.BindOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def exchange_bind(
        self,
        destination: str,
        source: str,
        routing_key: str = '',
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Bind an exchange to another exchange.

    :param destination: The destination exchange to bind
    :param source: The source exchange to bind to
    :param routing_key: The routing key to bind on
    :param arguments: Custom key/value pair arguments for the binding
    :raises ValueError:
    :returns: Deferred that fires on the Exchange.BindOk response
    """
    return self._wrap_channel_method('exchange_bind')(
        destination=destination,
        source=source,
        routing_key=routing_key,
        arguments=arguments,
    )

exchange_declare

exchange_declare(
    exchange: str,
    exchange_type: str | ExchangeType = direct,
    passive: bool = False,
    durable: bool = False,
    auto_delete: bool = False,
    internal: bool = False,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

This method creates an exchange if it does not already exist, and if the exchange exists, verifies that it is of the correct and expected class.

If passive set, the server will reply with Declare-Ok if the exchange already exists with the same name, and raise an error if not and if the exchange does not already exist, the server MUST raise a channel exception with reply code 404 (not found).

PARAMETER DESCRIPTION
exchange

The exchange name consists of a non-empty sequence of these characters: letters, digits, hyphen, underscore, period, or colon

TYPE: str

exchange_type

The exchange type to use

TYPE: str | ExchangeType DEFAULT: direct

passive

Perform a declare or just check to see if it exists

TYPE: bool DEFAULT: False

durable

Survive a reboot of RabbitMQ

TYPE: bool DEFAULT: False

auto_delete

Remove when no more queues are bound to it

TYPE: bool DEFAULT: False

internal

Can only be published to by other exchanges

TYPE: bool DEFAULT: False

arguments

Custom key/value pair arguments for the exchange

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Exchange.DeclareOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def exchange_declare(
        self,
        exchange: str,
        exchange_type: str | ExchangeType = ExchangeType.direct,
        passive: bool = False,
        durable: bool = False,
        auto_delete: bool = False,
        internal: bool = False,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    This method creates an exchange if it does not already exist, and if the exchange exists,
    verifies that it is of the correct and expected class.

    If passive set, the server will reply with Declare-Ok if the exchange
    already exists with the same name, and raise an error if not and if the
    exchange does not already exist, the server MUST raise a channel
    exception with reply code 404 (not found).

    :param exchange: The exchange name consists of a non-empty sequence
        of these characters: letters, digits, hyphen, underscore, period,
        or colon
    :param exchange_type: The exchange type to use
    :param passive: Perform a declare or just check to see if it
        exists
    :param durable: Survive a reboot of RabbitMQ
    :param auto_delete: Remove when no more queues are bound to it
    :param internal: Can only be published to by other exchanges
    :param arguments: Custom key/value pair arguments for the exchange
    :returns: Deferred that fires on the Exchange.DeclareOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('exchange_declare')(
        exchange=exchange,
        exchange_type=exchange_type,
        passive=passive,
        durable=durable,
        auto_delete=auto_delete,
        internal=internal,
        arguments=arguments,
    )

exchange_delete

exchange_delete(
    exchange: str | None = None, if_unused: bool = False
) -> Deferred[Any]

Delete the exchange.

PARAMETER DESCRIPTION
exchange

The exchange name

TYPE: str | None DEFAULT: None

if_unused

only delete if the exchange is unused

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Exchange.DeleteOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def exchange_delete(self,
                    exchange: str | None = None,
                    if_unused: bool = False) -> defer.Deferred[Any]:
    """
    Delete the exchange.

    :param exchange: The exchange name
    :param if_unused: only delete if the exchange is unused
    :returns: Deferred that fires on the Exchange.DeleteOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('exchange_delete')(
        exchange=exchange,
        if_unused=if_unused,
    )

exchange_unbind

exchange_unbind(
    destination: str,
    source: str,
    routing_key: str = '',
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Unbind an exchange from another exchange.

PARAMETER DESCRIPTION
destination

The destination exchange to unbind

TYPE: str

source

The source exchange to unbind from

TYPE: str

routing_key

The routing key to unbind

TYPE: str DEFAULT: ''

arguments

Custom key/value pair arguments for the binding

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Exchange.UnbindOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def exchange_unbind(
        self,
        destination: str,
        source: str,
        routing_key: str = '',
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Unbind an exchange from another exchange.

    :param destination: The destination exchange to unbind
    :param source: The source exchange to unbind from
    :param routing_key: The routing key to unbind
    :param arguments: Custom key/value pair arguments for the binding
    :returns: Deferred that fires on the Exchange.UnbindOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('exchange_unbind')(
        destination=destination,
        source=source,
        routing_key=routing_key,
        arguments=arguments,
    )

flow

flow(active: bool = True) -> Deferred[Any]

Turn Channel flow control off and on.

Returns a Deferred that will fire with a bool indicating the channel flow state. For more information, please reference:

https://www.rabbitmq.com/amqp-0-9-1-reference.html#channel.flow

PARAMETER DESCRIPTION
active

Turn flow on or off

TYPE: bool DEFAULT: True

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires with the channel flow state

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def flow(self, active: bool = True) -> defer.Deferred[Any]:
    """
    Turn Channel flow control off and on.

    Returns a Deferred that will fire with a bool indicating the channel flow state. For more
    information, please reference:

    https://www.rabbitmq.com/amqp-0-9-1-reference.html#channel.flow

    :param active: Turn flow on or off
    :returns: Deferred that fires with the channel flow state
    :raises ValueError:
    """
    return self._wrap_channel_method('flow')(active=active)

open

open() -> None

Open the channel.

Source code in pika/adapters/twisted_connection.py
def open(self) -> None:
    """Open the channel."""
    return self._channel.open()

queue_bind

queue_bind(
    queue: str,
    exchange: str,
    routing_key: str | None = None,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Bind the queue to the specified exchange.

PARAMETER DESCRIPTION
queue

The queue to bind to the exchange

TYPE: str

exchange

The source exchange to bind to

TYPE: str

routing_key

The routing key to bind on

TYPE: str | None DEFAULT: None

arguments

Custom key/value pair arguments for the binding

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.BindOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_bind(
        self,
        queue: str,
        exchange: str,
        routing_key: str | None = None,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Bind the queue to the specified exchange.

    :param queue: The queue to bind to the exchange
    :param exchange: The source exchange to bind to
    :param routing_key: The routing key to bind on
    :param arguments: Custom key/value pair arguments for the binding
    :returns: Deferred that fires on the Queue.BindOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('queue_bind')(
        queue=queue,
        exchange=exchange,
        routing_key=routing_key,
        arguments=arguments,
    )

queue_declare

queue_declare(
    queue: str,
    passive: bool = False,
    durable: bool = False,
    exclusive: bool = False,
    auto_delete: bool = False,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Declare queue, create if needed.

This method creates or checks a queue. When creating a new queue the client can specify various properties that control the durability of the queue and its contents, and the level of sharing for the queue.

Use an empty string as the queue name for the broker to auto-generate one

PARAMETER DESCRIPTION
queue

The queue name; if empty string, the broker will create a unique queue name

TYPE: str

passive

Only check to see if the queue exists

TYPE: bool DEFAULT: False

durable

Survive reboots of the broker

TYPE: bool DEFAULT: False

exclusive

Only allow access by the current connection

TYPE: bool DEFAULT: False

auto_delete

Delete after consumer cancels or disconnects

TYPE: bool DEFAULT: False

arguments

Custom key/value arguments for the queue

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.DeclareOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_declare(
        self,
        queue: str,
        passive: bool = False,
        durable: bool = False,
        exclusive: bool = False,
        auto_delete: bool = False,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Declare queue, create if needed.

    This method creates or checks a queue. When creating a new queue the client can specify various
    properties that control the durability of the queue and its contents,
    and the level of sharing for the queue.

    Use an empty string as the queue name for the broker to auto-generate
    one

    :param queue: The queue name; if empty string, the broker will
        create a unique queue name
    :param passive: Only check to see if the queue exists
    :param durable: Survive reboots of the broker
    :param exclusive: Only allow access by the current connection
    :param auto_delete: Delete after consumer cancels or disconnects
    :param arguments: Custom key/value arguments for the queue
    :returns: Deferred that fires on the Queue.DeclareOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('queue_declare')(
        queue=queue,
        passive=passive,
        durable=durable,
        exclusive=exclusive,
        auto_delete=auto_delete,
        arguments=arguments,
    )

queue_delete

queue_delete(
    queue: str,
    if_unused: bool = False,
    if_empty: bool = False,
) -> Deferred[Any]

Delete a queue from the broker.

This method wraps :meth:Channel.queue_delete <pika.channel.Channel.queue_delete>, and removes the reference to the queue object after it gets deleted on the server.

PARAMETER DESCRIPTION
queue

The queue to delete

TYPE: str

if_unused

only delete if it's unused

TYPE: bool DEFAULT: False

if_empty

only delete if the queue is empty

TYPE: bool DEFAULT: False

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.DeleteOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_delete(self,
                 queue: str,
                 if_unused: bool = False,
                 if_empty: bool = False) -> defer.Deferred[Any]:
    """
    Delete a queue from the broker.

    This method wraps :meth:`Channel.queue_delete <pika.channel.Channel.queue_delete>`, and
    removes the reference to the queue object after it gets deleted on the server.

    :param queue: The queue to delete
    :param if_unused: only delete if it's unused
    :param if_empty: only delete if the queue is empty
    :returns: Deferred that fires on the Queue.DeleteOk response
    :raises ValueError:
    """
    wrapped = self._wrap_channel_method('queue_delete')
    d = wrapped(queue=queue, if_unused=if_unused, if_empty=if_empty)

    def _clear_consumer(ret, queue_name):
        for consumer_tag in list(
                self._queue_name_to_consumer_tags.get(queue_name, set())):
            self._consumers[consumer_tag].close(
                exceptions.ConsumerCancelled(
                    f'Queue {queue_name} was deleted.'))
            del self._consumers[consumer_tag]
            self._queue_name_to_consumer_tags[queue_name].remove(
                consumer_tag)
        return ret

    return d.addCallback(_clear_consumer, queue)

queue_purge

queue_purge(queue: str) -> Deferred[Any]

Purge all of the messages from the specified queue.

PARAMETER DESCRIPTION
queue

The queue to purge

TYPE: str

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.PurgeOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_purge(self, queue: str) -> defer.Deferred[Any]:
    """
    Purge all of the messages from the specified queue.

    :param queue: The queue to purge
    :returns: Deferred that fires on the Queue.PurgeOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('queue_purge')(queue=queue)

queue_unbind

queue_unbind(
    queue: str,
    exchange: str | None,
    routing_key: str | None = None,
    arguments: dict[str, Any] | None = None,
) -> Deferred[Any]

Unbind a queue from an exchange.

PARAMETER DESCRIPTION
queue

The queue to unbind from the exchange

TYPE: str

exchange

The source exchange to bind from

TYPE: str | None

routing_key

The routing key to unbind

TYPE: str | None DEFAULT: None

arguments

Custom key/value pair arguments for the binding

TYPE: dict[str, Any] | None DEFAULT: None

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Queue.UnbindOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def queue_unbind(
        self,
        queue: str,
        exchange: str | None,
        routing_key: str | None = None,
        arguments: dict[str, Any] | None = None) -> defer.Deferred[Any]:
    """
    Unbind a queue from an exchange.

    :param queue: The queue to unbind from the exchange
    :param exchange: The source exchange to bind from
    :param routing_key: The routing key to unbind
    :param arguments: Custom key/value pair arguments for the binding
    :returns: Deferred that fires on the Queue.UnbindOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('queue_unbind')(
        queue=queue,
        exchange=exchange,
        routing_key=routing_key,
        arguments=arguments,
    )

tx_commit

tx_commit() -> Deferred[Any]

Commit a transaction.

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Tx.CommitOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def tx_commit(self) -> defer.Deferred[Any]:
    """
    Commit a transaction.

    :returns: Deferred that fires on the Tx.CommitOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('tx_commit')()

tx_rollback

tx_rollback() -> Deferred[Any]

Rollback a transaction.

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Tx.RollbackOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def tx_rollback(self) -> defer.Deferred[Any]:
    """
    Rollback a transaction.

    :returns: Deferred that fires on the Tx.RollbackOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('tx_rollback')()

tx_select

tx_select() -> Deferred[Any]

Select standard transaction mode.

This method sets the channel to use standard transactions. The client must use this method at least once on a channel before using the Commit or Rollback methods.

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires on the Tx.SelectOk response

RAISES DESCRIPTION
ValueError
Source code in pika/adapters/twisted_connection.py
def tx_select(self) -> defer.Deferred[Any]:
    """
    Select standard transaction mode.

    This method sets the channel to use standard transactions. The client must use this method
    at least once on a channel before using the Commit or Rollback methods.

    :returns: Deferred that fires on the Tx.SelectOk response
    :raises ValueError:
    """
    return self._wrap_channel_method('tx_select')()

ClosableDeferredQueue

Bases: DeferredQueue

Like the normal Twisted DeferredQueue, but after close() is called with an exception instance all pending Deferreds are errbacked and further attempts to call get() or put() return a Failure wrapping that exception.

Source code in pika/adapters/twisted_connection.py
class ClosableDeferredQueue(defer.DeferredQueue):
    """Like the normal Twisted DeferredQueue, but after close() is called with an exception instance
    all pending Deferreds are errbacked and further attempts to call get() or put() return a Failure
    wrapping that exception.
    """

    def __init__(self,
                 size: int | None = None,
                 backlog: int | None = None) -> None:
        self.closed: Exception | None = None
        super().__init__(size, backlog)

    @override
    def put(  # type: ignore[override]
            self, obj: Any) -> defer.Deferred[Any] | None:
        """
        Like the original :meth:`DeferredQueue.put` method, but returns an errback if the queue is
        closed.

        :param obj: Object to put into the queue
        :returns: Nothing on success, or a failed Deferred if the queue is closed
        """
        if self.closed:
            LOGGER.error('Impossible to put to the queue, it is closed.')
            return defer.fail(self.closed)
        # `DeferredQueue.put()` returns nothing, so the success path yields
        # `None`. The return type stays optional for the closed path above.
        defer.DeferredQueue.put(self, obj)
        return None

    @override
    def get(self) -> defer.Deferred[Any]:
        """
        Returns a Deferred that will fire with the next item in the queue, when it's available.

        The Deferred will errback if the queue is closed.

        :returns: Deferred that fires with the next item.
        """
        if self.closed:
            LOGGER.error('Impossible to get from the queue, it is closed.')
            return defer.fail(self.closed)
        return defer.DeferredQueue.get(self)

    def close(self, reason: Exception | None) -> None:
        """
        Closes the queue.

        Errback the pending calls to :meth:`get()`.

        :param reason: The reason for closing the queue
        """
        if self.closed:
            LOGGER.warning('Queue was already closed with reason: %s.',
                           self.closed)
        self.closed = reason
        while self.waiting:
            self.waiting.pop().errback(reason)
        self.pending = []

closed instance-attribute

closed: Exception | None = None

__init__

__init__(
    size: int | None = None, backlog: int | None = None
) -> None
Source code in pika/adapters/twisted_connection.py
def __init__(self,
             size: int | None = None,
             backlog: int | None = None) -> None:
    self.closed: Exception | None = None
    super().__init__(size, backlog)

close

close(reason: Exception | None) -> None

Closes the queue.

Errback the pending calls to :meth:get().

PARAMETER DESCRIPTION
reason

The reason for closing the queue

TYPE: Exception | None

Source code in pika/adapters/twisted_connection.py
def close(self, reason: Exception | None) -> None:
    """
    Closes the queue.

    Errback the pending calls to :meth:`get()`.

    :param reason: The reason for closing the queue
    """
    if self.closed:
        LOGGER.warning('Queue was already closed with reason: %s.',
                       self.closed)
    self.closed = reason
    while self.waiting:
        self.waiting.pop().errback(reason)
    self.pending = []

get

get() -> Deferred[Any]

Returns a Deferred that will fire with the next item in the queue, when it's available.

The Deferred will errback if the queue is closed.

RETURNS DESCRIPTION
Deferred[Any]

Deferred that fires with the next item.

Source code in pika/adapters/twisted_connection.py
@override
def get(self) -> defer.Deferred[Any]:
    """
    Returns a Deferred that will fire with the next item in the queue, when it's available.

    The Deferred will errback if the queue is closed.

    :returns: Deferred that fires with the next item.
    """
    if self.closed:
        LOGGER.error('Impossible to get from the queue, it is closed.')
        return defer.fail(self.closed)
    return defer.DeferredQueue.get(self)

put

put(obj: Any) -> Deferred[Any] | None

Like the original :meth:DeferredQueue.put method, but returns an errback if the queue is closed.

PARAMETER DESCRIPTION
obj

Object to put into the queue

TYPE: Any

RETURNS DESCRIPTION
Deferred[Any] | None

Nothing on success, or a failed Deferred if the queue is closed

Source code in pika/adapters/twisted_connection.py
@override
def put(  # type: ignore[override]
        self, obj: Any) -> defer.Deferred[Any] | None:
    """
    Like the original :meth:`DeferredQueue.put` method, but returns an errback if the queue is
    closed.

    :param obj: Object to put into the queue
    :returns: Nothing on success, or a failed Deferred if the queue is closed
    """
    if self.closed:
        LOGGER.error('Impossible to put to the queue, it is closed.')
        return defer.fail(self.closed)
    # `DeferredQueue.put()` returns nothing, so the success path yields
    # `None`. The return type stays optional for the closed path above.
    defer.DeferredQueue.put(self, obj)
    return None