Skip to content
JackSparrow414
Go back

Understanding Java NIO (Part 2): I/O Multiplexing and Reactor Servers in C

Table of contents

Open Table of contents

I/O-Bound and CPU-Bound Workloads

The previous post explained that, with blocking I/O, a process or thread sitting idle is a waste of resources. If every process or thread is doing useful work almost all the time, however, those resources are not being wasted, and performance might even be better. This brings us to a question: are the tasks performed by the server’s request-handling threads I/O-bound or CPU-bound?

Suppose all 10,000 connections in the final example from the previous post are doing useful work. The bottleneck is then more likely to be the server hardware itself. We may need to expand from one server to a cluster, or improve the hardware of the single server—for example, by adding CPU cores—to handle all 10,000 connections quickly. This is a CPU-bound workload. Typical examples include image compression, encryption, and decryption.

Conversely, suppose that after 10,000 connections are established, their threads spend most of their time on I/O, and only a few connections occasionally transfer data. This is an I/O-bound workload. Typical examples include an application server waiting for database query results, and uploading or downloading files.

Our assumed scenario: an I/O-bound workload with thousands or tens of thousands of connections, where only a few have data at any given moment and the rest are waiting for data.

The problem we want to solve in this scenario: how can we efficiently find the connections that actually have data?

One Thread per Connection vs. One Thread per I/O Event

The key characteristic of the blocking examples in the previous post was that one thread was bound to one connection. A thread was occupied or blocked while waiting for messages. Even the thread pool in the final version could quickly be exhausted. Threads should instead handle work after an I/O event occurs, so that concurrency is useful. In other words, associate threads with I/O events, occupying a thread only when there is actually data to process.

Multithreading with Nonblocking I/O Polling

Based on these points, we can write our first version. Start by decoupling connections from threads. Setting a connection’s file descriptor to nonblocking I/O prevents a process or thread from blocking on that connection.

int set_nonblocking(int sockfd)
{
    int flags = fcntl(sockfd, F_GETFL);
    if (flags == -1)
    {
        perror("fcntl");
        exit(EXIT_FAILURE);
    }
    flags |= O_NONBLOCK;
    if (fcntl(sockfd, F_SETFL, flags) == -1)
    {
        perror("fcntl");
        exit(EXIT_FAILURE);
    }
    printf("socket is nonblocking I/O now\n");
    return 0;
}

In nonblocking I/O mode, read returns immediately rather than blocking if the connection currently has no I/O event.

Only after the main process detects an I/O event on a connection does it pass that connection to the thread pool. The pool focuses on processing I/O events and the subsequent logic, without being bound to connections.

How does the main process detect an I/O event on a connection? It still uses read. If read reports EAGAIN, the connection currently has no I/O event. If read returns a value greater than 0, the connection has an I/O event and the main process has already read some data into a buffer. It can then hand the work over to the threads.

The pseudocode is as follows:

// List of all established connections
int connections[MAX_CONNECTIONS];
int conn_count = 0;

void main_process() {
    int server_fd = create_server_socket();
    set_nonblocking(server_fd);

    while (1) {

        /* ── 1. Establish new connections ── */
        int client_fd = accept(server_fd, NULL, NULL);
        if (client_fd > 0) {
            set_nonblocking(client_fd);
            connections[conn_count++] = client_fd;
        }

        /* ── 2. Poll established connections for I/O events ── */
        for (int i = 0; i < conn_count; i++) {
            int fd = connections[i];
            char buffer[BUFFER_SIZE];

            ssize_t n = read(fd, buffer, sizeof(buffer));

            if (n == -1 && errno == EAGAIN) {
                // No data on this connection yet; skip it
                continue;
            }

            if (n > 0) {
                // An I/O event occurred; the main process has read the data into buffer
                // Submit it to the thread pool asynchronously to avoid blocking the main process
                task_t *task = create_task(fd, buffer, n);
                thread_pool_submit(task);
            }

            if (n == 0) {
                // The peer closed the connection
                close(fd);
                remove_connection(connections, &conn_count, i);
            }
        }
    }
}

A problem now becomes apparent. The main process must establish connections, poll established connections for I/O events, and read their data. This is worse than one thread per connection because the code is essentially serial again. If many clients try to connect, the main process is busy polling or reading data—since checking readiness itself uses read—and cannot promptly return to the start of the loop to accept connections.

One improvement is to let the main process only establish connections and use a separate thread for polling. This prevents polling from delaying the main process’s ability to accept new connections. It still does not solve the problem of that single polling thread reading data and processing connections serially.

We can improve this further: if one polling thread is a bottleneck, why not use several? This makes sense in our scenario, where only a few of the thousands of connections have I/O events at the same time.

The server architecture now becomes: the main process only accepts new connections and distributes them across four threads using a load-balancing algorithm. Each of those four threads polls a group of established connections. When I/O is ready, it submits a task to a thread pool—say, 100 threads—that performs the actual business processing.

The pseudocode is as follows:

/* ================================================================
 * Architecture overview:
 *   Main process (1) ──accept──▶ Load balancing ──▶ I/O polling threads (4)
 *                                              │
 *                                         I/O readiness detected
 *                                              │
 *                                              ▼
 *                                        Thread pool (100) ── Business processing
 * ================================================================ */

#define IO_THREAD_COUNT   4
#define WORKER_COUNT      100
#define MAX_CONN_PER_IO   (MAX_CONNECTIONS / IO_THREAD_COUNT)

/* ── Connection group owned by each I/O polling thread ── */
typedef struct {
    int  fds[MAX_CONN_PER_IO];
    int  count;
    mutex_t lock;
} io_group_t;

io_group_t io_groups[IO_THREAD_COUNT];   // Four connection groups
thread_pool_t *worker_pool;              // 100 business-processing threads


/* ================================================================
 * Load balancing: assign each new connection to the I/O thread with the fewest connections
 * ================================================================ */
int pick_io_thread() {
    int target = 0;
    for (int i = 1; i < IO_THREAD_COUNT; i++) {
        if (io_groups[i].count < io_groups[target].count)
            target = i;
    }
    return target;
}


/* ================================================================
 * I/O polling threads (×4): poll their groups for I/O events and submit tasks when ready
 * ================================================================ */
void *io_thread(void *arg) {
    io_group_t *group = (io_group_t *)arg;

    while (1) {
        int io_ready_count = 0; // Count connections that are ready in this iteration
        mutex_lock(&group->lock);

        for (int i = 0; i < group->count; i++) {
            int fd = group->fds[i];
            char buffer[BUFFER_SIZE];

            ssize_t n = read(fd, buffer, sizeof(buffer));

            if (n == -1 && errno == EAGAIN) {
                // No data on this connection; poll the next one
                continue;
            }

            if (n > 0) {
                // I/O is ready: data is in buffer; submit it to the business thread pool
                task_t *task = create_task(fd, buffer, n);
                thread_pool_submit(worker_pool, task);
            }

            if (n == 0) {
                // The peer disconnected; clean up the connection
                close(fd);
                remove_fd(group, i);
                i--;  // Adjust the index
            }
        }

        mutex_unlock(&group->lock);
        /* ── Key point: if no connection was ready, sleep to yield the CPU ── */
        if (io_ready_count == 0) {
            usleep(POLL_INTERVAL_US);   // For example, sleep for 100–500 microseconds
        }
    }
}


/* ================================================================
 * Business thread pool (×100): perform the actual business logic
 * ================================================================ */
void *worker_thread(void *arg) {
    while (1) {
        task_t *task = thread_pool_fetch(worker_pool);  // Block while waiting for a task
        handle_business_logic(task->fd, task->buffer, task->len);
        free_task(task);
    }
}


/* ================================================================
 * Main process: only accept and distribute connections; perform no I/O polling
 * ================================================================ */
int main() {
    int server_fd = create_server_socket();
    set_nonblocking(server_fd);

    /* Initialize four I/O polling threads */
    for (int i = 0; i < IO_THREAD_COUNT; i++) {
        io_groups[i].count = 0;
        mutex_init(&io_groups[i].lock);
        thread_create(io_thread, &io_groups[i]);
    }

    /* Initialize a pool of 100 business-processing threads */
    worker_pool = thread_pool_init(WORKER_COUNT);

    /* Main loop: accept connections and leave everything else to the other threads */
    while (1) {
        int client_fd = accept(server_fd, NULL, NULL);
        if (client_fd < 0) continue;

        set_nonblocking(client_fd);

        // Load balancing: choose the I/O thread with the fewest connections
        int idx = pick_io_thread();

        mutex_lock(&io_groups[idx].lock);
        io_groups[idx].fds[io_groups[idx].count++] = client_fd;
        mutex_unlock(&io_groups[idx].lock);
    }
}

The overall architecture is shown below: Server architecture showing a main process, multiple I/O processes, and worker threads

┌─────────────────┬──────────────────────────────────────────────┐
│     Role        │                  Responsibility                       │
├─────────────────┼──────────────────────────────────────────────┤
│ Main process (×1) │ Accept connections; distribute by least connections │
├─────────────────┼──────────────────────────────────────────────┤
│ I/O polling threads (×4) │ Poll connections with read; read data → submit tasks │
├─────────────────┼──────────────────────────────────────────────┤
│ Business thread pool (×100) │ Wait for tasks; execute business logic     │
└─────────────────┴──────────────────────────────────────────────┘

I/O Multiplexing

In this architecture, each of the four threads checks a batch of connections. This is essentially what I/O multiplexing means.

Official definition: I/O multiplexing lets us simultaneously check multiple file descriptors to see whether I/O can be performed on any of them.

The Master–Worker Pattern

This design fully decouples connection management, I/O multiplexing, and business processing. There are a small number of I/O threads to avoid excessive context switching, and many business-processing threads to use the CPU’s parallel processing capabilities. This architecture is called the Master–Worker pattern. Its core idea is to separate task scheduling from execution: the master receives and distributes tasks, while multiple workers process subtasks in parallel, greatly increasing system throughput.

The Drawbacks of Polling

If the sleep interval is too long, I/O events are not handled promptly and latency increases. If it is too short, the program keeps looping without useful work and wastes CPU time.

Since the application cannot know immediately which connections have I/O events, let the kernel do this work. When an I/O event occurs, the kernel can notify the application.

I/O Multiplexing with select

select() blocks until one or more descriptors in the file descriptor sets become ready.

Q: When polling with nonblocking I/O, we use read to determine whether I/O can be performed on a file descriptor. How does select know that a file descriptor is ready for I/O? A: select() checks readiness in kernel space. It directly examines the kernel data structures associated with each fd, rather than probing by calling an I/O function. The kernel maintains corresponding data structures for each descriptor type, and readiness checks read their fields directly—for example, through tcp_poll and pipe_poll.

For example:

  1. Is a socket readable? → Check whether its receive buffer, sk_receive_queue, is nonempty.
  2. Is a socket writable? → Check whether enough space remains in its send buffer.
  3. Is a listening socket ready? → Check whether its accept queue is nonempty.
  4. Is a pipe readable? → Check whether its buffer contains data or its write end has closed.

Here is a simple single-process example of the select system call:

/*
	FD_ZERO() initializes the set pointed to by fdset as empty.
    FD_SET() adds the file descriptor fd to the set pointed to by fdset.
    FD_CLR() removes the file descriptor fd from the set pointed to by fdset.
*/
int server_that_can_process_requests_concurrently_using_io_multiplexing_by_select()
{
    int server_socket = bind_server_socket_to_port_and_listen();
    set_nonblocking(server_socket);
    fd_set active_fd_set, read_fd_set;
    FD_ZERO(&active_fd_set);
    FD_SET(server_socket, &active_fd_set);
    int max_fd = 0;
    if (server_socket >= max_fd)
    {
        max_fd = server_socket + 1;
    }
    while (1)
    {
        read_fd_set = active_fd_set;
        /*
        	The nfds argument must be one greater than the largest descriptor in the three file descriptor sets.
            This makes select() more efficient because the kernel need not check whether higher-numbered descriptors belong to the sets.

            readfds is the set of descriptors to check for input readiness.
            writefds is the set of descriptors to check for output readiness.
            exceptfds is the set of descriptors to check for exceptional conditions.

			The structures pointed to by readfds, writefds, and exceptfds hold the results.
            Before calling select(), initialize these structures with FD_ZERO() and FD_SET()
            to include the descriptors we are interested in.
            select() then modifies the structures; when it returns, they contain the ready descriptors.

			A positive return value means one or more descriptors are ready.
     		The return value is the number of ready descriptors.
        */
        int active_fd_num = select(max_fd, &read_fd_set, NULL, NULL, NULL);
        if (active_fd_num < 0)
        {
            perror("select error");
            exit(EXIT_FAILURE);
        }
        /*

        		Knowing the count is not enough; we need to identify the connections.
        		Check each returned descriptor set with FD_ISSET()
        		to determine which I/O events occurred.
        */
        for (int i = 0; i < max_fd; i++)
        {
            if (FD_ISSET(i, &read_fd_set))
            {
                if (i == server_socket)
                {
                    int new_client_socket = accept(server_socket, NULL, NULL);
                    if (new_client_socket < 0)
                    {
                        perror("accept");
                        exit(EXIT_FAILURE);
                    }
                    printf("new client connected\n");
                    set_nonblocking(new_client_socket);
                    FD_SET(new_client_socket, &active_fd_set);
                    if (new_client_socket >= max_fd)
                    {
                        max_fd = new_client_socket + 1;
                    }
                }
                else
                {
                    int result = read_from_client(i);
                    if (result <= 0)
                    {
                        close(i);
                        FD_CLR(i, &active_fd_set);
                    }
                    else if (result == READ_AGAIN)
                    {
                        print_time();
                        printf("no connection now, please connect after 3 seconds\n");
                        continue;
                    }
                }
            }
        }
    }
}

Understanding the fd_set Bitmap

The fd_set bitmap is an array of 16 unsigned long elements. On a 64-bit Linux machine, this type occupies 8 bytes, or 64 bits, so each element holds 64 bits. Each bit can be used to determine whether a file descriptor is in the bitmap, as follows:

fd_set = unsigned long fds_bits[16]
          └── 16 elements × 64 bits = 1024 bits of total capacity
Step 1: fd / 64 → Determine the array element (slot).
Step 2: fd % 64 → Determine the bit within that element (offset).

fd = 0:   0/64=0  → fds_bits[0],  0%64=0  → bit 0
fd = 5:   5/64=0  → fds_bits[0],  5%64=5  → bit 5
fd = 64:  64/64=1 → fds_bits[1],  64%64=0 → bit 0
fd = 100: 100/64=1→ fds_bits[1],  100%64=36→bit 36
fd = 200: 200/64=3→ fds_bits[3],  200%64=8 → bit 8

The size of this fd_set array is hard-coded. Changing it requires manually modifying the source and recompiling. Although it has 1024 bits, it may not support a particularly high level of concurrency in practice.

Usable capacity = 1024 - 3 reserved fds (0,1,2) - 1 (the server’s own fd) = 1020

Suppose each client connection opens three files, each with a distinct entry in the system-wide open file table. These files stay open for some time, so their descriptors cannot be released and reused. Ignoring other factors, the server can then support at most 1020/5=255 concurrent clients.

Furthermore, the bitmap cannot handle any fd whose value is ≥ 1024.

The Drawbacks of select

  1. The limitations of the fd_set bitmap prevent it from supporting very high concurrency.
  2. Kernel-side inefficiency: the first argument to select gives the maximum descriptor value to check. Even if only a few connections in [0-max_fd] have I/O events, the kernel must traverse the entire set.
  3. Application-side inefficiency: the kernel does not directly tell the application which descriptors have I/O events. The application must scan the entire set to find them, with O(n) complexity.

The later poll() system call improves on select()’s descriptor-count limit, but the performance issue remains. I did not study it in depth and moved directly to the final solution.

epoll

epoll returns the list of file descriptors with I/O events directly from the kernel. The application can use them immediately without checking descriptors one by one.

Here is a simple single-process example of the epoll system calls:

#include <sys/epoll.h>

#define MAX_EVENTS 10

/*
	EPOLL_CTL_ADD adds fd to the interest list of the epoll instance epfd. The events we are interested in are specified in the structure pointed to by ev.
*/
int server_that_can_process_requests_concurrently_using_io_multiplexing_by_event_poll()
{
    int epfd = epoll_create1(0);
    if (epfd == -1)
    {
        perror("epoll_create error");
        exit(1);
    }
    int server_socket = bind_server_socket_to_port_and_listen();
    set_nonblocking(server_socket);
    struct epoll_event ev;
    ev.events = EPOLLIN;
    ev.data.fd = server_socket;
    if (epoll_ctl(epfd, EPOLL_CTL_ADD, server_socket, &ev) == -1)
    {
        perror("epoll_ctl error");
        exit(1);
    }
    struct epoll_event evlist[MAX_EVENTS];
    while (1)
    {
        /*
        	The array of structures pointed to by evlist returns information about ready file descriptors.
        */
        int nready = epoll_wait(epfd, evlist, MAX_EVENTS, -1);
        if (nready == -1)
        {
            perror("epoll_wait error");
            exit(1);
        }
        for (int i = 0; i < nready; i++)
        {
            int fd = evlist[i].data.fd;
            if (fd == server_socket)
            {
                int new_client_socket = accept(server_socket, NULL, NULL);
                if (new_client_socket < 0)
                {
                    perror("accept error");
                    exit(EXIT_FAILURE);
                }
                printf("new epoll client connected\n");
                set_nonblocking(new_client_socket);
                struct epoll_event client_ev;
                client_ev.events = EPOLLIN;
                client_ev.data.fd = new_client_socket;
                if (epoll_ctl(epfd, EPOLL_CTL_ADD, new_client_socket, &client_ev) == -1)
                {
                    perror("epoll_ctl error");
                    exit(1);
                }
            }
            else if (evlist[i].events & (EPOLLIN | EPOLLHUP | EPOLLERR | EPOLLRDHUP))
            {
                int result = read_from_client(fd);
                if (result <= 0)
                {
                    // epoll_ctl(epfd, EPOLL_CTL_DEL, fd, NULL);
                    // After close, the kernel automatically removes fd from epoll; no explicit epoll_ctl is needed
                    close(fd);
                }
                else if (result == READ_AGAIN)
                {
                    print_time();
                    printf("no connection now, please connect after 3 seconds\n");
                    continue;
                }
            }
        }
    }
    close(server_socket);
    close(epfd);
    return 0;
}

Closing File Descriptors and epoll’s Behavior

As discussed in the previous post, the kernel only releases the underlying resources when the descriptor reference count for an entry in the system-wide open file table reaches zero. The same applies here. Once all file descriptors referring to an open file description have been closed, that open file description is removed from epoll’s interest list. If we duplicate a descriptor with dup() or a similar function, or through fork(), the open file is only removed after both the original descriptor and all its copies are closed.

The code above explicitly calls close(), ensuring that the reference count reaches zero, so the EPOLL_CTL_DEL operation is commented out.

Using Docker to Run epoll on macOS

epoll is Linux-specific. For a similar mechanism on macOS, use kqueue. I also considered libuv, which abstracts away platform differences, but decided to focus on epoll for now. I can explore the other two later.

I therefore compile the code in a Docker Linux image. Install the build-base package in the container; it includes gcc, make, and other tools.

services:
  hello-c:
    image: alpine:latest
    container_name: hello-c
    ports:
      - 18080:18080
    command: >
      sh -c "
        apk add --no-cache build-base &&
        echo 'build-base installed successfully' &&
        gcc --version &&
        make --version &&
        tail -f /dev/null
      "
    working_dir: /app
    volumes:
      - type: bind
        source: ./
        target: /app
        read_only: false
docker compose up -d
docker exec -it hello-c sh
make
./hello

Level-Triggered and Edge-Triggered Notifications

Suppose a file descriptor uses nonblocking I/O, an I/O event occurs, and the application reads only part of the available data:

  1. Level-triggered: until the application has read all the data, every evlist returned by epoll_wait includes that descriptor. Once all data is read, evlist no longer includes it. This shows the drawback: if the data cannot be read in one operation, level-triggered handling repeatedly calls epoll_wait, increasing overhead.

    epoll_wait()        ← System call 1: fd is ready and returned
    read(fd, buf, 40)   ← System call 2: read 40 bytes
    
    epoll_wait()        ← System call 3: buffer is nonempty; fd is returned again
    read(fd, buf, 40)   ← System call 4: read 40 bytes
    
    epoll_wait()        ← System call 5: buffer is nonempty; fd is returned again
    read(fd, buf, 20)   ← System call 6: read the remaining 20 bytes
    
    6 system calls in total
  2. Edge-triggered: after one read, the descriptor no longer appears in evlist until new I/O activity occurs on it.

    epoll_wait()        ← System call 1: fd is ready and returned
    
    // Keep reading until drained; otherwise, no further notification arrives
    read(fd, buf, 40)   ← System call 2
    read(fd, buf, 40)   ← System call 3
    read(fd, buf, 20)   ← System call 4: returns EAGAIN; exit the loop
    
    4 system calls in total

    To use edge-triggered notifications, specify EPOLLET when registering events.

    struct epoll_event client_ev;
    client_ev.events = EPOLLIN | EPOLLET;

    With edge-triggered handling, a connection is not fully handled until all its data has been read. A connection that continuously receives data can delay other connections in the same batch. I will not go into this further here; interested readers can look into the solutions.

The Multi-Reactor Pattern

With all these improvements in mind, return to our original scenario and problem:

A server serving many clients must monitor a large number of file descriptors, most of which are idle, with only a few ready at any given time.

Overall Architecture

Our final server architecture builds on the Master–Worker pattern and uses epoll with edge-triggered notifications for efficient handling. This is the Multi-Reactor pattern, shown below:

Multi-Reactor architecture with epoll: main Reactor accepts connections, sub-Reactors handle I/O, and workers process tasks

Implementing the Architecture

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <errno.h>
#include <fcntl.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <sys/epoll.h>
#include <sys/eventfd.h>
#include <pthread.h>

#include "nio_master_slave.h"

#define PORT 8080
#define MAX_EVENTS 100
#define SUB_REACTOR_COUNT 2   // Number of I/O workers
#define WORKER_THREAD_COUNT 4 // Number of business-processing workers
#define MAX_CLIENTS 10000

// --- 1. Connection context ---
typedef struct
{
    int fd;
    int sub_reactor_id;
    char out_buffer[8192]; // Send buffer
    int out_len;
    pthread_mutex_t lock; // Protect out_buffer
} Connection;

Connection *connections[MAX_CLIENTS];

// --- 2. Core component structures ---
typedef struct
{
    int id;
    int epoll_fd;
    int wakeup_fd;
    int fd_queue[1024];
    int queue_count;
    pthread_mutex_t lock;
    pthread_t thread_id;
} SubReactor;

SubReactor sub_reactors[SUB_REACTOR_COUNT];

typedef struct
{
    int client_fd;
    char *data;
} Task;

Task worker_queue[2048];
int worker_task_count = 0;
pthread_mutex_t worker_queue_lock = PTHREAD_MUTEX_INITIALIZER;
pthread_cond_t worker_queue_cond = PTHREAD_COND_INITIALIZER;

// --- Helper functions ---
static void set_nonblocking(int sockfd)
{
    int flags = fcntl(sockfd, F_GETFL, 0);
    fcntl(sockfd, F_SETFL, flags | O_NONBLOCK);
}

void update_epoll(int epoll_fd, int fd, int events)
{
    struct epoll_event ev;
    ev.events = events;
    ev.data.fd = fd;
    epoll_ctl(epoll_fd, EPOLL_CTL_MOD, fd, &ev);
}

// Clean up connection resources
void close_connection(int fd)
{
    Connection *conn = connections[fd];
    if (conn)
    {
        close(fd); // Automatically removed from epoll
        pthread_mutex_destroy(&conn->lock);
        free(conn);
        connections[fd] = NULL;
        printf("Client FD %d disconnected; resources cleaned up\n", fd);
    }
}

// --- 3. Business thread pool (dedicated processing workers) ---
void *worker_thread_loop(void *arg)
{
    while (1)
    {
        Task task;
        pthread_mutex_lock(&worker_queue_lock);
        while (worker_task_count == 0)
        {
            pthread_cond_wait(&worker_queue_cond, &worker_queue_lock);
        }
        task = worker_queue[--worker_task_count];
        pthread_mutex_unlock(&worker_queue_lock);

        // Simulate business processing (0.1 seconds)
        usleep(100000);

        Connection *conn = connections[task.client_fd];
        if (conn)
        {
            char response[2048];
            int len = snprintf(response, sizeof(response), "[Worker Processed] %s", task.data);

            pthread_mutex_lock(&conn->lock);
            // Safely append the result to the send buffer
            if (conn->out_len + len < sizeof(conn->out_buffer))
            {
                memcpy(conn->out_buffer + conn->out_len, response, len);
                conn->out_len += len;

                // Wake the corresponding sub-reactor to send the data
                int sub_epoll = sub_reactors[conn->sub_reactor_id].epoll_fd;
                update_epoll(sub_epoll, task.client_fd, EPOLLIN | EPOLLOUT | EPOLLET);
            }
            else
            {
                printf("Warning: send buffer for FD %d is full; dropping data\n", task.client_fd);
            }
            pthread_mutex_unlock(&conn->lock);
        }
        free(task.data); // Free memory allocated by strdup
    }
    return NULL;
}

// --- 4. Sub-reactor (dedicated I/O, strictly using ET mode) ---
void *sub_reactor_loop(void *arg)
{
    SubReactor *reactor = (SubReactor *)arg;
    struct epoll_event events[MAX_EVENTS];
    char buffer[1024];

    while (1)
    {
        int n = epoll_wait(reactor->epoll_fd, events, MAX_EVENTS, -1);

        for (int i = 0; i < n; i++)
        {
            int fd = events[i].data.fd;

            // --- Task 1: handle new connections passed by the main thread ---
            if (fd == reactor->wakeup_fd)
            {
                uint64_t val;
                read(reactor->wakeup_fd, &val, sizeof(val));

                pthread_mutex_lock(&reactor->lock);
                for (int j = 0; j < reactor->queue_count; j++)
                {
                    int new_client = reactor->fd_queue[j];
                    set_nonblocking(new_client);
                    connections[new_client]->sub_reactor_id = reactor->id;

                    struct epoll_event client_ev;
                    client_ev.events = EPOLLIN | EPOLLET; // Register edge-triggered notifications
                    client_ev.data.fd = new_client;
                    epoll_ctl(reactor->epoll_fd, EPOLL_CTL_ADD, new_client, &client_ev);
                }
                reactor->queue_count = 0;
                pthread_mutex_unlock(&reactor->lock);
            }
            // --- Task 2: handle client I/O ---
            else
            {
                Connection *conn = connections[fd];
                if (!conn)
                    continue;

                // [ET fix] Readable event: use a while loop to drain the kernel buffer!
                if (events[i].events & EPOLLIN)
                {
                    while (1)
                    {
                        ssize_t count = read(fd, buffer, sizeof(buffer) - 1);
                        if (count > 0)
                        {
                            buffer[count] = '\0';

                            // Submit each chunk of data to the business thread pool
                            pthread_mutex_lock(&worker_queue_lock);
                            worker_queue[worker_task_count].client_fd = fd;
                            worker_queue[worker_task_count].data = strdup(buffer);
                            worker_task_count++;
                            pthread_cond_signal(&worker_queue_cond);
                            pthread_mutex_unlock(&worker_queue_lock);
                        }
                        else if (count == -1)
                        {
                            if (errno == EAGAIN || errno == EWOULDBLOCK)
                            {
                                // All available data has been read; it is safe to leave the loop
                                break;
                            }
                            else
                            {
                                // A genuine error occurred
                                close_connection(fd);
                                break;
                            }
                        }
                        else if (count == 0)
                        {
                            // The client disconnected
                            close_connection(fd);
                            break;
                        }
                    }
                }

                // [ET fix] Writable event: send until finished or EAGAIN is encountered
                if (conn && (events[i].events & EPOLLOUT))
                {
                    pthread_mutex_lock(&conn->lock);

                    while (conn->out_len > 0)
                    {
                        ssize_t sent = send(fd, conn->out_buffer, conn->out_len, 0);
                        if (sent > 0)
                        {
                            conn->out_len -= sent;
                            if (conn->out_len > 0)
                            {
                                // Move the unsent data to the start of the buffer
                                memmove(conn->out_buffer, conn->out_buffer + sent, conn->out_len);
                            }
                        }
                        else if (sent == -1)
                        {
                            if (errno == EAGAIN || errno == EWOULDBLOCK)
                            {
                                // The kernel send buffer is full; stop and wait for the next EPOLLOUT
                                break;
                            }
                            else
                            {
                                perror("send error");
                                close_connection(fd);
                                break;
                            }
                        }
                    }

                    // All data has been sent; stop monitoring EPOLLOUT and retain EPOLLIN
                    if (conn && conn->out_len == 0)
                    {
                        update_epoll(reactor->epoll_fd, fd, EPOLLIN | EPOLLET);
                    }

                    pthread_mutex_unlock(&conn->lock);
                }
            }
        }
    }
    return NULL;
}

// --- 5. Main thread (dedicated to accept) ---
int multiple_level_reactor()
{
    // 1. Initialize resources
    for (int i = 0; i < MAX_CLIENTS; i++)
        connections[i] = NULL;

    for (int i = 0; i < WORKER_THREAD_COUNT; i++)
    {
        pthread_t tid;
        pthread_create(&tid, NULL, worker_thread_loop, NULL);
    }

    for (int i = 0; i < SUB_REACTOR_COUNT; i++)
    {
        sub_reactors[i].id = i;
        sub_reactors[i].queue_count = 0;
        pthread_mutex_init(&sub_reactors[i].lock, NULL);
        sub_reactors[i].epoll_fd = epoll_create1(0);
        sub_reactors[i].wakeup_fd = eventfd(0, EFD_NONBLOCK);

        struct epoll_event ev;
        // The Linux Programming Interface, 63.4.6: Edge-Triggered Notification
        // For edge-triggered notifications, specify EPOLLET in ev.events when calling epoll_ctl()
        ev.events = EPOLLIN | EPOLLET;
        ev.data.fd = sub_reactors[i].wakeup_fd;
        epoll_ctl(sub_reactors[i].epoll_fd, EPOLL_CTL_ADD, sub_reactors[i].wakeup_fd, &ev);

        pthread_create(&sub_reactors[i].thread_id, NULL, sub_reactor_loop, &sub_reactors[i]);
    }

    // 2. Create the server socket
    int server_fd = socket(AF_INET, SOCK_STREAM, 0);
    int opt = 1;
    setsockopt(server_fd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt));
    struct sockaddr_in address;
    address.sin_family = AF_INET;
    address.sin_addr.s_addr = INADDR_ANY;
    address.sin_port = htons(PORT);
    bind(server_fd, (struct sockaddr *)&address, sizeof(address));
    listen(server_fd, SOMAXCONN);
    set_nonblocking(server_fd); // Must be nonblocking

    // 3. Main reactor
    int main_epoll_fd = epoll_create1(0);
    struct epoll_event event;
    event.events = EPOLLIN | EPOLLET; // The main reactor also uses ET
    event.data.fd = server_fd;
    epoll_ctl(main_epoll_fd, EPOLL_CTL_ADD, server_fd, &event);

    printf(">>> Starting the final NIO architecture (strict ET mode) <<<\n");
    printf("PORT: %d | Main(1) + Sub(%d) + Worker(%d)\n", PORT, SUB_REACTOR_COUNT, WORKER_THREAD_COUNT);

    int next_reactor = 0;
    struct epoll_event events[MAX_EVENTS];

    while (1)
    {
        int n = epoll_wait(main_epoll_fd, events, MAX_EVENTS, -1);
        for (int i = 0; i < n; i++)
        {
            if (events[i].data.fd == server_fd)
            {
                int client_fd;
                // [ET fix] Accept must drain the queue!
                while ((client_fd = accept(server_fd, NULL, NULL)) > 0)
                {
                    if (client_fd >= MAX_CLIENTS)
                    {
                        printf("Maximum connection count reached; rejecting connection\n");
                        close(client_fd);
                        continue;
                    }

                    Connection *conn = malloc(sizeof(Connection));
                    conn->fd = client_fd;
                    conn->out_len = 0;
                    pthread_mutex_init(&conn->lock, NULL);
                    connections[client_fd] = conn;

                    int target = next_reactor % SUB_REACTOR_COUNT;
                    next_reactor++;

                    SubReactor *reactor = &sub_reactors[target];
                    pthread_mutex_lock(&reactor->lock);
                    reactor->fd_queue[reactor->queue_count++] = client_fd;
                    uint64_t one = 1;
                    write(reactor->wakeup_fd, &one, sizeof(one));
                    pthread_mutex_unlock(&reactor->lock);
                }

                // Handle accept errors
                if (client_fd == -1)
                {
                    if (errno != EAGAIN && errno != EWOULDBLOCK)
                    {
                        perror("accept error");
                    }
                }
            }
        }
    }
    return 0;
}

Notes

With these foundations, the next post will explain efficient multiplexing in Java with java.nio.channels.Selector and examine parts of other frameworks that use the Reactor pattern.


Share this post:

Continue this series

Understanding Java NIO

  1. Understanding Java NIO (Part 1): Blocking I/O Servers in C with Processes, Threads, and Thread Pools
  2. Understanding Java NIO (Part 2): I/O Multiplexing and Reactor Servers in CYou are here
  3. Understanding Java NIO (Part 3): I/O Multiplexing and the Reactor Pattern in Java, with Open-Source Framework Code Analysis

Comments

Questions, corrections, and experiences are welcome. Sign in with GitHub to comment; both language versions share this discussion.

Comments are available on the live site only.