One read is not one request

Framing that splits a byte stream into requests, then a server that handles many clients at once.

  • Part 2
  • beginner
  • about 40 minutes

You will build

framing that splits a byte stream into requests, then a server that handles many clients at once

You will understand

why your Part 1 server was quietly wrong, the accumulate-and-extract pattern every protocol uses, and what a blocked thread costs

  • roughly 35 lines changed
  • Stages 06 and 07

Where we are going

Two requests in one packet get two replies, and a slow client stops freezing everyone else.

$ printf 'hello\nworld\n' | nc localhost 6379
hello
world
flowchart LR
    A[Part 1<br/>echo bytes] --> B[Stage 06<br/>framing]
    B --> C[Stage 07<br/>thread per client]
    C --> D([Part 3<br/>RESP])

Stage 06, framing

Goal. Split the incoming byte stream into individual requests and answer each one.

The idea

Your Part 1 server passed a test it did not deserve to pass. Look at it again.

sent b'first\n'  -> got b'first\n'
sent b'second\n' -> got b'second\n'
sent b'third'    -> got b'third'

Three messages, three replies. Correct, and only because the client waited for a reply before sending the next message. Each message arrived in its own read(). The round trip did the splitting. TCP did nothing.

Take the pause away and it breaks:

$ printf 'hello\nworld\n' | nc localhost 6379
hello
world

That is what you want. What Part 1 gives you is one blob echoed back, because the server saw a single read of twelve bytes and answered once. Two requests in, one reply out. No error, no crash, no log line. Quietly wrong is the worst kind of wrong.

Framing is the rule that says where one request ends. TCP does not provide one. Every protocol built on TCP defines its own.

flowchart TD
    R["read() gives you<br/>whatever arrived"] --> P[append to a<br/>pending buffer]
    P --> Q{a complete<br/>request in there?}
    Q -->|yes| H[extract it, remove it,<br/>handle it] --> Q
    Q -->|no| W[wait for the next read,<br/>keep the leftovers]

Three properties make that correct, and each one is a bug if you get it wrong.

The buffer is per connection. Share it and two clients’ half-arrived requests merge into one corrupted stream, a bug that only appears under concurrent load and looks like a parser fault.

The extraction loop runs until nothing complete remains, because one read may hold several requests.

Incomplete leftovers are kept. A half-arrived request completes on a later read.

The code

Replace the body of the read loop:

// bytes that arrived but don't form a complete request yet
ByteArrayOutputStream pending = new ByteArrayOutputStream();

while ((bytesRead = inputStream.read(buffer)) != -1) {
    pending.write(buffer, 0, bytesRead);

    // one read can carry several requests, or none at all
    byte[] data = pending.toByteArray();
    int start = 0;
    for (int i = 0; i < data.length; i++) {
        if (data[i] != '\n') {
            continue;
        }

        String request = new String(data, start, i - start, StandardCharsets.UTF_8);
        System.out.println("Request: " + request.replace("\r", "\\r"));

        outputStream.write((request + "\n").getBytes(StandardCharsets.UTF_8));
        outputStream.flush();

        start = i + 1;
    }

    // whatever came after the last newline is an unfinished request
    pending.reset();
    pending.write(data, start, data.length - start);
}

Add import java.io.ByteArrayOutputStream;.

pending lives outside the read loop and inside the client loop. That placement is the whole design. Outside the read loop so leftovers survive to the next read(). Inside the client loop so each connection gets a fresh one.

Append before you inspect. Never decide anything about the stream from one read in isolation.

The for loop scans for every \n, not just the first, so hello\nworld\n in one read produces two replies.

start = i + 1 is the line that matters most. Without it start stays 0 forever: every request re-includes all the bytes before it, and the leftover write re-queues everything already answered. A loop that grows and repeats.

i - start as the length excludes the newline itself, so the request is the line content.

Run it

Two requests in one packet, the test Part 1 failed:

$ printf 'hello\nworld\n' | nc localhost 6379
hello
world

One request split across two reads:

$ python3 -c "
import socket, time
s = socket.create_connection(('127.0.0.1', 6379))
s.sendall(b'hel'); time.sleep(0.3); s.sendall(b'lo\n')
print(s.recv(1024))"
b'hello\n'

Trailing bytes with no newline:

$ printf 'hello\nworl' | nc localhost 6379
hello

Notice that

Three arrival patterns, three correct outcomes. In the second the server waited rather than answering half a request. In the third worl was never answered, because it is incomplete and correctly still pending when the client hung up.

The framing rule decided all three. The network decided none of them.

Try it yourself

  1. Delete start = i + 1 and send a\nb\nc\n. Watch the replies repeat and grow. This is the most instructive bug in the post. Put it back once you have seen it.
  2. Move pending outside the client loop, making it a field. Connect two clients and interleave partial requests from each. Their bytes merge into nonsense. Put it back.
  3. Send 10,000 bytes with no newline at all and watch memory grow. That is a one-line denial of service, and Part 3 explains what fixes it.

What usually goes wrong

Replies repeat and get longer. The cursor is not advancing.

Two clients corrupt each other. The pending buffer is shared.

A request is answered twice. You handled it and also left it in the leftovers.

Go deeper: the two ways to frame, and why delimiters lose

Delimiter framing ends a request on a marker byte such as \n. Simple, and the payload can never contain the marker. A client that never sends one grows your buffer forever.

Length-prefix framing declares the size up front, then reads exactly that many bytes. Binary safe, and you can reject an absurd length before allocating for it.

Redis uses length prefixes, which is why SET key <a JPEG> works. Part 3 builds it.


Stage 07, many clients at once

Goal. A second client is served immediately, not after the first leaves.

The idea

Since Part 1 the server has had a limitation I left in deliberately. Here it is in a log:

Client connected: :59871      ← A connects and stays silent
                              ← B connects here, and nothing prints
Client disconnected           ← A finally leaves
Client connected: :59872      ← only now is B served

B’s TCP handshake completed the instant it connected. The OS did that, as Part 1 showed. But accept() and the read loop sit on the same thread, and the read loop can block for hours.

sequenceDiagram
    participant A as Client A
    participant T as Your single thread
    participant B as Client B
    A->>T: connect
    T->>T: accept() returns
    A->>T: (says nothing)
    T->>T: blocked in read()
    B->>T: connect
    Note over B,T: handshake completes,<br/>but nobody calls accept()
    A->>T: disconnect
    T->>T: accept() again
    Note over B,T: B is finally served

The fix is not faster reads. Blocking is cheap. The problem is doing two unrelated waits on one thread.

The code

The accept loop shrinks to three lines:

while (true) {
    Socket clientSocket = serverSocket.accept();
    System.out.println("Client connected: " + clientSocket.getRemoteSocketAddress());

    // hand off and go straight back to accepting
    new Thread(() -> handleClient(clientSocket)).start();
}

Everything that was inside it moves to a method:

private static void handleClient(Socket clientSocket) {
    // try-with-resources so this worker owns closing its own socket
    try (Socket socket = clientSocket) {
        InputStream inputStream = socket.getInputStream();
        OutputStream outputStream = socket.getOutputStream();
        byte[] buffer = new byte[1024];
        int bytesRead;

        ByteArrayOutputStream pending = new ByteArrayOutputStream();

        while ((bytesRead = inputStream.read(buffer)) != -1) {
            // ... the framing loop from Stage 06, unchanged ...
        }

    } catch (IOException e) {
        // one client failing must not take the server down
        System.err.println("Client error: " + e.getMessage());
    }

    System.out.println("Client disconnected");
}

pending stayed a local inside the method, so each thread gets its own. Make it a field and two clients corrupt each other, under load only, which is the worst way to find out.

The worker owns the socket’s lifetime. main cannot close it, because main is already off accepting someone else. Reassigning to a local inside the method is only to satisfy try-with-resources, which requires an effectively final variable.

The catch is inside the worker. A client that dies mid-request throws there. Uncaught, it kills that thread with a stack trace and no cleanup. Caught, one connection ends and the server never notices.

Run it

Client connected: /127.0.0.1:61344     ← A connects, stays silent
Client connected: /127.0.0.1:61345     ← B accepted while A is idle
Request: from-b                        ← B served first
Request: from-a                        ← A served after, out of order

Five at once:

for i in 1 2 3 4 5; do (printf "client-$i\n" | nc localhost 6379 &) ; done

All five answered, server still listening afterwards.

Notice that

B was accepted while A sat idle, and served first. Compare with the Stage 04 log, where B’s connect line did not appear until A left. accept() no longer waits on anyone.

Try it yourself

  1. Reproduce the old behaviour by calling handleClient(clientSocket) directly instead of starting a thread. Two clients, and you are back to the queue. Undo it.

  2. Make one client sleep for 30 seconds mid-request. Others are unaffected.

  3. Swap in a virtual thread and confirm nothing changes:

    Thread.ofVirtual().start(() -> handleClient(clientSocket));

Why a virtual thread, and why write the other one first

Task 3 deserves its own note, because “nothing changes” is the interesting part.

A platform thread is a real OS thread. Roughly a megabyte of reserved stack, scheduled by the kernel, and blocking one blocks an OS thread. A few thousand is a lot.

A virtual thread (JEP 444, final in Java 21) is scheduled by the JVM onto a small pool of carrier threads. When it blocks on socket I/O the JVM parks it and the carrier goes off to run something else. Millions are affordable.

For a server whose threads spend their entire lives blocked on read(), virtual threads fit better. Write the platform version first so the swap means something.

One thing to avoid: CompletableFuture.runAsync(). It looks like the concise option and runs on the common ForkJoinPool, which is sized for short CPU-bound work. Parking those threads on blocking socket reads starves every other user of the pool. A well-known Redis clone I read for reference makes exactly this mistake.

What usually goes wrong

local variables referenced from a lambda must be final. Assign the socket to a local inside the method rather than reusing the loop variable.

The server dies when one client disconnects abruptly. The catch is outside the worker, or missing.

Two clients see each other’s data. The pending buffer became a field.

Go deeper: what real Redis does instead

Redis is single-threaded. It never blocks a thread on a client, because it never dedicates a thread to one. It uses an event loop, epoll on Linux or kqueue on macOS, which asks the OS which of thousands of sockets are ready and handles only those.

Thread per connection is the right choice for learning. It keeps the blocking code you already understand and makes the concurrency visible. It is also why blocking commands like BLPOP are simple for us in Part 7 and genuinely hard for Redis.


What you built

The transport layer is finished. Your server frames requests, serves many clients at once, and survives a client that vanishes mid-sentence.

It still has no idea what a request means. The framing rule is “up to the next newline”, which is about to stop being good enough, because real Redis values contain newlines.

Checkpoint

  1. Why did Part 1’s three-message test pass, and what was actually doing the framing?
  2. Why must the pending buffer be per connection?
  3. Why does the extraction loop have to run more than once per read?
  4. Why does blocking on read() prevent the next accept(), when the OS already completed that client’s handshake?
  5. Name the two failure modes of delimiter framing.

Resources

Next

Point the real redis-cli at this server and read, byte by byte, what a Redis client actually sends. Then build a parser for it.

Part 3: What redis-cli Actually Sends →