7#include <condition_variable>
23enum class JobScope { Heap, Frame };
27thread_local const void *tlsCurrentState =
nullptr;
28thread_local JobImpl *tlsCurrentJob =
nullptr;
30bool isDoneStatus(JobStatus
status) {
31 return status == JobStatus::Done ||
status == JobStatus::Failed;
45 FrameArena() =
default;
46 ~FrameArena() { reset(); }
54 void *alloc(
size_t size,
size_t align) {
55 if (
size > kArenaBlockSize)
56 throw eve::Exception(
"JobSystem arena allocation exceeds block size");
58 blocks_.emplace_back(kArenaBlockSize);
59 size_t aligned = (offset_ + align - 1) & ~(align - 1);
60 if (aligned +
size > blocks_.back().size()) {
61 blocks_.emplace_back(kArenaBlockSize);
64 void *
ptr = blocks_.back().data() + aligned;
65 offset_ = aligned +
size;
74 void track(
void *
ptr,
void (*
destroy)(
void *)) {
80 for (
auto &entry : tracked_)
87 static constexpr size_t kArenaBlockSize = 64 * 1024;
94 std::vector<std::vector<std::byte>> blocks_;
95 std::vector<Tracked> tracked_;
102 mutable std::mutex
mu;
103 std::condition_variable
cv;
123 void runJob(JobImpl *job);
137struct JobImpl final :
public Job {
141 ~JobImpl()
override {
149 void wait()
override;
150 bool isDone()
const override;
151 bool hasFailed()
const override;
152 std::string getError()
const override;
153 int getPendingDependencyCount()
const override;
174class TaskGroupImpl final :
public TaskGroup {
176 TaskGroupImpl(JobSystemThreadPool::State *owner, JobScope scope_)
178 ~TaskGroupImpl()
override =
default;
185 JobImpl *job =
nullptr;
187 std::lock_guard<std::mutex> lock(
state->mu);
191 state->scheduleLocked(job);
194 state->cv.notify_one();
198 void wait()
override {
199 std::vector<JobImpl *> snapshot;
201 std::lock_guard<std::mutex> lock(
state->mu);
204 for (JobImpl *child : snapshot)
206 std::lock_guard<std::mutex> lock(
state->mu);
210 int getPendingCount()
const override {
211 std::lock_guard<std::mutex> lock(
state->mu);
214 if (!isDoneStatus(
child->status))
220 JobSystemThreadPool::State *
state;
231 if (
scope == JobScope::Frame) {
232 void *mem =
arena.alloc(
sizeof(JobImpl),
alignof(JobImpl));
234 arena.track(job, [](
void *
ptr) {
static_cast<JobImpl *
>(
ptr)->~JobImpl(); });
241 job->scheduled =
true;
243 if (job->scope == JobScope::Frame)
245 if (job->depCount == 0)
252 job->enqueued =
true;
253 ready.push_back(job);
258 std::lock_guard<std::mutex> lock(
mu);
259 job->status = JobStatus::Running;
265 const void *previousState = tlsCurrentState;
266 if (previousState !=
this)
267 tlsCurrentState =
this;
277 recordError(job, e.
what());
278 }
catch (
const std::exception &e) {
280 recordError(job, e.
what());
283 recordError(job,
"unknown exception");
288 completeJob(job,
ok);
289 tlsCurrentState = previousState;
294 std::lock_guard<std::mutex> lock(
mu);
295 job->status =
ok ? JobStatus::Done : JobStatus::Failed;
313 std::lock_guard<std::mutex> lock(
mu);
315 if (job->scope == JobScope::Frame)
317 releaseDependents(job);
318 job->completionDone =
true;
321 std::lock_guard<std::mutex> lock(
mu);
328 if (!job->onComplete)
333 recordError(job, std::string(
"completion callback: ") + e.
what());
334 }
catch (
const std::exception &e) {
335 recordError(job, std::string(
"completion callback: ") + e.
what());
337 recordError(job,
"completion callback: unknown exception");
342 std::lock_guard<std::mutex> lock(
mu);
343 if (!job->error.empty())
353 for (JobImpl *dep : job->dependents) {
354 if (dep->depCount > 0)
356 if (dep->depCount == 0 && dep->scheduled && !dep->enqueued) {
361 job->dependents.clear();
367 std::unique_lock<std::mutex> lock(
mu);
368 if (job->completionDone)
370 if (tlsCurrentJob == job)
374 if (job->completionDone)
376 if (!ready.empty()) {
377 JobImpl *next = ready.front();
391 std::unique_lock<std::mutex> lock(
mu);
392 while (outstandingFrame > 0 || completing > 0) {
393 if (!ready.empty()) {
394 JobImpl *next = ready.front();
409 const int target = std::max(1, workerCount * 4);
414 tlsCurrentState =
state;
416 JobImpl *job =
nullptr;
418 std::unique_lock<std::mutex> lock(
state->mu);
422 job =
state->ready.front();
423 state->ready.pop_front();
428 tlsCurrentState =
nullptr;
433void JobImpl::wait() {
434 state->waitJob(
this);
437bool JobImpl::isDone()
const {
438 std::lock_guard<std::mutex> lock(
state->mu);
442bool JobImpl::hasFailed()
const {
443 std::lock_guard<std::mutex> lock(
state->mu);
444 return status == JobStatus::Failed;
447std::string JobImpl::getError()
const {
448 std::lock_guard<std::mutex> lock(
state->mu);
452int JobImpl::getPendingDependencyCount()
const {
453 std::lock_guard<std::mutex> lock(
state->mu);
459 if (!pred || pred->state != state)
460 throw eve::Exception(
"JobSystem::addDependency: job does not belong to this system");
462 throw eve::Exception(
"JobSystem::addDependency: a job cannot depend on itself");
464 std::lock_guard<std::mutex> lock(
state->mu);
466 throw eve::Exception(
"JobSystem::addDependency: job already scheduled");
467 if (isDoneStatus(pred->status))
470 pred->dependents.push_back(
this);
474 bool alreadyDone =
false;
476 std::lock_guard<std::mutex> lock(
state->mu);
477 if (isDoneStatus(
status)) {
488 state->recordError(
this, std::string(
"completion callback: ") + e.
what());
489 }
catch (
const std::exception &e) {
490 state->recordError(
this, std::string(
"completion callback: ") + e.what());
492 state->recordError(
this,
"completion callback: unknown exception");
501JobImpl *castJob(JobSystemThreadPool::State *state, Job *job) {
504 auto *
impl =
dynamic_cast<JobImpl *
>(job);
506 throw eve::Exception(
"JobSystem: job does not belong to this system");
513 if (workerCount <= 0) {
514 workerCount =
static_cast<int>(std::thread::hardware_concurrency());
515 if (workerCount <= 0)
518 state_ = std::make_shared<State>();
519 pool_ = std::make_unique<ThreadPool>(workerCount);
520 state_->workerCount = pool_->getWorkerCount();
521 poolTasks_.reserve(
static_cast<size_t>(state_->workerCount));
523 for (
int i = 0; i < state_->workerCount; ++i) {
525 poolTasks_.push_back(pool_->submit([
state] { State::workerLoop(state.get()); }));
529 std::lock_guard<std::mutex> lock(state_->mu);
530 state_->stopping =
true;
532 state_->cv.notify_all();
535 for (
Task *task : poolTasks_)
544 std::lock_guard<std::mutex> lifecycleLock(lifecycleMu_);
545 const bool fromWorker = tlsCurrentState == state_.get();
547 std::lock_guard<std::mutex> lock(state_->mu);
548 state_->stopping =
true;
550 state_->cv.notify_all();
557 for (
Task *task : poolTasks_)
565 for (
Task *task : poolTasks_)
573 return state_->workerCount;
577 std::lock_guard<std::mutex> lock(state_->mu);
578 return !state_->stopping;
582 std::lock_guard<std::mutex> lock(state_->mu);
583 return static_cast<int>(state_->ready.size());
587 std::lock_guard<std::mutex> lock(state_->mu);
588 return state_->outstanding;
592 JobImpl *job =
nullptr;
594 std::lock_guard<std::mutex> lock(state_->mu);
595 if (state_->stopping)
597 job = state_->createJobLocked(JobScope::Heap, std::move(
body),
JobFunc{});
598 state_->scheduleLocked(job);
600 state_->cv.notify_one();
605 std::lock_guard<std::mutex> lock(state_->mu);
606 return state_->createJobLocked(JobScope::Heap, std::move(
body),
JobFunc{});
610 JobImpl *
impl = castJob(state_.get(), job);
612 std::lock_guard<std::mutex> lock(state_->mu);
613 if (state_->stopping)
616 throw eve::Exception(
"JobSystem::schedule: job already scheduled");
617 state_->scheduleLocked(
impl);
619 state_->cv.notify_one();
623 return parallelForImpl(
first, last, std::move(
body), chunk,
false);
627 return new TaskGroupImpl(state_.get(), JobScope::Heap);
631 JobImpl *job =
nullptr;
633 std::lock_guard<std::mutex> lock(state_->mu);
634 if (state_->stopping)
636 job = state_->createJobLocked(JobScope::Frame, std::move(
body),
JobFunc{});
637 state_->scheduleLocked(job);
639 state_->cv.notify_one();
644 std::lock_guard<std::mutex> lock(state_->mu);
645 return state_->createJobLocked(JobScope::Frame, std::move(
body),
JobFunc{});
649 return parallelForImpl(
first, last, std::move(
body), chunk,
true);
653 std::lock_guard<std::mutex> lock(state_->mu);
654 void *mem = state_->arena.alloc(
sizeof(TaskGroupImpl),
alignof(TaskGroupImpl));
655 auto *
group =
new (mem) TaskGroupImpl(state_.get(), JobScope::Frame);
656 state_->arena.track(
group, [](
void *
ptr) {
static_cast<TaskGroupImpl *
>(
ptr)->~TaskGroupImpl(); });
661 state_->waitFrameJobs();
665 state_->waitFrameJobs();
669 if (tlsCurrentState == state_.get())
670 throw eve::Exception(
"JobSystem::waitAll cannot be called from its worker");
672 std::unique_lock<std::mutex> lock(
state->mu);
673 while (
state->outstanding > 0) {
674 if (!
state->ready.empty()) {
675 JobImpl *next =
state->ready.front();
676 state->ready.pop_front();
684 state->cv.wait(lock);
689 std::lock_guard<std::mutex> lifecycleLock(lifecycleMu_);
690 if (tlsCurrentState == state_.get())
691 throw eve::Exception(
"JobSystem::stop cannot be called from its worker");
693 std::lock_guard<std::mutex> lock(state_->mu);
694 state_->stopping =
true;
696 state_->cv.notify_all();
699 for (
Task *task : poolTasks_)
711 const JobScope
scope = frameScope ? JobScope::Frame : JobScope::Heap;
714 chunk = state_->autoChunk(
count);
715 const int taskCount = (
count + chunk - 1) / chunk;
717 std::lock_guard<std::mutex> lock(state_->mu);
718 if (state_->stopping)
722 join->depCount = taskCount;
723 if (
scope == JobScope::Heap)
724 join->ownedChildren.reserve(
static_cast<size_t>(taskCount));
727 children.reserve(
static_cast<size_t>(taskCount));
728 for (
int t = 0;
t < taskCount; ++
t) {
730 const int end = std::min(last,
begin + chunk);
733 child->dependents.push_back(join);
735 if (
scope == JobScope::Heap)
736 join->ownedChildren.push_back(child);
741 state_->scheduleLocked(join);
743 state_->scheduleLocked(
child);
744 state_->cv.notify_all();
std::vector< QuestEvent > pending
building::EdgeCurveGroup group
wgpu::PopErrorScopeStatus status
std::vector< JobImpl * > dependents
std::vector< JobImpl * > ownedChildren
graphics::Canvas * previous
std::map< Cell, Cell > predecessor
EVENGINE_API_FOUNDATION public API.
virtual const char * what() const
Returns a string containing reason for the exception.
Job * createJob(JobFunc body) override
Creates job.
void endFrame() override
Ends frame.
Job * submitFrame(JobFunc body) override
Submit frame.
Job * parallelFor(int first, int last, ParallelForBody body, int chunk=1) override
Parallel for.
Job * parallelForFrame(int first, int last, ParallelForBody body, int chunk=1) override
Parallel for frame.
TaskGroup * createTaskGroup() override
Creates task group.
int getPendingCount() const override
Returns the pending count.
JobSystemThreadPool(int workerCount)
Constructs a JobSystemThreadPool.
void waitAll() override
Waits all.
int getOutstandingCount() const override
Returns the outstanding count.
int getWorkerCount() const override
Returns the worker count.
Job * submit(JobFunc body) override
Submit.
void stop() override
Stops .
void schedule(Job *job) override
Schedule.
TaskGroup * createFrameTaskGroup() override
Creates frame task group.
bool isRunning() const override
True when running.
~JobSystemThreadPool() override
Releases JobSystemThreadPool resources.
Job * createFrameJob(JobFunc body) override
Creates frame job.
void beginFrame() override
Begins frame.
Handle to a single job inside a JobSystem.
Fork/join group: fork() spawns scheduled children, wait() joins them.
A job executed by a ThreadPool worker. Status strings (no enums): "pending" | "running" | "done" | "f...
std::function< void(int first, int last)> ParallelForBody
Body of a parallel_for subrange.
std::function< void()> JobFunc
A unit of work executed by a JobSystem worker.
WidgetDesc child(std::string id, std::vector< WidgetDesc > children, float width, float height)
Scrollable child region with an explicit size.
Build metadata (engine git commit, build time, third-party version).
SettlementPipeline::Stage fn
JobImpl * createJobLocked(JobScope scope, JobFunc body, JobFunc onComplete)
void enqueueLocked(JobImpl *job)
void scheduleLocked(JobImpl *job)
void waitJob(JobImpl *job)
void recordError(JobImpl *job, const std::string &message)
static void workerLoop(State *state)
std::condition_variable cv
void releaseDependents(JobImpl *job)
void runJob(JobImpl *job)
int autoChunk(int count) const
std::deque< JobImpl * > ready
void completeJob(JobImpl *job, bool ok)
void fireCompletion(JobImpl *job)