Foundation
Loading...
Searching...
No Matches
JobSystem.hpp
Go to the documentation of this file.
1#pragma once
2#include "ThreadPool.hpp"
3#include <functional>
4
5namespace Foundation::Core
6{
7 enum class JobStatus : uint8_t
8 {
10 Queued,
11 Running,
13 Failed,
15 };
16
18 {
19 size_t workerCount{};
20 size_t maxJobs{1024};
21 size_t maxBarriers{16};
22 size_t readyQueueSize{1024};
25 };
26
27 class JobSystem;
28 class JobHandle;
29 class JobDependency;
30 struct JobPool;
31 struct JobBarrierState;
32
34 {
35 friend class JobSystem;
36 size_t mWorkerId{};
38
39 JobContext(size_t workerId, Atomic<bool> const* cancellation) :
40 mWorkerId(workerId), mCancellation(cancellation)
41 {
42 }
43
44 public:
45 [[nodiscard]] size_t GetWorkerId() const noexcept { return mWorkerId; }
46 [[nodiscard]] bool IsCancellationRequested() const noexcept
47 {
48 return mCancellation && mCancellation->load(std::memory_order_acquire);
49 }
50 };
51
53 {
54 friend class JobSystem;
55 friend class JobBarrier;
56 friend class JobDependency;
57 friend struct JobPool;
59 uint32_t mIndex{UINT32_MAX};
60 uint32_t mGeneration{};
61
62 struct AdoptRef
63 {
64 };
65 JobHandle(SharedPtr<JobPool> pool, uint32_t index, uint32_t generation, AdoptRef) noexcept;
66 void AddRef() const noexcept;
67 void Release() noexcept;
68
69 public:
70 JobHandle() = default;
71 JobHandle(JobHandle const& other) noexcept;
72 JobHandle& operator=(JobHandle const& other) noexcept;
73 JobHandle(JobHandle&& other) noexcept;
74 JobHandle& operator=(JobHandle&& other) noexcept;
75 ~JobHandle();
76
77 [[nodiscard]] bool IsValid() const noexcept;
78 [[nodiscard]] bool IsDone() const noexcept;
79 [[nodiscard]] JobStatus Status() const noexcept;
80 void AddDependency(uint32_t count = 1) const;
81 void RemoveDependency(uint32_t count = 1) const;
82 };
83
84
86 {
87 friend class JobSystem;
90
92 mPool(std::move(pool)), mState(std::move(state))
93 {
94 }
95 void Reset() noexcept;
96
97 public:
98 JobBarrier() = default;
99 JobBarrier(JobBarrier const&) = delete;
100 JobBarrier& operator=(JobBarrier const&) = delete;
101 JobBarrier(JobBarrier&& other) noexcept;
102 JobBarrier& operator=(JobBarrier&& other) noexcept;
103 ~JobBarrier();
104
105 void Add(JobHandle const& job);
106 void Add(Span<const JobHandle> jobs);
107 [[nodiscard]] bool IsEmpty() const noexcept;
108 };
109
111 {
112 public:
113 virtual ~JobCallable() = default;
114 virtual void Execute(JobContext& context) = 0;
115 };
116
117 template <typename Fn>
118 class LambdaJobCallable final : public JobCallable
119 {
121
122 public:
123 explicit LambdaJobCallable(Fn&& function) : mFunction(std::move(function)) {}
124
125 void Execute(JobContext& context) override
126 {
127 if constexpr (std::is_invocable_v<Fn&, JobContext&>)
128 std::invoke(mFunction, context);
129 else
130 std::invoke(mFunction);
131 }
132 };
133
135 {
140
141 Atomic<bool> mAccepting{true};
144 size_t mSubmitting{};
145 Atomic<bool> mStopping{false};
146 Atomic<bool> mStopped{false};
147 Atomic<size_t> mWakeEpoch{0};
148 Atomic<size_t> mOutstanding{0};
150
155
156 static constexpr size_t PriorityIndex(JobPriority priority) noexcept
157 {
158 return static_cast<size_t>(priority);
159 }
160
161 bool BeginSubmit() noexcept;
162 void EndSubmit() noexcept;
163 JobHandle CreateJobInternal(StringView name, JobPriority priority, UniquePtr<JobCallable> callable,
164 uint32_t dependencyCount);
165 void Queue(JobHandle const& job) noexcept;
166 void ReleaseDependency(JobHandle const& job, uint32_t count, bool prerequisiteFailed);
167 void Finish(JobHandle const& job, JobStatus status) noexcept;
168 void RegisterSuccessor(JobHandle const& prerequisite, JobHandle const& dependent);
169 bool TryExecuteOne(size_t workerId);
170 void WorkerMain(size_t workerId);
171 void ReleaseBarrier() noexcept;
172 void CancelInternal(JobHandle const& job);
173
174 friend class JobHandle;
175 friend class JobBarrier;
176 friend class JobDependency;
177 public:
178 explicit JobSystem(JobSystemDesc const& desc);
179 JobSystem(JobSystem const&) = delete;
180 JobSystem& operator=(JobSystem const&) = delete;
181 ~JobSystem();
182
183 [[nodiscard]] size_t GetWorkerCount() const noexcept { return mThreads.size(); }
184 [[nodiscard]] size_t GetMaxConcurrency() const noexcept { return mThreads.size() + 1; }
185
186 template <typename Fn>
187 JobHandle CreateJob(StringView name, JobPriority priority, Fn&& function, uint32_t dependencyCount = 0)
188 {
189 bool const submitting = BeginSubmit();
190 CHECK_MSG(submitting, "JobSystem shutting down");
191 if (!submitting)
192 return {};
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);
197 EndSubmit();
198 return job;
199 }
200
201 template <typename Fn>
202 JobHandle CreateJob(StringView name, Fn&& function, uint32_t dependencyCount = 0)
203 {
204 return CreateJob(name, JobPriority::Normal, std::forward<Fn>(function), dependencyCount);
205 }
206
207 template <typename Fn>
209 Fn&& function)
210 {
211 bool const submitting = BeginSubmit();
212 CHECK_MSG(submitting, "JobSystem shutting down");
213 if (!submitting)
214 return {};
215 for (JobHandle const& prerequisite : prerequisites)
216 CHECK_MSG(prerequisite.IsValid() && prerequisite.mPool.get() == mPool.get(),
217 "Invalid or foreign job prerequisite");
218
219 using Function = std::decay_t<Fn>;
220 auto callable = ConstructUniqueBase<JobCallable, LambdaJobCallable<Function>>(
221 mAllocator, Function(std::forward<Fn>(function)));
222 JobHandle dependent = CreateJobInternal(
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();
227 EndSubmit();
228 return dependent;
229 }
230
231 template <typename Fn>
232 JobHandle CreateJobAfter(StringView name, Span<const JobHandle> prerequisites, Fn&& function)
233 {
234 return CreateJobAfter(name, JobPriority::Normal, prerequisites, std::forward<Fn>(function));
235 }
236
237 JobBarrier CreateBarrier();
238 void Wait(JobBarrier& barrier);
239 void Wait(JobBarrier&& barrier) { Wait(barrier); }
240 void Cancel(JobHandle const& job);
241 void Join();
242 void Shutdown();
243
244 template <typename Fn>
245 JobHandle Dispatch(StringView name, size_t count, size_t step, JobPriority priority, Fn&& function)
246 {
247 CHECK_MSG(step != 0, "Job dispatch step size must be non-zero");
248 if (step == 0)
249 return {};
250 if (count == 0)
251 return CreateJob(name, priority, [] {});
252
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;
256 Vector<JobHandle> chunks(mAllocator);
257 chunks.reserve(chunkCount);
258 for (size_t begin = 0; begin < count; begin += step)
259 {
260 size_t const end = std::min(begin + step, count);
261 chunks.push_back(CreateJob(
262 name, priority,
263 [sharedFunction, begin, end](JobContext& context)
264 {
265 std::invoke(*sharedFunction, begin, end, context);
266 }));
267 }
268 return CreateJobAfter(name, priority, Span<const JobHandle>{chunks.data(), chunks.size()}, [] {});
269 }
270
271 template <typename Fn>
272 JobBarrier ParallelFor(StringView name, size_t count, size_t step, Fn&& function)
273 {
274 CHECK_MSG(step != 0, "ParallelFor grain size must be non-zero");
275 if (step == 0)
276 return CreateBarrier();
277 if (count == 0)
278 return CreateBarrier();
279 JobBarrier barrier = CreateBarrier();
280 JobHandle completion =
281 Dispatch(name, count, step, JobPriority::Normal, std::forward<Fn>(function));
282 barrier.Add(completion);
283 return barrier;
284 }
285 };
286} // namespace Foundation::Core
#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