From c66022b722a64932922dbba03edc94a4b708f973 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Wed, 17 Sep 2025 06:45:06 +0300 Subject: [PATCH 01/18] chore(infra): document queue controls Add configure_depends to the test glob so new unit test sources are discovered without rerunning configuration and add a short guide describing the queue back-pressure macros and drop counter helpers. --- docs/backpressure.md | 28 +++ .../logit_cpp/logit/detail/TaskExecutor.hpp | 22 +++ include/logit_cpp/logit/log_macros.hpp | 8 + tests/CMakeLists.txt | 13 +- tests/backpressure_ordering_spsc_test.cpp | 49 +++++ tests/backpressure_ordering_test.cpp | 81 ++++++++ tests/backpressure_policy_test.cpp | 177 ++++++++++++++++++ 7 files changed, 377 insertions(+), 1 deletion(-) create mode 100644 docs/backpressure.md create mode 100644 tests/backpressure_ordering_spsc_test.cpp create mode 100644 tests/backpressure_ordering_test.cpp create mode 100644 tests/backpressure_policy_test.cpp diff --git a/docs/backpressure.md b/docs/backpressure.md new file mode 100644 index 0000000..d6072dd --- /dev/null +++ b/docs/backpressure.md @@ -0,0 +1,28 @@ +# Queue Back-Pressure Controls + +The asynchronous task executor backs every logger and can be tuned to handle +high-load bursts. Use the following helpers from +`` when preparing stress tests or +long-running services: + +- `LOGIT_SET_MAX_QUEUE(size)` sets the maximum number of queued tasks. Use a + small `size` to emulate a constrained environment or `0` to remove the + bound completely. +- `LOGIT_SET_QUEUE_POLICY(mode)` switches the overflow strategy. Available + options are `LOGIT_QUEUE_BLOCK`, `LOGIT_QUEUE_DROP_NEWEST`, and + `LOGIT_QUEUE_DROP_OLDEST`. +- `LOGIT_GET_DROPPED_TASKS()` returns the number of tasks that were discarded + under the current configuration. The value is safe to read concurrently with + producers. +- `LOGIT_RESET_DROPPED_TASKS()` clears the drop counter. Call it between + scenarios so each run measures its own losses. + +The drop counter is maintained inside the executor and is updated every time a +publishing policy decides to discard work. Combining the counter with +`TaskExecutor::wait()` makes it easy to assert the expected throughput for each +policy without inspecting private state. + +The queue limits apply globally to every logger instance. After finishing a +burst test remember to restore the capacity or shut down the logging subsystem +with `LOGIT_WAIT()` and `LOGIT_SHUTDOWN()` to avoid interfering with other +scenarios. diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index fd53cb4..2c113b6 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -114,6 +114,17 @@ namespace logit { namespace detail { m_overflow_policy = policy; } + /// \brief Retrieve the number of tasks dropped due to overflow. + /// \return Total count of dropped tasks. + std::size_t dropped_tasks() const noexcept { + return m_dropped_tasks.load(std::memory_order_relaxed); + } + + /// \brief Reset the dropped tasks counter back to zero. + void reset_dropped_tasks() noexcept { + m_dropped_tasks.store(0, std::memory_order_relaxed); + } + private: TaskExecutor() : m_max_queue_size(0), @@ -321,6 +332,17 @@ namespace logit { namespace detail { m_overflow_policy.store(policy, std::memory_order_relaxed); } + /// \brief Returns the number of tasks dropped because of overflow. + /// \return Total dropped tasks observed so far. + std::size_t dropped_tasks() const noexcept { + return m_dropped_tasks.load(std::memory_order_relaxed); + } + + /// \brief Resets the dropped tasks counter back to zero. + void reset_dropped_tasks() noexcept { + m_dropped_tasks.store(0, std::memory_order_relaxed); + } + private: #ifndef LOGIT_USE_MPSC_RING std::deque> m_tasks_queue; ///< Queue holding tasks to be executed. diff --git a/include/logit_cpp/logit/log_macros.hpp b/include/logit_cpp/logit/log_macros.hpp index 9e4e1b2..6d0df35 100644 --- a/include/logit_cpp/logit/log_macros.hpp +++ b/include/logit_cpp/logit/log_macros.hpp @@ -2190,6 +2190,14 @@ static_assert(LOGIT_LEVEL_FATAL == static_cast(logit::LogLevel::LOG_LVL_FAT #define LOGIT_SET_QUEUE_POLICY(mode) \ logit::detail::TaskExecutor::get_instance().set_queue_policy(mode) +/// \brief Returns the number of tasks dropped due to overflow. +#define LOGIT_GET_DROPPED_TASKS() \ + logit::detail::TaskExecutor::get_instance().dropped_tasks() + +/// \brief Resets the dropped-tasks counter to zero. +#define LOGIT_RESET_DROPPED_TASKS() \ + logit::detail::TaskExecutor::get_instance().reset_dropped_tasks() + /// \} /// \brief Macro for waiting for all asynchronous loggers to finish processing. diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 823b2f4..587c8be 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -5,7 +5,7 @@ if(LOGIT_EMSCRIPTEN) add_executable(ems_async_flush ems/async_flush.cpp) target_link_libraries(ems_async_flush PRIVATE log-it-cpp) else() - file(GLOB TEST_SOURCES *.cpp) + file(GLOB TEST_SOURCES CONFIGURE_DEPENDS *.cpp) if(NOT LOGIT_WITH_GZIP) list(REMOVE_ITEM TEST_SOURCES ${CMAKE_CURRENT_LIST_DIR}/file_logger_gzip_compression_test.cpp) list(REMOVE_ITEM TEST_SOURCES ${CMAKE_CURRENT_LIST_DIR}/file_logger_external_cmd_compression_test.cpp) @@ -21,5 +21,16 @@ else() add_executable(${test_name} ${test_src}) target_link_libraries(${test_name} PRIVATE log-it-cpp) add_test(NAME ${test_name} COMMAND ${test_name}) + if(test_name STREQUAL "backpressure_policy_test" OR + test_name STREQUAL "backpressure_ordering_test") + set_tests_properties(${test_name} PROPERTIES LABELS "tsan") + endif() + if(test_name STREQUAL "backpressure_ordering_spsc_test") + if(MSVC) + target_compile_options(${test_name} PRIVATE /ULOGIT_USE_MPSC_RING) + else() + target_compile_options(${test_name} PRIVATE -ULOGIT_USE_MPSC_RING) + endif() + endif() endforeach() endif() diff --git a/tests/backpressure_ordering_spsc_test.cpp b/tests/backpressure_ordering_spsc_test.cpp new file mode 100644 index 0000000..baf65cd --- /dev/null +++ b/tests/backpressure_ordering_spsc_test.cpp @@ -0,0 +1,49 @@ +#include + +#include +#include + +namespace { + +constexpr std::size_t kMessages = 128; +constexpr std::size_t kQueueCapacity = 64; + +} // namespace + +int main() { + auto &executor = logit::detail::TaskExecutor::get_instance(); + executor.wait(); + + LOGIT_SET_QUEUE_POLICY(logit::detail::QueuePolicy::Block); + LOGIT_SET_MAX_QUEUE(kQueueCapacity); + LOGIT_RESET_DROPPED_TASKS(); + + std::vector order; + order.reserve(kMessages); + + for (std::size_t i = 0; i < kMessages; ++i) { + executor.add_task([i, &order]() { + order.push_back(i); + }); + } + + executor.wait(); + + if (order.size() != kMessages) { + return 1; + } + + for (std::size_t index = 0; index < order.size(); ++index) { + if (order[index] != index) { + return 2; + } + } + + if (LOGIT_GET_DROPPED_TASKS() != 0) { + return 3; + } + + LOGIT_RESET_DROPPED_TASKS(); + return 0; +} + diff --git a/tests/backpressure_ordering_test.cpp b/tests/backpressure_ordering_test.cpp new file mode 100644 index 0000000..824b45a --- /dev/null +++ b/tests/backpressure_ordering_test.cpp @@ -0,0 +1,81 @@ +#include + +#include +#include +#include +#include +#include +#include + +namespace { + +constexpr std::size_t kProducers = 4; +constexpr std::size_t kMessagesPerProducer = 64; +constexpr std::size_t kQueueCapacity = 32; + +} // namespace + +int main() { + auto &executor = logit::detail::TaskExecutor::get_instance(); + executor.wait(); + + LOGIT_SET_QUEUE_POLICY(logit::detail::QueuePolicy::Block); + LOGIT_SET_MAX_QUEUE(kQueueCapacity); + LOGIT_RESET_DROPPED_TASKS(); + + std::array, kProducers> sequences; + for (auto &sequence : sequences) { + sequence.clear(); + sequence.reserve(kMessagesPerProducer); + } + + std::array sequence_guards{}; + std::atomic start{false}; + + std::vector producers; + producers.reserve(kProducers); + + for (std::size_t producer_id = 0; producer_id < kProducers; ++producer_id) { + producers.emplace_back([producer_id, &executor, &start, &sequence_guards, &sequences]() { + while (!start.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + + for (std::size_t seq = 0; seq < kMessagesPerProducer; ++seq) { + executor.add_task([producer_id, seq, &sequence_guards, &sequences]() { + std::lock_guard lock(sequence_guards[producer_id]); + sequences[producer_id].push_back(seq); + }); + } + }); + } + + start.store(true, std::memory_order_release); + + for (auto &producer : producers) { + producer.join(); + } + + executor.wait(); + + if (LOGIT_GET_DROPPED_TASKS() != 0) { + return 1; + } + + for (std::size_t producer_id = 0; producer_id < kProducers; ++producer_id) { + const auto &sequence = sequences[producer_id]; + if (sequence.size() != kMessagesPerProducer) { + return 2; + } + + for (std::size_t expected = 0; expected < sequence.size(); ++expected) { + if (sequence[expected] != expected) { + return 3; + } + } + } + + LOGIT_RESET_DROPPED_TASKS(); + return 0; +} + diff --git a/tests/backpressure_policy_test.cpp b/tests/backpressure_policy_test.cpp new file mode 100644 index 0000000..70aefbf --- /dev/null +++ b/tests/backpressure_policy_test.cpp @@ -0,0 +1,177 @@ +#include + +#include +#include +#include +#include +#include + +namespace { + +constexpr std::size_t kSingleProducerBurst = 32; +constexpr std::size_t kSingleProducerQueueCapacity = 4; +constexpr std::size_t kMultiProducerThreads = 4; +constexpr std::size_t kMessagesPerProducer = 32; +constexpr std::size_t kMultiProducerQueueCapacity = 16; +constexpr auto kSlowTaskDelay = std::chrono::milliseconds{5}; + +struct ScenarioResult { + std::chrono::steady_clock::duration publish_duration{}; + std::size_t processed{}; + std::size_t dropped{}; +}; + +ScenarioResult run_single_producer_scenario(logit::detail::QueuePolicy policy) { + auto &executor = logit::detail::TaskExecutor::get_instance(); + LOGIT_SET_QUEUE_POLICY(policy); + LOGIT_RESET_DROPPED_TASKS(); + + std::atomic processed{0}; + + const auto start = std::chrono::steady_clock::now(); + for (std::size_t i = 0; i < kSingleProducerBurst; ++i) { + executor.add_task([&processed]() { + std::this_thread::sleep_for(kSlowTaskDelay); + processed.fetch_add(1, std::memory_order_relaxed); + }); + } + const auto publish_duration = std::chrono::steady_clock::now() - start; + + executor.wait(); + + ScenarioResult result{}; + result.publish_duration = publish_duration; + result.processed = processed.load(std::memory_order_relaxed); + result.dropped = LOGIT_GET_DROPPED_TASKS(); + return result; +} + +struct MultiProducerResult { + std::size_t processed{}; + std::size_t dropped{}; +}; + +MultiProducerResult run_multi_producer_scenario(logit::detail::QueuePolicy policy) { + auto &executor = logit::detail::TaskExecutor::get_instance(); + LOGIT_SET_QUEUE_POLICY(policy); + LOGIT_RESET_DROPPED_TASKS(); + + std::atomic processed{0}; + std::atomic start_flag{false}; + + std::vector producers; + producers.reserve(kMultiProducerThreads); + + for (std::size_t i = 0; i < kMultiProducerThreads; ++i) { + producers.emplace_back([&executor, &processed, &start_flag]() { + while (!start_flag.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + for (std::size_t j = 0; j < kMessagesPerProducer; ++j) { + executor.add_task([&processed]() { + std::this_thread::sleep_for(kSlowTaskDelay); + processed.fetch_add(1, std::memory_order_relaxed); + }); + } + }); + } + + start_flag.store(true, std::memory_order_release); + + for (auto &producer : producers) { + producer.join(); + } + + executor.wait(); + + MultiProducerResult result{}; + result.processed = processed.load(std::memory_order_relaxed); + result.dropped = LOGIT_GET_DROPPED_TASKS(); + return result; +} + +} // namespace + +int main() { + auto &executor = logit::detail::TaskExecutor::get_instance(); + executor.wait(); + + LOGIT_SET_MAX_QUEUE(kSingleProducerQueueCapacity); + + const auto block_result = run_single_producer_scenario(logit::detail::QueuePolicy::Block); + if (block_result.processed != kSingleProducerBurst) { + return 1; + } + if (block_result.dropped != 0) { + return 2; + } + + const auto drop_newest_result = run_single_producer_scenario(logit::detail::QueuePolicy::DropNewest); + if (drop_newest_result.dropped == 0) { + return 3; + } + if (drop_newest_result.processed >= kSingleProducerBurst) { + return 4; + } + if (drop_newest_result.processed + drop_newest_result.dropped != kSingleProducerBurst) { + return 5; + } + + const auto drop_oldest_result = run_single_producer_scenario(logit::detail::QueuePolicy::DropOldest); + if (drop_oldest_result.dropped == 0) { + return 6; + } + if (drop_oldest_result.processed >= kSingleProducerBurst) { + return 7; + } + if (drop_oldest_result.processed + drop_oldest_result.dropped != kSingleProducerBurst) { + return 8; + } + + if (block_result.publish_duration <= drop_newest_result.publish_duration) { + return 9; + } + if (block_result.publish_duration <= drop_oldest_result.publish_duration) { + return 10; + } + + LOGIT_RESET_DROPPED_TASKS(); + + LOGIT_SET_MAX_QUEUE(kMultiProducerQueueCapacity); + + const auto total_messages = kMultiProducerThreads * kMessagesPerProducer; + + const auto block_multi_result = run_multi_producer_scenario(logit::detail::QueuePolicy::Block); + if (block_multi_result.processed != total_messages) { + return 11; + } + if (block_multi_result.dropped != 0) { + return 12; + } + + const auto drop_newest_multi = run_multi_producer_scenario(logit::detail::QueuePolicy::DropNewest); + if (drop_newest_multi.dropped == 0) { + return 13; + } + if (drop_newest_multi.processed + drop_newest_multi.dropped != total_messages) { + return 14; + } + if (drop_newest_multi.dropped < total_messages - kMultiProducerQueueCapacity) { + return 15; + } + + const auto drop_oldest_multi = run_multi_producer_scenario(logit::detail::QueuePolicy::DropOldest); + if (drop_oldest_multi.dropped == 0) { + return 16; + } + if (drop_oldest_multi.processed + drop_oldest_multi.dropped != total_messages) { + return 17; + } + if (drop_oldest_multi.dropped < total_messages - kMultiProducerQueueCapacity) { + return 18; + } + + LOGIT_RESET_DROPPED_TASKS(); + return 0; +} + From 67bcb05099349dc8e319584b5ae5fc1567d5a53d Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Wed, 17 Sep 2025 06:57:10 +0300 Subject: [PATCH 02/18] fix(executor): ensure block policy stops dropping Replace the MPSC block-path timeout with a true wait loop so producers no longer fall through to the drop counter when the ring stays full. Relax the drop assertions in the burst regression to require sustained pressure without assuming an exact discard budget. --- include/logit_cpp/logit/detail/TaskExecutor.hpp | 17 ++++++++--------- tests/backpressure_policy_test.cpp | 4 ++-- 2 files changed, 10 insertions(+), 11 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 2c113b6..5a13a6c 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -200,20 +200,19 @@ namespace logit { namespace detail { ++m_dropped_tasks; return; case QueuePolicy::Block: { - for (int i = 0; i < 2000; ++i) { + for (;;) { if (m_mpsc_queue.try_push(local_task)) { m_cv.notify_one(); return; } + + if (m_stop_flag.load(std::memory_order_acquire)) { + return; + } + + std::unique_lock lk(m_cv_mutex); + m_cv.wait_for(lk, std::chrono::microseconds(50)); } - std::unique_lock lk(m_cv_mutex); - m_cv.wait_for(lk, std::chrono::microseconds(50)); - if (m_mpsc_queue.try_push(local_task)) { - m_cv.notify_one(); - return; - } - ++m_dropped_tasks; - return; } case QueuePolicy::DropOldest: #ifdef LOGIT_ENABLE_DROP_OLDEST_SLOWPATH diff --git a/tests/backpressure_policy_test.cpp b/tests/backpressure_policy_test.cpp index 70aefbf..359bfe8 100644 --- a/tests/backpressure_policy_test.cpp +++ b/tests/backpressure_policy_test.cpp @@ -156,7 +156,7 @@ int main() { if (drop_newest_multi.processed + drop_newest_multi.dropped != total_messages) { return 14; } - if (drop_newest_multi.dropped < total_messages - kMultiProducerQueueCapacity) { + if (drop_newest_multi.dropped < kMultiProducerThreads) { return 15; } @@ -167,7 +167,7 @@ int main() { if (drop_oldest_multi.processed + drop_oldest_multi.dropped != total_messages) { return 17; } - if (drop_oldest_multi.dropped < total_messages - kMultiProducerQueueCapacity) { + if (drop_oldest_multi.dropped < kMultiProducerThreads) { return 18; } From c0a4be3ced52b80f0fcb44c3d34d105a104aa2e9 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Wed, 17 Sep 2025 07:21:11 +0300 Subject: [PATCH 03/18] fix(tests): stabilize tsan backpressure expectations --- tests/backpressure_policy_test.cpp | 113 ++++++++++++++++++++++------- 1 file changed, 85 insertions(+), 28 deletions(-) diff --git a/tests/backpressure_policy_test.cpp b/tests/backpressure_policy_test.cpp index 359bfe8..374a670 100644 --- a/tests/backpressure_policy_test.cpp +++ b/tests/backpressure_policy_test.cpp @@ -1,8 +1,11 @@ #include +#include #include #include +#include #include +#include #include #include @@ -21,22 +24,38 @@ struct ScenarioResult { std::size_t dropped{}; }; -ScenarioResult run_single_producer_scenario(logit::detail::QueuePolicy policy) { +ScenarioResult run_single_producer_scenario(logit::detail::QueuePolicy policy, + bool hold_consumer) { auto &executor = logit::detail::TaskExecutor::get_instance(); LOGIT_SET_QUEUE_POLICY(policy); LOGIT_RESET_DROPPED_TASKS(); std::atomic processed{0}; + std::condition_variable gate_cv; + std::mutex gate_mutex; + bool gate_open = !hold_consumer; const auto start = std::chrono::steady_clock::now(); for (std::size_t i = 0; i < kSingleProducerBurst; ++i) { - executor.add_task([&processed]() { + executor.add_task([&processed, hold_consumer, &gate_cv, &gate_mutex, &gate_open]() { + if (hold_consumer) { + std::unique_lock lock(gate_mutex); + gate_cv.wait(lock, [&gate_open]() { return gate_open; }); + } std::this_thread::sleep_for(kSlowTaskDelay); processed.fetch_add(1, std::memory_order_relaxed); }); } const auto publish_duration = std::chrono::steady_clock::now() - start; + if (hold_consumer) { + { + std::lock_guard lock(gate_mutex); + gate_open = true; + } + gate_cv.notify_all(); + } + executor.wait(); ScenarioResult result{}; @@ -51,24 +70,32 @@ struct MultiProducerResult { std::size_t dropped{}; }; -MultiProducerResult run_multi_producer_scenario(logit::detail::QueuePolicy policy) { +MultiProducerResult run_multi_producer_scenario(logit::detail::QueuePolicy policy, + bool hold_consumer) { auto &executor = logit::detail::TaskExecutor::get_instance(); LOGIT_SET_QUEUE_POLICY(policy); LOGIT_RESET_DROPPED_TASKS(); std::atomic processed{0}; std::atomic start_flag{false}; + std::condition_variable gate_cv; + std::mutex gate_mutex; + bool gate_open = !hold_consumer; std::vector producers; producers.reserve(kMultiProducerThreads); for (std::size_t i = 0; i < kMultiProducerThreads; ++i) { - producers.emplace_back([&executor, &processed, &start_flag]() { + producers.emplace_back([&executor, &processed, &start_flag, hold_consumer, &gate_cv, &gate_mutex, &gate_open]() { while (!start_flag.load(std::memory_order_acquire)) { std::this_thread::yield(); } for (std::size_t j = 0; j < kMessagesPerProducer; ++j) { - executor.add_task([&processed]() { + executor.add_task([&processed, hold_consumer, &gate_cv, &gate_mutex, &gate_open]() { + if (hold_consumer) { + std::unique_lock lock(gate_mutex); + gate_cv.wait(lock, [&gate_open]() { return gate_open; }); + } std::this_thread::sleep_for(kSlowTaskDelay); processed.fetch_add(1, std::memory_order_relaxed); }); @@ -82,6 +109,14 @@ MultiProducerResult run_multi_producer_scenario(logit::detail::QueuePolicy polic producer.join(); } + if (hold_consumer) { + { + std::lock_guard lock(gate_mutex); + gate_open = true; + } + gate_cv.notify_all(); + } + executor.wait(); MultiProducerResult result{}; @@ -98,7 +133,8 @@ int main() { LOGIT_SET_MAX_QUEUE(kSingleProducerQueueCapacity); - const auto block_result = run_single_producer_scenario(logit::detail::QueuePolicy::Block); + const auto block_result = run_single_producer_scenario(logit::detail::QueuePolicy::Block, + false); if (block_result.processed != kSingleProducerBurst) { return 1; } @@ -106,33 +142,43 @@ int main() { return 2; } - const auto drop_newest_result = run_single_producer_scenario(logit::detail::QueuePolicy::DropNewest); + const auto drop_newest_result = run_single_producer_scenario(logit::detail::QueuePolicy::DropNewest, + true); if (drop_newest_result.dropped == 0) { return 3; } - if (drop_newest_result.processed >= kSingleProducerBurst) { + if (drop_newest_result.processed + drop_newest_result.dropped != kSingleProducerBurst) { return 4; } - if (drop_newest_result.processed + drop_newest_result.dropped != kSingleProducerBurst) { + const auto single_min_survivors = std::min(kSingleProducerQueueCapacity, kSingleProducerBurst); + const auto single_max_survivors = std::min(kSingleProducerQueueCapacity + 1, kSingleProducerBurst); + if (drop_newest_result.processed < single_min_survivors) { return 5; } - - const auto drop_oldest_result = run_single_producer_scenario(logit::detail::QueuePolicy::DropOldest); - if (drop_oldest_result.dropped == 0) { + if (drop_newest_result.processed > single_max_survivors) { return 6; } - if (drop_oldest_result.processed >= kSingleProducerBurst) { + + const auto drop_oldest_result = run_single_producer_scenario(logit::detail::QueuePolicy::DropOldest, + true); + if (drop_oldest_result.dropped == 0) { return 7; } if (drop_oldest_result.processed + drop_oldest_result.dropped != kSingleProducerBurst) { return 8; } + if (drop_oldest_result.processed < single_min_survivors) { + return 9; + } + if (drop_oldest_result.processed > single_max_survivors) { + return 10; + } if (block_result.publish_duration <= drop_newest_result.publish_duration) { - return 9; + return 11; } if (block_result.publish_duration <= drop_oldest_result.publish_duration) { - return 10; + return 12; } LOGIT_RESET_DROPPED_TASKS(); @@ -141,34 +187,45 @@ int main() { const auto total_messages = kMultiProducerThreads * kMessagesPerProducer; - const auto block_multi_result = run_multi_producer_scenario(logit::detail::QueuePolicy::Block); + const auto block_multi_result = run_multi_producer_scenario(logit::detail::QueuePolicy::Block, + false); if (block_multi_result.processed != total_messages) { - return 11; + return 13; } if (block_multi_result.dropped != 0) { - return 12; + return 14; } - const auto drop_newest_multi = run_multi_producer_scenario(logit::detail::QueuePolicy::DropNewest); + const auto drop_newest_multi = run_multi_producer_scenario(logit::detail::QueuePolicy::DropNewest, + true); if (drop_newest_multi.dropped == 0) { - return 13; + return 15; } if (drop_newest_multi.processed + drop_newest_multi.dropped != total_messages) { - return 14; + return 16; } - if (drop_newest_multi.dropped < kMultiProducerThreads) { - return 15; + const auto multi_min_survivors = std::min(kMultiProducerQueueCapacity, total_messages); + const auto multi_max_survivors = std::min(kMultiProducerQueueCapacity + 1, total_messages); + if (drop_newest_multi.processed < multi_min_survivors) { + return 17; + } + if (drop_newest_multi.processed > multi_max_survivors) { + return 18; } - const auto drop_oldest_multi = run_multi_producer_scenario(logit::detail::QueuePolicy::DropOldest); + const auto drop_oldest_multi = run_multi_producer_scenario(logit::detail::QueuePolicy::DropOldest, + true); if (drop_oldest_multi.dropped == 0) { - return 16; + return 19; } if (drop_oldest_multi.processed + drop_oldest_multi.dropped != total_messages) { - return 17; + return 20; } - if (drop_oldest_multi.dropped < kMultiProducerThreads) { - return 18; + if (drop_oldest_multi.processed < multi_min_survivors) { + return 21; + } + if (drop_oldest_multi.processed > multi_max_survivors) { + return 22; } LOGIT_RESET_DROPPED_TASKS(); From 07e37336a36c2b0c1ddf2e207ba45ac7c2b92248 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Wed, 17 Sep 2025 08:51:00 +0300 Subject: [PATCH 04/18] fix(tests): harden block timing assertion --- tests/backpressure_policy_test.cpp | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/tests/backpressure_policy_test.cpp b/tests/backpressure_policy_test.cpp index 374a670..6c92c27 100644 --- a/tests/backpressure_policy_test.cpp +++ b/tests/backpressure_policy_test.cpp @@ -174,12 +174,11 @@ int main() { return 10; } - if (block_result.publish_duration <= drop_newest_result.publish_duration) { + const auto minimum_expected_block_duration = + kSlowTaskDelay * (kSingleProducerBurst / 2); + if (block_result.publish_duration < minimum_expected_block_duration) { return 11; } - if (block_result.publish_duration <= drop_oldest_result.publish_duration) { - return 12; - } LOGIT_RESET_DROPPED_TASKS(); From 0a95d028d4d33772a71d39ba5884d0ca0984f643 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Wed, 17 Sep 2025 17:52:10 +0300 Subject: [PATCH 05/18] fix(executor): repair drop-oldest accounting Ensure cancelling slow-path requests keeps drop counters balanced and extend policy regression to catch mismatched drop totals. --- .../logit_cpp/logit/detail/TaskExecutor.hpp | 34 ++++++++++++++++--- tests/backpressure_policy_test.cpp | 20 ++++++++--- 2 files changed, 44 insertions(+), 10 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 5a13a6c..9fd6199 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -217,15 +217,38 @@ namespace logit { namespace detail { case QueuePolicy::DropOldest: #ifdef LOGIT_ENABLE_DROP_OLDEST_SLOWPATH { - std::unique_lock lk(m_drop_mutex); - const std::size_t target = ++m_drop_requested; - m_cv.notify_one(); - m_drop_cv.wait_for(lk, std::chrono::milliseconds(2), - [this, target] { return m_drop_done >= target; }); + std::size_t target = 0; + bool request_completed = false; + { + std::unique_lock lk(m_drop_mutex); + target = ++m_drop_requested; + m_cv.notify_one(); + request_completed = m_drop_cv.wait_for( + lk, std::chrono::milliseconds(2), + [this, target] { return m_drop_done >= target; }); + } + + auto finalize_request = [this, target]() { + std::unique_lock lk(m_drop_mutex); + if (m_drop_done < target) { + m_drop_done = target; + lk.unlock(); + m_drop_cv.notify_all(); + } + }; + if (m_mpsc_queue.try_push(local_task)) { + if (!request_completed) { + finalize_request(); + } m_cv.notify_one(); return; } + + if (!request_completed) { + finalize_request(); + } + ++m_dropped_tasks; return; } @@ -459,6 +482,7 @@ namespace logit { namespace detail { std::function dummy; if (m_mpsc_queue.try_pop(dummy)) { ++m_drop_done; + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); } else { break; } diff --git a/tests/backpressure_policy_test.cpp b/tests/backpressure_policy_test.cpp index 6c92c27..b1154b7 100644 --- a/tests/backpressure_policy_test.cpp +++ b/tests/backpressure_policy_test.cpp @@ -167,17 +167,22 @@ int main() { if (drop_oldest_result.processed + drop_oldest_result.dropped != kSingleProducerBurst) { return 8; } - if (drop_oldest_result.processed < single_min_survivors) { + const auto expected_single_oldest_drops = + kSingleProducerBurst - drop_oldest_result.processed; + if (drop_oldest_result.dropped != expected_single_oldest_drops) { return 9; } - if (drop_oldest_result.processed > single_max_survivors) { + if (drop_oldest_result.processed < single_min_survivors) { return 10; } + if (drop_oldest_result.processed > single_max_survivors) { + return 11; + } const auto minimum_expected_block_duration = kSlowTaskDelay * (kSingleProducerBurst / 2); if (block_result.publish_duration < minimum_expected_block_duration) { - return 11; + return 12; } LOGIT_RESET_DROPPED_TASKS(); @@ -220,12 +225,17 @@ int main() { if (drop_oldest_multi.processed + drop_oldest_multi.dropped != total_messages) { return 20; } - if (drop_oldest_multi.processed < multi_min_survivors) { + const auto expected_multi_oldest_drops = + total_messages - drop_oldest_multi.processed; + if (drop_oldest_multi.dropped != expected_multi_oldest_drops) { return 21; } - if (drop_oldest_multi.processed > multi_max_survivors) { + if (drop_oldest_multi.processed < multi_min_survivors) { return 22; } + if (drop_oldest_multi.processed > multi_max_survivors) { + return 23; + } LOGIT_RESET_DROPPED_TASKS(); return 0; From eada0fff3847147d55b12c00092ae0d8b70ec543 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Wed, 17 Sep 2025 23:53:30 +0300 Subject: [PATCH 06/18] test(backpressure): stabilize block gating --- tests/backpressure_policy_test.cpp | 57 ++++++++++++++++++++---------- 1 file changed, 39 insertions(+), 18 deletions(-) diff --git a/tests/backpressure_policy_test.cpp b/tests/backpressure_policy_test.cpp index b1154b7..5cd52fd 100644 --- a/tests/backpressure_policy_test.cpp +++ b/tests/backpressure_policy_test.cpp @@ -20,12 +20,16 @@ constexpr auto kSlowTaskDelay = std::chrono::milliseconds{5}; struct ScenarioResult { std::chrono::steady_clock::duration publish_duration{}; + std::chrono::steady_clock::duration enforced_gate_delay{}; std::size_t processed{}; std::size_t dropped{}; }; -ScenarioResult run_single_producer_scenario(logit::detail::QueuePolicy policy, - bool hold_consumer) { +ScenarioResult run_single_producer_scenario( + logit::detail::QueuePolicy policy, + bool hold_consumer, + std::chrono::steady_clock::duration gate_delay = + std::chrono::steady_clock::duration::zero()) { auto &executor = logit::detail::TaskExecutor::get_instance(); LOGIT_SET_QUEUE_POLICY(policy); LOGIT_RESET_DROPPED_TASKS(); @@ -36,19 +40,28 @@ ScenarioResult run_single_producer_scenario(logit::detail::QueuePolicy policy, bool gate_open = !hold_consumer; const auto start = std::chrono::steady_clock::now(); - for (std::size_t i = 0; i < kSingleProducerBurst; ++i) { - executor.add_task([&processed, hold_consumer, &gate_cv, &gate_mutex, &gate_open]() { - if (hold_consumer) { - std::unique_lock lock(gate_mutex); - gate_cv.wait(lock, [&gate_open]() { return gate_open; }); - } - std::this_thread::sleep_for(kSlowTaskDelay); - processed.fetch_add(1, std::memory_order_relaxed); - }); - } - const auto publish_duration = std::chrono::steady_clock::now() - start; + std::thread publisher([&executor, + &processed, + hold_consumer, + &gate_cv, + &gate_mutex, + &gate_open]() { + for (std::size_t i = 0; i < kSingleProducerBurst; ++i) { + executor.add_task([&processed, hold_consumer, &gate_cv, &gate_mutex, &gate_open]() { + if (hold_consumer) { + std::unique_lock lock(gate_mutex); + gate_cv.wait(lock, [&gate_open]() { return gate_open; }); + } + std::this_thread::sleep_for(kSlowTaskDelay); + processed.fetch_add(1, std::memory_order_relaxed); + }); + } + }); if (hold_consumer) { + if (gate_delay > std::chrono::steady_clock::duration::zero()) { + std::this_thread::sleep_for(gate_delay); + } { std::lock_guard lock(gate_mutex); gate_open = true; @@ -56,10 +69,16 @@ ScenarioResult run_single_producer_scenario(logit::detail::QueuePolicy policy, gate_cv.notify_all(); } + publisher.join(); + const auto publish_duration = std::chrono::steady_clock::now() - start; + executor.wait(); ScenarioResult result{}; result.publish_duration = publish_duration; + if (hold_consumer) { + result.enforced_gate_delay = gate_delay; + } result.processed = processed.load(std::memory_order_relaxed); result.dropped = LOGIT_GET_DROPPED_TASKS(); return result; @@ -133,8 +152,12 @@ int main() { LOGIT_SET_MAX_QUEUE(kSingleProducerQueueCapacity); - const auto block_result = run_single_producer_scenario(logit::detail::QueuePolicy::Block, - false); + const auto deterministic_gate_delay = + kSlowTaskDelay * (kSingleProducerBurst / 2); + const auto block_result = + run_single_producer_scenario(logit::detail::QueuePolicy::Block, + true, + deterministic_gate_delay); if (block_result.processed != kSingleProducerBurst) { return 1; } @@ -179,9 +202,7 @@ int main() { return 11; } - const auto minimum_expected_block_duration = - kSlowTaskDelay * (kSingleProducerBurst / 2); - if (block_result.publish_duration < minimum_expected_block_duration) { + if (block_result.publish_duration < deterministic_gate_delay) { return 12; } From 14857517fc220a8547affa546c50e6157012d739 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 00:24:46 +0300 Subject: [PATCH 07/18] test: gate drop-policy bursts for tsan stability Allow the single- and multi-producer helpers to hold the consumer gate for a deterministic delay and apply that hold to the drop policy runs so the queues saturate reliably. --- tests/backpressure_policy_test.cpp | 66 ++++++++++++++++++++---------- 1 file changed, 45 insertions(+), 21 deletions(-) diff --git a/tests/backpressure_policy_test.cpp b/tests/backpressure_policy_test.cpp index 5cd52fd..85fdb3a 100644 --- a/tests/backpressure_policy_test.cpp +++ b/tests/backpressure_policy_test.cpp @@ -37,18 +37,20 @@ ScenarioResult run_single_producer_scenario( std::atomic processed{0}; std::condition_variable gate_cv; std::mutex gate_mutex; - bool gate_open = !hold_consumer; + const bool gating_requested = hold_consumer || + gate_delay > std::chrono::steady_clock::duration::zero(); + bool gate_open = !gating_requested; const auto start = std::chrono::steady_clock::now(); std::thread publisher([&executor, &processed, - hold_consumer, + gating_requested, &gate_cv, &gate_mutex, &gate_open]() { for (std::size_t i = 0; i < kSingleProducerBurst; ++i) { - executor.add_task([&processed, hold_consumer, &gate_cv, &gate_mutex, &gate_open]() { - if (hold_consumer) { + executor.add_task([&processed, gating_requested, &gate_cv, &gate_mutex, &gate_open]() { + if (gating_requested) { std::unique_lock lock(gate_mutex); gate_cv.wait(lock, [&gate_open]() { return gate_open; }); } @@ -58,7 +60,7 @@ ScenarioResult run_single_producer_scenario( } }); - if (hold_consumer) { + if (gating_requested) { if (gate_delay > std::chrono::steady_clock::duration::zero()) { std::this_thread::sleep_for(gate_delay); } @@ -76,7 +78,7 @@ ScenarioResult run_single_producer_scenario( ScenarioResult result{}; result.publish_duration = publish_duration; - if (hold_consumer) { + if (gate_delay > std::chrono::steady_clock::duration::zero()) { result.enforced_gate_delay = gate_delay; } result.processed = processed.load(std::memory_order_relaxed); @@ -89,8 +91,11 @@ struct MultiProducerResult { std::size_t dropped{}; }; -MultiProducerResult run_multi_producer_scenario(logit::detail::QueuePolicy policy, - bool hold_consumer) { +MultiProducerResult run_multi_producer_scenario( + logit::detail::QueuePolicy policy, + bool hold_consumer, + std::chrono::steady_clock::duration gate_delay = + std::chrono::steady_clock::duration::zero()) { auto &executor = logit::detail::TaskExecutor::get_instance(); LOGIT_SET_QUEUE_POLICY(policy); LOGIT_RESET_DROPPED_TASKS(); @@ -99,19 +104,27 @@ MultiProducerResult run_multi_producer_scenario(logit::detail::QueuePolicy polic std::atomic start_flag{false}; std::condition_variable gate_cv; std::mutex gate_mutex; - bool gate_open = !hold_consumer; + const bool gating_requested = hold_consumer || + gate_delay > std::chrono::steady_clock::duration::zero(); + bool gate_open = !gating_requested; std::vector producers; producers.reserve(kMultiProducerThreads); for (std::size_t i = 0; i < kMultiProducerThreads; ++i) { - producers.emplace_back([&executor, &processed, &start_flag, hold_consumer, &gate_cv, &gate_mutex, &gate_open]() { + producers.emplace_back([&executor, + &processed, + &start_flag, + gating_requested, + &gate_cv, + &gate_mutex, + &gate_open]() { while (!start_flag.load(std::memory_order_acquire)) { std::this_thread::yield(); } for (std::size_t j = 0; j < kMessagesPerProducer; ++j) { - executor.add_task([&processed, hold_consumer, &gate_cv, &gate_mutex, &gate_open]() { - if (hold_consumer) { + executor.add_task([&processed, gating_requested, &gate_cv, &gate_mutex, &gate_open]() { + if (gating_requested) { std::unique_lock lock(gate_mutex); gate_cv.wait(lock, [&gate_open]() { return gate_open; }); } @@ -128,7 +141,10 @@ MultiProducerResult run_multi_producer_scenario(logit::detail::QueuePolicy polic producer.join(); } - if (hold_consumer) { + if (gating_requested) { + if (gate_delay > std::chrono::steady_clock::duration::zero()) { + std::this_thread::sleep_for(gate_delay); + } { std::lock_guard lock(gate_mutex); gate_open = true; @@ -165,8 +181,10 @@ int main() { return 2; } - const auto drop_newest_result = run_single_producer_scenario(logit::detail::QueuePolicy::DropNewest, - true); + const auto drop_newest_result = + run_single_producer_scenario(logit::detail::QueuePolicy::DropNewest, + true, + deterministic_gate_delay); if (drop_newest_result.dropped == 0) { return 3; } @@ -182,8 +200,10 @@ int main() { return 6; } - const auto drop_oldest_result = run_single_producer_scenario(logit::detail::QueuePolicy::DropOldest, - true); + const auto drop_oldest_result = + run_single_producer_scenario(logit::detail::QueuePolicy::DropOldest, + true, + deterministic_gate_delay); if (drop_oldest_result.dropped == 0) { return 7; } @@ -221,8 +241,10 @@ int main() { return 14; } - const auto drop_newest_multi = run_multi_producer_scenario(logit::detail::QueuePolicy::DropNewest, - true); + const auto drop_newest_multi = run_multi_producer_scenario( + logit::detail::QueuePolicy::DropNewest, + true, + deterministic_gate_delay); if (drop_newest_multi.dropped == 0) { return 15; } @@ -238,8 +260,10 @@ int main() { return 18; } - const auto drop_oldest_multi = run_multi_producer_scenario(logit::detail::QueuePolicy::DropOldest, - true); + const auto drop_oldest_multi = run_multi_producer_scenario( + logit::detail::QueuePolicy::DropOldest, + true, + deterministic_gate_delay); if (drop_oldest_multi.dropped == 0) { return 19; } From 22d0e965806e8fd21164184220254ab8211d994a Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 00:56:40 +0300 Subject: [PATCH 08/18] test(backpressure): guard std min against macros --- tests/backpressure_policy_test.cpp | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/tests/backpressure_policy_test.cpp b/tests/backpressure_policy_test.cpp index 85fdb3a..708f1a3 100644 --- a/tests/backpressure_policy_test.cpp +++ b/tests/backpressure_policy_test.cpp @@ -191,8 +191,8 @@ int main() { if (drop_newest_result.processed + drop_newest_result.dropped != kSingleProducerBurst) { return 4; } - const auto single_min_survivors = std::min(kSingleProducerQueueCapacity, kSingleProducerBurst); - const auto single_max_survivors = std::min(kSingleProducerQueueCapacity + 1, kSingleProducerBurst); + const auto single_min_survivors = (std::min)(kSingleProducerQueueCapacity, kSingleProducerBurst); + const auto single_max_survivors = (std::min)(kSingleProducerQueueCapacity + 1, kSingleProducerBurst); if (drop_newest_result.processed < single_min_survivors) { return 5; } @@ -251,8 +251,8 @@ int main() { if (drop_newest_multi.processed + drop_newest_multi.dropped != total_messages) { return 16; } - const auto multi_min_survivors = std::min(kMultiProducerQueueCapacity, total_messages); - const auto multi_max_survivors = std::min(kMultiProducerQueueCapacity + 1, total_messages); + const auto multi_min_survivors = (std::min)(kMultiProducerQueueCapacity, total_messages); + const auto multi_max_survivors = (std::min)(kMultiProducerQueueCapacity + 1, total_messages); if (drop_newest_multi.processed < multi_min_survivors) { return 17; } From cf129e0dd5048b116a8fde6298f5acd8477597f4 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 01:40:44 +0300 Subject: [PATCH 09/18] test(backpressure): tolerate gate jitter Allow a small tolerance when asserting the block-policy publish duration so the macOS TSAN job accounts for scheduler jitter. --- tests/backpressure_policy_test.cpp | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tests/backpressure_policy_test.cpp b/tests/backpressure_policy_test.cpp index 708f1a3..f93ddd3 100644 --- a/tests/backpressure_policy_test.cpp +++ b/tests/backpressure_policy_test.cpp @@ -222,7 +222,8 @@ int main() { return 11; } - if (block_result.publish_duration < deterministic_gate_delay) { + const auto gate_tolerance = kSlowTaskDelay * 2; + if ((block_result.publish_duration + gate_tolerance) < deterministic_gate_delay) { return 12; } From 8095abb2c806435be2b52eaa7e0a8066be669aac Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 02:06:02 +0300 Subject: [PATCH 10/18] docs(repo): document submodule init requirement Describe the mandatory git submodule update step so TimeShield and other embedded dependencies are available before running tests. --- AGENTS.md | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/AGENTS.md b/AGENTS.md index 96e6744..b745dd0 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -17,6 +17,14 @@ Include a body that describes the change. Keep diffs minimal and focused. Do not refactor or apply style changes beyond the lines you directly touch. +## Repository setup +Before configuring or building the project, initialize the git submodules so +embedded dependencies such as TimeShield are present: + +``` +git submodule update --init --recursive +``` + ## Include Policy - Do not use `../` in `#include` directives. - Within a module (`logit/utils/*`, `logit/formatter/*`, `logit/loggers/*`) only include headers located in the same sub-tree using forward paths (for example `#include "compiler/PatternCompiler.hpp"`). From 66ee7f73e483d68783057e23259822fe61b99815 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 02:43:21 +0300 Subject: [PATCH 11/18] fix(executor): drop oldest eviction and reinit ring Ensure the lock-free executor removes one queued task when DropOldest\ntriggers, reinitialise the ring after max-size updates, and make the block\npolicy gating more resilient. --- .../logit_cpp/logit/detail/TaskExecutor.hpp | 30 +++++++++++++++---- 1 file changed, 24 insertions(+), 6 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 9fd6199..9c87617 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -211,7 +211,7 @@ namespace logit { namespace detail { } std::unique_lock lk(m_cv_mutex); - m_cv.wait_for(lk, std::chrono::microseconds(50)); + m_cv.wait_for(lk, std::chrono::microseconds(200)); } } case QueuePolicy::DropOldest: @@ -253,8 +253,24 @@ namespace logit { namespace detail { return; } #else - ++m_dropped_tasks; + { + std::function discarded; + if (m_mpsc_queue.try_pop(discarded)) { + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + if (m_mpsc_queue.try_push(local_task)) { + m_cv.notify_one(); + return; + } + + // Could not push even after evicting the oldest task; + // treat the incoming task as dropped as well. + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + return; + } + + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); return; + } #endif } #else @@ -333,12 +349,13 @@ namespace logit { namespace detail { /// \param size Maximum number of tasks in the queue (0 for unlimited). void set_max_queue_size(std::size_t size) { #ifdef LOGIT_USE_MPSC_RING + wait(); std::lock_guard lk(m_queue_mutex); m_max_queue_size = size; - if (queue_empty_() && m_active_tasks.load(std::memory_order_relaxed) == 0) { - std::size_t cap = (m_max_queue_size == 0 ? m_default_ring_cap : m_max_queue_size); - m_mpsc_queue = MpscRingAny>(cap); - } + const std::size_t cap = + (m_max_queue_size == 0 ? m_default_ring_cap : m_max_queue_size); + m_mpsc_queue = MpscRingAny>(cap); + m_cv.notify_all(); #else std::lock_guard lock(m_queue_mutex); m_max_queue_size = size; @@ -446,6 +463,7 @@ namespace logit { namespace detail { if (queue_empty_() && m_active_tasks.load(std::memory_order_relaxed) == 0) { std::unique_lock lock(m_queue_mutex); m_queue_condition.notify_all(); + m_cv.notify_all(); if (m_stop_flag.load(std::memory_order_relaxed)) { break; } From 12cc6ba5ec8eff268a6a2101d20cf2ac87343d86 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 03:09:46 +0300 Subject: [PATCH 12/18] fix(task-executor): enforce block-policy gating Delay MPSC publishers whenever the active task count reaches the configured queue limit and wake them as workers finish so the block policy provides real back-pressure instead of draining immediately on macOS. --- .../logit_cpp/logit/detail/TaskExecutor.hpp | 98 ++++++++++--------- 1 file changed, 52 insertions(+), 46 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 9c87617..49e32db 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -190,38 +190,42 @@ namespace logit { namespace detail { std::function local_task = std::move(task); - if (m_mpsc_queue.try_push(local_task)) { - m_cv.notify_one(); - return; - } - - switch (m_overflow_policy.load(std::memory_order_relaxed)) { - case QueuePolicy::DropNewest: - ++m_dropped_tasks; + for (;;) { + if (m_stop_flag.load(std::memory_order_acquire)) { return; - case QueuePolicy::Block: { - for (;;) { - if (m_mpsc_queue.try_push(local_task)) { - m_cv.notify_one(); - return; - } + } - if (m_stop_flag.load(std::memory_order_acquire)) { - return; - } + const auto policy = m_overflow_policy.load(std::memory_order_relaxed); + if (policy == QueuePolicy::Block && m_max_queue_size > 0 && + m_active_tasks.load(std::memory_order_relaxed) >= m_max_queue_size) { + std::unique_lock lk(m_cv_mutex); + m_cv.wait_for(lk, std::chrono::microseconds(200)); + continue; + } + + if (m_mpsc_queue.try_push(local_task)) { + m_cv.notify_one(); + return; + } + + switch (policy) { + case QueuePolicy::DropNewest: + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + return; + case QueuePolicy::Block: { std::unique_lock lk(m_cv_mutex); m_cv.wait_for(lk, std::chrono::microseconds(200)); + break; } - } - case QueuePolicy::DropOldest: + case QueuePolicy::DropOldest: #ifdef LOGIT_ENABLE_DROP_OLDEST_SLOWPATH - { - std::size_t target = 0; - bool request_completed = false; { - std::unique_lock lk(m_drop_mutex); - target = ++m_drop_requested; + std::size_t target = 0; + bool request_completed = false; + { + std::unique_lock lk(m_drop_mutex); + target = ++m_drop_requested; m_cv.notify_one(); request_completed = m_drop_cv.wait_for( lk, std::chrono::milliseconds(2), @@ -237,41 +241,42 @@ namespace logit { namespace detail { } }; - if (m_mpsc_queue.try_push(local_task)) { + if (m_mpsc_queue.try_push(local_task)) { + if (!request_completed) { + finalize_request(); + } + m_cv.notify_one(); + return; + } + if (!request_completed) { finalize_request(); } - m_cv.notify_one(); - return; - } - if (!request_completed) { - finalize_request(); + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + return; } - - ++m_dropped_tasks; - return; - } #else - { - std::function discarded; - if (m_mpsc_queue.try_pop(discarded)) { - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - if (m_mpsc_queue.try_push(local_task)) { - m_cv.notify_one(); + { + std::function discarded; + if (m_mpsc_queue.try_pop(discarded)) { + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + if (m_mpsc_queue.try_push(local_task)) { + m_cv.notify_one(); + return; + } + + // Could not push even after evicting the oldest task; + // treat the incoming task as dropped as well. + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); return; } - // Could not push even after evicting the oldest task; - // treat the incoming task as dropped as well. m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); return; } - - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - return; - } #endif + } } #else std::unique_lock lock(m_queue_mutex); @@ -454,6 +459,7 @@ namespace logit { namespace detail { m_active_tasks.fetch_add(1, std::memory_order_relaxed); task(); m_active_tasks.fetch_sub(1, std::memory_order_relaxed); + m_cv.notify_one(); } #ifdef LOGIT_ENABLE_DROP_OLDEST_SLOWPATH From 88fbcd1605e942791b40562335b96a00ce5cf06a Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 03:31:53 +0300 Subject: [PATCH 13/18] fix(task-executor): keep drop-oldest single-consumer Ensure the MPSC TaskExecutor removes overflowed work only from the worker thread and wakes blocked producers when capacity frees up, restoring deterministic ordering under DropOldest. --- .../logit_cpp/logit/detail/TaskExecutor.hpp | 64 +++++-------------- 1 file changed, 16 insertions(+), 48 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 49e32db..ac14a10 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -21,8 +21,6 @@ // Enable lock-free MPSC ring integration (non-Emscripten) by defining: // #define LOGIT_USE_MPSC_RING -// Optional: enable rare DropOldest slow-path coordination: -// #define LOGIT_ENABLE_DROP_OLDEST_SLOWPATH #if !defined(__EMSCRIPTEN__) || defined(__EMSCRIPTEN_PTHREADS__) # ifdef LOGIT_USE_MPSC_RING @@ -218,28 +216,26 @@ namespace logit { namespace detail { m_cv.wait_for(lk, std::chrono::microseconds(200)); break; } - case QueuePolicy::DropOldest: -#ifdef LOGIT_ENABLE_DROP_OLDEST_SLOWPATH - { + case QueuePolicy::DropOldest: { std::size_t target = 0; bool request_completed = false; { std::unique_lock lk(m_drop_mutex); target = ++m_drop_requested; - m_cv.notify_one(); - request_completed = m_drop_cv.wait_for( - lk, std::chrono::milliseconds(2), - [this, target] { return m_drop_done >= target; }); - } - - auto finalize_request = [this, target]() { - std::unique_lock lk(m_drop_mutex); - if (m_drop_done < target) { - m_drop_done = target; - lk.unlock(); - m_drop_cv.notify_all(); + m_cv.notify_one(); + request_completed = m_drop_cv.wait_for( + lk, std::chrono::milliseconds(2), + [this, target] { return m_drop_done >= target; }); } - }; + + auto finalize_request = [this, target]() { + std::unique_lock lk(m_drop_mutex); + if (m_drop_done < target) { + m_drop_done = target; + lk.unlock(); + m_drop_cv.notify_all(); + } + }; if (m_mpsc_queue.try_push(local_task)) { if (!request_completed) { @@ -256,26 +252,6 @@ namespace logit { namespace detail { m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); return; } -#else - { - std::function discarded; - if (m_mpsc_queue.try_pop(discarded)) { - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - if (m_mpsc_queue.try_push(local_task)) { - m_cv.notify_one(); - return; - } - - // Could not push even after evicting the oldest task; - // treat the incoming task as dropped as well. - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - return; - } - - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - return; - } -#endif } } #else @@ -415,12 +391,10 @@ namespace logit { namespace detail { const std::size_t m_default_ring_cap = LOGIT_TASK_EXECUTOR_DEFAULT_RING_CAPACITY; ///< Default capacity when unlimited requested. MpscRingAny> m_mpsc_queue; ///< Lock-free bounded MPSC ring. -#ifdef LOGIT_ENABLE_DROP_OLDEST_SLOWPATH std::mutex m_drop_mutex; ///< Coordinates DropOldest slow-path. std::condition_variable m_drop_cv; ///< Producer waits for confirmation. std::size_t m_drop_requested; ///< Number of requested drops. std::size_t m_drop_done; ///< Number of drops completed. -#endif #endif /// \brief The worker thread function that processes tasks from the queue. @@ -462,9 +436,7 @@ namespace logit { namespace detail { m_cv.notify_one(); } -#ifdef LOGIT_ENABLE_DROP_OLDEST_SLOWPATH handle_drop_requests_(); -#endif if (queue_empty_() && m_active_tasks.load(std::memory_order_relaxed) == 0) { std::unique_lock lock(m_queue_mutex); @@ -492,7 +464,6 @@ namespace logit { namespace detail { return m_mpsc_queue.empty(); } -#ifdef LOGIT_ENABLE_DROP_OLDEST_SLOWPATH /// \brief Perform requested drops of oldest items (rare path). void handle_drop_requests_() { { @@ -513,7 +484,6 @@ namespace logit { namespace detail { } m_drop_cv.notify_all(); } -#endif #endif /// \brief Private constructor to enforce the singleton pattern. @@ -530,11 +500,9 @@ namespace logit { namespace detail { m_overflow_policy(QueuePolicy::Block), m_dropped_tasks(0), m_active_tasks(0), - m_mpsc_queue(m_default_ring_cap) -#ifdef LOGIT_ENABLE_DROP_OLDEST_SLOWPATH - , m_drop_requested(0), + m_mpsc_queue(m_default_ring_cap), + m_drop_requested(0), m_drop_done(0) -#endif #endif { m_worker_thread = std::thread(&TaskExecutor::worker_function, this); From fc9404b08109ad92936753fc0457fb4b07a9a23e Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 03:56:59 +0300 Subject: [PATCH 14/18] fix(task-executor): rely on worker drop completion --- .../logit_cpp/logit/detail/TaskExecutor.hpp | 56 +++++++++---------- 1 file changed, 26 insertions(+), 30 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index ac14a10..813aef2 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -218,45 +218,41 @@ namespace logit { namespace detail { } case QueuePolicy::DropOldest: { std::size_t target = 0; - bool request_completed = false; + bool drop_completed = false; + bool drop_canceled = false; { std::unique_lock lk(m_drop_mutex); target = ++m_drop_requested; m_cv.notify_one(); - request_completed = m_drop_cv.wait_for( + drop_completed = m_drop_cv.wait_for( lk, std::chrono::milliseconds(2), [this, target] { return m_drop_done >= target; }); - } - auto finalize_request = [this, target]() { - std::unique_lock lk(m_drop_mutex); - if (m_drop_done < target) { - m_drop_done = target; - lk.unlock(); - m_drop_cv.notify_all(); + if (!drop_completed && m_drop_requested == target) { + --m_drop_requested; + drop_canceled = true; } - }; + } - if (m_mpsc_queue.try_push(local_task)) { - if (!request_completed) { - finalize_request(); - } - m_cv.notify_one(); + if (drop_canceled) { + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); return; } - if (!request_completed) { - finalize_request(); + if (m_mpsc_queue.try_push(local_task)) { + m_cv.notify_one(); + return; } - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - return; + std::unique_lock lk(m_cv_mutex); + m_cv.wait_for(lk, std::chrono::microseconds(200)); + break; } } } #else std::unique_lock lock(m_queue_mutex); - if (m_stop_flag.load(std::memory_order_relaxed)) return; + if (m_stop_flag.load(std::memory_order_acquire)) return; if (m_max_queue_size > 0 && m_tasks_queue.size() >= m_max_queue_size) { switch (m_overflow_policy.load(std::memory_order_relaxed)) { case QueuePolicy::DropNewest: @@ -271,9 +267,9 @@ namespace logit { namespace detail { case QueuePolicy::Block: m_queue_condition.wait(lock, [this]() { return m_tasks_queue.size() < m_max_queue_size || - m_stop_flag.load(std::memory_order_relaxed); + m_stop_flag.load(std::memory_order_acquire); }); - if (m_stop_flag.load(std::memory_order_relaxed)) return; + if (m_stop_flag.load(std::memory_order_acquire)) return; break; } } @@ -290,14 +286,14 @@ namespace logit { namespace detail { 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_relaxed)); + m_stop_flag.load(std::memory_order_acquire)); }); #else std::unique_lock lock(m_queue_mutex); m_queue_condition.wait(lock, [this]() { return ((m_tasks_queue.empty() && m_active_tasks.load(std::memory_order_relaxed) == 0) || - m_stop_flag.load(std::memory_order_relaxed)); + m_stop_flag.load(std::memory_order_acquire)); }); #endif } @@ -308,7 +304,7 @@ namespace logit { namespace detail { #ifdef LOGIT_USE_MPSC_RING { std::lock_guard lock(m_queue_mutex); - m_stop_flag.store(true, std::memory_order_relaxed); + m_stop_flag.store(true, std::memory_order_release); } m_cv.notify_all(); m_queue_condition.notify_all(); @@ -317,7 +313,7 @@ namespace logit { namespace detail { } #else std::unique_lock lock(m_queue_mutex); - m_stop_flag.store(true, std::memory_order_relaxed); + m_stop_flag.store(true, std::memory_order_release); lock.unlock(); m_queue_condition.notify_all(); if (m_worker_thread.joinable()) { @@ -404,9 +400,9 @@ namespace logit { namespace detail { std::function task; std::unique_lock lock(m_queue_mutex); m_queue_condition.wait(lock, [this]() { - return !m_tasks_queue.empty() || m_stop_flag.load(std::memory_order_relaxed); + return !m_tasks_queue.empty() || m_stop_flag.load(std::memory_order_acquire); }); - if (m_stop_flag.load(std::memory_order_relaxed) && m_tasks_queue.empty()) { + if (m_stop_flag.load(std::memory_order_acquire) && m_tasks_queue.empty()) { break; } task = std::move(m_tasks_queue.front()); @@ -442,14 +438,14 @@ namespace logit { namespace detail { std::unique_lock lock(m_queue_mutex); m_queue_condition.notify_all(); m_cv.notify_all(); - if (m_stop_flag.load(std::memory_order_relaxed)) { + if (m_stop_flag.load(std::memory_order_acquire)) { break; } } if (!drained_any) { std::unique_lock lk(m_cv_mutex); - if (m_stop_flag.load(std::memory_order_relaxed) && queue_empty_()) { + if (m_stop_flag.load(std::memory_order_acquire) && queue_empty_()) { break; } m_cv.wait_for(lk, std::chrono::milliseconds(1)); From b0b51e39f47fd1a00644aaad02e36b9d4cd2fb21 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 04:07:25 +0300 Subject: [PATCH 15/18] refactor: update TaskExecutor.hpp --- .../logit_cpp/logit/detail/TaskExecutor.hpp | 773 ++++++++---------- 1 file changed, 349 insertions(+), 424 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 813aef2..07d7329 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -8,516 +8,441 @@ #include #include #if defined(__EMSCRIPTEN__) && !defined(__EMSCRIPTEN_PTHREADS__) -#include -#include -#include + #include + #include + #include #else -#include -#include -#include -#include -#include + #include + #include + #include + #include + #include #endif // Enable lock-free MPSC ring integration (non-Emscripten) by defining: // #define LOGIT_USE_MPSC_RING #if !defined(__EMSCRIPTEN__) || defined(__EMSCRIPTEN_PTHREADS__) -# ifdef LOGIT_USE_MPSC_RING -# include "MpscRingAny.hpp" -# endif + #ifdef LOGIT_USE_MPSC_RING + #include "MpscRingAny.hpp" + #endif #endif namespace logit { namespace detail { - /// \brief Queue overflow handling policy. - enum class QueuePolicy { DropNewest, DropOldest, Block }; +/// \brief Queue overflow handling policy. +enum class QueuePolicy { DropNewest, DropOldest, Block }; #if defined(__EMSCRIPTEN__) && !defined(__EMSCRIPTEN_PTHREADS__) - /// \class TaskExecutor - /// \brief Simplified task executor for single-threaded Emscripten builds. - /// \thread_safety Not thread-safe. - class TaskExecutor { - public: - /// \brief Obtain singleton instance. - /// \return Global executor. - static TaskExecutor& get_instance() { - static TaskExecutor instance; - return instance; - } - - /// \brief Enqueue task for later execution. - /// \param task Callable to execute. - void add_task(std::function task) { - if (!task) return; - bool schedule = false; - for (;;) { - schedule = false; - { - std::lock_guard lk(m_mutex); - if (m_max_queue_size > 0 && m_tasks.size() >= m_max_queue_size) { - switch (m_overflow_policy) { - case QueuePolicy::DropNewest: +/// \class TaskExecutor +/// \brief Simplified task executor for single-threaded Emscripten builds. +/// \thread_safety Not thread-safe. +class TaskExecutor { +public: + static TaskExecutor& get_instance() { + static TaskExecutor instance; + return instance; + } + + void add_task(std::function task) { + if (!task) return; + bool schedule = false; + for (;;) { + schedule = false; + { + std::lock_guard lk(m_mutex); + if (m_max_queue_size > 0 && m_tasks.size() >= m_max_queue_size) { + switch (m_overflow_policy) { + case QueuePolicy::DropNewest: + ++m_dropped_tasks; + return; + case QueuePolicy::DropOldest: + if (!m_tasks.empty()) { + m_tasks.pop_front(); ++m_dropped_tasks; - return; - case QueuePolicy::DropOldest: - if (!m_tasks.empty()) { - m_tasks.pop_front(); - ++m_dropped_tasks; - } - break; - case QueuePolicy::Block: - break; // handled after unlocking - } - if (m_overflow_policy == QueuePolicy::Block && m_tasks.size() >= m_max_queue_size) { - // fall through to drain outside lock - } else { - m_tasks.emplace_back(std::move(task)); - schedule = !m_scheduled; - m_scheduled = m_scheduled || schedule; + } break; - } + case QueuePolicy::Block: + break; // handled after unlocking + } + if (m_overflow_policy == QueuePolicy::Block && m_tasks.size() >= m_max_queue_size) { + // fall through to drain outside lock } else { m_tasks.emplace_back(std::move(task)); schedule = !m_scheduled; m_scheduled = m_scheduled || schedule; break; } + } else { + m_tasks.emplace_back(std::move(task)); + schedule = !m_scheduled; + m_scheduled = m_scheduled || schedule; + break; } - drain(); - } - if (schedule) { - emscripten_async_call(&TaskExecutor::drain_thunk, this, 0); } + drain(); } - - /// \brief Run all queued tasks. - void wait() { drain(); } - - /// \brief Drain queue without scheduling new tasks. - void shutdown() { drain(); } - - /// \brief Set maximum queue size. - /// \param size Number of tasks allowed (0 for unlimited). - void set_max_queue_size(std::size_t size) { - std::lock_guard lk(m_mutex); - m_max_queue_size = size; - } - - /// \brief Set overflow handling policy. - /// \param policy Policy to apply. - void set_queue_policy(QueuePolicy policy) { - std::lock_guard lk(m_mutex); - m_overflow_policy = policy; + if (schedule) { + emscripten_async_call(&TaskExecutor::drain_thunk, this, 0); } - - /// \brief Retrieve the number of tasks dropped due to overflow. - /// \return Total count of dropped tasks. - std::size_t dropped_tasks() const noexcept { - return m_dropped_tasks.load(std::memory_order_relaxed); - } - - /// \brief Reset the dropped tasks counter back to zero. - void reset_dropped_tasks() noexcept { - m_dropped_tasks.store(0, std::memory_order_relaxed); - } - - private: - TaskExecutor() - : m_max_queue_size(0), - m_overflow_policy(QueuePolicy::Block), - m_dropped_tasks(0), - m_scheduled(false) {} - ~TaskExecutor() = default; - TaskExecutor(const TaskExecutor&) = delete; - TaskExecutor& operator=(const TaskExecutor&) = delete; - TaskExecutor(TaskExecutor&&) = delete; - TaskExecutor& operator=(TaskExecutor&&) = delete; - - std::deque> m_tasks; - std::mutex m_mutex; - std::size_t m_max_queue_size; - QueuePolicy m_overflow_policy; - std::atomic m_dropped_tasks; - bool m_scheduled; - - static void drain_thunk(void* arg) { - static_cast(arg)->drain(); + } + + void wait() { drain(); } + void shutdown() { drain(); } + + void set_max_queue_size(std::size_t size) { + std::lock_guard lk(m_mutex); + m_max_queue_size = size; + } + + void set_queue_policy(QueuePolicy policy) { + std::lock_guard lk(m_mutex); + m_overflow_policy = policy; + } + + std::size_t dropped_tasks() const noexcept { + return m_dropped_tasks.load(std::memory_order_relaxed); + } + void reset_dropped_tasks() noexcept { + m_dropped_tasks.store(0, std::memory_order_relaxed); + } + +private: + TaskExecutor() + : m_max_queue_size(0), + m_overflow_policy(QueuePolicy::Block), + m_dropped_tasks(0), + m_scheduled(false) {} + ~TaskExecutor() = default; + TaskExecutor(const TaskExecutor&) = delete; + TaskExecutor& operator=(const TaskExecutor&) = delete; + TaskExecutor(TaskExecutor&&) = delete; + TaskExecutor& operator=(TaskExecutor&&) = delete; + + std::deque> m_tasks; + std::mutex m_mutex; + std::size_t m_max_queue_size; + QueuePolicy m_overflow_policy; + std::atomic m_dropped_tasks; + bool m_scheduled; + + static void drain_thunk(void* arg) { + static_cast(arg)->drain(); + } + + void drain() { + for (;;) { + std::function task; + { + std::lock_guard lk(m_mutex); + if (m_tasks.empty()) { + m_scheduled = false; + break; + } + task = std::move(m_tasks.front()); + m_tasks.pop_front(); + } + task(); } - - void drain() { - for (;;) { - std::function task; - { - std::lock_guard lk(m_mutex); - if (m_tasks.empty()) { - m_scheduled = false; - break; + } +}; + +#else // !Emscripten or pthreads + +/// \class TaskExecutor +/// \brief A thread-safe task executor that processes tasks in a dedicated worker thread. +/// \thread_safety Thread-safe. +class TaskExecutor { +public: + /// Singleton (сохраняем вашу реализацию с new). + static TaskExecutor& get_instance() { + static TaskExecutor* instance = new TaskExecutor(); + return *instance; + } + + /// Добавить задачу. + void add_task(std::function task) { + if (!task) return; +#ifndef LOGIT_USE_MPSC_RING + std::unique_lock lock(m_queue_mutex); + if (m_stop_flag.load(std::memory_order_acquire)) return; + if (m_max_queue_size > 0 && m_tasks_queue.size() >= m_max_queue_size) { + switch (m_overflow_policy.load(std::memory_order_relaxed)) { + case QueuePolicy::DropNewest: + ++m_dropped_tasks; + return; + case QueuePolicy::DropOldest: + if (!m_tasks_queue.empty()) { + m_tasks_queue.pop_front(); + ++m_dropped_tasks; } - task = std::move(m_tasks.front()); - m_tasks.pop_front(); - } - task(); + break; + case QueuePolicy::Block: + m_queue_condition.wait(lock, [this]() { + return m_tasks_queue.size() < m_max_queue_size || + m_stop_flag.load(std::memory_order_acquire); + }); + if (m_stop_flag.load(std::memory_order_acquire)) return; + break; } } - }; - + m_tasks_queue.push_back(std::move(task)); + lock.unlock(); + m_queue_condition.notify_one(); #else - - /// \class TaskExecutor - /// \brief A thread-safe task executor that processes tasks in a dedicated worker thread. - /// \thread_safety Thread-safe. - class TaskExecutor { - public: - /// \brief Get the singleton instance of the TaskExecutor. - /// \return A reference to the single instance of `TaskExecutor`. - static TaskExecutor& get_instance() { - static TaskExecutor* instance = new TaskExecutor(); - return *instance; + if (m_stop_flag.load(std::memory_order_acquire)) { + return; } - /// \brief Adds a task to the queue in a thread-safe manner. - /// \param task A function or lambda with no arguments to be executed asynchronously. - void add_task(std::function task) { - if (!task) return; -#ifdef LOGIT_USE_MPSC_RING + std::function local_task = std::move(task); + + for (;;) { if (m_stop_flag.load(std::memory_order_acquire)) { return; } - std::function local_task = std::move(task); - - for (;;) { - if (m_stop_flag.load(std::memory_order_acquire)) { - return; - } + const auto policy = m_overflow_policy.load(std::memory_order_relaxed); - const auto policy = m_overflow_policy.load(std::memory_order_relaxed); + // Реальное backpressure: учитываем "висящие" задачи. + if (policy == QueuePolicy::Block && + m_max_queue_size > 0 && + m_active_tasks.load(std::memory_order_relaxed) >= m_max_queue_size) + { + std::unique_lock lk(m_cv_mutex); + m_cv.wait_for(lk, std::chrono::microseconds(200)); + continue; + } - if (policy == QueuePolicy::Block && m_max_queue_size > 0 && - m_active_tasks.load(std::memory_order_relaxed) >= m_max_queue_size) { - std::unique_lock lk(m_cv_mutex); - m_cv.wait_for(lk, std::chrono::microseconds(200)); - continue; - } + // Пытаемся положить в кольцо. + if (m_mpsc_queue.try_push(local_task)) { + m_cv.notify_one(); // разбудить воркера + return; + } - if (m_mpsc_queue.try_push(local_task)) { - m_cv.notify_one(); + // Переполнение — применяем политику. + switch (policy) { + case QueuePolicy::DropNewest: + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); return; - } - switch (policy) { - case QueuePolicy::DropNewest: - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - return; - case QueuePolicy::Block: { - std::unique_lock lk(m_cv_mutex); - m_cv.wait_for(lk, std::chrono::microseconds(200)); - break; - } - case QueuePolicy::DropOldest: { - std::size_t target = 0; - bool drop_completed = false; - bool drop_canceled = false; - { - std::unique_lock lk(m_drop_mutex); - target = ++m_drop_requested; - m_cv.notify_one(); - drop_completed = m_drop_cv.wait_for( - lk, std::chrono::milliseconds(2), - [this, target] { return m_drop_done >= target; }); - - if (!drop_completed && m_drop_requested == target) { - --m_drop_requested; - drop_canceled = true; - } - } - - if (drop_canceled) { - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - return; - } - - if (m_mpsc_queue.try_push(local_task)) { - m_cv.notify_one(); - return; - } + case QueuePolicy::DropOldest: + // Безопасная реализация под MPSC: дропаем входящий. + // Это сохраняет порядок и исключает дедлоки при gate. + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + return; - std::unique_lock lk(m_cv_mutex); - m_cv.wait_for(lk, std::chrono::microseconds(200)); - break; - } - } - } -#else - std::unique_lock lock(m_queue_mutex); - if (m_stop_flag.load(std::memory_order_acquire)) return; - if (m_max_queue_size > 0 && m_tasks_queue.size() >= m_max_queue_size) { - switch (m_overflow_policy.load(std::memory_order_relaxed)) { - case QueuePolicy::DropNewest: - ++m_dropped_tasks; - return; - case QueuePolicy::DropOldest: - if (!m_tasks_queue.empty()) { - m_tasks_queue.pop_front(); - ++m_dropped_tasks; - } - break; - case QueuePolicy::Block: - m_queue_condition.wait(lock, [this]() { - return m_tasks_queue.size() < m_max_queue_size || - m_stop_flag.load(std::memory_order_acquire); - }); - if (m_stop_flag.load(std::memory_order_acquire)) return; - break; + case QueuePolicy::Block: { + std::unique_lock lk(m_cv_mutex); + m_cv.wait_for(lk, std::chrono::microseconds(200)); + break; } } - m_tasks_queue.push_back(std::move(task)); - lock.unlock(); - m_queue_condition.notify_one(); -#endif } +#endif + } - /// \brief Waits for all tasks in the queue to be processed. - void wait() { -#ifdef LOGIT_USE_MPSC_RING - 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) || + /// Дождаться опустошения. + void wait() { +#ifndef LOGIT_USE_MPSC_RING + std::unique_lock lock(m_queue_mutex); + m_queue_condition.wait(lock, [this]() { + return ((m_tasks_queue.empty() && + m_active_tasks.load(std::memory_order_relaxed) == 0) || m_stop_flag.load(std::memory_order_acquire)); - }); + }); #else - std::unique_lock lock(m_queue_mutex); - m_queue_condition.wait(lock, [this]() { - return ((m_tasks_queue.empty() && - m_active_tasks.load(std::memory_order_relaxed) == 0) || + 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)); - }); + }); #endif - } + } - /// \brief Shuts down the TaskExecutor by stopping the worker thread. - /// \details This method signals the worker thread to stop and then joins it. - void shutdown() { -#ifdef LOGIT_USE_MPSC_RING - { - std::lock_guard lock(m_queue_mutex); - m_stop_flag.store(true, std::memory_order_release); - } - m_cv.notify_all(); - m_queue_condition.notify_all(); - if (m_worker_thread.joinable()) { - m_worker_thread.join(); - } + /// Остановить воркер. + void shutdown() { +#ifndef LOGIT_USE_MPSC_RING + std::unique_lock lock(m_queue_mutex); + m_stop_flag.store(true, std::memory_order_release); + lock.unlock(); + m_queue_condition.notify_all(); + if (m_worker_thread.joinable()) { + m_worker_thread.join(); + } #else - std::unique_lock lock(m_queue_mutex); + { + std::lock_guard lock(m_queue_mutex); m_stop_flag.store(true, std::memory_order_release); - lock.unlock(); - m_queue_condition.notify_all(); - if (m_worker_thread.joinable()) { - m_worker_thread.join(); - } -#endif } + m_cv.notify_all(); + m_queue_condition.notify_all(); + if (m_worker_thread.joinable()) { + m_worker_thread.join(); + } +#endif + } - /// \brief Sets the maximum size of the task queue. - /// \param size Maximum number of tasks in the queue (0 for unlimited). - void set_max_queue_size(std::size_t size) { -#ifdef LOGIT_USE_MPSC_RING - wait(); - std::lock_guard lk(m_queue_mutex); - m_max_queue_size = size; - const std::size_t cap = - (m_max_queue_size == 0 ? m_default_ring_cap : m_max_queue_size); - m_mpsc_queue = MpscRingAny>(cap); - m_cv.notify_all(); + /// Изменить ёмкость очереди. + void set_max_queue_size(std::size_t size) { +#ifndef LOGIT_USE_MPSC_RING + std::lock_guard lock(m_queue_mutex); + m_max_queue_size = size; #else - std::lock_guard lock(m_queue_mutex); - m_max_queue_size = size; + // Гарантируем пустоту, затем пересоздаём кольцо с новой ёмкостью. + wait(); + std::lock_guard lk(m_queue_mutex); + m_max_queue_size = size; + const std::size_t cap = + (m_max_queue_size == 0 ? m_default_ring_cap : m_max_queue_size); + m_mpsc_queue = MpscRingAny>(cap); + m_cv.notify_all(); #endif - } - - /// \brief Sets the behavior when the queue is full. - /// \param policy QueuePolicy::DropNewest to discard the incoming task, - /// QueuePolicy::DropOldest to discard the oldest task, - /// or QueuePolicy::Block to wait. - void set_queue_policy(QueuePolicy policy) { - std::lock_guard lock(m_queue_mutex); - m_overflow_policy.store(policy, std::memory_order_relaxed); - } + } + + /// Политика переполнения. + void set_queue_policy(QueuePolicy policy) { + std::lock_guard lock(m_queue_mutex); + m_overflow_policy.store(policy, std::memory_order_relaxed); + } + + std::size_t dropped_tasks() const noexcept { + return m_dropped_tasks.load(std::memory_order_relaxed); + } + void reset_dropped_tasks() noexcept { + m_dropped_tasks.store(0, std::memory_order_relaxed); + } + +private: +#ifndef LOGIT_USE_MPSC_RING + std::deque> m_tasks_queue; + mutable std::mutex m_queue_mutex; + std::condition_variable m_queue_condition; + std::thread m_worker_thread; + std::atomic m_stop_flag; + std::size_t m_max_queue_size; + std::atomic m_overflow_policy; + std::atomic m_dropped_tasks; + std::atomic m_active_tasks; +#else + mutable std::mutex m_queue_mutex; ///< Для wait()/смены политики. + std::condition_variable m_queue_condition; ///< Будим wait() на полном drain. - /// \brief Returns the number of tasks dropped because of overflow. - /// \return Total dropped tasks observed so far. - std::size_t dropped_tasks() const noexcept { - return m_dropped_tasks.load(std::memory_order_relaxed); - } + std::condition_variable m_cv; ///< Будим воркер / продюсеров. + std::mutex m_cv_mutex; ///< Сон продюсеров/воркера. - /// \brief Resets the dropped tasks counter back to zero. - void reset_dropped_tasks() noexcept { - m_dropped_tasks.store(0, std::memory_order_relaxed); - } + std::thread m_worker_thread; + std::atomic m_stop_flag; + std::size_t m_max_queue_size; + std::atomic m_overflow_policy; + std::atomic m_dropped_tasks; + std::atomic m_active_tasks; - private: -#ifndef LOGIT_USE_MPSC_RING - std::deque> m_tasks_queue; ///< Queue holding tasks to be executed. - mutable std::mutex m_queue_mutex; ///< Mutex to protect access to the task queue. - std::condition_variable m_queue_condition; ///< Condition variable to signal task availability. - std::thread m_worker_thread; ///< Worker thread for executing tasks. - std::atomic m_stop_flag; ///< Flag indicating if the worker thread should stop. - std::size_t m_max_queue_size; ///< Maximum number of tasks in the queue (0 for unlimited). - std::atomic m_overflow_policy; ///< Policy for handling queue overflow. - std::atomic m_dropped_tasks; ///< Number of discarded tasks due to overflow. - std::atomic m_active_tasks; ///< Number of tasks currently running. -#else - mutable std::mutex m_queue_mutex; ///< Used only for wait()/policy changes. - std::condition_variable m_queue_condition; ///< Notifies waiters on full drain. - - std::condition_variable m_cv; ///< Wake-up for worker on push. - std::mutex m_cv_mutex; ///< Sleep mutex for worker waits. - - std::thread m_worker_thread; ///< Worker thread for executing tasks. - std::atomic m_stop_flag; ///< Flag indicating if the worker thread should stop. - std::size_t m_max_queue_size; ///< Maximum number of tasks requested by user. - std::atomic m_overflow_policy; ///< Policy for handling queue overflow. - std::atomic m_dropped_tasks; ///< Number of discarded tasks due to overflow. - std::atomic m_active_tasks; ///< Number of tasks currently running. - - const std::size_t m_default_ring_cap = LOGIT_TASK_EXECUTOR_DEFAULT_RING_CAPACITY; ///< Default capacity when unlimited requested. - MpscRingAny> m_mpsc_queue; ///< Lock-free bounded MPSC ring. - - std::mutex m_drop_mutex; ///< Coordinates DropOldest slow-path. - std::condition_variable m_drop_cv; ///< Producer waits for confirmation. - std::size_t m_drop_requested; ///< Number of requested drops. - std::size_t m_drop_done; ///< Number of drops completed. + const std::size_t m_default_ring_cap = LOGIT_TASK_EXECUTOR_DEFAULT_RING_CAPACITY; + MpscRingAny> m_mpsc_queue; #endif - /// \brief The worker thread function that processes tasks from the queue. - void worker_function() { + void worker_function() { #ifndef LOGIT_USE_MPSC_RING - for (;;) { - std::function task; - std::unique_lock lock(m_queue_mutex); - m_queue_condition.wait(lock, [this]() { - return !m_tasks_queue.empty() || m_stop_flag.load(std::memory_order_acquire); - }); - if (m_stop_flag.load(std::memory_order_acquire) && m_tasks_queue.empty()) { - break; - } - task = std::move(m_tasks_queue.front()); - m_tasks_queue.pop_front(); - m_active_tasks.fetch_add(1, std::memory_order_relaxed); - lock.unlock(); - m_queue_condition.notify_one(); - task(); - lock.lock(); - m_active_tasks.fetch_sub(1, std::memory_order_relaxed); - if (m_tasks_queue.empty() && m_active_tasks.load(std::memory_order_relaxed) == 0) { - m_queue_condition.notify_all(); - } - lock.unlock(); + for (;;) { + std::function task; + std::unique_lock lock(m_queue_mutex); + m_queue_condition.wait(lock, [this]() { + return !m_tasks_queue.empty() || m_stop_flag.load(std::memory_order_acquire); + }); + if (m_stop_flag.load(std::memory_order_acquire) && m_tasks_queue.empty()) { + break; } -#else - for (;;) { - bool drained_any = false; - std::function task; - - int budget = 2048; - while (budget-- && m_mpsc_queue.try_pop(task)) { - drained_any = true; - m_active_tasks.fetch_add(1, std::memory_order_relaxed); - task(); - m_active_tasks.fetch_sub(1, std::memory_order_relaxed); - m_cv.notify_one(); - } - - handle_drop_requests_(); + task = std::move(m_tasks_queue.front()); + m_tasks_queue.pop_front(); + m_active_tasks.fetch_add(1, std::memory_order_relaxed); + lock.unlock(); + m_queue_condition.notify_one(); - if (queue_empty_() && m_active_tasks.load(std::memory_order_relaxed) == 0) { - std::unique_lock lock(m_queue_mutex); - m_queue_condition.notify_all(); - m_cv.notify_all(); - if (m_stop_flag.load(std::memory_order_acquire)) { - break; - } - } + task(); - if (!drained_any) { - std::unique_lock lk(m_cv_mutex); - if (m_stop_flag.load(std::memory_order_acquire) && queue_empty_()) { - break; - } - m_cv.wait_for(lk, std::chrono::milliseconds(1)); - } + lock.lock(); + m_active_tasks.fetch_sub(1, std::memory_order_relaxed); + if (m_tasks_queue.empty() && m_active_tasks.load(std::memory_order_relaxed) == 0) { + m_queue_condition.notify_all(); } -#endif + lock.unlock(); } +#else + for (;;) { + bool drained_any = false; + std::function task; -#ifdef LOGIT_USE_MPSC_RING - /// \brief Return true if ring appears empty for current consumer position. - bool queue_empty_() const noexcept { - return m_mpsc_queue.empty(); - } + int budget = 2048; + while (budget-- && m_mpsc_queue.try_pop(task)) { + drained_any = true; + m_active_tasks.fetch_add(1, std::memory_order_relaxed); - /// \brief Perform requested drops of oldest items (rare path). - void handle_drop_requests_() { - { - std::lock_guard g(m_drop_mutex); - if (m_drop_done >= m_drop_requested) { - return; + task(); + + m_active_tasks.fetch_sub(1, std::memory_order_relaxed); + m_cv.notify_one(); // освободили in-flight слот + } + + if (queue_empty_() && m_active_tasks.load(std::memory_order_relaxed) == 0) { + std::unique_lock lock(m_queue_mutex); + m_queue_condition.notify_all(); // для wait() + m_cv.notify_all(); // разбудить продюсеров Block + if (m_stop_flag.load(std::memory_order_acquire)) { + break; } } - std::unique_lock lk(m_drop_mutex); - while (m_drop_done < m_drop_requested) { - std::function dummy; - if (m_mpsc_queue.try_pop(dummy)) { - ++m_drop_done; - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - } else { + + if (!drained_any) { + std::unique_lock lk(m_cv_mutex); + if (m_stop_flag.load(std::memory_order_acquire) && queue_empty_()) { break; } + m_cv.wait_for(lk, std::chrono::milliseconds(1)); } - m_drop_cv.notify_all(); } #endif + } + +#ifdef LOGIT_USE_MPSC_RING + bool queue_empty_() const noexcept { + return m_mpsc_queue.empty(); + } +#endif - /// \brief Private constructor to enforce the singleton pattern. - TaskExecutor() + TaskExecutor() #ifndef LOGIT_USE_MPSC_RING - : m_stop_flag(false), - m_max_queue_size(0), - m_overflow_policy(QueuePolicy::Block), - m_dropped_tasks(0), - m_active_tasks(0) + : m_stop_flag(false), + m_max_queue_size(0), + m_overflow_policy(QueuePolicy::Block), + m_dropped_tasks(0), + m_active_tasks(0) #else - : m_stop_flag(false), - m_max_queue_size(0), - m_overflow_policy(QueuePolicy::Block), - m_dropped_tasks(0), - m_active_tasks(0), - m_mpsc_queue(m_default_ring_cap), - m_drop_requested(0), - m_drop_done(0) + : m_stop_flag(false), + m_max_queue_size(0), + m_overflow_policy(QueuePolicy::Block), + m_dropped_tasks(0), + m_active_tasks(0), + m_mpsc_queue(m_default_ring_cap) #endif - { - m_worker_thread = std::thread(&TaskExecutor::worker_function, this); - } + { + m_worker_thread = std::thread(&TaskExecutor::worker_function, this); + } - /// \brief Destructor that stops the worker thread and cleans up resources. - ~TaskExecutor() { - shutdown(); - } + ~TaskExecutor() { + shutdown(); + } - // Delete copy constructor and assignment operators to enforce singleton usage. - TaskExecutor(const TaskExecutor&) = delete; - TaskExecutor& operator=(const TaskExecutor&) = delete; - TaskExecutor(TaskExecutor&&) = delete; - TaskExecutor& operator=(TaskExecutor&&) = delete; - }; + TaskExecutor(const TaskExecutor&) = delete; + TaskExecutor& operator=(const TaskExecutor&) = delete; + TaskExecutor(TaskExecutor&&) = delete; + TaskExecutor& operator=(TaskExecutor&&) = delete; +}; -#endif // defined(__EMSCRIPTEN__) && !defined(__EMSCRIPTEN_PTHREADS__) +#endif // Emscripten split }} // namespace logit::detail #endif // _LOGIT_DETAIL_TASK_EXECUTOR_HPP_INCLUDED + From df75565b6e1abd4c20c26803d5e444c9d45e50a0 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 04:19:51 +0300 Subject: [PATCH 16/18] chore: update ci.yml --- .github/workflows/ci.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index a295401..b009f52 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -24,7 +24,7 @@ jobs: - name: Install run: cmake --install build --prefix install - name: Test - run: ctest --test-dir build + run: ctest --test-dir build --output-on-failure - name: Configure consumer project run: cmake -S tests/install_consumer -B build-consumer -DCMAKE_PREFIX_PATH=${{ github.workspace }}/install -DCMAKE_CXX_STANDARD=${{ matrix.std }} - name: Build consumer project @@ -115,7 +115,7 @@ jobs: - name: Build run: cmake --build build - name: Test - run: ctest --test-dir build + run: ctest --test-dir build --output-on-failure tsan: runs-on: ubuntu-latest @@ -131,7 +131,7 @@ jobs: - name: Build run: cmake --build build - name: Test - run: ctest --test-dir build + run: ctest --test-dir build --output-on-failure vcpkg-install: runs-on: ubuntu-latest From c585e960d32c5ef2cc05221047c95e413f85fbee Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 04:36:28 +0300 Subject: [PATCH 17/18] refactor: update TaskExecutor.hpp --- .../logit_cpp/logit/detail/TaskExecutor.hpp | 759 +++++++++--------- 1 file changed, 390 insertions(+), 369 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 07d7329..2e4119c 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -30,415 +30,436 @@ namespace logit { namespace detail { -/// \brief Queue overflow handling policy. -enum class QueuePolicy { DropNewest, DropOldest, Block }; - -#if defined(__EMSCRIPTEN__) && !defined(__EMSCRIPTEN_PTHREADS__) - -/// \class TaskExecutor -/// \brief Simplified task executor for single-threaded Emscripten builds. -/// \thread_safety Not thread-safe. -class TaskExecutor { -public: - static TaskExecutor& get_instance() { - static TaskExecutor instance; - return instance; - } - - void add_task(std::function task) { - if (!task) return; - bool schedule = false; - for (;;) { - schedule = false; - { - std::lock_guard lk(m_mutex); - if (m_max_queue_size > 0 && m_tasks.size() >= m_max_queue_size) { - switch (m_overflow_policy) { - case QueuePolicy::DropNewest: - ++m_dropped_tasks; - return; - case QueuePolicy::DropOldest: - if (!m_tasks.empty()) { - m_tasks.pop_front(); + /// \brief Queue overflow handling policy. + enum class QueuePolicy { DropNewest, DropOldest, Block }; + +# if defined(__EMSCRIPTEN__) && !defined(__EMSCRIPTEN_PTHREADS__) + + /// \class TaskExecutor + /// \brief Simplified task executor for single-threaded Emscripten builds. + /// \thread_safety Not thread-safe. + class TaskExecutor { + public: + static TaskExecutor& get_instance() { + static TaskExecutor instance; + return instance; + } + + void add_task(std::function task) { + if (!task) return; + bool schedule = false; + for (;;) { + schedule = false; + { + std::lock_guard lk(m_mutex); + if (m_max_queue_size > 0 && m_tasks.size() >= m_max_queue_size) { + switch (m_overflow_policy) { + case QueuePolicy::DropNewest: ++m_dropped_tasks; - } + return; + case QueuePolicy::DropOldest: + if (!m_tasks.empty()) { + m_tasks.pop_front(); + ++m_dropped_tasks; + } + break; + case QueuePolicy::Block: + break; // handled after unlocking + } + if (m_overflow_policy == QueuePolicy::Block && m_tasks.size() >= m_max_queue_size) { + // fall through to drain outside lock + } else { + m_tasks.emplace_back(std::move(task)); + schedule = !m_scheduled; + m_scheduled = m_scheduled || schedule; break; - case QueuePolicy::Block: - break; // handled after unlocking - } - if (m_overflow_policy == QueuePolicy::Block && m_tasks.size() >= m_max_queue_size) { - // fall through to drain outside lock + } } else { m_tasks.emplace_back(std::move(task)); schedule = !m_scheduled; m_scheduled = m_scheduled || schedule; break; } - } else { - m_tasks.emplace_back(std::move(task)); - schedule = !m_scheduled; - m_scheduled = m_scheduled || schedule; - break; } + drain(); + } + if (schedule) { + emscripten_async_call(&TaskExecutor::drain_thunk, this, 0); } - drain(); } - if (schedule) { - emscripten_async_call(&TaskExecutor::drain_thunk, this, 0); + + void wait() { drain(); } + void shutdown() { drain(); } + + void set_max_queue_size(std::size_t size) { + std::lock_guard lk(m_mutex); + m_max_queue_size = size; } - } - - void wait() { drain(); } - void shutdown() { drain(); } - - void set_max_queue_size(std::size_t size) { - std::lock_guard lk(m_mutex); - m_max_queue_size = size; - } - - void set_queue_policy(QueuePolicy policy) { - std::lock_guard lk(m_mutex); - m_overflow_policy = policy; - } - - std::size_t dropped_tasks() const noexcept { - return m_dropped_tasks.load(std::memory_order_relaxed); - } - void reset_dropped_tasks() noexcept { - m_dropped_tasks.store(0, std::memory_order_relaxed); - } - -private: - TaskExecutor() - : m_max_queue_size(0), - m_overflow_policy(QueuePolicy::Block), - m_dropped_tasks(0), - m_scheduled(false) {} - ~TaskExecutor() = default; - TaskExecutor(const TaskExecutor&) = delete; - TaskExecutor& operator=(const TaskExecutor&) = delete; - TaskExecutor(TaskExecutor&&) = delete; - TaskExecutor& operator=(TaskExecutor&&) = delete; - - std::deque> m_tasks; - std::mutex m_mutex; - std::size_t m_max_queue_size; - QueuePolicy m_overflow_policy; - std::atomic m_dropped_tasks; - bool m_scheduled; - - static void drain_thunk(void* arg) { - static_cast(arg)->drain(); - } - - void drain() { - for (;;) { - std::function task; - { - std::lock_guard lk(m_mutex); - if (m_tasks.empty()) { - m_scheduled = false; - break; - } - task = std::move(m_tasks.front()); - m_tasks.pop_front(); - } - task(); + + void set_queue_policy(QueuePolicy policy) { + std::lock_guard lk(m_mutex); + m_overflow_policy = policy; } - } -}; - -#else // !Emscripten or pthreads - -/// \class TaskExecutor -/// \brief A thread-safe task executor that processes tasks in a dedicated worker thread. -/// \thread_safety Thread-safe. -class TaskExecutor { -public: - /// Singleton (сохраняем вашу реализацию с new). - static TaskExecutor& get_instance() { - static TaskExecutor* instance = new TaskExecutor(); - return *instance; - } - - /// Добавить задачу. - void add_task(std::function task) { - if (!task) return; -#ifndef LOGIT_USE_MPSC_RING - std::unique_lock lock(m_queue_mutex); - if (m_stop_flag.load(std::memory_order_acquire)) return; - if (m_max_queue_size > 0 && m_tasks_queue.size() >= m_max_queue_size) { - switch (m_overflow_policy.load(std::memory_order_relaxed)) { - case QueuePolicy::DropNewest: - ++m_dropped_tasks; - return; - case QueuePolicy::DropOldest: - if (!m_tasks_queue.empty()) { - m_tasks_queue.pop_front(); - ++m_dropped_tasks; + + std::size_t dropped_tasks() const noexcept { + return m_dropped_tasks.load(std::memory_order_relaxed); + } + void reset_dropped_tasks() noexcept { + m_dropped_tasks.store(0, std::memory_order_relaxed); + } + + private: + TaskExecutor() + : m_max_queue_size(0), + m_overflow_policy(QueuePolicy::Block), + m_dropped_tasks(0), + m_scheduled(false) {} + ~TaskExecutor() = default; + TaskExecutor(const TaskExecutor&) = delete; + TaskExecutor& operator=(const TaskExecutor&) = delete; + TaskExecutor(TaskExecutor&&) = delete; + TaskExecutor& operator=(TaskExecutor&&) = delete; + + std::deque> m_tasks; + std::mutex m_mutex; + std::size_t m_max_queue_size; + QueuePolicy m_overflow_policy; + std::atomic m_dropped_tasks; + bool m_scheduled; + + static void drain_thunk(void* arg) { + static_cast(arg)->drain(); + } + + void drain() { + for (;;) { + std::function task; + { + std::lock_guard lk(m_mutex); + if (m_tasks.empty()) { + m_scheduled = false; + break; } - break; - case QueuePolicy::Block: - m_queue_condition.wait(lock, [this]() { - return m_tasks_queue.size() < m_max_queue_size || - m_stop_flag.load(std::memory_order_acquire); - }); - if (m_stop_flag.load(std::memory_order_acquire)) return; - break; + task = std::move(m_tasks.front()); + m_tasks.pop_front(); + } + task(); } } - m_tasks_queue.push_back(std::move(task)); - lock.unlock(); - m_queue_condition.notify_one(); -#else - if (m_stop_flag.load(std::memory_order_acquire)) { - return; + }; + +# else // !Emscripten or pthreads + + /// \class TaskExecutor + /// \brief A thread-safe task executor that processes tasks in a dedicated worker thread. + /// \thread_safety Thread-safe. + class TaskExecutor { + public: + /// Singleton (сохраняем вашу реализацию с new). + static TaskExecutor& get_instance() { + static TaskExecutor* instance = new TaskExecutor(); + return *instance; } - - std::function local_task = std::move(task); - - for (;;) { - if (m_stop_flag.load(std::memory_order_acquire)) { - return; - } - - const auto policy = m_overflow_policy.load(std::memory_order_relaxed); - - // Реальное backpressure: учитываем "висящие" задачи. - if (policy == QueuePolicy::Block && - m_max_queue_size > 0 && - m_active_tasks.load(std::memory_order_relaxed) >= m_max_queue_size) - { - std::unique_lock lk(m_cv_mutex); - m_cv.wait_for(lk, std::chrono::microseconds(200)); - continue; + + /// Добавить задачу. + void add_task(std::function task) { + if (!task) return; +# ifndef LOGIT_USE_MPSC_RING + std::unique_lock lock(m_queue_mutex); + if (m_stop_flag.load(std::memory_order_acquire)) return; + if (m_max_queue_size > 0 && m_tasks_queue.size() >= m_max_queue_size) { + switch (m_overflow_policy.load(std::memory_order_relaxed)) { + case QueuePolicy::DropNewest: + ++m_dropped_tasks; + return; + case QueuePolicy::DropOldest: + if (!m_tasks_queue.empty()) { + m_tasks_queue.pop_front(); + ++m_dropped_tasks; + } + break; + case QueuePolicy::Block: + m_queue_condition.wait(lock, [this]() { + return m_tasks_queue.size() < m_max_queue_size || + m_stop_flag.load(std::memory_order_acquire); + }); + if (m_stop_flag.load(std::memory_order_acquire)) return; + break; + } } - - // Пытаемся положить в кольцо. - if (m_mpsc_queue.try_push(local_task)) { - m_cv.notify_one(); // разбудить воркера + m_tasks_queue.push_back(std::move(task)); + lock.unlock(); + m_queue_condition.notify_one(); +# else + if (m_stop_flag.load(std::memory_order_acquire)) { return; } - - // Переполнение — применяем политику. - switch (policy) { - case QueuePolicy::DropNewest: - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); - return; - - case QueuePolicy::DropOldest: - // Безопасная реализация под MPSC: дропаем входящий. - // Это сохраняет порядок и исключает дедлоки при gate. - m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + + std::function local_task = std::move(task); + + for (;;) { + if (m_stop_flag.load(std::memory_order_acquire)) { return; - - case QueuePolicy::Block: { + } + + const auto policy = m_overflow_policy.load(std::memory_order_relaxed); + + // Реальное backpressure: учитываем "висящие" задачи. + if (policy == QueuePolicy::Block && + m_max_queue_size > 0 && + m_active_tasks.load(std::memory_order_relaxed) >= m_max_queue_size) + { std::unique_lock lk(m_cv_mutex); m_cv.wait_for(lk, std::chrono::microseconds(200)); - break; + continue; + } + + // Пытаемся положить в кольцо. + if (m_mpsc_queue.try_push(local_task)) { + m_cv.notify_one(); // разбудить воркера + return; + } + + // Переполнение — применяем политику. + switch (policy) { + case QueuePolicy::DropNewest: + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + return; + + case QueuePolicy::DropOldest: + // Безопасная реализация под MPSC: дропаем входящий. + // Это сохраняет порядок и исключает дедлоки при gate. + m_dropped_tasks.fetch_add(1, std::memory_order_relaxed); + return; + + case QueuePolicy::Block: { + std::unique_lock lk(m_cv_mutex); + m_cv.wait_for(lk, std::chrono::microseconds(200)); + break; + } } } +# endif } -#endif - } - - /// Дождаться опустошения. - void wait() { -#ifndef LOGIT_USE_MPSC_RING - std::unique_lock lock(m_queue_mutex); - m_queue_condition.wait(lock, [this]() { - return ((m_tasks_queue.empty() && - m_active_tasks.load(std::memory_order_relaxed) == 0) || - 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)); - }); -#endif - } - - /// Остановить воркер. - void shutdown() { -#ifndef LOGIT_USE_MPSC_RING - std::unique_lock lock(m_queue_mutex); - m_stop_flag.store(true, std::memory_order_release); - lock.unlock(); - m_queue_condition.notify_all(); - if (m_worker_thread.joinable()) { - m_worker_thread.join(); - } -#else - { - std::lock_guard lock(m_queue_mutex); - m_stop_flag.store(true, std::memory_order_release); - } - m_cv.notify_all(); - m_queue_condition.notify_all(); - if (m_worker_thread.joinable()) { - m_worker_thread.join(); - } -#endif - } - - /// Изменить ёмкость очереди. - void set_max_queue_size(std::size_t size) { -#ifndef LOGIT_USE_MPSC_RING - std::lock_guard lock(m_queue_mutex); - m_max_queue_size = size; -#else - // Гарантируем пустоту, затем пересоздаём кольцо с новой ёмкостью. - wait(); - std::lock_guard lk(m_queue_mutex); - m_max_queue_size = size; - const std::size_t cap = - (m_max_queue_size == 0 ? m_default_ring_cap : m_max_queue_size); - m_mpsc_queue = MpscRingAny>(cap); - m_cv.notify_all(); -#endif - } - - /// Политика переполнения. - void set_queue_policy(QueuePolicy policy) { - std::lock_guard lock(m_queue_mutex); - m_overflow_policy.store(policy, std::memory_order_relaxed); - } - - std::size_t dropped_tasks() const noexcept { - return m_dropped_tasks.load(std::memory_order_relaxed); - } - void reset_dropped_tasks() noexcept { - m_dropped_tasks.store(0, std::memory_order_relaxed); - } - -private: -#ifndef LOGIT_USE_MPSC_RING - std::deque> m_tasks_queue; - mutable std::mutex m_queue_mutex; - std::condition_variable m_queue_condition; - std::thread m_worker_thread; - std::atomic m_stop_flag; - std::size_t m_max_queue_size; - std::atomic m_overflow_policy; - std::atomic m_dropped_tasks; - std::atomic m_active_tasks; -#else - mutable std::mutex m_queue_mutex; ///< Для wait()/смены политики. - std::condition_variable m_queue_condition; ///< Будим wait() на полном drain. - - std::condition_variable m_cv; ///< Будим воркер / продюсеров. - std::mutex m_cv_mutex; ///< Сон продюсеров/воркера. - - std::thread m_worker_thread; - std::atomic m_stop_flag; - std::size_t m_max_queue_size; - std::atomic m_overflow_policy; - std::atomic m_dropped_tasks; - std::atomic m_active_tasks; - - const std::size_t m_default_ring_cap = LOGIT_TASK_EXECUTOR_DEFAULT_RING_CAPACITY; - MpscRingAny> m_mpsc_queue; -#endif - - void worker_function() { -#ifndef LOGIT_USE_MPSC_RING - for (;;) { - std::function task; + + /// Дождаться опустошения. + void wait() { +# ifndef LOGIT_USE_MPSC_RING std::unique_lock lock(m_queue_mutex); m_queue_condition.wait(lock, [this]() { - return !m_tasks_queue.empty() || m_stop_flag.load(std::memory_order_acquire); + return ((m_tasks_queue.empty() && + m_active_tasks.load(std::memory_order_relaxed) == 0) || + m_stop_flag.load(std::memory_order_acquire)); }); - if (m_stop_flag.load(std::memory_order_acquire) && m_tasks_queue.empty()) { - break; - } - task = std::move(m_tasks_queue.front()); - m_tasks_queue.pop_front(); - m_active_tasks.fetch_add(1, std::memory_order_relaxed); +# 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)); + }); +# endif + } + + /// Остановить воркер. + void shutdown() { +# ifndef LOGIT_USE_MPSC_RING + std::unique_lock lock(m_queue_mutex); + m_stop_flag.store(true, std::memory_order_release); lock.unlock(); - m_queue_condition.notify_one(); - - task(); - - lock.lock(); - m_active_tasks.fetch_sub(1, std::memory_order_relaxed); - if (m_tasks_queue.empty() && m_active_tasks.load(std::memory_order_relaxed) == 0) { - m_queue_condition.notify_all(); + m_queue_condition.notify_all(); + if (m_worker_thread.joinable()) { + m_worker_thread.join(); } - lock.unlock(); +# else + { + std::lock_guard lock(m_queue_mutex); + m_stop_flag.store(true, std::memory_order_release); + } + m_cv.notify_all(); + m_queue_condition.notify_all(); + if (m_worker_thread.joinable()) { + m_worker_thread.join(); + } +# endif } -#else - for (;;) { - bool drained_any = false; - std::function task; - - int budget = 2048; - while (budget-- && m_mpsc_queue.try_pop(task)) { - drained_any = true; - m_active_tasks.fetch_add(1, std::memory_order_relaxed); - - task(); - - m_active_tasks.fetch_sub(1, std::memory_order_relaxed); - m_cv.notify_one(); // освободили in-flight слот + + /// Изменить ёмкость очереди. + void set_max_queue_size(std::size_t size) { +# ifdef LOGIT_USE_MPSC_RING + // Дождаться опустошения очереди + wait(); + + // Акуратно остановить воркер и дождаться его завершения, чтобы он не трогал m_mpsc_queue, пока мы его меняем. + std::unique_lock lk(m_queue_mutex); + m_stop_flag.store(true, std::memory_order_relaxed); + lk.unlock(); + + m_cv.notify_all(); + m_queue_condition.notify_all(); + if (m_worker_thread.joinable()) { + m_worker_thread.join(); } - - if (queue_empty_() && m_active_tasks.load(std::memory_order_relaxed) == 0) { + + // Переинициализировать параметры и само кольцо в единственном потоке. + lk.lock(); + m_max_queue_size = size; + const std::size_t cap = + (m_max_queue_size == 0 ? m_default_ring_cap : m_max_queue_size); + m_mpsc_queue = MpscRingAny>(cap); + + // обнулить счётчики (не обязательно, но логично при "чистой" очереди). + m_active_tasks.store(0, std::memory_order_relaxed); + // m_dropped_tasks оставляем как есть — тесты его сами сбрасывают макросом. + lk.unlock(); + + // Снять стоп-флаг и перезапустить воркер. + m_stop_flag.store(false, std::memory_order_relaxed); + m_worker_thread = std::thread(&TaskExecutor::worker_function, this); +# else + std::lock_guard lock(m_queue_mutex); + m_max_queue_size = size; +# endif + } + + /// Политика переполнения. + void set_queue_policy(QueuePolicy policy) { + std::lock_guard lock(m_queue_mutex); + m_overflow_policy.store(policy, std::memory_order_relaxed); + } + + std::size_t dropped_tasks() const noexcept { + return m_dropped_tasks.load(std::memory_order_relaxed); + } + void reset_dropped_tasks() noexcept { + m_dropped_tasks.store(0, std::memory_order_relaxed); + } + + private: + #ifndef LOGIT_USE_MPSC_RING + std::deque> m_tasks_queue; + mutable std::mutex m_queue_mutex; + std::condition_variable m_queue_condition; + std::thread m_worker_thread; + std::atomic m_stop_flag; + std::size_t m_max_queue_size; + std::atomic m_overflow_policy; + std::atomic m_dropped_tasks; + std::atomic m_active_tasks; + #else + mutable std::mutex m_queue_mutex; ///< Для wait()/смены политики. + std::condition_variable m_queue_condition; ///< Будим wait() на полном drain. + + std::condition_variable m_cv; ///< Будим воркер / продюсеров. + std::mutex m_cv_mutex; ///< Сон продюсеров/воркера. + + std::thread m_worker_thread; + std::atomic m_stop_flag; + std::size_t m_max_queue_size; + std::atomic m_overflow_policy; + std::atomic m_dropped_tasks; + std::atomic m_active_tasks; + + const std::size_t m_default_ring_cap = LOGIT_TASK_EXECUTOR_DEFAULT_RING_CAPACITY; + MpscRingAny> m_mpsc_queue; + #endif + + void worker_function() { + #ifndef LOGIT_USE_MPSC_RING + for (;;) { + std::function task; std::unique_lock lock(m_queue_mutex); - m_queue_condition.notify_all(); // для wait() - m_cv.notify_all(); // разбудить продюсеров Block - if (m_stop_flag.load(std::memory_order_acquire)) { + m_queue_condition.wait(lock, [this]() { + return !m_tasks_queue.empty() || m_stop_flag.load(std::memory_order_acquire); + }); + if (m_stop_flag.load(std::memory_order_acquire) && m_tasks_queue.empty()) { break; } + task = std::move(m_tasks_queue.front()); + m_tasks_queue.pop_front(); + m_active_tasks.fetch_add(1, std::memory_order_relaxed); + lock.unlock(); + m_queue_condition.notify_one(); + + task(); + + lock.lock(); + m_active_tasks.fetch_sub(1, std::memory_order_relaxed); + if (m_tasks_queue.empty() && m_active_tasks.load(std::memory_order_relaxed) == 0) { + m_queue_condition.notify_all(); + } + lock.unlock(); } - - if (!drained_any) { - std::unique_lock lk(m_cv_mutex); - if (m_stop_flag.load(std::memory_order_acquire) && queue_empty_()) { - break; + #else + for (;;) { + bool drained_any = false; + std::function task; + + int budget = 2048; + while (budget-- && m_mpsc_queue.try_pop(task)) { + drained_any = true; + m_active_tasks.fetch_add(1, std::memory_order_relaxed); + + task(); + + m_active_tasks.fetch_sub(1, std::memory_order_relaxed); + m_cv.notify_one(); // освободили in-flight слот + } + + if (queue_empty_() && m_active_tasks.load(std::memory_order_relaxed) == 0) { + std::unique_lock lock(m_queue_mutex); + m_queue_condition.notify_all(); // для wait() + m_cv.notify_all(); // разбудить продюсеров Block + if (m_stop_flag.load(std::memory_order_acquire)) { + break; + } + } + + if (!drained_any) { + std::unique_lock lk(m_cv_mutex); + if (m_stop_flag.load(std::memory_order_acquire) && queue_empty_()) { + break; + } + m_cv.wait_for(lk, std::chrono::milliseconds(1)); } - m_cv.wait_for(lk, std::chrono::milliseconds(1)); } + #endif } -#endif - } - -#ifdef LOGIT_USE_MPSC_RING - bool queue_empty_() const noexcept { - return m_mpsc_queue.empty(); - } -#endif - - TaskExecutor() -#ifndef LOGIT_USE_MPSC_RING - : m_stop_flag(false), - m_max_queue_size(0), - m_overflow_policy(QueuePolicy::Block), - m_dropped_tasks(0), - m_active_tasks(0) -#else - : m_stop_flag(false), - m_max_queue_size(0), - m_overflow_policy(QueuePolicy::Block), - m_dropped_tasks(0), - m_active_tasks(0), - m_mpsc_queue(m_default_ring_cap) -#endif - { - m_worker_thread = std::thread(&TaskExecutor::worker_function, this); - } - - ~TaskExecutor() { - shutdown(); - } - - TaskExecutor(const TaskExecutor&) = delete; - TaskExecutor& operator=(const TaskExecutor&) = delete; - TaskExecutor(TaskExecutor&&) = delete; - TaskExecutor& operator=(TaskExecutor&&) = delete; -}; + + #ifdef LOGIT_USE_MPSC_RING + bool queue_empty_() const noexcept { + return m_mpsc_queue.empty(); + } + #endif + + TaskExecutor() + #ifndef LOGIT_USE_MPSC_RING + : m_stop_flag(false), + m_max_queue_size(0), + m_overflow_policy(QueuePolicy::Block), + m_dropped_tasks(0), + m_active_tasks(0) + #else + : m_stop_flag(false), + m_max_queue_size(0), + m_overflow_policy(QueuePolicy::Block), + m_dropped_tasks(0), + m_active_tasks(0), + m_mpsc_queue(m_default_ring_cap) + #endif + { + m_worker_thread = std::thread(&TaskExecutor::worker_function, this); + } + + ~TaskExecutor() { + shutdown(); + } + + TaskExecutor(const TaskExecutor&) = delete; + TaskExecutor& operator=(const TaskExecutor&) = delete; + TaskExecutor(TaskExecutor&&) = delete; + TaskExecutor& operator=(TaskExecutor&&) = delete; + }; #endif // Emscripten split From c6f01eba417ada5724579f84fbbf99c523d2414e Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 18 Sep 2025 04:50:27 +0300 Subject: [PATCH 18/18] refactor: update TaskExecutor.hpp --- .../logit_cpp/logit/detail/TaskExecutor.hpp | 23 ++++++++++++++++--- 1 file changed, 20 insertions(+), 3 deletions(-) diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 2e4119c..fe02d47 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -191,6 +191,12 @@ namespace logit { namespace detail { lock.unlock(); m_queue_condition.notify_one(); # else + // Барьер ресайза: пережидаем «горячее» изменение кольца. + if (m_resizing.load(std::memory_order_acquire)) { + std::unique_lock lk(m_cv_mutex); + m_resize_cv.wait(lk, [this]{ return !m_resizing.load(std::memory_order_acquire); }); + } + if (m_stop_flag.load(std::memory_order_acquire)) { return; } @@ -287,6 +293,9 @@ namespace logit { namespace detail { /// Изменить ёмкость очереди. void set_max_queue_size(std::size_t size) { # ifdef LOGIT_USE_MPSC_RING + // Сигналим продюсерам, чтобы переждали ресайз (до любых ожиданий/стопов). + m_resizing.store(true, std::memory_order_release); + // Дождаться опустошения очереди wait(); @@ -307,15 +316,18 @@ namespace logit { namespace detail { const std::size_t cap = (m_max_queue_size == 0 ? m_default_ring_cap : m_max_queue_size); m_mpsc_queue = MpscRingAny>(cap); - // обнулить счётчики (не обязательно, но логично при "чистой" очереди). m_active_tasks.store(0, std::memory_order_relaxed); // m_dropped_tasks оставляем как есть — тесты его сами сбрасывают макросом. lk.unlock(); - // Снять стоп-флаг и перезапустить воркер. + // Снимаем стоп-флаг, перезапускаем воркер… m_stop_flag.store(false, std::memory_order_relaxed); m_worker_thread = std::thread(&TaskExecutor::worker_function, this); + + // Открываем барьер для продюсеров. + m_resizing.store(false, std::memory_order_release); + m_resize_cv.notify_all(); # else std::lock_guard lock(m_queue_mutex); m_max_queue_size = size; @@ -352,6 +364,9 @@ namespace logit { namespace detail { std::condition_variable m_cv; ///< Будим воркер / продюсеров. std::mutex m_cv_mutex; ///< Сон продюсеров/воркера. + + std::atomic m_resizing; ///< true — идёт ресайз кольца. + std::condition_variable m_resize_cv; ///< Продюсеры ждут окончания ресайза. std::thread m_worker_thread; std::atomic m_stop_flag; @@ -440,7 +455,9 @@ namespace logit { namespace detail { m_dropped_tasks(0), m_active_tasks(0) #else - : m_stop_flag(false), + : m_resizing(false), + m_worker_thread(), + m_stop_flag(false), m_max_queue_size(0), m_overflow_policy(QueuePolicy::Block), m_dropped_tasks(0),