diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index d2d5f69..1e90d45 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_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 @@ -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 diff --git a/bench/CMakeLists.txt b/bench/CMakeLists.txt index de6221a..0d285af 100644 --- a/bench/CMakeLists.txt +++ b/bench/CMakeLists.txt @@ -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}) diff --git a/bench/task_ownership_profile.cpp b/bench/task_ownership_profile.cpp new file mode 100644 index 0000000..2721976 --- /dev/null +++ b/bench/task_ownership_profile.cpp @@ -0,0 +1,314 @@ +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include + +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 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::stoull(value)); + } catch (...) { + } + } + return fallback; +} + +template +double measure(Loop loop, Cleanup cleanup, std::size_t warmup, + std::size_t total, std::size_t repeats) { + loop(warmup); + cleanup(warmup); + + std::vector 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::steady_clock::now() - start).count(); + cleanup(total); + samples.push_back(static_cast(elapsed)); + } + + std::sort(samples.begin(), samples.end()); + return static_cast(samples[samples.size() / 2]) / + static_cast(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> tasks; + tasks.reserve(count); + for (std::size_t i = 0; i < count; ++i) { + tasks.emplace_back([]() { ++g_observer; }); + } + + logit::detail::MpscRingAny> ring(kQueueCapacity); + std::atomic started{false}; + std::atomic done{false}; + std::atomic ready{0}; + std::atomic 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 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::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(elapsed); + }; + + run(warmup); + std::vector 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(samples[samples.size() / 2]) / + static_cast(total); +} + +template +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 ring(kQueueCapacity); + std::atomic started{false}; + std::atomic done{false}; + std::atomic ready{0}; + std::atomic 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::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(elapsed); + }; + + run(warmup); + std::vector 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(samples[samples.size() / 2]) / + static_cast(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( + 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 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 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 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>( + [](auto& ring, std::size_t count) { + for (std::size_t i = 0; i < count; ++i) { + std::function value = []() { ++g_observer; }; + while (!ring.try_push(std::move(value))) { + std::this_thread::yield(); + } + } + }, + [](std::function& 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( + [](auto& ring, std::size_t count) { + for (std::size_t i = 0; i < count; ++i) { + while (!ring.try_push(static_cast(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; +} diff --git a/docs/benchmarks.md b/docs/benchmarks.md index 469cf7a..3accc55 100644 --- a/docs/benchmarks.md +++ b/docs/benchmarks.md @@ -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`, a prebuilt `std::function`, 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