Class: Omnizip::Parallel::WorkerPool

Inherits:
Object
  • Object
show all
Defined in:
lib/omnizip/parallel/worker_pool.rb

Overview

Worker pool wrapper for Fractor-based parallel processing

Manages a pool of Fractor workers for parallel compression/extraction. Handles job distribution, result collection, and graceful shutdown.

Examples:

Create and use worker pool

pool = Omnizip::Parallel::WorkerPool.new(
  worker_class: CompressionWorker,
  num_workers: 4
)
pool.start
pool.submit(work_item)
results = pool.results
pool.shutdown

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(worker_class:, num_workers: nil, continuous: false) ⇒ WorkerPool

Initialize worker pool

Parameters:

  • worker_class (Class)

    Fractor::Worker subclass

  • num_workers (Integer) (defaults to: nil)

    number of worker threads

  • continuous (Boolean) (defaults to: false)

    continuous mode for long-running tasks



39
40
41
42
43
44
45
46
47
48
49
# File 'lib/omnizip/parallel/worker_pool.rb', line 39

def initialize(worker_class:, num_workers: nil, continuous: false)
  @worker_class = worker_class
  @num_workers = num_workers || detect_cpu_count
  @continuous = continuous
  @supervisor = nil
  @results = []
  @errors = []
  @running = false
  @result_mutex = Mutex.new
  @work_queue = nil
end

Instance Attribute Details

#errorsArray (readonly)

Returns collected errors.

Returns:

  • (Array)

    collected errors



29
30
31
# File 'lib/omnizip/parallel/worker_pool.rb', line 29

def errors
  @errors
end

#resultsArray (readonly)

Returns collected results.

Returns:

  • (Array)

    collected results



26
27
28
# File 'lib/omnizip/parallel/worker_pool.rb', line 26

def results
  @results
end

#runningBoolean (readonly)

Returns whether pool is running.

Returns:

  • (Boolean)

    whether pool is running



32
33
34
# File 'lib/omnizip/parallel/worker_pool.rb', line 32

def running
  @running
end

#supervisorFractor::Supervisor (readonly)

Returns underlying Fractor supervisor.

Returns:

  • (Fractor::Supervisor)

    underlying Fractor supervisor



23
24
25
# File 'lib/omnizip/parallel/worker_pool.rb', line 23

def supervisor
  @supervisor
end

Instance Method Details

#complete?Boolean

Check if pool has completed all work

Returns:

  • (Boolean)

    true if complete



179
180
181
182
183
184
185
186
# File 'lib/omnizip/parallel/worker_pool.rb', line 179

def complete?
  return false if @continuous
  return false unless @supervisor

  # In batch mode, check if all work is processed
  result_aggregator = @supervisor.results
  result_aggregator && !@supervisor.work_queue.empty?
end

#failed_resultsArray

Get failed results

Returns:

  • (Array)

    array of errors



156
157
158
# File 'lib/omnizip/parallel/worker_pool.rb', line 156

def failed_results
  @result_mutex.synchronize { @errors.dup }
end

#runvoid

This method returns an undefined value.

Run the pool in batch mode and wait for completion



118
119
120
121
122
123
124
125
126
# File 'lib/omnizip/parallel/worker_pool.rb', line 118

def run
  raise "Can only run in batch mode" if @continuous
  raise "Worker pool not started" unless @running

  @supervisor.run

  # Collect results
  collect_results
end

#shutdown(timeout: 30) ⇒ void

This method returns an undefined value.

Shutdown the worker pool

Parameters:

  • timeout (Numeric) (defaults to: 30)

    timeout in seconds



132
133
134
135
136
137
138
139
140
141
142
143
144
# File 'lib/omnizip/parallel/worker_pool.rb', line 132

def shutdown(timeout: 30)
  return unless @running

  if @continuous
    @supervisor.stop
    @supervisor_thread&.join(timeout)
  end

  # Collect final results
  collect_results

  @running = false
end

#startvoid

This method returns an undefined value.

Start the worker pool



54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
# File 'lib/omnizip/parallel/worker_pool.rb', line 54

def start
  return if @running

  # Create Fractor supervisor with worker pool configuration
  @supervisor = Fractor::Supervisor.new(
    worker_pools: [
      {
        worker_class: @worker_class,
        num_workers: @num_workers,
      },
    ],
    continuous_mode: @continuous,
  )

  # For continuous mode, set up work queue
  if @continuous
    @work_queue = Fractor::WorkQueue.new
    @work_queue.register_with_supervisor(@supervisor)
  end

  @running = true

  # Start supervisor in background thread for continuous mode
  if @continuous
    @supervisor_thread = Thread.new do
      @supervisor.run
    rescue StandardError => e
      @result_mutex.synchronize do
        @errors << { error: e, message: "Supervisor error: #{e.message}" }
      end
    end
  else
    # For batch mode, don't start yet - wait for work items
    @supervisor.start_workers
  end
end

#statsHash

Get pool statistics

Returns:

  • (Hash)

    statistics hash



163
164
165
166
167
168
169
170
171
172
173
174
# File 'lib/omnizip/parallel/worker_pool.rb', line 163

def stats
  return {} unless @supervisor

  {
    workers: @num_workers,
    running: @running,
    continuous: @continuous,
    results: @results.size,
    errors: @errors.size,
    total_processed: @results.size + @errors.size,
  }
end

#submit(work) ⇒ void

This method returns an undefined value.

Submit work item to the pool

Parameters:

  • work (Fractor::Work)

    work item to process



95
96
97
98
99
100
101
102
103
104
105
# File 'lib/omnizip/parallel/worker_pool.rb', line 95

def submit(work)
  raise "Worker pool not started" unless @running

  if @continuous
    # In continuous mode, add to work queue
    @work_queue << work
  else
    # In batch mode, add to supervisor
    @supervisor.add_work_item(work)
  end
end

#submit_batch(works) ⇒ void

This method returns an undefined value.

Submit multiple work items

Parameters:

  • works (Array<Fractor::Work>)

    array of work items



111
112
113
# File 'lib/omnizip/parallel/worker_pool.rb', line 111

def submit_batch(works)
  works.each { |work| submit(work) }
end

#successful_resultsArray

Get successful results

Returns:

  • (Array)

    array of successful results



149
150
151
# File 'lib/omnizip/parallel/worker_pool.rb', line 149

def successful_results
  @result_mutex.synchronize { @results.dup }
end