From 3c33f5f4557ae5572d16684da6139efcb01fee0a Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Mon, 21 Sep 2026 18:19:13 +0300 Subject: [PATCH 1/4] fix(executor): make wait a completion barrier Track MPSC submissions before publication and complete each ticket when the task runs or is rejected. Make wait and resize drain against the completion protocol, and add blocked-task, nested-submission, and stress regressions for the documented wait contract. --- docs/TaskExecutor.md | 17 +- .../logit_cpp/logit/detail/TaskExecutor.hpp | 72 +++++++-- tests/CMakeLists.txt | 1 + tests/task_executor_wait_barrier_test.cpp | 151 ++++++++++++++++++ 4 files changed, 223 insertions(+), 18 deletions(-) create mode 100644 tests/task_executor_wait_barrier_test.cpp diff --git a/docs/TaskExecutor.md b/docs/TaskExecutor.md index 7d000ce..9682d1d 100644 --- a/docs/TaskExecutor.md +++ b/docs/TaskExecutor.md @@ -29,9 +29,13 @@ integrations rely on. * Synchronisation primitives: * `m_cv` + `m_cv_mutex` coordinate sleepers for both the worker and producers that wait for capacity during `QueuePolicy::Block`. - * `m_queue_condition` wakes `wait()` callers once the queue drains. + * `m_queue_condition` coordinates worker/lifecycle drain notifications. + * A completion ticket is reserved before each MPSC submission is published; + `wait()` snapshots those tickets and waits for every ticket to complete or + be rejected. This closes the worker-side race between an empty-ring check + and the next `try_pop()` attempt. * `m_active_tasks` tracks in-flight work so that `Block` limits concurrent - execution and `wait()` can determine quiescence. + execution and lifecycle resize checks can observe quiescence. * `m_stop_flag` terminates the worker and stops accepting new tasks. * Enables very low producer overhead while maintaining FIFO ordering on the consumer side. @@ -103,10 +107,11 @@ mutate the stopped singleton worker. accepted by the consumer. * When the ring build is enabled, `DropNewest` and `DropOldest` both drop the incoming task; accepted tasks keep their order. -* `wait()` returns once the queue is empty and `m_active_tasks == 0`, or when a - shutdown is requested. In MPSC builds the worker marks a pop attempt active - before removing a task, so `wait()` cannot return in the narrow window between - a dequeued cell becoming free and the task body starting. +* `wait()` returns once all submissions observed when it entered have either + completed or been rejected, or when a shutdown is requested. In MPSC builds + this is enforced by completion tickets rather than by an independent + empty-ring/active-task observation, so `wait()` cannot return in the narrow + window between a dequeued cell becoming free and the task body starting. * `shutdown()` blocks until the worker thread terminates. It is safe to call multiple times. diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 428ecc5..14a723c 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -213,12 +213,22 @@ namespace logit { namespace detail { m_queue_condition.notify_one(); # else enter_producer_(); + + // Reserve a completion ticket before attempting publication. A + // waiter that observes this ticket must wait for either the task + // to run or the submission to be rejected, so it cannot race the + // worker between an empty-ring observation and try_pop(). + { + std::lock_guard completion_lock(m_completion_mutex); + ++m_submitted_tasks; + } std::function local_task = std::move(task); bool done = false; while (!done) { if (m_stop_flag.load(std::memory_order_acquire)) { + complete_task_(); break; } @@ -248,6 +258,7 @@ namespace logit { namespace detail { switch (policy) { case QueuePolicy::DropNewest: m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + complete_task_(); done = true; break; @@ -255,6 +266,7 @@ namespace logit { namespace detail { // Safe MPSC behaviour: drop the incoming task. // Preserves ordering and avoids producer/consumer deadlocks. m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + complete_task_(); done = true; break; @@ -279,12 +291,18 @@ namespace logit { namespace detail { m_stop_flag.load(std::memory_order_acquire)); }); # else - std::unique_lock lock(m_queue_mutex); - m_queue_condition.wait(lock, [this]() { - return ((queue_empty_() && - m_active_tasks.load(std::memory_order_relaxed) == 0) || - m_stop_flag.load(std::memory_order_acquire)); - }); + std::unique_lock completion_lock(m_completion_mutex); + for (;;) { + const auto target = m_submitted_tasks; + m_completion_cv.wait(completion_lock, [this, target]() { + return m_completed_tasks >= target || + m_stop_flag.load(std::memory_order_acquire); + }); + if (m_stop_flag.load(std::memory_order_acquire) || + m_completed_tasks >= m_submitted_tasks) { + break; + } + } # endif } @@ -306,6 +324,7 @@ namespace logit { namespace detail { } m_cv.notify_all(); m_queue_condition.notify_all(); + m_completion_cv.notify_all(); if (m_worker_thread.joinable()) { m_worker_thread.join(); } @@ -413,6 +432,11 @@ namespace logit { namespace detail { std::condition_variable m_cv; ///< Wakes the worker or producers. std::mutex m_cv_mutex; ///< Protects producer/worker sleeps. + std::mutex m_completion_mutex; ///< Protects completion snapshots. + std::condition_variable m_completion_cv; ///< Notifies completion waiters. + std::size_t m_submitted_tasks; ///< Reserved submission tickets. + std::size_t m_completed_tasks; ///< Completed submission tickets. + std::atomic m_resizing; ///< true while a hot resize is in flight. std::condition_variable m_resize_cv; ///< Producers wait here during a resize. std::atomic m_active_producers; ///< Producers currently touching the ring. @@ -473,6 +497,8 @@ namespace logit { namespace detail { drained_any = true; task(); + + complete_task_(); m_active_tasks.fetch_sub(1, std::memory_order_relaxed); m_cv.notify_one(); // freed an in-flight slot @@ -504,12 +530,32 @@ namespace logit { namespace detail { } bool wait_until_idle_(std::chrono::steady_clock::time_point deadline) { - std::unique_lock lock(m_queue_mutex); - return m_queue_condition.wait_until(lock, deadline, [this]() { - return ((queue_empty_() && - m_active_tasks.load(std::memory_order_relaxed) == 0) || - m_stop_flag.load(std::memory_order_acquire)); - }); + std::unique_lock lock(m_completion_mutex); + for (;;) { + const auto target = m_submitted_tasks; + if (m_completed_tasks < target && + !m_completion_cv.wait_until(lock, deadline, [this, target]() { + return m_completed_tasks >= target || + m_stop_flag.load(std::memory_order_acquire); + })) { + return false; + } + if (m_stop_flag.load(std::memory_order_acquire) || + m_completed_tasks >= m_submitted_tasks) { + return true; + } + if (std::chrono::steady_clock::now() >= deadline) { + return false; + } + } + } + + void complete_task_() { + { + std::lock_guard lock(m_completion_mutex); + ++m_completed_tasks; + } + m_completion_cv.notify_all(); } void enter_producer_() { @@ -561,6 +607,8 @@ namespace logit { namespace detail { m_overflow_policy(QueuePolicy::Block), m_dropped_tasks(0), m_active_tasks(0), + m_submitted_tasks(0), + m_completed_tasks(0), m_mpsc_queue(m_default_ring_cap) #endif { diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 7d9713b..7e9fcfa 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -83,6 +83,7 @@ else() scope_timer_test.cpp single_thread_executor_test.cpp task_executor_resize_race_test.cpp + task_executor_wait_barrier_test.cpp unique_file_logger_file_api_test.cpp unique_file_logger_set_queue_config_test.cpp windows_debug_logger_set_queue_config_test.cpp diff --git a/tests/task_executor_wait_barrier_test.cpp b/tests/task_executor_wait_barrier_test.cpp new file mode 100644 index 0000000..cdc6a6f --- /dev/null +++ b/tests/task_executor_wait_barrier_test.cpp @@ -0,0 +1,151 @@ +#include + +#include +#include +#include +#include +#include +#include +#include +#include + +namespace { + +bool test_wait_blocks_for_running_task(logit::detail::TaskExecutor& executor) { + std::mutex gate_mutex; + std::condition_variable gate_cv; + bool task_started = false; + bool release_task = false; + + executor.add_task([&]() { + std::unique_lock lock(gate_mutex); + task_started = true; + gate_cv.notify_all(); + gate_cv.wait(lock, [&]() { return release_task; }); + }); + + { + std::unique_lock lock(gate_mutex); + if (!gate_cv.wait_for(lock, std::chrono::seconds(2), [&]() { + return task_started; + })) { + return false; + } + } + + std::mutex completion_mutex; + std::condition_variable completion_cv; + bool waiter_started = false; + bool waiter_done = false; + + std::thread waiter([&]() { + { + std::lock_guard lock(completion_mutex); + waiter_started = true; + } + completion_cv.notify_all(); + + executor.wait(); + + { + std::lock_guard lock(completion_mutex); + waiter_done = true; + } + completion_cv.notify_all(); + }); + + bool blocked = false; + { + std::unique_lock lock(completion_mutex); + if (completion_cv.wait_for(lock, std::chrono::seconds(2), [&]() { + return waiter_started; + })) { + blocked = !completion_cv.wait_for(lock, std::chrono::milliseconds(100), [&]() { + return waiter_done; + }); + } + } + + { + std::lock_guard lock(gate_mutex); + release_task = true; + } + gate_cv.notify_all(); + waiter.join(); + + return blocked && waiter_done; +} + +bool test_wait_drains_mpsc_submissions(logit::detail::TaskExecutor& executor) { + constexpr std::size_t kRounds = 200; + constexpr std::size_t kProducers = 4; + constexpr std::size_t kTasksPerProducer = 64; + + for (std::size_t round = 0; round < kRounds; ++round) { + std::atomic start(false); + std::atomic completed(0); + std::vector producers; + producers.reserve(kProducers); + + for (std::size_t producer = 0; producer < kProducers; ++producer) { + producers.emplace_back([&]() { + while (!start.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + for (std::size_t task = 0; task < kTasksPerProducer; ++task) { + executor.add_task([&completed]() { + completed.fetch_add(1, std::memory_order_relaxed); + }); + } + }); + } + + start.store(true, std::memory_order_release); + for (auto& producer : producers) { + producer.join(); + } + + executor.wait(); + if (completed.load(std::memory_order_relaxed) != + kProducers * kTasksPerProducer) { + return false; + } + } + + return true; +} + +bool test_wait_drains_nested_submission(logit::detail::TaskExecutor& executor) { + std::atomic completed(0); + executor.add_task([&]() { + completed.fetch_add(1, std::memory_order_relaxed); + executor.add_task([&completed]() { + completed.fetch_add(1, std::memory_order_relaxed); + }); + }); + executor.wait(); + return completed.load(std::memory_order_relaxed) == 2; +} + +} // namespace + +int main() { + auto& executor = logit::detail::TaskExecutor::get_instance(); + executor.wait(); + executor.set_queue_policy(logit::detail::QueuePolicy::Block); + executor.set_max_queue_size(0); + executor.reset_dropped_tasks(); + + const bool running_task = test_wait_blocks_for_running_task(executor); + executor.wait(); + const bool mpsc_submissions = test_wait_drains_mpsc_submissions(executor); + const bool nested_submission = test_wait_drains_nested_submission(executor); + + executor.wait(); + executor.reset_dropped_tasks(); + + const bool ok = running_task && mpsc_submissions && nested_submission; + std::cout << (ok ? "PASS" : "FAIL") + << ": task_executor_wait_barrier" << std::endl; + return ok ? 0 : 1; +} From d62b8ff6fd0931d363ddc18a6e5f7d5702f75e10 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Mon, 21 Sep 2026 18:24:14 +0300 Subject: [PATCH 2/4] perf(executor): keep completion tracking lock-free Use atomic submission and completion counters so the MPSC producer path does not acquire a new mutex. Keep the wait condition variable only for blocking callers while preserving the completion barrier. --- .../logit_cpp/logit/detail/TaskExecutor.hpp | 36 +++++++++---------- 1 file changed, 16 insertions(+), 20 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 14a723c..bc24c4a 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -218,10 +218,7 @@ namespace logit { namespace detail { // waiter that observes this ticket must wait for either the task // to run or the submission to be rejected, so it cannot race the // worker between an empty-ring observation and try_pop(). - { - std::lock_guard completion_lock(m_completion_mutex); - ++m_submitted_tasks; - } + m_submitted_tasks.fetch_add(1, std::memory_order_release); std::function local_task = std::move(task); bool done = false; @@ -291,15 +288,16 @@ namespace logit { namespace detail { m_stop_flag.load(std::memory_order_acquire)); }); # else - std::unique_lock completion_lock(m_completion_mutex); + std::unique_lock completion_lock(m_completion_wait_mutex); for (;;) { - const auto target = m_submitted_tasks; + const auto target = m_submitted_tasks.load(std::memory_order_acquire); m_completion_cv.wait(completion_lock, [this, target]() { - return m_completed_tasks >= target || + return m_completed_tasks.load(std::memory_order_acquire) >= target || m_stop_flag.load(std::memory_order_acquire); }); if (m_stop_flag.load(std::memory_order_acquire) || - m_completed_tasks >= m_submitted_tasks) { + m_completed_tasks.load(std::memory_order_acquire) >= + m_submitted_tasks.load(std::memory_order_acquire)) { break; } } @@ -432,10 +430,10 @@ namespace logit { namespace detail { std::condition_variable m_cv; ///< Wakes the worker or producers. std::mutex m_cv_mutex; ///< Protects producer/worker sleeps. - std::mutex m_completion_mutex; ///< Protects completion snapshots. + std::mutex m_completion_wait_mutex; ///< Serializes completion waiters. std::condition_variable m_completion_cv; ///< Notifies completion waiters. - std::size_t m_submitted_tasks; ///< Reserved submission tickets. - std::size_t m_completed_tasks; ///< Completed submission tickets. + std::atomic m_submitted_tasks; ///< Reserved submission tickets. + std::atomic m_completed_tasks; ///< Completed submission tickets. std::atomic m_resizing; ///< true while a hot resize is in flight. std::condition_variable m_resize_cv; ///< Producers wait here during a resize. @@ -530,18 +528,19 @@ namespace logit { namespace detail { } bool wait_until_idle_(std::chrono::steady_clock::time_point deadline) { - std::unique_lock lock(m_completion_mutex); + std::unique_lock lock(m_completion_wait_mutex); for (;;) { - const auto target = m_submitted_tasks; - if (m_completed_tasks < target && + const auto target = m_submitted_tasks.load(std::memory_order_acquire); + if (m_completed_tasks.load(std::memory_order_acquire) < target && !m_completion_cv.wait_until(lock, deadline, [this, target]() { - return m_completed_tasks >= target || + return m_completed_tasks.load(std::memory_order_acquire) >= target || m_stop_flag.load(std::memory_order_acquire); })) { return false; } if (m_stop_flag.load(std::memory_order_acquire) || - m_completed_tasks >= m_submitted_tasks) { + m_completed_tasks.load(std::memory_order_acquire) >= + m_submitted_tasks.load(std::memory_order_acquire)) { return true; } if (std::chrono::steady_clock::now() >= deadline) { @@ -551,10 +550,7 @@ namespace logit { namespace detail { } void complete_task_() { - { - std::lock_guard lock(m_completion_mutex); - ++m_completed_tasks; - } + m_completed_tasks.fetch_add(1, std::memory_order_release); m_completion_cv.notify_all(); } From 22e8bab41ce3a49287d773f3bb291685ac1bcaa9 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Mon, 21 Sep 2026 19:20:30 +0300 Subject: [PATCH 3/4] fix(executor): synchronize completion notifications Update the completion state while holding the wait mutex before notifying condition-variable waiters. This prevents a lost wake-up while keeping the MPSC submission path free of locks. --- include/logit_cpp/logit/detail/TaskExecutor.hpp | 17 ++++++++++------- 1 file changed, 10 insertions(+), 7 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index bc24c4a..cee0465 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -292,11 +292,11 @@ namespace logit { namespace detail { for (;;) { const auto target = m_submitted_tasks.load(std::memory_order_acquire); m_completion_cv.wait(completion_lock, [this, target]() { - return m_completed_tasks.load(std::memory_order_acquire) >= target || + return m_completed_tasks >= target || m_stop_flag.load(std::memory_order_acquire); }); if (m_stop_flag.load(std::memory_order_acquire) || - m_completed_tasks.load(std::memory_order_acquire) >= + m_completed_tasks >= m_submitted_tasks.load(std::memory_order_acquire)) { break; } @@ -433,7 +433,7 @@ namespace logit { namespace detail { std::mutex m_completion_wait_mutex; ///< Serializes completion waiters. std::condition_variable m_completion_cv; ///< Notifies completion waiters. std::atomic m_submitted_tasks; ///< Reserved submission tickets. - std::atomic m_completed_tasks; ///< Completed submission tickets. + std::size_t m_completed_tasks; ///< Completed submission tickets. std::atomic m_resizing; ///< true while a hot resize is in flight. std::condition_variable m_resize_cv; ///< Producers wait here during a resize. @@ -531,15 +531,15 @@ namespace logit { namespace detail { std::unique_lock lock(m_completion_wait_mutex); for (;;) { const auto target = m_submitted_tasks.load(std::memory_order_acquire); - if (m_completed_tasks.load(std::memory_order_acquire) < target && + if (m_completed_tasks < target && !m_completion_cv.wait_until(lock, deadline, [this, target]() { - return m_completed_tasks.load(std::memory_order_acquire) >= target || + return m_completed_tasks >= target || m_stop_flag.load(std::memory_order_acquire); })) { return false; } if (m_stop_flag.load(std::memory_order_acquire) || - m_completed_tasks.load(std::memory_order_acquire) >= + m_completed_tasks >= m_submitted_tasks.load(std::memory_order_acquire)) { return true; } @@ -550,7 +550,10 @@ namespace logit { namespace detail { } void complete_task_() { - m_completed_tasks.fetch_add(1, std::memory_order_release); + { + std::lock_guard lock(m_completion_wait_mutex); + ++m_completed_tasks; + } m_completion_cv.notify_all(); } From 362c09929758bf49f2bbc8466e6c8dc80de9e096 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Mon, 21 Sep 2026 19:26:44 +0300 Subject: [PATCH 4/4] test(executor): bound wait barrier regression Give the TaskExecutor wait regression a bounded CTest timeout so a future completion protocol regression fails the suite instead of hanging CI indefinitely. --- tests/CMakeLists.txt | 3 +++ 1 file changed, 3 insertions(+) diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 7e9fcfa..7b195a5 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -138,6 +138,9 @@ else() if(test_name STREQUAL "logger_legacy_targeted_path_test") target_compile_definitions(${test_name} PRIVATE LOGIT_BENCH_LEGACY_REGISTRY=1) endif() + if(test_name STREQUAL "task_executor_wait_barrier_test") + set_tests_properties(${test_name} PROPERTIES TIMEOUT 30) + endif() if(LOGIT_WITH_OTLP AND test_name MATCHES "^otlp_http_logger_(integration|callback|gzip|zstd)_test$") target_include_directories(${test_name} PRIVATE "${CMAKE_CURRENT_SOURCE_DIR}/../external/kurlyk/external/Simple-Web-Server")