Skip to content
Open
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
9 changes: 9 additions & 0 deletions README-RU.md
Original file line number Diff line number Diff line change
Expand Up @@ -961,6 +961,15 @@ 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` добавляет одинаковую для обеих библиотек
управляемую offered load. Эти результаты являются исследованием admission,
backpressure и backlog, а не универсальным рейтингом скорости библиотек.


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

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()
47 changes: 47 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,48 @@ enum class AsyncPayloadMode {
FullMessage,
};

// Benchmark-only, library-neutral pipeline counters. Adapters report the
// same two observable events so the research harness can compare admission
// and drain behaviour without reaching into either library's private queue.
struct BenchmarkTelemetry {
std::atomic<std::uint64_t> submitted{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_submitted() {
const auto current = submitted.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 = submitted.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 +80,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