Async IO with epoll, io_uring, and C++ coroutines

Andreas Hohmann July 26, 2026 #io_uring #Linux #async #c++ #coroutines #epoll

Introduction

Why does the C++ coroutine API look the way it does? This was the main question on my mind when I started the explorations documented in this post. Most C++ coroutine explanations focus on the mechanics of the API, the many customization points, and maybe a few examples. This post approaches C++ coroutines from one of its motivating use cases: concurrent IO with event loops. I find that it's much easier to memorize a complex API when I understand its origin.

To make things concrete, I picked a simple HTTP client as the running example. We will implement many versions of this client:

Moving up this ladder we will see why state machines are a good model for async IO and how coroutines automate the construction of these state machines.

At each stage, we will wrap the APIs in C++ classes in order to simplify the following code. We will, for example, add wrapper classes for sockets, epoll instances, and io_uring.

Besides coroutines we can take advantage of a few more newer C++ features, in particular std::expected, std::println, and std::source_location. The std::expected API in particular allows us to avoid the complexity of exceptions in async code.

Disclaimer: The code in this post is for demonstration purposes only. The code is published under the MIT license. I have made no attempt to make it safe or optimal and skipped some of the best practices (such as a top-level namespace). If you are looking for a complete, production-ready, and highly efficient async C++ library, start with Asio. I used it to implement the "adder" HTTP service that our HTTP clients will talk to.

Project setup

To get started, let's set up the build using cmake. This will use the default C/C++ compiler available on your system (gcc 14 in my case).

cmake_minimum_required(VERSION 3.20)
project(http_client CXX)

set(CMAKE_CXX_STANDARD 23)
set(CMAKE_CXX_STANDARD_REQUIRED ON)
set(CMAKE_EXPORT_COMPILE_COMMANDS ON)

add_executable(http_client
  main.cpp)

if (NOT CMAKE_SYSTEM_NAME STREQUAL "Linux")
    message(FATAL_ERROR "io_uring requires a Linux environment!")
endif ()

find_package(PkgConfig REQUIRED)
pkg_check_modules(URING REQUIRED IMPORTED_TARGET liburing)

target_link_libraries(http_client PRIVATE PkgConfig::URING)

The CMAKE_EXPORT_COMPILE_COMMANDS option lets cmake generate the compile_commands.json file used by LSP (language server protocol) servers to learn how the project is compiled. As the language server looks for this JSON file in the base directory, we add a symbolic link to the generated file in the build directory.

$ ln -s build/compile_commands.json .

We assume that liburing has been installed on the system and add its compiler and linker flags via pkg-config.

The initial file contains only the single main.cpp source file. As we proceed, we have to add the new .cpp files to the add_executable command. To check that we can indeed compile C++23 code, let's define a "modern" C++ hello world program:

#include <print>

int main() {
  std::println("Starting HTTP client");
}

We build and run it with the usual cmake sequence using build as the build directory:

$ cmake -B build -S .
$ cmake --build build
$ ./build/http_client
Starting HTTP client

Errors

We are going to add a thin C++ layer on top of the C APIs such as socket IO, epoll, and io_uring. Most functions will return a std::expected<T, IoError> for some return type T (which may be void) and a simple error type IoError wrapping a Linux errno. The older C APIs rely on the global errno value while newer ones such as io_uring return the error number explicitly (e.g., in the res field of the completion queue entry).

#pragma once

#include <errno.h>   // for errno

#include <format>    // for std::formatter, std::format_context
#include <optional>
#include <source_location>

#include <string.h>  // for strerror

/// Error object used for all IO functions.
struct IoError {
  /// Application-specific code identifying the error; 0 for generic errno.
  int code = 0;

  /// OS error number.
  std::optional<int> os_errno;

  /// Location in application source code where the error originated.
  std::source_location location;

  /// Optional human-readable message.
  std::string message;

  /// Returns an `IoError` with the given `code` and current (caller) location.
  static inline IoError
  from_code(int code, const std::source_location& location = std::source_location::current()) {
    return {
      .code = code,
      .location = location,
    };
  }

  /// Returns an `IoError` with the current errno and (caller) location.
  static inline IoError
  from_errno(const std::source_location& location = std::source_location::current()) {
    return {
      .os_errno = errno,
      .location = location,
    };
  }

  /// Returns an `IoError` with the `err` error code and current (caller) location.
  static inline IoError
  from_errno(int err, const std::source_location& location = std::source_location::current()) {
    return {
      .os_errno = err,
      .location = location,
    };
  }
};

/// Defines formatter for IoError.
template <>
struct std::formatter<IoError> : std::formatter<std::string> {
  auto format(const IoError& error, format_context& ctx) const {
    auto out = std::format_to(ctx.out(), "IoError: code={}", error.code);
    if (error.os_errno.has_value()) {
      const char* msg = strerror(*error.os_errno);
      out = std::format_to(out, " errno={} errno_message='{}'", *error.os_errno, msg);
    }
    out = std::format_to(out, " file={} line={} function={}",
                         error.location.file_name(),
                         error.location.line(),
                         error.location.function_name());
    return out;
  }
};

std::expected doesn't support an error propagation operator such as Rust's ? or Zig's try yet, but we can define a few macros that turn the error check and propagation into one statement. The much shorter code outweighs the cost of the hidden control flow in this case.

In anticipation of our coroutine application, we define the macros with a regular return and in a CO_ variant using co_return.

#pragma once

#include "io_error.h"

#define _EXPECT(return_stmt, var)		   \
  if (!(var)) {				   \
    return_stmt std::unexpected(var.error());  \
  }

#define _EXPECT_VOID(return_stmt, expr)            \
  if (auto _result = (expr); !_result) {           \
    return_stmt std::unexpected(_result.error());  \
  }

#define _EXPECT_SUCCESS(return_stmt, expr)               \
  if (int _err = (expr); _err < 0) {                    \
    return_stmt std::unexpected(IoError::from_errno()); \
  }

#define _EXPECT_SUCCESS_OR_NEG_ERRNO(return_stmt, expr)      \
  if (int _err = (expr); _err < 0) {			 \
    return_stmt std::unexpected(IoError::from_errno(-_err)); \
  }

#define EXPECT(var) _EXPECT(return, var)
#define CO_EXPECT(var) _EXPECT(co_return, var)

#define EXPECT_VOID(expr) _EXPECT_VOID(return, expr)
#define CO_EXPECT_VOID(expr) _EXPECT_VOID(co_return, expr)

#define EXPECT_SUCCESS(expr) _EXPECT_SUCCESS(return, expr)
#define CO_EXPECT_SUCCESS(expr) _EXPECT_SUCCESS(co_return, expr)

#define EXPECT_SUCCESS_OR_NEG_ERRNO(expr) _EXPECT_SUCCESS_OR_NEG_ERRNO(return, expr)
#define CO_EXPECT_SUCCESS_OR_NEG_ERRNO(expr) _EXPECT_SUCCESS_OR_NEG_ERRNO(co_return, expr)

EXPECT_SUCCESS and EXPECT_SUCCESS_OR_NET_ERRNO take care of the two patterns for error codes in C APIs mentioned above.

File descriptors and sockets

Another key property of the wrappers will be resource handling via RAII. A file descriptor, for example, can be wrapped in a class that automatically closes the file descriptor in the destructor (unless the object is moved).

#pragma once

#include <expected>
#include <utility>

#include <unistd.h>     // for close

/// Thin wrapper of a file descriptor ensuring that it is closed exactly once.
class Fd {
 public:
  explicit Fd(int fd) : fd_(fd) {}

  ~Fd() {
    close_file();
  }

  Fd(Fd&& other) : fd_(std::exchange(other.fd_, -1)) {}

  Fd& operator=(Fd&& other) {
    if (this != &other) {
      close_file();
      fd_ = other.fd_;
      other.fd_ = -1;
    }
    return *this;
  }

  int fd() const { return fd_; }

private:
  /// Linux file descriptor or -1 if moved from.
  int fd_;

  void close_file() {
    if (fd_ >= 0) {
      close(fd_);
    }
  }
};

Next, we can define a Socket class that wraps a file descriptor returned by a successful call to the socket function. The only way to create a socket is the create method which may fail.

#pragma once

#include <expected>
#include <utility>

#include <unistd.h>     // for close

#include "fd.h"
#include "io_error.h"

/// Wrapper of a socket file descriptor.
///
/// The `Socket` is neither connected (client socket) nor listening (server
/// socket) yet.
class Socket {
public:
  struct Options {
    /// If true, the socket will be created in non-blocking mode.
    bool async = false;
  };

  const Options& options() const { return options_; }

  int fd() const { return fd_.fd(); }

  std::expected<int, IoError> error_status() const;

  /// Tries to create a new `Socket`.
  static std::expected<Socket, IoError>
  create(Socket::Options options={.async = false});

private:
  Socket(Options options, int fd) : options_(options), fd_(fd) {}

  Options options_;
  Fd fd_;
};

We also added an error_status function that retrieves a socket's error status using the socket option API (getsockop). We will need this function later to check the result of the connect function in the asynchronous case.

#include "socket.h"

#include <sys/socket.h> // for socket

#include "expected_macros.h"

std::expected<Socket, IoError> Socket::create(Socket::Options options) {
  int type = SOCK_STREAM | (options.async ? SOCK_NONBLOCK : 0);
  int fd = socket(AF_INET, type, 0);
  EXPECT_SUCCESS(fd);
  return Socket(options, fd);
}

std::expected<int, IoError> Socket::error_status() const {
  int err;
  socklen_t len = sizeof(err);
  int result = getsockopt(fd(), SOL_SOCKET, SO_ERROR, &err, &len);
  EXPECT_SUCCESS(result);
  return err;
}

Let's test the construction and moving of sockets in our main program. Note that C++ automatically creates the move operations based on the underlying file descriptor wrapper and deletes the copy operations as desired.

#include <print>

#include "socket.h"

int main() {
  std::expected<Socket, IoError> socket = Socket::create();
  if (!socket) {
    return -1;
  }
  std::println("created socket: fd={}", socket->fd());

  Socket socket2 = std::move(*socket);
  std::println("socket2: fd={} socket: fd={}", socket2.fd(), socket->fd());

  // syntax error (use of deleted function 'Socket::Socket(const Socket&)')
  // Socket socket3 = socket2;
}

Before we create a client socket connected to some (TCP) server, there is another API worth wrapping, namely the construction of an IP socket address. The following InetAddress class wraps an internet socket address (sockaddr_in) and adds parsing and printing.

#pragma once

#include <expected>
#include <format>
#include <string>

#include <netinet/in.h>

#include "io_error.h"

/// Wraps an internet (TCP/UDP) socket address.
class InetAddress {
 public:
  /// Tries to create a socket address for the given `host` and `port`.
  ///
  /// The `host` should be an IP address in string format, for example "127.0.0.1"
  /// or "::1".
  static std::expected<InetAddress, IoError> create(const char* host, int port);

  const struct sockaddr_in& addr() const {
    return addr_;
  }

  std::string to_string() const;

  enum ErrorCode {
    InvalidNetworkAddressString = 200,
  };
 private:
  explicit InetAddress(struct sockaddr_in addr): addr_(addr) {}

  struct sockaddr_in addr_;
};

template <>
struct std::formatter<InetAddress> : std::formatter<std::string> {
  auto format(const InetAddress& addr, std::format_context& ctx) const {
    return std::formatter<std::string>::format(addr.to_string(), ctx);
  }
};

Here is the implementation:

#include "inet_address.h"

#include <cstring>
#include <arpa/inet.h>  // for inet_pton, inet_ntop

std::string InetAddress::to_string() const {
  char buf[INET_ADDRSTRLEN];
  if (inet_ntop(AF_INET, &addr_.sin_addr, buf, sizeof(buf)) == nullptr) {
    return "unknown";
  }
  return std::format("{}:{}", buf, ntohs(addr_.sin_port));
}

std::expected<InetAddress, IoError> InetAddress::create(const char* host, int port) {
  struct sockaddr_in addr;
  memset(&addr, 0, sizeof(addr));
  addr.sin_family = AF_INET;
  addr.sin_port = htons(port);
  int result = inet_pton(AF_INET, host, &addr.sin_addr);
  if (result == 0) {
    return std::unexpected(IoError::from_code(InetAddress::InvalidNetworkAddressString));
  }
  return InetAddress(addr);
}

Next up is a class representing a socket connected to a server. We call it ClientSocket. A client socket is a socket so that inheritance is appropriate.

#pragma once

#include <expected>
#include <utility>

#include "inet_address.h"
#include "io_error.h"
#include "socket.h"

/// Wraps a socket that is connected to a TCP server.
class ClientSocket : public Socket {
 public:
  enum ErrorCode {
    ExpectedBlockingSocket = 100,
    ExpectedNonBlockingSocket = 101,
  };

  /// Tries to connect the `socket` to the `server_address`.
  static std::expected<ClientSocket, IoError>
  connect(Socket socket, InetAddress server_address);

  /// Tries to create a client socket from a `Socket` that has been connected
  /// asynchronously.
  static std::expected<ClientSocket, IoError> from_connected(Socket socket);
 private:
  explicit ClientSocket(Socket socket) : Socket(std::move(socket)) {}
};

The connect method is the blocking factory method while from_connected constructs a ClientSocket from a non-blocking socket that has been connected asynchronously (with or without success).

To be able to share the connect logic later, we define the socket utility functions separately.

#pragma once

#include <expected>

#include "inet_address.h"
#include "socket.h"

std::expected<void, IoError>
socket_connect(const Socket& socket, const InetAddress& server_address);

This function works for both blocking and non-blocking sockets. The non-blocking sockets will be used with epoll and io_uring to connect asynchronously to the server.

#include "socket_util.h"

#include <cerrno>    // for errno and EAGAIN

#include <sys/socket.h>  // for connect

std::expected<void, IoError>
socket_connect(const Socket& socket, const InetAddress& server_address) {
  struct sockaddr* addr = (struct sockaddr*)&server_address.addr();
  int result = connect(socket.fd(), addr, sizeof(struct sockaddr_in));
  if (result < 0 && !(socket.options().async && errno == EINPROGRESS)) {
    return std::unexpected(IoError::from_errno());
  }
  return {};
}

With this helper function, we can now implement the Socket wrapper. It would be safer to have different classes for blocking and non-blocking sockets to avoid the runtime checks, but the resulting proliferation of types looks like the worse tradeoff in this case.

#include "client_socket.h"

#include <sys/socket.h>  // for connect

#include "expected_macros.h"
#include "socket_util.h"

std::expected<ClientSocket, IoError>
ClientSocket::connect(Socket socket, InetAddress server_address) {
  if (socket.options().async) {
    return std::unexpected(IoError::from_code(ClientSocket::ExpectedBlockingSocket));
  }
  std::expected<void, IoError> result = socket_connect(socket, server_address);
  EXPECT(result);
  return ClientSocket(std::move(socket));
}

std::expected<ClientSocket, IoError> ClientSocket::from_connected(Socket socket) {
  if (!socket.options().async) {
    return std::unexpected(IoError::from_code(ClientSocket::ExpectedNonBlockingSocket));
  }
  std::expected<int, IoError> error_status = socket.error_status();
  EXPECT(error_status);
  EXPECT_SUCCESS(*error_status);
  return ClientSocket(std::move(socket));
}

For non-blocking sockets, the connect call returns immediately with the EINPROGRESS error code. When the connect completes (whether successful or not), the socket transitions to the writable state (for select or epoll) and one has to check the socket's error state explicitly by retrieving the SO_ERROR socket option with getsockopt.

We can now connect to a server (albeit without sending any data yet):

#include <print>

#include "client_socket.h"
#include "expected_macros.h"
#include "inet_address.h"

std::expected<void, IoError> test_socket() {
  std::expected<Socket, IoError> socket = Socket::create();
  EXPECT(socket);
  std::expected<InetAddress, IoError> address = InetAddress::create("127.0.0.1", 8080);
  EXPECT(address);
  std::expected<ClientSocket, IoError> client_socket =
    ClientSocket::connect(std::move(*socket), *address);
  EXPECT(client_socket);
  std::println("connected to server: {}", *address);
  return {};
}

int main() {
  std::expected<void, IoError> result = test_socket();
  if (!result) {
    std::println("error: {}", result.error());
    return -1;
  }
  return 0;
}

HTTP client interface and buffers

After these preparations we can finally turn to the HTTP clients. Let's start with the interface. The send method takes the request as a string and returns a std::expected containing the response returned by the server (also as a string). The method is synchronous. The async implementations must run the event loop in this method until the HTTP client is done.

#pragma once

#include <expected>
#include <string>

#include "inet_address.h"
#include "io_error.h"

class HttpClient {
 public:
  virtual ~HttpClient() = default;

  virtual std::expected<std::string, IoError> send(std::string data) = 0;

  const InetAddress& server_address() const { return server_address_; }
 protected:
  HttpClient(InetAddress server_address): server_address_(server_address) {}

 private:
  InetAddress server_address_;
};

We picked std::string as input and output to keep the API simple. A more efficient implementation would use a streaming API that takes a stream of data as input and returns the response as a stream as data chunks. The string input and output puts the burden on the HTTP client to send the data in blocks and collect the response blocks.

All IO APIs use buffers (memory blocks) in some fashion. The C APIs take a void pointer (const for a read-only buffer) and a length. In C++, we can combine pointer and length in a std::span and use std::span<const char> for read-only buffers and std::span<char> for read-write buffers. It would be more appropriate to use std::byte instead of char for the buffers, but this leads to too much casting when working with strings. Most IO functions also come in scatter-gather flavors that allow for reading into and writing from multiple buffer fragments with a single API call (see struct iovec, readv and writev). We will stick to the simple single-buffer APIs in our examples.

The memory used by buffers has to be managed. In C++, the easiest solution is a std::unique_ptr<char[]> which takes care of the deallocation via RAII.

The following BufferVec class helps the HTTP clients to collect the response data.

#pragma once

#include <memory>
#include <span>
#include <string>
#include <vector>

class BufferVec {
 public:
  /// Creates a new `BufferVec` whose buffers will be of size `buffer_size`.
  BufferVec(size_t buffer_size = 1024): buffer_size_(buffer_size) {}

  /// Returns the next buffer to write to.
  std::span<char> next_read_buffer();

  /// Advance the write position by `n` bytes.
  void advance(size_t n);

  /// Creates a single string containing the data collected in this `BufferVec`.
  std::string to_string() const;
 private:
  /// Size of the memory blocks.
  size_t buffer_size_;

  /// Position where the next byte should be written.
  size_t pos_ = 0;

  /// Vector of memory blocks (buffers).
  std::vector<std::unique_ptr<char[]>> buffers_;

  size_t capacity() const {
    return buffers_.size() * buffer_size_;
  }
};

The next_read_buffer gives us the next memory block to write to. The IO APIs tells us how many bytes were received which we can pass to the advance method. At the end, we can call to_string to get the collected data as one string.

#include "buffer_vec.h"

#include <algorithm>  // for std::min

std::span<char> BufferVec::next_read_buffer() {
  while (capacity() <= pos_) {
    buffers_.emplace_back(std::make_unique<char[]>(buffer_size_));
  }
  int i = pos_ / buffer_size_;
  size_t offset = pos_ % buffer_size_;
  size_t length = (i + 1) * buffer_size_ - pos_;
  return std::span(buffers_[i].get() + offset, length);
}

void BufferVec::advance(size_t n) {
  pos_ += n;
}

std::string BufferVec::to_string() const {
  std::string s;
  s.reserve(pos_);
  size_t n = 0;
  while (n < pos_) {
    int i = n / buffer_size_;
    size_t length = std::min(buffer_size_, pos_ - i * buffer_size_);
    s.append(buffers_[i].get(), length);
    n += length;
  }
  return s;
}

Blocking HTTP client

The first implementation calls the blocking IO APIs. Here is the declaration of this BlockingHttpClient. We will omit these declarations for the following HTTP clients and only show their implementations.

#pragma once

#include "http_client.h"

class BlockingHttpClient : public HttpClient {
 public:
  explicit BlockingHttpClient(InetAddress server_address) :
    HttpClient(server_address) {}

  std::expected<std::string, IoError> send(std::string data) override;
};

The wrapper libraries and std::expected error handling result in a fairly concise and clear implementation. The typical loops handle the potential partial reads and writes.

#include "blocking_http_client.h"

#include <unistd.h>  // for read and write

#include "buffer_vec.h"
#include "client_socket.h"
#include "expected_macros.h"

std::expected<std::string, IoError> BlockingHttpClient::send(std::string data) {
  std::expected<Socket, IoError> socket = Socket::create();
  EXPECT(socket);

  std::expected<ClientSocket, IoError> client_socket =
    ClientSocket::connect(std::move(*socket), server_address());
  EXPECT(client_socket);

  int size = data.size();
  int sent = 0;
  while (sent < size) {
    int result = write(client_socket->fd(), data.data() + sent, size - sent);
    EXPECT_SUCCESS(result);
    if (result == 0) {
      break;
    }
    sent += result;
  }

  BufferVec buffer_vec;
  int n;
  do {
    std::span<char> buffer = buffer_vec.next_read_buffer();
    n = read(client_socket->fd(), buffer.data(), buffer.size());
    EXPECT_SUCCESS(n);
    buffer_vec.advance(n);
  } while(n > 0);
  return buffer_vec.to_string();
}

Next, let's test the client from our main program. The above mentioned "adder" HTTP service adds two integers given as the query parameters a and b. Our main program constructs an HTTP 1.1 request asking the service to add 4 and 5.

#include <print>

#include "blocking_http_client.h"
#include "expected_macros.h"
#include "inet_address.h"

const std::string request =
  "GET /add?a=4&b=5 HTTP/1.1\r\nHost: localhost\r\nConnection: "
  "close\r\n\r\n";

std::expected<std::string, IoError> run_http_client() {
  std::expected<InetAddress, IoError> address = InetAddress::create("127.0.0.1", 8080);
  EXPECT(address);

  BlockingHttpClient http_client(*address);
  return http_client.send(request);
}

int main() {
  std::expected<std::string, IoError> response = run_http_client();
  if (response) {
    std::println("received response: {}", *response);
  } else {
    std::println("error: {}", response.error());
    return -1;
  }
  return 0;
}

Here is the resulting output:

$ cmake --build build
...
$ build/http_client
received response:
HTTP/1.1 200 OK
Content-Type: text/plain
Connection: close
Content-Length: 10

4 + 5 = 9

This may look like a lot of code for so little functionality, but most of the pieces are completely reusable and the final program is short and easy to understand.

Epoll HTTP client

Our next goal is an asynchronous version of the HTTP client using epoll. An epoll instance is also identified by a file descriptor. Let's wrap this with a few nice methods similar to the client socket.

#pragma once

#include <expected>

#include <sys/epoll.h>

#include "fd.h"
#include "io_error.h"

/// Wrapper of an epoll file descriptor.
class Epoll {
 public:
  explicit Epoll(int fd) : fd_(fd) {}

  static std::expected<Epoll, IoError> create();

  int fd() const { return fd_.fd(); }

  /// Tells epoll to observe events on the socket with the `socket_fd`.
  std::expected<void, IoError> add(int events, int socket_fd) {
    return control(EPOLL_CTL_ADD, events, socket_fd);
  }

  /// Modifies the events that epoll observes for the `socket_fd`.
  std::expected<void, IoError> modify(int events, int socket_fd) {
    return control(EPOLL_CTL_MOD, events, socket_fd);
  }

  /// Waits (blocking) until a new event occurs.
  //
  // Returns the number of events that are ready.
  std::expected<int, IoError> wait(struct epoll_event* events, int events_size);

  /// Thin wrapper of `epoll_ctl` that translates the error result.
  std::expected<void, IoError> control(int operation,
                                       int socket_fd,
                                       struct epoll_event* event);

 private:
  /// File descriptor of the epoll object.
  ///
  /// Automatically closes the epoll object on destruction.
  Fd fd_;

  /// Thin wrapper of `epoll_ctl`.
  std::expected<void, IoError> control(int operation,
                                       int events,
                                       int socket_fd);
};

The implementation follows the now familiar pattern of translating the result codes to IoErrors. Epoll allows us to register our interest in IO "readiness" events such as "read to read from fd" or "ready to write to fd". The functions to add, modify, or delete the interest are fast (not blocking on any IO). Only the epoll_wait function blocks and returns when one of the registered readiness events triggers. We will call the wait function in the event loop.

In typical C-API fashion, the epoll API uses a single flexible epoll_ctl function for all event manipulations. We split them into multiple methods in the Epoll class for easier use.

#include "epoll.h"

#include "expected_macros.h"

std::expected<void, IoError> Epoll::control(int operation,
                                            int socket_fd,
                                            struct epoll_event* event) {
  EXPECT_SUCCESS(epoll_ctl(fd(), operation, socket_fd, event));
  return {};
}

std::expected<void, IoError> Epoll::control(int operation,
                                            int events,
                                            int socket_fd) {
  struct epoll_event ev;
  ev.events = events;
  ev.data.fd = socket_fd;  // set user data to the socket file descriptor
  return control(operation, socket_fd, &ev);
}

std::expected<int, IoError> Epoll::wait(struct epoll_event* events,
                                        int events_size) {
  int event_count = epoll_wait(fd(), events, events_size, -1);
  EXPECT_SUCCESS(event_count);
  return event_count;
}

std::expected<Epoll, IoError> Epoll::create() {
  int fd = epoll_create1(0);
  EXPECT_SUCCESS(fd);
  return Epoll(fd);
}

Our HTTP client performs three IO operations: connect, write, and read. To use epoll with the asynchronous IO operations, we need to create a non-blocking socket with the SOCK_NONBLOCK option. We prepared this already by adding the async parameter to our Socket::Options. As a result, the connect function will return immediately.

Next, we can wait with epoll until the socket is writable and then check the outcome of the connection attempt with the socket's error_status.

The easiest way to manage the events and event loop is a little state machine.

#include "epoll_http_client.h"

#include <cerrno>    // for errno and EAGAIN
#include <unistd.h>  // for read and write

#include "buffer_vec.h"
#include "client_socket.h"
#include "epoll.h"
#include "expected_macros.h"
#include "socket_util.h"

namespace {

enum class EpollClientState {
  Initial,
  Connecting,
  Writing,
  Reading,
  Done,
};

}  // namespace

std::expected<std::string, IoError> EpollHttpClient::send(std::string data) {
  EpollClientState state = EpollClientState::Initial;

  std::expected<Socket, IoError> socket = Socket::create(Socket::Options{.async = true});
  EXPECT(socket);
  int fd = socket->fd();

  std::expected<Epoll, IoError> epoll = Epoll::create();
  EXPECT(epoll);

  // Initially track write readiness for completion of "connect".
  EXPECT_VOID(epoll->add(EPOLLOUT, fd));

  EXPECT_VOID(socket_connect(*socket, server_address()));
  state = EpollClientState::Connecting;

  // Number of (request) bytes that have been written.
  size_t sent = 0;

  // Array of events we can receive (only one).
  struct epoll_event events[1];

  std::optional<ClientSocket> client_socket;

  BufferVec buffers;

  // epoll event loop.
  while (state != EpollClientState::Done) {
    // Blocking call waiting for readiness of the observed file descriptor.
    std::expected<int, IoError> event_count = epoll->wait(events, 1);
    EXPECT(event_count);

    // Process the event (at most one in our case).
    for (int i = 0; i < *event_count; ++i) {
      struct epoll_event& event = events[i];
      int fd = event.data.fd;

      // Processing depends on the event and the current state.
      if (event.events & EPOLLOUT) {
        if (state == EpollClientState::Connecting) {
          std::expected<ClientSocket, IoError> connect_result =
            ClientSocket::from_connected(std::move(*socket));
          EXPECT(connect_result);

          client_socket = std::move(*connect_result);
          state = EpollClientState::Writing;
        }

        // Write remaining request bytes and check for error.
        int n = write(fd, data.data() + sent, data.size() - sent);
        if (n > 0) {
          sent += n;
          // If whole request was sent, switch epoll to read events and
          // transition to `Reading` state.
          if (sent >= data.size()) {
            EXPECT_VOID(epoll->modify(EPOLLIN, fd));
            state = EpollClientState::Reading;
          }
        } else if (n < 0 && errno != EAGAIN) {
          return std::unexpected(IoError::from_errno());
        }
      }
      if (event.events & EPOLLIN) {
        std::span<char> buffer = buffers.next_read_buffer();
        int n = read(fd, buffer.data(), buffer.size());
        if (n > 0) {
          buffers.advance(n);
        } else if (n == 0) {
          state = EpollClientState::Done;
          break;
        } else if (errno != EAGAIN) {
          return std::unexpected(IoError::from_errno());
        }
      }
    }
  }

  return buffers.to_string();
}

Let's extend the main program so that we can switch between the two HTTP client implementations.

#include <cstdio>
#include <functional>
#include <map>
#include <memory>
#include <optional>
#include <print>
#include <string>

#include "blocking_http_client.h"
#include "collection_util.h"
#include "epoll_http_client.h"
#include "expected_macros.h"
#include "inet_address.h"

const std::string request =
  "GET /add?a=4&b=5 HTTP/1.1\r\nHost: localhost\r\nConnection: "
  "close\r\n\r\n";

using HttpClientFactory = std::function<std::unique_ptr<HttpClient>(InetAddress)>;

template <typename T>
std::function<std::unique_ptr<HttpClient>(InetAddress)> create_http_client_factory() {
  return [](InetAddress server_address) {
    return std::make_unique<T>(server_address);
  };
}

std::map<std::string_view, HttpClientFactory> http_client_factories {
  {"blocking", create_http_client_factory<BlockingHttpClient>()},
  {"epoll", create_http_client_factory<EpollHttpClient>()},
};


std::expected<std::string, IoError>
run_http_client(HttpClientFactory& http_client_factory) {
  std::expected<InetAddress, IoError> address = InetAddress::create("127.0.0.1", 8080);
  EXPECT(address);

  std::unique_ptr<HttpClient> http_client = http_client_factory(*address);
  return http_client->send(request);
}

int main(int argc, const char* argv[]) {
  std::string_view http_client_name = argc > 1 ? argv[1] : "blocking";
  std::optional<HttpClientFactory> http_client_factory =
    find_as_optional(http_client_factories, http_client_name);

  if (http_client_factory == std::nullopt) {
    std::println(stderr, "no such http client factory: {}", http_client_name);
    return -1;
  }
  std::println("http client: {}", http_client_name);
  std::expected<std::string, IoError> response = run_http_client(*http_client_factory);
  if (response) {
    std::println("received response:\n{}", *response);
  } else {
    std::println(stderr, "error: {}", response.error());
    return -1;
  }
  return 0;
}

The find_as_optional function is a little utility function compensating for C++'s awkward collection API (that's a story for another day).

#pragma once

#include <functional>  // for std::reference_wrapper
#include <optional>

template <typename Map, typename Key>
auto find_as_optional(Map& m, const Key& key)
  -> std::optional<std::reference_wrapper<typename Map::mapped_type>>
{
  auto it = m.find(key);
  if (it != m.end()) {
      return it->second;
  }
  return std::nullopt;
}

io_uring

Epoll is fairly limited. It does not support regular files, for example. After a few more attempt of async IO APIs, it looks like Linux has now settled on io_uring as the universal API for asynchronous operations.

io_uring follows the "proactor model": We hand buffers over to the operating system, and the operating system informs us when the operation such as reading or writing the buffer is complete (successful or not). That's in contrast to epoll's "reactor model" where epoll informs us about the readiness of the socket for input or output, and we have to call the actual read and write functions afterwards. Both approaches have pros and cons. io_uring's approach is more efficient especially when using registered buffers. However, the memory management of io_uring buffers is trickier because we have to hand the buffers to the operating system and can only reclaim them once we have received the associated completion event. Using epoll, the buffers are only borrowed by the operating system during the synchronous (non-blocking) calls of the read and write functions.

Similar to our epoll approach, we'll first wrap the io_uring functions that we will need for the HTTP client in a C++ API. The IoUring class wraps the io_uring struct of the liburing library. We use a std::unique_ptr to make sure that this structure is not moved.

The io_uring API relies on two queues (ring buffers), the submission queue to which we submit new operation requests in the form of submission queue entries (SQEs), and the completion queue where we'll find the completion queue entries (CQEs) telling us the result of these operations. A CQE has a res (result) field. If negative, it contains the negation of the error code. If non-negative, it contains the result of the operation such as the number of bytes read or written.

A submission queue entry has a user_data field (64 bit) that is echoed in the completion queue entry. This field can be used as a correlation ID between the request (submission) and response (completion) but can also carry any other data such as, for example, the pointer to a callback function or object.

The IoUringCqe wrapper in the following interface uses RAII to return the completion entry back to io_uring.

#pragma once

#include <expected>
#include <functional>
#include <memory>
#include <span>

#include <liburing.h>

#include "inet_address.h"
#include "io_error.h"

/// Wrapper of an io_uring completion queue entry (CQE) pointer that marks
/// the entry as seen in the destructor.
class IoUringCqe {
 public:
  IoUringCqe(struct io_uring* ring, struct io_uring_cqe* cqe) :
    ring_(ring), cqe_(cqe) {}

  IoUringCqe(IoUringCqe&& other) : ring_(other.ring_), cqe_(other.cqe_) {
    other.cqe_ = nullptr;
  }

  ~IoUringCqe() {
    mark_seen();
  }

  IoUringCqe& operator=(IoUringCqe&& other) {
    if (this != &other) {
      mark_seen();
      ring_ = other.ring_;
      cqe_ = other.cqe_;
      other.cqe_ = nullptr;
    }
    return *this;
  }

  const io_uring_cqe& cqe() const { return *cqe_; }
 private:
  struct io_uring* ring_;
  struct io_uring_cqe* cqe_;

  void mark_seen() {
    if (cqe_ != nullptr) {
      io_uring_cqe_seen(ring_, cqe_);
    }
  }
};

/// Thin wrapper around a liburing io_uring struct.
///
/// The wrapper makes sure that the destructed properly.
class IoUring {
 public:
  ~IoUring();

  IoUring(IoUring&&) = default;

  static std::expected<IoUring, IoError> create();

  struct io_uring* ring() { return ring_.get(); }

  /// Submits a request to connect socket `fd` to the `server_address`.
  ///
  /// The `user_data` will be echoed in the associated CQE.
  std::expected<void, IoError>
  connect(int fd, const InetAddress& server_address, int64_t user_data);

  /// Submits a request to send `data` to the socket `fd`.
  std::expected<void, IoError>
  send(int fd, std::span<const char> data, int64_t user_data);

  /// Submits a request to receive data from the socket `fd` into the `buffer`.
  std::expected<void, IoError>
  recv(int fd, std::span<char> buffer, int64_t user_data);

  /// Waits for the next completion entry.
  std::expected<IoUringCqe, IoError> wait();

  enum ErrorCode {
    IoUringSubmissionQueueIsFull = 300,
  };

 private:
  IoUring(std::unique_ptr<struct io_uring> ring) : ring_(std::move(ring)) {}

  std::expected<struct io_uring_sqe*, IoError> get_sqe();

  std::unique_ptr<struct io_uring> ring_;
};

Each operation has an associated "prep" function in the liburing library such as io_uring_prep_connect to connect a socket or io_uring_prep_send to send data to a socket. This leads to a short and consistent implementation of our wrapper:

#include "io_uring.h"

#include <netinet/in.h>

#include "expected_macros.h"
#include "socket.h"

std::expected<IoUring, IoError> IoUring::create() {
  std::unique_ptr<struct io_uring> ring = std::make_unique<struct io_uring>();
  int result = io_uring_queue_init(8, ring.get(), 0);
  EXPECT_SUCCESS_OR_NEG_ERRNO(result);
  return IoUring(std::move(ring));
}

IoUring::~IoUring() {
  if (ring_.get()) {
    io_uring_queue_exit(ring_.get());
  }
}

std::expected<struct io_uring_sqe*, IoError> IoUring::get_sqe() {
  struct io_uring_sqe* sqe = io_uring_get_sqe(ring());
  if (sqe == nullptr) {
    return std::unexpected(IoError{
        .code = IoUringSubmissionQueueIsFull,
        .location = std::source_location::current()
      });
  }
  return sqe;
}

std::expected<void, IoError> IoUring::connect(
    int fd,
    const InetAddress& server_address,
    int64_t user_data) {
  struct sockaddr* addr = (struct sockaddr*)&server_address.addr();
  std::expected<struct io_uring_sqe*, IoError> sqe = get_sqe();
  EXPECT(sqe);
  io_uring_prep_connect(*sqe, fd, addr, sizeof(struct sockaddr_in));
  (*sqe)->user_data = user_data;
  io_uring_submit(ring());
  return {};
}

std::expected<void, IoError>
IoUring::send(int fd, std::span<const char> data, int64_t user_data) {
  std::expected<struct io_uring_sqe*, IoError> sqe = get_sqe();
  EXPECT(sqe);
  io_uring_prep_send(*sqe, fd, data.data(), data.size(), 0);
  (*sqe)->user_data = user_data;
  io_uring_submit(ring());
  return {};
}

std::expected<void, IoError>
IoUring::recv(int fd, std::span<char> buffer, int64_t user_data) {
  std::expected<struct io_uring_sqe*, IoError> sqe = get_sqe();
  EXPECT(sqe);
  io_uring_prep_recv(*sqe, fd, buffer.data(), buffer.size(), 0);
  (*sqe)->user_data = user_data;
  io_uring_submit(ring());
  return {};
}

std::expected<IoUringCqe, IoError> IoUring::wait() {
  struct io_uring_cqe* cqe;
  int result = io_uring_wait_cqe(ring(), &cqe);
  EXPECT_SUCCESS_OR_NEG_ERRNO(result);
  return IoUringCqe(ring(), cqe);
}

io_uring HTTP client

Using the IoUring wrapper, we can implement the io_uring HTTP client. Similar to the epoll implementation, we use a little state machine again. This time, we name the states after the three phases of the HTTP client: connect, write, and read.

#include "io_uring_http_client.h"

#include <algorithm>  // for std::min

#include <arpa/inet.h>
#include <liburing.h>
#include <netinet/in.h>

#include "buffer_vec.h"
#include "expected_macros.h"
#include "io_uring.h"
#include "socket.h"

namespace {

  constexpr size_t send_buffer_size = 1024;

  enum class Op { Connect, Write, Read };

}  // namespace

std::expected<std::string, IoError> IoUringHttpClient::send(std::string data) {
  std::expected<IoUring, IoError> io_uring = IoUring::create();
  EXPECT(io_uring);

  std::expected<Socket, IoError> socket =
    Socket::create(Socket::Options{.async = true});
  EXPECT(socket);
  int fd = socket->fd();

  EXPECT_VOID(io_uring->connect(fd, server_address(), static_cast<int64_t>(Op::Connect)));

  int sent = 0;
  BufferVec buffers;

  auto next_send_data = [&]() -> std::span<const char> {
    return { data.data() + sent, std::min(send_buffer_size, data.size() - sent) };
  };

  for (bool done = false; !done;) {
    std::expected<IoUringCqe, IoError> cqe = io_uring->wait();
    EXPECT(cqe);
    int res = cqe->cqe().res;
    EXPECT_SUCCESS_OR_NEG_ERRNO(res);
    Op op = static_cast<Op>(cqe->cqe().user_data);
    switch (op) {
    case Op::Connect:
      EXPECT_VOID(io_uring->send(fd, next_send_data(), static_cast<int64_t>(Op::Write)));
      break;
    case Op::Write:
      sent += res;
      if (sent < data.size()) {
        EXPECT_VOID(io_uring->send(fd, next_send_data(), static_cast<int64_t>(Op::Write)));
      } else {
        std::span<char> buffer = buffers.next_read_buffer();
        EXPECT_VOID(io_uring->recv(fd, buffer, static_cast<int64_t>(Op::Read)));
      }
      break;
    case Op::Read:
      if (res == 0) {  // EOF, server closed connection
        done = true;
      } else {
        buffers.advance(res);
        std::span<char> buffer = buffers.next_read_buffer();
        EXPECT_VOID(io_uring->recv(fd, buffer, static_cast<int64_t>(Op::Read)));
      }
      break;
    }
  }
  return buffers.to_string();
}

The main loop calls the io_uring wait function with a callback that handles the completion queue event according to the current state. In every case, we first check if the completion was successful and return an IoError otherwise.

In the Connect state, the completion means that the connection was established and we can start sending data. In the Write state, the completion result tells us how many bytes were sent. We update the sent counter accordingly and keep sending data if there is still data to be sent or switch to reading otherwise. The Read state similarly keeps reading and collecting the response data until it's done.

io_uring with "manual" state machine

The io_uring HTTP client is a relatively short program. There is a problem, however. The application logic is intertwined with the main event loop (calling wait). Imagine what happens if we add more functionality to this loop such as calling multiple clients. If we want to split the state machine from the event loop, we need to define an explicit state machine object. Let's first try this "by hand". In the next sections, we will use coroutines instead.

Our HTTP client task will implement the following interface. It has a single resume method that takes the io_uring completion result of the last operation (or zero for first resumption).

#pragma once

#include <expected>

#include "io_error.h"

class IoUringTask {
 public:
  virtual std::expected<void, IoError> resume(int res) = 0;
};

If we store a pointer to such a task in the io_uring user_data, we can resume the task in response to a completion (CQE). Let's create another io_uring wrapper using these task callbacks.

#pragma once

#include <expected>
#include <span>
#include <utility>

#include "inet_address.h"
#include "io_error.h"
#include "io_uring.h"
#include "io_uring_task.h"

class TaskIoUring {
 public:
  TaskIoUring(IoUring io_uring) : io_uring_(std::move(io_uring)) {}

  /// Connects the socket `fd` to the `server_address` and resumes the
  /// `task` on completion.
  std::expected<void, IoError>
  connect(int fd, InetAddress server_address, IoUringTask* task);

  /// Sends the `data` to the socket `fd` and resumes the `task` on completion.
  std::expected<void, IoError>
  send(int fd, std::span<const char> data, IoUringTask* task);

  /// Reads data from the socket `fd` into the `buffer` and resumes the
  /// `task` on completion.
  std::expected<void, IoError>
  recv(int fd, std::span<char> buffer, IoUringTask* task);

  /// Waits for the next completion (CQE) and resumes the associated task.
  std::expected<void, IoError> wait();
 private:
  IoUring io_uring_;
};

The implementation calls the IoUring methods with the task pointers as user data. The wait method handling the completion entries then calls the resume method of these tasks.

#include "task_io_uring.h"

#include <cerrno>
#include <concepts>  // for std::invocable

#include <netinet/in.h>

#include "expected_macros.h"

namespace {

std::expected<void, IoError>
submit(IoUring& io_uring, IoUringTask* task, std::invocable<io_uring_sqe&> auto command) {
  struct io_uring_sqe* sqe = io_uring_get_sqe(io_uring.ring());
  if (sqe == nullptr) {
    return std::unexpected(IoError::from_errno(EBUSY));
  }
  command(*sqe);
  sqe->user_data = reinterpret_cast<int64_t>(task);
  io_uring_submit(io_uring.ring());
  return {};
}

}

std::expected<void, IoError>
TaskIoUring::connect(int fd, InetAddress server_address, IoUringTask* task) {
  return submit(io_uring_, task, [&](io_uring_sqe& sqe) {
    struct sockaddr* addr = (struct sockaddr*)&server_address.addr();
    io_uring_prep_connect(&sqe, fd, addr, sizeof(struct sockaddr_in));
  });
}

std::expected<void, IoError>
TaskIoUring::send(int fd, std::span<const char> data, IoUringTask* task) {
  return submit(io_uring_, task, [&](io_uring_sqe& sqe) {
    io_uring_prep_send(&sqe, fd, data.data(), data.size(), 0);
  });
}

std::expected<void, IoError>
TaskIoUring::recv(int fd, std::span<char> buffer, IoUringTask* task) {
  return submit(io_uring_, task, [&](io_uring_sqe& sqe) {
    io_uring_prep_recv(&sqe, fd, buffer.data(), buffer.size(), 0);
  });
}

std::expected<void, IoError> TaskIoUring::wait() {
  std::expected<IoUringCqe, IoError> wait_cqe = io_uring_.wait();
  EXPECT(wait_cqe);
  const struct io_uring_cqe& cqe = wait_cqe->cqe();
  IoUringTask* task = reinterpret_cast<IoUringTask*>(cqe.user_data);
  EXPECT_VOID(task->resume(cqe.res));
  return {};
}

Now we can turn the handler of the completion events in the io_uring HTTP client function shown above into a state machine implementing the IoUringTask interface.

#pragma once

#include <optional>
#include <span>
#include <string>

#include "buffer_vec.h"
#include "inet_address.h"
#include "io_uring_task.h"
#include "socket.h"
#include "task_io_uring.h"

class HttpClientIoUringTask : public IoUringTask {
 public:

  enum class State { Initial, Connecting, Writing, Reading, Done };

  HttpClientIoUringTask(TaskIoUring& task_io_uring,
                        InetAddress server_address,
                        std::string_view data):
    task_io_uring_(task_io_uring), server_address_(server_address), data_(data) {}

  std::expected<void, IoError> resume(int res) override;

  State state() const { return state_; }
  std::string response() const { return buffers_.to_string(); }

  bool complete() const { return state_ == State::Done; }

 private:
  TaskIoUring& task_io_uring_;
  InetAddress server_address_;
  std::string_view data_;
  State state_ = State::Initial;

  std::optional<Socket> socket_ = std::nullopt;
  int sent_ = 0;
  BufferVec buffers_;

  std::span<const char> next_send_data(size_t buffer_size_max) {
    return { data_.data() + sent_, std::min(buffer_size_max, data_.size() - sent_) };
  }
};

The local variables of the function must become fields of the task because they have to survive multiple calls of the resume method. However, we have to somehow handle the lifetime of the variables and create the objects at the right time. That's why we use a std::optional for the socket. The BufferVec solves this problem for the read buffer.

#include "http_client_io_uring_task.h"

#include "expected_macros.h"
#include "socket.h"

namespace {

constexpr size_t send_buffer_size = 1024;

}

std::expected<void, IoError> HttpClientIoUringTask::resume(int res) {
  EXPECT_SUCCESS_OR_NEG_ERRNO(res);
  switch (state_) {
  case State::Initial: {
    std::expected<Socket, IoError> s = Socket::create();
    EXPECT(s);
    socket_ = std::move(*s);
    state_ = State::Connecting;
    EXPECT_VOID(task_io_uring_.connect(socket_->fd(), server_address_, this));
    break;
  }
  case State::Connecting:
    state_ = State::Writing;
    EXPECT_VOID(task_io_uring_.send(socket_->fd(), next_send_data(send_buffer_size), this));
    break;
  case State::Writing:
    sent_ += res;
    if (sent_ < data_.size()) {
      EXPECT_VOID(task_io_uring_.send(socket_->fd(), next_send_data(send_buffer_size), this));
    } else {
      state_ = State::Reading;
      EXPECT_VOID(task_io_uring_.recv(socket_->fd(), buffers_.next_read_buffer(), this));
    }
    break;
  case State::Reading:
    if (res == 0) {  // EOF, server closed connection
      state_ = State::Done;
    } else {
      buffers_.advance(res);
      EXPECT_VOID(task_io_uring_.recv(socket_->fd(), buffers_.next_read_buffer(), this));
    }
    break;
  case State::Done:
    break;
  }
  return {};
}

Having implemented the task and the task-oriented io_uring wrapper, we can combine the two to implement the IoUringTaskHttpClient.

#include "io_uring_task_http_client.h"

#include "expected_macros.h"
#include "http_client_io_uring_task.h"
#include "io_uring.h"

std::expected<std::string, IoError> IoUringTaskHttpClient::send(std::string data) {
  std::expected<IoUring, IoError> io_uring = IoUring::create();
  EXPECT(io_uring);

  TaskIoUring task_io_uring(std::move(*io_uring));
  HttpClientIoUringTask task(task_io_uring, server_address(), data);
  task.resume(0);
  while (!task.complete()) {
    EXPECT_VOID(task_io_uring.wait());
  }
  return task.response();
}

Somewhat more interestingly, we can run two clients concurrently:

#include "io_uring_task_http_client.h"

#include "expected_macros.h"
#include "http_client_io_uring_task.h"
#include "io_uring.h"

std::expected<std::string, IoError> IoUringTaskHttpClient::send(std::string data) {
  std::expected<IoUring, IoError> io_uring = IoUring::create();
  EXPECT(io_uring);

  TaskIoUring task_io_uring(std::move(*io_uring));

  HttpClientIoUringTask task1(task_io_uring, server_address(), data);
  HttpClientIoUringTask task2(task_io_uring, server_address(), data);
  task1.resume(0);
  task2.resume(0);
  while (!task1.complete() || !task2.complete()) {
    EXPECT_VOID(task_io_uring.wait());
  }
  return task1.response() + task2.response();
}

As we have seen, writing these state machines by hand gets tricky in C++ because we have to manage the lifetimes of the objects manually. In a function such as the send method above, we can rely on the local variables and their lifetime (bound to their lexical scope). What we would like to have is a combination of the two approaches: A function with local variables that can be suspended and resumed, that is, an automatic mechanism that turns a normal function into a state machine. That's exactly what C++ coroutines provide.

Coroutines

The description of the C++ coroutine API can fill many pages that make it hard to see the forest for all the trees.

It took me a while to realize that C++ coroutines are a lot less magic than they might appear at first sight. Their main purpose is to instruct the compiler to implement a function as a state machine. If the compiler sees one of the co_ keywords (co_await, co_return, and co_yield) in the body of a function definition, it generates a state machine (the coroutine) based on the body of the function and returns an object (of the type declared in the function's signature) that (typically) references this coroutine object. The state machine allows us to suspend and resume the function. There is a suspension point at the beginning and at the end of the coroutine and at every co_await and co_yield (the explicit suspension points).

The function declaration of a coroutine is a normal function declaration with a normal return type. There is no special async keyword as in other languages that changes how a function is called. In this sense, there is no "function coloring". If the function is declared as a virtual method, one implementation can implement it "manually" and another implementation can be a coroutine. If we want the compiler to implement the state machine, the compiler needs to be able to derive the so-called promise type (see below) from the return type. Apart from this, the return type is completely up to us.

The second aspect of C++ coroutines I had to get used to is the API style. The coroutine API relies on "static duck typing". The compiler looks for methods with certain names and signatures, often going through several options until a matching function or type is found. Besides the duck typing, the API has an object-oriented flavor in the sense that the functionality is spread evenly across multiple (stateful) objects whose interactions result in the desired behavior. This makes it very flexible, but also harder to learn because one has to understand all the objects and their interactions before the whole API makes sense.

The coroutine API uses two main types: The promise type allows us to customize the coroutine as a whole, and the awaitable type allows us to customize each suspension point. Each coroutine definition has its promise type, and each suspension point has its awaitable type. At runtime, there is exactly one instance of the promise type per coroutine instance and one instance of an awaitable for each (potential) suspension.

Because of the single promise type instance per coroutine instance, the promise can be used as a data channel between the calling program, the coroutine, and the awaitables. We can move data from the coroutine into the promise and then via an awaitable or the coroutine's return value from the promise to the calling program. That's similar to how promises work in other future/promise APIs (that is, the name "promise" makes sense).

Similarly, we can pass data to the coroutine at each suspension point by storing the data in the awaiter and letting the awaiter return the data on resumption to the coroutine.

Coroutine API

Let's start with the awaitable API. The expression following the co_await keyword must be an "awaitable" which is something that can be converted to an "awaiter". An "awaiter" object must have three methods:

The await_ready method does not take any arguments and always return bool (true when the suspension should be skipped). The await_suspend method is more complicated. First and crucially, it gets a handle of the calling coroutine as its single argument. As we will see, we can go back and forth between the coroutine handle and the associated promise. The return type is (confusingly) flexible. It can be void (just suspend), a bool (true meaning "suspend" - the opposite of the await_ready method), or even another coroutine handle. The await_resume method takes no arguments and returns whatever co_await should return (which may be void).

As an example, the built-in awaiter std::suspend_always could be implemented as follows:

struct MySuspendAlways {
  bool await_ready() const noexcept { return false; }
  void await_suspend(std::coroutine_handle<> handle) const noexcept {}
  void await_resume() const noexcept {}
};

We are not skipping the suspension (await_ready return false), we are actually suspending (await_suspend returns void), and we are not returning anything when resuming.

A more interesting example is the following awaiter that allows us to get the handle to the current coroutine.

template <typename PromiseType>
struct GetCoroutineHandleAwaiter {
  std::coroutine_handle<PromiseType> handle;

  bool await_ready() const noexcept { return false; }

  bool await_suspend(std::coroutine_handle<PromiseType> handle) noexcept {
    this->handle = handle;
    return false;
  }

  std::coroutine_handle<PromiseType> await_resume() const noexcept { return handle; }
};

The awaiter instance is used as a temporary holder for the coroutine handle that the compiler passed to the await_suspend method. This method only copies the handle and returns true to directly continue. The await_resume can then return the captured handle.

auto handle = co_await GetCoroutineHandleAwaiter<MyPromiseType>{};

The coroutine API probably deserves a more direct way to get the handle of the current coroutine, but this "trick" definitely demonstrates how flexible the API is.

Awaiters control the suspension points, but to control the whole coroutine we need the promise type. The compiler derives the promise type from the coroutine's return type. The return type is also called the "coroutine interface" because it's the interface that the caller uses to interact with the coroutine. In the simplest case, the return type has a nested promise_type struct or class, but one can also specify a trait to map the return type to a promise type.

The promise type is the centerpiece of the coroutine API. Let's see what the main customization points are:

The code generated by the compiler for the coroutine performs the following steps:

Because the compiler knows the memory layout and relative position of the promise instance in the coroutine frame, it can derive the pointer to the coroutine from the pointer to the promise. The coroutine handle is a reified coroutine frame pointer. We can therefore get the coroutine handle from the promise pointer and the reference to the promise from the coroutine handle. This gives us access to the coroutine handle from every promise method. We can use it to store the handle in the return object when constructing the return object in get_return_object (which is almost always done). We can pass it to the awaiters returned by initial_suspend and final_suspend.

An async coroutine type (Z)

The following coroutine return type demonstrates how these extension points work together to provide an async API. Despite all the machinery, the coroutine return type is relatively small. Our heavily commented Z type is less than 200 lines long.

#pragma once

#include <cassert>
#include <coroutine>
#include <optional>
#include <utility>

/// Coroutine return type capturing a task that produces a value of type `T`.
///
/// `Z` (as in "Zukunft" - German for "future") wraps the coroutine handle and
/// thereby the coroutine frame and the promise contained in it. The promise
/// acts as a "normal" promise in async APIs, that is, as the data exchange from
/// producer to consumer of some value of type `T`.
template <typename T>
class Z {
public:
  class promise_type;
  using Handle = std::coroutine_handle<promise_type>;

  Z(std::coroutine_handle<promise_type> handle) : handle_(handle) {
  }

  Z(Z&& other) noexcept : handle_(std::exchange(other.handle_, nullptr)) {}

  ~Z() {
    if (handle_) {
      handle_.destroy();
    }
  }

  // Returns an immediately-completed task.
  static Z<T> completed(T value) {
    co_return value;
  }

  Z& operator=(Z&& other) noexcept {
    if (this != &other) {
      if (handle_) {
        handle_.destroy();
      }
      handle_ = std::exchange(other.handle_, nullptr);
    }
    return *this;
  }

  /// Promise type for coroutines producing a value of type `T`.
  ///
  /// The value is moved into the promise when the coroutine returns
  /// with `co_return` and move out of the promise with the custom
  /// `result` method.
  class promise_type {
  public:
    T result() {
      assert(result_.has_value() && "Result was already consumed");
      T val = std::move(*result_);
      result_.reset();
      return val;
    }

    /// Constructs the coroutine return object for this promise type.
    ///
    /// This method is called automatically after the coroutine frame has been
    /// created with the promise contained in it.
    Z get_return_object() {
      auto handle = std::coroutine_handle<promise_type>::from_promise(*this);
      return Z(handle);
    }

    /// Returns the awaitable that's awaited automatically when the coroutine is
    /// started.
    ///
    /// For now, we always suspend.
    std::suspend_always initial_suspend() {
      return {};
    }

    /// Awaiter for the `final_suspend`.
    ///
    /// Resumes the parent if available.
    struct FinalAwaiter {
      /// Don't skip suspension.
      bool await_ready() noexcept { return false; }

      /// If there is a parent handle, perform symmetric transfer to this parent
      /// coroutine.
      std::coroutine_handle<> await_suspend(Handle handle) noexcept {
        auto& promise = handle.promise();
        if (promise.parent_handle) {
          return promise.parent_handle;
        } else {
          return std::noop_coroutine();
        }
      }
      /// Do nothing on resumption (will never be called).
      void await_resume() noexcept {}
    };

    /// Returns the awaitable that's awaited automatically just before the
    /// coroutine completes.
    FinalAwaiter final_suspend() noexcept {
      return FinalAwaiter();
    }

    /// Handles any exceptions thrown during the coroutine execution.
    ///
    /// For now, we avoid exceptions (and use std::expected instead).
    void unhandled_exception() {}

    /// Handles the `co_return` statements with a value.
    ///
    /// The value is moved to the `result_` that can then be extracted by the
    /// caller of the coroutine with the `result` method.
    void return_value(T result) {
      result_.emplace(std::move(result));
    }

    /// Handle of the coroutine that awaits the coroutine of this promise (when
    /// the return value is used as an awaitable).
    std::coroutine_handle<> parent_handle;

  private:
    std::optional<T> result_;
  };

  // Allows functions returning Z to be co_awaited.
  auto operator co_await() && noexcept {
    struct Awaiter {
      Handle handle;

      /// Returns true (to skip the suspension) if the coroutine is already done.
      bool await_ready() const noexcept {
        return handle.done();
      }

      /// Saves the calling coroutine's handle in the `parent_handle` of the promise.
      void await_suspend(std::coroutine_handle<> awaiting) noexcept {
        handle.promise().parent_handle = awaiting;
      }

      /// Returns the coroutine's result on resumption.
      T await_resume() {
        return handle.promise().result();
      }
    };
    return Awaiter{handle_};
  }

  /// Resumes the associated coroutine.
  ///
  /// Returns true if the coroutine is done.
  void resume() {
    if (!handle_.done()) {
      handle_.resume();
    }
  }

  T result() { return std::move(handle_.promise().result()); }

  bool done() const {
    return handle_.done();
  }

private:
  /// Handle of the coroutine being suspended.
  Handle handle_;
};

The FinalAwaiter is an example of how the promise type and awaiters work together. The promise keeps the parent_handle and the FinalAwaiter uses it to resume the parent coroutine (if there is any). Note that this resumption happens without a function call. The C++ coroutine API guarantees that the resumption of another coroutine via the return value of await_suspend happens directly on the state machine level without a nested function call. This "symmetric transfer" prevents stack overflows when chaining coroutines.

Coroutine handle awaiter

The following awaiter implements the "trick" for obtaining the coroutine from the coroutine itself. We will use this below to register the coroutine address with io_uring (in the user_data of the submission queue entry).

#pragma once

#include <coroutine>

/// Awaiter that allows for getting the coroutine handle from within a
/// coroutine using
///
///    auto handle = co_await co_get_coroutine_handle<MyPromiseType>();
///
/// The awaiter gets the handle passed to the `await_suspend` method and then
/// does not suspend and returns the captured handle via `await_resume`.
template <typename PromiseType>
struct GetCoroutineHandleAwaiter {
  std::coroutine_handle<PromiseType> handle;

  // Returns false to force the compiler to call await_suspend.
  bool await_ready() const noexcept { return false; }

  // Intercepts the typed coroutine handle.
  bool await_suspend(std::coroutine_handle<PromiseType> handle) noexcept {
    this->handle = handle;
    return false;  // Return false to immediately resume the coroutine
  }

  // Returns the coroutine handle to the co_await expression
  std::coroutine_handle<PromiseType> await_resume() const noexcept { return handle; }
};

/// Helper function to simplify the syntax for getting the coroutine handle
/// from within the coroutine itself.
template <typename PromiseType>
auto get_coroutine_handle() {
  return GetCoroutineHandleAwaiter<PromiseType>{};
}

io_uring event loop with coroutines

Let's define an io_uring wrapper for io_uring programs that use coroutines that suspend when awaiting for completions. The first thing we need is an awaiter for these suspension points. The awaiter submits the IO operation in await_suspend and resumes with the result of the io_uring completion.

#pragma once

#include <coroutine>
#include <functional>

#include "io_uring.h"

using IoUringPrep = std::move_only_function<void(io_uring_sqe&)>;

/// Awaiter for io_uring operations.
///
/// Submits the IO operation during suspension and resumes with the io_uring
/// completion result.
class IoUringAwaiter {
 public:
  IoUringAwaiter(IoUring& io_uring, IoUringPrep prep):
    io_uring_(io_uring), prep_(std::move(prep)) {}

  /// Always attempts to suspend and submit the io_uring operation.
  bool await_ready() noexcept { return false; }

  /// Submits the SQE prepared with `prep` to io_uring on suspension.
  ///
  /// If the submission queue is full, the coroutine will not suspend and
  /// instead immediately resume with -EBUSY.
  bool await_suspend(std::coroutine_handle<> handle) noexcept;

  /// Resumes with the io_uring completion result.
  int await_resume() noexcept { return result_; }

  /// Resume with the io_uring completion `result`.
  void resume(int result) {
    result_ = result;
    handle_.resume();
  }

 private:
  /// io_uring instance the submission entries are submitted to.
  IoUring& io_uring_;

  /// Function preparing the submission queue entries.
  std::move_only_function<void(io_uring_sqe&)> prep_;

  /// Handle of the suspended coroutine.
  std::coroutine_handle<> handle_;

  /// Result from the io_uring completion entry.
  int result_;
};

To be able to submit the IO operation, we need the io_uring structure (in the form of our IoUring wrapper) and the function preparing the submission queue entry (SQE). We use a std::move_only_function for that.

The await_suspend performs these steps. It also handles the case when the submission queue is full.

#include "io_uring_awaiter.h"

#include <cerrno>

bool IoUringAwaiter::await_suspend(std::coroutine_handle<> handle) noexcept {
  struct io_uring_sqe* sqe = io_uring_get_sqe(io_uring_.ring());
  if (sqe == nullptr) {
    // Submission queue is full. Set error and prevent suspension.
    result_ = -EBUSY;
    return false;
  }
  handle_ = handle;
  prep_(*sqe);

  // Pass pointer to this awaiter as user_data to io_uring.
  sqe->user_data = reinterpret_cast<int64_t>(this);
  io_uring_submit(io_uring_.ring());

  return true;
}

The only thing left is the handling of the completion entries associated with these awaiters. Let's define yet another io_uring wrapper that does that. It wraps our low-level IoUring wrapper and simplifies the creation of the awaiters for io_uring operations.

#pragma once

#include <expected>

#include "io_error.h"
#include "io_uring.h"
#include "io_uring_awaiter.h"
#include "inet_address.h"

/// Low-level coroutine wrapper around io_uring.
///
/// The methods such as `connect`, `write`, and `read` return an `IoUringAwaiter`
/// that submits the associated io_uring SQE and awaits its completion.
///
/// The overall state machine is moved forward by the `wait` method that waits
/// for the next completion.
class ZIoUring {
 public:
  ZIoUring(IoUring io_uring) : io_uring_(std::move(io_uring)) {}

  /// Returns an awaiter that submits an io_uring operation prepared with
  /// `prep` on suspension and resumes with the io_uring completion result.
  IoUringAwaiter submit(IoUringPrep prep) {
    return IoUringAwaiter(io_uring_, std::move(prep));
  }

  /// Connects the socket `fd` to the `server_address`.
  IoUringAwaiter connect(int fd, InetAddress server_address);

  /// Sends the `data` to the socket `fd`.
  ///
  /// Returns the number of bytes read or the negative errno.
  IoUringAwaiter send(int fd, std::span<const char> data);

  /// Reads data from socket `fd` into `buffer`.
  ///
  /// Returns the number of bytes read or the negative errno.
  IoUringAwaiter recv(int fd, std::span<char> buffer);

  /// Waits for the next completion (completion queue entry - CQE) and
  /// resumes its awaiter with the io_uring completion result.
  std::expected<void, IoError> wait();

 private:
  /// Thin io_uring wrapper.
  IoUring io_uring_;
};

The interesting part is the wait method. It recovers the IoUringAwaiter from the SQE's user_data and calls its resume method with the io_uring result. The resume method stores the result in the awaiter itself and returns it to the suspended coroutine via the await_resume method of the coroutine API.

#include "z_io_uring.h"

#include <concepts>

#include <netinet/in.h>

#include "coroutine_handle_awaiter.h"
#include "expected_macros.h"
#include "socket.h"

IoUringAwaiter ZIoUring::connect(int fd, InetAddress server_address) {
  return submit([fd, server_address](io_uring_sqe& sqe) {
    struct sockaddr* addr = (struct sockaddr*)&server_address.addr();
    io_uring_prep_connect(&sqe, fd, addr, sizeof(struct sockaddr_in));
  });
}

IoUringAwaiter ZIoUring::send(int fd, std::span<const char> data) {
  return submit([fd, data](io_uring_sqe& sqe) {
    io_uring_prep_send(&sqe, fd, data.data(), data.size(), 0);
  });
}

IoUringAwaiter ZIoUring::recv(int fd, std::span<char> buffer) {
  return submit([fd, buffer](io_uring_sqe& sqe) {
    io_uring_prep_recv(&sqe, fd, buffer.data(), buffer.size(), 0);
  });
}

std::expected<void, IoError> ZIoUring::wait() {
  // Wait for the next completion entry (CQE).
  std::expected<IoUringCqe, IoError> wait_cqe = io_uring_.wait();
  EXPECT(wait_cqe);
  const struct io_uring_cqe& cqe = wait_cqe->cqe();

  // Get the awaiter from the CQE's `user_data` and resume it with the io_uring result.
  IoUringAwaiter* awaiter = reinterpret_cast<IoUringAwaiter*>(cqe.user_data);
  awaiter->resume(cqe.res);
  return {};
}

HTTP client with io_uring and coroutines

Finally! An HTTP client using io_uring and coroutines. I hope you appreciate the clarity of the resulting code after all this work and enjoy the insights you gained in the process.

#include "z_io_uring_http_client.h"

#include <expected>
#include <span>
#include <string>

#include "buffer_vec.h"
#include "expected_macros.h"
#include "io_error.h"
#include "io_uring.h"
#include "socket.h"
#include "z_io_uring.h"

Z<std::expected<std::string, IoError>>
ZIoUringHttpClient::send_async(ZIoUring& z_io_uring, std::span<const char> data) {
  std::expected<Socket, IoError> s = Socket::create();
  CO_EXPECT(s);
  int fd = s->fd();
  int res = co_await z_io_uring.connect(fd, server_address());
  CO_EXPECT_SUCCESS_OR_NEG_ERRNO(res);

  int sent = 0;
  do {
    res = co_await z_io_uring.send(fd, data.subspan(sent, data.size() - sent));
    CO_EXPECT_SUCCESS_OR_NEG_ERRNO(res);
    sent += res;
  } while (sent < data.size());

  BufferVec response;
  do {
    res = co_await z_io_uring.recv(fd, response.next_read_buffer());
    CO_EXPECT_SUCCESS_OR_NEG_ERRNO(res);
    response.advance(res);
  } while (res > 0);
  co_return response.to_string();
}

std::expected<std::string, IoError> ZIoUringHttpClient::send(std::string data) {
  std::expected<IoUring, IoError> io_uring = IoUring::create();
  EXPECT(io_uring);

  ZIoUring z_io_uring(std::move(*io_uring));

  Z<std::expected<std::string, IoError>> z = send_async(z_io_uring, data);

  // start the coroutine (after the initial suspension)
  z.resume();

  // run the io_uring event loop until done
  while (!z.done()) {
    z_io_uring.wait();
  }

  return z.result();
}

Combining coroutines

The coroutine-based HTTP client is still just a single one-shot HTTP client that requires a lot of boilerplate code. The current example does not even exercise the nested coroutine logic in the Z type (the co_await operator).

To get a glimpse of what's possible with these building blocks, let's add another example. Out new main program, concurrent_main.cpp performs two calls to the adder service in parallel from an outer coroutine using a generic zip coroutine combinator.

#include <expected>
#include <print>
#include <utility>
#include <string>
#include <expected>

#include "expected_macros.h"
#include "inet_address.h"
#include "io_error.h"
#include "io_uring.h"
#include "z_io_uring.h"
#include "z_io_uring_http_client.h"
#include "z.h"

using HttpResponse = std::expected<std::string, IoError>;

/// Concurrent combinator that runs two tasks concurrently and returns their paired results.
template <typename T1, typename T2>
Z<std::pair<T1, T2>> zip(Z<T1> lhs, Z<T2> rhs) {
  // Start both tasks after the initial resumption.
  lhs.resume();
  rhs.resume();

  // Await both tasks. If the task isn't done, the parent (zip) suspends.
  // When the task completes, its FinalAwaiter resumes this zip coroutine.
  std::println("awaiting lhs");
  T1 r1 = co_await std::move(lhs);
  std::println("awaiting rhs");
  T2 r2 = co_await std::move(rhs);

  co_return std::make_pair(std::move(r1), std::move(r2));
}

std::expected<void, IoError> send_requests() {
  std::expected<IoUring, IoError> io_uring = IoUring::create();
  EXPECT(io_uring);
  ZIoUring z_io_uring(std::move(*io_uring));
  std::expected<InetAddress, IoError> address = InetAddress::create("127.0.0.1", 8080);
  EXPECT(address);
  ZIoUringHttpClient client(*address);

  std::string req1 = "GET /add?a=4&b=5 HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n";
  std::string req2 = "GET /add?a=10&b=20 HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n";

  Z<HttpResponse> task1 = client.send_async(z_io_uring, req1);
  Z<HttpResponse> task2 = client.send_async(z_io_uring, req2);

  auto zipped_task = zip(std::move(task1), std::move(task2));

  zipped_task.resume();

  while (!zipped_task.done()) {
    EXPECT_VOID(z_io_uring.wait());
  }

  auto [res1, res2] = zipped_task.result();

  if (res1) std::println("Response 1:\n{}", *res1);
  else std::println(stderr, "Response 1 Error: {}", res1.error());

  if (res2) std::println("Response 2:\n{}", *res2);
  else std::println(stderr, "Response 2 Error: {}", res2.error());

  return {};
}

int main() {
  std::expected<void, IoError> result = send_requests();
  if (!result) {
    std::println(stderr, "error: {}", result.error());
    return -1;
  }
}

This example demonstrates how we can construct generic building blocks base on the single coroutine type that can be combined in many ways to implement an async application. The example also raises the next questions: How should we combined the errors of the combined tasks? Should we fail if one of them fails? If so, how can we cancel the other coroutine that may still be running? Those questions sound like interesting topics for another post covering structured concurrency and the latest developments such as C++26 senders and receivers.