Class: TypedBus::DeadLetterQueue

Inherits:
Object
  • Object
show all
Includes:
Enumerable
Defined in:
lib/typed_bus/dead_letter_queue.rb

Overview

Stores deliveries that failed — either explicitly NACKed or timed out.

Each Channel owns a DLQ. Failed deliveries accumulate here for inspection, retry, or manual drain.

Instance Method Summary collapse

Constructor Details

#initialize(&on_dead_letter) ⇒ DeadLetterQueue

Returns a new instance of DeadLetterQueue.

Parameters:

  • on_dead_letter (Proc, nil)

    optional callback fired on each push



10
11
12
13
# File 'lib/typed_bus/dead_letter_queue.rb', line 10

def initialize(&on_dead_letter)
  @entries = []
  @on_dead_letter = on_dead_letter
end

Instance Method Details

#clear!Object



48
49
50
51
52
# File 'lib/typed_bus/dead_letter_queue.rb', line 48

def clear!
  count = @entries.size
  @entries.clear
  log(:info, "cleared #{count} entries") if count > 0
end

#drainArray<Delivery>

Yield each dead letter and remove it from the queue.

Returns:

  • (Array<Delivery>)

    the drained entries



40
41
42
43
44
45
46
# File 'lib/typed_bus/dead_letter_queue.rb', line 40

def drain
  drained = @entries.dup
  @entries.clear
  log(:info, "drained #{drained.size} entries")
  drained.each { |d| yield d } if block_given?
  drained
end

#each(&block) ⇒ Object

Iterate over dead letters without removing them.



32
33
34
# File 'lib/typed_bus/dead_letter_queue.rb', line 32

def each(&block)
  @entries.each(&block)
end

#empty?Boolean

Returns:

  • (Boolean)


27
28
29
# File 'lib/typed_bus/dead_letter_queue.rb', line 27

def empty?
  @entries.empty?
end

#on_dead_letter(&block) ⇒ Object

Register or replace the callback fired when entries arrive.



55
56
57
# File 'lib/typed_bus/dead_letter_queue.rb', line 55

def on_dead_letter(&block)
  @on_dead_letter = block
end

#push(delivery) ⇒ Object

Add a failed delivery.



16
17
18
19
20
21
# File 'lib/typed_bus/dead_letter_queue.rb', line 16

def push(delivery)
  @entries << delivery
  reason = delivery.timed_out? ? "timeout" : "nack"
  log(:warn, "entry added (channel=:#{delivery.channel_name}, subscriber=##{delivery.subscriber_id}, reason=#{reason}, total=#{@entries.size})")
  @on_dead_letter&.call(delivery)
end

#sizeObject



23
24
25
# File 'lib/typed_bus/dead_letter_queue.rb', line 23

def size
  @entries.size
end