Module: Kettle::Family::Concurrency

Defined in:
lib/kettle/family/concurrency.rb

Overview

Allocates CPU capacity between concurrent family members and optional member-internal workers. A member consumes one wave slot before any command-specific workers are considered.

Class Method Summary collapse

Class Method Details

.default_wave_jobs(cpu_count: Etc.nprocessors) ⇒ Object



13
14
15
# File 'lib/kettle/family/concurrency.rb', line 13

def default_wave_jobs(cpu_count: Etc.nprocessors)
  [cpu_count / 2, 1].max
end

.internal_worker_limit(wave_jobs:, cpu_count: Etc.nprocessors) ⇒ Object

Raises:

  • (ArgumentError)


24
25
26
27
28
29
30
# File 'lib/kettle/family/concurrency.rb', line 24

def internal_worker_limit(wave_jobs:, cpu_count: Etc.nprocessors)
  raise ArgumentError, "wave jobs must be positive" unless wave_jobs.to_i.positive?

  half_cores = default_wave_jobs(cpu_count: cpu_count)
  available_per_member = (cpu_count - wave_jobs.to_i) / wave_jobs.to_i
  available_per_member.clamp(1, half_cores)
end

.test_process_ceiling(wave_jobs:, cpu_count: Etc.nprocessors) ⇒ Object

Test runners count their primary worker as part of the requested process pool. Reserve one process for each wave member, then divide the remaining capacity across them. No member receives more than half of the detected CPUs, including that primary process.



36
37
38
# File 'lib/kettle/family/concurrency.rb', line 36

def test_process_ceiling(wave_jobs:, cpu_count: Etc.nprocessors)
  [default_wave_jobs(cpu_count: cpu_count), 1 + internal_worker_limit(wave_jobs: wave_jobs, cpu_count: cpu_count)].min
end

.wave_jobs(requested:, item_count:, cpu_count: Etc.nprocessors) ⇒ Object



17
18
19
20
21
22
# File 'lib/kettle/family/concurrency.rb', line 17

def wave_jobs(requested:, item_count:, cpu_count: Etc.nprocessors)
  return 0 if item_count.zero?

  count = requested.nil? ? default_wave_jobs(cpu_count: cpu_count) : requested.to_i
  count.clamp(1, item_count)
end