#include #include #include #include #include #include #include using namespace kpn; using namespace std::chrono_literals; // ── basic execution ─────────────────────────────────────────────────────────── TEST_CASE("scheduler runs submitted tasks", "[scheduler]") { ThreadPool pool(2); pool.start(); std::atomic counter{0}; for (int i = 0; i < 100; ++i) pool.submit([&counter]{ counter.fetch_add(1, std::memory_order_relaxed); }); pool.drain(); REQUIRE(counter.load() == 100); pool.stop(); } TEST_CASE("scheduler single thread executes all tasks", "[scheduler]") { ThreadPool pool(1); pool.start(); std::atomic counter{0}; for (int i = 0; i < 50; ++i) pool.submit([&counter]{ counter.fetch_add(1, std::memory_order_relaxed); }); pool.drain(); REQUIRE(counter.load() == 50); pool.stop(); } // ── drain ───────────────────────────────────────────────────────────────────── TEST_CASE("drain returns immediately when pool is idle", "[scheduler]") { ThreadPool pool(2); pool.start(); pool.drain(); // nothing submitted — should return immediately pool.stop(); } TEST_CASE("drain waits for all tasks to complete", "[scheduler]") { ThreadPool pool(4); pool.start(); std::atomic counter{0}; constexpr int N = 200; for (int i = 0; i < N; ++i) { pool.submit([&counter]{ std::this_thread::sleep_for(1ms); counter.fetch_add(1, std::memory_order_relaxed); }); } pool.drain(); REQUIRE(counter.load() == N); pool.stop(); } TEST_CASE("drain is safe to call multiple times", "[scheduler]") { ThreadPool pool(2); pool.start(); std::atomic counter{0}; pool.submit([&counter]{ counter.fetch_add(1, std::memory_order_relaxed); }); pool.drain(); REQUIRE(counter.load() == 1); pool.submit([&counter]{ counter.fetch_add(1, std::memory_order_relaxed); }); pool.drain(); REQUIRE(counter.load() == 2); pool.stop(); } // ── priority ordering ───────────────────────────────────────────────────────── TEST_CASE("higher priority tasks run before lower priority on single thread", "[scheduler]") { // Single thread guarantees serial execution — we can observe order. ThreadPool pool(1); pool.start(); // Pause the worker so we can fill the queue before it drains. std::mutex gate; gate.lock(); pool.submit([&gate]{ std::lock_guard lg(gate); }); // blocks worker std::vector order; std::mutex order_mx; for (float p : {0.1f, 0.9f, 0.5f, 0.8f, 0.2f}) { pool.submit([p, &order, &order_mx]{ std::lock_guard lg(order_mx); order.push_back(p); }, p); } gate.unlock(); // release the blocking task pool.drain(); pool.stop(); // order should be descending by priority REQUIRE(order.size() == 5); for (std::size_t i = 1; i < order.size(); ++i) REQUIRE(order[i - 1] >= order[i]); } TEST_CASE("equal priority tasks execute in FIFO order on single thread", "[scheduler]") { ThreadPool pool(1); pool.start(); std::mutex gate; gate.lock(); pool.submit([&gate]{ std::lock_guard lg(gate); }); std::vector order; std::mutex order_mx; for (int i = 0; i < 5; ++i) { pool.submit([i, &order, &order_mx]{ std::lock_guard lg(order_mx); order.push_back(i); }, 0.5f); // all same priority } gate.unlock(); pool.drain(); pool.stop(); REQUIRE(order == std::vector{0, 1, 2, 3, 4}); } // ── total_ / active_ accounting ─────────────────────────────────────────────── TEST_CASE("snapshot queue depth and active counts are consistent", "[scheduler]") { ThreadPool pool(2); pool.start(); // While tasks are running, active should be > 0 and total >= active. std::atomic running{false}; std::mutex gate; gate.lock(); for (int i = 0; i < 4; ++i) { pool.submit([&gate, &running]{ running.store(true, std::memory_order_relaxed); std::lock_guard lg(gate); }); } // Spin until at least one task has started while (!running.load(std::memory_order_relaxed)) std::this_thread::yield(); auto snap = pool.snapshot("test"); REQUIRE(snap.active_count > 0); REQUIRE(snap.queue_depth + snap.active_count > 0); gate.unlock(); pool.drain(); auto snap2 = pool.snapshot("test"); REQUIRE(snap2.active_count == 0); REQUIRE(snap2.queue_depth == 0); pool.stop(); } TEST_CASE("submitted and completed counters are accurate", "[scheduler]") { ThreadPool pool(3); pool.start(); constexpr int N = 60; for (int i = 0; i < N; ++i) pool.submit([]{ std::this_thread::yield(); }); pool.drain(); auto snap = pool.snapshot("test"); REQUIRE(snap.tasks_submitted == static_cast(N)); REQUIRE(snap.tasks_completed == static_cast(N)); pool.stop(); } // ── work stealing ───────────────────────────────────────────────────────────── TEST_CASE("work stealing: all tasks complete with uneven initial distribution", "[scheduler]") { // 4-thread pool. Submit a burst to ensure some threads start empty and must steal. ThreadPool pool(4); pool.start(); std::atomic counter{0}; constexpr int N = 400; for (int i = 0; i < N; ++i) pool.submit([&counter]{ std::this_thread::sleep_for(100us); counter.fetch_add(1, std::memory_order_relaxed); }); pool.drain(); REQUIRE(counter.load() == N); pool.stop(); } TEST_CASE("work stealing: tasks complete with more threads than initial queue targets", "[scheduler]") { // With round-robin, some threads may get no tasks initially and must steal. constexpr std::size_t THREADS = 8; ThreadPool pool(THREADS); pool.start(); std::atomic counter{0}; // Submit fewer tasks than threads so most threads must steal for (int i = 0; i < 4; ++i) pool.submit([&counter]{ counter.fetch_add(1, std::memory_order_relaxed); }); pool.drain(); REQUIRE(counter.load() == 4); pool.stop(); } // Regression: submitting to a stopped pool must be a no-op, not a segfault. // // stop() ends with queues_.clear(), and submit() went straight to // queues_[target] with no check — so a submission arriving after stop indexed // an empty vector. // // This is not a contrived teardown ordering; it happens on a normal path. A // node's space callback fires from whichever thread drained the channel, and // that thread belongs to the *consumer*. Stop the producer first — which a // sources-first shutdown does by design — and the consumer keeps draining its // backlog, firing the producer's space callback into a pool that has already // been torn down. Before this fix the static-network shutdown case crashed // about 12 runs in 20. // // Checking stopped_ without the lock would not be enough: the window between // reading the flag and indexing the vector is exactly where clear() runs. TEST_CASE("submitting to a stopped pool is refused, not fatal", "[scheduler]") { ThreadPool pool(2); pool.start(); pool.stop(); std::atomic ran{0}; for (int i = 0; i < 10; ++i) pool.submit([&] { ran.fetch_add(1, std::memory_order_relaxed); }); CHECK(ran.load(std::memory_order_relaxed) == 0); CHECK(pool.rejected() == 10); } TEST_CASE("submitting while the pool stops does not crash", "[scheduler]") { // The racing form of the case above: a producer thread submitting // continuously while stop() runs underneath it. Nothing is asserted about // how many tasks run — the point is that every submission either enqueues // or is refused, and none touches a destroyed queue. for (int rep = 0; rep < 20; ++rep) { ThreadPool pool(4); pool.start(); std::atomic go{false}; std::atomic ran{0}; std::thread submitter([&] { while (!go.load(std::memory_order_acquire)) {} for (int i = 0; i < 2000; ++i) pool.submit([&] { ran.fetch_add(1, std::memory_order_relaxed); }); }); go.store(true, std::memory_order_release); std::this_thread::sleep_for(std::chrono::microseconds(200)); pool.stop(); submitter.join(); // Everything submitted was either executed or refused; nothing vanished // into a queue that no longer existed. 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); }