Class: Mutineer::WorkerPool

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

Constructor Details

#initialize(size) ⇒ WorkerPool

Builds a pool.

Parameters:

  • size (Integer) —

    desired pool size.



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.

Parameters:

  • data (String) —

    marshaled payload.

Returns:



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.message}")
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.message}")
          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.

Parameters:

  • items (Array<Array>) —

    work items.

  • stop_when (Proc, nil) (defaults to: nil) —

    called with each collected Result; when it returns truthy, no further items are scheduled and the run drains and returns early (--fail-fast). Unscheduled slots stay nil.

  • on_result (Proc, nil) (defaults to: nil) —

    called in the parent with each collected Result as it is reaped, in finish order (progress reporting); its return value is ignored. Must be fast and non-blocking: it runs on the pool's single reap thread, so a slow callback stalls draining the other in-flight children's pipes (the #4 deadlock discipline).

Yield Parameters:

  • item (Array) —

    one work item.

Returns:

  • (Array<Mutineer::Result>) —

    results in input order (nil for any item left unscheduled by an early stop).



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