From cc708e6deaa7951a9a0293794d9bfb104e767fb4 Mon Sep 17 00:00:00 2001 From: Aster Seker Date: Thu, 4 Dec 2025 03:08:01 +0300 Subject: [PATCH] fix(bench): retain spdlog payloads during flush Prevent MeasuringSink from deleting payloads while the async thread pool may still dereference them by moving completed entries into a retired list. --- .github/workflows/ci.yml | 31 +++++++ bench/LatencyRecorder.hpp | 20 ++++- bench/adapters/ILoggerAdapter.hpp | 3 + bench/adapters/SpdlogAdapter.cpp | 88 ++++++++++++++++--- bench/adapters/SpdlogAdapter.hpp | 3 + bench/logit_bench.cpp | 76 ++++++++++++++-- docs/TaskExecutor.md | 8 +- .../logit_cpp/logit/detail/MpscRingAny.hpp | 51 ++++++----- .../logit_cpp/logit/detail/TaskExecutor.hpp | 35 ++++++-- 9 files changed, 261 insertions(+), 54 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 1f35483..a512f11 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -147,6 +147,37 @@ jobs: - name: Test run: ctest --test-dir build --output-on-failure + bench-asan: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + with: + submodules: true + - run: git submodule update --init --recursive + - name: Configure benchmarks (asan) + run: | + cmake -S . -B build-bench-asan \ + -DLOGIT_BENCH_ENABLE=ON \ + -DLOGIT_BENCH_WITH_SPDLOG=ON \ + -DCMAKE_BUILD_TYPE=RelWithDebInfo \ + -DCMAKE_CXX_STANDARD=17 \ + -DCMAKE_CXX_FLAGS='-fsanitize=address,undefined -fno-omit-frame-pointer -g' \ + -DCMAKE_EXE_LINKER_FLAGS='-fsanitize=address,undefined' + - name: Build benchmarks (asan) + run: cmake --build build-bench-asan --target logit_bench + - name: Run spdlog async null bench (asan) + timeout-minutes: 10 + env: + LOGIT_BENCH_FILTER_LIB: spdlog + LOGIT_BENCH_FILTER_ASYNC: "1" + LOGIT_BENCH_FILTER_SINK: null + LOGIT_BENCH_FILTER_PRODUCERS: "4" + LOGIT_BENCH_FILTER_BYTES: "40" + LOGIT_BENCH_TOTAL: 200 + LOGIT_BENCH_WARMUP: 20 + LOGIT_BENCH_TIMEOUT_SEC: 120 + run: ./build-bench-asan/logit_bench + vcpkg-install: runs-on: ubuntu-latest env: diff --git a/bench/LatencyRecorder.hpp b/bench/LatencyRecorder.hpp index f25a84e..a9b65b5 100644 --- a/bench/LatencyRecorder.hpp +++ b/bench/LatencyRecorder.hpp @@ -4,7 +4,9 @@ #include #include #include +#include #include +#include #include #include #include @@ -36,7 +38,8 @@ class LatencyRecorder { explicit LatencyRecorder(std::size_t total) : m_values(total), m_expected(total), - m_next_slot(0) {} + m_next_slot(0), + m_completed(0) {} /** * Reserve a slot (if record==true) and capture t0 using steady_clock. @@ -61,14 +64,24 @@ class LatencyRecorder { if (!token.active) return; const auto t1_ns = now(); m_values[token.slot] = t1_ns - token.t0_ns; // distinct slots -> no data race + const auto done = m_completed.fetch_add(1, std::memory_order_acq_rel) + 1; + if (done == m_expected) { + std::lock_guard lk(m_wait_mx); + m_wait_cv.notify_all(); + } } std::size_t recorded() const { return m_next_slot.load(std::memory_order_relaxed); } + void wait_for_all() const { + std::unique_lock lk(m_wait_mx); + m_wait_cv.wait(lk, [&]{ return m_completed.load(std::memory_order_acquire) >= m_expected; }); + } + Summary finalize() const { - if (recorded() != m_expected) { + if (recorded() != m_expected || m_completed.load(std::memory_order_acquire) != m_expected) { throw std::runtime_error("Incomplete latency capture"); } std::vector sorted = m_values; @@ -102,6 +115,9 @@ class LatencyRecorder { std::vector m_values; // preallocated; no reallocation const std::size_t m_expected; // total messages to record std::atomic m_next_slot; + std::atomic m_completed; + mutable std::condition_variable m_wait_cv; + mutable std::mutex m_wait_mx; }; } // namespace logit_bench diff --git a/bench/adapters/ILoggerAdapter.hpp b/bench/adapters/ILoggerAdapter.hpp index 508a4ad..da469c0 100644 --- a/bench/adapters/ILoggerAdapter.hpp +++ b/bench/adapters/ILoggerAdapter.hpp @@ -1,5 +1,6 @@ #pragma once +#include #include #include "../LatencyRecorder.hpp" @@ -18,6 +19,8 @@ class ILoggerAdapter { virtual void log(const LatencyRecorder::Token& token, std::string_view message) = 0; virtual void flush() = 0; + + virtual void set_recorder_handle(std::shared_ptr) {} }; } // namespace logit_bench diff --git a/bench/adapters/SpdlogAdapter.cpp b/bench/adapters/SpdlogAdapter.cpp index 5cb22c1..b312f00 100644 --- a/bench/adapters/SpdlogAdapter.cpp +++ b/bench/adapters/SpdlogAdapter.cpp @@ -5,9 +5,11 @@ #include #include #include +#include #include #include #include +#include #include #include @@ -32,6 +34,11 @@ class SpdlogAdapter::MeasuringSink : public spdlog::sinks::sink { void configure(const Scenario& scenario, LatencyRecorder& recorder) { m_sink = scenario.sink; m_recorder = &recorder; + { + std::lock_guard lock(m_pending_mx); + m_pending.clear(); + m_retired.clear(); + } if (m_sink == SinkKind::File) { std::filesystem::create_directories("bench/results"); std::lock_guard lock(m_mutex); @@ -43,14 +50,20 @@ class SpdlogAdapter::MeasuringSink : public spdlog::sinks::sink { } } + void track_token(const LatencyRecorder::Token& token, std::unique_ptr payload) { + std::lock_guard lock(m_pending_mx); + m_pending.push_back(Pending{std::move(payload), token}); + } + void log(const spdlog::details::log_msg& msg) override { - const auto* payload_ptr = reinterpret_cast(msg.source.funcname); - if (!payload_ptr) { - return; + const char* func = msg.source.funcname; + if (msg.payload.size() == 0 || !func || *func == '\0') { + return; // Flush/control messages have no payload attached. } + + const auto* payload_ptr = reinterpret_cast(func); auto* payload = const_cast(payload_ptr); - consume(*payload); - delete payload; + consume(*payload, payload); } void set_pattern(const std::string&) override {} @@ -64,15 +77,50 @@ class SpdlogAdapter::MeasuringSink : public spdlog::sinks::sink { } } + void complete_pending() { + std::vector pending; + { + std::lock_guard lock(m_pending_mx); + pending.swap(m_pending); + } + std::vector> retired; + retired.reserve(pending.size()); + for (auto& entry : pending) { + if (entry.token.active && m_recorder) { + m_recorder->complete(entry.token); + } + if (entry.payload) { + retired.push_back(std::move(entry.payload)); + } + } + m_retired.insert(m_retired.end(), + std::make_move_iterator(retired.begin()), + std::make_move_iterator(retired.end())); + } + private: - void consume(const MessagePayload& payload) { - if (payload.token.active && m_recorder) { - m_recorder->complete(payload.token); + void consume(const MessagePayload& payload, MessagePayload* payload_ptr) { + LatencyRecorder::Token token = payload.token; + std::unique_ptr owned; + { + std::lock_guard lock(m_pending_mx); + auto it = std::find_if(m_pending.begin(), m_pending.end(), [&](const Pending& p){ return p.payload.get() == payload_ptr; }); + if (it != m_pending.end()) { + token = it->token; + owned = std::move(it->payload); + m_pending.erase(it); + } + } + + const MessagePayload& msg = owned ? *owned : payload; + + if (token.active && m_recorder) { + m_recorder->complete(token); } if (m_sink == SinkKind::File) { std::lock_guard lock(m_mutex); if (m_file.is_open()) { - m_file << payload.text << '\n'; + m_file << msg.text << '\n'; } } } @@ -81,6 +129,13 @@ class SpdlogAdapter::MeasuringSink : public spdlog::sinks::sink { LatencyRecorder* m_recorder = nullptr; std::ofstream m_file; std::mutex m_mutex; + struct Pending { + std::unique_ptr payload; + LatencyRecorder::Token token; + }; + std::vector m_pending; + std::vector> m_retired; + std::mutex m_pending_mx; }; SpdlogAdapter::SpdlogAdapter() = default; @@ -90,6 +145,10 @@ SpdlogAdapter::~SpdlogAdapter() { spdlog::shutdown(); } +void SpdlogAdapter::set_recorder_handle(std::shared_ptr recorder) { + m_recorder_handle = std::move(recorder); +} + void SpdlogAdapter::prepare(const Scenario& scenario, LatencyRecorder& recorder) { m_logger.reset(); m_sink.reset(); @@ -123,11 +182,15 @@ void SpdlogAdapter::log(const LatencyRecorder::Token& token, std::string_view me if (!m_logger) { return; } - auto* payload = new MessagePayload(); + auto payload = std::make_unique(); payload->token = token; payload->text.assign(message.data(), message.size()); - spdlog::source_loc loc{nullptr, 0, reinterpret_cast(payload)}; - m_logger->log(loc, spdlog::level::info, spdlog::string_view_t(payload->text)); + MessagePayload* payload_ptr = payload.get(); + if (m_sink) { + m_sink->track_token(token, std::move(payload)); + } + spdlog::source_loc loc{nullptr, 0, reinterpret_cast(payload_ptr)}; + m_logger->log(loc, spdlog::level::info, spdlog::string_view_t(payload_ptr->text)); } void SpdlogAdapter::flush() { @@ -135,6 +198,7 @@ void SpdlogAdapter::flush() { m_logger->flush(); } if (m_sink) { + m_sink->complete_pending(); m_sink->flush(); } } diff --git a/bench/adapters/SpdlogAdapter.hpp b/bench/adapters/SpdlogAdapter.hpp index b34bf68..d6f75ec 100644 --- a/bench/adapters/SpdlogAdapter.hpp +++ b/bench/adapters/SpdlogAdapter.hpp @@ -24,11 +24,14 @@ class SpdlogAdapter : public ILoggerAdapter { void flush() override; + void set_recorder_handle(std::shared_ptr recorder) override; + private: class MeasuringSink; std::shared_ptr m_logger; std::shared_ptr m_sink; + std::shared_ptr m_recorder_handle; bool m_async = false; }; diff --git a/bench/logit_bench.cpp b/bench/logit_bench.cpp index 98b2974..d35ccc3 100644 --- a/bench/logit_bench.cpp +++ b/bench/logit_bench.cpp @@ -9,6 +9,7 @@ #include #include #include +#include #include #include #include @@ -48,6 +49,50 @@ std::size_t get_env_size_t(const char* name, std::size_t def) { return def; } +struct BenchFilter { + std::optional library; + std::optional async; + std::optional sink; + std::optional producers; + std::optional bytes; + + bool matches(const std::string& lib, + bool async_mode, + SinkKind sink_kind, + std::size_t producer_count, + std::size_t msg_bytes) const { + if (library && *library != lib) return false; + if (async && *async != async_mode) return false; + if (sink && *sink != sink_kind) return false; + if (producers && *producers != producer_count) return false; + if (bytes && *bytes != msg_bytes) return false; + return true; + } +}; + +BenchFilter load_filter() { + BenchFilter filter; + if (const char* v = std::getenv("LOGIT_BENCH_FILTER_LIB")) { + filter.library = std::string(v); + } + if (const char* v = std::getenv("LOGIT_BENCH_FILTER_ASYNC")) { + filter.async = std::string(v) == "1"; + } + if (const char* v = std::getenv("LOGIT_BENCH_FILTER_SINK")) { + std::string s(v); + std::transform(s.begin(), s.end(), s.begin(), [](unsigned char c){ return static_cast(std::tolower(c)); }); + if (s == "null") filter.sink = SinkKind::Null; + if (s == "file") filter.sink = SinkKind::File; + } + if (const char* v = std::getenv("LOGIT_BENCH_FILTER_PRODUCERS")) { + filter.producers = get_env_size_t("LOGIT_BENCH_FILTER_PRODUCERS", 0); + } + if (const char* v = std::getenv("LOGIT_BENCH_FILTER_BYTES")) { + filter.bytes = get_env_size_t("LOGIT_BENCH_FILTER_BYTES", 0); + } + return filter; +} + std::uint64_t steady_now_ns() { const auto now_tp = std::chrono::steady_clock::now().time_since_epoch(); return std::chrono::duration_cast(now_tp).count(); @@ -119,6 +164,7 @@ std::chrono::nanoseconds run_workload( // Barrier to start together. std::mutex start_mx; std::condition_variable start_cv; + std::condition_variable ready_cv; bool start_flag = false; std::size_t ready = 0; @@ -132,7 +178,7 @@ std::chrono::nanoseconds run_workload( { std::unique_lock lk(start_mx); ++ready; - if (ready == scenario.producers) start_cv.notify_one(); + if (ready == scenario.producers) ready_cv.notify_one(); start_cv.wait(lk, [&]{ return start_flag; }); } for (std::size_t n = 0; n < per_thread[i]; ++n) { @@ -150,7 +196,7 @@ std::chrono::nanoseconds run_workload( std::chrono::steady_clock::time_point t0; { std::unique_lock lk(start_mx); - start_cv.wait(lk, [&]{ return ready == scenario.producers; }); + ready_cv.wait(lk, [&]{ return ready == scenario.producers; }); if (measure_duration) t0 = std::chrono::steady_clock::now(); start_flag = true; start_cv.notify_all(); @@ -176,10 +222,12 @@ ScenarioResult execute_scenario( const Scenario& scenario, std::size_t warmup_messages) { - LatencyRecorder recorder(scenario.total_messages); + auto recorder = std::make_shared(scenario.total_messages); + + adapter.set_recorder_handle(recorder); // Adapter should keep a pointer/ref to recorder and call complete(token) from its sink. - adapter.prepare(scenario, recorder); + adapter.prepare(scenario, *recorder); // Warm-up (no recording, no duration). { @@ -192,7 +240,7 @@ ScenarioResult execute_scenario( << " total=" << warmup_messages; log_info(oss.str()); } - run_workload(adapter, recorder, scenario, warmup_messages, false, false); + run_workload(adapter, *recorder, scenario, warmup_messages, false, false); { std::ostringstream oss; oss << "Warm-up completed lib=" << adapter.library_name() @@ -214,7 +262,7 @@ ScenarioResult execute_scenario( << " total=" << scenario.total_messages; log_info(oss.str()); } - const auto dur = run_workload(adapter, recorder, scenario, scenario.total_messages, true, true); + const auto dur = run_workload(adapter, *recorder, scenario, scenario.total_messages, true, true); { std::ostringstream oss; oss << "Measure completed lib=" << adapter.library_name() @@ -225,13 +273,22 @@ ScenarioResult execute_scenario( log_info(oss.str()); } - const auto sum = recorder.finalize(); + // Ensure async pipelines (e.g., spdlog thread pool) are fully drained before + // destroying the recorder referenced by sinks. + adapter.flush(); + + recorder->wait_for_all(); + + const auto sum = recorder->finalize(); double thr = 0.0; if (dur.count() > 0) { const double sec = static_cast(dur.count()) / 1'000'000'000.0; thr = static_cast(scenario.total_messages) / sec; } + + adapter.set_recorder_handle(nullptr); + return ScenarioResult{sum, thr, dur}; } @@ -313,6 +370,8 @@ int main() { const std::size_t warmup_messages = get_env_size_t("LOGIT_BENCH_WARMUP", 4096); const std::size_t timeout_seconds = get_env_size_t("LOGIT_BENCH_TIMEOUT_SEC", 1200); + const BenchFilter filter = load_filter(); + LOGIT_SET_MAX_QUEUE(total_messages); if (timeout_seconds > 0) { @@ -338,6 +397,9 @@ int main() { for (auto sink : sinks) { for (std::size_t producers : producer_counts) { for (std::size_t msg_bytes : message_sizes) { + if (!filter.matches(adapter->library_name(), async_mode, sink, producers, msg_bytes)) { + continue; + } Scenario scenario; scenario.async = async_mode; scenario.sink = sink; diff --git a/docs/TaskExecutor.md b/docs/TaskExecutor.md index f93392f..154c3bb 100644 --- a/docs/TaskExecutor.md +++ b/docs/TaskExecutor.md @@ -191,7 +191,13 @@ const auto lost = LOGIT_GET_DROPPED_TASKS(); `set_queue_policy()`. * The hot-resize barrier uses `m_resizing` and `m_resize_cv` so producers never touch a ring buffer that is being rebuilt. This eliminates the data races that - TSAN previously reported on `try_pop()` vs. buffer assignment. + TSAN previously reported on `try_pop()` vs. buffer assignment. The barrier + only drops once the worker thread fully stops and the queue drains; if a sink + blocks the worker or `QueuePolicy::Block` keeps `m_active_tasks` above the + limit for more than one second, `set_max_queue_size()` abandons the hot + resize, clears `m_resizing`, and leaves the existing ring untouched so + producers cannot wait indefinitely. Non-MPSC builds perform the resize as an + atomic update of `m_max_queue_size`, so they are not subject to this stall. * Non-MPSC builds rely solely on mutexes and had no known data races. * The Emscripten path is single-threaded and should not be used concurrently. diff --git a/include/logit_cpp/logit/detail/MpscRingAny.hpp b/include/logit_cpp/logit/detail/MpscRingAny.hpp index a4c206a..1a747d7 100644 --- a/include/logit_cpp/logit/detail/MpscRingAny.hpp +++ b/include/logit_cpp/logit/detail/MpscRingAny.hpp @@ -112,32 +112,35 @@ namespace logit { namespace detail { /// \brief Try to dequeue value into out. Non-blocking. /// \return true on success; false if queue is empty. bool try_pop(T& out) noexcept { - std::size_t pos = m_dequeue_pos.load(std::memory_order_relaxed); - Cell& c = m_cells[pos % m_cap]; - std::size_t seq = c.m_seq.load(std::memory_order_acquire); - - // When ready, seq == pos + 1 - std::intptr_t diff = - static_cast(seq) - static_cast(pos + 1); - - if (diff == 0) { - if (!m_dequeue_pos.compare_exchange_strong( - pos, pos + 1, - std::memory_order_relaxed, - std::memory_order_relaxed)) { - return false; // Single consumer: should be rare. + for (;;) { + std::size_t pos = m_dequeue_pos.load(std::memory_order_relaxed); + Cell& c = m_cells[pos % m_cap]; + std::size_t seq = c.m_seq.load(std::memory_order_acquire); + + // When ready, seq == pos + 1 + std::intptr_t diff = + static_cast(seq) - static_cast(pos + 1); + + if (diff == 0) { + if (!m_dequeue_pos.compare_exchange_weak( + pos, pos + 1, + std::memory_order_relaxed, + std::memory_order_relaxed)) { + // Spurious failure: retry until we own the slot. + continue; + } + + T* p = reinterpret_cast(&c.m_storage); + out = std::move(*p); + p->~T(); + + // Mark cell free for next cycle. + c.m_seq.store(pos + m_cap, std::memory_order_release); + return true; } - - T* p = reinterpret_cast(&c.m_storage); - out = std::move(*p); - p->~T(); - - // Mark cell free for next cycle. - c.m_seq.store(pos + m_cap, std::memory_order_release); - return true; + + return false; // Empty or not yet published. } - - return false; // Empty or not yet published. } /// \brief Lightweight emptiness check for current consumer position. diff --git a/include/logit_cpp/logit/detail/TaskExecutor.hpp b/include/logit_cpp/logit/detail/TaskExecutor.hpp index 3af0e6e..da813c0 100644 --- a/include/logit_cpp/logit/detail/TaskExecutor.hpp +++ b/include/logit_cpp/logit/detail/TaskExecutor.hpp @@ -12,12 +12,12 @@ #include #include #include -#else - #include - #include - #include - #include - #include + #else + #include + #include + #include + #include + #include #endif // Enable lock-free MPSC ring integration (non-Emscripten) by defining: @@ -323,8 +323,18 @@ namespace logit { namespace detail { // Tell producers to pause before any wait()/stop conditions run. m_resizing.store(true, std::memory_order_release); - // Drain the queue completely. - wait(); + // Drain the queue completely, but do not wait forever if the worker + // is stalled (e.g., blocked sink or backpressure keeping + // m_active_tasks > 0). If we fail to drain before the deadline, + // abort the resize and re-open the barrier so producers can + // continue. + const auto deadline = std::chrono::steady_clock::now() + + std::chrono::seconds(1); + if (!wait_until_idle_(deadline)) { + m_resizing.store(false, std::memory_order_release); + m_resize_cv.notify_all(); + return; + } // Stop the worker so it cannot touch m_mpsc_queue during the resize. std::unique_lock lk(m_queue_mutex); @@ -474,6 +484,15 @@ namespace logit { namespace detail { bool queue_empty_() const noexcept { return m_mpsc_queue.empty(); } + + bool wait_until_idle_(std::chrono::steady_clock::time_point deadline) { + std::unique_lock lock(m_queue_mutex); + return m_queue_condition.wait_until(lock, deadline, [this]() { + return ((queue_empty_() && + m_active_tasks.load(std::memory_order_relaxed) == 0) || + m_stop_flag.load(std::memory_order_acquire)); + }); + } #endif TaskExecutor()