class Crabbit::Consumer

Overview

Recovering pull or callback consumer for one RabbitMQ stream.

Pull mode uses #receive, #receive?, or #each. Callback mode is created by Environment#consumer(stream, options) { |delivery| ... } and marks each delivery processed after the handler returns. A consumer reconnects and resumes from the latest contiguous broker-delivery prefix until closed.

Included Modules

Defined in:

crabbit/consumer.cr

Constructors

Instance Method Summary

Constructor Detail

def self.new(environment : Environment, stream : String, options : ConsumerOptions, handler : Proc(Delivery, Nil) | Nil = nil) #

Creates and subscribes a consumer.

Applications normally use the Environment#consumer overload matching pull or callback mode.


[View source]

Instance Method Detail

def close : Nil #

Idempotently unsubscribes, stores any pending automatic offset, and closes the delivery queue.


[View source]
def closed? : Bool #

Returns whether the consumer was permanently closed.


[View source]
def each(&block : Delivery -> ) : Nil #

Yields deliveries until the consumer closes.

Enumerable consumption acknowledges after the block returns, including when it raises. Use #receive for explicit acknowledgement control.


[View source]
def open? : Bool #

Returns whether the consumer currently has an active subscription.


[View source]
def options : ConsumerOptions #

Returns the immutable consumer options.


[View source]
def receive : Delivery #

Receives one delivery, waiting while the queue is empty.

The caller owns acknowledgement and must invoke Delivery#processed! when processing is complete. Raises ResourceClosedError after close.


[View source]
def receive? : Delivery | Nil #

Receives one delivery, waiting while the queue is empty, or returns nil after the consumer closes and the queue drains.

The caller must invoke Delivery#processed! for every returned delivery.


[View source]
def state : ResourceState #

Returns the current lifecycle state.


[View source]
def store_offset(offset : UInt64) : Nil #

Stores an absolute offset for this named consumer.

Raises ConfigurationError when ConsumerOptions#name is absent.


[View source]
def store_offset(delivery : Delivery) : Nil #

Stores the offset of delivery for this named consumer.

Raises ArgumentError when the delivery belongs to another stream.


[View source]
def stored_offset : UInt64 | Nil #

Returns this named consumer's stored offset, or nil when absent.


[View source]
def stream : String #

Returns the subscribed stream name.


[View source]