31 struct JobBarrierState;
66 void AddRef() const noexcept;
77 [[nodiscard]]
bool IsValid() const noexcept;
78 [[nodiscard]]
bool IsDone() const noexcept;
92 mPool(std::move(pool)), mState(std::move(state))
95 void Reset() noexcept;
107 [[nodiscard]]
bool IsEmpty() const noexcept;
117 template <
typename Fn>
127 if constexpr (std::is_invocable_v<Fn&, JobContext&>)
128 std::invoke(mFunction, context);
130 std::invoke(mFunction);
144 size_t mSubmitting{};
158 return static_cast<size_t>(priority);
161 bool BeginSubmit() noexcept;
162 void EndSubmit() noexcept;
164 uint32_t dependencyCount);
166 void ReleaseDependency(
JobHandle const& job, uint32_t count,
bool prerequisiteFailed);
169 bool TryExecuteOne(
size_t workerId);
170 void WorkerMain(
size_t workerId);
171 void ReleaseBarrier() noexcept;
172 void CancelInternal(
JobHandle const& job);
183 [[nodiscard]]
size_t GetWorkerCount() const noexcept {
return mThreads.size(); }
186 template <
typename Fn>
189 bool const submitting = BeginSubmit();
190 CHECK_MSG(submitting,
"JobSystem shutting down");
193 using Function = std::decay_t<Fn>;
194 auto callable = ConstructUniqueBase<JobCallable, LambdaJobCallable<Function>>(
195 mAllocator, Function(std::forward<Fn>(function)));
196 JobHandle job = CreateJobInternal(name, priority, std::move(callable), dependencyCount);
201 template <
typename Fn>
207 template <
typename Fn>
211 bool const submitting = BeginSubmit();
212 CHECK_MSG(submitting,
"JobSystem shutting down");
215 for (
JobHandle const& prerequisite : prerequisites)
216 CHECK_MSG(prerequisite.IsValid() && prerequisite.mPool.get() ==
mPool.get(),
217 "Invalid or foreign job prerequisite");
219 using Function = std::decay_t<Fn>;
220 auto callable = ConstructUniqueBase<JobCallable, LambdaJobCallable<Function>>(
221 mAllocator, Function(std::forward<Fn>(function)));
223 name, priority, std::move(callable),
static_cast<uint32_t
>(prerequisites.size() + 1));
224 for (
JobHandle const& prerequisite : prerequisites)
225 RegisterSuccessor(prerequisite, dependent);
226 dependent.RemoveDependency();
231 template <
typename Fn>
234 return CreateJobAfter(name,
JobPriority::Normal, prerequisites, std::forward<Fn>(function));
244 template <
typename Fn>
247 CHECK_MSG(step != 0,
"Job dispatch step size must be non-zero");
251 return CreateJob(name, priority, [] {});
253 using Function = std::decay_t<Fn>;
254 auto sharedFunction = ConstructShared<Function>(mAllocator, std::forward<Fn>(function));
255 size_t const chunkCount = 1 + (count - 1) / step;
257 chunks.reserve(chunkCount);
258 for (
size_t begin = 0; begin < count; begin += step)
260 size_t const end = std::min(begin + step, count);
261 chunks.push_back(CreateJob(
263 [sharedFunction, begin, end](
JobContext& context)
265 std::invoke(*sharedFunction, begin, end, context);
271 template <
typename Fn>
274 CHECK_MSG(step != 0,
"ParallelFor grain size must be non-zero");
276 return CreateBarrier();
278 return CreateBarrier();
282 barrier.
Add(completion);
#define CHECK_MSG(expr, format_str,...)
Allocator interface (noexcept)
Definition Allocator.hpp:29
Definition JobSystem.hpp:86
JobBarrier(SharedPtr< JobPool > pool, SharedPtr< JobBarrierState > state)
Definition JobSystem.hpp:91
SharedPtr< JobBarrierState > mState
Definition JobSystem.hpp:89
void Add(JobHandle const &job)
Definition JobSystem.cpp:330
SharedPtr< JobPool > mPool
Definition JobSystem.hpp:88
Definition JobSystem.hpp:111
virtual ~JobCallable()=default
virtual void Execute(JobContext &context)=0
Definition JobSystem.hpp:34
size_t GetWorkerId() const noexcept
Definition JobSystem.hpp:45
JobContext(size_t workerId, Atomic< bool > const *cancellation)
Definition JobSystem.hpp:39
bool IsCancellationRequested() const noexcept
Definition JobSystem.hpp:46
size_t mWorkerId
Definition JobSystem.hpp:36
Atomic< bool > const * mCancellation
Definition JobSystem.hpp:37
Definition JobSystem.hpp:53
bool IsValid() const noexcept
Definition JobSystem.cpp:269
void RemoveDependency(uint32_t count=1) const
Definition JobSystem.cpp:294
uint32_t mIndex
Definition JobSystem.hpp:59
JobStatus Status() const noexcept
Definition JobSystem.cpp:271
uint32_t mGeneration
Definition JobSystem.hpp:60
void AddRef() const noexcept
Definition JobSystem.cpp:254
friend class JobDependency
Definition JobSystem.hpp:56
bool IsDone() const noexcept
Definition JobSystem.cpp:278
void AddDependency(uint32_t count=1) const
Definition JobSystem.cpp:280
SharedPtr< JobPool > mPool
Definition JobSystem.hpp:58
void Release() noexcept
Definition JobSystem.cpp:260
Definition JobSystem.hpp:135
CondVar mSubmitCV
Definition JobSystem.hpp:143
String mName
Definition JobSystem.hpp:137
Allocator * mAllocator
Definition JobSystem.hpp:136
Vector< Thread > mThreads
Definition JobSystem.hpp:154
SharedPtr< JobPool > mPool
Definition JobSystem.hpp:139
void Wait(JobBarrier &&barrier)
Definition JobSystem.hpp:239
Mutex mCallerMutex
Definition JobSystem.hpp:149
size_t GetMaxConcurrency() const noexcept
Definition JobSystem.hpp:184
JobHandle CreateJob(StringView name, JobPriority priority, Fn &&function, uint32_t dependencyCount=0)
Definition JobSystem.hpp:187
JobHandle CreateJob(StringView name, Fn &&function, uint32_t dependencyCount=0)
Definition JobSystem.hpp:202
ReadyQueues mReady
Definition JobSystem.hpp:153
JobHandle CreateJobAfter(StringView name, JobPriority priority, Span< const JobHandle > prerequisites, Fn &&function)
Definition JobSystem.hpp:208
size_t mMaxBarriers
Definition JobSystem.hpp:138
JobHandle CreateJobAfter(StringView name, Span< const JobHandle > prerequisites, Fn &&function)
Definition JobSystem.hpp:232
Mutex mSubmitMutex
Definition JobSystem.hpp:142
JobHandle Dispatch(StringView name, size_t count, size_t step, JobPriority priority, Fn &&function)
Definition JobSystem.hpp:245
JobBarrier ParallelFor(StringView name, size_t count, size_t step, Fn &&function)
Definition JobSystem.hpp:272
static constexpr size_t PriorityIndex(JobPriority priority) noexcept
Definition JobSystem.hpp:156
Array< ReadyQueue, kJobPriorityCount > ReadyQueues
Definition JobSystem.hpp:152
Definition JobSystem.hpp:119
Fn mFunction
Definition JobSystem.hpp:120
LambdaJobCallable(Fn &&function)
Definition JobSystem.hpp:123
void Execute(JobContext &context) override
Definition JobSystem.hpp:125
Atomic, bounded multi-producer multi-consumer FIFO ring buffer with a fixed maximum size.
Definition AtomicQueue.hpp:85
Lock-free atomic primitives and implementations of data structures.
Definition Allocator.hpp:6
std::vector< T, StlAllocator< T > > Vector
std::vector with explicit Foundation::Core::StlAllocator constructor
Definition Container.hpp:149
std::mutex Mutex
Definition Thread.hpp:10
JobPriority
Definition ThreadPool.hpp:13
JobStatus
Definition JobSystem.hpp:8
std::queue< T, Container > Queue
std::queue with explicit Foundation::Core::StlAllocator constructor
Definition Container.hpp:240
std::shared_ptr< T > SharedPtr
std::shared_ptr with custom deleter that uses a Foundation::Core::Allocator to deallocate memory.
Definition Allocator.hpp:209
std::condition_variable CondVar
Definition Thread.hpp:9
std::atomic< T > Atomic
Alias of std::atomic<T>.
Definition Atomic.hpp:26
std::basic_string< char, std::char_traits< char >, StlDefaultAllocator< char > > String
Alias for std::basic_string<char>, without an explicit allocator constructor.
Definition Container.hpp:120
std::array< T, Size > Array
Alias for std::array
Definition Container.hpp:45
std::basic_string_view< char > StringView
Alias for std::basic_string_view<char>
Definition Container.hpp:56
std::unique_ptr< T, Deleter > UniquePtr
std::unique_ptr with custom deleter that uses a Foundation::Core::Allocator to deallocate memory.
Definition Allocator.hpp:180
std::span< T > Span
Alias for std::span
Definition Container.hpp:62
Definition JobSystem.hpp:63
Definition JobSystem.cpp:93
Definition JobSystem.hpp:18
size_t workerCount
Definition JobSystem.hpp:19
size_t maxBarriers
Definition JobSystem.hpp:21
size_t maxJobs
Definition JobSystem.hpp:20
StringView name
Definition JobSystem.hpp:24
size_t readyQueueSize
Definition JobSystem.hpp:22
Allocator * allocator
Definition JobSystem.hpp:23