diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index fb9c998..b682160 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -30,7 +30,7 @@ jobs: run: cmake -S . -B build-bench -DLOGIT_BENCH_ENABLE=ON -DLOGIT_BENCH_WITH_SPDLOG=ON -DCMAKE_BUILD_TYPE=Release -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 logit_bench_async_contract logit_bench_async_payload_contract_test logit_bench_flush_test logit_public_macro_bench logit_public_macro_formatted_bench logit_hotpath_bench logit_hotpath_bench_legacy logit_exec_mx_bench logit_exec_mx_bench_concurrent benchmark_validation_test + run: cmake --build build-bench --target logit_bench logit_bench_async_contract logit_bench_async_payload_contract_test logit_bench_flush_test logit_bench_pipeline_research logit_public_macro_bench logit_public_macro_formatted_bench logit_hotpath_bench logit_hotpath_bench_legacy logit_exec_mx_bench logit_exec_mx_bench_concurrent benchmark_validation_test - name: Run spdlog async flush regression run: ./build-bench/logit_bench_flush_test - name: Run public macro benchmark smoke @@ -80,6 +80,51 @@ jobs: LOGIT_BENCH_TIMEOUT_SEC: 120 LOGIT_BENCH_OUTPUT: bench/results/latency-async-contract.csv run: ./build-bench/logit_bench_async_contract + - name: Run pipeline research matrix smoke + if: matrix.std == 17 + env: + LOGIT_BENCH_TOTAL: 1000 + LOGIT_BENCH_WARMUP: 100 + LOGIT_BENCH_REPEATS: 1 + LOGIT_BENCH_PRODUCERS: 1 + LOGIT_BENCH_QUEUE_CAPACITIES: 1024 + LOGIT_BENCH_RESEARCH_CSV: build-bench/pipeline-research-matrix.csv + LOGIT_BENCH_RESEARCH_JSONL: build-bench/pipeline-research-matrix.jsonl + LOGIT_BENCH_RESEARCH_AGGREGATE: build-bench/pipeline-research-matrix-aggregate.csv + run: | + ./build-bench/logit_bench_pipeline_research + test -s build-bench/pipeline-research-matrix.csv + test -s build-bench/pipeline-research-matrix.jsonl + python3 - <<'PY' + import csv + with open("build-bench/pipeline-research-matrix.csv", newline="") as stream: + rows = list(csv.DictReader(stream)) + assert rows, "matrix research receipt is empty" + assert all(int(row["issued"]) == int(row["sink_completed"]) == int(row["total"]) for row in rows) + PY + - name: Run pipeline research rate smoke + if: matrix.std == 17 + env: + LOGIT_BENCH_RESEARCH_MODE: rate + LOGIT_BENCH_TOTAL: 1000 + LOGIT_BENCH_WARMUP: 100 + LOGIT_BENCH_REPEATS: 1 + LOGIT_BENCH_PRODUCERS: 1 + LOGIT_BENCH_RATES: 100000 + LOGIT_BENCH_RESEARCH_CSV: build-bench/pipeline-research-rate.csv + LOGIT_BENCH_RESEARCH_JSONL: build-bench/pipeline-research-rate.jsonl + LOGIT_BENCH_RESEARCH_AGGREGATE: build-bench/pipeline-research-rate-aggregate.csv + run: | + ./build-bench/logit_bench_pipeline_research + test -s build-bench/pipeline-research-rate.csv + test -s build-bench/pipeline-research-rate.jsonl + python3 - <<'PY' + import csv + with open("build-bench/pipeline-research-rate.csv", newline="") as stream: + rows = list(csv.DictReader(stream)) + assert rows, "rate research receipt is empty" + assert all(int(row["issued"]) == int(row["sink_completed"]) == int(row["total"]) for row in rows) + PY - 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/README-RU.md b/README-RU.md index 7be0e39..60ec7aa 100644 --- a/README-RU.md +++ b/README-RU.md @@ -961,6 +961,20 @@ CSV `bench/results/latency-async-contract.csv` (или в путь из Детали `LatencyRecorder` собраны в [`docs/benchmarks.md`](docs/benchmarks.md). +Для отдельного исследовательского прогона async pipeline используйте target +`logit_bench_pipeline_research` при включённых `LOGIT_BENCH_ENABLE=ON` и +`LOGIT_BENCH_WITH_SPDLOG=ON`. Он сохраняет индивидуальные CSV/JSONL receipts +и aggregate CSV; подробные параметры запуска и определения метрик описаны в +[`docs/benchmarks.md`](docs/benchmarks.md). Режим +`LOGIT_BENCH_RESEARCH_MODE=rate` добавляет одинаковую для обеих библиотек +управляемую target rate. В rate mode дополнительно записываются realized +submission rate и schedule lag; finite producer threads могут отставать от +target из-за blocking admission. Эти результаты являются исследованием +admission, backpressure и benchmark outstanding, а не прямым измерением +внутренней очереди или универсальным рейтингом скорости библиотек. Sink +latency в этом target — отдельная instrumented metric: она использует общий +call-start timestamp, но включает небольшой overhead reservation/telemetry. + ## Матрица бэкендов diff --git a/bench/CMakeLists.txt b/bench/CMakeLists.txt index c7f697c..e0a5cf7 100644 --- a/bench/CMakeLists.txt +++ b/bench/CMakeLists.txt @@ -138,4 +138,26 @@ if(LOGIT_BENCH_WITH_SPDLOG) ) endforeach() add_test(NAME logit_bench_flush_test COMMAND logit_bench_flush_test) + + add_executable(logit_bench_pipeline_research + pipeline_research.cpp + adapters/LogItAdapter.cpp + adapters/SpdlogAdapter.cpp + ) + target_include_directories(logit_bench_pipeline_research PRIVATE ${CMAKE_CURRENT_SOURCE_DIR}) + target_compile_definitions(logit_bench_pipeline_research PRIVATE + LOGIT_BENCH_HAVE_SPDLOG=1 + LOGIT_BENCH_CONTRACT_MATCHED_ASYNC=1 + LOGIT_BENCH_PIPELINE_RESEARCH=1 + LOGIT_BENCH_BUILD_TYPE="$" + ) + target_compile_features(logit_bench_pipeline_research PRIVATE cxx_std_17) + target_link_libraries(logit_bench_pipeline_research PRIVATE + log-it-cpp::log-it-cpp spdlog::spdlog) + set_target_properties(logit_bench_pipeline_research PROPERTIES + RUNTIME_OUTPUT_DIRECTORY ${CMAKE_BINARY_DIR}) + foreach(config IN ITEMS DEBUG RELEASE RELWITHDEBINFO MINSIZEREL) + set_target_properties(logit_bench_pipeline_research PROPERTIES + RUNTIME_OUTPUT_DIRECTORY_${config} ${CMAKE_BINARY_DIR}) + endforeach() endif() diff --git a/bench/LatencyRecorder.hpp b/bench/LatencyRecorder.hpp index 4eb23c0..d0a75ad 100644 --- a/bench/LatencyRecorder.hpp +++ b/bench/LatencyRecorder.hpp @@ -43,6 +43,7 @@ namespace logit_bench { std::uint64_t p50_ns = 0; std::uint64_t p99_ns = 0; std::uint64_t p999_ns = 0; + std::uint64_t max_ns = 0; }; explicit LatencyRecorder(std::size_t total) @@ -68,6 +69,11 @@ namespace logit_bench { * Also stores t0 into internal t0 array, enabling complete_slot(slot). */ Token begin(bool record) { + return begin_at(record, record ? now() : 0); + } + + /// Reserve a slot using a caller-provided timestamp. + Token begin_at(bool record, std::uint64_t t0_ns) { Token token; token.active = record; if (!record) return token; @@ -78,7 +84,7 @@ namespace logit_bench { } token.slot = static_cast(slot); - token.t0_ns = now(); + token.t0_ns = t0_ns; // Store t0 for "slot-only" completion path (spdlog-friendly). // Relaxed is fine: slot is unique per begin(); consumer uses the same slot. @@ -90,7 +96,13 @@ namespace logit_bench { /// Capture t1 and store (t1 - t0) into the reserved slot (deduplicated). void complete(const Token& token) { if (!token.active) return; - complete_impl(token.slot, token.t0_ns); + complete_at(token, now()); + } + + /// Complete a slot using a caller-provided end timestamp. + void complete_at(const Token& token, std::uint64_t t1_ns) { + if (!token.active) return; + complete_impl(token.slot, token.t0_ns, t1_ns); } /** @@ -103,7 +115,7 @@ namespace logit_bench { throw std::out_of_range("LatencyRecorder capacity exceeded"); } const auto t0 = m_t0_ns[static_cast(slot)]; - complete_impl(slot, t0); + complete_impl(slot, t0, now()); } std::size_t recorded() const { @@ -134,6 +146,7 @@ namespace logit_bench { summary.p50_ns = pick(sorted, 0.50); summary.p99_ns = pick(sorted, 0.99); summary.p999_ns = pick(sorted, 0.999); + summary.max_ns = sorted.back(); return summary; } @@ -158,7 +171,8 @@ namespace logit_bench { return data[idx]; } - void complete_impl(std::uint64_t slot_u64, std::uint64_t t0_ns) { + void complete_impl(std::uint64_t slot_u64, std::uint64_t t0_ns, + std::uint64_t t1_ns) { if (slot_u64 >= m_expected) { throw std::out_of_range("LatencyRecorder capacity exceeded"); } @@ -172,7 +186,6 @@ namespace logit_bench { return; // duplicate completion -> ignore } - const auto t1_ns = now(); m_values[slot] = t1_ns - t0_ns; const auto done = m_completed.fetch_add(1, std::memory_order_acq_rel) + 1; diff --git a/bench/Scenario.hpp b/bench/Scenario.hpp index 9811979..e2aab08 100644 --- a/bench/Scenario.hpp +++ b/bench/Scenario.hpp @@ -1,7 +1,11 @@ #pragma once #include +#include +#include +#include #include +#include #include #include @@ -17,6 +21,49 @@ enum class AsyncPayloadMode { FullMessage, }; +// Benchmark-only, library-neutral pipeline counters. The harness records an +// issued call before entering adapter.log() and a sink completion after the +// adapter callback. This compares admission and drain behaviour without +// reaching into either library's private queue. +struct BenchmarkTelemetry { + std::atomic issued{0}; + std::atomic sink_completed{0}; + std::atomic high_water{0}; + std::atomic last_sink_entry_ns{0}; + + static std::uint64_t now_ns() { + return static_cast(std::chrono::duration_cast( + std::chrono::steady_clock::now().time_since_epoch()).count()); + } + + void on_issued() { + const auto current = issued.fetch_add(1, std::memory_order_acq_rel) + 1; + const auto completed_now = sink_completed.load(std::memory_order_acquire); + const auto outstanding = current > completed_now ? current - completed_now : 0; + auto observed = high_water.load(std::memory_order_relaxed); + while (observed < outstanding && + !high_water.compare_exchange_weak(observed, outstanding, + std::memory_order_relaxed, + std::memory_order_relaxed)) {} + } + + void on_sink_entry() { + const auto timestamp = now_ns(); + auto previous = last_sink_entry_ns.load(std::memory_order_relaxed); + while (previous < timestamp && + !last_sink_entry_ns.compare_exchange_weak(previous, timestamp, + std::memory_order_relaxed, + std::memory_order_relaxed)) {} + sink_completed.fetch_add(1, std::memory_order_acq_rel); + } + + std::uint64_t outstanding() const { + const auto accepted = issued.load(std::memory_order_acquire); + const auto completed = sink_completed.load(std::memory_order_acquire); + return accepted > completed ? accepted - completed : 0; + } +}; + inline std::string sink_name(SinkKind sink) { switch (sink) { case SinkKind::Null: return "null"; @@ -34,6 +81,7 @@ struct Scenario { std::size_t message_bytes = 0; std::size_t total_messages = 0; std::size_t queue_capacity = 0; + std::shared_ptr telemetry; }; } // namespace logit_bench diff --git a/bench/adapters/LogItAdapter.cpp b/bench/adapters/LogItAdapter.cpp index ee7d06d..adbfedb 100644 --- a/bench/adapters/LogItAdapter.cpp +++ b/bench/adapters/LogItAdapter.cpp @@ -40,6 +40,7 @@ namespace logit_bench { m_sink = scenario.sink; m_async_payload = scenario.async_payload; m_async_payload_observer = scenario.async_payload_observer; + m_telemetry = scenario.telemetry; m_recorder = &recorder; if (m_sink == SinkKind::File) { @@ -134,6 +135,9 @@ namespace logit_bench { if (slot_line >= 0 && m_recorder) { m_recorder->complete_slot(static_cast(slot_line)); } + if (m_telemetry) { + m_telemetry->on_sink_entry(); + } if (m_async_payload_observer) { m_async_payload_observer(text); @@ -156,6 +160,7 @@ namespace logit_bench { SinkKind m_sink = SinkKind::Null; AsyncPayloadMode m_async_payload = AsyncPayloadMode::MarkerOnly; std::function m_async_payload_observer; + std::shared_ptr m_telemetry; LatencyRecorder* m_recorder = nullptr; std::ofstream m_file; diff --git a/bench/adapters/SpdlogAdapter.cpp b/bench/adapters/SpdlogAdapter.cpp index b6170a6..f5fad47 100644 --- a/bench/adapters/SpdlogAdapter.cpp +++ b/bench/adapters/SpdlogAdapter.cpp @@ -34,6 +34,7 @@ namespace logit_bench { void configure(const Scenario& scenario, std::shared_ptr recorder) { m_sink = scenario.sink; m_recorder = std::move(recorder); + m_telemetry = scenario.telemetry; m_delay_ms = 0; if (const char* delay = std::getenv("LOGIT_BENCH_SPDLOG_SINK_DELAY_MS")) { try { @@ -65,6 +66,9 @@ namespace logit_bench { if (line >= 0 && m_recorder) { m_recorder->complete_slot(static_cast(line)); } + if (m_telemetry) { + m_telemetry->on_sink_entry(); + } if (m_sink == SinkKind::File) { std::lock_guard lock(m_mutex); @@ -106,6 +110,7 @@ namespace logit_bench { SinkKind m_sink = SinkKind::Null; std::shared_ptr m_recorder; + std::shared_ptr m_telemetry; std::size_t m_delay_ms = 0; std::ofstream m_file; diff --git a/bench/pipeline_research.cpp b/bench/pipeline_research.cpp new file mode 100644 index 0000000..1bce68c --- /dev/null +++ b/bench/pipeline_research.cpp @@ -0,0 +1,508 @@ +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "BenchmarkMetadata.hpp" +#include "LatencyRecorder.hpp" +#include "Scenario.hpp" +#include "adapters/ILoggerAdapter.hpp" +#include "adapters/LogItAdapter.hpp" +#ifdef LOGIT_BENCH_HAVE_SPDLOG +#include "adapters/SpdlogAdapter.hpp" +#endif + +namespace logit_bench { +namespace { + +using Clock = std::chrono::steady_clock; + +struct Options { + std::string mode = "matrix"; + std::size_t total = 200000; + std::size_t warmup = 4096; + std::size_t repeats = 5; + std::size_t bytes = 200; + std::vector producers{1, 2, 4, 8}; + std::vector queues{1024, 8192, 65536, 400000}; + std::vector rates{100000, 250000, 500000, 750000, 1000000}; + std::filesystem::path csv = "bench/results/pipeline-research.csv"; + std::filesystem::path jsonl = "bench/results/pipeline-research.jsonl"; + std::filesystem::path aggregate = "bench/results/pipeline-research-aggregate.csv"; +}; + +std::size_t env_size(const char* name, std::size_t fallback) { + if (const char* value = std::getenv(name)) { + try { return static_cast(std::stoull(value)); } + catch (...) {} + } + return fallback; +} + +std::string env_string(const char* name, const char* fallback) { + if (const char* value = std::getenv(name); value && *value) return value; + return fallback; +} + +std::vector parse_list(const char* name, + std::vector fallback) { + const char* value = std::getenv(name); + if (!value || !*value) return fallback; + std::vector result; + std::stringstream input(value); + std::string item; + while (std::getline(input, item, ',')) { + try { result.push_back(static_cast(std::stoull(item))); } + catch (...) { throw std::invalid_argument(std::string("Invalid list in ") + name); } + } + if (result.empty()) throw std::invalid_argument(std::string("Empty list in ") + name); + return result; +} + +std::uint64_t now_ns() { + return static_cast(std::chrono::duration_cast( + Clock::now().time_since_epoch()).count()); +} + +std::string message_for(std::size_t bytes) { + return std::string(bytes, 'X'); +} + +struct StageResult { + LatencyRecorder::Summary producer_call; + LatencyRecorder::Summary sink_entry; + LatencyRecorder::Summary schedule_lag; + std::uint64_t producer_phase_ns = 0; + std::uint64_t drain_tail_ns = 0; + std::uint64_t total_wall_ns = 0; + double throughput = 0.0; + double realized_submission_rate = 0.0; + std::uint64_t outstanding_high_water = 0; + std::uint64_t outstanding_at_producer_done = 0; + std::uint64_t issued = 0; + std::uint64_t sink_completed = 0; +}; + +struct RunRecord { + std::string mode; + std::size_t repeat = 0; + std::size_t run_index = 0; + std::string run_order; + std::string library; + std::size_t producers = 0; + std::size_t queue = 0; + std::size_t total = 0; + std::size_t warmup = 0; + std::size_t bytes = 0; + std::size_t target_rate = 0; + StageResult result; +}; + +void validate_csv(const std::filesystem::path& path, const std::string& header) { + namespace fs = std::filesystem; + if (path.parent_path() != fs::path()) fs::create_directories(path.parent_path()); + if (!fs::exists(path) || fs::file_size(path) == 0) return; + std::ifstream input(path); + std::string actual; + if (!input || !std::getline(input, actual) || actual != header) { + throw std::runtime_error("Unsupported pipeline research CSV schema; rename or remove the existing file"); + } +} + +const char* csv_header() { + return "mode,repeat,run_index,run_order,library,producers,queue_capacity,total,warmup,msg_bytes,target_rate," + "realized_submission_rate,schedule_lag_p50_ns,schedule_lag_p99_ns,schedule_lag_max_ns," + "producer_p50_ns,producer_p99_ns,producer_p999_ns,sink_p50_ns,sink_p99_ns,sink_p999_ns," + "producer_phase_ns,drain_tail_ns,total_wall_ns,throughput,outstanding_high_water," + "outstanding_at_producer_done,issued,sink_completed,source_commit,compiler,compiler_version," + "toolchain,cxx_standard,platform,build_type,architecture,machine_id,cpu_model,queue_policy," + "latency_completion,flush_barrier,workload_contract"; +} + +void append_record(const Options& options, const RunRecord& r, + const BenchmarkMetadata& metadata) { + validate_csv(options.csv, csv_header()); + const bool header = !std::filesystem::exists(options.csv) || + std::filesystem::file_size(options.csv) == 0; + std::ofstream out(options.csv, std::ios::app); + if (!out) throw std::runtime_error("Failed to open pipeline research CSV"); + if (header) out << csv_header() << '\n'; + const auto& p = r.result.producer_call; + const auto& s = r.result.sink_entry; + out << r.mode << ',' << r.repeat << ',' << r.run_index << ',' << r.run_order << ',' + << r.library << ',' << r.producers << ',' << r.queue << ',' << r.total << ',' + << r.warmup << ',' << r.bytes << ',' << r.target_rate << ',' + << std::fixed << std::setprecision(2) << r.result.realized_submission_rate << ',' + << r.result.schedule_lag.p50_ns << ',' << r.result.schedule_lag.p99_ns << ',' + << r.result.schedule_lag.max_ns << ',' + << p.p50_ns << ',' << p.p99_ns << ',' << p.p999_ns << ',' + << s.p50_ns << ',' << s.p99_ns << ',' << s.p999_ns << ',' + << r.result.producer_phase_ns << ',' << r.result.drain_tail_ns << ',' + << r.result.total_wall_ns << ',' << std::fixed << std::setprecision(2) + << r.result.throughput << ',' << r.result.outstanding_high_water << ',' + << r.result.outstanding_at_producer_done << ',' << r.result.issued << ',' + << r.result.sink_completed << ',' << metadata.source_commit << ',' << metadata.compiler << ',' + << metadata.compiler_version << ',' << metadata.toolchain << ',' << metadata.cxx_standard << ',' + << metadata.platform << ',' << metadata.build_type << ',' << metadata.architecture << ',' + << metadata.machine_id << ',' << metadata.cpu_model << ',' << metadata.queue_policy << ',' + << metadata.latency_completion << ',' << metadata.flush_barrier << ',' + << metadata.workload_contract << '\n'; + + if (options.jsonl.parent_path() != std::filesystem::path()) + std::filesystem::create_directories(options.jsonl.parent_path()); + std::ofstream json(options.jsonl, std::ios::app); + if (!json) throw std::runtime_error("Failed to open pipeline research JSONL"); + json << "{\"fixture_version\":2,\"source_commit\":\"" << metadata.source_commit + << "\",\"compiler\":\"" << metadata.compiler + << "\",\"compiler_version\":\"" << metadata.compiler_version + << "\",\"toolchain\":\"" << metadata.toolchain + << "\",\"cxx_standard\":\"" << metadata.cxx_standard + << "\",\"platform\":\"" << metadata.platform + << "\",\"build_type\":\"" << metadata.build_type + << "\",\"architecture\":\"" << metadata.architecture + << "\",\"machine_id\":\"" << metadata.machine_id + << "\",\"cpu_model\":\"" << metadata.cpu_model + << "\",\"queue_policy\":\"" << metadata.queue_policy + << "\",\"latency_completion\":\"" << metadata.latency_completion + << "\",\"flush_barrier\":\"" << metadata.flush_barrier + << "\",\"workload_contract\":\"" << metadata.workload_contract + << "\",\"mode\":\"" << r.mode << "\",\"repeat\":" + << r.repeat << ",\"run_index\":" << r.run_index << ",\"run_order\":\"" + << r.run_order << "\",\"library\":\"" << r.library << "\",\"producers\":" + << r.producers << ",\"queue_capacity\":" << r.queue << ",\"total\":" + << r.total << ",\"warmup\":" << r.warmup << ",\"msg_bytes\":" << r.bytes + << ",\"target_rate\":" << r.target_rate + << ",\"realized_submission_rate\":" << std::fixed << std::setprecision(2) + << r.result.realized_submission_rate + << ",\"schedule_lag_p50_ns\":" << r.result.schedule_lag.p50_ns + << ",\"schedule_lag_p99_ns\":" << r.result.schedule_lag.p99_ns + << ",\"schedule_lag_max_ns\":" << r.result.schedule_lag.max_ns + << ",\"producer_p50_ns\":" << p.p50_ns << ",\"producer_p99_ns\":" << p.p99_ns + << ",\"producer_p999_ns\":" << p.p999_ns << ",\"sink_p50_ns\":" << s.p50_ns + << ",\"sink_p99_ns\":" << s.p99_ns << ",\"sink_p999_ns\":" << s.p999_ns + << ",\"producer_phase_ns\":" << r.result.producer_phase_ns + << ",\"drain_tail_ns\":" << r.result.drain_tail_ns + << ",\"total_wall_ns\":" << r.result.total_wall_ns + << ",\"throughput\":" << std::fixed << std::setprecision(2) << r.result.throughput + << ",\"outstanding_high_water\":" << r.result.outstanding_high_water + << ",\"outstanding_at_producer_done\":" << r.result.outstanding_at_producer_done + << ",\"issued\":" << r.result.issued << ",\"sink_completed\":" + << r.result.sink_completed << "}\n"; +} + +void run_warmup(ILoggerAdapter& adapter, const Scenario& scenario, + std::string_view message, std::size_t count) { + std::mutex mx; + std::condition_variable cv; + std::size_t ready = 0; + bool released = false; + std::vector threads; + threads.reserve(scenario.producers); + for (std::size_t producer = 0; producer < scenario.producers; ++producer) { + threads.emplace_back([&, producer]() { + const auto share = count / scenario.producers + + (producer < count % scenario.producers ? 1 : 0); + { + std::unique_lock lock(mx); + ++ready; + cv.notify_all(); + cv.wait(lock, [&] { return released; }); + } + for (std::size_t n = 0; n < share; ++n) { + adapter.log(LatencyRecorder::Token{}, message); + } + }); + } + { + std::unique_lock lock(mx); + cv.wait(lock, [&] { return ready == scenario.producers; }); + released = true; + cv.notify_all(); + } + for (auto& thread : threads) thread.join(); +} + +StageResult run_once(ILoggerAdapter& adapter, Scenario scenario, + std::size_t warmup, std::size_t target_rate) { + auto warmup_recorder = std::make_shared(1); + scenario.telemetry.reset(); + adapter.prepare(scenario, *warmup_recorder); + const auto message = message_for(scenario.message_bytes); + + run_warmup(adapter, scenario, message, warmup); + adapter.flush(); + + auto sink_recorder = std::make_shared(scenario.total_messages); + auto producer_recorder = std::make_shared(scenario.total_messages); + auto telemetry = std::make_shared(); + scenario.telemetry = telemetry; + adapter.prepare(scenario, *sink_recorder); + std::uint64_t start = 0; + std::atomic last_producer_done{0}; + std::mutex producer_done_mx; + std::uint64_t outstanding_at_done_snapshot = 0; + std::mutex mx; + std::condition_variable cv; + std::size_t ready = 0; + bool released = false; + std::atomic next_ticket{0}; + const auto interval = target_rate ? (1'000'000'000ULL / target_rate) : 0; + auto schedule_recorder = target_rate + ? std::make_shared(scenario.total_messages) + : std::shared_ptr(); + std::vector threads; + threads.reserve(scenario.producers); + for (std::size_t producer = 0; producer < scenario.producers; ++producer) { + threads.emplace_back([&, producer]() { + const auto share = scenario.total_messages / scenario.producers + + (producer < scenario.total_messages % scenario.producers ? 1 : 0); + { + std::unique_lock lock(mx); + ++ready; + cv.notify_all(); + cv.wait(lock, [&] { return released; }); + } + for (std::size_t n = 0; n < share; ++n) { + std::uint64_t scheduled_start = 0; + if (target_rate) { + const auto ticket = next_ticket.fetch_add(1, std::memory_order_relaxed); + scheduled_start = start + ticket * interval; + for (;;) { + const auto current = now_ns(); + if (current >= scheduled_start) break; + const auto remaining = scheduled_start - current; + if (remaining > 200'000) std::this_thread::sleep_for( + std::chrono::nanoseconds(remaining / 2)); + else std::this_thread::yield(); + } + } + const auto call_start = now_ns(); + auto sink_token = sink_recorder->begin_at(true, call_start); + auto producer_token = producer_recorder->begin_at(true, call_start); + auto schedule_token = schedule_recorder + ? schedule_recorder->begin_at(true, scheduled_start) + : LatencyRecorder::Token{}; + telemetry->on_issued(); + adapter.log(sink_token, message); + const auto returned = now_ns(); + if (schedule_recorder) { + schedule_recorder->complete_at(schedule_token, call_start); + } + if (n + 1 == share) { + std::lock_guard lock(producer_done_mx); + if (returned > last_producer_done.load(std::memory_order_relaxed)) { + last_producer_done.store(returned, std::memory_order_relaxed); + outstanding_at_done_snapshot = telemetry->outstanding(); + } + } + producer_recorder->complete(producer_token); + } + }); + } + { + std::unique_lock lock(mx); + cv.wait(lock, [&] { return ready == scenario.producers; }); + start = now_ns(); + released = true; + cv.notify_all(); + } + for (auto& thread : threads) thread.join(); + const auto producer_done = last_producer_done.load(std::memory_order_acquire); + const auto outstanding_at_done = outstanding_at_done_snapshot; + adapter.flush(); + const auto end = now_ns(); + sink_recorder->wait_for_all(); + producer_recorder->wait_for_all(); + + StageResult result; + result.producer_call = producer_recorder->finalize(); + result.sink_entry = sink_recorder->finalize(); + if (schedule_recorder) { + result.schedule_lag = schedule_recorder->finalize(); + } + result.producer_phase_ns = producer_done > start ? producer_done - start : 0; + const auto last_sink = telemetry->last_sink_entry_ns.load(std::memory_order_acquire); + result.drain_tail_ns = last_sink > producer_done ? last_sink - producer_done : 0; + result.total_wall_ns = end > start ? end - start : 0; + result.throughput = result.total_wall_ns + ? static_cast(scenario.total_messages) * 1'000'000'000.0 / + static_cast(result.total_wall_ns) : 0.0; + result.realized_submission_rate = result.producer_phase_ns + ? static_cast(scenario.total_messages) * 1'000'000'000.0 / + static_cast(result.producer_phase_ns) : 0.0; + result.outstanding_high_water = telemetry->high_water.load(std::memory_order_acquire); + result.outstanding_at_producer_done = outstanding_at_done; + result.issued = telemetry->issued.load(std::memory_order_acquire); + result.sink_completed = telemetry->sink_completed.load(std::memory_order_acquire); + if (result.issued != scenario.total_messages || + result.sink_completed != scenario.total_messages) { + throw std::runtime_error("pipeline telemetry did not drain all messages"); + } + return result; +} + +double median(std::vector values) { + if (values.empty()) return 0.0; + std::sort(values.begin(), values.end()); + return values[values.size() / 2]; +} + +void write_aggregate(const Options& options, const std::vector& records) { + struct Group { + std::string mode, library; + std::size_t producers{}, queue{}, rate{}; + std::vector throughput, realized_rate, producer_p50, producer_p99; + std::vector sink_p50, sink_p99, schedule_lag_p50, schedule_lag_p99, schedule_lag_max; + std::vector producer_phase, drain_tail, total_wall, high_water, at_done; + }; + std::vector groups; + for (const auto& record : records) { + auto it = std::find_if(groups.begin(), groups.end(), [&](const Group& g) { + return g.mode == record.mode && g.library == record.library && + g.producers == record.producers && g.queue == record.queue && + g.rate == record.target_rate; + }); + if (it == groups.end()) { + groups.push_back(Group{record.mode, record.library, record.producers, + record.queue, record.target_rate}); + it = groups.end() - 1; + } + it->throughput.push_back(record.result.throughput); + it->realized_rate.push_back(record.result.realized_submission_rate); + it->producer_p50.push_back(static_cast(record.result.producer_call.p50_ns)); + it->producer_p99.push_back(static_cast(record.result.producer_call.p99_ns)); + it->sink_p50.push_back(static_cast(record.result.sink_entry.p50_ns)); + it->sink_p99.push_back(static_cast(record.result.sink_entry.p99_ns)); + it->schedule_lag_p50.push_back(static_cast(record.result.schedule_lag.p50_ns)); + it->schedule_lag_p99.push_back(static_cast(record.result.schedule_lag.p99_ns)); + it->schedule_lag_max.push_back(static_cast(record.result.schedule_lag.max_ns)); + it->producer_phase.push_back(static_cast(record.result.producer_phase_ns)); + it->drain_tail.push_back(static_cast(record.result.drain_tail_ns)); + it->total_wall.push_back(static_cast(record.result.total_wall_ns)); + it->high_water.push_back(static_cast(record.result.outstanding_high_water)); + it->at_done.push_back(static_cast(record.result.outstanding_at_producer_done)); + } + if (options.aggregate.parent_path() != std::filesystem::path()) + std::filesystem::create_directories(options.aggregate.parent_path()); + std::ofstream out(options.aggregate); + out << "mode,library,producers,queue_capacity,target_rate,repeats," + "median_realized_submission_rate,median_schedule_lag_p50_ns,median_schedule_lag_p99_ns," + "median_schedule_lag_max_ns," + "median_producer_p50_ns,median_producer_p99_ns,median_sink_p50_ns,median_sink_p99_ns," + "median_producer_phase_ns,median_drain_tail_ns,median_total_wall_ns,median_throughput," + "median_outstanding_high_water,median_outstanding_at_producer_done\n"; + for (const auto& group : groups) { + out << group.mode << ',' << group.library << ',' << group.producers << ',' << group.queue + << ',' << group.rate << ',' << group.throughput.size() << ',' + << std::fixed << std::setprecision(2) + << median(group.realized_rate) << ',' << median(group.schedule_lag_p50) << ',' + << median(group.schedule_lag_p99) << ',' << median(group.schedule_lag_max) << ',' + << median(group.producer_p50) << ',' + << median(group.producer_p99) << ',' << median(group.sink_p50) << ',' + << median(group.sink_p99) << ',' << median(group.producer_phase) << ',' + << median(group.drain_tail) << ',' << median(group.total_wall) << ',' + << median(group.throughput) << ',' << median(group.high_water) << ',' + << median(group.at_done) << '\n'; + } +} + +} // namespace +} // namespace logit_bench + +int main() { + using namespace logit_bench; +#ifndef LOGIT_BENCH_HAVE_SPDLOG + std::cerr << "pipeline research requires LOGIT_BENCH_WITH_SPDLOG=ON\n"; + return 2; +#else + try { + Options options; + options.mode = env_string("LOGIT_BENCH_RESEARCH_MODE", "matrix"); + options.total = env_size("LOGIT_BENCH_TOTAL", options.total); + options.warmup = env_size("LOGIT_BENCH_WARMUP", options.warmup); + options.repeats = env_size("LOGIT_BENCH_REPEATS", options.repeats); + options.bytes = env_size("LOGIT_BENCH_BYTES", options.bytes); + options.producers = parse_list("LOGIT_BENCH_PRODUCERS", options.producers); + options.queues = parse_list("LOGIT_BENCH_QUEUE_CAPACITIES", options.queues); + options.rates = parse_list("LOGIT_BENCH_RATES", options.rates); + options.csv = env_string("LOGIT_BENCH_RESEARCH_CSV", options.csv.string().c_str()); + options.jsonl = env_string("LOGIT_BENCH_RESEARCH_JSONL", options.jsonl.string().c_str()); + options.aggregate = env_string("LOGIT_BENCH_RESEARCH_AGGREGATE", options.aggregate.string().c_str()); + if (options.repeats == 0 || options.total == 0 || options.producers.empty()) + throw std::invalid_argument("research totals, repeats, and producers must be non-zero"); + + const auto metadata = make_benchmark_metadata( + "matrix", "block", "sink-entry", "all-prior-work-drained", + "prepared-message/async-full-message"); + validate_comparable_metadata(metadata); + std::vector records; + auto logit_adapter = std::make_unique(); + auto spdlog_adapter = std::make_unique(); + const auto run_mode = options.mode == "rate" ? std::string("rate") : std::string("matrix"); + const auto rates = run_mode == "rate" ? options.rates : std::vector{0}; + const auto queues = run_mode == "rate" ? std::vector{400000} : options.queues; + for (const auto rate : rates) { + for (const auto producers : options.producers) { + for (const auto queue : queues) { + // Keep the two libraries adjacent for each matched key point; + // only the order within that point alternates by repeat. + for (std::size_t repeat = 1; repeat <= options.repeats; ++repeat) { + std::vector order{"log-it-cpp", "spdlog"}; + if ((repeat % 2) == 0) std::swap(order[0], order[1]); + for (const auto& library : order) { + ILoggerAdapter& adapter = library == "log-it-cpp" + ? static_cast(*logit_adapter) + : static_cast(*spdlog_adapter); + Scenario scenario; + scenario.async = true; + scenario.sink = SinkKind::Null; + scenario.async_payload = AsyncPayloadMode::FullMessage; + scenario.producers = producers; + scenario.message_bytes = options.bytes; + scenario.total_messages = options.total; + scenario.queue_capacity = queue; + LOGIT_SET_MAX_QUEUE(queue); + LOGIT_SET_QUEUE_POLICY(LOGIT_QUEUE_BLOCK); + auto result = run_once(adapter, scenario, options.warmup, rate); + RunRecord record{run_mode, repeat, records.size() + 1, + (repeat % 2) ? "logit-then-spdlog" : "spdlog-then-logit", + library, producers, queue, options.total, + options.warmup, options.bytes, rate, result}; + append_record(options, record, metadata); + records.push_back(std::move(record)); + std::cout << library << " mode=" << run_mode << " repeat=" << repeat + << " producers=" << producers << " queue=" << queue + << " target_rate=" << rate << " realized_rate=" + << result.realized_submission_rate << " sink_p50=" + << result.sink_entry.p50_ns + << " producer_p50=" << result.producer_call.p50_ns + << " throughput=" << std::fixed << std::setprecision(2) + << result.throughput << " outstanding_high_water=" + << result.outstanding_high_water << '\n'; + } + } + } + } + } + write_aggregate(options, records); + return 0; + } catch (const std::exception& ex) { + std::cerr << "Pipeline research failed: " << ex.what() << '\n'; + return 1; + } +#endif +} diff --git a/docs/benchmarks.md b/docs/benchmarks.md index d33c8d8..5696d86 100644 --- a/docs/benchmarks.md +++ b/docs/benchmarks.md @@ -118,6 +118,70 @@ a performance measurement. It configures the matched null-sink path with a 200-byte payload and verifies the exact payload observed by the worker-side sink callback. +## Async pipeline research target + +`logit_bench_pipeline_research` is a separate, benchmark-only instrumented +target for decomposing the matched `async/null/full-message` pipeline. Build it +with `LOGIT_BENCH_ENABLE=ON` and `LOGIT_BENCH_WITH_SPDLOG=ON`. The default +matrix uses 200-byte messages, 200,000 measured messages, 4,096 warmup +messages, producers `1,2,4,8`, queue capacities `1024,8192,65536,400000`, +five repeats, and alternates library order on each matched key point. A +rate-controlled run is selected with `LOGIT_BENCH_RESEARCH_MODE=rate`; its +default target rates are 100k, 250k, 500k, 750k, and 1M messages/second and it uses a +400,000-entry queue. + +For example, a short exploratory run is: + +```powershell +$env:LOGIT_BENCH_TOTAL = "200000" +$env:LOGIT_BENCH_WARMUP = "4096" +$env:LOGIT_BENCH_REPEATS = "5" +$env:LOGIT_BENCH_PRODUCERS = "1,2,4,8" +$env:LOGIT_BENCH_QUEUE_CAPACITIES = "1024,8192,65536,400000" +./build/logit_bench_pipeline_research +``` + +The target writes individual rows to `pipeline-research.csv` and JSONL +receipts to `pipeline-research.jsonl`, plus a median-per-key-point +`pipeline-research-aggregate.csv`. Paths can be overridden with +`LOGIT_BENCH_RESEARCH_CSV`, `LOGIT_BENCH_RESEARCH_JSONL`, and +`LOGIT_BENCH_RESEARCH_AGGREGATE`. Each row retains fixture metadata and the +deterministic `run_order`; the aggregate also retains producer/sink p50/p99, +producer phase, drain tail, total wall time, realized rate, schedule lag, and +both benchmark-outstanding summaries. + +The metrics are deliberately library-neutral. `producer_p50/p99/p999_ns` +measure only the `adapter.log()` call; rate limiting is outside that timed +region. Both producer and sink recorders use the same benchmark call-start +timestamp. This gives the two metrics a common origin, but the research +measurement still includes the small recorder reservation and telemetry +overhead between that timestamp and entering `adapter.log()`. Therefore +`sink_p50/p99/p999_ns` is an **instrumented sink-entry metric** and must not be +compared as an identical absolute quantity with older matched-benchmark runs. +`producer_phase_ns` runs from the shared producer release barrier to the last +producer return, `drain_tail_ns` runs from the last producer return to the +final sink entry, and `total_wall_ns` ends after the drain barrier. +`throughput` is measured messages/second over that total interval. + +In rate mode, `target_rate` is a pacing schedule, not a guaranteed external +open-loop arrival rate: a finite producer can receive its next ticket only +after its previous blocking `adapter.log()` returns. `realized_submission_rate` +is the issued count divided by the producer phase. `schedule_lag_p50/p99/max_ns` +records how far actual call start falls behind its scheduled start; growing lag +means the finite producer set cannot keep up with the target schedule. + +`issued - sink_completed` is named **benchmark outstanding**. It is a common +benchmark counter, not a direct queue-depth measurement: it can include calls +in admission/backpressure and must not be presented as either library's +private queue size. `outstanding_high_water` and +`outstanding_at_producer_done` use this benchmark-issued semantics. + +These results describe admission, backpressure, worker scheduling, and queue +backlog under the selected workload. They must not be reported as intrinsic +library latency or as a universal speed ranking. In particular, a lower +sink-entry p50 can coexist with lower throughput when a producer-side path +applies stronger admission pressure. + `logit_exec_mx_bench` and `logit_exec_mx_bench_concurrent` are a guarded lock-elision experiment. They use the same prepared `LogRecord` and a small thread-safe counting backend/formatter pair; the first target keeps the default