class Riffer::Runner::Fibers
rbs_inline: enabled
Public Class Methods
Source
# File lib/riffer/runner/fibers.rb, line 9 def initialize(max_concurrency: nil) super() depends_on "async" depends_on "async/semaphore" if max_concurrency @max_concurrency = max_concurrency end
Calls superclass method
Public Instance Methods
Source
# File lib/riffer/runner/fibers.rb, line 18 def map(items, context:, &block) return [] if items.empty? results = Array.new(items.size) errors = Array.new(items.size) barrier = Async::Barrier.new max = @max_concurrency parent = if max Async::Semaphore.new(max, parent: barrier) else barrier end # Sync joins the running reactor task if there is one, otherwise starts its own. Sync do items.each_with_index do |item, index| parent.async do results[index] = yield(item) rescue StandardError => e errors[index] = e end end barrier.wait ensure barrier.stop end first_error = errors.compact.first raise first_error if first_error results end