Class: Mutineer::WorkerPool
- Inherits:
-
Object
- Object
- Mutineer::WorkerPool
- Defined in:
- lib/mutineer/worker_pool.rb
Overview
Fixed-size fork pool. run forks up to size children at once; each child
runs the block on one work item, marshals its Result to a private pipe, and
exits. The parent drains pipes with IO.select and reaps each finished child
by known pid (never wait2(-1), which would steal the host suite's children),
opening exactly one slot per reap, then refills. Results are returned in the
SAME ORDER as items regardless of finish order, so verdicts are identical
to a serial run and downstream output is stable.
The block is run inside the child via yield(*items[i]); whatever it
returns (a Result) is the marshaled payload. Per-mutant timeout is handled
one level down by Isolation. The pool adds no separate wall clock.
Instance Method Summary collapse
-
#decode(data) ⇒ Mutineer::Result
private
A partial/garbage Marshal stream (dead worker) must not crash the pool.
-
#fill(items, queue, running) ⇒ Object
private
The child must ALWAYS hard-exit.
-
#initialize(size) ⇒ WorkerPool
constructor
Builds a pool.
-
#reap(results, running) ⇒ Object
private
Drain pipes with IO.select and reap a child only on EOF.
-
#run(items, stop_when: nil, on_result: nil) {|item| ... } ⇒ Array<Mutineer::Result>
Runs the work items through the pool.
Constructor Details
#initialize(size) ⇒ WorkerPool
Builds a pool.
21 22 23 |
# File 'lib/mutineer/worker_pool.rb', line 21 def initialize(size) @size = [size.to_i, 1].max end |
Instance Method Details
#decode(data) ⇒ Mutineer::Result (private)
A partial/garbage Marshal stream (dead worker) must not crash the pool. Degrade to an error Result.
137 138 139 140 141 142 143 |
# File 'lib/mutineer/worker_pool.rb', line 137 def decode(data) return Result.error("worker produced no result") if data.empty? Marshal.load(data) rescue StandardError => e Result.error("worker result unreadable: #{e.class}: #{e.}") end |
#fill(items, queue, running) ⇒ Object (private)
The child must ALWAYS hard-exit. If yield raises, marshal an error Result
and exit! in ensure. Otherwise the child unwinds normally and our
Minitest at_exit autorun re-runs the parent suite inside the worker,
losing the real error.
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 90 91 92 93 94 95 96 97 98 99 100 101 |
# File 'lib/mutineer/worker_pool.rb', line 64 def fill(items, queue, running) while running.size < @size && !queue.empty? idx = queue.shift rd, wr = IO.pipe rd.binmode # Marshal output is binary: keep the pipe byte-exact wr.binmode begin pid = fork do rd.close payload = begin yield(*items[idx]) rescue Exception => e # rubocop:disable Lint/RescueException Result.error("worker crashed: #{e.class}: #{e.}") end begin wr.write(Marshal.dump(payload)) rescue StandardError # rubocop:disable Lint/SuppressedException # pipe gone; parent will record "no result" ensure wr.close exit!(0) end end rescue Errno::EAGAIN # Process table is full. Put the item back and reap before retrying; # if nothing is running we cannot make progress, so re-raise. rd.close wr.close raise if running.empty? queue.unshift(idx) return end wr.close running[pid] = [idx, rd, +""] # idx, read end, accumulated bytes end end |
#reap(results, running) ⇒ Object (private)
Drain pipes with IO.select and reap a child only on EOF. The old code
reaped first and read after, but a child whose payload exceeds the OS pipe
buffer (~64KB) blocks on write before it can exit, so it was never reaped
and the pool deadlocked. Reading concurrently keeps the pipe drained so the
child can finish and exit; EOF means it closed its write end (done writing).
We waitpid only OUR known pids (never wait2(-1), which would steal the host
suite's children).
110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 |
# File 'lib/mutineer/worker_pool.rb', line 110 def reap(results, running) return if running.empty? loop do rds = running.values.map { |v| v[1] } ready, = IO.select(rds, nil, nil) ready.each do |rd| pid, (idx, _io, buf) = running.find { |_, v| v[1].equal?(rd) } chunk = rd.read_nonblock(65_536, exception: false) next if chunk == :wait_readable if chunk.nil? # EOF: child closed wr; it is done writing and exiting rd.close Process.waitpid(pid) # reap the now-finished child (no zombie) running.delete(pid) # Return the collected Result so the caller's stop_when can see it. return results[idx] = decode(buf) end buf << chunk end end end |
#run(items, stop_when: nil, on_result: nil) {|item| ... } ⇒ Array<Mutineer::Result>
Runs the work items through the pool.
39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 |
# File 'lib/mutineer/worker_pool.rb', line 39 def run(items, stop_when: nil, on_result: nil) results = Array.new(items.size) queue = (0...items.size).to_a running = {} # pid => [index, read_io, buffer] stopping = false until queue.empty? && running.empty? fill(items, queue, running) { |*args| yield(*args) } unless stopping result = reap(results, running) on_result&.call(result) if result if !stopping && stop_when && result && stop_when.call(result) stopping = true queue.clear # schedule no more; let in-flight workers drain end end results end |