coroutines truly are cursed
This commit is contained in:
+34
-10
@@ -46,42 +46,53 @@ ThreadPool::~ThreadPool()
|
||||
void ThreadPool::addJob(Job&& job)
|
||||
{
|
||||
std::unique_lock lock(jobQueueLock);
|
||||
jobQueue.add(job);
|
||||
//std::cout << "Adding job " << job << std::endl;
|
||||
jobQueue.add(std::move(job));
|
||||
jobQueueCV.notify_one();
|
||||
}
|
||||
void ThreadPool::addJob(MainJob&& job)
|
||||
{
|
||||
std::unique_lock lock(mainJobLock);
|
||||
mainJobs.add(job);
|
||||
//std::cout << "Adding main job " << job << std::endl;
|
||||
mainJobs.add(std::move(job));
|
||||
mainJobCV.notify_one();
|
||||
}
|
||||
void ThreadPool::enqueueWaiting(Event* event, Job&& job)
|
||||
{
|
||||
//std::cout << job << " waiting for event " << event->name << std::endl;
|
||||
std::unique_lock lock(waitingLock);
|
||||
waitingJobs[event].add(std::move(job));
|
||||
waitingJobs[*event].add(std::move(job));
|
||||
}
|
||||
void ThreadPool::enqueueWaiting(Event* event, MainJob&& job)
|
||||
{
|
||||
//std::cout << job << " main waiting for event " << event->name << std::endl;
|
||||
std::unique_lock lock(waitingMainLock);
|
||||
waitingMainJobs[event].add(std::move(job));
|
||||
waitingMainJobs[*event].add(std::move(job));
|
||||
}
|
||||
void ThreadPool::notify(Event* event)
|
||||
{
|
||||
{
|
||||
std::unique_lock lock(jobQueueLock);
|
||||
std::unique_lock lock2(waitingLock);
|
||||
for(auto&& job : waitingJobs[event])
|
||||
for(Job& job : waitingJobs[*event])
|
||||
{
|
||||
jobQueue.add(job);
|
||||
//std::cout << "Waking up job " << job << std::endl;
|
||||
jobQueue.add(std::move(job));
|
||||
jobQueueCV.notify_one();
|
||||
}
|
||||
waitingJobs[event].clear();
|
||||
waitingJobs[*event].clear();
|
||||
}
|
||||
{
|
||||
std::unique_lock lock(mainJobLock);
|
||||
std::unique_lock lock2(waitingMainLock);
|
||||
for(auto&& job : waitingMainJobs[event])
|
||||
auto&& jobsToQueue = std::move(waitingMainJobs[*event]);
|
||||
waitingMainJobs.erase(*event);
|
||||
for(MainJob& job : jobsToQueue)
|
||||
{
|
||||
mainJobs.add(job);
|
||||
//std::cout << "Waking up main job " << job << std::endl;
|
||||
mainJobs.add(std::move(job));
|
||||
mainJobCV.notify_one();
|
||||
}
|
||||
waitingMainJobs[event].clear();
|
||||
}
|
||||
event->reset();
|
||||
}
|
||||
@@ -95,9 +106,22 @@ void ThreadPool::threadLoop(const bool mainThread)
|
||||
if(!mainJobs.empty())
|
||||
{
|
||||
MainJob&& job = mainJobs.retrieve();
|
||||
lock.unlock();
|
||||
job.resume();
|
||||
}
|
||||
}
|
||||
std::unique_lock lock(jobQueueLock);
|
||||
if(jobQueue.empty())
|
||||
{
|
||||
jobQueueCV.wait(lock);
|
||||
}
|
||||
if(!jobQueue.empty())
|
||||
{
|
||||
Job&& job = jobQueue.retrieve();
|
||||
lock.unlock();
|
||||
//std::cout << "Starting job " << job << std::endl;
|
||||
job.resume();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user