Basic mutability framework

This commit is contained in:
Dynamitos
2022-01-12 14:40:26 +01:00
parent b2718dde70
commit 6d48267ec2
72 changed files with 1781 additions and 1341 deletions
+51 -43
View File
@@ -6,7 +6,7 @@ using namespace Seele;
std::atomic_uint64_t Seele::globalCounter;
Event::Event()
: flag(new std::atomic_bool())
: flag(std::make_shared<std::atomic_bool>())
{
}
@@ -16,7 +16,8 @@ Event::Event(nullptr_t)
}
Event::Event(const std::string &name)
: name(name), flag(new std::atomic_bool())
: name(name)
, flag(std::make_shared<std::atomic_bool>())
{
}
@@ -48,75 +49,80 @@ ThreadPool::ThreadPool(uint32 threadCount)
ThreadPool::~ThreadPool()
{
running.store(false);
{
std::scoped_lock lock(jobQueueLock);
jobQueueCV.notify_all();
}
{
std::scoped_lock lock(mainJobLock);
mainJobCV.notify_all();
}
for (auto &thread : workers)
{
thread.join();
}
workers.clear();
waitingJobs.clear();
waitingMainJobs.clear();
}
void ThreadPool::enqueueWaiting(Event &event, Promise* job)
{
assert(!job->done());
if(event == nullptr)
{
std::unique_lock lock(jobQueueLock);
//std::cout << "Queueing job " << job->finishedEvent.name << std::endl;
jobQueue.add(job);
jobQueueCV.notify_one();
}
else
{
std::unique_lock lock(waitingLock);
//std::cout << "Job " << job->finishedEvent.name << " waiting on event " << event.name << std::endl;
waitingJobs[event].add(job);
}
std::scoped_lock lock(waitingLock);
//std::cout << "Job " << job->finishedEvent.name << " waiting on event " << event.name << std::endl;
waitingJobs[event].push_back(job);
}
void ThreadPool::enqueueWaiting(Event &event, MainPromise* job)
{
assert(!job->done());
if(event == nullptr || event)
{
std::unique_lock lock(mainJobLock);
//std::cout << "Queueing job " << job->finishedEvent.name << std::endl;
mainJobs.add(job);
mainJobCV.notify_one();
}
else
{
std::unique_lock lock(waitingMainLock);
//std::cout << job->finishedEvent.name << " waiting on event " << event.name << std::endl;
waitingMainJobs[event].add(job);
}
std::scoped_lock lock(waitingMainLock);
//std::cout << job->finishedEvent.name << " waiting on event " << event.name << std::endl;
waitingMainJobs[event].push_back(job);
}
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();
}
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();
}
void ThreadPool::notify(Event &event)
{
//std::cout << "Event " << event.name << " raised" << std::endl;
{
std::unique_lock lock(jobQueueLock);
std::unique_lock lock2(waitingLock);
List<Promise*> &jobs = waitingJobs[event];
std::scoped_lock lock(jobQueueLock, waitingLock);
std::list<Promise*> jobs = std::move(waitingJobs[event]);
waitingJobs.erase(event);
for (auto &job : jobs)
{
//assert(job.id != -1ull);
//std::cout << "Waking up " << job->finishedEvent.name << std::endl;
job->state = Promise::State::SCHEDULED;
jobQueue.add(job);
jobQueue.push_back(job);
jobQueueCV.notify_one();
}
waitingJobs.erase(event);
}
{
std::unique_lock lock(mainJobLock);
std::unique_lock lock2(waitingMainLock);
List<MainPromise*> &jobs = waitingMainJobs[event];
std::scoped_lock lock(mainJobLock, waitingMainLock);
std::list<MainPromise*> jobs = std::move(waitingMainJobs[event]);
waitingMainJobs.erase(event);
for (auto &job : jobs)
{
//assert(job.id != -1ull);
//std::cout << "Waking up main " << job->finishedEvent.name << std::endl;
job->state = MainPromise::State::SCHEDULED;
mainJobs.add(job);
mainJobs.push_back(job);
mainJobCV.notify_one();
}
waitingMainJobs.erase(event);
}
}
void ThreadPool::threadLoop(const bool mainThread)
@@ -134,7 +140,8 @@ void ThreadPool::threadLoop(const bool mainThread)
}
if (!mainJobs.empty())
{
job = mainJobs.retrieve();
job = mainJobs.front();
mainJobs.pop_front();
}
else
{
@@ -144,7 +151,7 @@ void ThreadPool::threadLoop(const bool mainThread)
job->resume();
if (job->done())
{
job->raise();
job->finalize();
}
}
else
@@ -158,7 +165,8 @@ void ThreadPool::threadLoop(const bool mainThread)
}
if (!jobQueue.empty())
{
job = jobQueue.retrieve();
job = jobQueue.front();
jobQueue.pop_front();
}
else
{
@@ -167,9 +175,9 @@ void ThreadPool::threadLoop(const bool mainThread)
}
//std::cout << "Starting job " << job.id << std::endl;
job->resume();
if (job->done())
if(job->done())
{
job->raise();
job->finalize();
}
}
}