Class: TypedBus::DeliveryTracker

Inherits:
Object
  • Object
show all
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

Instance Method Summary collapse

Constructor Details

#initialize(message, channel_name:, subscriber_ids:) ⇒ DeliveryTracker

Returns a new instance of DeliveryTracker.

Parameters:

  • message (Object)

    the published payload

  • channel_name (Symbol)
  • subscriber_ids (Array<Integer>)

    set of subscribers for this message



15
16
17
18
19
20
21
22
23
# File 'lib/typed_bus/delivery_tracker.rb', line 15

def initialize(message, channel_name:, subscriber_ids:)
  @message      = 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_nameObject (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

#messageObject (readonly)

Returns the value of attribute message.



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

def message
  @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.

Returns:

  • (Boolean)


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).

Returns:

  • (Boolean)


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_countObject

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