http: replace WorkQueue and single threads handling for ThreadPool#33689
Conversation
|
The following sections might be updated with supplementary metadata relevant to reviewers and maintainers. Code Coverage & BenchmarksFor details see: https://corecheck.dev/bitcoin/bitcoin/pulls/33689. ReviewsSee the guideline for information on the review process.
If your review is incorrectly listed, please copy-paste ConflictsReviewers, this pull request conflicts with the following ones:
If you consider this pull request important, please also help to review the conflicting pull requests. Ideally, start with the one that should be merged first. LLM Linter (✨ experimental)Possible typos and grammar issues:
2026-01-30 21:18:09 |
|
concept ACK :-) will be reviewing this |
|
Adding myself as i wrote the original shitty code. |
|
Can't we use an already existing open source library instead of reinventing the wheel? |
That's a good question. It's usually where we all start. That being said, for the changes introduced in this PR, can argue that we’re encapsulating, documenting, and unit + fuzz testing code that wasn’t covered before, while also improving separation of responsibilities. We’re not adding anything more complex or that behaves radically differently from what we currently have. |
|
Approach ACK |
bb71b0e to
195a962
Compare
ismaelsadeeq
left a comment
There was a problem hiding this comment.
Concept ACK, thanks for the seperation from the initial index PR.
pinheadmz
left a comment
There was a problem hiding this comment.
ACK 195a96258f970c384ce180d57e73616904ef5fa1
Built and tested on macos/arm64 and debian/x86. Reviewed each commit and left a few comments.
I also tested the branch against other popular software that consumes the HTTP interface by running their CI:
The http worker threads are clearly labeled and visible in htop!
Show Signature
-----BEGIN PGP SIGNED MESSAGE-----
Hash: SHA256
ACK 195a96258f970c384ce180d57e73616904ef5fa1
-----BEGIN PGP SIGNATURE-----
iQIzBAEBCAAdFiEE5hdzzW4BBA4vG9eM5+KYS2KJyToFAmkCMdwACgkQ5+KYS2KJ
yTriXA//RRnzezUHdmzRlKmoSDA+ZBHz0RY3z2LGk63izb/YdYaJY9JBZL2Y9BA8
K2nyexdSDC/DFOm4H56ddEe6ChlB7w+uZc92SgFSLSvavInpZ80KEJRk07vgoIL7
hwuyyevWyOOU32iz1NE3q316TMaJmzVsPhRGwbdmTXNwJLtUX6g4czfh28ajW1DC
Y9ULKwT36rFHRcKwC1YZYuBJUNBZWQgVBcydcmS1UEykY4mBnCW9knrATwn/29b7
2AYPV+yuaiy9OpDEOJ8iKtZOPGR36NrIUleMUqruq4Sy2/TnJtm3AKNK0336/Fxu
MqVyPKSyusg7kBA6f2h/2+NHpbyLoboYhjZew+HvED/aFfi+Jla+nxybkUYXfciL
pzbND22TTuRGB1fKU7AwPD0TO9JwOTU385iEdpoGq8rbT3EpgPr31N4TeDQDJx5t
jPzzWZYj43JMuIc3bm/K5S2HYSdFZUEDQC81kbND+jOLF6YkJwS9794anLO38tvi
fip/eLK8Nw4pmWnW63/9lc+Y/1gLpgLDMxxhA1NKJyytk7z2IRo7vJKck6TqAjcZ
nU8Wv9/ful5ndDJfLIKuYT9jqRk9ORohVwv+P+ppuO8jhFjhuswxPlFJKkMXTeZn
hGi8QCrAUuibvVuLfLKVExJqSmmeUAkTVKp5ZTWLwNB7IkZrSoo=
=hZYm
-----END PGP SIGNATURE-----
pinheadmz's public key is on openpgp.org
195a962 to
f5eca8d
Compare
|
Oh, awesome extensive test and review @pinheadmz! you rock! |
pinheadmz
left a comment
There was a problem hiding this comment.
ACK f5eca8d
minor changes since last ack
Show Signature
-----BEGIN PGP SIGNED MESSAGE-----
Hash: SHA256
ACK f5eca8d5252a68f0fbc8b5efec021c75744555b6
-----BEGIN PGP SIGNATURE-----
iQIzBAEBCAAdFiEE5hdzzW4BBA4vG9eM5+KYS2KJyToFAmkCWuYACgkQ5+KYS2KJ
yTpb9RAAgxNlDQlqkVfnIK2Xt8BQmRsiG+1KHQGJA2d3RHoxjFg2129fj47y5i+t
GlatlcsXEcibI2D0F22sKSzY1u0LWy4GYLI1d12JAIewA99n/lRr1ktDk3v64pp8
kldvYpd8UBs7DCHSJhPs4HbOPgwIILPASdSbQTb+6m48X5jg+cu+yMBuANVq3sbB
9I8rUlbUgRva5voy3EGRkXsGTuwcaCoLDNHnnjrqQkJMYTym47rMF/xTZJJ11isF
QrpzFu+P2tFy1oj0bGd4e5EhqNk++qfRCuw9sGLgLL6YEHzC/ihVH/L5NAxoeJX4
L0yX09BG5wDrbUQvZn4w1uaQz39Uor5mb8+tZyDUyRmO9MrG5AHomUtfdbiYLwSO
007bLD+WX4Hxode47xCdhUhqflXqbHD2mhkf5ESyUgGu3smSHSQQosuzIJiVmTn9
b2UJG6see/tckFvart6YLn1AIu9uCyUMmB+hZIcSasZv95oEHv1YB3E5dcKUmq0d
A1cUHIhNvU4i/hNy2xgijgSjDdpVQMFbLYUt9y4ZnzpJXFd7mnbpHVUVM1P+Onvp
ScAtgvosXfbd5PjnAQWsMWgSmVHUtSc4t9GPjGlRYaVpHFkEbqWpWeKrYqMkU3id
9YiFP9BueWNbNV7OBlB0LcfCpwO2WbvfJOW93/jxDovc0tgohbU=
=c8aD
-----END PGP SIGNATURE-----
pinheadmz's public key is on openpgp.org
sedited
left a comment
There was a problem hiding this comment.
ACK f5eca8d
I have fuzzed the thread pool for some time again, and did not find any new issues. Also the previously observed memory leak during fuzzing was not observed anymore. The changes to the http workers look sane to me.
Can you run the commits through clang-format-diff.py? There are a bunch of formatting issues in the first two commits.
f5eca8d to
435e8b5
Compare
Pushed. |
There was a problem hiding this comment.
I have started reviewing this PR but have only finished the first commit (e1eb4cd3a5eb192cd6d9ee5d255688c06ab2089a).
The depth of tests gives me hope that this transition will be smooth, the tests cover most of the functionality, even a few corner cases I haven't thought of!
It seemed, however, that I had more to say about that than anticipated. I didn't even get to reviewing the threadpool properly, just most of the tests.
I know receiving this amount of feedback can be daunting, but we both want this to succeed. I have spent a lot of time meticulously going through the details. I hope you will take it as I meant it: to make absolutely sure this won't cause any problems and that our test suite is rock solid.
I will continue the review, but wanted to make sure you're aware of the progress and thought we can synchronize more often this way (pun intended).
In general I think adding a threadpool can untangle complicated code as such, but we have to find a balance between IO and CPU bound tasks. I found that oversubscription is usually a smaller problem, so we can likely solve the IO/CPU contention problem by creating dedicated ThreadPools for each major task (http, script verification, input fetcher, compaction, etc). This way it won't be a general resource guardian (which can be a ThreadPool's main purpose), just a tool to avoid the overhead of starting/stopping threads.
As hinted in the original thread, I find the current structure to be harder to follow, I would appreciate if you would consider doing the extraction in smaller steps:
- first step would be to use the new
ThreadPool's structure, but extract the old logic into it; - add the tests (unit and fuzz, no need to split it) against the old impl in a separate commit;
- in a last step swap it out with the new logic, keeping the same structure and tests.
This way we could debug and understand the old code before we jump into the new one - and the tests would guide us in guaranteeing that the behavior stayed basically the same.
I understand that's extra work, but I'm of course willing to help in whatever you need to be able to do a less risky transition.
I glanced quickly to the actual ThreadPool implementation, looks okay, but I want to investigate as part of my next wave of reviews why we need so much locking here and why we're relying on old primitives here instead of the C++20 concurrency tools (such as latches and barriers) and whether we can create and destroy the pool via RAII instead of manual start/stops.
Note that I have done an IBD before and after this change to see if everything was still working and I didn't see any regression!
here are all the changes I did locally while reviewing the change
diff --git a/src/test/threadpool_tests.cpp b/src/test/threadpool_tests.cpp
index e8200533cd..052784db37 100644
--- a/src/test/threadpool_tests.cpp
+++ b/src/test/threadpool_tests.cpp
@@ -2,270 +2,208 @@
// Distributed under the MIT software license, see the accompanying
// file COPYING or http://www.opensource.org/licenses/mit-license.php.
+#include <common/system.h>
+#include <test/util/setup_common.h>
#include <util/string.h>
#include <util/threadpool.h>
+#include <util/time.h>
#include <boost/test/unit_test.hpp>
-BOOST_AUTO_TEST_SUITE(threadpool_tests)
+#include <latch>
+#include <semaphore>
-constexpr auto TIMEOUT_SECS = std::chrono::seconds(120);
+using namespace std::chrono;
+constexpr auto TIMEOUT = seconds(120);
-template <typename T>
-void WaitFor(const std::vector<std::future<T>>& futures, const std::string& context)
+void WaitFor(std::span<const std::future<void>> futures)
{
- for (size_t i = 0; i < futures.size(); ++i) {
- if (futures[i].wait_for(TIMEOUT_SECS) != std::future_status::ready) {
- throw std::runtime_error("Timeout waiting for: " + context + ", task index " + util::ToString(i));
- }
+ for (const auto& f : futures) {
+ BOOST_REQUIRE(f.wait_for(TIMEOUT) == std::future_status::ready);
}
}
// Block a number of worker threads by submitting tasks that wait on `blocker_future`.
// Returns the futures of the blocking tasks, ensuring all have started and are waiting.
-std::vector<std::future<void>> BlockWorkers(ThreadPool& threadPool, std::shared_future<void>& blocker_future, int num_of_threads_to_block, const std::string& context)
+std::vector<std::future<void>> BlockWorkers(ThreadPool& threadPool, std::counting_semaphore<>& release_sem, size_t num_of_threads_to_block)
{
- // Per-thread ready promises to ensure all workers are actually blocked
- std::vector<std::promise<void>> ready_promises(num_of_threads_to_block);
- std::vector<std::future<void>> ready_futures;
- ready_futures.reserve(num_of_threads_to_block);
- for (auto& p : ready_promises) ready_futures.emplace_back(p.get_future());
-
- // Fill all workers with blocking tasks
- std::vector<std::future<void>> blocking_tasks;
- for (int i = 0; i < num_of_threads_to_block; i++) {
- std::promise<void>& ready = ready_promises[i];
- blocking_tasks.emplace_back(threadPool.Submit([blocker_future, &ready]() {
- ready.set_value();
- blocker_future.wait();
- }));
- }
+ assert(threadPool.WorkersCount() >= num_of_threads_to_block);
+ std::latch ready{std::ptrdiff_t(num_of_threads_to_block)};
- // Wait until all threads are actually blocked
- WaitFor(ready_futures, context);
+ std::vector<std::future<void>> blocking_tasks(num_of_threads_to_block);
+ for (auto& f : blocking_tasks) f = threadPool.Submit([&] {
+ ready.count_down();
+ release_sem.acquire();
+ });
+
+ ready.wait();
return blocking_tasks;
}
-BOOST_AUTO_TEST_CASE(threadpool_basic)
+BOOST_FIXTURE_TEST_SUITE(threadpool_tests, BasicTestingSetup)
+
+const size_t NUM_WORKERS_DEFAULT{size_t(GetNumCores()) + 1}; // we need to make sure there's *some* contention
+
+BOOST_AUTO_TEST_CASE(submit_to_non_started_pool_throws)
{
- // Test Cases
- // 0) Submit task to a non-started pool.
- // 1) Submit tasks and verify completion.
- // 2) Maintain all threads busy except one.
- // 3) Wait for work to finish.
- // 4) Wait for result object.
- // 5) The task throws an exception, catch must be done in the consumer side.
- // 6) Busy workers, help them by processing tasks from outside.
- // 7) Recursive submission of tasks.
- // 8) Submit task when all threads are busy, stop pool and verify the task gets executed.
-
- const int NUM_WORKERS_DEFAULT = 3;
- const std::string POOL_NAME = "test";
-
- // Test case 0, submit task to a non-started pool
- {
- ThreadPool threadPool(POOL_NAME);
- bool err = false;
- try {
- threadPool.Submit([]() { return false; });
- } catch (const std::runtime_error&) {
- err = true;
- }
- BOOST_CHECK(err);
+ ThreadPool threadPool{"not_started"};
+ BOOST_CHECK_EXCEPTION(threadPool.Submit([] { return 0; }), std::runtime_error, HasReason{"No active workers"});
+}
+
+BOOST_AUTO_TEST_CASE(submit_and_verify_completion)
+{
+ ThreadPool threadPool{"completion"};
+ threadPool.Start(NUM_WORKERS_DEFAULT);
+
+ const auto num_tasks{1 + m_rng.randrange<size_t>(50)};
+ std::atomic_size_t counter{0};
+
+ std::vector<std::future<void>> futures(num_tasks);
+ for (size_t i{0}; i < num_tasks; ++i) {
+ futures[i] = threadPool.Submit([&counter, i] { counter.fetch_add(i, std::memory_order_relaxed); });
}
- // Test case 1, submit tasks and verify completion.
- {
- int num_tasks = 50;
+ WaitFor(futures);
+ BOOST_CHECK_EQUAL(counter.load(std::memory_order_relaxed), (num_tasks - 1) * num_tasks / 2);
+ BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), 0);
+}
- ThreadPool threadPool(POOL_NAME);
- threadPool.Start(NUM_WORKERS_DEFAULT);
- std::atomic<int> counter = 0;
+BOOST_AUTO_TEST_CASE(limited_free_workers_processes_all_task)
+{
+ ThreadPool threadPool{"block_counts"};
+ threadPool.Start(NUM_WORKERS_DEFAULT);
+ const auto num_tasks{1 + m_rng.randrange<size_t>(20)};
- // Store futures to ensure completion before checking counter.
- std::vector<std::future<void>> futures;
- futures.reserve(num_tasks);
+ for (size_t free{1}; free < NUM_WORKERS_DEFAULT; ++free) {
+ BOOST_TEST_MESSAGE("Testing with " << free << " available workers");
+ std::counting_semaphore sem{0};
+ const auto blocking_tasks{BlockWorkers(threadPool, sem, free)};
- for (int i = 1; i <= num_tasks; i++) {
- futures.emplace_back(threadPool.Submit([&counter, i]() {
- counter.fetch_add(i);
- }));
- }
+ size_t counter{0};
+ std::vector<std::future<void>> futures(num_tasks);
+ for (auto& f : futures) f = threadPool.Submit([&counter] { ++counter; });
- // Wait for all tasks to finish
- WaitFor(futures, /*context=*/"test1 task");
- int expected_value = (num_tasks * (num_tasks + 1)) / 2; // Gauss sum.
- BOOST_CHECK_EQUAL(counter.load(), expected_value);
+ WaitFor(futures);
BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), 0);
- }
- // Test case 2, maintain all threads busy except one.
- {
- ThreadPool threadPool(POOL_NAME);
- threadPool.Start(NUM_WORKERS_DEFAULT);
- // Single blocking future for all threads
- std::promise<void> blocker;
- std::shared_future<void> blocker_future(blocker.get_future());
- const auto& blocking_tasks = BlockWorkers(threadPool, blocker_future, NUM_WORKERS_DEFAULT - 1, /*context=*/"test2 blocking tasks enabled");
-
- // Now execute tasks on the single available worker
- // and check that all the tasks are executed.
- int num_tasks = 15;
- int counter = 0;
-
- // Store futures to wait on
- std::vector<std::future<void>> futures;
- futures.reserve(num_tasks);
- for (int i = 0; i < num_tasks; i++) {
- futures.emplace_back(threadPool.Submit([&counter]() {
- counter += 1;
- }));
+ if (free == 1) {
+ BOOST_CHECK_EQUAL(counter, num_tasks);
+ } else {
+ BOOST_CHECK_LE(counter, num_tasks); // unsynchronized update from multiple threads doesn't guarantee consistency
}
- WaitFor(futures, /*context=*/"test2 tasks");
- BOOST_CHECK_EQUAL(counter, num_tasks);
-
- blocker.set_value();
- WaitFor(blocking_tasks, /*context=*/"test2 blocking tasks disabled");
- threadPool.Stop();
- BOOST_CHECK_EQUAL(threadPool.WorkersCount(), 0);
+ sem.release(free);
+ WaitFor(blocking_tasks);
}
- // Test case 3, wait for work to finish.
- {
- ThreadPool threadPool(POOL_NAME);
- threadPool.Start(NUM_WORKERS_DEFAULT);
- std::atomic<bool> flag = false;
- std::future<void> future = threadPool.Submit([&flag]() {
- std::this_thread::sleep_for(std::chrono::milliseconds{200});
- flag.store(true);
- });
- future.wait();
- BOOST_CHECK(flag.load());
- }
+ threadPool.Stop();
+}
- // Test case 4, obtain result object.
- {
- ThreadPool threadPool(POOL_NAME);
- threadPool.Start(NUM_WORKERS_DEFAULT);
- std::future<bool> future_bool = threadPool.Submit([]() {
- return true;
- });
- BOOST_CHECK(future_bool.get());
+BOOST_AUTO_TEST_CASE(future_wait_blocks_until_task_completes)
+{
+ ThreadPool threadPool{"wait_test"};
+ threadPool.Start(NUM_WORKERS_DEFAULT);
+ const auto num_tasks{1 + m_rng.randrange<size_t>(20)};
- std::future<std::string> future_str = threadPool.Submit([]() {
- return std::string("true");
- });
- std::string result = future_str.get();
- BOOST_CHECK_EQUAL(result, "true");
+ const auto start{steady_clock::now()};
+
+ std::vector<std::future<void>> futures(num_tasks + 1);
+ for (size_t i{0}; i <= num_tasks; ++i) {
+ futures[i] = threadPool.Submit([i] { UninterruptibleSleep(milliseconds{i}); });
}
+ WaitFor(futures);
- // Test case 5, throw exception and catch it on the consumer side.
- {
- ThreadPool threadPool(POOL_NAME);
- threadPool.Start(NUM_WORKERS_DEFAULT);
-
- int ROUNDS = 5;
- std::string err_msg{"something wrong happened"};
- std::vector<std::future<void>> futures;
- futures.reserve(ROUNDS);
- for (int i = 0; i < ROUNDS; i++) {
- futures.emplace_back(threadPool.Submit([err_msg, i]() {
- throw std::runtime_error(err_msg + util::ToString(i));
- }));
- }
+ const size_t elapsed_ms{size_t(duration_cast<milliseconds>(steady_clock::now() - start).count())};
+ BOOST_CHECK(elapsed_ms >= num_tasks);
+}
- for (int i = 0; i < ROUNDS; i++) {
- try {
- futures.at(i).get();
- BOOST_FAIL("Expected exception not thrown");
- } catch (const std::runtime_error& e) {
- BOOST_CHECK_EQUAL(e.what(), err_msg + util::ToString(i));
- }
- }
- }
+BOOST_AUTO_TEST_CASE(future_get_returns_task_result)
+{
+ ThreadPool threadPool{"result_test"};
+ threadPool.Start(NUM_WORKERS_DEFAULT);
- // Test case 6, all workers are busy, help them by processing tasks from outside.
- {
- ThreadPool threadPool(POOL_NAME);
- threadPool.Start(NUM_WORKERS_DEFAULT);
-
- std::promise<void> blocker;
- std::shared_future<void> blocker_future(blocker.get_future());
- const auto& blocking_tasks = BlockWorkers(threadPool, blocker_future, NUM_WORKERS_DEFAULT, /*context=*/"test6 blocking tasks enabled");
-
- // Now submit tasks and check that none of them are executed.
- int num_tasks = 20;
- std::atomic<int> counter = 0;
- for (int i = 0; i < num_tasks; i++) {
- threadPool.Submit([&counter]() {
- counter.fetch_add(1);
- });
- }
- std::this_thread::sleep_for(std::chrono::milliseconds{100});
- BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), 20);
+ BOOST_CHECK_EQUAL(threadPool.Submit([] { return true; }).get(), true);
+ BOOST_CHECK_EQUAL(threadPool.Submit([] { return 42; }).get(), 42);
+ BOOST_CHECK_EQUAL(threadPool.Submit([] { return std::string{"true"}; }).get(), "true");
+}
- // Now process manually
- for (int i = 0; i < num_tasks; i++) {
- threadPool.ProcessTask();
- }
- BOOST_CHECK_EQUAL(counter.load(), num_tasks);
- BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), 0);
- blocker.set_value();
- threadPool.Stop();
- WaitFor(blocking_tasks, "Failure waiting for test6 blocking task futures");
+BOOST_AUTO_TEST_CASE(task_exception_propagated_to_future)
+{
+ ThreadPool threadPool{"exception_test"};
+ threadPool.Start(NUM_WORKERS_DEFAULT);
+
+ const auto err{[&](size_t n) { return strprintf("error on thread #%s", n); }};
+
+ const auto num_tasks{1 + m_rng.randrange<size_t>(20)};
+ for (size_t i{0}; i < num_tasks; ++i) {
+ BOOST_CHECK_EXCEPTION(threadPool.Submit([&] { throw std::runtime_error(err(i)); }).get(), std::runtime_error, HasReason{err(i)});
}
+}
+
+BOOST_AUTO_TEST_CASE(process_task_manually_when_workers_busy)
+{
+ ThreadPool threadPool{"manual_process"};
+ threadPool.Start(NUM_WORKERS_DEFAULT);
+ const auto num_tasks{1 + m_rng.randrange<size_t>(20)};
- // Test case 7, recursive submission of tasks.
- {
- ThreadPool threadPool(POOL_NAME);
- threadPool.Start(NUM_WORKERS_DEFAULT);
+ std::counting_semaphore sem{0};
+ const auto blocking_tasks{BlockWorkers(threadPool, sem, NUM_WORKERS_DEFAULT)};
- std::promise<void> signal;
- threadPool.Submit([&]() {
- threadPool.Submit([&]() {
- signal.set_value();
- });
- });
+ std::atomic_size_t counter{0};
+ std::vector<std::future<void>> futures(num_tasks);
+ for (auto& f : futures) f = threadPool.Submit([&counter] { counter.fetch_add(1, std::memory_order_relaxed); });
- signal.get_future().wait();
- threadPool.Stop();
- }
+ UninterruptibleSleep(milliseconds{100});
+ BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), num_tasks);
- // Test case 8, submit a task when all threads are busy and then stop the pool.
- {
- ThreadPool threadPool(POOL_NAME);
- threadPool.Start(NUM_WORKERS_DEFAULT);
+ for (size_t i{0}; i < num_tasks; ++i) {
+ threadPool.ProcessTask();
+ }
- std::promise<void> blocker;
- std::shared_future<void> blocker_future(blocker.get_future());
- const auto& blocking_tasks = BlockWorkers(threadPool, blocker_future, NUM_WORKERS_DEFAULT, /*context=*/"test8 blocking tasks enabled");
+ BOOST_CHECK_EQUAL(counter.load(), num_tasks);
+ BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), 0);
- // Submit an extra task that should execute once a worker is free
- std::future<bool> future = threadPool.Submit([]() { return true; });
+ sem.release(NUM_WORKERS_DEFAULT);
+ WaitFor(blocking_tasks);
+}
- // At this point, all workers are blocked, and the extra task is queued
- BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), 1);
+BOOST_AUTO_TEST_CASE(recursive_task_submission)
+{
+ ThreadPool threadPool{"recursive"};
+ threadPool.Start(NUM_WORKERS_DEFAULT);
- // Wait a short moment before unblocking the threads to mimic a concurrent shutdown
- std::thread thread_unblocker([&blocker]() {
- std::this_thread::sleep_for(std::chrono::milliseconds{300});
- blocker.set_value();
+ std::promise<void> signal;
+ threadPool.Submit([&threadPool, &signal] {
+ threadPool.Submit([&signal] {
+ signal.set_value();
});
+ });
- // Stop the pool while the workers are still blocked
- threadPool.Stop();
+ signal.get_future().wait();
+}
- // Expect the submitted task to complete
- BOOST_CHECK(future.get());
- thread_unblocker.join();
+BOOST_AUTO_TEST_CASE(stop_completes_queued_tasks_gracefully)
+{
+ ThreadPool threadPool{"graceful_stop"};
+ threadPool.Start(NUM_WORKERS_DEFAULT);
- // Obviously all the previously blocking tasks should be completed at this point too
- WaitFor(blocking_tasks, "Failure waiting for test8 blocking task futures");
+ std::counting_semaphore sem{0};
+ const auto blocking_tasks{BlockWorkers(threadPool, sem, NUM_WORKERS_DEFAULT)};
- // Pool should be stopped and no workers remaining
- BOOST_CHECK_EQUAL(threadPool.WorkersCount(), 0);
- }
+ auto future{threadPool.Submit([] { return true; })};
+ BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), 1);
+
+ std::thread thread_unblocker{[&sem] {
+ std::this_thread::sleep_for(milliseconds{300});
+ sem.release(NUM_WORKERS_DEFAULT);
+ }};
+
+ threadPool.Stop();
+
+ BOOST_CHECK(future.get());
+ thread_unblocker.join();
+ WaitFor(blocking_tasks);
+ BOOST_CHECK_EQUAL(threadPool.WorkersCount(), 0);
}
BOOST_AUTO_TEST_SUITE_END()
diff --git a/src/util/threadpool.h b/src/util/threadpool.h
index 5d9884086e..c89fda37c2 100644
--- a/src/util/threadpool.h
+++ b/src/util/threadpool.h
@@ -24,6 +24,8 @@
#include <utility>
#include <vector>
+#include <tinyformat.h>
+
/**
* @brief Fixed-size thread pool for running arbitrary tasks concurrently.
*
@@ -62,16 +64,9 @@ private:
for (;;) {
std::packaged_task<void()> task;
{
- // Wait only if needed; avoid sleeping when a new task was submitted while we were processing another one.
- if (!m_interrupt && m_work_queue.empty()) {
- // Block until the pool is interrupted or a task is available.
- m_cv.wait(wait_lock, [&]() EXCLUSIVE_LOCKS_REQUIRED(m_mutex) { return m_interrupt || !m_work_queue.empty(); });
- }
-
- // If stopped and no work left, exit worker
- if (m_interrupt && m_work_queue.empty()) {
- return;
- }
+ // Block until the pool is interrupted or a task is available.
+ m_cv.wait(wait_lock, [&]() EXCLUSIVE_LOCKS_REQUIRED(m_mutex) { return m_interrupt || !m_work_queue.empty(); });
+ if (m_interrupt && m_work_queue.empty()) return;
task = std::move(m_work_queue.front());
m_work_queue.pop();
@@ -101,17 +96,16 @@ public:
*
* Must be called from a controller (non-worker) thread.
*/
- void Start(int num_workers) EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
+ void Start(size_t num_workers) EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
{
assert(num_workers > 0);
LOCK(m_mutex);
if (!m_workers.empty()) throw std::runtime_error("Thread pool already started");
m_interrupt = false; // Reset
- // Create workers
m_workers.reserve(num_workers);
- for (int i = 0; i < num_workers; i++) {
- m_workers.emplace_back(&util::TraceThread, m_name + "_pool_" + util::ToString(i), [this] { WorkerThread(); });
+ for (size_t i{0}; i < num_workers; i++) {
+ m_workers.emplace_back(&util::TraceThread, strprintf("%s_pool_%d", m_name, i), [this] { WorkerThread(); });
}
}
@@ -179,12 +173,6 @@ public:
task();
}
- void Interrupt() EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
- {
- WITH_LOCK(m_mutex, m_interrupt = true);
- m_cv.notify_all();
- }
-
size_t WorkQueueSize() EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
{
return WITH_LOCK(m_mutex, return m_work_queue.size());|
|
||
| static void setup_threadpool_test() | ||
| { | ||
| // Disable logging entirely. It seems to cause memory leaks. |
There was a problem hiding this comment.
It is a known issue. The pcp fuzz test does it too for the same reason (in an undocumented manner). I only document it properly so we don't forget this exists. Feel free to investigate it in a separate issue/PR. I'm not planning to do it. It is not an issue of the code introduced in this PR.
There was a problem hiding this comment.
k, thanks, please resolve the comment
| // Create workers | ||
| m_workers.reserve(num_workers); | ||
| for (int i = 0; i < num_workers; i++) { | ||
| m_workers.emplace_back(&util::TraceThread, m_name + "_pool_" + util::ToString(i), [this] { WorkerThread(); }); |
There was a problem hiding this comment.
nit: we could use strprintf in a few more places:
| m_workers.emplace_back(&util::TraceThread, m_name + "_pool_" + util::ToString(i), [this] { WorkerThread(); }); | |
| m_workers.emplace_back(&util::TraceThread, strprintf("%s_pool_%d", m_name, i), [this] { WorkerThread(); }); |
There was a problem hiding this comment.
I'm not sure about including extra dependencies on a low-level class for tiny readability improvements. But I'm not really opposed to this one, so done as suggested.
| void Stop() EXCLUSIVE_LOCKS_REQUIRED(!m_mutex) | ||
| { | ||
| // Notify workers and join them. | ||
| std::vector<std::thread> threads_to_join; |
There was a problem hiding this comment.
I'm not sure this is necessarry
| ThreadPool threadPool(POOL_NAME); | ||
| threadPool.Start(NUM_WORKERS_DEFAULT); | ||
|
|
||
| int ROUNDS = 5; |
There was a problem hiding this comment.
in other cases we called this num_tasks
|
|
||
| int ROUNDS = 5; | ||
| std::string err_msg{"something wrong happened"}; | ||
| std::vector<std::future<void>> futures; |
There was a problem hiding this comment.
do we need a vector here of can we just consume whatever we just submitted?
BOOST_AUTO_TEST_CASE(task_exception_propagated_to_future)
{
ThreadPool threadPool{"exception_test"};
threadPool.Start(NUM_WORKERS_DEFAULT);
const auto make_err{[&](size_t n) { return strprintf("error on thread #%s", n); }};
const auto num_tasks{1 + m_rng.randrange<size_t>(20)};
for (size_t i{0}; i < num_tasks; ++i) {
BOOST_CHECK_EXCEPTION(threadPool.Submit([&] { throw std::runtime_error(make_err(i)); }).get(), std::runtime_error, HasReason{make_err(i)});
}
}There was a problem hiding this comment.
do we need a vector here of can we just consume whatever we just submitted?
Consuming right away would wait for the task to be executed; we want to exercise some concurrency too.
There was a problem hiding this comment.
we want to exercise some concurrency too
You're right, that's better indeed.
My suggestion updated to keep the vector
// Test 5, throw exceptions and catch it on the consumer side
BOOST_AUTO_TEST_CASE(task_exception_propagates_to_future)
{
ThreadPool threadPool("exception_test");
threadPool.Start(NUM_WORKERS_DEFAULT);
const auto make_err{[&](size_t n) { return strprintf("error on thread #%s", n); }};
const auto num_tasks{1 + m_rng.randrange<size_t>(20)};
std::vector<std::future<void>> futures(num_tasks);
for (size_t i{0}; i < num_tasks; ++i) {
futures[i] = threadPool.Submit([&make_err, i] { throw std::runtime_error(make_err(i)); });
}
for (size_t i{0}; i < num_tasks; ++i) {
BOOST_CHECK_EXCEPTION(futures[i].get(), std::runtime_error, HasReason{make_err(i)});
}
}| } | ||
| } | ||
|
|
||
| // Test case 6, all workers are busy, help them by processing tasks from outside. |
There was a problem hiding this comment.
We can split out the remaining 3 tests as well:
BOOST_AUTO_TEST_CASE(process_task_manually_when_workers_busy)
{
ThreadPool threadPool{"manual_process"};
threadPool.Start(NUM_WORKERS_DEFAULT);
const auto num_tasks{1 + m_rng.randrange<size_t>(20)};
std::counting_semaphore sem{0};
const auto blocking_tasks{BlockWorkers(threadPool, sem, NUM_WORKERS_DEFAULT)};
std::atomic_size_t counter{0};
std::vector<std::future<void>> futures(num_tasks);
for (auto& f : futures) f = threadPool.Submit([&counter] { counter.fetch_add(1, std::memory_order_relaxed); });
UninterruptibleSleep(milliseconds{100});
BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), num_tasks);
for (size_t i{0}; i < num_tasks; ++i) {
threadPool.ProcessTask();
}
BOOST_CHECK_EQUAL(counter.load(), num_tasks);
BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), 0);
sem.release(NUM_WORKERS_DEFAULT);
WaitFor(blocking_tasks);
}
BOOST_AUTO_TEST_CASE(recursive_task_submission)
{
ThreadPool threadPool{"recursive"};
threadPool.Start(NUM_WORKERS_DEFAULT);
std::promise<void> signal;
threadPool.Submit([&threadPool, &signal] {
threadPool.Submit([&signal] {
signal.set_value();
});
});
signal.get_future().wait();
}
BOOST_AUTO_TEST_CASE(stop_completes_queued_tasks_gracefully)
{
ThreadPool threadPool{"graceful_stop"};
threadPool.Start(NUM_WORKERS_DEFAULT);
std::counting_semaphore sem{0};
const auto blocking_tasks{BlockWorkers(threadPool, sem, NUM_WORKERS_DEFAULT)};
auto future{threadPool.Submit([] { return true; })};
BOOST_CHECK_EQUAL(threadPool.WorkQueueSize(), 1);
std::thread thread_unblocker{[&sem] {
std::this_thread::sleep_for(milliseconds{300});
sem.release(NUM_WORKERS_DEFAULT);
}};
threadPool.Stop();
BOOST_CHECK(future.get());
thread_unblocker.join();
WaitFor(blocking_tasks);
BOOST_CHECK_EQUAL(threadPool.WorkersCount(), 0);
}| @@ -49,83 +50,6 @@ using common::InvalidPortErrMsg; | |||
| /** Maximum size of http request (request line + headers) */ | |||
| static const size_t MAX_HEADERS_SIZE = 8192; | |||
|
|
|||
| /** HTTP request work item */ | |||
| class HTTPWorkItem final : public HTTPClosure | |||
There was a problem hiding this comment.
Is there a way to do this final threadpool migration in smaller steps?
There was a problem hiding this comment.
I guess you wrote this first and then answered yourself in your final comment; which is basically a long version of the PR's "Note 2" description.
There was a problem hiding this comment.
Not exactly sure what you mean by that - Note 2 does seem like a good idea to me and it's probably similar to what I meant. It would help the review process to gain more confidence in the implementation if we migrated away in smaller steps, documenting how the tests we add for the old implementation also pass when we switch over to the new implementation.
There was a problem hiding this comment.
I imagine the problem is that you haven't compared the current http code with the ThreadPool code, which is pretty much the same but properly documented and test covered. That's one of the reasons behind the not modernization of the underlying sync mechanism in this PR, and why I added the "Note 2" paragraph as well as wrote #33689 (comment).
Look at the current WorkQueue:
void Run() EXCLUSIVE_LOCKS_REQUIRED(!cs)
{
while (true) {
std::unique_ptr<WorkItem> i;
{
WAIT_LOCK(cs, lock);
while (running && queue.empty())
cond.wait(lock);
if (!running && queue.empty())
break;
i = std::move(queue.front());
queue.pop_front();
}
(*i)();
}
}And this is the ThreadPool (stripping all comments)
void WorkerThread() EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)
{
WAIT_LOCK(m_mutex, wait_lock);
for (;;) {
std::packaged_task<void()> task;
{
if (!m_interrupt && m_work_queue.empty()) {
m_cv.wait(wait_lock, [&]() EXCLUSIVE_LOCKS_REQUIRED(m_mutex) { return m_interrupt || !m_work_queue.empty(); });
}
if (m_interrupt && m_work_queue.empty()) return;
task = std::move(m_work_queue.front());
m_work_queue.pop();
}
REVERSE_LOCK(wait_lock, m_mutex);
task();
}
}In any case, since most reviewers seem ok to proceed with the current approach, I think it’s more productive to keep working on it rather than continue circling around this topic, even if it’s not to your taste.
There was a problem hiding this comment.
This is a big change, other reviewers can just review the unified view if they don't want smaller commits, but I currently cannot view this in small steps, so I insist that we should split it into smaller changes. I want to help, that's why I'm spending so much time with the details, I think this is a really risky change, I want to make sure it's correct. Let me know how I can help.
|
I've been playing around with the fuzz test and wanted to update here, even if it's still a bit inconclusive. I'm running into non-reproducible timeouts when using I have "fixed" this by changing edit: This should be reproducible in the sense that if you run with |
Eunovo
left a comment
There was a problem hiding this comment.
Concept ACK 435e8b5
The PR looks good already, but I think we can block users from calling Threadpool::Start() and Threadpool::Stop() inside Worker threads; We can use a thread local variable to identify worker threads and reject the operation.
| * Creates and launches `num_workers` threads that begin executing tasks | ||
| * from the queue. If the pool is already started, throws. | ||
| * | ||
| * Must be called from a controller (non-worker) thread. |
There was a problem hiding this comment.
Can we enforce this rule using a thread-local variable?
There was a problem hiding this comment.
i remember that in the past we had tons of issues with thread-local variables. They've been a pain in the ass on some platforms, and decided to not use them except for one case. Not sure if this is still the case, but if not, this needs to be updated:
https://github.com/bitcoin/bitcoin/blob/master/src/util/threadnames.cpp#L39
There was a problem hiding this comment.
There's also our thread_local clang-tidy plugin: https://github.com/bitcoin/bitcoin/tree/master/contrib/devtools/bitcoin-tidy.
There was a problem hiding this comment.
There's also our
thread_localclang-tidy plugin: https://github.com/bitcoin/bitcoin/tree/master/contrib/devtools/bitcoin-tidy.
I don't see why we would need a thread local variable with a non-trivial desctructor to implement this.
i remember that in the past we had tons of issues with thread-local variables. They've been a pain in the ass on some platforms, and decided to not use them except for one case. Not sure if this is still the case, but if not, this needs to be updated:
https://github.com/bitcoin/bitcoin/blob/master/src/util/threadnames.cpp#L39
I think it would be fine in this case because the use is limited. A simple thread-local boolean and assertion should be enough.
I don't have a strong opinion here; I just think it would be better to programmatically enforce the rule that certain functions should not be called from Worker Threads.
There was a problem hiding this comment.
Can we enforce this rule using a thread-local variable?
As m_workers is guarded, we could enforce it in a simple way:
for (const auto& worker : m_workers) assert(worker.get_id() != std::this_thread::get_id());Replace the HTTP server's WorkQueue implementation and single threads handling code with ThreadPool for processing HTTP requests. The ThreadPool class encapsulates all this functionality on a reusable class, properly unit and fuzz tested (the previous code was not unit nor fuzz tested at all). This cleanly separates responsibilities: The HTTP server now focuses solely on receiving and dispatching requests, while ThreadPool handles concurrency, queuing, and execution. It simplifies init, shutdown and requests tracking. This also allows us to experiment with further performance improvements at the task queuing and execution level, such as a lock-free structure, task prioritization or any other performance improvement in the future, without having to deal with HTTP code that lives on a different layer.
ad8fb89 to
38fd85c
Compare
|
Updated per feedback. Thanks Eunovo. |
pinheadmz
left a comment
There was a problem hiding this comment.
ACK 38fd85c
Re-reviewed all commits since a lot has changed since last review. Built and tested on arm64/macos and x86/debian.
Tried breaking it with stuff like -rpcthreads=9000 (expect init failure).
Also tried -rpcthreads=1 with RPC-pummeling scripts.
I roughly rebased #32061 on this branch and tested. Conflicts were extremely minimal.
I'm excited to use this feature to parallelize other tasks in bitcoin!
Several questions below, mostly just about style. No blockers there. RFM 🚀
/home/zip/bitcoin/build/bin/bitcoind -regtest -rpcthreads=1000
├─ b-scheduler
├─ b-http_pool_0
├─ b-http_pool_1
├─ b-http_pool_2
├─ b-http_pool_3
├─ b-http_pool_4
├─ b-http_pool_5
...
Show Signature
-----BEGIN PGP SIGNED MESSAGE-----
Hash: SHA256
ACK 38fd85c676a072ebf256e806beda9d7533790baa
-----BEGIN PGP SIGNATURE-----
iQIzBAEBCAAdFiEE5hdzzW4BBA4vG9eM5+KYS2KJyToFAmmLdj0ACgkQ5+KYS2KJ
yTpkNQ//Z8464/u6+Zw6tfIBEtcNCWLi2m3pk2xnJE8q2HIzk2GT6NUzu5f6yRsV
fkB1AJqxwO9aKf+tmF85c+f6PoKhjyBDeA1nFa39ms8Madkn32W3su6fPcfqjeCy
ROS3T18QueM7MkH4qezCWZ4TkYYeZu/yEMQdfZzhLJF14DX/8HdL+/fanqyB0CHc
V8mckYkX9vszKmQz4HZqSSxL9RUvv/ipbLKSU7BgO+EOWHkLAHNr5V27UwySM/9s
uDDDQGDCsnQqtkajRGE+gg230Xq70Fwp12+japQ5yvpn17HjkkDnCsvoLK01BHAO
WHASMt6m+RSZGlIwv7WXd0vM2gqIAcrJwkS3VzoDlkVI0fLdcsyacPn2zKL0cOBg
j7OjaW8KV+MjS1KvRuJR0CAE5akOPTgO9aCi5Aly+LZyyLPoDoQOZYTGDq4TJB3n
REgN6LuvjBoctWQGYM3HTTMDTchQqd+WWTaF7aKoCM8VTsffN8d00c/Y4ADIQ4Fd
xRv6jC0QBw434Nb023h15OHNYVgMNr/xqulfGizZUzXOcrcfEuBNYTyK7zBi4njm
gjEmI2rEoIrCzrgRcwinAF0XpVMca9YyGao4bBtpYY5lmc+LeU7lz2ac1wk3kLGf
SXfHL6jqetA1QWSVuEZ7ZuFJJSxRRqwjuNOl0YtoJ1kHTH8tiSc=
=dFzb
-----END PGP SIGNATURE-----
pinheadmz's public key is on openpgp.org
| req->WriteHeader("Connection", "close"); | ||
| // TODO: Implement specific error formatting for the REST and JSON-RPC servers responses. | ||
| req->WriteReply(HTTP_INTERNAL_SERVER_ERROR, err_msg); | ||
| return false; |
There was a problem hiding this comment.
nit: HTTPRequestHandler is defined as returning a bool but I don't think we ever use its return value, can probably be void. No change needed at this time, I just noticed this code either returns a hard coded bool, or the value from some arbitrary function, whose type is not so easy to see on the same page.
There was a problem hiding this comment.
Yeah, the first commit fixes the possible crash using the existing HTTPRequestHandler, which returns a bool. In the removal commit, when HTTPRequestHandler gets dropped, the bool return value is dropped as well.
| self.log.exception(f"Called Process failed with stdout='{e.stdout}'; stderr='{e.stderr}';") | ||
| self.success = TestStatus.FAILED | ||
| except JSONRPCException as e: | ||
| self.log.exception(f"Failure during setup: error={e.error}, http_status={e.http_status}") |
There was a problem hiding this comment.
This is good but could even be improved (doesn't have to be here). For example with the diff below and the suggested diff in the first commit message, you can actually get all the information logged out:
current:
test_framework.authproxy.JSONRPCException: non-JSON HTTP response with '500 Internal Server Error' from server (-342)
with diff below:
test_framework.authproxy.JSONRPCException: non-JSON HTTP response with '500 Internal Server Error' from server: error from json rpc handler (-342)
diff --git a/test/functional/test_framework/authproxy.py b/test/functional/test_framework/authproxy.py
index 9b2fc0f7f9..051423928d 100644
--- a/test/functional/test_framework/authproxy.py
+++ b/test/functional/test_framework/authproxy.py
@@ -195,7 +195,7 @@ class AuthServiceProxy():
content_type = http_response.getheader('Content-Type')
if content_type != 'application/json':
raise JSONRPCException(
- {'code': -342, 'message': 'non-JSON HTTP response with \'%i %s\' from server' % (http_response.status, http_response.reason)},
+ {'code': -342, 'message': 'non-JSON HTTP response with \'%i %s\' from server: %s' % (http_response.status, http_response.reason, http_response.read().decode())},
http_response.status)
data = http_response.read()There was a problem hiding this comment.
The commit message seems wrong, because Unexpected exception will also log all details that are available.
So this seems like the wrong fix. The correct fix would be to add any missing details to the exception string. Fixed in #34575
| assert(num_workers > 0); | ||
| LOCK(m_mutex); | ||
| if (!m_workers.empty()) throw std::runtime_error("Thread pool already started"); | ||
| m_interrupt = false; // Reset |
There was a problem hiding this comment.
What scenario would we stop and then re-start an instantiated worker pool?
There was a problem hiding this comment.
What scenario would we stop and then re-start an instantiated worker pool?
I was thinking about single-shot tasks that require a large number of threads, where keeping those threads alive for the entire lifetime of the software would be unnecessary overhead. E.g. RPC commands such as scanblocks, rescanblockchain, dumputxoset, etc.
Right now, the thread pool code is simple enough and not very configurable, but we will likely add more features and tuning options over time. Those settings will likely only be available during the node's initialization, so we need to construct the pool early with them (note: we really don't want to go back to the global ArgsManager dependency). And, at the same time, we don’t want to keep a large number of threads running when there's no work for them.
| template <class F> [[nodiscard]] EXCLUSIVE_LOCKS_REQUIRED(!m_mutex) | ||
| auto Submit(F&& fn) |
There was a problem hiding this comment.
Just curious about your style choice here, I'd've written like this
template <class F>
[[nodiscard]] auto Submit(F&& fn) EXCLUSIVE_LOCKS_REQUIRED(!m_mutex)There was a problem hiding this comment.
I'm pretty sure I had it like that, and there was a suggestion to change it to the current form. Since it was merely an aesthetic change and there were so many comments, I didn't put much thought into it.
| BOOST_CHECK_EXCEPTION((void)threadPool.Submit([]{ return false; }), std::runtime_error, [&](const std::runtime_error& e) { | ||
| BOOST_CHECK_EQUAL(e.what(), "No active workers; cannot accept new tasks"); |
There was a problem hiding this comment.
Why not use HasReason()?
#include <test/util/setup_common.h>
...
BOOST_CHECK_EXCEPTION(
(void)threadPool.Submit([]{ return false; }),
std::runtime_error,
HasReason{"No active workers; cannot accept new tasks"
});Drahtbot should be complaining about this ;-)
see #34242 (comment)
and in task_exception_propagates_to_future and interrupt_blocks_new_submissions below
There was a problem hiding this comment.
I think I answered this in a previous comment.
I'm happy using HasReason, but not adding a dependency to the entire unit test framework machinery (<test/util/setup_common.h>). The rationale is that the thread pool lives at a lower level and shouldn't get access to the chainstate, feerate, or any other higher-level concepts. It is also easier for anyone who wants to experiment with the code to be able to pull the thread pool alone into a separate project. This last sentence comes from experience.. it would have been so nice to have constructed the script interpreter and other primitives in this way, not so tied to the whole project machinery.. we would have found many issues and improvements way faster.
So, in short, I agree with using it. I just think we should first move the HasReason class to a general utility file, to beak the general framework dependency. It could be done in a quick follow-up, happy to do it or ACK it if someone wants to tackle it too.
| LogWarning("Request rejected because http work queue depth exceeded, it can be increased with the -rpcworkqueue= setting"); | ||
| item->req->WriteReply(HTTP_SERVICE_UNAVAILABLE, "Work queue depth exceeded"); | ||
| } | ||
| [[maybe_unused]] auto _{g_threadpool_http.Submit(std::move(item))}; |
There was a problem hiding this comment.
Related to my comment on the first commit about how these workers dont return any thing. So now in the above try/catch wrapper the return values are removed which makes sense to me.
My guess is that for this use case of threadpool, we don't need to examine the return value of the future from the task, but other uses of threadpool will not discard it?
There was a problem hiding this comment.
My guess is that for this use case of threadpool, we don't need to examine the return value of the future from the task, but other uses of threadpool will not discard it?
yeah, any parallelizable user single-shot locking command would wait for it to finish. The examples I have in my head are RPC-related like scanblocks, rescanblockchain, dumputxoset, etc. But we could surely find more just by looking at the REST server or the GUI.
|
This seems to have introduced a bug: #34573. |
|
post merge ACK 38fd85c 🧑🏭 |
|
while reviewing #35182 I looked at the existing queue-overload path that was introduced here #33689 ( is it intentional that we do not set coz on #35182, HTTP/1.1 replies default to keep-alive in WriteReply(), so overload 503 responses may leave reusable keep-alive connections open until i used regtest to test results below shows no connection close header HTTP/1.1 503 Service Unavailable
Date: Sun, 24 May 2026 20:51:35 GMT
Content-Length: 25
Content-Type: text/html; charset=ISO-8859-1 |
it was not introduced here. This PR only made the behavior explicit. See
Why? The server being busy doesn't seem to be a strong reason for closing the connection on the server side. Forcing a close just means the client has to redo the handshake on the next request, which doesn't help either side. |
This has been a recent discovery; the general thread pool class created for #26966, cleanly
integrates into the HTTP server. It simplifies init, shutdown and requests execution logic.
Replacing code that was never unit tested for code that is properly unit and fuzz tested.
Although our functional test framework extensively uses this RPC interface (that’s how
we’ve been ensuring its correct behavior so far - which is not the best).
This clearly separates the responsibilities:
The HTTP server now focuses solely on receiving and dispatching requests, while ThreadPool handles
concurrency, queuing, and execution.
This will also allows us to experiment with further performance improvements at the task queuing and
execution level, such as a lock-free structure or task prioritization or any other implementation detail
like coroutines in the future, without having to deal with HTTP code that lives on a different layer.
Note:
The rationale behind introducing the ThreadPool first is to be able to easily cherry-pick it across different
working paths. Some of the ones that are benefited from it are #26966 for the parallelization of the indexes
initial sync, #31132 for the parallelization of the inputs fetching procedure, #32061 for the libevent replacement,
the kernel API #30595 (#30595 (comment)) to avoid blocking validation among others use cases not publicly available.
Note 2:
I could have created a wrapper around the existing code and replaced the
WorkQueuein a subsequentcommit, but it didn’t seem worth the extra commits and review effort. The
ThreadPoolimplementsessentially the same functionality in a more modern and cleaner way.