mirror of
https://github.com/verilator/verilator.git
synced 2026-09-04 00:30:34 +02:00
Thread pool rewrite (#5161)
Signed-off-by: Krzysztof Bieganski <[email protected]> Signed-off-by: Bartłomiej Chmiel <[email protected]> Signed-off-by: Arkadiusz Kozdra <[email protected]> Co-authored-by: Krzysztof Bieganski <[email protected]> Co-authored-by: Arkadiusz Kozdra <[email protected]> Co-authored-by: Wilson Snyder <[email protected]>
This commit is contained in:
co-authored by
Krzysztof Bieganski
Arkadiusz Kozdra
Wilson Snyder
parent
eb8bbcda05
commit
ffe76717c6
+97
-160
@@ -14,139 +14,78 @@
|
||||
//
|
||||
//*************************************************************************
|
||||
|
||||
#include "config_build.h"
|
||||
|
||||
#include "V3ThreadPool.h"
|
||||
|
||||
#include "V3Error.h"
|
||||
#include "V3Global.h"
|
||||
#include "V3Mutex.h"
|
||||
|
||||
// c++11 requires definition of static constexpr as well as declaration
|
||||
constexpr unsigned int V3ThreadPool::FUTUREWAITFOR_MS;
|
||||
V3ThreadPool::V3ThreadPool(int numThreads) {
|
||||
numThreads = std::max(numThreads, 1);
|
||||
if (numThreads == 1) return;
|
||||
for (int i = 0; i < numThreads; ++i) {
|
||||
m_workers.emplace_back(&V3ThreadPool::startWorker, this);
|
||||
}
|
||||
}
|
||||
|
||||
void V3ThreadPool::resize(unsigned n) VL_MT_UNSAFE VL_EXCLUDES(m_mutex)
|
||||
VL_EXCLUDES(m_stoppedJobsMutex) VL_EXCLUDES(V3MtDisabledLock::instance()) {
|
||||
// At least one thread (main)
|
||||
n = std::max(1u, n);
|
||||
if (n == (m_workers.size() + 1)) return;
|
||||
// This function is not thread-safe and can result in race between threads
|
||||
UASSERT(V3MutexConfig::s().lockConfig(),
|
||||
"Mutex config needs to be locked before starting ThreadPool");
|
||||
V3ThreadPool::~V3ThreadPool() {
|
||||
{
|
||||
V3LockGuard stoppedJobsLock{m_stoppedJobsMutex};
|
||||
V3LockGuard lock{m_mutex};
|
||||
|
||||
UASSERT(m_queue.empty(), "Resizing busy thread pool");
|
||||
// Shut down old threads
|
||||
const V3LockGuard lock{m_mutex};
|
||||
m_shutdown = true;
|
||||
m_stoppedJobs = 0;
|
||||
m_cv.notify_all();
|
||||
m_stoppedJobsCV.notify_all();
|
||||
m_exclusiveAccessThreadCV.notify_all();
|
||||
}
|
||||
while (!m_workers.empty()) {
|
||||
m_workers.front().join();
|
||||
m_workers.pop_front();
|
||||
}
|
||||
if (n > 1) {
|
||||
V3LockGuard lock{m_mutex};
|
||||
// Start new threads
|
||||
m_shutdown = false;
|
||||
for (unsigned int i = 1; i < n; ++i) {
|
||||
m_workers.emplace_back(&V3ThreadPool::startWorker, this, i);
|
||||
}
|
||||
}
|
||||
m_cv.notify_all();
|
||||
wait();
|
||||
}
|
||||
|
||||
void V3ThreadPool::suspendMultithreading() VL_MT_SAFE VL_EXCLUDES(m_mutex)
|
||||
VL_EXCLUDES(m_stoppedJobsMutex) {
|
||||
V3LockGuard stoppedJobsLock{m_stoppedJobsMutex};
|
||||
if (!m_workers.empty()) stopOtherThreads();
|
||||
|
||||
if (!m_mutex.try_lock()) {
|
||||
v3fatal("Tried to suspend thread pool when other thread uses it.");
|
||||
}
|
||||
V3LockGuard lock{m_mutex, std::adopt_lock_t{}};
|
||||
|
||||
UASSERT(m_queue.empty(), "Thread pool has pending jobs");
|
||||
UASSERT(m_jobsInProgress == 0, "Thread pool has jobs in progress");
|
||||
m_exclusiveAccess = true;
|
||||
m_multithreadingSuspended = true;
|
||||
}
|
||||
|
||||
void V3ThreadPool::resumeMultithreading() VL_MT_SAFE VL_EXCLUDES(m_mutex)
|
||||
VL_EXCLUDES(m_stoppedJobsMutex) {
|
||||
if (!m_mutex.try_lock()) v3fatal("Tried to resume thread pool when other thread uses it.");
|
||||
{
|
||||
V3LockGuard lock{m_mutex, std::adopt_lock_t{}};
|
||||
UASSERT(m_multithreadingSuspended, "Multithreading is not suspended");
|
||||
m_multithreadingSuspended = false;
|
||||
m_exclusiveAccess = false;
|
||||
}
|
||||
if (!m_workers.empty()) {
|
||||
V3LockGuard stoppedJobsLock{m_stoppedJobsMutex};
|
||||
resumeOtherThreads();
|
||||
}
|
||||
}
|
||||
|
||||
void V3ThreadPool::startWorker(V3ThreadPool* selfThreadp, int id) VL_MT_SAFE {
|
||||
selfThreadp->workerJobLoop(id);
|
||||
}
|
||||
|
||||
void V3ThreadPool::workerJobLoop(int id) VL_MT_SAFE {
|
||||
while (true) {
|
||||
// Wait for a notification
|
||||
waitIfStopRequested();
|
||||
VAnyPackagedTask job;
|
||||
void V3ThreadPool::enqueue(std::function<void()>&& f) {
|
||||
if (m_workers.empty()) {
|
||||
f();
|
||||
} else {
|
||||
{
|
||||
V3LockGuard lock(m_mutex);
|
||||
m_cv.wait(m_mutex, [&]() VL_REQUIRES(m_mutex) {
|
||||
return !m_queue.empty() || m_shutdown || m_stopRequested;
|
||||
});
|
||||
if (m_shutdown) return; // Terminate if requested
|
||||
if (stopRequested()) continue;
|
||||
// Get the job
|
||||
UASSERT(!m_queue.empty(), "Job should be available");
|
||||
const V3LockGuard lock{m_mutex};
|
||||
m_queue.push(std::move(f));
|
||||
}
|
||||
m_pendingJobs.fetch_add(1, std::memory_order_release);
|
||||
m_cv.notify_one();
|
||||
}
|
||||
}
|
||||
|
||||
void V3ThreadPool::wait() {
|
||||
while (m_pendingJobs.load(std::memory_order_acquire) > 0 && !m_shutdown) {
|
||||
std::this_thread::yield();
|
||||
}
|
||||
if (m_shutdown) {
|
||||
for (auto& worker : m_workers) worker.join();
|
||||
}
|
||||
}
|
||||
|
||||
void V3ThreadPool::startWorker(V3ThreadPool* selfThreadp) { selfThreadp->workerJobLoop(); }
|
||||
|
||||
void V3ThreadPool::workerJobLoop() {
|
||||
while (true) {
|
||||
std::function<void()> job;
|
||||
{
|
||||
// Locking before `condition_variable::wait` is required to ensure that the
|
||||
// `m_cv` condition will be executed under a lock. Taking a lock
|
||||
// before `condition_variable::wait` may lead to missed
|
||||
// `condition_variable::notify_all` notification (when a thread waits at the lock
|
||||
// before) but, according to C++ standard, the `condition_variable::wait` first checks
|
||||
// the condition and then waits for the notification, thus even when notification is
|
||||
// missed, the condition still will be checked.
|
||||
const V3LockGuard lock{m_mutex};
|
||||
m_cv.wait(m_mutex,
|
||||
[&]() VL_REQUIRES(m_mutex) { return !m_queue.empty() || m_shutdown; });
|
||||
if (m_shutdown) return; // Terminate if requested
|
||||
UASSERT(!m_queue.empty(), "Job should be available");
|
||||
if (m_queue.empty()) continue;
|
||||
job = std::move(m_queue.front());
|
||||
m_queue.pop();
|
||||
++m_jobsInProgress;
|
||||
}
|
||||
|
||||
// Execute the job
|
||||
job();
|
||||
// Note that a context switch can happen here. This means `m_jobsInProgress` could still
|
||||
// contain old value even after the job promise has been fulfilled.
|
||||
--m_jobsInProgress;
|
||||
m_pendingJobs.fetch_sub(1, std::memory_order_release);
|
||||
}
|
||||
}
|
||||
|
||||
bool V3ThreadPool::waitIfStopRequested() VL_MT_SAFE VL_EXCLUDES(m_stoppedJobsMutex) {
|
||||
if (!stopRequested()) return false;
|
||||
V3LockGuard stoppedJobLock(m_stoppedJobsMutex);
|
||||
waitForResumeRequest();
|
||||
return true;
|
||||
}
|
||||
|
||||
void V3ThreadPool::waitForResumeRequest() VL_REQUIRES(m_stoppedJobsMutex) {
|
||||
++m_stoppedJobs;
|
||||
m_exclusiveAccessThreadCV.notify_one();
|
||||
m_stoppedJobsCV.wait(m_stoppedJobsMutex,
|
||||
[&]() VL_REQUIRES(m_stoppedJobsMutex) { return !m_stopRequested; });
|
||||
--m_stoppedJobs;
|
||||
}
|
||||
|
||||
void V3ThreadPool::stopOtherThreads() VL_MT_SAFE_EXCLUDES(m_mutex)
|
||||
VL_REQUIRES(m_stoppedJobsMutex) {
|
||||
m_stopRequested = true;
|
||||
{
|
||||
V3LockGuard lock{m_mutex};
|
||||
m_cv.notify_all();
|
||||
}
|
||||
m_exclusiveAccessThreadCV.wait(m_stoppedJobsMutex, [&]() VL_REQUIRES(m_stoppedJobsMutex) {
|
||||
return m_stoppedJobs == m_workers.size();
|
||||
});
|
||||
}
|
||||
|
||||
void V3ThreadPool::selfTestMtDisabled() {
|
||||
// empty
|
||||
}
|
||||
@@ -157,7 +96,7 @@ void V3ThreadPool::selfTest() {
|
||||
|
||||
auto firstJob = [&](int sleep) -> void {
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds{sleep});
|
||||
const V3ThreadPool::ScopedExclusiveAccess exclusiveAccess;
|
||||
const V3LockGuard lock{commonMutex};
|
||||
commonValue = 10;
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds{sleep + 10});
|
||||
UASSERT(commonValue == 10, "unexpected commonValue = " << commonValue);
|
||||
@@ -165,7 +104,6 @@ void V3ThreadPool::selfTest() {
|
||||
auto secondJob = [&](int sleep) -> void {
|
||||
commonMutex.lock();
|
||||
commonMutex.unlock();
|
||||
s().waitIfStopRequested();
|
||||
V3LockGuard lock{commonMutex};
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds{sleep});
|
||||
commonValue = 1000;
|
||||
@@ -175,57 +113,56 @@ void V3ThreadPool::selfTest() {
|
||||
V3LockGuard lock{commonMutex};
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds{sleep});
|
||||
}
|
||||
s().requestExclusiveAccess([&]() { firstJob(sleep); });
|
||||
firstJob(sleep);
|
||||
V3LockGuard lock{commonMutex};
|
||||
commonValue = 1000;
|
||||
commonValue = 100;
|
||||
};
|
||||
std::list<std::future<void>> futures;
|
||||
|
||||
futures.push_back(s().enqueue(std::bind(firstJob, 100)));
|
||||
futures.push_back(s().enqueue(std::bind(secondJob, 100)));
|
||||
futures.push_back(s().enqueue(std::bind(firstJob, 100)));
|
||||
futures.push_back(s().enqueue(std::bind(secondJob, 100)));
|
||||
futures.push_back(s().enqueue(std::bind(secondJob, 200)));
|
||||
futures.push_back(s().enqueue(std::bind(firstJob, 200)));
|
||||
futures.push_back(s().enqueue(std::bind(firstJob, 300)));
|
||||
while (!futures.empty()) {
|
||||
s().waitForFuture(futures.front());
|
||||
futures.pop_front();
|
||||
}
|
||||
futures.push_back(s().enqueue(std::bind(thirdJob, 100)));
|
||||
futures.push_back(s().enqueue(std::bind(thirdJob, 100)));
|
||||
V3ThreadPool::waitForFutures(futures);
|
||||
|
||||
s().waitIfStopRequested();
|
||||
s().requestExclusiveAccess(std::bind(firstJob, 100));
|
||||
|
||||
auto forthJob = [&]() -> int { return 1234; };
|
||||
std::list<std::future<int>> futuresInt;
|
||||
futuresInt.push_back(s().enqueue(forthJob));
|
||||
auto result = V3ThreadPool::waitForFutures(futuresInt);
|
||||
UASSERT(result.back() == 1234, "unexpected future result = " << result.back());
|
||||
{
|
||||
const V3MtDisabledLockGuard mtDisabler{v3MtDisabledLock()};
|
||||
selfTestMtDisabled();
|
||||
{
|
||||
V3LockGuard lock{V3ThreadPool::s().m_mutex};
|
||||
UASSERT(V3ThreadPool::s().m_multithreadingSuspended,
|
||||
"Multithreading should be suspended at this point");
|
||||
}
|
||||
V3ThreadScope scope;
|
||||
|
||||
scope.enqueue(std::bind(firstJob, 100));
|
||||
scope.enqueue(std::bind(secondJob, 100));
|
||||
scope.enqueue(std::bind(firstJob, 100));
|
||||
scope.enqueue(std::bind(secondJob, 100));
|
||||
scope.enqueue(std::bind(secondJob, 200));
|
||||
scope.enqueue(std::bind(firstJob, 200));
|
||||
scope.enqueue(std::bind(firstJob, 300));
|
||||
scope.wait();
|
||||
|
||||
UASSERT(commonValue == 1000 || commonValue == 10,
|
||||
"unexpected common value = " << commonValue);
|
||||
|
||||
scope.enqueue(std::bind(thirdJob, 100));
|
||||
scope.enqueue(std::bind(thirdJob, 100));
|
||||
}
|
||||
|
||||
UASSERT(commonValue == 100, "unexpected common value = " << commonValue);
|
||||
|
||||
{
|
||||
V3LockGuard lock{V3ThreadPool::s().m_mutex};
|
||||
UASSERT(!V3ThreadPool::s().m_multithreadingSuspended,
|
||||
"Multithreading should not be suspended at this point");
|
||||
V3ThreadScope scope;
|
||||
scope.enqueue(std::bind(firstJob, 100));
|
||||
}
|
||||
|
||||
UASSERT(commonValue == 10, "unexpected common value = " << commonValue);
|
||||
|
||||
{
|
||||
int result = 0;
|
||||
auto forthJob = [&]() -> void { result = 1234; };
|
||||
|
||||
V3ThreadScope scope;
|
||||
scope.enqueue(forthJob);
|
||||
scope.wait();
|
||||
UASSERT(result == 1234, "unexpected job result = " << result);
|
||||
}
|
||||
selfTestMtDisabled();
|
||||
}
|
||||
|
||||
V3MtDisabledLock V3MtDisabledLock::s_mtDisabledLock;
|
||||
|
||||
void V3MtDisabledLock::lock() VL_ACQUIRE() VL_MT_SAFE {
|
||||
V3ThreadPool::s().suspendMultithreading();
|
||||
V3ThreadScope::V3ThreadScope() {
|
||||
UASSERT(v3Global.threadPoolp(), "ThreadPool must be initialized before ThreadScope.");
|
||||
m_pool = v3Global.threadPoolp();
|
||||
wait();
|
||||
}
|
||||
|
||||
void V3MtDisabledLock::unlock() VL_RELEASE() VL_MT_SAFE {
|
||||
V3ThreadPool::s().resumeMultithreading();
|
||||
}
|
||||
void V3ThreadScope::enqueue(std::function<void()>&& f) { m_pool->enqueue(std::move(f)); }
|
||||
|
||||
void V3ThreadScope::wait() { m_pool->wait(); }
|
||||
|
||||
Reference in New Issue
Block a user