diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index b009f52..7592ab4 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -2,7 +2,7 @@ name: CI on: push: - branches: [ main ] + branches: [ main, stable ] pull_request: branches: [ main ] @@ -25,6 +25,20 @@ jobs: run: cmake --install build --prefix install - name: Test run: ctest --test-dir build --output-on-failure + - name: Configure benchmarks + if: ${{ github.event_name == 'pull_request' || (github.event_name == 'push' && github.ref == 'refs/heads/stable') }} + run: cmake -S . -B build-bench -DLOGIT_BENCH_ENABLE=ON -DLOGIT_BENCH_WITH_SPDLOG=ON -DCMAKE_CXX_STANDARD=${{ matrix.std }} -DLOGIT_WITH_SYSLOG=ON -DLOGIT_WITH_WIN_EVENT_LOG=OFF + - name: Build benchmarks + if: ${{ github.event_name == 'pull_request' || (github.event_name == 'push' && github.ref == 'refs/heads/stable') }} + run: cmake --build build-bench --target logit_bench + - name: Run latency benchmarks + if: ${{ github.event_name == 'pull_request' || (github.event_name == 'push' && github.ref == 'refs/heads/stable') }} + timeout-minutes: 20 + env: + LOGIT_BENCH_TIMEOUT_SEC: 900 + LOGIT_BENCH_TOTAL: 20000 + LOGIT_BENCH_WARMUP: 2000 + run: ./build-bench/logit_bench - 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 diff --git a/CMakeLists.txt b/CMakeLists.txt index 8ba99e4..6975864 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -3,6 +3,8 @@ project(log-it-cpp VERSION 1.0.0 LANGUAGES CXX) option(LOGIT_CPP_BUILD_TESTS "Build log-it-cpp tests" ${PROJECT_IS_TOP_LEVEL}) option(LOGIT_CPP_BUILD_EXAMPLES "Build log-it-cpp examples" OFF) +option(LOGIT_BENCH_ENABLE "Build log-it-cpp benchmarks" OFF) +option(LOGIT_BENCH_WITH_SPDLOG "Enable spdlog comparison benchmarks" OFF) option(LOGIT_WITH_GZIP "Enable gzip via zlib" OFF) option(LOGIT_WITH_ZSTD "Enable zstd" OFF) option(LOGIT_WITH_FMT "Enable fmt support" OFF) @@ -140,6 +142,10 @@ if(LOGIT_CPP_BUILD_EXAMPLES) add_subdirectory(examples) endif() +if(LOGIT_BENCH_ENABLE) + add_subdirectory(bench) +endif() + include(CMakePackageConfigHelpers) install(DIRECTORY include/ DESTINATION include) diff --git a/README.md b/README.md index d6b054b..6c48578 100644 --- a/README.md +++ b/README.md @@ -736,6 +736,20 @@ When building with Emscripten the library runs without threads. Console logging works as usual while file-based loggers are replaced by stubs that warn when used. +## Benchmarks + +Latency and throughput benchmarks live under `bench/`. Enable them during configuration and optionally pull in the spdlog +adapters: + +```bash +cmake -S . -B build -DLOGIT_BENCH_ENABLE=ON -DLOGIT_BENCH_WITH_SPDLOG=ON +cmake --build build --target logit_bench +``` + +Run `./build/bench/logit_bench` to record the full matrix (sync/async × null/file × producer counts × message sizes). Results +are appended to `bench/results/latency.csv` with one row per library/combination. Override the workload via `LOGIT_BENCH_TOTAL` +and `LOGIT_BENCH_WARMUP` environment variables if you need a lighter run. + --- ## Documentation diff --git a/bench/CMakeLists.txt b/bench/CMakeLists.txt new file mode 100644 index 0000000..ebffc7c --- /dev/null +++ b/bench/CMakeLists.txt @@ -0,0 +1,39 @@ +set(LOGIT_BENCH_SOURCES + logit_bench.cpp + adapters/LogItAdapter.cpp +) + +if(LOGIT_BENCH_WITH_SPDLOG) + list(APPEND LOGIT_BENCH_SOURCES adapters/SpdlogAdapter.cpp) +endif() + +add_executable(logit_bench ${LOGIT_BENCH_SOURCES}) + +target_include_directories(logit_bench PRIVATE ${CMAKE_CURRENT_SOURCE_DIR}) + +target_compile_features(logit_bench PRIVATE cxx_std_17) + +set_target_properties(logit_bench PROPERTIES + RUNTIME_OUTPUT_DIRECTORY ${CMAKE_BINARY_DIR} +) + +foreach(config IN ITEMS DEBUG RELEASE RELWITHDEBINFO MINSIZEREL) + set_target_properties(logit_bench PROPERTIES + RUNTIME_OUTPUT_DIRECTORY_${config} ${CMAKE_BINARY_DIR} + ) +endforeach() + +target_link_libraries(logit_bench PRIVATE log-it-cpp::log-it-cpp) + +if(LOGIT_BENCH_WITH_SPDLOG) + target_compile_definitions(logit_bench PRIVATE LOGIT_BENCH_HAVE_SPDLOG=1) + if(NOT TARGET spdlog::spdlog) + include(FetchContent) + FetchContent_Declare(spdlog + GIT_REPOSITORY https://github.com/gabime/spdlog.git + GIT_TAG v1.12.0 + ) + FetchContent_MakeAvailable(spdlog) + endif() + target_link_libraries(logit_bench PRIVATE spdlog::spdlog) +endif() diff --git a/bench/LatencyRecorder.hpp b/bench/LatencyRecorder.hpp new file mode 100644 index 0000000..f25a84e --- /dev/null +++ b/bench/LatencyRecorder.hpp @@ -0,0 +1,107 @@ +#pragma once + +#include +#include +#include +#include +#include +#include +#include +#include + +namespace logit_bench { + +/** + * Lock-free recorder for latency samples: + * - begin(record=true) returns a Token with an assigned slot and t0_ns (steady_clock). + * - complete(token) stores (t1-t0) in that slot. + * - finalize() returns p50/p99/p99.9 using nearest-rank (ceil) on a sorted copy. + * + * Thread-safety: concurrent writers store into distinct preallocated slots. + */ +class LatencyRecorder { +public: + struct Token { + std::uint64_t slot = invalid_slot(); + std::uint64_t t0_ns = 0; + bool active = false; + }; + + struct Summary { + std::uint64_t p50_ns = 0; + std::uint64_t p99_ns = 0; + std::uint64_t p999_ns = 0; + }; + + explicit LatencyRecorder(std::size_t total) + : m_values(total), + m_expected(total), + m_next_slot(0) {} + + /** + * Reserve a slot (if record==true) and capture t0 using steady_clock. + * We take t0 **after** the slot reservation to minimize skew before log(). + */ + Token begin(bool record) { + Token token; + token.active = record; + if (record) { + const auto slot = m_next_slot.fetch_add(1, std::memory_order_relaxed); + if (slot >= m_expected) { + throw std::out_of_range("LatencyRecorder capacity exceeded"); + } + token.slot = static_cast(slot); + token.t0_ns = now(); + } + return token; + } + + /// Capture t1 and store (t1 - t0) into the reserved slot. + void complete(const Token& token) { + if (!token.active) return; + const auto t1_ns = now(); + m_values[token.slot] = t1_ns - token.t0_ns; // distinct slots -> no data race + } + + std::size_t recorded() const { + return m_next_slot.load(std::memory_order_relaxed); + } + + Summary finalize() const { + if (recorded() != m_expected) { + throw std::runtime_error("Incomplete latency capture"); + } + std::vector sorted = m_values; + std::sort(sorted.begin(), sorted.end()); + Summary summary; + summary.p50_ns = pick(sorted, 0.50); + summary.p99_ns = pick(sorted, 0.99); + summary.p999_ns = pick(sorted, 0.999); + return summary; + } + + static std::uint64_t invalid_slot() { + return std::numeric_limits::max(); + } + + static std::uint64_t now() { + const auto now_tp = std::chrono::steady_clock::now().time_since_epoch(); + return std::chrono::duration_cast(now_tp).count(); + } + +private: + // Nearest-rank percentile with ceil(p * N), clamped to [0..N-1]. + static std::uint64_t pick(const std::vector& data, double p) { + if (data.empty()) return 0; + const double r = std::ceil(p * static_cast(data.size())); + std::size_t idx = (r <= 1.0) ? 0 : static_cast(r) - 1; + if (idx >= data.size()) idx = data.size() - 1; + return data[idx]; + } + + std::vector m_values; // preallocated; no reallocation + const std::size_t m_expected; // total messages to record + std::atomic m_next_slot; +}; + +} // namespace logit_bench diff --git a/bench/Scenario.hpp b/bench/Scenario.hpp new file mode 100644 index 0000000..3fe052a --- /dev/null +++ b/bench/Scenario.hpp @@ -0,0 +1,29 @@ +#pragma once + +#include +#include + +namespace logit_bench { + +enum class SinkKind { + Null, + File, +}; + +inline std::string sink_name(SinkKind sink) { + switch (sink) { + case SinkKind::Null: return "null"; + case SinkKind::File: return "file"; + } + return "unknown"; +} + +struct Scenario { + bool async = false; + SinkKind sink = SinkKind::Null; + std::size_t producers = 1; + std::size_t message_bytes = 0; + std::size_t total_messages = 0; +}; + +} // namespace logit_bench diff --git a/bench/adapters/ILoggerAdapter.hpp b/bench/adapters/ILoggerAdapter.hpp new file mode 100644 index 0000000..508a4ad --- /dev/null +++ b/bench/adapters/ILoggerAdapter.hpp @@ -0,0 +1,23 @@ +#pragma once + +#include + +#include "../LatencyRecorder.hpp" +#include "../Scenario.hpp" + +namespace logit_bench { + +class ILoggerAdapter { +public: + virtual ~ILoggerAdapter() = default; + + virtual const char* library_name() const = 0; + + virtual void prepare(const Scenario& scenario, LatencyRecorder& recorder) = 0; + + virtual void log(const LatencyRecorder::Token& token, std::string_view message) = 0; + + virtual void flush() = 0; +}; + +} // namespace logit_bench diff --git a/bench/adapters/LogItAdapter.cpp b/bench/adapters/LogItAdapter.cpp new file mode 100644 index 0000000..7ddd835 --- /dev/null +++ b/bench/adapters/LogItAdapter.cpp @@ -0,0 +1,210 @@ +#include "LogItAdapter.hpp" + +#include +#include +#include +#include +#include +#include + +#include + +namespace logit_bench { +namespace { +constexpr const char* kFilePath = "bench/results/logit_sink.log"; +constexpr std::size_t kSlotIndex = 0; +constexpr std::size_t kT0Index = 1; +constexpr std::size_t kActiveIndex = 2; +} // namespace + +class PassthroughFormatter : public logit::ILogFormatter { +public: + void set_timestamp_offset(int64_t) override {} + + std::string format(const logit::LogRecord& record) const override { + return record.format; + } +}; + +class MeasuringSink : public logit::ILogger { +public: + MeasuringSink() = default; + + void configure(const Scenario& scenario, LatencyRecorder& recorder) { + m_async = scenario.async; + m_sink = scenario.sink; + m_recorder = &recorder; + if (m_sink == SinkKind::File) { + std::filesystem::create_directories("bench/results"); + std::lock_guard lock(m_file_mutex); + m_file.close(); + m_file.open(kFilePath, std::ios::out | std::ios::trunc); + } else { + std::lock_guard lock(m_file_mutex); + m_file.close(); + } + } + + void log(const logit::LogRecord& record, const std::string& message) override { + LatencyRecorder::Token token = extract(record); + if (!m_async) { + consume(token, message); + return; + } + AsyncPayload payload; + payload.token = token; + payload.text = message; + logit::detail::TaskExecutor::get_instance().add_task([this, payload = std::move(payload)]() mutable { + consume(payload.token, payload.text); + }); + } + + std::string get_string_param(const logit::LoggerParam&) const override { return std::string(); } + int64_t get_int_param(const logit::LoggerParam&) const override { return 0; } + double get_float_param(const logit::LoggerParam&) const override { return 0.0; } + + void set_log_level(logit::LogLevel level) override { + m_level.store(static_cast(level), std::memory_order_relaxed); + } + + logit::LogLevel get_log_level() const override { + return static_cast(m_level.load(std::memory_order_relaxed)); + } + + void wait() override { + if (m_async) { + logit::detail::TaskExecutor::get_instance().wait(); + } + std::lock_guard lock(m_file_mutex); + if (m_file.is_open()) { + m_file.flush(); + } + } + +private: + struct AsyncPayload { + LatencyRecorder::Token token; + std::string text; + }; + + static LatencyRecorder::Token extract(const logit::LogRecord& record) { + LatencyRecorder::Token token; + if (record.args_array.size() <= kActiveIndex) { + return token; + } + const auto& slot = record.args_array[kSlotIndex]; + const auto& t0 = record.args_array[kT0Index]; + const auto& active = record.args_array[kActiveIndex]; + token.slot = read_u64(slot); + token.t0_ns = read_u64(t0); + token.active = read_u64(active) != 0; + return token; + } + + static std::uint64_t read_u64(const logit::VariableValue& value) { + using VT = logit::VariableValue::ValueType; + switch (value.type) { + case VT::UINT64_VAL: + return value.pod_value.uint64_value; + case VT::INT64_VAL: + return static_cast(value.pod_value.int64_value); + case VT::UINT32_VAL: + return value.pod_value.uint32_value; + case VT::INT32_VAL: + return static_cast(value.pod_value.int32_value); + default: + break; + } + return 0; + } + + void consume(const LatencyRecorder::Token& token, std::string_view text) { + if (token.active && m_recorder) { + m_recorder->complete(token); + } + if (m_sink == SinkKind::File) { + std::lock_guard lock(m_file_mutex); + if (m_file.is_open()) { + m_file << text << '\n'; + } + } + } + + bool m_async = false; + SinkKind m_sink = SinkKind::Null; + LatencyRecorder* m_recorder = nullptr; + std::ofstream m_file; + mutable std::mutex m_file_mutex; + std::atomic m_level{static_cast(logit::LogLevel::LOG_LVL_TRACE)}; +}; + +class LogItAdapter::Impl { +public: + Impl() + : logger(logit::Logger::get_instance()) { + auto sink_ptr = std::make_unique(); + sink = sink_ptr.get(); + auto formatter = std::unique_ptr(new PassthroughFormatter()); + logger.add_logger(std::move(sink_ptr), std::move(formatter)); + } + + void prepare(const Scenario& scenario, LatencyRecorder& recorder) { + if (sink) { + sink->configure(scenario, recorder); + } + } + + void log(const LatencyRecorder::Token& token, std::string_view message) { + std::string text(message); + logit::LogRecord record( + logit::LogLevel::LOG_LVL_INFO, + 0, + std::string(), + 0, + std::string(), + text, + std::string(), + -1, + false, + false); + record.args_array.reserve(3); + record.args_array.emplace_back("slot", static_cast(token.slot)); + record.args_array.emplace_back("t0", static_cast(token.t0_ns)); + record.args_array.emplace_back("active", static_cast(token.active ? 1 : 0)); + logger.log(record); + } + + void flush() { + if (sink) { + sink->wait(); + } + } + + logit::Logger& logger; + MeasuringSink* sink = nullptr; +}; + +LogItAdapter::LogItAdapter() + : m_impl(std::make_unique()) {} + +LogItAdapter::~LogItAdapter() = default; + +void LogItAdapter::prepare(const Scenario& scenario, LatencyRecorder& recorder) { + if (m_impl) { + m_impl->prepare(scenario, recorder); + } +} + +void LogItAdapter::log(const LatencyRecorder::Token& token, std::string_view message) { + if (m_impl) { + m_impl->log(token, message); + } +} + +void LogItAdapter::flush() { + if (m_impl) { + m_impl->flush(); + } +} + +} // namespace logit_bench diff --git a/bench/adapters/LogItAdapter.hpp b/bench/adapters/LogItAdapter.hpp new file mode 100644 index 0000000..bc5467b --- /dev/null +++ b/bench/adapters/LogItAdapter.hpp @@ -0,0 +1,29 @@ +#pragma once + +#include +#include +#include + +#include "ILoggerAdapter.hpp" + +namespace logit_bench { + +class LogItAdapter : public ILoggerAdapter { +public: + LogItAdapter(); + ~LogItAdapter() override; + + const char* library_name() const override { return "log-it-cpp"; } + + void prepare(const Scenario& scenario, LatencyRecorder& recorder) override; + + void log(const LatencyRecorder::Token& token, std::string_view message) override; + + void flush() override; + +private: + class Impl; + std::unique_ptr m_impl; +}; + +} // namespace logit_bench diff --git a/bench/adapters/SpdlogAdapter.cpp b/bench/adapters/SpdlogAdapter.cpp new file mode 100644 index 0000000..5cb22c1 --- /dev/null +++ b/bench/adapters/SpdlogAdapter.cpp @@ -0,0 +1,144 @@ +#include "SpdlogAdapter.hpp" + +#ifdef LOGIT_BENCH_HAVE_SPDLOG + +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include + +namespace logit_bench { +namespace { +constexpr const char* kFilePath = "bench/results/spdlog_sink.log"; +constexpr std::size_t kDefaultQueue = 8192; + +struct MessagePayload { + LatencyRecorder::Token token; + std::string text; +}; +} // namespace + +class SpdlogAdapter::MeasuringSink : public spdlog::sinks::sink { +public: + MeasuringSink() = default; + + void configure(const Scenario& scenario, LatencyRecorder& recorder) { + m_sink = scenario.sink; + m_recorder = &recorder; + if (m_sink == SinkKind::File) { + std::filesystem::create_directories("bench/results"); + std::lock_guard lock(m_mutex); + m_file.close(); + m_file.open(kFilePath, std::ios::out | std::ios::trunc); + } else { + std::lock_guard lock(m_mutex); + m_file.close(); + } + } + + void log(const spdlog::details::log_msg& msg) override { + const auto* payload_ptr = reinterpret_cast(msg.source.funcname); + if (!payload_ptr) { + return; + } + auto* payload = const_cast(payload_ptr); + consume(*payload); + delete payload; + } + + void set_pattern(const std::string&) override {} + + void set_formatter(std::unique_ptr) override {} + + void flush() override { + std::lock_guard lock(m_mutex); + if (m_file.is_open()) { + m_file.flush(); + } + } + +private: + void consume(const MessagePayload& payload) { + if (payload.token.active && m_recorder) { + m_recorder->complete(payload.token); + } + if (m_sink == SinkKind::File) { + std::lock_guard lock(m_mutex); + if (m_file.is_open()) { + m_file << payload.text << '\n'; + } + } + } + + SinkKind m_sink = SinkKind::Null; + LatencyRecorder* m_recorder = nullptr; + std::ofstream m_file; + std::mutex m_mutex; +}; + +SpdlogAdapter::SpdlogAdapter() = default; + +SpdlogAdapter::~SpdlogAdapter() { + flush(); + spdlog::shutdown(); +} + +void SpdlogAdapter::prepare(const Scenario& scenario, LatencyRecorder& recorder) { + m_logger.reset(); + m_sink.reset(); + spdlog::shutdown(); + + m_sink = std::make_shared(); + m_sink->configure(scenario, recorder); + m_async = scenario.async; + + std::string logger_name = m_async ? "logit_bench_async" : "logit_bench_sync"; + if (m_async) { + const std::size_t queue_size = std::max(kDefaultQueue, scenario.total_messages * 2); + spdlog::init_thread_pool(queue_size, 1); + auto async_logger = std::make_shared( + logger_name, + m_sink, + spdlog::thread_pool(), + spdlog::async_overflow_policy::block); + async_logger->set_level(spdlog::level::trace); + async_logger->set_pattern("%v"); + m_logger = std::move(async_logger); + } else { + auto logger = std::make_shared(logger_name, m_sink); + logger->set_level(spdlog::level::trace); + logger->set_pattern("%v"); + m_logger = std::move(logger); + } +} + +void SpdlogAdapter::log(const LatencyRecorder::Token& token, std::string_view message) { + if (!m_logger) { + return; + } + auto* payload = new MessagePayload(); + 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)); +} + +void SpdlogAdapter::flush() { + if (m_logger) { + m_logger->flush(); + } + if (m_sink) { + m_sink->flush(); + } +} + +} // namespace logit_bench + +#endif // LOGIT_BENCH_HAVE_SPDLOG diff --git a/bench/adapters/SpdlogAdapter.hpp b/bench/adapters/SpdlogAdapter.hpp new file mode 100644 index 0000000..b34bf68 --- /dev/null +++ b/bench/adapters/SpdlogAdapter.hpp @@ -0,0 +1,37 @@ +#pragma once + +#ifdef LOGIT_BENCH_HAVE_SPDLOG + +#include +#include + +#include + +#include "ILoggerAdapter.hpp" + +namespace logit_bench { + +class SpdlogAdapter : public ILoggerAdapter { +public: + SpdlogAdapter(); + ~SpdlogAdapter() override; + + const char* library_name() const override { return "spdlog"; } + + void prepare(const Scenario& scenario, LatencyRecorder& recorder) override; + + void log(const LatencyRecorder::Token& token, std::string_view message) override; + + void flush() override; + +private: + class MeasuringSink; + + std::shared_ptr m_logger; + std::shared_ptr m_sink; + bool m_async = false; +}; + +} // namespace logit_bench + +#endif // LOGIT_BENCH_HAVE_SPDLOG diff --git a/bench/logit_bench.cpp b/bench/logit_bench.cpp new file mode 100644 index 0000000..98b2974 --- /dev/null +++ b/bench/logit_bench.cpp @@ -0,0 +1,378 @@ +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "LatencyRecorder.hpp" +#include "Scenario.hpp" +#include "adapters/LogItAdapter.hpp" + +#ifdef LOGIT_BENCH_HAVE_SPDLOG +#include "adapters/SpdlogAdapter.hpp" +#endif + +namespace logit_bench { +namespace { + +std::atomic* g_watchdog_progress = nullptr; +constexpr std::size_t k_watchdog_stride = 256; + +std::string make_message(std::size_t bytes, std::size_t index) { + if (bytes == 0) return {}; + const char fill = static_cast('A' + static_cast(index % 26)); + return std::string(bytes, fill); +} + +std::size_t get_env_size_t(const char* name, std::size_t def) { + if (const char* v = std::getenv(name)) { + try { + return static_cast(std::stoull(v)); + } catch (...) { + // fallthrough + } + } + return def; +} + +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(); +} + +std::string format_timestamp() { + const auto now = std::chrono::system_clock::now(); + const auto time = std::chrono::system_clock::to_time_t(now); + std::tm tm{}; +#ifdef _WIN32 + localtime_s(&tm, &time); +#else + localtime_r(&time, &tm); +#endif + const auto ms = std::chrono::duration_cast( + now.time_since_epoch()) % 1000; + std::ostringstream oss; + oss << std::put_time(&tm, "%Y-%m-%d %H:%M:%S") + << '.' << std::setw(3) << std::setfill('0') << ms.count(); + return oss.str(); +} + +void touch_watchdog() { + if (g_watchdog_progress) { + g_watchdog_progress->store(steady_now_ns(), std::memory_order_relaxed); + } +} + +void log_info(const std::string& message) { + std::cout << "[logit_bench " << format_timestamp() << "] " << message << std::endl; + touch_watchdog(); +} + +void log_error(const std::string& message) { + std::cerr << "[logit_bench " << format_timestamp() << "] " << message << std::endl; + touch_watchdog(); +} + +/** + * Run a workload: + * - producers start together (barrier), + * - each producer logs its portion of total_messages, + * - LatencyRecorder::begin(record) captures t0 and slot, + * - adapter.log(token, message) must eventually call recorder.complete(token) from sink/consumer, + * - returns total wall duration (for throughput). + */ +std::chrono::nanoseconds run_workload( + ILoggerAdapter& adapter, + LatencyRecorder& recorder, + const Scenario& scenario, + std::size_t total_messages, + bool record_latency, + bool measure_duration) +{ + if (scenario.producers == 0) { + adapter.flush(); + return std::chrono::nanoseconds(0); + } + + // Distribute messages across producers. + std::vector per_thread(scenario.producers, 0); + const std::size_t base = total_messages / scenario.producers; + std::size_t rem = total_messages % scenario.producers; + for (std::size_t i = 0; i < scenario.producers; ++i) { + per_thread[i] = base + (rem ? 1 : 0); + if (rem) --rem; + } + + // Barrier to start together. + std::mutex start_mx; + std::condition_variable start_cv; + bool start_flag = false; + std::size_t ready = 0; + + std::vector threads; + threads.reserve(scenario.producers); + + for (std::size_t i = 0; i < scenario.producers; ++i) { + threads.emplace_back([&, i]() { + std::string message = make_message(scenario.message_bytes, i); + std::size_t watchdog_counter = 0; + { + std::unique_lock lk(start_mx); + ++ready; + if (ready == scenario.producers) start_cv.notify_one(); + start_cv.wait(lk, [&]{ return start_flag; }); + } + for (std::size_t n = 0; n < per_thread[i]; ++n) { + auto token = recorder.begin(record_latency); + adapter.log(token, message); + ++watchdog_counter; + if ((watchdog_counter & (k_watchdog_stride - 1)) == 0) { + touch_watchdog(); + } + } + touch_watchdog(); + }); + } + + std::chrono::steady_clock::time_point t0; + { + std::unique_lock lk(start_mx); + start_cv.wait(lk, [&]{ return ready == scenario.producers; }); + if (measure_duration) t0 = std::chrono::steady_clock::now(); + start_flag = true; + start_cv.notify_all(); + } + + for (auto& th : threads) th.join(); + adapter.flush(); + touch_watchdog(); + + if (!measure_duration) return std::chrono::nanoseconds(0); + auto t1 = std::chrono::steady_clock::now(); + return std::chrono::duration_cast(t1 - t0); +} + +struct ScenarioResult { + LatencyRecorder::Summary summary; + double throughput = 0.0; + std::chrono::nanoseconds duration{0}; +}; + +ScenarioResult execute_scenario( + ILoggerAdapter& adapter, + const Scenario& scenario, + std::size_t warmup_messages) +{ + LatencyRecorder recorder(scenario.total_messages); + + // Adapter should keep a pointer/ref to recorder and call complete(token) from its sink. + adapter.prepare(scenario, recorder); + + // Warm-up (no recording, no duration). + { + std::ostringstream oss; + oss << "Warm-up start lib=" << adapter.library_name() + << " async=" << (scenario.async ? '1' : '0') + << " sink=" << sink_name(scenario.sink) + << " producers=" << scenario.producers + << " bytes=" << scenario.message_bytes + << " total=" << warmup_messages; + log_info(oss.str()); + } + run_workload(adapter, recorder, scenario, warmup_messages, false, false); + { + std::ostringstream oss; + oss << "Warm-up completed lib=" << adapter.library_name() + << " async=" << (scenario.async ? '1' : '0') + << " sink=" << sink_name(scenario.sink) + << " producers=" << scenario.producers + << " bytes=" << scenario.message_bytes; + log_info(oss.str()); + } + + // Measured run. + { + std::ostringstream oss; + oss << "Measure start lib=" << adapter.library_name() + << " async=" << (scenario.async ? '1' : '0') + << " sink=" << sink_name(scenario.sink) + << " producers=" << scenario.producers + << " bytes=" << scenario.message_bytes + << " total=" << scenario.total_messages; + log_info(oss.str()); + } + const auto dur = run_workload(adapter, recorder, scenario, scenario.total_messages, true, true); + { + std::ostringstream oss; + oss << "Measure completed lib=" << adapter.library_name() + << " async=" << (scenario.async ? '1' : '0') + << " sink=" << sink_name(scenario.sink) + << " producers=" << scenario.producers + << " bytes=" << scenario.message_bytes; + log_info(oss.str()); + } + + 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; + } + return ScenarioResult{sum, thr, dur}; +} + +void append_csv( + const std::string& library, + const Scenario& scenario, + const LatencyRecorder::Summary& summary, + double throughput) +{ + namespace fs = std::filesystem; + const fs::path csv_path{"bench/results/latency.csv"}; + fs::create_directories(csv_path.parent_path()); + + const bool write_header = !fs::exists(csv_path) || fs::file_size(csv_path) == 0; + + std::ofstream out(csv_path, std::ios::app); + if (!out) throw std::runtime_error("Failed to open latency.csv for writing"); + + if (write_header) { + out << "lib,async,sink,producers,msg_bytes,total,p50_ns,p99_ns,p999_ns,throughput\n"; + } + out << library << ',' + << (scenario.async ? 1 : 0) << ',' + << sink_name(scenario.sink) << ',' + << scenario.producers << ',' + << scenario.message_bytes << ',' + << scenario.total_messages << ',' + << summary.p50_ns << ',' + << summary.p99_ns << ',' + << summary.p999_ns << ',' + << std::fixed << std::setprecision(2) << throughput << '\n'; +} + +void print_summary( + const std::string& library, + const Scenario& scenario, + const ScenarioResult& result) +{ + std::ostringstream oss; + oss << library + << " async=" << (scenario.async ? '1' : '0') + << " sink=" << sink_name(scenario.sink) + << " producers=" << scenario.producers + << " bytes=" << scenario.message_bytes + << " total=" << scenario.total_messages + << " p50=" << result.summary.p50_ns + << "ns p99=" << result.summary.p99_ns + << "ns p999=" << result.summary.p999_ns + << "ns throughput=" << std::fixed << std::setprecision(2) + << result.throughput << " msg/s"; + log_info(oss.str()); +} + +} // namespace +} // namespace logit_bench + +int main() { + using namespace logit_bench; + std::atomic watchdog_done{false}; + std::thread watchdog; + std::atomic watchdog_progress{steady_now_ns()}; + g_watchdog_progress = &watchdog_progress; + + try { + std::vector> adapters; + adapters.emplace_back(std::make_unique()); +#ifdef LOGIT_BENCH_HAVE_SPDLOG + adapters.emplace_back(std::make_unique()); +#endif + + // Matrix + const std::array async_modes{false, true}; + const std::array sinks{SinkKind::Null, SinkKind::File}; + const std::array producer_counts{1, 4, 16}; + const std::array message_sizes{40, 200, 1024}; + + // Totals (can be overridden by env): + const std::size_t total_messages = get_env_size_t("LOGIT_BENCH_TOTAL", 200000); + 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); + + LOGIT_SET_MAX_QUEUE(total_messages); + + if (timeout_seconds > 0) { + watchdog = std::thread([timeout_seconds, &watchdog_done, &watchdog_progress]() { + const auto timeout = std::chrono::seconds(timeout_seconds); + while (!watchdog_done.load(std::memory_order_relaxed)) { + const auto last_ns = watchdog_progress.load(std::memory_order_relaxed); + const auto last_tp = std::chrono::steady_clock::time_point(std::chrono::nanoseconds(last_ns)); + if (std::chrono::steady_clock::now() - last_tp >= timeout) { + log_error(std::string("Timeout reached after ") + std::to_string(timeout_seconds) + + " seconds without progress. Terminating benchmark."); + + std::cerr.flush(); + std::_Exit(124); + } + std::this_thread::sleep_for(std::chrono::milliseconds(200)); + } + }); + } + + for (auto& adapter : adapters) { + for (bool async_mode : async_modes) { + for (auto sink : sinks) { + for (std::size_t producers : producer_counts) { + for (std::size_t msg_bytes : message_sizes) { + Scenario scenario; + scenario.async = async_mode; + scenario.sink = sink; + scenario.producers = producers; + scenario.message_bytes = msg_bytes; + scenario.total_messages = total_messages; + + { + std::ostringstream oss; + oss << "Scenario start lib=" << adapter->library_name() + << " async=" << (scenario.async ? '1' : '0') + << " sink=" << sink_name(scenario.sink) + << " producers=" << scenario.producers + << " bytes=" << scenario.message_bytes + << " total=" << scenario.total_messages; + log_info(oss.str()); + } + + auto result = execute_scenario(*adapter, scenario, warmup_messages); + append_csv(adapter->library_name(), scenario, result.summary, result.throughput); + print_summary(adapter->library_name(), scenario, result); + } + } + } + } + } + watchdog_done.store(true, std::memory_order_relaxed); + if (watchdog.joinable()) watchdog.join(); + } catch (const std::exception& ex) { + watchdog_done.store(true, std::memory_order_relaxed); + if (watchdog.joinable()) watchdog.join(); + log_error(std::string("Benchmark failed: ") + ex.what()); + g_watchdog_progress = nullptr; + return 1; + } + g_watchdog_progress = nullptr; + return 0; +} diff --git a/bench/results/.gitignore b/bench/results/.gitignore new file mode 100644 index 0000000..9bae93b --- /dev/null +++ b/bench/results/.gitignore @@ -0,0 +1 @@ +latency.csv diff --git a/tests/backpressure_ordering_test.cpp b/tests/backpressure_ordering_test.cpp index 824b45a..3b3a6b4 100644 --- a/tests/backpressure_ordering_test.cpp +++ b/tests/backpressure_ordering_test.cpp @@ -23,6 +23,8 @@ int main() { LOGIT_SET_MAX_QUEUE(kQueueCapacity); LOGIT_RESET_DROPPED_TASKS(); + // Векторы лежат в стеке main-потока, но их изменяют воркеры. + // Пишем/читаем их ТОЛЬКО под одним и тем же per-producer мьютексом. std::array, kProducers> sequences; for (auto &sequence : sequences) { sequence.clear(); @@ -37,6 +39,7 @@ int main() { 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(); } @@ -56,13 +59,16 @@ int main() { producer.join(); } - executor.wait(); + executor.wait(); // гарантируем завершение всех задач if (LOGIT_GET_DROPPED_TASKS() != 0) { return 1; } + // Читаем под тем же мьютексом — это устраняет data race в TSAN for (std::size_t producer_id = 0; producer_id < kProducers; ++producer_id) { + std::lock_guard lock(sequence_guards[producer_id]); + const auto &sequence = sequences[producer_id]; if (sequence.size() != kMessagesPerProducer) { return 2;