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

Class rpc_uring55_server

less_slow.cpp:7028–7125  ·  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. */

Source from the content-addressed store, hash-verified

7026 * to setup the UDP socket, and process many requests concurrently.
7027 */
7028class 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

Callers

nothing calls this directly

Calls

no outgoing calls

Tested by

no test coverage detected