Class: TypedBus::DeliveryTracker
- Inherits:
-
Object
- Object
- TypedBus::DeliveryTracker
- Defined in:
- lib/typed_bus/delivery_tracker.rb
Overview
Tracks the delivery state of a single published message across N subscribers.
Created per publish call. Each subscriber gets a Delivery envelope;
the tracker aggregates their ack/nack responses and fires callbacks
when the message is fully delivered or needs dead-lettering.
Instance Attribute Summary collapse
-
#channel_name ⇒ Object
readonly
Returns the value of attribute channel_name.
-
#message ⇒ Object
readonly
Returns the value of attribute message.
Instance Method Summary collapse
-
#ack(subscriber_id) ⇒ Object
Record an ACK from a subscriber.
-
#fully_delivered? ⇒ Boolean
True when every subscriber has ACKed.
-
#fully_resolved? ⇒ Boolean
True when all subscribers have resolved (acked or nacked).
-
#initialize(message, channel_name:, subscriber_ids:) ⇒ DeliveryTracker
constructor
A new instance of DeliveryTracker.
-
#nack(subscriber_id) ⇒ Object
Record a NACK from a subscriber.
-
#on_complete(&block) ⇒ Object
Register a callback fired when all subscribers have resolved and every one ACKed (successful delivery).
-
#on_dead_letter(&block) ⇒ Object
Register a callback fired for each NACK, receiving the subscriber_id.
-
#on_resolved(&block) ⇒ Object
Register a callback fired when all subscribers have resolved, regardless of outcome (acked or nacked).
-
#pending_count ⇒ Object
Number of subscribers still pending.
Constructor Details
#initialize(message, channel_name:, subscriber_ids:) ⇒ DeliveryTracker
Returns a new instance of DeliveryTracker.
15 16 17 18 19 20 21 22 23 |
# File 'lib/typed_bus/delivery_tracker.rb', line 15 def initialize(, channel_name:, subscriber_ids:) @message = @channel_name = channel_name @states = subscriber_ids.each_with_object({}) { |id, h| h[id] = :pending } @on_complete = nil @on_resolved = nil @on_dead_letter = nil @resolved = false end |
Instance Attribute Details
#channel_name ⇒ Object (readonly)
Returns the value of attribute channel_name.
10 11 12 |
# File 'lib/typed_bus/delivery_tracker.rb', line 10 def channel_name @channel_name end |
#message ⇒ Object (readonly)
Returns the value of attribute message.
10 11 12 |
# File 'lib/typed_bus/delivery_tracker.rb', line 10 def @message end |
Instance Method Details
#ack(subscriber_id) ⇒ Object
Record an ACK from a subscriber.
26 27 28 29 |
# File 'lib/typed_bus/delivery_tracker.rb', line 26 def ack(subscriber_id) set_state(subscriber_id, :acked) check_resolution end |
#fully_delivered? ⇒ Boolean
True when every subscriber has ACKed.
39 40 41 |
# File 'lib/typed_bus/delivery_tracker.rb', line 39 def fully_delivered? @states.values.all? { |s| s == :acked } end |
#fully_resolved? ⇒ Boolean
True when all subscribers have resolved (acked or nacked).
44 45 46 |
# File 'lib/typed_bus/delivery_tracker.rb', line 44 def fully_resolved? @states.values.none? { |s| s == :pending } end |
#nack(subscriber_id) ⇒ Object
Record a NACK from a subscriber.
32 33 34 35 36 |
# File 'lib/typed_bus/delivery_tracker.rb', line 32 def nack(subscriber_id) set_state(subscriber_id, :nacked) fire_dead_letter(subscriber_id) check_resolution end |
#on_complete(&block) ⇒ Object
Register a callback fired when all subscribers have resolved and every one ACKed (successful delivery).
55 56 57 |
# File 'lib/typed_bus/delivery_tracker.rb', line 55 def on_complete(&block) @on_complete = block end |
#on_dead_letter(&block) ⇒ Object
Register a callback fired for each NACK, receiving the subscriber_id.
66 67 68 |
# File 'lib/typed_bus/delivery_tracker.rb', line 66 def on_dead_letter(&block) @on_dead_letter = block end |
#on_resolved(&block) ⇒ Object
Register a callback fired when all subscribers have resolved, regardless of outcome (acked or nacked). Use for cleanup.
61 62 63 |
# File 'lib/typed_bus/delivery_tracker.rb', line 61 def on_resolved(&block) @on_resolved = block end |
#pending_count ⇒ Object
Number of subscribers still pending.
49 50 51 |
# File 'lib/typed_bus/delivery_tracker.rb', line 49 def pending_count @states.count { |_, s| s == :pending } end |