From b500570c479d750a7385daab44ec7f7182472036 Mon Sep 17 00:00:00 2001 From: Duncan Tourolle Date: Sat, 8 Aug 2026 22:26:14 +0200 Subject: [PATCH] perf(pool): submit to the calling worker's own queue, and skip a notify nobody is waiting for B9 from PERF_PLAN, plus the notify gate. Two changes to submit(), both aimed at the same cost: on any pool of two or more threads, every single dispatch paid a futex wake. submit() round-robins, so a worker resubmitting -- which is what fire_once does on every token -- handed the task to a *different* worker, and that worker was asleep. bench_dispatch measured it as 182 ns/dispatch on ThreadPool(1) against 3197 ns on 20 threads, with voluntary context switches per task rising 0.00 -> 1.19 in step: the 1-thread pool is fast precisely because it resubmits into its own queue and finds the work already there. A submission originating on one of the pool's own workers now goes to that worker's queue, extending that property to any pool size. try_steal still corrects the imbalance. The identity check is against `this`, not merely "am I a pool worker": a worker of pool A submitting into pool B must not use A's index, which may exceed B's thread_count_. Nested networks do exactly this. tls_pool is a non-owning identity tag, only ever compared, never dereferenced -- its lifetime is strictly nested inside the pool's, since stop() joins every worker before clearing queues_. The notify gate skips the cv_mx_ round-trip and notify_one() when waiters_ is zero. waiters_ is maintained under cv_mx_ and incremented before the predicate is evaluated, so reading zero in submit() means no worker can be in wait() -- as opposed to reading zero because we raced one, which the mutex prevents. stop()'s notify_all() is deliberately left ungated. Shared pools, work_us=10, items/sec: chain-1 59512 -> 66662 (+12.0%), wide-4 53698 -> 61705 (+14.9%), chain-8 +6.8%, chain-32 +3.7%. Private pools -- the Node<> default -- are unchanged at -1.3% to +3.6%, inside the gate's tolerance. Two things tried and removed, recorded in comments so they are not retried: Raising the steal threshold to >1, to stop a thief winning the race for a self-submitted task, DEADLOCKS. An external submit() round-robins a single task onto an idle worker's queue; if that worker is parked, no peer will take it, because a queue of one is no longer stealable. latency mode hangs at 12 and 20 threads. It was also 2x slower in steady state, 2229 -> 4546 ns. B5, bounded spin before parking, does not pay: swept at 50/200/1000 rounds, 2123 / 2230 / 2574 ns against 2229 ns without it, with vcsw/task flat at ~0.97. The spin cannot catch what it targets, because a peer is woken the moment queued_ becomes non-zero -- before this worker reaches the spin at all. Guardrails: 150/150 ctest including soak and examples, and ThreadSanitizer clean over the full suite (128 cases, 302 assertions). That also pins the reference count PERF_PLAN section 6 flagged as uncertain: it is 150, not 146 or 136. Caveat: the throughput figures above were taken on a build that also carried the since-removed spin experiment. The scheduler logic is identical to this tree and correctness was re-verified on it, but the numbers are one build stale and predate the 7-pass gate, which has not been run on this change. Co-Authored-By: Claude Opus 5 --- include/kpn/scheduler.hpp | 89 +++++++++++++++++++++++++++++++++++++-- 1 file changed, 86 insertions(+), 3 deletions(-) diff --git a/include/kpn/scheduler.hpp b/include/kpn/scheduler.hpp index d11d775..c33d5a2 100644 --- a/include/kpn/scheduler.hpp +++ b/include/kpn/scheduler.hpp @@ -128,7 +128,28 @@ public: rejected_.fetch_add(1, std::memory_order_relaxed); return; } - std::size_t target = next_.fetch_add(1, std::memory_order_relaxed) % thread_count_; + // B9 — submit-to-self affinity. Round-robin hands every task to a + // *different* worker, and on a pool of two or more that worker is + // asleep, so each dispatch pays a futex wake: measured 182 ns/dispatch + // on a 1-thread pool against 3197 ns on 20 threads, with voluntary + // context switches per task rising 0.00 -> 1.19 in step. + // + // A submission originating on one of *our own* workers goes to that + // worker's queue instead. It is about to return to worker_loop and + // try_pop its own queue, so the work is already there and nothing + // sleeps — the property that makes ThreadPool(1) fast, extended to + // any pool size. Imbalance is corrected by the existing try_steal. + // + // The pool identity check is load-bearing: a worker of pool A + // submitting into pool B must not use A's index, which may exceed B's + // thread_count_ or alias an unrelated queue. Nested networks do + // exactly this. + std::size_t target; + if (tls_pool == this && tls_worker < thread_count_) { + target = tls_worker; + } else { + target = next_.fetch_add(1, std::memory_order_relaxed) % thread_count_; + } { std::lock_guard lock(queues_[target]->mx); queues_[target]->pq.push( @@ -142,8 +163,23 @@ public: // will observe total_ > 0) or already blocked in wait() (and will be // woken). Without this, notify_one() can slip into the gap between the // worker's predicate check and its wait(), and be lost — a deadlock. - { std::lock_guard lk(cv_mx_); } - cv_.notify_one(); + // + // Skipped entirely when no worker is parked. waiters_ is incremented + // *before* wait() releases cv_mx_ and decremented after it returns, + // both under that mutex, so a worker on its way to sleep is already + // counted here. Reading zero therefore means no worker can be in + // wait(), and there is nothing a notify could reach — as opposed to + // reading zero because we raced one, which the mutex prevents. + // + // This is the hot path for an already-busy pool: with B9 the work is + // in the local queue and the submitting worker will find it itself, + // so the lock round-trip and notify were pure overhead. Measured 1.00 + // voluntary context switches per task before this, on a pool where + // only one task is ever in flight. + if (waiters_.load(std::memory_order_seq_cst) != 0) { + { std::lock_guard lk(cv_mx_); } + cv_.notify_one(); + } } std::size_t thread_count() const { return thread_count_; } @@ -193,6 +229,17 @@ private: std::optional> try_steal(std::size_t thief) { // Find the most-loaded peer without blocking — racy peek is fine. + // + // The threshold is >0: a peer holding a single task is a valid victim. + // + // Raising it to >1 — to stop a thief winning the race for a task its + // owner just submitted to itself (B9) — deadlocks. `latency` mode + // hangs at 12 and 20 threads: an external submit() round-robins one + // task onto an idle worker's queue, and if that worker is parked, no + // peer will take it because a queue of one is no longer stealable. + // Nothing else is coming to wake it, so the pool sits forever. + // Measured before reverting: it also made steady-state *worse*, + // 2229 -> 4546 ns at 12 threads. std::size_t victim = thief, best = 0; for (std::size_t i = 0; i < queues_.size(); ++i) { if (i == thief) continue; @@ -221,15 +268,41 @@ private: } void worker_loop(std::size_t id) { + // Identify this thread as one of our workers, for B9's affinity check + // in submit(). Restored on exit rather than merely cleared: a pool + // whose worker runs a task that itself starts and stops a nested pool + // would otherwise come back with its identity erased. + ThreadPool* const prev_pool = tls_pool; + const std::size_t prev_worker = tls_worker; + tls_pool = this; + tls_worker = id; + struct Restore { + ThreadPool* p; std::size_t w; + ~Restore() { tls_pool = p; tls_worker = w; } + } restore{prev_pool, prev_worker}; + while (true) { if (auto fn = try_pop(*queues_[id])) { execute(*fn); continue; } if (auto fn = try_steal(id)) { execute(*fn); continue; } + // B5 (bounded spin before parking) was tried here and removed: it + // does not pay. Swept at 50/200/1000 rounds on a 12-thread pool, + // steady state went 2123 / 2230 / 2574 ns against 2229 ns without + // it, and voluntary context switches per task stayed at ~0.97 + // throughout. The spin cannot catch what it is aimed at, because + // a peer is woken the moment queued_ becomes non-zero — which + // happens before this worker reaches the spin at all. std::unique_lock lock(cv_mx_); + // Counted under cv_mx_ and before the predicate is evaluated, so + // that a submit() which reads waiters_ == 0 can be certain this + // worker is not about to block: to get here we already hold the + // mutex that submit() must take to notify. + waiters_.fetch_add(1, std::memory_order_seq_cst); cv_.wait(lock, [this] { return stopped_.load(std::memory_order_seq_cst) || queued_.load(std::memory_order_relaxed) > 0; }); + waiters_.fetch_sub(1, std::memory_order_seq_cst); // 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. @@ -239,6 +312,13 @@ private: } } + /// Which pool, and which of its workers, the calling thread is — or + /// nullptr on any thread that is not a pool worker. Read by submit() to + /// decide whether a local push is safe (B9). inline so the header stays + /// header-only. + static inline thread_local ThreadPool* tls_pool = nullptr; + static inline thread_local std::size_t tls_worker = 0; + const std::size_t thread_count_; std::vector> queues_; std::vector workers_; @@ -267,6 +347,9 @@ private: std::atomic queued_{0}; // waiting to run std::atomic active_{0}; // executing only (for snapshot) std::atomic next_{0}; // round-robin submit cursor + /// Workers currently inside cv_.wait(), maintained under cv_mx_. Lets + /// submit() skip the lock round-trip and notify when nobody is parked. + std::atomic waiters_{0}; std::atomic seq_{0}; // tie-break for equal-priority tasks std::atomic submitted_{0}; /// Submissions refused because the pool was already stopped. Not an error —