Class: Raif::Evals::WorkerPool
- Inherits:
-
Object
- Object
- Raif::Evals::WorkerPool
- Defined in:
- lib/raif/evals/worker_pool.rb
Overview
Runs a list of work items across a bounded pool of forked worker processes.
An eval run is almost entirely waiting on provider HTTP responses, so overlapping the waiting is where the wall clock goes. Processes rather than threads, because each worker needs a database of its own: every execution runs in a transaction that is rolled back, and two executions sharing a database block on any unique index both of them write the same key to until the first one rolls back. Seeded reference rows and a repeated case do that on every execution, so on one shared database those executions run one at a time.
The work runs in a child, and everything that records it runs here in the parent: work
returns a value that crosses the process boundary with Marshal, and collect receives it.
Defined Under Namespace
Classes: Dispatch, WorkerError
Instance Attribute Summary collapse
-
#concurrency ⇒ Object
readonly
Returns the value of attribute concurrency.
Instance Method Summary collapse
-
#initialize(concurrency: 1, setup_worker: nil) ⇒ WorkerPool
constructor
A new instance of WorkerPool.
-
#run(items, work:, collect:, dispatched: nil) ⇒ Object
Calls
workonce per item in a worker andcollectwith each item and its value in this process, in completion order.
Constructor Details
#initialize(concurrency: 1, setup_worker: nil) ⇒ WorkerPool
Returns a new instance of WorkerPool.
25 26 27 28 |
# File 'lib/raif/evals/worker_pool.rb', line 25 def initialize(concurrency: 1, setup_worker: nil) @concurrency = [concurrency.to_i, 1].max @setup_worker = setup_worker end |
Instance Attribute Details
#concurrency ⇒ Object (readonly)
Returns the value of attribute concurrency.
21 22 23 |
# File 'lib/raif/evals/worker_pool.rb', line 21 def concurrency @concurrency end |
Instance Method Details
#run(items, work:, collect:, dispatched: nil) ⇒ Object
Calls work once per item in a worker and collect with each item and its value in this
process, in completion order. dispatched is called here with each item as a worker takes
it. Both callbacks only ever run on the calling thread, so they need no locking.
On Ctrl-C, workers stop taking items and the in-flight ones are waited for rather than
killed, so their results still reach collect. The Interrupt is then re-raised for the
caller to report on. A second Ctrl-C stops waiting and kills the workers.
37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 |
# File 'lib/raif/evals/worker_pool.rb', line 37 def run(items, work:, collect:, dispatched: nil) return items if items.empty? # Only at concurrency 1, not for a lone item at a higher one: an item run here uses this # process's database rather than a worker's, and its result must not depend on how many # items happen to be left. if concurrency == 1 return items.each do |item| dispatched&.call(item) collect.call(item, work.call(item)) end end Dispatch.new(items: items, work: work, collect: collect, dispatched: dispatched, setup_worker: @setup_worker) .run(worker_count: [concurrency, items.size].min) end |