MCPcopy Create free account
hub / github.com/ashvardanian/less_slow.cpp / rpc_uring60_server

Class rpc_uring60_server

less_slow.cpp:7353–7460  ·  view source on GitHub ↗

* @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 */

Source from the content-addressed store, hash-verified

7351 * - reduces the number of receive operations, using multi-shot receive
7352 */
7353class 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_);

Callers

nothing calls this directly

Calls

no outgoing calls

Tested by

no test coverage detected