89 friend class JobManager;
95 static constexpr uint32_t MaxWorkers = JobState::DEP_BITS;
98 std::thread::id m_mainThreadId;
101 std::atomic_bool m_stop{};
107 MpmcQueue<JobHandle, 1024> m_jobQueue[JobPriorityCnt];
109 MpmcQueue<JobHandle, 1024> m_jobQueueBackground;
111 uint32_t m_frameWorkersCnt = 0;
113 uint32_t m_backgroundWorkersCnt = 0;
115 uint32_t m_workersCnt[JobPriorityCnt]{};
120 uint32_t m_workerThreadsCnt[JobPriorityCnt]{};
127 std::atomic_uint32_t m_blockedInWorkUntil;
130 JobManager m_jobManager;
136 GAIA_PROF_MUTEX(
SpinLock, m_jobAllocMtx);
144 const auto hwThreads = hw_thread_cnt();
145 const auto hwEffThreads = hw_efficiency_cores_cnt();
146 uint32_t hiPrioWorkers = hwThreads;
147 if (hwEffThreads < hwThreads)
148 hiPrioWorkers -= hwEffThreads;
150 set_max_workers(hwThreads, hiPrioWorkers);
172 m_mainThreadId = std::this_thread::get_id();
173 if (!m_workersCtx.
empty())
174 detail::tl_workerCtx = &m_workersCtx[0];
180 return m_frameWorkersCnt;
186 return m_backgroundWorkersCnt;
196 const auto maxFrameWorkers = MaxWorkers - m_backgroundWorkersCnt;
197 const auto workersCnt = core::get_max(core::get_min(maxFrameWorkers, count), 1U);
198 countHighPrio = core::get_min(countHighPrio, workersCnt);
204 for (
auto& ctx: m_workersCtx)
207 m_frameWorkersCnt = workersCnt - 1;
211 m_workersCtx.
resize(workersCnt + m_backgroundWorkersCnt);
213 m_workers.
resize(m_frameWorkersCnt + m_backgroundWorkersCnt);
218 detail::tl_workerCtx = m_workersCtx.
data();
219 m_workersCtx[0].tp =
this;
220 m_workersCtx[0].workerIdx = 0;
221 m_workersCtx[0].prio = JobPriority::High;
224 for (
auto& worker: m_workers)
228 uint32_t workerIdx = 1;
229 set_workers_high_prio_inter(workerIdx, countHighPrio);
230 create_background_worker_threads(workerIdx);
238 count = gaia::core::get_min(count, m_frameWorkersCnt);
239 m_workerThreadsCnt[0] = count;
240 m_workerThreadsCnt[1] = m_frameWorkersCnt - count;
241 m_workersCnt[0] = count + 1;
242 m_workersCnt[1] = m_workerThreadsCnt[1];
245 create_worker_threads(workerIdx, JobPriority::High, m_workerThreadsCnt[0]);
246 create_worker_threads(workerIdx, JobPriority::Low, m_workerThreadsCnt[1]);
254 const uint32_t realCnt = gaia::core::get_min(count, m_frameWorkersCnt);
255 m_workerThreadsCnt[0] = m_frameWorkersCnt - realCnt;
256 m_workerThreadsCnt[1] = realCnt;
257 m_workersCnt[0] = m_workerThreadsCnt[0] + 1;
258 m_workersCnt[1] = m_workerThreadsCnt[1];
261 create_worker_threads(workerIdx, JobPriority::High, m_workerThreadsCnt[0]);
262 create_worker_threads(workerIdx, JobPriority::Low, m_workerThreadsCnt[1]);
270 detail::tl_workerCtx = m_workersCtx.
data();
272 uint32_t workerIdx = 1;
273 set_workers_high_prio_inter(workerIdx, count);
274 create_background_worker_threads(workerIdx);
282 detail::tl_workerCtx = m_workersCtx.
data();
284 uint32_t workerIdx = 1;
285 set_workers_low_prio_inter(workerIdx, count);
286 create_background_worker_threads(workerIdx);
295 const auto maxBackgroundWorkers = MaxWorkers - 1;
296 count = core::get_min(maxBackgroundWorkers, count);
298 const auto frameWorkersCntOld = m_frameWorkersCnt;
299 const auto highWorkersCntOld = m_workerThreadsCnt[0];
304 m_backgroundWorkersCnt = count;
306 const auto maxFrameWorkers = MaxWorkers - m_backgroundWorkersCnt - 1;
307 m_frameWorkersCnt = core::get_min(frameWorkersCntOld, maxFrameWorkers);
309 for (
auto& ctx: m_workersCtx)
312 m_workersCtx.
resize(m_frameWorkersCnt + 1 + m_backgroundWorkersCnt);
313 m_workers.
resize(m_frameWorkersCnt + m_backgroundWorkersCnt);
315 detail::tl_workerCtx = m_workersCtx.
data();
316 m_workersCtx[0].tp =
this;
317 m_workersCtx[0].workerIdx = 0;
318 m_workersCtx[0].prio = JobPriority::High;
320 for (
auto& worker: m_workers)
323 uint32_t workerIdx = 1;
324 set_workers_high_prio_inter(workerIdx, highWorkersCntOld);
325 create_background_worker_threads(workerIdx);
334 GAIA_ASSERT(main_thread());
336 m_jobManager.dep(std::span(&jobFirst, 1), jobSecond);
345 GAIA_ASSERT(main_thread());
347 m_jobManager.dep(jobsFirst, jobSecond);
358 GAIA_ASSERT(main_thread());
360 m_jobManager.dep_refresh(std::span(&jobFirst, 1), jobSecond);
371 GAIA_ASSERT(main_thread());
373 m_jobManager.dep_refresh(jobsFirst, jobSecond);
383 template <
typename TJob>
385 GAIA_ASSERT(main_thread());
387 job.priority = final_prio(job);
389 auto& mtx = GAIA_PROF_EXTRACT_MUTEX(m_jobAllocMtx);
391 GAIA_PROF_LOCK_MARK(m_jobAllocMtx);
393 return m_jobManager.alloc_job(GAIA_MOV(job));
397 void add_n(JobPriority prio, std::span<JobHandle> jobHandles) {
398 GAIA_ASSERT(main_thread());
399 GAIA_ASSERT(!jobHandles.empty());
401 auto& mtx = GAIA_PROF_EXTRACT_MUTEX(m_jobAllocMtx);
403 GAIA_PROF_LOCK_MARK(m_jobAllocMtx);
405 for (
auto& jobHandle: jobHandles)
406 jobHandle = m_jobManager.alloc_job({{}, prio, JobCreationFlags::Default});
409 GAIA_NODISCARD ParallelCallbackHandle add_parallel_callback(JobArgsFunc callback, uint32_t refs) {
410 auto& mtx = GAIA_PROF_EXTRACT_MUTEX(m_jobAllocMtx);
411 core::lock_scope lock(mtx);
412 GAIA_PROF_LOCK_MARK(m_jobAllocMtx);
414 return m_jobManager.alloc_parallel_callback(GAIA_MOV(callback), refs);
417 void release_parallel_callback(ParallelCallbackHandle handle) {
418 if (!m_jobManager.release_parallel_callback_ref(handle))
421 auto& mtx = GAIA_PROF_EXTRACT_MUTEX(m_jobAllocMtx);
422 core::lock_scope lock(mtx);
423 GAIA_PROF_LOCK_MARK(m_jobAllocMtx);
425 m_jobManager.free_parallel_callback(handle);
428 void release_job(JobHandle jobHandle) {
429 auto& mtx = GAIA_PROF_EXTRACT_MUTEX(m_jobAllocMtx);
430 core::lock_scope lock(mtx);
431 GAIA_PROF_LOCK_MARK(m_jobAllocMtx);
433 m_jobManager.free_job(jobHandle);
444 auto& mtx = GAIA_PROF_EXTRACT_MUTEX(m_jobAllocMtx);
446 GAIA_PROF_LOCK_MARK(m_jobAllocMtx);
447 if (!m_jobManager.valid(jobHandle))
450#if GAIA_ASSERT_ENABLED
452 const auto& jobData = m_jobManager.data(jobHandle);
453 GAIA_ASSERT(jobData.state == 0 || m_jobManager.done(jobData));
457 m_jobManager.free_job(jobHandle);
466 void submit(std::span<JobHandle> jobHandles) {
467 if (jobHandles.empty())
470 GAIA_PROF_SCOPE(tp::submitn);
475 for (
auto handle: jobHandles) {
478 auto& jobData = m_jobManager.data(handle);
479 if GAIA_UNLIKELY (jobData.data.gen != handle.gen())
482 const auto state = m_jobManager.submit(jobData) & JobState::DEP_BITS_MASK;
488 pHandles[cnt++] = handle;
491 auto* ctx = detail::tl_workerCtx;
492 process(std::span(pHandles, cnt), ctx);
503 GAIA_PROF_SCOPE(tp::submit);
505 auto& jobData = m_jobManager.data(jobHandle);
506 if GAIA_UNLIKELY (jobData.data.gen != jobHandle.
gen())
509 const auto state = m_jobManager.submit(jobData) & JobState::DEP_BITS_MASK;
513 auto* ctx = detail::tl_workerCtx;
514 process(std::span(&jobHandle, 1), ctx);
520 if (jobHandles.empty())
523 GAIA_PROF_SCOPE(tp::reset);
525 for (
auto handle: jobHandles) {
526 auto& jobData = m_jobManager.data(handle);
527 m_jobManager.reset_state(jobData);
534 reset_state(std::span(&jobHandle, 1));
540 void reset(std::span<JobHandle> jobHandles) {
541 if (jobHandles.empty())
544 GAIA_ASSERT(main_thread());
545 GAIA_PROF_SCOPE(tp::reset_wait);
548 for (
auto handle: jobHandles) {
554 for (
auto handle: jobHandles) {
558 auto& jobData = m_jobManager.data(handle);
559 const auto state = jobData.state.load() & JobState::STATE_BITS_MASK;
561 if (state == JobState::Released)
564 m_jobManager.reset_state(jobData);
571 reset(std::span(&jobHandle, 1));
580 JobHandle jobHandle = add(GAIA_MOV(job));
594 job.
flags = (JobCreationFlags)((uint8_t)job.
flags | (uint8_t)JobCreationFlags::Background);
595 JobHandle jobHandle = add(GAIA_MOV(job));
607 JobHandle jobHandle = add(GAIA_MOV(job));
608 dep(dependsOn, jobHandle);
621 GAIA_ASSERT(main_thread());
624 GAIA_ASSERT(itemsToProcess != 0);
625 if (itemsToProcess == 0)
629 if GAIA_UNLIKELY (m_stop)
633 const auto prio = job.
priority = final_prio(job);
636 if (groupSize == 0) {
637 const auto cntWorkers = core::get_max(1U, m_workersCnt[(uint32_t)prio]);
638 groupSize = itemsToProcess / cntWorkers + (itemsToProcess % cntWorkers != 0);
646 constexpr uint32_t maxUnitsOfWorkPerGroup = 8;
647 groupSize = groupSize / maxUnitsOfWorkPerGroup;
652 const auto jobs = itemsToProcess / groupSize + (itemsToProcess % groupSize != 0);
659 const uint32_t groupJobIdxEnd = groupSize < itemsToProcess ? groupSize : itemsToProcess;
660 auto groupFunc = GAIA_MOV(job.
func);
661 auto groupJobFunc = [func = GAIA_MOV(groupFunc), groupJobIdxEnd]()
mutable {
664 args.
idxEnd = groupJobIdxEnd;
668 auto handle = add(
Job{GAIA_MOV(groupJobFunc), prio, JobCreationFlags::Default});
675 auto callbackHandle = add_parallel_callback(GAIA_MOV(job.
func), jobs);
678 std::span<JobHandle> handles(pHandles, jobs + 1);
680 add_n(prio, handles);
682#if GAIA_ASSERT_ENABLED
683 for (
auto jobHandle: handles)
684 GAIA_ASSERT(m_jobManager.is_clear(jobHandle));
688 for (uint32_t jobIndex = 0; jobIndex < jobs; ++jobIndex) {
689 const uint32_t groupJobIdxStart = jobIndex * groupSize;
690 const uint32_t groupJobIdxEnd =
691 core::get_min(groupSize, itemsToProcess - groupJobIdxStart) + groupJobIdxStart;
693 auto groupJobFunc = [
this, callbackHandle, groupJobIdxStart, groupJobIdxEnd]() {
696 args.
idxEnd = groupJobIdxEnd;
697 m_jobManager.invoke_parallel_callback(callbackHandle, args);
698 release_parallel_callback(callbackHandle);
701 auto& jobData = m_jobManager.data(pHandles[jobIndex]);
707 auto& jobData = m_jobManager.data(pHandles[jobs]);
712 dep(handles.subspan(0, jobs), pHandles[jobs]);
717 return pHandles[jobs];
728 GAIA_ASSERT(main_thread());
729 GAIA_ASSERT(job.
pCtx !=
nullptr);
730 GAIA_ASSERT(job.
invoke !=
nullptr);
732 GAIA_ASSERT(itemsToProcess != 0);
733 if (itemsToProcess == 0)
736 if GAIA_UNLIKELY (m_stop)
739 const auto prio = job.
priority = final_prio(job);
741 if (groupSize == 0) {
742 const auto cntWorkers = core::get_max(1U, m_workersCnt[(uint32_t)prio]);
743 groupSize = itemsToProcess / cntWorkers + (itemsToProcess % cntWorkers != 0);
745 constexpr uint32_t maxUnitsOfWorkPerGroup = 8;
746 groupSize = groupSize / maxUnitsOfWorkPerGroup;
751 const auto jobs = itemsToProcess / groupSize + (itemsToProcess % groupSize != 0);
754 const uint32_t groupJobIdxEnd = groupSize < itemsToProcess ? groupSize : itemsToProcess;
755 auto* pCtx = job.
pCtx;
757 auto groupJobFunc = [pCtx, invoke, groupJobIdxEnd]() {
760 args.
idxEnd = groupJobIdxEnd;
764 auto handle = add(
Job{GAIA_MOV(groupJobFunc), prio, JobCreationFlags::Default});
770 std::span<JobHandle> handles(pHandles, jobs + 1);
772 add_n(prio, handles);
774#if GAIA_ASSERT_ENABLED
775 for (
auto jobHandle: handles)
776 GAIA_ASSERT(m_jobManager.is_clear(jobHandle));
779 for (uint32_t jobIndex = 0; jobIndex < jobs; ++jobIndex) {
780 const uint32_t groupJobIdxStart = jobIndex * groupSize;
781 const uint32_t groupJobIdxEnd =
782 core::get_min(groupSize, itemsToProcess - groupJobIdxStart) + groupJobIdxStart;
784 auto* pCtx = job.
pCtx;
786 auto groupJobFunc = [pCtx, invoke, groupJobIdxStart, groupJobIdxEnd]() {
789 args.
idxEnd = groupJobIdxEnd;
793 auto& jobData = m_jobManager.data(pHandles[jobIndex]);
798 auto& jobData = m_jobManager.data(pHandles[jobs]);
802 dep(handles.subspan(0, jobs), pHandles[jobs]);
804 return pHandles[jobs];
815 GAIA_PROF_SCOPE(tp::wait);
817 GAIA_ASSERT(main_thread());
823 auto* ctx = detail::tl_workerCtx;
824 auto& jobData = m_jobManager.data(jobHandle);
825 const bool waitBackground = is_background(jobData);
826 auto state = jobData.state.load(std::memory_order_acquire);
829 GAIA_ASSERT(state != 0);
832 for (; (state & JobState::STATE_BITS_MASK) < JobState::Done;
833 state = jobData.state.load(std::memory_order_acquire)) {
836 const bool canHelpBackground = waitBackground && m_backgroundWorkersCnt == 0;
837 const bool hasBackgroundJob = canHelpBackground && try_fetch_background_job(otherJobHandle);
838 const bool hasJob = hasBackgroundJob || try_fetch_job(*ctx, otherJobHandle);
840 if (run(otherJobHandle, ctx))
846 if ((state & JobState::STATE_BITS_MASK) == JobState::Executing) {
847 const auto workerId = (state & JobState::DEP_BITS_MASK);
848 auto* jobDoneEvent = &m_workersCtx[workerId].event;
849 jobDoneEvent->wait();
856 const auto workerBit = 1U << ctx->workerIdx;
857 const auto oldBlockedMask = m_blockedInWorkUntil.fetch_or(workerBit);
858 const auto newState = jobData.state.load();
859 if (newState == state)
860 Futex::wait(&m_blockedInWorkUntil, oldBlockedMask | workerBit, detail::WaitMaskAny);
861 m_blockedInWorkUntil.fetch_and(~workerBit);
868 GAIA_ASSERT(main_thread());
875 auto hwThreads = (uint32_t)std::thread::hardware_concurrency();
876 return core::get_max(1U, hwThreads);
882 uint32_t efficiencyCores = 0;
883#if GAIA_PLATFORM_APPLE
884 size_t size =
sizeof(efficiencyCores);
885 if (sysctlbyname(
"hw.perflevel1.logicalcpu", &efficiencyCores, &size,
nullptr, 0) != 0)
887#elif GAIA_PLATFORM_FREEBSD
891 size_t size =
sizeof(coreType);
893 GAIA_STRFMT(oidName,
sizeof(oidName),
"dev.cpu.%d.coretype", cpuIndex);
894 if (sysctlbyname(oidName, &coreType, &size,
nullptr, 0) != 0)
904#elif GAIA_PLATFORM_WINDOWS
908 if (!GetLogicalProcessorInformationEx(RelationProcessorCore,
nullptr, &length))
912 auto* pBuffer = (SYSTEM_LOGICAL_PROCESSOR_INFORMATION_EX*)malloc(length);
913 if (pBuffer ==
nullptr)
917 if (!GetLogicalProcessorInformationEx(RelationProcessorCore, pBuffer, &length)) {
922 uint32_t heterogenousCnt = 0;
935 for (
char* ptr = (
char*)pBuffer; ptr < (
char*)pBuffer + length;
936 ptr += ((SYSTEM_LOGICAL_PROCESSOR_INFORMATION_EX*)ptr)->Size) {
937 auto* entry = (SYSTEM_LOGICAL_PROCESSOR_INFORMATION_EX*)ptr;
938 if (entry->Relationship == RelationProcessorCore) {
939 if (entry->Processor.EfficiencyClass == 0)
946 if (heterogenousCnt == 0)
950#elif GAIA_PLATFORM_LINUX
953 DIR* dir = opendir(
"/sys/devices/cpu_atom/cpus/");
958 while ((entry = readdir(dir)) !=
nullptr) {
959 if (strncmp(entry->d_name,
"cpu", 3) == 0 && entry->d_name[3] >=
'0' && entry->d_name[3] <=
'9')
966 if (efficiencyCores == 0) {
982 return efficiencyCores;
986 static void* thread_func(
void* pCtx) {
989 detail::tl_workerCtx = &ctx;
994 ctx.tp->set_thread_name(ctx.workerIdx, ctx.prio);
997 ctx.tp->set_thread_priority(ctx.workerIdx, ctx.prio);
1000 ctx.tp->worker_loop(ctx);
1002 detail::tl_workerCtx =
nullptr;
1011 void create_thread(uint32_t workerIdx, JobPriority prio,
bool background) {
1013 GAIA_ASSERT(workerIdx > 0);
1015 auto& ctx = m_workersCtx[workerIdx];
1017 ctx.workerIdx = workerIdx;
1019 ctx.background = background;
1020 ctx.threadCreated =
false;
1022#if GAIA_THREAD_PLATFORM == GAIA_THREAD_STD
1023 m_workers[workerIdx - 1] = std::thread([&ctx]() {
1024 thread_func((
void*)&ctx);
1027 pthread_attr_t attr{};
1028 int ret = pthread_attr_init(&attr);
1030 GAIA_LOG_W(
"pthread_attr_init failed for worker thread %u. ErrCode = %d", workerIdx, ret);
1051 if (prio == JobPriority::Low) {
1052 #if GAIA_PLATFORM_APPLE
1053 ret = pthread_attr_set_qos_class_np(&attr, QOS_CLASS_USER_INTERACTIVE, -9);
1056 "pthread_attr_set_qos_class_np failed for worker thread %u [prio=%u]. ErrCode = %d", workerIdx,
1057 (uint32_t)prio, ret);
1060 ret = pthread_attr_setschedpolicy(&attr, SCHED_OTHER);
1063 "pthread_attr_setschedpolicy SCHED_RR failed for worker thread %u [prio=%u]. ErrCode = %d", workerIdx,
1064 (uint32_t)prio, ret);
1067 int prioMax = core::get_min(38, sched_get_priority_max(SCHED_OTHER));
1068 int prioMin = core::get_min(prioMax, sched_get_priority_min(SCHED_OTHER));
1069 int prioUse = core::get_min(prioMin + 5, prioMax);
1070 prioUse = core::get_max(prioUse, prioMin);
1071 sched_param param{};
1072 param.sched_priority = prioUse;
1074 ret = pthread_attr_setschedparam(&attr, ¶m);
1077 "pthread_attr_setschedparam %d failed for worker thread %u [prio=%u]. ErrCode = %d",
1078 param.sched_priority, workerIdx, (uint32_t)prio, ret);
1082 ret = pthread_attr_setschedpolicy(&attr, SCHED_RR);
1085 "pthread_attr_setschedpolicy SCHED_RR failed for worker thread %u [prio=%u]. ErrCode = %d", workerIdx,
1086 (uint32_t)prio, ret);
1089 int prioMax = core::get_min(41, sched_get_priority_max(SCHED_RR));
1090 int prioMin = core::get_min(prioMax, sched_get_priority_min(SCHED_RR));
1091 int prioUse = core::get_max(prioMax - 5, prioMin);
1092 prioUse = core::get_min(prioUse, prioMax);
1093 sched_param param{};
1094 param.sched_priority = prioUse;
1096 ret = pthread_attr_setschedparam(&attr, ¶m);
1099 "pthread_attr_setschedparam %d failed for worker thread %u [prio=%u]. ErrCode = %d",
1100 param.sched_priority, workerIdx, (uint32_t)prio, ret);
1105 ret = pthread_create(&m_workers[workerIdx - 1], &attr, thread_func, (
void*)&ctx);
1107 GAIA_LOG_W(
"pthread_create failed for worker thread %u. ErrCode = %d", workerIdx, ret);
1109 ctx.threadCreated =
true;
1112 pthread_attr_destroy(&attr);
1116 set_thread_affinity(workerIdx);
1121 void join_thread(uint32_t workerIdx) {
1122 if GAIA_UNLIKELY (workerIdx > m_workers.
size())
1125#if GAIA_THREAD_PLATFORM == GAIA_THREAD_STD
1126 auto& t = m_workers[workerIdx - 1];
1130 auto& ctx = m_workersCtx[workerIdx];
1131 if (!ctx.threadCreated)
1134 auto& t = m_workers[workerIdx - 1];
1135 pthread_join(t,
nullptr);
1136 ctx.threadCreated =
false;
1140 void create_worker_threads(uint32_t& workerIdx, JobPriority prio, uint32_t count) {
1141 for (uint32_t i = 0; i < count; ++i)
1142 create_thread(workerIdx++, prio,
false);
1145 void create_background_worker_threads(uint32_t& workerIdx) {
1146 for (uint32_t i = 0; i < m_backgroundWorkersCnt; ++i)
1147 create_thread(workerIdx++, JobPriority::Low,
true);
1150 void set_thread_priority([[maybe_unused]] uint32_t workerIdx, [[maybe_unused]] JobPriority priority) {
1151#if GAIA_PLATFORM_WINDOWS
1152 HANDLE nativeHandle = (HANDLE)m_workers[workerIdx - 1].native_handle();
1154 THREAD_POWER_THROTTLING_STATE state{};
1155 state.Version = THREAD_POWER_THROTTLING_CURRENT_VERSION;
1156 if (priority == JobPriority::High) {
1160 state.ControlMask = THREAD_POWER_THROTTLING_EXECUTION_SPEED;
1161 state.StateMask = 0;
1166 state.ControlMask = THREAD_POWER_THROTTLING_EXECUTION_SPEED;
1167 state.StateMask = THREAD_POWER_THROTTLING_EXECUTION_SPEED;
1170 BOOL ret = SetThreadInformation(nativeHandle, ThreadPowerThrottling, &state,
sizeof(state));
1172 GAIA_LOG_W(
"SetThreadInformation failed for thread %u", workerIdx);
1180 void set_thread_affinity([[maybe_unused]] uint32_t workerIdx) {
1220 void set_thread_name(uint32_t workerIdx, JobPriority prio) {
1221 const bool background = m_workersCtx[workerIdx].background;
1222 const char* workerKind = background ?
"BG" : prio == JobPriority::High ?
"HI" :
"LO";
1223#if GAIA_PROF_USE_PROFILER_THREAD_NAME
1224 char threadName[16]{};
1225 GAIA_STRFMT(threadName, 16,
"worker_%s_%u", workerKind, workerIdx);
1226 GAIA_PROF_THREAD_NAME(threadName);
1227#elif GAIA_PLATFORM_WINDOWS
1228 auto nativeHandle = (HANDLE)m_workers[workerIdx - 1].native_handle();
1229 const wchar_t* workerKindW = background ? L
"BG" : prio == JobPriority::High ? L
"HI" : L
"LO";
1231 TOSApiFunc_SetThreadDescription pSetThreadDescFunc =
nullptr;
1232 if (
auto* pModule = GetModuleHandleA(
"kernel32.dll")) {
1233 auto* pFunc = GetProcAddress(pModule,
"SetThreadDescription");
1234 pSetThreadDescFunc =
reinterpret_cast<TOSApiFunc_SetThreadDescription
>(
reinterpret_cast<void*
>(pFunc));
1236 if (pSetThreadDescFunc !=
nullptr) {
1237 wchar_t threadName[16]{};
1238 swprintf_s(threadName, L
"worker_%s_%u", workerKindW, workerIdx);
1240 auto hr = pSetThreadDescFunc(nativeHandle, threadName);
1242 GAIA_LOG_W(
"Issue setting name for worker %s thread %u!", workerKind, workerIdx);
1245 #if defined _MSC_VER
1246 char threadName[16]{};
1247 GAIA_STRFMT(threadName, 16,
"worker_%s_%u", workerKind, workerIdx);
1249 THREADNAME_INFO info{};
1250 info.dwType = 0x1000;
1251 info.szName = threadName;
1252 info.dwThreadID = GetThreadId(nativeHandle);
1255 RaiseException(0x406D1388, 0,
sizeof(info) /
sizeof(ULONG_PTR), (ULONG_PTR*)&info);
1256 } __except (EXCEPTION_EXECUTE_HANDLER) {
1260#elif GAIA_PLATFORM_APPLE
1261 char threadName[16]{};
1262 GAIA_STRFMT(threadName, 16,
"worker_%s_%u", workerKind, workerIdx);
1263 auto ret = pthread_setname_np(threadName);
1265 GAIA_LOG_W(
"Issue setting name for worker %s thread %u!", workerKind, workerIdx);
1266#elif GAIA_PLATFORM_LINUX || GAIA_PLATFORM_FREEBSD
1267 auto nativeHandle = m_workers[workerIdx - 1];
1269 char threadName[16]{};
1270 GAIA_STRFMT(threadName, 16,
"worker_%s_%u", workerKind, workerIdx);
1271 GAIA_PROF_THREAD_NAME(threadName);
1272 auto ret = pthread_setname_np(nativeHandle, threadName);
1274 GAIA_LOG_W(
"Issue setting name for worker %s thread %u!", workerKind, workerIdx);
1280 GAIA_NODISCARD
bool main_thread()
const {
1281 return std::this_thread::get_id() == m_mainThreadId;
1286 void main_thread_tick() {
1287 auto& ctx = *detail::tl_workerCtx;
1291 JobHandle jobHandle;
1292 if (!try_fetch_job(ctx, jobHandle))
1295 (void)run(jobHandle, &ctx);
1304 GAIA_NODISCARD
bool try_steal_job(ThreadCtx& ctx, JobPriority prio, JobHandle& jobHandle) {
1305 const auto workerCnt = m_workersCtx.
size();
1306 for (uint32_t i = 0; i < workerCnt;) {
1308 if (i == ctx.workerIdx || m_workersCtx[i].background || m_workersCtx[i].prio != prio) {
1313 const auto res = m_workersCtx[i].jobQueue.try_steal(jobHandle);
1321 if (jobHandle != (JobHandle)JobNull_t{})
1335 GAIA_NODISCARD
bool try_fetch_prio(ThreadCtx& ctx, JobPriority prio, JobHandle& jobHandle) {
1336 if (m_jobQueue[(uint32_t)prio].try_pop(jobHandle))
1339 return try_steal_job(ctx, prio, jobHandle);
1345 GAIA_NODISCARD
bool try_fetch_background_job(JobHandle& jobHandle) {
1346 return m_jobQueueBackground.try_pop(jobHandle);
1353 GAIA_NODISCARD
bool try_fetch_job(ThreadCtx& ctx, JobHandle& jobHandle) {
1355 return try_fetch_background_job(jobHandle);
1358 if (ctx.jobQueue.try_pop(jobHandle))
1362 if (ctx.workerIdx == 0) {
1363 if (try_fetch_prio(ctx, JobPriority::High, jobHandle))
1366 return try_fetch_prio(ctx, JobPriority::Low, jobHandle);
1369 return try_fetch_prio(ctx, ctx.prio, jobHandle);
1376 GAIA_NODISCARD
bool can_run_inline(
const ThreadCtx* ctx,
const JobContainer& jobData)
const {
1377 const bool background = is_background(jobData);
1379 return (ctx !=
nullptr && ctx->background) || m_backgroundWorkersCnt == 0;
1382 if (ctx ==
nullptr || ctx->workerIdx == 0)
1386 if (!ctx->background && ctx->prio == jobData.prio)
1391 return m_workerThreadsCnt[(uint32_t)jobData.prio] == 0;
1397 void wait_for_queue_space(ThreadCtx& ctx,
const JobContainer& jobData) {
1398 const bool background = is_background(jobData);
1403 if (m_backgroundWorkersCnt != 0)
1406 const auto prioIdx = (uint32_t)jobData.prio;
1407 if (m_workerThreadsCnt[prioIdx] != 0)
1412 JobHandle otherJobHandle;
1413 const bool hasWork =
1414 ctx.background ? try_fetch_background_job(otherJobHandle) : try_fetch_prio(ctx, ctx.prio, otherJobHandle);
1416 (void)run(otherJobHandle, &ctx);
1420 std::this_thread::yield();
1426 void worker_loop(ThreadCtx& ctx) {
1430 m_semBackground.
wait();
1432 m_sem[(uint32_t)ctx.prio].
wait();
1436 JobHandle jobHandle;
1437 if (!try_fetch_job(ctx, jobHandle))
1440 (void)run(jobHandle, detail::tl_workerCtx);
1444 const bool stop = m_stop.load();
1452 if (m_workers.
empty())
1459 GAIA_FOR(JobPriorityCnt) {
1460 if (m_workerThreadsCnt[i] != 0)
1461 m_sem[i].
release((int32_t)m_workerThreadsCnt[i]);
1463 if (m_backgroundWorkersCnt != 0)
1464 m_semBackground.
release((int32_t)m_backgroundWorkersCnt);
1466 auto* ctx = detail::tl_workerCtx;
1467 if (ctx ==
nullptr) {
1470 ctx = &m_workersCtx[0];
1471 detail::tl_workerCtx = ctx;
1475 JobHandle jobHandle;
1476 while (try_fetch_job(*ctx, jobHandle)) {
1477 run(jobHandle, ctx);
1479 while (try_fetch_background_job(jobHandle)) {
1480 run(jobHandle, ctx);
1483 detail::tl_workerCtx =
nullptr;
1486 GAIA_FOR(m_workers.
size()) join_thread(i + 1);
1489 m_stop.store(false);
1493 JobPriority final_frame_prio(JobPriority priority) {
1494 const auto cntWorkers = m_workersCnt[(uint32_t)priority];
1495 return cntWorkers > 0
1499 : (JobPriority)(((uint32_t)priority + 1U) % (uint32_t)JobPriorityCnt);
1503 JobPriority final_prio(
const Job& job) {
1504 if ((job.flags & JobCreationFlags::Background) != 0U)
1505 return job.priority;
1507 return final_frame_prio(job.priority);
1511 template <
typename TJob>
1512 JobPriority final_prio(
const TJob& job) {
1513 return final_frame_prio(job.priority);
1519 GAIA_NODISCARD
static bool is_background(
const JobContainer& jobData) {
1520 return (jobData.flags & JobCreationFlags::Background) != 0U;
1523 uint32_t signal_edges(JobContainer& jobData, JobHandle* pReadyHandles) {
1524 const auto max = jobData.edges.depCnt;
1532 auto depHandle = jobData.edges.dep;
1533#if GAIA_LOG_JOB_STATES
1534 GAIA_LOG_N(
"SIGNAL %u.%u -> %u.%u", jobData.idx, jobData.gen, depHandle.id(), depHandle.gen());
1538 auto& depData = m_jobManager.data(depHandle);
1539 if (!JobManager::signal_edge(depData))
1542 pReadyHandles[0] = depHandle;
1547 GAIA_ASSERT(jobData.edges.pDeps !=
nullptr);
1551 auto depHandle = jobData.edges.pDeps[i];
1554 auto& depData = m_jobManager.data(depHandle);
1555 if (!JobManager::signal_edge(depData))
1558 pReadyHandles[cnt++] = depHandle;
1567 void process(std::span<JobHandle> jobHandles, ThreadCtx* ctx) {
1568 auto* pHandles = (JobHandle*)alloca(
sizeof(JobHandle) * jobHandles.size());
1569 uint32_t handlesCnt = 0;
1571 for (
auto handle: jobHandles) {
1572 auto& jobData = m_jobManager.data(handle);
1573 m_jobManager.processing(jobData);
1581 if (!jobData.func.operator
bool())
1582 (void)run(handle, ctx);
1584 pHandles[handlesCnt++] = handle;
1587 std::span handles(pHandles, handlesCnt);
1588 while (!handles.empty()) {
1590 uint32_t pushed = 0;
1591 uint32_t released[JobPriorityCnt]{};
1592 uint32_t backgroundReleased = 0;
1593 for (; pushed < handles.size(); ++pushed) {
1594 const auto handle = handles[pushed];
1595 const auto& jobData = m_jobManager.data(handle);
1596 if (is_background(jobData)) {
1597 if (!m_jobQueueBackground.try_push(handle))
1600 ++backgroundReleased;
1604 const auto prio = jobData.prio;
1608 const bool useLocalQueue = ctx !=
nullptr && !ctx->background && ctx->workerIdx != 0 && ctx->prio == prio;
1610 useLocalQueue ? ctx->jobQueue.try_push(handle) : m_jobQueue[(uint32_t)prio].try_push(handle);
1614 released[(uint32_t)prio]++;
1617 GAIA_FOR(JobPriorityCnt) {
1620 const auto cnt = core::get_min(released[i], m_workerThreadsCnt[i]);
1622 m_sem[i].
release((int32_t)cnt);
1624 const auto backgroundCnt = core::get_min(backgroundReleased, m_backgroundWorkersCnt);
1625 if (backgroundCnt != 0)
1626 m_semBackground.
release((int32_t)backgroundCnt);
1628 handles = handles.subspan(pushed);
1629 if (!handles.empty()) {
1630 const auto handle = handles[0];
1631 const auto& jobData = m_jobManager.data(handle);
1633 if (can_run_inline(ctx, jobData)) {
1637 handles = handles.subspan(1);
1639 GAIA_ASSERT(ctx !=
nullptr);
1640 wait_for_queue_space(*ctx, jobData);
1646 bool run(JobHandle jobHandle, ThreadCtx* ctx) {
1647 if (jobHandle == (JobHandle)JobNull_t{})
1650 auto& jobData = m_jobManager.data(jobHandle);
1651 const bool manualDelete = (jobData.flags & JobCreationFlags::ManualDelete) != 0U;
1652 const bool canWait = (jobData.flags & JobCreationFlags::CanWait) != 0U;
1654 m_jobManager.executing(jobData, ctx->workerIdx);
1656 if (m_blockedInWorkUntil.load() != 0) {
1657 const auto blockedCnt = m_blockedInWorkUntil.exchange(0);
1658 if (blockedCnt != 0)
1659 Futex::wake(&m_blockedInWorkUntil, detail::WaitMaskAll);
1662 GAIA_ASSERT(jobData.idx != (uint32_t)-1 && jobData.data.gen != (uint32_t)-1);
1665 m_jobManager.run(jobData);
1667 if (jobData.edges.depCnt == 0) {
1668 JobManager::finalize(jobData);
1671 auto* pReadyHandles = (JobHandle*)alloca(
sizeof(JobHandle) * jobData.edges.depCnt);
1672 const auto readyHandlesCnt = signal_edges(jobData, pReadyHandles);
1674 JobManager::free_edges(jobData);
1675 JobManager::finalize(jobData);
1678 process(std::span(pReadyHandles, readyHandlesCnt), ctx);
1684 const auto* pFutexValue = &jobData.state;
1689 release_job(jobHandle);