* @brief A minimal RPC @b server using @b `io_uring` functionality * to setup the UDP socket, and process many requests concurrently. * * Unlike the `rpc_uring55_server`, this version: * - registers buffers and off-loads buffer selection to the kernel * - reduces the number of receive operations, using multi-shot receive */
| 7351 | * - reduces the number of receive operations, using multi-shot receive |
| 7352 | */ |
| 7353 | class rpc_uring60_server { |
| 7354 | |
| 7355 | int socket_descriptor_; |
| 7356 | sockaddr_in server_address_; |
| 7357 | std::atomic_bool should_stop_; |
| 7358 | io_uring ring_; |
| 7359 | |
| 7360 | // Pre-allocated resources |
| 7361 | mmap_array<message_t> messages_; |
| 7362 | std::size_t max_concurrency_; |
| 7363 | |
| 7364 | public: |
| 7365 | using status_t = message_t::message_status_t; |
| 7366 | |
| 7367 | rpc_uring60_server(std::string const &server_address_str, std::uint16_t port, std::size_t max_concurrency) |
| 7368 | : should_stop_(false), messages_(max_concurrency * 2), max_concurrency_(max_concurrency) { |
| 7369 | |
| 7370 | auto [socket_descriptor, server_address] = rpc_server_socket(port, server_address_str); |
| 7371 | socket_descriptor_ = socket_descriptor; |
| 7372 | server_address_ = server_address; |
| 7373 | |
| 7374 | // Zero copy operations would require more socket options |
| 7375 | int const one = 1; |
| 7376 | if (setsockopt(socket_descriptor_, SOL_SOCKET, SO_ZEROCOPY, &one, sizeof(one)) < 0) |
| 7377 | raise_system_error("Failed to enable zero-copy on socket"); |
| 7378 | |
| 7379 | // Initialize `io_uring` with one slot for each receive/send operation |
| 7380 | // TODO: |= IORING_SETUP_COOP_TASKRUN | IORING_SETUP_SINGLE_ISSUER | IORING_SETUP_SUBMIT_ALL |
| 7381 | auto io_uring_setup_flags = 0; |
| 7382 | if (io_uring_queue_init(max_concurrency * 2, &ring_, io_uring_setup_flags) < 0) |
| 7383 | raise_system_error("Failed to initialize io_uring 6.0 server"); |
| 7384 | if (io_uring_register_files(&ring_, &socket_descriptor_, 1) < 0) |
| 7385 | raise_system_error("Failed to register file descriptor with io_uring 6.0 server"); |
| 7386 | |
| 7387 | // Initialize message resources |
| 7388 | for (message_t &message : messages_) { |
| 7389 | memset(&message.header, 0, sizeof(message.header)); |
| 7390 | message.header.msg_name = &message.peer_address; |
| 7391 | message.header.msg_namelen = sizeof(sockaddr_in); |
| 7392 | // Each message will be made of just one buffer |
| 7393 | message.header.msg_iov = &message.io_vec; |
| 7394 | message.header.msg_iovlen = 1; |
| 7395 | // ... and that buffer is a member of our `message` |
| 7396 | message.io_vec.iov_base = message.buffer.data(); |
| 7397 | message.io_vec.iov_len = message.buffer.size(); |
| 7398 | message.status = status_t::pending_k; |
| 7399 | } |
| 7400 | |
| 7401 | // Let's register all of those with `IORING_REGISTER_BUFFERS` |
| 7402 | std::vector<struct iovec> iovecs_to_register; |
| 7403 | for (message_t &message : messages_) iovecs_to_register.push_back(message.io_vec); |
| 7404 | if (io_uring_register_buffers(&ring_, iovecs_to_register.data(), iovecs_to_register.size()) < 0) |
| 7405 | raise_system_error("Failed to register buffers with io_uring 6.0 server"); |
| 7406 | } |
| 7407 | |
| 7408 | ~rpc_uring60_server() noexcept {} |
| 7409 | void close() noexcept { |
| 7410 | ::close(socket_descriptor_); |
nothing calls this directly
no outgoing calls
no test coverage detected