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
- .default_wave_jobs(cpu_count: Etc.nprocessors) ⇒ Object
- .internal_worker_limit(wave_jobs:, cpu_count: Etc.nprocessors) ⇒ Object
-
.test_process_ceiling(wave_jobs:, cpu_count: Etc.nprocessors) ⇒ Object
Test runners count their primary worker as part of the requested process pool.
- .wave_jobs(requested:, item_count:, cpu_count: Etc.nprocessors) ⇒ Object
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
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 |