fix: 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. 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.
This commit is contained in:
@@ -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<Task&>(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<bool> stopped_{true};
|
||||
std::atomic<size_t> total_{0}; // queued + executing
|
||||
std::atomic<size_t> 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<size_t> queued_{0}; // waiting to run
|
||||
std::atomic<size_t> active_{0}; // executing only (for snapshot)
|
||||
std::atomic<size_t> next_{0}; // round-robin submit cursor
|
||||
std::atomic<uint64_t> seq_{0}; // tie-break for equal-priority tasks
|
||||
|
||||
Reference in New Issue
Block a user