From 97e3a3e2552a796934430d9a7679782c20a1278a Mon Sep 17 00:00:00 2001 From: Jacob Nilsson Date: Mon, 17 Aug 2026 15:09:54 +0200 Subject: [PATCH] add kqueue io_backend for macOS and BSD systems --- CMakeLists.txt | 6 + include/netkit/definitions.hpp | 6 + include/netkit/io/bsd/io_backend.hpp | 53 ++++ include/netkit/io/fallback/io_backend.hpp | 4 + include/netkit/io/io_backend.hpp | 3 +- include/netkit/netkit.hpp | 1 + src/io/bsd/io_backend.cpp | 321 ++++++++++++++++++++++ src/io/fallback/io_backend.cpp | 4 + 8 files changed, 397 insertions(+), 1 deletion(-) create mode 100644 include/netkit/io/bsd/io_backend.hpp create mode 100644 src/io/bsd/io_backend.cpp diff --git a/CMakeLists.txt b/CMakeLists.txt index 1a2e0db..cb5adba 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -8,6 +8,7 @@ option(NETKIT_ENABLE_FALLBACK_CA "Enable fallback CA certficate (required for De option(NETKIT_ENABLE_SOCK_CUSTOM_RESOLVER "Enable using netkit's DNS resolver instead of the system resolver" OFF) option(NETKIT_ENABLE_EPOLL "Enable using epoll backend (Linux only)" ON) option(NETKIT_ENABLE_WSAPOLL "Enable using wsapoll backend (Windows only)" ON) +option(NETKIT_ENABLE_KQUEUE "Enable using kqueue backend (macOS, BSD only)" ON) option(NETKIT_ENABLE_HTTP "Enable netkit's HTTP abstractions" ON) option(NETKIT_ENABLE_DNS "Enable netkit's DNS features" ON) option(NETKIT_WOLFSSL_DEBUG "Enable WolfSSL debugging" OFF) @@ -75,6 +76,7 @@ add_library(netkit ${NETKIT_LIB_TYPE} src/body/multipart_part_body.cpp src/io/io_awaitable.cpp src/io/linux/io_backend.cpp + src/io/bsd/io_backend.cpp src/socket/native/native_async_socket.cpp src/socket/native/native_sync_listener.cpp src/socket/native/native_async_listener.cpp @@ -215,6 +217,10 @@ if (NETKIT_ENABLE_WSAPOLL) target_compile_definitions(netkit PUBLIC NETKIT_WSAPOLL) endif() +if (NETKIT_ENABLE_KQUEUE) + target_compile_definitions(netkit PUBLIC NETKIT_KQUEUE) +endif() + if (WIN32) add_definitions(-DNOMINMAX) endif() diff --git a/include/netkit/definitions.hpp b/include/netkit/definitions.hpp index 7238ea2..6cb9013 100644 --- a/include/netkit/definitions.hpp +++ b/include/netkit/definitions.hpp @@ -32,6 +32,12 @@ #define NETKIT_LINUX 1 #endif +#if !defined(NETKIT_LINUX) && \ + defined(NETKIT_UNIX) && !defined(NETKIT_DKP) && \ + !defined(NETKIT_MACOS) +#define NETKIT_BSD +#endif + #ifndef NETKIT_FALLBACK_IPV4_DNS_1 #define NETKIT_FALLBACK_IPV4_DNS_1 "8.8.8.8" #endif diff --git a/include/netkit/io/bsd/io_backend.hpp b/include/netkit/io/bsd/io_backend.hpp new file mode 100644 index 0000000..d78fae3 --- /dev/null +++ b/include/netkit/io/bsd/io_backend.hpp @@ -0,0 +1,53 @@ +#pragma once + +#include +#include + +#if (defined(NETKIT_MACOS) || defined(NETKIT_BSD)) \ + && defined(NETKIT_KQUEUE) + +#include + +#include +#include +#include + +namespace netkit::io { + +class io_backend : public basic_io_backend { +public: + io_backend(); + ~io_backend() override; + + void wake() override; + + void update_state( + io_handle_t fd, + const io_handle_state& state + ) override; + + void register_waiter( + io_handle_t fd, + io_event ev, + std::coroutine_handle<> h + ) override; + + void run() override; + void stop() override; + + void poll(int timeout_ms) override; + void poll() override; + +private: + static constexpr uintptr_t wake_ident = 1; + + io_handle_t kqueue_fd_ = -1; + bool running_ = true; + + std::unordered_map fd_map_; + std::unordered_set registered_fds_; +}; + +} // namespace netkit::io + +#endif diff --git a/include/netkit/io/fallback/io_backend.hpp b/include/netkit/io/fallback/io_backend.hpp index 8e0bb07..6f429a2 100644 --- a/include/netkit/io/fallback/io_backend.hpp +++ b/include/netkit/io/fallback/io_backend.hpp @@ -5,6 +5,8 @@ #if !defined(NETKIT_LINUX) || !defined(NETKIT_EPOLL) #if !defined(NETKIT_WINDOWS) || !defined(NETKIT_WSAPOLL) +#if !defined(NETKIT_BSD) || !defined(NETKIT_KQUEUE) +#if !defined(NETKIT_MACOS) || !defined(NETKIT_KQUEUE) #include #include @@ -54,5 +56,7 @@ class NETKIT_API io_backend : public basic_io_backend { } // namespace netkit::io +#endif +#endif #endif #endif \ No newline at end of file diff --git a/include/netkit/io/io_backend.hpp b/include/netkit/io/io_backend.hpp index e91ec49..26f229b 100644 --- a/include/netkit/io/io_backend.hpp +++ b/include/netkit/io/io_backend.hpp @@ -2,4 +2,5 @@ #include #include -#include \ No newline at end of file +#include +#include \ No newline at end of file diff --git a/include/netkit/netkit.hpp b/include/netkit/netkit.hpp index 3d5f13f..a3322ce 100644 --- a/include/netkit/netkit.hpp +++ b/include/netkit/netkit.hpp @@ -85,6 +85,7 @@ #include #include #include +#include #include #include #include diff --git a/src/io/bsd/io_backend.cpp b/src/io/bsd/io_backend.cpp new file mode 100644 index 0000000..47a166e --- /dev/null +++ b/src/io/bsd/io_backend.cpp @@ -0,0 +1,321 @@ +#include + +#if (defined(NETKIT_MACOS) || defined(NETKIT_BSD)) \ + && defined(NETKIT_KQUEUE) + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include + +netkit::io::io_backend::io_backend() { + kqueue_fd_ = kqueue(); + + if (kqueue_fd_ == -1) + throw std::runtime_error("kqueue failed"); + + struct kevent ev{}; + + EV_SET( + &ev, + wake_ident, + EVFILT_USER, + EV_ADD | EV_CLEAR, + 0, + 0, + nullptr + ); + + if (kevent(kqueue_fd_, &ev, 1, nullptr, 0, nullptr) == -1) { + close(kqueue_fd_); + kqueue_fd_ = -1; + + throw std::runtime_error("failed to add kqueue wake event"); + } +} + +netkit::io::io_backend::~io_backend() { + if (kqueue_fd_ != -1) + close(kqueue_fd_); +} + +void netkit::io::io_backend::wake() { + struct kevent ev{}; + + EV_SET( + &ev, + wake_ident, + EVFILT_USER, + 0, + NOTE_TRIGGER, + 0, + nullptr + ); + + kevent(kqueue_fd_, &ev, 1, nullptr, 0, nullptr); +} + +void netkit::io::io_backend::update_state( + io_handle_t fd, + const io_handle_state& state +) { + const bool wants_read = !state.read_waiters.empty(); + const bool wants_write = !state.write_waiters.empty(); + + if (wants_read) { + struct kevent ev{}; + + EV_SET( + &ev, + static_cast(fd), + EVFILT_READ, + EV_ADD | EV_ENABLE, + 0, + 0, + nullptr + ); + + if (kevent(kqueue_fd_, &ev, 1, nullptr, 0, nullptr) == -1) + throw std::runtime_error( + "failed to add kqueue read filter: " + + std::string(std::strerror(errno)) + ); + } else { + struct kevent ev{}; + + EV_SET( + &ev, + static_cast(fd), + EVFILT_READ, + EV_DELETE, + 0, + 0, + nullptr + ); + + if (kevent(kqueue_fd_, &ev, 1, nullptr, 0, nullptr) == -1 && + errno != ENOENT) { + throw std::runtime_error("failed to remove kqueue read filter: " + std::string(std::strerror(errno))); + } + } + + if (wants_write) { + struct kevent ev{}; + + EV_SET( + &ev, + static_cast(fd), + EVFILT_WRITE, + EV_ADD | EV_ENABLE, + 0, + 0, + nullptr + ); + + if (kevent(kqueue_fd_, &ev, 1, nullptr, 0, nullptr) == -1) + throw std::runtime_error( + "failed to add kqueue write filter: " + + std::string(std::strerror(errno)) + ); + } else { + struct kevent ev{}; + + EV_SET( + &ev, + static_cast(fd), + EVFILT_WRITE, + EV_DELETE, + 0, + 0, + nullptr + ); + + if (kevent(kqueue_fd_, &ev, 1, nullptr, 0, nullptr) == -1 && + errno != ENOENT) { + throw std::runtime_error( + "failed to remove kqueue write filter: " + + std::string(std::strerror(errno)) + ); + } + } + + if (wants_read || wants_write) + registered_fds_.insert(fd); + else + registered_fds_.erase(fd); +} + +void netkit::io::io_backend::register_waiter( + io_handle_t fd, + io_event ev, + std::coroutine_handle<> h +) { + auto& state = fd_map_[fd]; + + if (ev == io_event::read) + state.read_waiters.push_back(h); + else + state.write_waiters.push_back(h); + + update_state(fd, state); +} + +void netkit::io::io_backend::run() { + constexpr int MAX_EVENTS = 64; + + struct kevent events[MAX_EVENTS]; + + while (running_) { + int n = kevent( + kqueue_fd_, + nullptr, + 0, + events, + MAX_EVENTS, + nullptr + ); + + if (n == -1) { + if (errno == EINTR) + continue; + + throw std::runtime_error("kevent failed"); + } + + for (int i = 0; i < n; ++i) { + const auto& event = events[i]; + + if (event.filter == EVFILT_USER && + event.ident == wake_ident) { + continue; + } + + const io_handle_t fd = + static_cast(event.ident); + + auto it = fd_map_.find(fd); + + if (it == fd_map_.end()) + continue; + + auto& state = it->second; + + const bool error = + (event.flags & (EV_EOF | EV_ERROR)) != 0; + + if (event.filter == EVFILT_READ || error) { + auto waiters = + std::move(state.read_waiters); + + state.read_waiters.clear(); + + for (auto h : waiters) + h.resume(); + } + + if (event.filter == EVFILT_WRITE || error) { + auto waiters = + std::move(state.write_waiters); + + state.write_waiters.clear(); + + for (auto h : waiters) + h.resume(); + } + + update_state(fd, state); + } + } +} + +void netkit::io::io_backend::stop() { + running_ = false; + + wake(); +} + +void netkit::io::io_backend::poll(int timeout_ms) { + struct kevent events[64]; + + timespec timeout{}; + + if (timeout_ms >= 0) { + timeout.tv_sec = timeout_ms / 1000; + timeout.tv_nsec = (timeout_ms % 1000) * 1000000; + } + + int n = kevent( + kqueue_fd_, + nullptr, + 0, + events, + 64, + timeout_ms < 0 ? nullptr : &timeout + ); + + if (n == -1) { + if (errno == EINTR) + return; + + throw std::runtime_error("kevent failed"); + } + + for (int i = 0; i < n; ++i) { + const auto& event = events[i]; + + if (event.filter == EVFILT_USER && + event.ident == wake_ident) { + continue; + } + + const auto fd = + static_cast(event.ident); + + auto it = fd_map_.find(fd); + + if (it == fd_map_.end()) + continue; + + auto& state = it->second; + + const bool error = + (event.flags & (EV_EOF | EV_ERROR)) != 0; + + if (event.filter == EVFILT_READ || error) { + auto waiters = + std::move(state.read_waiters); + + state.read_waiters.clear(); + + for (auto h : waiters) + h.resume(); + } + + if (event.filter == EVFILT_WRITE || error) { + auto waiters = + std::move(state.write_waiters); + + state.write_waiters.clear(); + + for (auto h : waiters) + h.resume(); + } + + update_state(fd, state); + } +} + +void netkit::io::io_backend::poll() { + poll(-1); +} + +#endif \ No newline at end of file diff --git a/src/io/fallback/io_backend.cpp b/src/io/fallback/io_backend.cpp index b91b618..6baa662 100644 --- a/src/io/fallback/io_backend.cpp +++ b/src/io/fallback/io_backend.cpp @@ -4,6 +4,8 @@ #if !defined(NETKIT_LINUX) || !defined(NETKIT_EPOLL) #if !defined(NETKIT_WINDOWS) || !defined(NETKIT_WSAPOLL) +#if !defined(NETKIT_BSD) || !defined(NETKIT_KQUEUE) +#if !defined(NETKIT_MACOS) || !defined(NETKIT_KQUEUE) void netkit::io::io_backend::worker() { while (running_) { @@ -102,5 +104,7 @@ void netkit::io::io_backend::poll() { poll(-1); } +#endif +#endif #endif #endif \ No newline at end of file