Class: TransactionalOutbox::Relay::WorkerSet

Inherits:
Object
  • Object
show all
Defined in:
lib/transactional_outbox/relay/worker_set.rb,
lib/transactional_outbox/relay/worker_set/worker.rb,
lib/transactional_outbox/relay/worker_set/processor.rb

Defined Under Namespace

Classes: Processor, Worker

Instance Method Summary collapse

Constructor Details

#initializeWorkerSet

Returns a new instance of WorkerSet.



8
9
10
# File 'lib/transactional_outbox/relay/worker_set.rb', line 8

def initialize
  @workers = {}
end

Instance Method Details

#add_worker(queue) ⇒ Object



14
15
16
17
18
# File 'lib/transactional_outbox/relay/worker_set.rb', line 14

def add_worker(queue)
  return if workers.key?(queue)

  create_worker(queue)
end

#all_stopped?Boolean

Returns:

  • (Boolean)


29
# File 'lib/transactional_outbox/relay/worker_set.rb', line 29

def all_stopped? = workers.values.all?(&:stopped?)

#get_worker(queue) ⇒ Object



12
# File 'lib/transactional_outbox/relay/worker_set.rb', line 12

def get_worker(queue) = workers[queue]

#stop_workersObject



28
# File 'lib/transactional_outbox/relay/worker_set.rb', line 28

def stop_workers = workers.each_value(&:shutdown)

#try_to_recover_worker(queue) ⇒ Object



20
21
22
23
24
25
26
# File 'lib/transactional_outbox/relay/worker_set.rb', line 20

def try_to_recover_worker(queue)
  worker = get_worker(queue)

  return if !worker.nil? && (worker.shutting_down? || !worker.stopped?)

  create_worker(queue)
end