2021-10-12 14:20:30 +02:00
|
|
|
#include "ThreadPool.h"
|
2021-12-02 13:00:03 +01:00
|
|
|
#include <memory_resource>
|
2021-10-12 14:20:30 +02:00
|
|
|
|
|
|
|
|
using namespace Seele;
|
|
|
|
|
|
2021-12-27 15:04:53 +01:00
|
|
|
std::atomic_uint64_t Seele::globalCounter;
|
|
|
|
|
|
2021-11-01 20:25:16 +01:00
|
|
|
Event::Event()
|
2022-02-14 16:29:26 +01:00
|
|
|
: flag(std::make_shared<StateStore>())
|
2021-12-15 00:05:42 +01:00
|
|
|
{
|
|
|
|
|
}
|
2021-10-23 00:22:35 +02:00
|
|
|
|
2021-12-02 13:00:03 +01:00
|
|
|
Event::Event(nullptr_t)
|
|
|
|
|
: flag(nullptr)
|
2021-12-15 00:05:42 +01:00
|
|
|
{
|
|
|
|
|
}
|
2021-12-02 13:00:03 +01:00
|
|
|
|
2021-12-15 00:05:42 +01:00
|
|
|
Event::Event(const std::string &name)
|
2022-01-12 14:40:26 +01:00
|
|
|
: name(name)
|
2022-02-14 16:29:26 +01:00
|
|
|
, flag(std::make_shared<StateStore>())
|
2021-12-15 00:05:42 +01:00
|
|
|
{
|
|
|
|
|
}
|
2021-11-01 20:25:16 +01:00
|
|
|
|
|
|
|
|
void Event::raise()
|
2021-10-23 00:22:35 +02:00
|
|
|
{
|
2022-02-14 16:29:26 +01:00
|
|
|
std::scoped_lock lock(flag->lock);
|
|
|
|
|
flag->data = 1;
|
2021-11-24 12:10:23 +01:00
|
|
|
getGlobalThreadPool().notify(*this);
|
2021-11-01 20:25:16 +01:00
|
|
|
}
|
|
|
|
|
void Event::reset()
|
2021-12-15 00:05:42 +01:00
|
|
|
{
|
2022-02-14 16:29:26 +01:00
|
|
|
std::scoped_lock lock(flag->lock);
|
|
|
|
|
flag->data = 0;
|
2021-10-23 00:22:35 +02:00
|
|
|
}
|
2021-10-12 14:20:30 +02:00
|
|
|
|
2021-11-01 20:25:16 +01:00
|
|
|
bool Event::await_ready()
|
|
|
|
|
{
|
2022-02-14 16:29:26 +01:00
|
|
|
flag->lock.lock();
|
|
|
|
|
bool result = flag->data;
|
|
|
|
|
if(result)
|
|
|
|
|
{
|
|
|
|
|
flag->lock.unlock();
|
|
|
|
|
}
|
|
|
|
|
return result;
|
2021-10-12 14:20:30 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
ThreadPool::ThreadPool(uint32 threadCount)
|
|
|
|
|
: workers(threadCount)
|
|
|
|
|
{
|
2021-11-01 20:25:16 +01:00
|
|
|
running.store(true);
|
2021-12-15 00:05:42 +01:00
|
|
|
for (uint32 i = 0; i < threadCount; ++i)
|
2021-11-01 20:25:16 +01:00
|
|
|
{
|
2021-12-02 13:00:03 +01:00
|
|
|
workers[i] = std::thread(&ThreadPool::threadLoop, this, false);
|
2021-11-01 20:25:16 +01:00
|
|
|
}
|
2021-10-12 14:20:30 +02:00
|
|
|
}
|
|
|
|
|
|
2021-10-23 00:22:35 +02:00
|
|
|
ThreadPool::~ThreadPool()
|
2021-10-12 14:20:30 +02:00
|
|
|
{
|
2022-02-14 16:29:26 +01:00
|
|
|
cleanup();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
void ThreadPool::cleanup()
|
|
|
|
|
{
|
|
|
|
|
bool temp = true;
|
|
|
|
|
running.compare_exchange_strong(temp, false);
|
|
|
|
|
if(!temp)
|
|
|
|
|
return;
|
2022-01-12 14:40:26 +01:00
|
|
|
{
|
|
|
|
|
std::scoped_lock lock(jobQueueLock);
|
|
|
|
|
jobQueueCV.notify_all();
|
|
|
|
|
}
|
|
|
|
|
{
|
|
|
|
|
std::scoped_lock lock(mainJobLock);
|
|
|
|
|
mainJobCV.notify_all();
|
|
|
|
|
}
|
2021-12-15 00:05:42 +01:00
|
|
|
for (auto &thread : workers)
|
2021-10-23 00:22:35 +02:00
|
|
|
{
|
|
|
|
|
thread.join();
|
|
|
|
|
}
|
2022-01-12 14:40:26 +01:00
|
|
|
workers.clear();
|
|
|
|
|
waitingJobs.clear();
|
|
|
|
|
waitingMainJobs.clear();
|
2021-10-23 00:22:35 +02:00
|
|
|
}
|
2021-12-15 00:05:42 +01:00
|
|
|
void ThreadPool::enqueueWaiting(Event &event, Promise* job)
|
2021-11-01 20:25:16 +01:00
|
|
|
{
|
2021-12-15 00:05:42 +01:00
|
|
|
assert(!job->done());
|
2022-01-12 14:40:26 +01:00
|
|
|
std::scoped_lock lock(waitingLock);
|
|
|
|
|
//std::cout << "Job " << job->finishedEvent.name << " waiting on event " << event.name << std::endl;
|
|
|
|
|
waitingJobs[event].push_back(job);
|
2022-02-14 16:29:26 +01:00
|
|
|
job->addRef();
|
2021-11-01 20:25:16 +01:00
|
|
|
}
|
2021-12-15 00:05:42 +01:00
|
|
|
void ThreadPool::enqueueWaiting(Event &event, MainPromise* job)
|
2021-11-01 20:25:16 +01:00
|
|
|
{
|
2021-12-15 00:05:42 +01:00
|
|
|
assert(!job->done());
|
2022-01-12 14:40:26 +01:00
|
|
|
std::scoped_lock lock(waitingMainLock);
|
|
|
|
|
//std::cout << job->finishedEvent.name << " waiting on event " << event.name << std::endl;
|
|
|
|
|
waitingMainJobs[event].push_back(job);
|
2022-02-14 16:29:26 +01:00
|
|
|
job->addRef();
|
2022-01-12 14:40:26 +01:00
|
|
|
}
|
|
|
|
|
void ThreadPool::scheduleJob(Promise* job)
|
|
|
|
|
{
|
|
|
|
|
assert(!job->done());
|
|
|
|
|
std::scoped_lock lock(jobQueueLock);
|
|
|
|
|
//std::cout << "Queueing job " << job->finishedEvent.name << std::endl;
|
|
|
|
|
jobQueue.push_back(job);
|
|
|
|
|
jobQueueCV.notify_one();
|
2022-02-14 16:29:26 +01:00
|
|
|
job->addRef();
|
2022-01-12 14:40:26 +01:00
|
|
|
}
|
|
|
|
|
void ThreadPool::scheduleJob(MainPromise* job)
|
|
|
|
|
{
|
|
|
|
|
assert(!job->done());
|
|
|
|
|
std::scoped_lock lock(mainJobLock);
|
|
|
|
|
//std::cout << "Queueing job " << job->finishedEvent.name << std::endl;
|
|
|
|
|
mainJobs.push_back(job);
|
|
|
|
|
mainJobCV.notify_one();
|
2022-02-14 16:29:26 +01:00
|
|
|
job->addRef();
|
2021-11-01 20:25:16 +01:00
|
|
|
}
|
2021-12-15 00:05:42 +01:00
|
|
|
void ThreadPool::notify(Event &event)
|
2021-11-01 20:25:16 +01:00
|
|
|
{
|
2021-12-22 11:42:07 +01:00
|
|
|
//std::cout << "Event " << event.name << " raised" << std::endl;
|
2021-11-01 20:25:16 +01:00
|
|
|
{
|
2022-01-12 14:40:26 +01:00
|
|
|
std::scoped_lock lock(jobQueueLock, waitingLock);
|
|
|
|
|
std::list<Promise*> jobs = std::move(waitingJobs[event]);
|
|
|
|
|
waitingJobs.erase(event);
|
2021-12-15 00:05:42 +01:00
|
|
|
for (auto &job : jobs)
|
2021-11-01 20:25:16 +01:00
|
|
|
{
|
2021-11-24 12:10:23 +01:00
|
|
|
//assert(job.id != -1ull);
|
2021-12-22 11:42:07 +01:00
|
|
|
//std::cout << "Waking up " << job->finishedEvent.name << std::endl;
|
2021-12-27 15:04:53 +01:00
|
|
|
job->state = Promise::State::SCHEDULED;
|
2022-01-12 14:40:26 +01:00
|
|
|
jobQueue.push_back(job);
|
2021-11-11 20:12:50 +01:00
|
|
|
jobQueueCV.notify_one();
|
2021-11-01 20:25:16 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
{
|
2022-01-12 14:40:26 +01:00
|
|
|
std::scoped_lock lock(mainJobLock, waitingMainLock);
|
|
|
|
|
std::list<MainPromise*> jobs = std::move(waitingMainJobs[event]);
|
|
|
|
|
waitingMainJobs.erase(event);
|
2021-12-15 00:05:42 +01:00
|
|
|
for (auto &job : jobs)
|
2021-11-01 20:25:16 +01:00
|
|
|
{
|
2021-11-24 12:10:23 +01:00
|
|
|
//assert(job.id != -1ull);
|
2021-12-22 11:42:07 +01:00
|
|
|
//std::cout << "Waking up main " << job->finishedEvent.name << std::endl;
|
2021-12-27 15:04:53 +01:00
|
|
|
job->state = MainPromise::State::SCHEDULED;
|
2022-01-12 14:40:26 +01:00
|
|
|
mainJobs.push_back(job);
|
2021-11-11 20:12:50 +01:00
|
|
|
mainJobCV.notify_one();
|
2021-11-01 20:25:16 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
2021-12-15 00:05:42 +01:00
|
|
|
void ThreadPool::threadLoop(const bool mainThread)
|
2021-11-19 15:08:56 +01:00
|
|
|
{
|
2021-12-15 00:05:42 +01:00
|
|
|
while (running.load())
|
2021-11-19 15:08:56 +01:00
|
|
|
{
|
2021-12-15 00:05:42 +01:00
|
|
|
if (mainThread)
|
2021-11-19 15:08:56 +01:00
|
|
|
{
|
2021-12-15 00:05:42 +01:00
|
|
|
MainPromise* job;
|
|
|
|
|
{
|
|
|
|
|
std::unique_lock lock(mainJobLock);
|
|
|
|
|
if(mainJobs.empty())
|
|
|
|
|
{
|
|
|
|
|
mainJobCV.wait(lock);
|
|
|
|
|
}
|
|
|
|
|
if (!mainJobs.empty())
|
|
|
|
|
{
|
2022-01-12 14:40:26 +01:00
|
|
|
job = mainJobs.front();
|
|
|
|
|
mainJobs.pop_front();
|
2021-12-15 00:05:42 +01:00
|
|
|
}
|
|
|
|
|
else
|
|
|
|
|
{
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
job->resume();
|
2022-02-14 16:29:26 +01:00
|
|
|
job->removeRef();
|
2021-11-19 15:08:56 +01:00
|
|
|
}
|
|
|
|
|
else
|
|
|
|
|
{
|
2021-12-15 00:05:42 +01:00
|
|
|
Promise* job;
|
2021-11-01 20:25:16 +01:00
|
|
|
{
|
2021-12-15 00:05:42 +01:00
|
|
|
std::unique_lock lock(jobQueueLock);
|
|
|
|
|
if (jobQueue.empty())
|
|
|
|
|
{
|
|
|
|
|
jobQueueCV.wait(lock);
|
|
|
|
|
}
|
|
|
|
|
if (!jobQueue.empty())
|
|
|
|
|
{
|
2022-01-12 14:40:26 +01:00
|
|
|
job = jobQueue.front();
|
|
|
|
|
jobQueue.pop_front();
|
2021-12-15 00:05:42 +01:00
|
|
|
}
|
|
|
|
|
else
|
|
|
|
|
{
|
|
|
|
|
continue;
|
|
|
|
|
}
|
2021-11-19 15:08:56 +01:00
|
|
|
}
|
2021-12-15 00:05:42 +01:00
|
|
|
//std::cout << "Starting job " << job.id << std::endl;
|
|
|
|
|
job->resume();
|
2022-02-14 16:29:26 +01:00
|
|
|
job->removeRef();
|
2021-11-24 12:10:23 +01:00
|
|
|
}
|
2021-11-01 20:25:16 +01:00
|
|
|
}
|
|
|
|
|
}
|
2021-12-15 00:05:42 +01:00
|
|
|
ThreadPool &Seele::getGlobalThreadPool()
|
2021-10-23 00:22:35 +02:00
|
|
|
{
|
|
|
|
|
static ThreadPool threadPool;
|
|
|
|
|
return threadPool;
|
2021-10-12 14:20:30 +02:00
|
|
|
}
|