载入中...
搜索中...
未找到
JobSystemThreadPool.cpp
浏览该文件的文档.
2
3#include "common/Exception.h"
4#include "thread/ThreadPool.h"
5
6#include <algorithm>
7#include <condition_variable>
8#include <cstddef>
9#include <deque>
10#include <mutex>
11#include <new>
12#include <string>
13#include <thread>
14#include <utility>
15#include <vector>
16
17namespace eve {
18namespace thread {
19
20namespace {
21
22enum class JobStatus { Pending, Scheduled, Running, Done, Failed };
23enum class JobScope { Heap, Frame };
24
25struct JobImpl;
26
27thread_local const void *tlsCurrentState = nullptr;
28thread_local JobImpl *tlsCurrentJob = nullptr;
29
30bool isDoneStatus(JobStatus status) {
31 return status == JobStatus::Done || status == JobStatus::Failed;
32}
33
34} // namespace
35
36namespace {
37
43class FrameArena {
44public:
45 FrameArena() = default;
46 ~FrameArena() { reset(); }
47
54 void *alloc(size_t size, size_t align) {
55 if (size > kArenaBlockSize)
56 throw eve::Exception("JobSystem arena allocation exceeds block size");
57 if (blocks_.empty())
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);
62 aligned = 0;
63 }
64 void *ptr = blocks_.back().data() + aligned;
65 offset_ = aligned + size;
66 return ptr;
67 }
68
74 void track(void *ptr, void (*destroy)(void *)) {
75 tracked_.push_back({ptr, destroy});
76 }
77
79 void reset() {
80 for (auto &entry : tracked_)
81 entry.destroy(entry.ptr);
82 tracked_.clear();
83 offset_ = 0;
84 }
85
86private:
87 static constexpr size_t kArenaBlockSize = 64 * 1024;
88
89 struct Tracked {
90 void *ptr;
91 void (*destroy)(void *);
92 };
93
94 std::vector<std::vector<std::byte>> blocks_;
95 std::vector<Tracked> tracked_;
96 size_t offset_ = 0;
97};
98
99} // namespace
100
102 mutable std::mutex mu;
103 std::condition_variable cv;
104 std::deque<JobImpl *> ready;
105 int outstanding = 0;
107 // Number of jobs currently between status=Done/Failed and the tail of
108 // completeJob (fireCompletion / releaseDependents / counter decrement).
109 // waitFrameJobs() must not reset the frame arena while this is > 0:
110 // a child releases its join before its own counters are decremented, so
111 // another thread can finish the join, drive outstandingFrame to 0 and
112 // reset the arena while this worker is still touching the child (a
113 // use-after-free that shows up as a null std::function call in
114 // fireCompletion under concurrent GPU work).
115 int completing = 0;
116 int workerCount = 0;
117 bool stopping = false;
118 FrameArena arena;
119
120 JobImpl *createJobLocked(JobScope scope, JobFunc body, JobFunc onComplete);
121 void scheduleLocked(JobImpl *job);
122 void enqueueLocked(JobImpl *job);
123 void runJob(JobImpl *job);
124 void completeJob(JobImpl *job, bool ok);
125 void fireCompletion(JobImpl *job);
126 void recordError(JobImpl *job, const std::string &message);
127 void releaseDependents(JobImpl *job);
128 void waitJob(JobImpl *job);
129 void waitFrameJobs();
130 int autoChunk(int count) const;
131
132 static void workerLoop(State *state);
133};
134
135namespace {
136
137struct JobImpl final : public Job {
138 JobImpl(JobSystemThreadPool::State *owner, JobScope scope_, JobFunc fn, JobFunc onDone)
139 : state(owner), scope(scope_), body(std::move(fn)), onComplete(std::move(onDone)) {}
140
141 ~JobImpl() override {
142 // parallel_for children are owned by their join job and are guaranteed
143 // complete before the join finishes, so deleting them here is safe.
144 // Frame-scope jobs never populate ownedChildren (the arena owns them).
145 for (JobImpl *child : ownedChildren)
146 delete child;
147 }
148
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;
154 void addDependency(Job *predecessor) override;
155 void setCompletionCallback(JobFunc callback) override;
156
157 JobSystemThreadPool::State *state;
158 JobScope scope;
161 int depCount = 0;
162 std::vector<JobImpl *> dependents;
163 std::vector<JobImpl *> ownedChildren;
164 JobStatus status = JobStatus::Pending;
165 bool scheduled = false;
166 bool enqueued = false;
167 // Set under mu at the very end of completeJob. waitJob() waits on this
168 // (not just status) so a caller that deletes a join job never frees a
169 // child while a worker is still in the child's completeJob tail.
170 bool completionDone = false;
171 std::string error;
172};
173
174class TaskGroupImpl final : public TaskGroup {
175public:
176 TaskGroupImpl(JobSystemThreadPool::State *owner, JobScope scope_)
177 : state(owner), scope(scope_) {}
178 ~TaskGroupImpl() override = default;
179
180 Job *fork(JobFunc body) override { return fork(std::move(body), JobFunc{}); }
181
182 Job *fork(JobFunc body, JobFunc onComplete) override {
183 if (!body)
184 throw eve::Exception("TaskGroup::fork: null function");
185 JobImpl *job = nullptr;
186 {
187 std::lock_guard<std::mutex> lock(state->mu);
188 if (state->stopping)
189 throw eve::Exception("JobSystem is stopped");
190 job = state->createJobLocked(scope, std::move(body), std::move(onComplete));
191 state->scheduleLocked(job);
192 children.push_back(job);
193 }
194 state->cv.notify_one();
195 return job;
196 }
197
198 void wait() override {
199 std::vector<JobImpl *> snapshot;
200 {
201 std::lock_guard<std::mutex> lock(state->mu);
202 snapshot = children;
203 }
204 for (JobImpl *child : snapshot)
205 state->waitJob(child);
206 std::lock_guard<std::mutex> lock(state->mu);
207 children.clear();
208 }
209
210 int getPendingCount() const override {
211 std::lock_guard<std::mutex> lock(state->mu);
212 int pending = 0;
213 for (JobImpl *child : children) {
214 if (!isDoneStatus(child->status))
215 ++pending;
216 }
217 return pending;
218 }
219
220 JobSystemThreadPool::State *state;
221 JobScope scope;
222 std::vector<JobImpl *> children;
223};
224
225} // namespace
226
227// ---- State ----
228
231 if (scope == JobScope::Frame) {
232 void *mem = arena.alloc(sizeof(JobImpl), alignof(JobImpl));
233 auto *job = new (mem) JobImpl(this, scope, std::move(body), std::move(onComplete));
234 arena.track(job, [](void *ptr) { static_cast<JobImpl *>(ptr)->~JobImpl(); });
235 return job;
236 }
237 return new JobImpl(this, scope, std::move(body), std::move(onComplete));
238}
239
241 job->scheduled = true;
242 ++outstanding;
243 if (job->scope == JobScope::Frame)
244 ++outstandingFrame;
245 if (job->depCount == 0)
246 enqueueLocked(job);
247}
248
250 if (job->enqueued)
251 return;
252 job->enqueued = true;
253 ready.push_back(job);
254}
255
257 {
258 std::lock_guard<std::mutex> lock(mu);
259 job->status = JobStatus::Running;
260 }
261
262 // A job body that calls waitAll()/stop() must be treated as a worker even
263 // when it is help-executed inline by a waiting thread (fork/join path),
264 // otherwise those calls would re-enter the wait they are running inside.
265 const void *previousState = tlsCurrentState;
266 if (previousState != this)
267 tlsCurrentState = this;
268
269 JobImpl *previous = tlsCurrentJob;
270 tlsCurrentJob = job;
271 bool ok = true;
272 if (job->body) {
273 try {
274 job->body();
275 } catch (const eve::Exception &e) {
276 ok = false;
277 recordError(job, e.what());
278 } catch (const std::exception &e) {
279 ok = false;
280 recordError(job, e.what());
281 } catch (...) {
282 ok = false;
283 recordError(job, "unknown exception");
284 }
285 }
286 tlsCurrentJob = previous;
287
288 completeJob(job, ok);
289 tlsCurrentState = previousState;
290}
291
293 {
294 std::lock_guard<std::mutex> lock(mu);
295 job->status = ok ? JobStatus::Done : JobStatus::Failed;
296 ++completing;
297 }
298 cv.notify_all();
299
300 // Completion callback runs before dependents are released so it can
301 // publish results downstream jobs consume.
302 fireCompletion(job);
303
304 // Decrement the outstanding counters while the job is still pinned by
305 // `completing` (waitFrameJobs waits on it too), release dependents, then
306 // mark the job's own memory as finished. All of this stays inside one
307 // critical section: a waiter can only return once completionDone is set,
308 // so it cannot delete the job (heap jobs are owned by the caller) while
309 // this worker is still iterating job->dependents. The dependents are
310 // released before completionDone so a join can never become runnable (and
311 // later be deleted) while this job is still writing its own fields.
312 {
313 std::lock_guard<std::mutex> lock(mu);
314 --outstanding;
315 if (job->scope == JobScope::Frame)
316 --outstandingFrame;
317 releaseDependents(job);
318 job->completionDone = true;
319 }
320 {
321 std::lock_guard<std::mutex> lock(mu);
322 --completing;
323 }
324 cv.notify_all();
325}
326
328 if (!job->onComplete)
329 return;
330 try {
331 job->onComplete();
332 } catch (const eve::Exception &e) {
333 recordError(job, std::string("completion callback: ") + e.what());
334 } catch (const std::exception &e) {
335 recordError(job, std::string("completion callback: ") + e.what());
336 } catch (...) {
337 recordError(job, "completion callback: unknown exception");
338 }
339}
340
341void JobSystemThreadPool::State::recordError(JobImpl *job, const std::string &message) {
342 std::lock_guard<std::mutex> lock(mu);
343 if (!job->error.empty())
344 job->error += "; ";
345 job->error += message;
346}
347
349 // Caller holds mu. (Kept as a helper so completeJob reads as a single
350 // critical section; the enqueue + depCount updates below are lock-free
351 // only because they are performed under that same lock.)
352 bool woke = false;
353 for (JobImpl *dep : job->dependents) {
354 if (dep->depCount > 0)
355 --dep->depCount;
356 if (dep->depCount == 0 && dep->scheduled && !dep->enqueued) {
357 enqueueLocked(dep);
358 woke = true;
359 }
360 }
361 job->dependents.clear();
362 if (woke)
363 cv.notify_all();
364}
365
367 std::unique_lock<std::mutex> lock(mu);
368 if (job->completionDone)
369 return;
370 if (tlsCurrentJob == job)
371 throw eve::Exception("JobSystem: a job cannot wait on itself");
372
373 for (;;) {
374 if (job->completionDone)
375 return;
376 if (!ready.empty()) {
377 JobImpl *next = ready.front();
378 ready.pop_front();
379 lock.unlock();
380 runJob(next);
381 lock.lock();
382 continue;
383 }
384 if (stopping)
385 return;
386 cv.wait(lock);
387 }
388}
389
391 std::unique_lock<std::mutex> lock(mu);
392 while (outstandingFrame > 0 || completing > 0) {
393 if (!ready.empty()) {
394 JobImpl *next = ready.front();
395 ready.pop_front();
396 lock.unlock();
397 runJob(next);
398 lock.lock();
399 continue;
400 }
401 if (stopping)
402 return;
403 cv.wait(lock);
404 }
405 arena.reset();
406}
407
409 const int target = std::max(1, workerCount * 4);
410 return std::max(1, (count + target - 1) / target);
411}
412
414 tlsCurrentState = state;
415 for (;;) {
416 JobImpl *job = nullptr;
417 {
418 std::unique_lock<std::mutex> lock(state->mu);
419 state->cv.wait(lock, [state] { return state->stopping || !state->ready.empty(); });
420 if (state->stopping && state->ready.empty())
421 break;
422 job = state->ready.front();
423 state->ready.pop_front();
424 }
425 if (job)
426 state->runJob(job);
427 }
428 tlsCurrentState = nullptr;
429}
430
431// ---- JobImpl ----
432
433void JobImpl::wait() {
434 state->waitJob(this);
435}
436
437bool JobImpl::isDone() const {
438 std::lock_guard<std::mutex> lock(state->mu);
439 return completionDone;
440}
441
442bool JobImpl::hasFailed() const {
443 std::lock_guard<std::mutex> lock(state->mu);
444 return status == JobStatus::Failed;
445}
446
447std::string JobImpl::getError() const {
448 std::lock_guard<std::mutex> lock(state->mu);
449 return error;
450}
451
452int JobImpl::getPendingDependencyCount() const {
453 std::lock_guard<std::mutex> lock(state->mu);
454 return depCount;
455}
456
457void JobImpl::addDependency(Job *predecessor) {
458 auto *pred = dynamic_cast<JobImpl *>(predecessor);
459 if (!pred || pred->state != state)
460 throw eve::Exception("JobSystem::addDependency: job does not belong to this system");
461 if (pred == this)
462 throw eve::Exception("JobSystem::addDependency: a job cannot depend on itself");
463
464 std::lock_guard<std::mutex> lock(state->mu);
465 if (scheduled)
466 throw eve::Exception("JobSystem::addDependency: job already scheduled");
467 if (isDoneStatus(pred->status))
468 return; // Predecessor already finished; nothing to wait for.
469 ++depCount;
470 pred->dependents.push_back(this);
471}
472
473void JobImpl::setCompletionCallback(JobFunc callback) {
474 bool alreadyDone = false;
475 {
476 std::lock_guard<std::mutex> lock(state->mu);
477 if (isDoneStatus(status)) {
478 alreadyDone = true;
479 } else {
480 onComplete = std::move(callback);
481 return;
482 }
483 }
484 if (alreadyDone && callback) {
485 try {
486 callback();
487 } catch (const eve::Exception &e) {
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());
491 } catch (...) {
492 state->recordError(this, "completion callback: unknown exception");
493 }
494 }
495}
496
497// ---- JobSystemThreadPool ----
498
499namespace {
500
501JobImpl *castJob(JobSystemThreadPool::State *state, Job *job) {
502 if (!job)
503 throw eve::Exception("JobSystem: null job");
504 auto *impl = dynamic_cast<JobImpl *>(job);
505 if (!impl || impl->state != state)
506 throw eve::Exception("JobSystem: job does not belong to this system");
507 return impl;
508}
509
510} // namespace
511
513 if (workerCount <= 0) {
514 workerCount = static_cast<int>(std::thread::hardware_concurrency());
515 if (workerCount <= 0)
516 workerCount = 1;
517 }
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));
522 try {
523 for (int i = 0; i < state_->workerCount; ++i) {
524 auto state = state_;
525 poolTasks_.push_back(pool_->submit([state] { State::workerLoop(state.get()); }));
526 }
527 } catch (...) {
528 {
529 std::lock_guard<std::mutex> lock(state_->mu);
530 state_->stopping = true;
531 }
532 state_->cv.notify_all();
533 if (pool_)
534 pool_->stop();
535 for (Task *task : poolTasks_)
536 delete task;
537 poolTasks_.clear();
538 throw;
539 }
540}
541
543 try {
544 std::lock_guard<std::mutex> lifecycleLock(lifecycleMu_);
545 const bool fromWorker = tlsCurrentState == state_.get();
546 {
547 std::lock_guard<std::mutex> lock(state_->mu);
548 state_->stopping = true;
549 }
550 state_->cv.notify_all();
551
552 if (fromWorker) {
553 // Destroyed from inside one of our jobs: joining any worker could
554 // deadlock (another worker may wait on this task). Mirror
555 // ThreadPool's worker-caller path — detach; worker loops only
556 // touch shared State, which outlives them through their capture.
557 for (Task *task : poolTasks_)
558 delete task;
559 poolTasks_.clear();
560 pool_.reset();
561 return;
562 }
563 if (pool_)
564 pool_->stop();
565 for (Task *task : poolTasks_)
566 delete task;
567 poolTasks_.clear();
568 } catch (...) {
569 }
570}
571
573 return state_->workerCount;
574}
575
577 std::lock_guard<std::mutex> lock(state_->mu);
578 return !state_->stopping;
579}
580
582 std::lock_guard<std::mutex> lock(state_->mu);
583 return static_cast<int>(state_->ready.size());
584}
585
587 std::lock_guard<std::mutex> lock(state_->mu);
588 return state_->outstanding;
589}
590
592 JobImpl *job = nullptr;
593 {
594 std::lock_guard<std::mutex> lock(state_->mu);
595 if (state_->stopping)
596 throw eve::Exception("JobSystem is stopped");
597 job = state_->createJobLocked(JobScope::Heap, std::move(body), JobFunc{});
598 state_->scheduleLocked(job);
599 }
600 state_->cv.notify_one();
601 return job;
602}
603
605 std::lock_guard<std::mutex> lock(state_->mu);
606 return state_->createJobLocked(JobScope::Heap, std::move(body), JobFunc{});
607}
608
610 JobImpl *impl = castJob(state_.get(), job);
611 {
612 std::lock_guard<std::mutex> lock(state_->mu);
613 if (state_->stopping)
614 throw eve::Exception("JobSystem is stopped");
615 if (impl->scheduled)
616 throw eve::Exception("JobSystem::schedule: job already scheduled");
617 state_->scheduleLocked(impl);
618 }
619 state_->cv.notify_one();
620}
621
623 return parallelForImpl(first, last, std::move(body), chunk, false);
624}
625
627 return new TaskGroupImpl(state_.get(), JobScope::Heap);
628}
629
631 JobImpl *job = nullptr;
632 {
633 std::lock_guard<std::mutex> lock(state_->mu);
634 if (state_->stopping)
635 throw eve::Exception("JobSystem is stopped");
636 job = state_->createJobLocked(JobScope::Frame, std::move(body), JobFunc{});
637 state_->scheduleLocked(job);
638 }
639 state_->cv.notify_one();
640 return job;
641}
642
644 std::lock_guard<std::mutex> lock(state_->mu);
645 return state_->createJobLocked(JobScope::Frame, std::move(body), JobFunc{});
646}
647
649 return parallelForImpl(first, last, std::move(body), chunk, true);
650}
651
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(); });
657 return group;
658}
659
661 state_->waitFrameJobs();
662}
663
665 state_->waitFrameJobs();
666}
667
669 if (tlsCurrentState == state_.get())
670 throw eve::Exception("JobSystem::waitAll cannot be called from its worker");
671 State *state = state_.get();
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();
677 lock.unlock();
678 state->runJob(next);
679 lock.lock();
680 continue;
681 }
682 if (state->stopping)
683 return;
684 state->cv.wait(lock);
685 }
686}
687
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");
692 {
693 std::lock_guard<std::mutex> lock(state_->mu);
694 state_->stopping = true;
695 }
696 state_->cv.notify_all();
697 if (pool_)
698 pool_->stop();
699 for (Task *task : poolTasks_)
700 delete task;
701 poolTasks_.clear();
702}
703
704Job *JobSystemThreadPool::parallelForImpl(int first, int last, ParallelForBody body, int chunk,
705 bool frameScope) {
706 if (!body)
707 throw eve::Exception("JobSystem::parallelFor: null function");
708 if (last <= first)
709 return frameScope ? submitFrame(JobFunc{}) : submit(JobFunc{});
710
711 const JobScope scope = frameScope ? JobScope::Frame : JobScope::Heap;
712 const int count = last - first;
713 if (chunk <= 0)
714 chunk = state_->autoChunk(count);
715 const int taskCount = (count + chunk - 1) / chunk;
716
717 std::lock_guard<std::mutex> lock(state_->mu);
718 if (state_->stopping)
719 throw eve::Exception("JobSystem is stopped");
720
721 JobImpl *join = state_->createJobLocked(scope, JobFunc{}, JobFunc{});
722 join->depCount = taskCount;
723 if (scope == JobScope::Heap)
724 join->ownedChildren.reserve(static_cast<size_t>(taskCount));
725
726 std::vector<JobImpl *> children;
727 children.reserve(static_cast<size_t>(taskCount));
728 for (int t = 0; t < taskCount; ++t) {
729 const int begin = first + t * chunk;
730 const int end = std::min(last, begin + chunk);
731 JobImpl *child =
732 state_->createJobLocked(scope, [body, begin, end] { body(begin, end); }, JobFunc{});
733 child->dependents.push_back(join);
734 children.push_back(child);
735 if (scope == JobScope::Heap)
736 join->ownedChildren.push_back(child);
737 }
738
739 // Schedule the join first so it sits waiting, then the children; the join
740 // becomes ready when the last child completes.
741 state_->scheduleLocked(join);
742 for (JobImpl *child : children)
743 state_->scheduleLocked(child);
744 state_->cv.notify_all();
745 return join;
746}
747
748} // namespace thread
749} // namespace eve
LogicalId target
std::vector< QuestEvent > pending
building::EdgeCurveGroup group
void * impl
std::string message
wgpu::PopErrorScopeStatus status
std::int32_t first
bool enqueued
std::vector< JobImpl * > dependents
JobFunc onComplete
std::vector< JobImpl * > ownedChildren
bool completionDone
JobScope scope
void(* destroy)(void *)
void * ptr
bool scheduled
graphics::Canvas * previous
std::string error
Definition Package.cpp:60
float begin
float t
std::uint32_t count
std::map< Cell, Cell > predecessor
int children
Definition TreeMesh.cpp:295
float size
Definition TreeMesh.cpp:156
std::string body
ViewPreparation callback
EVENGINE_API_FOUNDATION public API.
Definition Exception.h:13
virtual const char * what() const
Returns a string containing reason for the exception.
Definition Exception.h:24
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.
int getOutstandingCount() const override
Returns the outstanding count.
int getWorkerCount() const override
Returns the worker count.
Job * submit(JobFunc body) override
Submit.
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.
Definition JobSystem.h:43
Fork/join group: fork() spawns scheduled children, wait() joins them.
Definition JobSystem.h:109
A job executed by a ThreadPool worker. Status strings (no enums): "pending" | "running" | "done" | "f...
Definition Task.h:18
std::function< void(int first, int last)> ParallelForBody
Body of a parallel_for subrange.
Definition JobSystem.h:24
std::function< void()> JobFunc
A unit of work executed by a JobSystem worker.
Definition JobSystem.h:17
WidgetDesc child(std::string id, std::vector< WidgetDesc > children, float width, float height)
Scrollable child region with an explicit size.
Definition Widget.cpp:635
Build metadata (engine git commit, build time, third-party version).
Definition Build.cpp:16
SettlementPipeline::Stage fn
JobImpl * createJobLocked(JobScope scope, JobFunc body, JobFunc onComplete)
void recordError(JobImpl *job, const std::string &message)
std::mutex mu
Definition Graphics.cpp:181