diff --git a/udf-runner-cpp/v2/BUILD.bazel b/udf-runner-cpp/v2/BUILD.bazel index cd8d76b..b58b2d2 100644 --- a/udf-runner-cpp/v2/BUILD.bazel +++ b/udf-runner-cpp/v2/BUILD.bazel @@ -212,3 +212,68 @@ cc_test( target_compatible_with = ["@platforms//os:linux"], deps = [":json_schema"], ) + +cc_library( + name = "moodycamel_queues", + hdrs = [ + "include/exasol/udf/v2/mpmc_queue.hpp", + "include/exasol/udf/v2/spsc_queue.hpp", + ], + includes = ["include"], + deps = [ + "@v2_concurrentqueue//:concurrentqueue", + "@v2_readerwriterqueue//:readerwriterqueue", + ], +) + +cc_test( + name = "moodycamel_queues_test", + srcs = ["moodycamel_queues_test.cc"], + copts = ["-std=c++20"], + deps = [":moodycamel_queues"], +) + +cc_binary( + name = "moodycamel_queues_shared", + srcs = ["moodycamel_queues_shared.cc"], + copts = ["-std=c++20"], + linkshared = 1, + deps = [":moodycamel_queues"], +) + +cc_test( + name = "moodycamel_symbol_leak_test", + srcs = ["moodycamel_symbol_leak_test.cc"], + copts = ["-std=c++20"], + data = [":moodycamel_queues_shared"], + args = ["$(location :moodycamel_queues_shared)"], + target_compatible_with = ["@platforms//os:linux"], + deps = [":moodycamel_queues"], +) + +cc_library( + name = "waitable_queue", + hdrs = ["include/exasol/udf/v2/waitable_queue.hpp"], + includes = ["include"], + deps = [":moodycamel_queues"], + target_compatible_with = ["@platforms//os:linux"], +) + +cc_test( + name = "waitable_queue_test", + srcs = ["waitable_queue_test.cc"], + copts = ["-std=c++20"], + deps = [":waitable_queue"], + target_compatible_with = ["@platforms//os:linux"], +) + +cc_binary( + name = "waitable_queue_benchmark", + srcs = ["waitable_queue_benchmark.cc"], + copts = ["-std=c++20"], + deps = [ + ":waitable_queue", + "@google_benchmark//:benchmark_main", + ], + target_compatible_with = ["@platforms//os:linux"], +) diff --git a/udf-runner-cpp/v2/MODULE.bazel b/udf-runner-cpp/v2/MODULE.bazel index 1372496..5371491 100644 --- a/udf-runner-cpp/v2/MODULE.bazel +++ b/udf-runner-cpp/v2/MODULE.bazel @@ -6,6 +6,27 @@ module( bazel_dep(name = "rules_cc", version = "0.2.17") bazel_dep(name = "platforms", version = "1.0.0") bazel_dep(name = "flatbuffers", version = "25.2.10") +bazel_dep(name = "google_benchmark", version = "1.9.5") + +# FlatBuffers currently selects versions of these transitive build tools that +# still use the removed incompatible_use_toolchain_transition rule attribute. +# Override them so the v2 module can be analyzed by Bazel 9. +single_version_override( + module_name = "aspect_bazel_lib", + version = "2.19.4", +) +single_version_override( + module_name = "rules_foreign_cc", + version = "0.15.1", +) +single_version_override( + module_name = "rules_go", + version = "0.63.0", +) +single_version_override( + module_name = "aspect_rules_esbuild", + version = "0.24.0", +) http_archive = use_repo_rule( "@bazel_tools//tools/build_defs/repo:http.bzl", @@ -80,3 +101,39 @@ cc_library( ) """, ) + +http_archive( + name = "v2_readerwriterqueue", + urls = ["https://github.com/cameron314/readerwriterqueue/archive/refs/tags/v1.0.7.tar.gz"], + sha256 = "532224ed052bcd5f4c6be0ed9bb2b8c88dfe7e26e3eb4dd9335303b059df6691", + strip_prefix = "readerwriterqueue-1.0.7", + build_file_content = """ +load("@rules_cc//cc:cc_library.bzl", "cc_library") + +package(default_visibility = ["//visibility:public"]) + +cc_library( + name = "readerwriterqueue", + hdrs = glob(["*.h"]), + includes = ["."], +) +""", +) + +http_archive( + name = "v2_concurrentqueue", + urls = ["https://github.com/cameron314/concurrentqueue/archive/refs/tags/v1.0.5.tar.gz"], + sha256 = "4d6368a27492d86011fde5ca0cf386dce7c49cd425aa3d9b063ca6ec373a6ef3", + strip_prefix = "concurrentqueue-1.0.5", + build_file_content = """ +load("@rules_cc//cc:cc_library.bzl", "cc_library") + +package(default_visibility = ["//visibility:public"]) + +cc_library( + name = "concurrentqueue", + hdrs = glob(["*.h"]), + includes = ["."], +) +""", +) diff --git a/udf-runner-cpp/v2/include/exasol/udf/v2/mpmc_queue.hpp b/udf-runner-cpp/v2/include/exasol/udf/v2/mpmc_queue.hpp new file mode 100644 index 0000000..da0006a --- /dev/null +++ b/udf-runner-cpp/v2/include/exasol/udf/v2/mpmc_queue.hpp @@ -0,0 +1,19 @@ +#pragma once + +// Keep all upstream moodycamel symbols below the project-owned namespace. Do +// not include these upstream headers directly in project or consumer code. +#define moodycamel exasol::udf::v2::third_party::moodycamel +#include +#include +#undef moodycamel + +namespace exasol::udf::v2 { + +template +using MpmcQueue = third_party::moodycamel::ConcurrentQueue; + +template +using BlockingMpmcQueue = + third_party::moodycamel::BlockingConcurrentQueue; + +} // namespace exasol::udf::v2 diff --git a/udf-runner-cpp/v2/include/exasol/udf/v2/spsc_queue.hpp b/udf-runner-cpp/v2/include/exasol/udf/v2/spsc_queue.hpp new file mode 100644 index 0000000..1a2cb51 --- /dev/null +++ b/udf-runner-cpp/v2/include/exasol/udf/v2/spsc_queue.hpp @@ -0,0 +1,20 @@ +#pragma once + +// Moodycamel is header-only. Rewrite its namespace while parsing the upstream +// headers so the resulting symbols cannot collide with an unrelated copy that +// may be loaded into the same process. +#define moodycamel exasol::udf::v2::third_party::moodycamel +#include +#include +#undef moodycamel + +namespace exasol::udf::v2 { + +template +using SpscQueue = third_party::moodycamel::ReaderWriterQueue; + +template +using SpscCircularBuffer = + third_party::moodycamel::BlockingReaderWriterCircularBuffer; + +} // namespace exasol::udf::v2 diff --git a/udf-runner-cpp/v2/include/exasol/udf/v2/waitable_queue.hpp b/udf-runner-cpp/v2/include/exasol/udf/v2/waitable_queue.hpp new file mode 100644 index 0000000..0ca7c50 --- /dev/null +++ b/udf-runner-cpp/v2/include/exasol/udf/v2/waitable_queue.hpp @@ -0,0 +1,170 @@ +#pragma once + +#if !defined(__linux__) +#error "exasol::udf::v2::WaitableQueue requires Linux eventfd" +#endif + +#include +#include + +#include +#include +#include +#include +#include + +#include +#include + +namespace exasol::udf::v2 { + +// Adds an epoll-compatible readiness descriptor to a queue. The descriptor +// signals that one or more queue elements may be available; it is not a +// one-to-one mapping between eventfd counter values and queue elements. +template +class WaitableQueue { + public: + using queue_type = Queue; + + WaitableQueue() + : notification_fd_(::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC)) { + if (notification_fd_ == -1) { + throw std::system_error(errno, std::generic_category(), + "eventfd"); + } + } + + explicit WaitableQueue(Queue queue) + : queue_(std::move(queue)), + notification_fd_(::eventfd(0, EFD_NONBLOCK | EFD_CLOEXEC)) { + if (notification_fd_ == -1) { + throw std::system_error(errno, std::generic_category(), + "eventfd"); + } + } + + ~WaitableQueue() { + if (notification_fd_ != -1) { + ::close(notification_fd_); + } + } + + WaitableQueue(const WaitableQueue&) = delete; + WaitableQueue& operator=(const WaitableQueue&) = delete; + + WaitableQueue(WaitableQueue&& other) noexcept + : queue_(std::move(other.queue_)), + notification_fd_(std::exchange(other.notification_fd_, -1)) {} + + WaitableQueue& operator=(WaitableQueue&& other) noexcept { + if (this != &other) { + if (notification_fd_ != -1) { + ::close(notification_fd_); + } + queue_ = std::move(other.queue_); + notification_fd_ = std::exchange(other.notification_fd_, -1); + } + return *this; + } + + [[nodiscard]] int native_handle() const noexcept { + return notification_fd_; + } + + template + [[nodiscard]] bool enqueue(T&& value) { + if (!queue_.enqueue(std::forward(value))) { + return false; + } + notify(); + return true; + } + + template + std::size_t enqueue_batch(InputIt first, InputIt last) { + std::size_t enqueued = 0; + for (; first != last; ++first) { + if (!queue_.enqueue(*first)) { + break; + } + ++enqueued; + } + if (enqueued != 0) { + notify(); + } + return enqueued; + } + + template + [[nodiscard]] bool try_dequeue(Output& value) { + return queue_.try_dequeue(value); + } + + // Drains all eventfd notifications and returns their accumulated count. + // Callers should then dequeue until the queue is empty and recheck it + // before going back to epoll_wait(). + std::uint64_t drain_notifications() { + std::uint64_t total = 0; + for (;;) { + std::uint64_t value = 0; + const ssize_t result = ::read(notification_fd_, &value, + sizeof(value)); + if (result == sizeof(value)) { + total += value; + continue; + } + if (result == -1 && errno == EINTR) { + continue; + } + if (result == -1 && errno == EAGAIN) { + return total; + } + if (result == -1) { + throw std::system_error(errno, std::generic_category(), + "read eventfd"); + } + throw std::system_error(EIO, std::generic_category(), + "short read from eventfd"); + } + } + + Queue& queue() noexcept { return queue_; } + const Queue& queue() const noexcept { return queue_; } + + private: + void notify() { + constexpr std::uint64_t signal = 1; + for (;;) { + const ssize_t result = + ::write(notification_fd_, &signal, sizeof(signal)); + if (result == sizeof(signal)) { + return; + } + if (result == -1 && errno == EINTR) { + continue; + } + // A saturated eventfd is already readable. The queue item remains + // available, so no additional notification is needed. + if (result == -1 && errno == EAGAIN) { + return; + } + if (result == -1) { + throw std::system_error(errno, std::generic_category(), + "write eventfd"); + } + throw std::system_error(EIO, std::generic_category(), + "short write to eventfd"); + } + } + + Queue queue_; + int notification_fd_; +}; + +template +using WaitableSpscQueue = WaitableQueue>; + +template +using WaitableMpmcQueue = WaitableQueue>; + +} // namespace exasol::udf::v2 diff --git a/udf-runner-cpp/v2/moodycamel_queues_shared.cc b/udf-runner-cpp/v2/moodycamel_queues_shared.cc new file mode 100644 index 0000000..bb89512 --- /dev/null +++ b/udf-runner-cpp/v2/moodycamel_queues_shared.cc @@ -0,0 +1,10 @@ +#include +#include + +extern "C" void exasol_udf_v2_moodycamel_queue_anchor() { + exasol::udf::v2::SpscQueue spsc; + spsc.enqueue(1); + + exasol::udf::v2::MpmcQueue mpmc; + mpmc.enqueue(2); +} diff --git a/udf-runner-cpp/v2/moodycamel_queues_test.cc b/udf-runner-cpp/v2/moodycamel_queues_test.cc new file mode 100644 index 0000000..15b6f07 --- /dev/null +++ b/udf-runner-cpp/v2/moodycamel_queues_test.cc @@ -0,0 +1,27 @@ +#include + +#include +#include + +int main() { + exasol::udf::v2::SpscQueue spsc; + assert(spsc.enqueue(7)); + int value = 0; + assert(spsc.try_dequeue(value)); + assert(value == 7); + + exasol::udf::v2::SpscCircularBuffer circular(2); + assert(circular.try_enqueue(8)); + assert(circular.try_dequeue(value)); + assert(value == 8); + + exasol::udf::v2::MpmcQueue mpmc; + assert(mpmc.enqueue(9)); + assert(mpmc.try_dequeue(value)); + assert(value == 9); + + exasol::udf::v2::BlockingMpmcQueue blocking; + assert(blocking.enqueue(10)); + assert(blocking.try_dequeue(value)); + assert(value == 10); +} diff --git a/udf-runner-cpp/v2/moodycamel_symbol_leak_test.cc b/udf-runner-cpp/v2/moodycamel_symbol_leak_test.cc new file mode 100644 index 0000000..cfebd9a --- /dev/null +++ b/udf-runner-cpp/v2/moodycamel_symbol_leak_test.cc @@ -0,0 +1,108 @@ +#include + +#include +#include +#include +#include +#include +#include +#include +#include + +#include + +namespace { + +constexpr std::string_view kGlobalNamespacePrefix = "_ZN10moodycamel"; +constexpr std::string_view kIsolatedNamespacePrefix = + "_ZN6exasol3udf2v211third_party10moodycamel"; + +[[noreturn]] void fail(const std::string& message) { + throw std::runtime_error(message); +} + +template +T read_object(const std::vector& file, std::size_t offset) { + if (offset > file.size() || file.size() - offset < sizeof(T)) { + fail("ELF file is truncated"); + } + T result; + std::memcpy(&result, file.data() + offset, sizeof(result)); + return result; +} + +std::string read_string(const std::vector& file, std::size_t offset, + std::size_t maximum_size) { + if (offset > file.size() || file.size() - offset < maximum_size) { + fail("ELF string table is truncated"); + } + const char* begin = file.data() + offset; + const void* end = std::memchr(begin, '\0', maximum_size); + if (end == nullptr) { + fail("ELF symbol name is not terminated"); + } + return std::string(begin, static_cast(end)); +} + +void verify_symbols(const std::string& path) { + std::ifstream input(path, std::ios::binary); + if (!input) { + fail("cannot read queue library: " + path); + } + const std::vector file{std::istreambuf_iterator(input), + std::istreambuf_iterator()}; + const Elf64_Ehdr header = read_object(file, 0); + if (std::memcmp(header.e_ident, ELFMAG, SELFMAG) != 0 || + header.e_ident[EI_CLASS] != ELFCLASS64 || + header.e_ident[EI_DATA] != ELFDATA2LSB || + header.e_shentsize != sizeof(Elf64_Shdr)) { + fail("queue library is not a little-endian ELF64 file"); + } + + Elf64_Shdr dynamic_symbols{}; + bool found_dynamic_symbols = false; + for (std::size_t i = 0; i < header.e_shnum; ++i) { + const Elf64_Shdr section = read_object( + file, header.e_shoff + i * sizeof(Elf64_Shdr)); + if (section.sh_type == SHT_DYNSYM) { + dynamic_symbols = section; + found_dynamic_symbols = true; + break; + } + } + if (!found_dynamic_symbols || dynamic_symbols.sh_link >= header.e_shnum || + dynamic_symbols.sh_entsize != sizeof(Elf64_Sym)) { + fail("ELF dynamic symbol table is invalid"); + } + + const Elf64_Shdr string_table = read_object( + file, header.e_shoff + dynamic_symbols.sh_link * sizeof(Elf64_Shdr)); + bool found_isolated_symbol = false; + for (std::size_t offset = 0; offset < dynamic_symbols.sh_size; + offset += sizeof(Elf64_Sym)) { + const Elf64_Sym symbol = read_object( + file, dynamic_symbols.sh_offset + offset); + if (symbol.st_name >= string_table.sh_size) { + continue; + } + const std::string name = read_string( + file, string_table.sh_offset + symbol.st_name, + string_table.sh_size - symbol.st_name); + if (name.starts_with(kGlobalNamespacePrefix)) { + fail("queue library exports a global moodycamel symbol: " + name); + } + if (name.starts_with(kIsolatedNamespacePrefix)) { + found_isolated_symbol = true; + } + } + if (!found_isolated_symbol) { + fail("queue library does not export an isolated moodycamel symbol"); + } +} + +} // namespace + +int main(int argc, char** argv) { + assert(argc == 2); + verify_symbols(argv[1]); +} diff --git a/udf-runner-cpp/v2/waitable_queue_benchmark.cc b/udf-runner-cpp/v2/waitable_queue_benchmark.cc new file mode 100644 index 0000000..5d238cd --- /dev/null +++ b/udf-runner-cpp/v2/waitable_queue_benchmark.cc @@ -0,0 +1,252 @@ +#include + +#include +#include + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include + +namespace { + +void benchmark_check(bool condition, const char* message) { + if (!condition) { + std::fprintf(stderr, "waitable queue benchmark failure: %s\n", message); + std::abort(); + } +} + +// These benchmarks compare raw Moodycamel SPSC behavior with the blocking +// circular-buffer baseline and the eventfd-backed waitable queue. Results are +// informational: build mode, CPU frequency, scheduler activity, and system +// load can materially affect them. +struct TimedItem { + std::uint64_t sequence; + std::chrono::steady_clock::time_point sent; +}; + +// Includes raw enqueue and dequeue only; this is the queue-operation baseline +// for the waitable round-trip benchmark. +void BM_RawSpscRoundTrip(benchmark::State& state) { + exasol::udf::v2::SpscQueue queue(1024); + for (auto _ : state) { + int value = 0; + benchmark::DoNotOptimize(queue.enqueue(1)); + benchmark::DoNotOptimize(queue.try_dequeue(value)); + benchmark::DoNotOptimize(value); + } + state.SetItemsProcessed(state.iterations()); +} + +// Includes enqueue, eventfd notification draining, and dequeue. The eventfd +// write is part of the measured round trip. +void BM_WaitableSpscRoundTrip(benchmark::State& state) { + exasol::udf::v2::WaitableSpscQueue queue( + exasol::udf::v2::SpscQueue(1024)); + for (auto _ : state) { + int value = 0; + benchmark::DoNotOptimize(queue.enqueue(1)); + benchmark::DoNotOptimize(queue.drain_notifications()); + benchmark::DoNotOptimize(queue.try_dequeue(value)); + benchmark::DoNotOptimize(value); + } + state.SetItemsProcessed(state.iterations()); +} + +// Includes wait_enqueue and dequeue on a non-full blocking SPSC queue. The +// benchmark measures the uncontended fast path rather than intentional waits. +void BM_BlockingSpscRoundTrip(benchmark::State& state) { + exasol::udf::v2::SpscCircularBuffer queue(1024); + for (auto _ : state) { + int value = 0; + queue.wait_enqueue(1); + benchmark::DoNotOptimize(queue.try_dequeue(value)); + benchmark::DoNotOptimize(value); + } + state.SetItemsProcessed(state.iterations()); +} + +// Enqueue-only benchmarks pause timing while removing the item so the queue +// remains empty for the next iteration. The waitable case includes its +// eventfd write; notification draining is cleanup and is not timed. +void BM_RawSpscEnqueueLatency(benchmark::State& state) { + exasol::udf::v2::SpscQueue queue(1024); + for (auto _ : state) { + benchmark::DoNotOptimize(queue.enqueue(1)); + state.PauseTiming(); + int value = 0; + benchmark::DoNotOptimize(queue.try_dequeue(value)); + benchmark::DoNotOptimize(value); + state.ResumeTiming(); + } + state.SetItemsProcessed(state.iterations()); +} + +void BM_WaitableSpscEnqueueLatency(benchmark::State& state) { + exasol::udf::v2::WaitableSpscQueue queue( + exasol::udf::v2::SpscQueue(1024)); + for (auto _ : state) { + benchmark::DoNotOptimize(queue.enqueue(1)); + state.PauseTiming(); + int value = 0; + benchmark::DoNotOptimize(queue.try_dequeue(value)); + benchmark::DoNotOptimize(queue.drain_notifications()); + benchmark::DoNotOptimize(value); + state.ResumeTiming(); + } + state.SetItemsProcessed(state.iterations()); +} + +void BM_BlockingSpscEnqueueLatency(benchmark::State& state) { + exasol::udf::v2::SpscCircularBuffer queue(1024); + for (auto _ : state) { + queue.wait_enqueue(1); + state.PauseTiming(); + int value = 0; + benchmark::DoNotOptimize(queue.try_dequeue(value)); + benchmark::DoNotOptimize(value); + state.ResumeTiming(); + } + state.SetItemsProcessed(state.iterations()); +} + +// Raw and waitable batch benchmarks use the same batch sizes. The waitable +// queue emits one eventfd notification after the entire batch, exposing how +// batching amortizes notification overhead. +void BM_RawSpscBatch(benchmark::State& state) { + const auto batch_size = static_cast(state.range(0)); + const std::vector batch(batch_size, 1); + exasol::udf::v2::SpscQueue queue(batch_size); + + for (auto _ : state) { + for (int value : batch) { + benchmark::DoNotOptimize(queue.enqueue(value)); + } + int value = 0; + for (std::size_t i = 0; i < batch_size; ++i) { + benchmark::DoNotOptimize(queue.try_dequeue(value)); + } + benchmark::DoNotOptimize(value); + } + state.SetItemsProcessed(state.iterations() * batch_size); +} + +void BM_WaitableSpscBatch(benchmark::State& state) { + const auto batch_size = static_cast(state.range(0)); + const std::vector batch(batch_size, 1); + exasol::udf::v2::WaitableSpscQueue queue{ + exasol::udf::v2::SpscQueue(batch_size)}; + + for (auto _ : state) { + benchmark::DoNotOptimize(queue.enqueue_batch(batch.begin(), batch.end())); + benchmark::DoNotOptimize(queue.drain_notifications()); + int value = 0; + for (std::size_t i = 0; i < batch_size; ++i) { + benchmark::DoNotOptimize(queue.try_dequeue(value)); + } + benchmark::DoNotOptimize(value); + } + state.SetItemsProcessed(state.iterations() * batch_size); +} + +// Measures producer timestamp through enqueue, eventfd readiness, epoll_wait, +// and dequeue. One item is outstanding at a time, so the result measures +// wakeup latency rather than latency caused by queue backlog. The producer +// handshake is outside the manually recorded interval. +void BM_WaitableSpscEpollLatency(benchmark::State& state) { + exasol::udf::v2::WaitableSpscQueue queue{ + exasol::udf::v2::SpscQueue(8)}; + const int epoll_fd = ::epoll_create1(EPOLL_CLOEXEC); + benchmark_check(epoll_fd != -1, "epoll_create1 failed"); + + epoll_event queue_event{}; + queue_event.events = EPOLLIN; + queue_event.data.fd = queue.native_handle(); + benchmark_check(::epoll_ctl(epoll_fd, EPOLL_CTL_ADD, queue.native_handle(), + &queue_event) == 0, + "epoll_ctl failed"); + + std::atomic requested{0}; + std::atomic completed{0}; + std::atomic stop{false}; + std::thread producer([&] { + std::uint64_t sequence = 0; + while (!stop.load(std::memory_order_acquire)) { + while (requested.load(std::memory_order_acquire) <= sequence && + !stop.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + if (stop.load(std::memory_order_acquire)) { + break; + } + + TimedItem item{sequence++, std::chrono::steady_clock::now()}; + benchmark_check(queue.enqueue(std::move(item)), "queue enqueue failed"); + + while (completed.load(std::memory_order_acquire) < sequence && + !stop.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + } + }); + + std::uint64_t expected_sequence = 0; + for (auto _ : state) { + requested.fetch_add(1, std::memory_order_release); + + epoll_event event{}; + int event_count = 0; + do { + event_count = ::epoll_wait(epoll_fd, &event, 1, -1); + } while (event_count == -1 && errno == EINTR); + benchmark_check(event_count == 1, "epoll_wait failed"); + benchmark_check(event.data.fd == queue.native_handle(), + "unexpected epoll event"); + + benchmark::DoNotOptimize(queue.drain_notifications()); + TimedItem item{}; + benchmark_check(queue.try_dequeue(item), "queue dequeue failed"); + benchmark_check(item.sequence == expected_sequence, + "unexpected item sequence"); + ++expected_sequence; + const auto elapsed = std::chrono::steady_clock::now() - item.sent; + benchmark::DoNotOptimize(item); + state.SetIterationTime( + std::chrono::duration(elapsed).count()); + completed.store(expected_sequence, std::memory_order_release); + } + + stop.store(true, std::memory_order_release); + requested.fetch_add(1, std::memory_order_release); + producer.join(); + ::close(epoll_fd); +} + +} // namespace + +BENCHMARK(BM_RawSpscRoundTrip); +BENCHMARK(BM_WaitableSpscRoundTrip); +BENCHMARK(BM_BlockingSpscRoundTrip); +BENCHMARK(BM_RawSpscEnqueueLatency); +BENCHMARK(BM_WaitableSpscEnqueueLatency); +BENCHMARK(BM_BlockingSpscEnqueueLatency); +BENCHMARK(BM_RawSpscBatch) + ->Args({1}) + ->Args({8}) + ->Args({64}) + ->Args({256}); +BENCHMARK(BM_WaitableSpscBatch) + ->Args({1}) + ->Args({8}) + ->Args({64}) + ->Args({256}); +BENCHMARK(BM_WaitableSpscEpollLatency)->UseManualTime(); diff --git a/udf-runner-cpp/v2/waitable_queue_test.cc b/udf-runner-cpp/v2/waitable_queue_test.cc new file mode 100644 index 0000000..f3b35f6 --- /dev/null +++ b/udf-runner-cpp/v2/waitable_queue_test.cc @@ -0,0 +1,96 @@ +#include +#include +#include + +#include +#include +#include +#include +#include +#include +#include + +#include + +namespace { + +void test_check(bool condition, const char* message) { + if (!condition) { + std::fprintf(stderr, "waitable queue test failure: %s\n", message); + std::abort(); + } +} + +void add_to_epoll(int epoll_fd, int fd, std::uint32_t events) { + epoll_event event{}; + event.events = events; + event.data.fd = fd; + test_check(::epoll_ctl(epoll_fd, EPOLL_CTL_ADD, fd, &event) == 0, + "epoll_ctl failed"); +} + +void close_pair(const std::array& sockets) { + ::close(sockets[0]); + ::close(sockets[1]); +} + +} // namespace + +int main() { + exasol::udf::v2::WaitableSpscQueue queue; + const int epoll_fd = ::epoll_create1(EPOLL_CLOEXEC); + test_check(epoll_fd != -1, "epoll_create1 failed"); + + std::array sockets{}; + test_check(::socketpair(AF_UNIX, SOCK_STREAM | SOCK_CLOEXEC, 0, + sockets.data()) == 0, + "socketpair failed"); + add_to_epoll(epoll_fd, queue.native_handle(), EPOLLIN); + add_to_epoll(epoll_fd, sockets[1], EPOLLIN); + + test_check(queue.enqueue(42), "queue enqueue failed"); + const char byte = 'x'; + test_check(::write(sockets[0], &byte, sizeof(byte)) == sizeof(byte), + "socket write failed"); + + std::array events{}; + const int event_count = ::epoll_wait(epoll_fd, events.data(), + events.size(), 1000); + test_check(event_count == 2, "epoll_wait did not report both descriptors"); + + bool queue_ready = false; + bool socket_ready = false; + for (int i = 0; i < event_count; ++i) { + queue_ready |= events[i].data.fd == queue.native_handle(); + socket_ready |= events[i].data.fd == sockets[1]; + } + test_check(queue_ready, "queue descriptor was not ready"); + test_check(socket_ready, "socket descriptor was not ready"); + + test_check(queue.drain_notifications() == 1, + "unexpected queue notification count"); + int value = 0; + test_check(queue.try_dequeue(value), "queue dequeue failed"); + test_check(value == 42, "unexpected dequeued value"); + + const std::vector batch{1, 2, 3}; + test_check(queue.enqueue_batch(batch.begin(), batch.end()) == batch.size(), + "batch enqueue failed"); + test_check(queue.drain_notifications() == 1, + "unexpected batch notification count"); + for (int expected : batch) { + test_check(queue.try_dequeue(value), "batch dequeue failed"); + test_check(value == expected, "unexpected batch value"); + } + test_check(!queue.try_dequeue(value), "queue should be empty"); + + exasol::udf::v2::WaitableMpmcQueue mpmc; + test_check(mpmc.enqueue(7), "MPMC queue enqueue failed"); + test_check(mpmc.drain_notifications() == 1, + "unexpected MPMC notification count"); + test_check(mpmc.try_dequeue(value), "MPMC queue dequeue failed"); + test_check(value == 7, "unexpected MPMC value"); + + close_pair(sockets); + ::close(epoll_fd); +}