Class: TypedBus::DeadLetterQueue
- Inherits:
-
Object
- Object
- TypedBus::DeadLetterQueue
- 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
- #clear! ⇒ Object
-
#drain ⇒ Array<Delivery>
Yield each dead letter and remove it from the queue.
-
#each(&block) ⇒ Object
Iterate over dead letters without removing them.
- #empty? ⇒ Boolean
-
#initialize(&on_dead_letter) ⇒ DeadLetterQueue
constructor
A new instance of DeadLetterQueue.
-
#on_dead_letter(&block) ⇒ Object
Register or replace the callback fired when entries arrive.
-
#push(delivery) ⇒ Object
Add a failed delivery.
- #size ⇒ Object
Constructor Details
#initialize(&on_dead_letter) ⇒ DeadLetterQueue
Returns a new instance of DeadLetterQueue.
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 |
#drain ⇒ Array<Delivery>
Yield each dead letter and remove it from the queue.
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
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 |
#size ⇒ Object
23 24 25 |
# File 'lib/typed_bus/dead_letter_queue.rb', line 23 def size @entries.size end |