From 06c290853d002d6393345f34880a9e0512c22a0b Mon Sep 17 00:00:00 2001 From: Morgan Caron Date: Sun, 20 Sep 2026 15:58:08 +0200 Subject: [PATCH] fix(execution): eliminate lost wakeup and signal confusion in EventQueue --- modules/Execution/EventQueue.mpp | 29 ++++++++++++++++++----------- 1 file changed, 18 insertions(+), 11 deletions(-) diff --git a/modules/Execution/EventQueue.mpp b/modules/Execution/EventQueue.mpp index ebf8fbc3..7aa40b4a 100644 --- a/modules/Execution/EventQueue.mpp +++ b/modules/Execution/EventQueue.mpp @@ -18,7 +18,7 @@ export namespace CppUtils::Execution inline EventQueue(): m_worker{ [this] { workerThread(); }, - [this] { m_condition.notify_all(); }} + [this] { m_workCondition.notify_all(); }} { m_worker.start(); } @@ -26,6 +26,7 @@ export namespace CppUtils::Execution inline ~EventQueue() { waitUntilFinished(); + m_worker.stop(); } EventQueue(const EventQueue&) = delete; @@ -77,10 +78,10 @@ export namespace CppUtils::Execution auto predicate = [this, &accessor] { return std::ranges::empty(accessor.value()) and not m_isTaskRunning; }; if (timeout == Duration::zero()) { - m_condition.wait(accessor.getLockGuard(), predicate); + m_finishedCondition.wait(accessor.getLockGuard(), predicate); return true; } - return m_condition.wait_for(accessor.getLockGuard(), timeout, predicate); + return m_finishedCondition.wait_for(accessor.getLockGuard(), timeout, predicate); } private: @@ -90,7 +91,7 @@ export namespace CppUtils::Execution auto accessor = m_queue.access(); accessor.value().push(std::move(task)); } - m_condition.notify_one(); + m_workCondition.notify_one(); } inline auto workerThread() -> void @@ -98,7 +99,7 @@ export namespace CppUtils::Execution auto task = std::function{}; { auto accessor = m_queue.access(); - m_condition.wait(accessor.getLockGuard(), [this, &accessor] { + m_workCondition.wait(accessor.getLockGuard(), [this, &accessor] { return not std::ranges::empty(accessor.value()) or m_worker.isStopRequested(); }); @@ -109,18 +110,24 @@ export namespace CppUtils::Execution task = std::move(accessor.value().front()); accessor.value().pop(); } + + auto _ = ScopeGuard{[this] { + { + auto accessor = m_queue.access(); + m_isTaskRunning = false; + } + m_finishedCondition.notify_all(); + }}; + if (task) - { - auto _ = ScopeGuard{[this] { m_isTaskRunning = false; }}; task(); - } - m_condition.notify_all(); } EventDispatcher m_dispatcher; Thread::UniqueLocker>> m_queue; - std::condition_variable m_condition; + std::condition_variable m_workCondition; + std::condition_variable m_finishedCondition; Thread::ThreadLoop m_worker; - std::atomic m_isTaskRunning = false; + bool m_isTaskRunning = false; }; }