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
10 changes: 9 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_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 logit_producer_profile 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 logit_producer_profile logit_task_ownership_profile benchmark_validation_test
- name: Run spdlog async flush regression
run: ./build-bench/logit_bench_flush_test
- name: Run public macro benchmark smoke
Expand All @@ -54,6 +54,14 @@ jobs:
LOGIT_PRODUCER_PROFILE_WARMUP: 200
LOGIT_PRODUCER_PROFILE_REPEATS: 1
run: ./build-bench/logit_producer_profile
- name: Run task ownership profiling smoke
if: matrix.std == 17
env:
LOGIT_TASK_OWNERSHIP_TOTAL: 200
LOGIT_TASK_OWNERSHIP_WARMUP: 20
LOGIT_TASK_OWNERSHIP_REPEATS: 1
LOGIT_TASK_OWNERSHIP_ROUNDTRIP_TOTAL: 50
run: ./build-bench/logit_task_ownership_profile
- name: Run async payload contract regression
run: ./build-bench/logit_bench_async_payload_contract_test
- name: Run logger hot-path A/B smoke
Expand Down
5 changes: 5 additions & 0 deletions bench/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,11 @@ target_compile_features(logit_producer_profile PRIVATE cxx_std_17)
target_link_libraries(logit_producer_profile PRIVATE log-it-cpp::log-it-cpp)
set_target_properties(logit_producer_profile PROPERTIES RUNTIME_OUTPUT_DIRECTORY ${CMAKE_BINARY_DIR})

add_executable(logit_task_ownership_profile task_ownership_profile.cpp)
target_compile_features(logit_task_ownership_profile PRIVATE cxx_std_17)
target_link_libraries(logit_task_ownership_profile PRIVATE log-it-cpp::log-it-cpp)
set_target_properties(logit_task_ownership_profile PROPERTIES RUNTIME_OUTPUT_DIRECTORY ${CMAKE_BINARY_DIR})

add_executable(benchmark_validation_test benchmark_validation_test.cpp)
target_compile_features(benchmark_validation_test PRIVATE cxx_std_17)
set_target_properties(benchmark_validation_test PROPERTIES RUNTIME_OUTPUT_DIRECTORY ${CMAKE_BINARY_DIR})
Expand Down
314 changes: 314 additions & 0 deletions bench/task_ownership_profile.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,314 @@
#include <algorithm>
#include <atomic>
#include <chrono>
#include <cstddef>
#include <cstdint>
#include <cstdlib>
#include <functional>
#include <iostream>
#include <string>
#include <thread>
#include <vector>

#include <logit.hpp>
#include <logit/detail/MpscRingAny.hpp>

namespace {

constexpr std::size_t kDefaultTotal = 100000;
constexpr std::size_t kDefaultWarmup = 10000;
constexpr std::size_t kDefaultRepeats = 5;
constexpr std::size_t kQueueCapacity = 262144;
const std::string kMessage(200, 'x');

#if defined(LOGIT_USE_MPSC_RING)
constexpr const char* kQueueBackend = "mpsc_ring";
#else
constexpr const char* kQueueBackend = "mutex_deque";
#endif

std::uint64_t g_observer = 0;
std::atomic<std::size_t> g_completed{0};

std::size_t env_size(const char* name, std::size_t fallback) {
if (const char* value = std::getenv(name)) {
try {
return static_cast<std::size_t>(std::stoull(value));
} catch (...) {
}
}
return fallback;
}

template <typename Loop, typename Cleanup>
double measure(Loop loop, Cleanup cleanup, std::size_t warmup,
std::size_t total, std::size_t repeats) {
loop(warmup);
cleanup(warmup);

std::vector<std::uint64_t> samples;
samples.reserve(repeats);
for (std::size_t repeat = 0; repeat < repeats; ++repeat) {
const auto start = std::chrono::steady_clock::now();
loop(total);
const auto elapsed = std::chrono::duration_cast<std::chrono::nanoseconds>(
std::chrono::steady_clock::now() - start).count();
cleanup(total);
samples.push_back(static_cast<std::uint64_t>(elapsed));
}

std::sort(samples.begin(), samples.end());
return static_cast<double>(samples[samples.size() / 2]) /
static_cast<double>(total);
}

#if defined(LOGIT_USE_MPSC_RING)
double measure_ring_prebuilt_std_function(std::size_t warmup,
std::size_t total,
std::size_t repeats) {
auto run = [&](std::size_t count) {
std::vector<std::function<void()>> tasks;
tasks.reserve(count);
for (std::size_t i = 0; i < count; ++i) {
tasks.emplace_back([]() { ++g_observer; });
}

logit::detail::MpscRingAny<std::function<void()>> ring(kQueueCapacity);
std::atomic<bool> started{false};
std::atomic<bool> done{false};
std::atomic<std::size_t> ready{0};
std::atomic<std::size_t> consumed{0};
std::thread worker([&]() {
ready.store(1, std::memory_order_release);
while (!started.load(std::memory_order_acquire)) {
std::this_thread::yield();
}
std::function<void()> task;
while (!done.load(std::memory_order_acquire) || !ring.empty()) {
if (ring.try_pop(task)) {
task();
consumed.fetch_add(1, std::memory_order_relaxed);
} else {
std::this_thread::yield();
}
}
});
while (ready.load(std::memory_order_acquire) == 0) {
std::this_thread::yield();
}

started.store(true, std::memory_order_release);
const auto start = std::chrono::steady_clock::now();
for (auto& task : tasks) {
while (!ring.try_push(std::move(task))) {
std::this_thread::yield();
}
}
const auto elapsed = std::chrono::duration_cast<std::chrono::nanoseconds>(
std::chrono::steady_clock::now() - start).count();
done.store(true, std::memory_order_release);
worker.join();
if (consumed.load(std::memory_order_relaxed) != count) {
std::cerr << "prebuilt MPSC ring completion mismatch\n";
std::exit(6);
}
return static_cast<std::uint64_t>(elapsed);
};

run(warmup);
std::vector<std::uint64_t> samples;
samples.reserve(repeats);
for (std::size_t repeat = 0; repeat < repeats; ++repeat) {
samples.push_back(run(total));
}
std::sort(samples.begin(), samples.end());
return static_cast<double>(samples[samples.size() / 2]) /
static_cast<double>(total);
}

template <typename T, typename Producer, typename Consumer>
double measure_ring(Producer producer, Consumer consumer, std::size_t warmup,
std::size_t total, std::size_t repeats) {
auto run = [&](std::size_t count) {
logit::detail::MpscRingAny<T> ring(kQueueCapacity);
std::atomic<bool> started{false};
std::atomic<bool> done{false};
std::atomic<std::size_t> ready{0};
std::atomic<std::size_t> consumed{0};
std::thread worker([&]() {
ready.store(1, std::memory_order_release);
while (!started.load(std::memory_order_acquire)) {
std::this_thread::yield();
}
T task;
while (!done.load(std::memory_order_acquire) || !ring.empty()) {
if (ring.try_pop(task)) {
consumer(task);
consumed.fetch_add(1, std::memory_order_relaxed);
} else {
std::this_thread::yield();
}
}
});
while (ready.load(std::memory_order_acquire) == 0) {
std::this_thread::yield();
}

started.store(true, std::memory_order_release);
const auto start = std::chrono::steady_clock::now();
producer(ring, count);
const auto elapsed = std::chrono::duration_cast<std::chrono::nanoseconds>(
std::chrono::steady_clock::now() - start).count();
done.store(true, std::memory_order_release);
worker.join();
if (consumed.load(std::memory_order_relaxed) != count) {
std::cerr << "MPSC ring completion mismatch\n";
std::exit(5);
}
return static_cast<std::uint64_t>(elapsed);
};

run(warmup);
std::vector<std::uint64_t> samples;
samples.reserve(repeats);
for (std::size_t repeat = 0; repeat < repeats; ++repeat) {
samples.push_back(run(total));
}
std::sort(samples.begin(), samples.end());
return static_cast<double>(samples[samples.size() / 2]) /
static_cast<double>(total);
}
#endif

} // namespace

int main() {
const std::size_t total = env_size("LOGIT_TASK_OWNERSHIP_TOTAL", kDefaultTotal);
const std::size_t warmup = env_size("LOGIT_TASK_OWNERSHIP_WARMUP", kDefaultWarmup);
const std::size_t repeats = env_size("LOGIT_TASK_OWNERSHIP_REPEATS", kDefaultRepeats);
const std::size_t roundtrip_total = std::min<std::size_t>(
env_size("LOGIT_TASK_OWNERSHIP_ROUNDTRIP_TOTAL", 1000), total);
const std::size_t roundtrip_warmup = std::min(roundtrip_total, warmup);
if (total == 0 || repeats == 0) {
std::cerr << "total and repeats must be positive\n";
return 2;
}

auto& executor = logit::detail::TaskExecutor::get_instance();
executor.set_queue_policy(logit::QueuePolicy::Block);
executor.set_max_queue_size(kQueueCapacity);
std::function<void()> prebuilt_task = []() {
g_completed.fetch_add(1, std::memory_order_relaxed);
};

std::cout << "task-ownership-profile total=" << total
<< " warmup=" << warmup
<< " repeats=" << repeats
<< " message_bytes=" << kMessage.size()
<< " queue_capacity=" << kQueueCapacity
<< " roundtrip_total=" << roundtrip_total
<< " queue_backend=" << kQueueBackend << '\n';

const double function_payload_ns = measure(
[](std::size_t count) {
for (std::size_t i = 0; i < count; ++i) {
std::string payload = kMessage;
std::function<void()> task =
[payload = std::move(payload)]() mutable {
g_observer += payload.size();
};
task();
}
},
[](std::size_t) {}, warmup, total, repeats);
std::cout << "case=function_construct_payload_invoke ns_per_call="
<< function_payload_ns << '\n';

const double function_copy_ns = measure(
[&prebuilt_task](std::size_t count) {
for (std::size_t i = 0; i < count; ++i) {
std::function<void()> copy = prebuilt_task;
copy();
}
},
[](std::size_t) {}, warmup, total, repeats);
std::cout << "case=function_copy_invoke ns_per_call="
<< function_copy_ns << '\n';
g_completed.store(0, std::memory_order_relaxed);

const double executor_batch_ns = measure(
[&executor, &prebuilt_task](std::size_t count) {
for (std::size_t i = 0; i < count; ++i) {
executor.add_task(prebuilt_task);
}
},
[&executor](std::size_t expected) {
executor.wait();
if (g_completed.load(std::memory_order_relaxed) != expected) {
std::cerr << "TaskExecutor batch completion mismatch\n";
std::exit(3);
}
g_completed.store(0, std::memory_order_relaxed);
}, warmup, total, repeats);
std::cout << "case=taskexecutor_batch_prebuilt ns_per_call="
<< executor_batch_ns << '\n';

const double executor_roundtrip_ns = measure(
[&executor, &prebuilt_task](std::size_t count) {
for (std::size_t i = 0; i < count; ++i) {
executor.add_task(prebuilt_task);
executor.wait();
}
},
[&executor](std::size_t expected) {
executor.wait();
if (g_completed.load(std::memory_order_relaxed) != expected) {
std::cerr << "TaskExecutor roundtrip completion mismatch\n";
std::exit(4);
}
g_completed.store(0, std::memory_order_relaxed);
}, roundtrip_warmup, roundtrip_total, repeats);
std::cout << "case=taskexecutor_roundtrip_prebuilt ns_per_call="
<< executor_roundtrip_ns << '\n';

#if defined(LOGIT_USE_MPSC_RING)
const double ring_construct_ns = measure_ring<std::function<void()>>(
[](auto& ring, std::size_t count) {
for (std::size_t i = 0; i < count; ++i) {
std::function<void()> value = []() { ++g_observer; };
while (!ring.try_push(std::move(value))) {
std::this_thread::yield();
}
}
},
[](std::function<void()>& task) { task(); },
warmup, total, repeats);
std::cout << "case=mpsc_ring_construct_and_publish_std_function ns_per_call="
<< ring_construct_ns << '\n';

const double ring_prebuilt_ns = measure_ring_prebuilt_std_function(
warmup, total, repeats);
std::cout << "case=mpsc_ring_publish_prebuilt_std_function ns_per_call="
<< ring_prebuilt_ns << '\n';

const double ring_uint64_ns = measure_ring<std::uint64_t>(
[](auto& ring, std::size_t count) {
for (std::size_t i = 0; i < count; ++i) {
while (!ring.try_push(static_cast<std::uint64_t>(i))) {
std::this_thread::yield();
}
}
},
[](std::uint64_t& value) { g_observer += value; },
warmup, total, repeats);
std::cout << "case=mpsc_ring_publish_uint64_control ns_per_call="
<< ring_uint64_ns << '\n';
#else
std::cout << "case=mpsc_ring_construct_and_publish_std_function ns_per_call=unavailable\n";
std::cout << "case=mpsc_ring_publish_prebuilt_std_function ns_per_call=unavailable\n";
std::cout << "case=mpsc_ring_publish_uint64_control ns_per_call=unavailable\n";
#endif

std::cout << "observer=" << g_observer << '\n';
return 0;
}
25 changes: 25 additions & 0 deletions docs/benchmarks.md
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,31 @@ $env:LOGIT_PRODUCER_PROFILE_REPEATS = "5"
./build/logit_producer_profile
```

`logit_task_ownership_profile` is a narrower follow-up experiment for the
async ownership path. It compares local `std::function` construction/copy,
batched versus one-at-a-time `TaskExecutor` admission, and (when
`LOGIT_USE_MPSC_RING` is enabled) direct `MpscRingAny` publication for a
constructed `std::function<void()>`, a prebuilt `std::function<void()>`, and a
`uint64_t` control payload. The constructed case includes callable construction,
ownership move, publication, and any retry/yield. The prebuilt case prepares all
callables before the timed producer region, so its timed work is only ownership
move/publication plus any retry/yield. Both direct-ring cases run with an active
concurrent consumer and therefore measure producer-side publication under
cache-line/coherence interaction, not intrinsic single-operation `try_push`
latency. They omit `TaskExecutor` condition-variable notification and are not
production API measurements. The `uint64_t` row is a control payload, not a
matched ownership contract. Every queue case validates that all submitted tasks
were consumed; the output records the queue backend.

This target is an A/B attribution tool, not a replacement for the prepared
producer benchmark. Its rows still have different ownership and scheduling
contracts and must not be added or subtracted as exact component costs.
Configure it with `LOGIT_TASK_OWNERSHIP_TOTAL`,
`LOGIT_TASK_OWNERSHIP_WARMUP`, and `LOGIT_TASK_OWNERSHIP_REPEATS`. The
round-trip case defaults to at most 1,000 samples because each sample waits
for worker completion; use `LOGIT_TASK_OWNERSHIP_ROUNDTRIP_TOTAL` to override
that independent sample count.

The flush regression target uses an intentionally delayed asynchronous sink and
asserts that `flush()` does not return before every queued message has reached
that sink. The fixture records latency completion and the flush barrier as
Expand Down
Loading