class Async::Barrier

@public Since *Async v1*.
A general purpose synchronisation primitive, which allows one task to wait for a number of other tasks to complete. It can be used in conjunction with {Semaphore}.

def async(*arguments, parent: (@parent or Task.current), **options, &block)

@returns [Task] The task which was created to execute the block.
@asynchronous Executes the given block concurrently.
Execute a child task and add it to the barrier.
def async(*arguments, parent: (@parent or Task.current), **options, &block)
	raise "Barrier is stopped!" if @finished.closed?
	
	waiting = nil
	
	task = parent.async(*arguments, **options) do |task, *arguments|
		# Create a new list node for the task and add it to the list of waiting tasks:
		node = TaskNode.new(task)
		@tasks.append(node)
		
		# Signal the outer async block that we have added the task to the list of waiting tasks, and that it can now wait for it to finish:
		waiting = node
		@condition.signal
		
		# Invoke the block, which may raise an error. If it does, we will still signal that the task has finished:
		block.call(task, *arguments)
	ensure
		# Signal that the task has finished, which will unblock the waiting task:
		@finished.signal(node) unless @finished.closed?
	end
	
	# `parent.async` may yield before the child block executes, so we wait here until the child has appended itself to `@tasks`, ensuring `wait` cannot return early and miss tracking it:
	@condition.wait while waiting.nil?
	
	return task
end

def cancel

@asynchronous May wait for tasks to finish executing.
Cancel all tasks held by the barrier.
def cancel
	@tasks.each do |waiting|
		waiting.task.cancel
	end
	
	@finished.close
end

def empty?

@returns [Boolean]
Whether there are any tasks being held by the barrier.
def empty?
	@tasks.empty?
end

def initialize(parent: nil)

@public Since *Async v1*.
@parameter parent [Task | Semaphore | Nil] The parent for holding any children tasks.
Initialize the barrier.
def initialize(parent: nil)
	@tasks = List.new
	@finished = Queue.new
	@condition = Condition.new
	
	@parent = parent
end

def size

Number of tasks being held by the barrier.
def size
	@tasks.size
end

def stop

Deprecated:
  • Use {#cancel} instead.
def stop
	cancel
end

def wait

@asynchronous Will wait for tasks to finish executing.

@returns [Integer | Nil] The number of tasks which were waited for, or `nil` if there were no tasks to wait for.
@yields {|task| ...} If a block is given, the unwaited task is yielded. You must invoke {Task#wait} yourself. In addition, you may `break` if you have captured enough results.

Wait for all tasks to complete by invoking {Task#wait} on each waiting task, which may raise an error. As long as the task has completed, it will be removed from the barrier.
def wait
	return nil if @tasks.empty?
	count = 0
	
	while true
		# Wait for a task to finish (we get the task node):
		break unless waiting = @finished.wait
		count += 1
		
		# Remove the task as it is now finishing:
		@tasks.remove?(waiting)
		
		# Get the task:
		task = waiting.task
		
		# If a block is given, the user can implement their own behaviour:
		if block_given?
			yield task
		else
			# Wait for it to either complete or raise an error:
			task.wait
		end
		
		break if @tasks.empty?
	end
	
	return count
end