Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 46 additions & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
14 changes: 14 additions & 0 deletions README-RU.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.


## Матрица бэкендов

Expand Down
22 changes: 22 additions & 0 deletions bench/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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="$<CONFIG>"
)
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()
23 changes: 18 additions & 5 deletions bench/LatencyRecorder.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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;
Expand All @@ -78,7 +84,7 @@ namespace logit_bench {
}

token.slot = static_cast<std::uint64_t>(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.
Expand All @@ -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);
}

/**
Expand All @@ -103,7 +115,7 @@ namespace logit_bench {
throw std::out_of_range("LatencyRecorder capacity exceeded");
}
const auto t0 = m_t0_ns[static_cast<std::size_t>(slot)];
complete_impl(slot, t0);
complete_impl(slot, t0, now());
}

std::size_t recorded() const {
Expand Down Expand Up @@ -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;
}

Expand All @@ -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");
}
Expand All @@ -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;
Expand Down
48 changes: 48 additions & 0 deletions bench/Scenario.hpp
Original file line number Diff line number Diff line change
@@ -1,7 +1,11 @@
#pragma once

#include <cstddef>
#include <atomic>
#include <chrono>
#include <cstdint>
#include <functional>
#include <memory>
#include <string>
#include <string_view>

Expand All @@ -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<std::uint64_t> issued{0};
std::atomic<std::uint64_t> sink_completed{0};
std::atomic<std::uint64_t> high_water{0};
std::atomic<std::uint64_t> last_sink_entry_ns{0};

static std::uint64_t now_ns() {
return static_cast<std::uint64_t>(std::chrono::duration_cast<std::chrono::nanoseconds>(
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";
Expand All @@ -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<BenchmarkTelemetry> telemetry;
};

} // namespace logit_bench
5 changes: 5 additions & 0 deletions bench/adapters/LogItAdapter.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -134,6 +135,9 @@ namespace logit_bench {
if (slot_line >= 0 && m_recorder) {
m_recorder->complete_slot(static_cast<std::uint64_t>(slot_line));
}
if (m_telemetry) {
m_telemetry->on_sink_entry();
}

if (m_async_payload_observer) {
m_async_payload_observer(text);
Expand All @@ -156,6 +160,7 @@ namespace logit_bench {
SinkKind m_sink = SinkKind::Null;
AsyncPayloadMode m_async_payload = AsyncPayloadMode::MarkerOnly;
std::function<void(std::string_view)> m_async_payload_observer;
std::shared_ptr<BenchmarkTelemetry> m_telemetry;
LatencyRecorder* m_recorder = nullptr;

std::ofstream m_file;
Expand Down
5 changes: 5 additions & 0 deletions bench/adapters/SpdlogAdapter.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ namespace logit_bench {
void configure(const Scenario& scenario, std::shared_ptr<LatencyRecorder> 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 {
Expand Down Expand Up @@ -65,6 +66,9 @@ namespace logit_bench {
if (line >= 0 && m_recorder) {
m_recorder->complete_slot(static_cast<std::uint64_t>(line));
}
if (m_telemetry) {
m_telemetry->on_sink_entry();
}

if (m_sink == SinkKind::File) {
std::lock_guard<std::mutex> lock(m_mutex);
Expand Down Expand Up @@ -106,6 +110,7 @@ namespace logit_bench {

SinkKind m_sink = SinkKind::Null;
std::shared_ptr<LatencyRecorder> m_recorder;
std::shared_ptr<BenchmarkTelemetry> m_telemetry;
std::size_t m_delay_ms = 0;

std::ofstream m_file;
Expand Down
Loading
Loading