diff --git a/FEXCore/Source/CMakeLists.txt b/FEXCore/Source/CMakeLists.txt index 26458cad9..7db99c5a7 100644 --- a/FEXCore/Source/CMakeLists.txt +++ b/FEXCore/Source/CMakeLists.txt @@ -70,6 +70,7 @@ set(SRCS Utils/LongJump.cpp Utils/Telemetry.cpp Utils/Threads.cpp + Utils/WorkQueueThread.cpp Utils/Profiler.cpp) if (ARCHITECTURE_arm64) diff --git a/FEXCore/Source/Utils/WorkQueueThread.cpp b/FEXCore/Source/Utils/WorkQueueThread.cpp new file mode 100644 index 000000000..1ca2eccfb --- /dev/null +++ b/FEXCore/Source/Utils/WorkQueueThread.cpp @@ -0,0 +1,54 @@ +// SPDX-License-Identifier: MIT +#include +#include + +namespace FEXCore { + +WorkQueueThread::WorkQueueThread() { + Thread = FEXCore::Threads::Thread::Create(ThreadEntry, this); +} + +WorkQueueThread::~WorkQueueThread() { + { + std::unique_lock lk {Mutex}; + Stop = true; + } + CV.notify_one(); + + if (Thread && Thread->joinable()) { + Thread->join(nullptr); + } +} + +void WorkQueueThread::QueueWork(fextl::unique_ptr Work) { + { + std::unique_lock lk {Mutex}; + Queue.push_back(std::move(Work)); + } + CV.notify_one(); +} + +void WorkQueueThread::ThreadProc() { + while (true) { + fextl::unique_ptr Work; + { + std::unique_lock lk {Mutex}; + while (!(Stop || !Queue.empty())) { + CV.wait(lk); + } + if (Queue.empty()) { + // nothing to do? must be stopping + LOGMAN_THROW_A_FMT(Stop, "WorkQueueThread wakes up empty but no Stop?"); + return; + } + + Work = std::move(Queue.front()); + Queue.pop_front(); + } + + Work->Run(); + // Work is destroyed here + } +} + +} // namespace FEXCore diff --git a/FEXCore/include/FEXCore/Utils/WorkQueueThread.h b/FEXCore/include/FEXCore/Utils/WorkQueueThread.h new file mode 100644 index 000000000..e9ff114a9 --- /dev/null +++ b/FEXCore/include/FEXCore/Utils/WorkQueueThread.h @@ -0,0 +1,41 @@ +// SPDX-License-Identifier: MIT +#pragma once + +#include +#include +#include + +#include +#include + +namespace FEXCore { + +class WorkQueueThread { +public: + // destroyed after Run() + struct WorkItem { + virtual ~WorkItem() = default; + virtual void Run() = 0; + }; + + WorkQueueThread(); + ~WorkQueueThread(); + + void QueueWork(fextl::unique_ptr Work); + +private: + // static function for the ::Thread to refer to + static void* ThreadEntry(void* Self) { + static_cast(Self)->ThreadProc(); + return nullptr; + } + void ThreadProc(); + + std::mutex Mutex; + std::condition_variable CV; + fextl::deque> Queue; + bool Stop = false; + + fextl::unique_ptr Thread; +}; +} // namespace FEXCore