From b9698fae60dc246dea67af28d0570df2416e582c Mon Sep 17 00:00:00 2001 From: Duncan Tourolle Date: Wed, 5 Aug 2026 15:34:33 +0200 Subject: [PATCH] fix: idle workers must sleep while another worker is busy MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The wait predicate was `stopped_ || total_ > 0`, and total_ counts queued *plus executing*. So while any one task ran, every other worker's predicate was true: wait() returned immediately and the worker spun through try_pop, try_steal and back to wait at full speed, try_lock-ing every peer queue on each pass. Measured on this tree, 8 workers and one 300 ms task: 1991 ms of CPU and 19205 voluntary context switches, against 0.4 ms and 10 with the fix. Call it six and a half cores burned for the duration of one sleeping task. Sleeping requires "no work is *waiting*", which total_ cannot express, so queued_ is now tracked separately: incremented on submit, decremented when a task leaves a queue, and adjusted for the tasks stop() discards. total_ stays as it was for drain(), which genuinely does need to know about executing work. The worker exit condition moves to queued_ for the same reason — waiting for total_ to reach zero meant waiting for someone else's task to finish, which a worker cannot help with and would spin through until it did. PoolSnapshot's queue depth stops being an estimate as a side effect. Latent for this pipeline, where each node owns a private single-thread pool and there is no idle peer to spin. Any use of a shared pool, which make_pool_node exists for, hits it immediately. Reproducing it needs the right trigger, and the test says so, because my first attempt got it wrong and passed against the bug: a worker that has never been woken stays blocked in wait() and never re-evaluates the predicate. The spin only appears once a worker *finishes* something and re-enters the loop while a peer is still busy, so the case submits one long task plus a trivial one per remaining worker. Submitting only the long task measures nothing. Verified in both directions: 1969 ms of CPU before, 0.35 ms after, against a 300 ms threshold. Full suite 142/142. --- include/kpn/scheduler.hpp | 27 ++++++++++++++----- tests/test_scheduler.cpp | 56 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 77 insertions(+), 6 deletions(-) diff --git a/include/kpn/scheduler.hpp b/include/kpn/scheduler.hpp index 61c57bd..54bbb5b 100644 --- a/include/kpn/scheduler.hpp +++ b/include/kpn/scheduler.hpp @@ -74,6 +74,7 @@ public: std::size_t discarded = q->pq.size(); while (!q->pq.empty()) q->pq.pop(); total_.fetch_sub(discarded, std::memory_order_relaxed); + queued_.fetch_sub(discarded, std::memory_order_relaxed); } // Lock cv_mx_ before notifying so the stop signal can't be lost in the // gap between a worker's predicate check and its wait() (see submit()). @@ -121,6 +122,7 @@ public: {std::move(task), priority, seq_.fetch_add(1, std::memory_order_relaxed)}); } total_.fetch_add(1, std::memory_order_relaxed); + queued_.fetch_add(1, std::memory_order_relaxed); submitted_.fetch_add(1, std::memory_order_relaxed); // Synchronize with worker_loop's predicate evaluation: taking cv_mx_ // here guarantees a worker is either before its predicate check (and @@ -140,11 +142,10 @@ public: PoolSnapshot snapshot(const std::string& name) const override { std::size_t a = active_.load(std::memory_order_relaxed); - std::size_t t = total_.load(std::memory_order_relaxed); return { name, thread_count_, - t > a ? t - a : 0, // queued (approximate) - a, // executing + queued_.load(std::memory_order_relaxed), // queued (exact) + a, // executing submitted_.load(std::memory_order_relaxed), completed_.load(std::memory_order_relaxed), }; @@ -173,6 +174,7 @@ private: if (q.pq.empty()) return std::nullopt; auto fn = std::move(const_cast(q.pq.top()).fn); q.pq.pop(); + queued_.fetch_sub(1, std::memory_order_relaxed); return fn; } @@ -213,10 +215,13 @@ private: std::unique_lock lock(cv_mx_); cv_.wait(lock, [this] { return stopped_.load(std::memory_order_seq_cst) - || total_.load(std::memory_order_relaxed) > 0; + || queued_.load(std::memory_order_relaxed) > 0; }); + // Exit on queued_, not total_: waiting for total_ to reach zero + // meant waiting for someone else's task to finish, which this + // worker cannot help with and would spin through until it did. if (stopped_.load(std::memory_order_seq_cst) - && total_.load(std::memory_order_relaxed) == 0) + && queued_.load(std::memory_order_relaxed) == 0) return; } } @@ -236,7 +241,17 @@ private: std::condition_variable drain_cv_; std::atomic stopped_{true}; - std::atomic total_{0}; // queued + executing + std::atomic total_{0}; // queued + executing (drain() waits on this) + /// Queued only — never counts a task that is already executing. + /// + /// The wait predicate used total_, which includes running tasks, so while + /// any one task ran every *other* worker's predicate was true: wait() + /// returned instantly and the worker spun through try_pop / try_steal / + /// wait at full speed, try_lock-ing every peer queue on each pass. One slow + /// task therefore pinned every other core and contended the very mutexes + /// the working thread needed. Sleeping requires "no work is *waiting*", + /// which is this. + std::atomic queued_{0}; // waiting to run std::atomic active_{0}; // executing only (for snapshot) std::atomic next_{0}; // round-robin submit cursor std::atomic seq_{0}; // tie-break for equal-priority tasks diff --git a/tests/test_scheduler.cpp b/tests/test_scheduler.cpp index d011129..9af89b9 100644 --- a/tests/test_scheduler.cpp +++ b/tests/test_scheduler.cpp @@ -284,3 +284,59 @@ TEST_CASE("submitting while the pool stops does not crash", "[scheduler]") { CHECK(pool.rejected() + pool.snapshot("p").tasks_completed <= 2000); } } + +// Regression: idle workers must sleep while another worker is busy. +// +// The wait predicate was `stopped_ || total_ > 0`, and total_ counts queued +// *plus executing*. So while any one task ran, every other worker's predicate +// was true: wait() returned immediately and the worker spun through try_pop, +// try_steal and back to wait at full speed — try_lock-ing every peer queue on +// each pass. One slow task pinned every other core and contended the very +// mutexes the working thread needed to make progress. +// +// That is the shape of this pipeline's load exactly: a handful of nodes whose +// work is tens of milliseconds of ONNX inference. It was latent only because +// each node currently owns a private single-thread pool, where there is no +// idle peer to spin. Any use of a shared pool — which make_pool_node exists +// for — hits it immediately. +// +// Measured as CPU time rather than wall time, because the bug does not make +// anything slower to finish; it makes seven cores burn while one works. A +// sleeping task consumes no CPU, so with workers correctly asleep the whole +// pool should account for almost none. +TEST_CASE("idle workers do not spin while one task runs", "[scheduler]") { + auto cpu_ms = [] { + struct timespec ts{}; + clock_gettime(CLOCK_PROCESS_CPUTIME_ID, &ts); + return ts.tv_sec * 1000.0 + ts.tv_nsec / 1e6; + }; + + constexpr int kThreads = 8; + constexpr int kWorkMs = 300; + + ThreadPool pool(kThreads); + pool.start(); + std::this_thread::sleep_for(20ms); // let workers reach the wait + + const double before = cpu_ms(); + + // One long task plus a trivial one per remaining worker. The trivial ones + // matter: a worker that has never been woken stays blocked in wait() and + // never re-evaluates the predicate, so the spin only appears once a worker + // *finishes* something and re-enters the loop while a peer is still busy. + // Submitting only the long task does not reproduce it. + pool.submit([&] { std::this_thread::sleep_for(std::chrono::milliseconds(kWorkMs)); }); + for (int i = 0; i < kThreads - 1; ++i) pool.submit([] {}); + + pool.drain(); + const double used = cpu_ms() - before; + pool.stop(); + + // Measured on this tree: 1991 ms of CPU with the total_ predicate against + // 0.4 ms with queued_, and 19205 voluntary context switches against 10 — + // roughly (kThreads - 1) cores burned for the duration of one sleeping + // task. The threshold sits far from both so the case is not sensitive to + // how loaded the machine is. + INFO("cpu " << used << " ms over " << kWorkMs << " ms of sleeping work"); + CHECK(used < kWorkMs); +}