Works, but with memory leaks
This commit is contained in:
+75
-94
@@ -3,55 +3,69 @@
|
||||
|
||||
using namespace Seele;
|
||||
|
||||
|
||||
Event::Event()
|
||||
: flag(std::make_shared<StateStore>())
|
||||
{
|
||||
}
|
||||
std::mutex Seele::promisesLock;
|
||||
List<Promise*> Seele::promises;
|
||||
|
||||
Event::Event(nullptr_t)
|
||||
: flag(nullptr)
|
||||
{
|
||||
}
|
||||
|
||||
Event::Event(const std::string &name)
|
||||
: flag(std::make_shared<StateStore>())
|
||||
Event::Event(const std::string &name, const std::source_location &location)
|
||||
: name(name)
|
||||
, location(location)
|
||||
{
|
||||
flag->name = name;
|
||||
}
|
||||
|
||||
|
||||
Event::Event(const std::source_location &location)
|
||||
: flag(std::make_shared<StateStore>())
|
||||
: name(location.function_name())
|
||||
, location(location)
|
||||
{
|
||||
flag->name = location.function_name();
|
||||
flag->location = location;
|
||||
}
|
||||
|
||||
|
||||
void Event::raise()
|
||||
{
|
||||
std::scoped_lock lock(flag->lock);
|
||||
flag->data = 1;
|
||||
getGlobalThreadPool().notify(*this);
|
||||
std::scoped_lock lock(eventLock);
|
||||
data = true;
|
||||
if(waitingJobs.size() > 0)
|
||||
{
|
||||
getGlobalThreadPool().scheduleBatch(waitingJobs);
|
||||
}
|
||||
if(waitingMainJobs.size() > 0)
|
||||
{
|
||||
getGlobalThreadPool().scheduleBatch(waitingMainJobs);
|
||||
}
|
||||
}
|
||||
void Event::reset()
|
||||
{
|
||||
std::scoped_lock lock(flag->lock);
|
||||
flag->data = 0;
|
||||
std::scoped_lock lock(eventLock);
|
||||
data = false;
|
||||
}
|
||||
|
||||
bool Event::await_ready()
|
||||
{
|
||||
flag->lock.lock();
|
||||
bool result = flag->data;
|
||||
eventLock.lock();
|
||||
bool result = data;
|
||||
if(result)
|
||||
{
|
||||
flag->lock.unlock();
|
||||
eventLock.unlock();
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
void Event::await_suspend(std::coroutine_handle<JobPromiseBase<false>> h)
|
||||
{
|
||||
//h.promise().enqueue(this);
|
||||
waitingJobs.add(JobBase<false>(&h.promise()));
|
||||
eventLock.unlock();
|
||||
}
|
||||
|
||||
void Event::await_suspend(std::coroutine_handle<JobPromiseBase<true>> h)
|
||||
{
|
||||
//h.promise().enqueue(this);
|
||||
waitingMainJobs.add(JobBase<true>(&h.promise()));
|
||||
eventLock.unlock();
|
||||
}
|
||||
|
||||
ThreadPool::ThreadPool(uint32 threadCount)
|
||||
: workers(threadCount)
|
||||
{
|
||||
@@ -59,13 +73,24 @@ ThreadPool::ThreadPool(uint32 threadCount)
|
||||
for (uint32 i = 0; i < threadCount; ++i)
|
||||
{
|
||||
workers[i] = std::thread(&ThreadPool::threadLoop, this);
|
||||
workers[i].detach();
|
||||
}
|
||||
}
|
||||
|
||||
ThreadPool::~ThreadPool()
|
||||
{
|
||||
workers.clear();
|
||||
running.store(false);
|
||||
{
|
||||
std::unique_lock lock(mainJobLock);
|
||||
mainJobCV.notify_all();
|
||||
}
|
||||
{
|
||||
std::unique_lock lock(jobQueueLock);
|
||||
jobQueueCV.notify_all();
|
||||
}
|
||||
for(auto& worker : workers)
|
||||
{
|
||||
worker.join();
|
||||
}
|
||||
}
|
||||
|
||||
void ThreadPool::waitIdle()
|
||||
@@ -73,108 +98,63 @@ void ThreadPool::waitIdle()
|
||||
while(true)
|
||||
{
|
||||
std::unique_lock lock(numIdlingLock);
|
||||
numIdlingIncr.wait(lock);
|
||||
if(numIdling == workers.size())
|
||||
{
|
||||
return;
|
||||
}
|
||||
numIdlingIncr.wait(lock);
|
||||
}
|
||||
}
|
||||
void ThreadPool::enqueueWaiting(Event &event, Promise* job)
|
||||
void ThreadPool::scheduleJob(Job job)
|
||||
{
|
||||
assert(!job->done());
|
||||
std::scoped_lock lock(waitingLock);
|
||||
//std::cout << "Job " << job->finishedEvent.name << " waiting on event " << event.name << std::endl;
|
||||
waitingJobs[event].add(job);
|
||||
job->addRef();
|
||||
}
|
||||
void ThreadPool::enqueueWaiting(Event &event, MainPromise* job)
|
||||
{
|
||||
assert(!job->done());
|
||||
std::scoped_lock lock(waitingMainLock);
|
||||
//std::cout << job->finishedEvent.name << " waiting on event " << event.name << std::endl;
|
||||
waitingMainJobs[event].add(job);
|
||||
job->addRef();
|
||||
}
|
||||
void ThreadPool::scheduleJob(Promise* job)
|
||||
{
|
||||
assert(!job->done());
|
||||
assert(!job.done());
|
||||
std::scoped_lock lock(jobQueueLock);
|
||||
//std::cout << "Queueing job " << job->finishedEvent << std::endl;
|
||||
jobQueue.add(job);
|
||||
jobQueue.add(std::move(job));
|
||||
jobQueueCV.notify_one();
|
||||
job->addRef();
|
||||
}
|
||||
void ThreadPool::scheduleJob(MainPromise* job)
|
||||
void ThreadPool::scheduleJob(MainJob job)
|
||||
{
|
||||
assert(!job->done());
|
||||
assert(!job.done());
|
||||
std::scoped_lock lock(mainJobLock);
|
||||
//std::cout << "Queueing job " << job->finishedEvent << std::endl;
|
||||
mainJobs.add(job);
|
||||
mainJobs.add(std::move(job));
|
||||
mainJobCV.notify_one();
|
||||
job->addRef();
|
||||
}
|
||||
void ThreadPool::notify(Event &event)
|
||||
{
|
||||
//std::cout << "Event " << event.name << " raised" << std::endl;
|
||||
{
|
||||
std::scoped_lock lock(jobQueueLock, waitingLock);
|
||||
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);
|
||||
jobQueueCV.notify_one();
|
||||
}
|
||||
}
|
||||
{
|
||||
std::scoped_lock lock(mainJobLock, waitingMainLock);
|
||||
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);
|
||||
mainJobCV.notify_one();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void ThreadPool::mainLoop()
|
||||
{
|
||||
while(true)
|
||||
while(running.load())
|
||||
{
|
||||
MainPromise* job;
|
||||
MainJob job;
|
||||
{
|
||||
std::unique_lock lock(mainJobLock);
|
||||
if(mainJobs.empty())
|
||||
{
|
||||
mainJobCV.wait(lock);
|
||||
}
|
||||
job = mainJobs.front();
|
||||
mainJobs.popFront();
|
||||
[[likely]]
|
||||
if(!mainJobs.empty())
|
||||
{
|
||||
job = mainJobs.front();
|
||||
mainJobs.popFront();
|
||||
}
|
||||
else
|
||||
{
|
||||
continue;
|
||||
}
|
||||
}
|
||||
job->resume();
|
||||
job->removeRef();
|
||||
job.resume();
|
||||
}
|
||||
}
|
||||
|
||||
void ThreadPool::threadLoop()
|
||||
{
|
||||
List<Promise*> localQueue;
|
||||
while (true)
|
||||
List<Job> localQueue;
|
||||
while (running.load())
|
||||
{
|
||||
[[likely]]
|
||||
if(!localQueue.empty())
|
||||
{
|
||||
Promise* job = localQueue.retrieve();
|
||||
job->resume();
|
||||
job->removeRef();
|
||||
Job job = localQueue.retrieve();
|
||||
job.resume();
|
||||
}
|
||||
else
|
||||
{
|
||||
@@ -194,7 +174,8 @@ void ThreadPool::threadLoop()
|
||||
}
|
||||
// take 1/numThreads jobs, maybe make this a parameter that
|
||||
// adjusts based on past workload
|
||||
uint32 numTaken = std::max(jobQueue.size() / workers.size(), 1ull);
|
||||
uint32 partitionedWorkload = (uint32)(jobQueue.size() / workers.size());
|
||||
uint32 numTaken = std::clamp(partitionedWorkload, 1u, localQueueSize);
|
||||
while (!jobQueue.empty() && localQueue.size() < numTaken)
|
||||
{
|
||||
localQueue.add(jobQueue.retrieve());
|
||||
|
||||
Reference in New Issue
Block a user