* @brief A minimal RPC @b server using @b `io_uring` functionality * to setup the UDP socket, and process many requests concurrently. */
| 7026 | * to setup the UDP socket, and process many requests concurrently. |
| 7027 | */ |
| 7028 | class rpc_uring55_server { |
| 7029 | |
| 7030 | int socket_descriptor_; |
| 7031 | sockaddr_in server_address_; |
| 7032 | std::atomic_bool should_stop_; |
| 7033 | io_uring ring_; |
| 7034 | |
| 7035 | // Pre-allocated resources |
| 7036 | mmap_array<message_t> messages_; |
| 7037 | std::size_t max_concurrency_; |
| 7038 | |
| 7039 | public: |
| 7040 | using status_t = message_t::message_status_t; |
| 7041 | |
| 7042 | rpc_uring55_server(std::string const &server_address_str, std::uint16_t port, std::size_t max_concurrency) |
| 7043 | : should_stop_(false), messages_(max_concurrency * 2), max_concurrency_(max_concurrency) { |
| 7044 | |
| 7045 | auto [socket_descriptor, server_address] = rpc_server_socket(port, server_address_str); |
| 7046 | socket_descriptor_ = socket_descriptor; |
| 7047 | server_address_ = server_address; |
| 7048 | |
| 7049 | // Initialize `io_uring` with one slot for each receive/send operation |
| 7050 | if (io_uring_queue_init(max_concurrency * 2, &ring_, 0) < 0) |
| 7051 | raise_system_error("Failed to initialize io_uring 5.5 server"); |
| 7052 | if (io_uring_register_files(&ring_, &socket_descriptor_, 1) < 0) |
| 7053 | raise_system_error("Failed to register file descriptor with io_uring 5.5 server"); |
| 7054 | |
| 7055 | // Initialize message resources |
| 7056 | for (message_t &message : messages_) { |
| 7057 | memset(&message.header, 0, sizeof(message.header)); |
| 7058 | message.header.msg_name = &message.peer_address; |
| 7059 | message.header.msg_namelen = sizeof(sockaddr_in); |
| 7060 | // Each message will be made of just one buffer |
| 7061 | message.header.msg_iov = &message.io_vec; |
| 7062 | message.header.msg_iovlen = 1; |
| 7063 | // ... and that buffer is a member of our `message` |
| 7064 | message.io_vec.iov_base = message.buffer.data(); |
| 7065 | message.io_vec.iov_len = message.buffer.size(); |
| 7066 | message.status = status_t::pending_k; |
| 7067 | } |
| 7068 | |
| 7069 | // Let's register all of those with `IORING_REGISTER_BUFFERS` |
| 7070 | std::vector<struct iovec> iovecs_to_register; |
| 7071 | for (message_t &message : messages_) iovecs_to_register.push_back(message.io_vec); |
| 7072 | if (io_uring_register_buffers(&ring_, iovecs_to_register.data(), iovecs_to_register.size()) < 0) |
| 7073 | raise_system_error("Failed to register buffers with io_uring 5.5 server"); |
| 7074 | } |
| 7075 | |
| 7076 | ~rpc_uring55_server() noexcept {} |
| 7077 | void close() noexcept { |
| 7078 | ::close(socket_descriptor_); |
| 7079 | io_uring_queue_exit(&ring_); |
| 7080 | } |
| 7081 | |
| 7082 | void stop() noexcept { should_stop_.store(true, std::memory_order_seq_cst); } |
| 7083 | |
| 7084 | void operator()() noexcept { |
| 7085 | // Submit initial receive operations |
nothing calls this directly
no outgoing calls
no test coverage detected