In the previous lesson you wrote CatalogClient, which speaks BTCP/1 perfectly. Its only problem is that it has nobody to talk to: the server was you, typing replies by hand into an nc -l 9090. This lesson builds the other end.

ServerSocket is the class that lets a program wait for connections instead of starting them. It is a complete reversal of the role: the client knows where it is going and takes the initiative; the server does not know who will come or when, it sits listening on a well-known port and reacts.

And this is where module 8 stops being theory. A server that serves one client at a time is useless: while Marta consults the catalogue, Diego waits. The solution —a bounded pool of threads, with names for the logs, a two-phase shutdown and shared state protected by concurrent structures— is exactly what you learned in module 8, applied to the problem it was designed for. The ConcurrentCatalog with its ConcurrentHashMap had been waiting two modules for this moment.

By the end, BiblioTech's catalogue server will be running: it will accept Marta, Diego and Nuria at the same time, validate everything arriving over the network, reply with codes, evict idle clients, limit simultaneous connections and shut down gracefully. And you will be able to talk to it from telnet.

Contents

  1. ServerSocket: listening on a port
  2. The bind, the backlog and the listen address
  3. accept(): the method that hands over connections
  4. BindException: Address already in use and setReuseAddress
  5. The sequential server and the demonstration of its limit
  6. The concurrent server: one thread per connection
  7. The correct solution: ExecutorService
  8. Closing the ServerSocket to unblock the accept
  9. Shared state between connections
  10. BiblioTech: the complete catalogue server
  11. Input validation: never trust the network
  12. Testing it with telnet and with nc
  13. Graceful shutdown with a shutdown hook
  14. Limits of the model and what comes next
  15. Common Mistakes and Tips
  16. Exercises

  1. ServerSocket: listening on a port

A ServerSocket is not a conversation socket. You neither read from it nor write to it. Its only job is to wait for incoming connections and, for each one, manufacture a normal Socket —one of those from 09-02— which is what you actually converse over.

import java.net.ServerSocket;
import java.net.Socket;

// The constructor does the "bind": it reserves port 9090 for this process.
try (ServerSocket server = new ServerSocket(9090)) {

    System.out.println("Listening on port " + server.getLocalPort());

    // accept() BLOCKS until a connection arrives.
    // When it returns, it gives back a Socket ALREADY CONNECTED to the client.
    Socket connection = server.accept();

    System.out.println("Client connected from " + connection.getRemoteSocketAddress());

    // From here on, 'connection' is exactly the Socket of 09-02:
    // getInputStream(), getOutputStream(), setSoTimeout(), close()...
}

The distinction between the two objects is the key to the whole lesson:

ServerSocket Socket (the one accept returns)
What it is for Waiting for connections Conversing with one particular client
How many there are One per service One per connected client
Port The well-known one (9090) The same 9090 on the local side; the client has its own ephemeral one
Streams It has none getInputStream() / getOutputStream()
It is closed When the service stops When each conversation ends
sequenceDiagram
    participant OS as Operating system
    participant SS as ServerSocket (9090)
    participant Srv as Server code
    participant C as Client
    Srv->>SS: new ServerSocket(9090)
    Note over SS,OS: bind: port 9090 is now reserved<br/>listen: the OS starts accepting SYN
    Srv->>SS: accept()
    Note over Srv: BLOCKED waiting
    C->>OS: SYN
    OS->>C: SYN + ACK
    C->>OS: ACK
    Note over OS: Connection already established.<br/>It goes to the pending queue (backlog).
    OS-->>SS: there is a connection in the queue
    SS-->>Srv: returns a connected Socket
    Note over Srv: From here on, normal I/O
    Srv->>C: 200 BIBLIOTECH BTCP/1
    C->>Srv: QUERY 978-0000000001
    Srv->>C: 200 OK ...
    C->>Srv: QUIT
    Srv->>C: 221 BYE
    Srv->>Srv: socket.close()
    Note over Srv: Back to accept() for the next one

Notice one extremely important detail of the diagram: the three-way handshake is completed by the operating system, not by your code. When your accept() returns, the connection had already been established beforehand. That means a client can connect successfully even if your server is busy and slow to call accept() — it will wait in a queue. That queue is the backlog.

  1. The bind, the backlog and the listen address

ServerSocket has several constructors, and their parameters are exactly the three decisions to be made:

// 1. Port. The bare minimum.
new ServerSocket(9090);

// 2. Port + backlog: size of the queue of pending connections.
new ServerSocket(9090, 100);

// 3. Port + backlog + listen address.
new ServerSocket(9090, 100, InetAddress.getByName("127.0.0.1"));

// 4. Unbound, to configure before the bind (we will see it in section 4).
ServerSocket s = new ServerSocket();
s.setReuseAddress(true);
s.bind(new InetSocketAddress(9090), 100);

The backlog

It is the size of the queue where the operating system holds the connections already established that your code has not yet picked up with accept().

                 pending queue (backlog = 3)
  new       -->  [ conn ][ conn ][ conn ]  -->  accept()  -->  your code
  clients                                                      (one at a time)

  If the queue is full, the OS refuses new connections
  and the client receives ConnectException: Connection refused.
  • Default value: 50.
  • If your code serves quickly and returns to accept() right away, the queue almost never fills up.
  • If your code takes a long time per connection (a sequential server), the queue fills up and clients start being refused as if the server did not exist.
  • The operating system may trim your value: on Linux it is capped by net.core.somaxconn. Asking for 10,000 does not guarantee 10,000.

The backlog is a shock absorber for bursts, not a solution to a slow server. A backlog of 1000 on a sequential server only achieves that a thousand clients wait a long time instead of nine hundred and fifty getting a quick error.

The listen address

The third parameter decides through which interfaces connections are accepted, and it is a security decision:

Address Effect
Omitted or 0.0.0.0 Listens on all interfaces: reachable from the network
127.0.0.1 Listens only on loopback: only from the machine itself
192.168.1.50 Listens only through that particular interface

During development, binding to 127.0.0.1 is a good habit: it guarantees that nobody on the network can touch your half-finished server. In production, BiblioTech will listen on 0.0.0.0 so that Marta and Diego can reach it from their desks. We will make it configurable with the Configuration/BusinessRules that already reads bibliotech.properties.

Port 0

try (ServerSocket server = new ServerSocket(0)) {
    int realPort = server.getLocalPort();         // for example, 43127
    System.out.println("The system gave me port " + realPort);
}

Asking for port 0 makes the system assign any free one. It is the standard technique in automated tests: each test brings up its server on a port that is certainly free, without clashing with other tests or with the developer's processes. You will see it in module 11 with JUnit.

  1. accept(): the method that hands over connections

Socket connection = server.accept();

Four things to know about this method:

  1. It blocks indefinitely until a connection is available. A thread parked in accept() consumes no CPU, but it is parked.
  2. It returns an already-connected Socket. There is no need to call connect() or anything else: the conversation can start.
  3. It is interruptible only by closing the ServerSocket. Thread.interrupt() does not unblock it. This detail governs the whole graceful shutdown, and we deal with it in section 8.
  4. It accepts a time limit with serverSocket.setSoTimeout(ms), which makes it throw SocketTimeoutException instead of waiting forever. Just as in 09-02, the ServerSocket remains valid afterwards.
server.setSoTimeout(1000);          // poll once a second
while (running) {
    try {
        Socket connection = server.accept();
        serve(connection);
    } catch (SocketTimeoutException e) {
        // One second with no connections: we take the chance to check
        // the stop flag and go back to waiting.
        continue;
    }
}

This pattern —the same as the setSoTimeout one of 09-02— is an alternative to closing the ServerSocket for shutdown. It is gentler but less immediate: up to a second of delay when stopping. Both techniques are used; in BiblioTech we will use the close, which is instantaneous.

  1. BindException: Address already in use and setReuseAddress

This is the first error you are going to meet, guaranteed, as soon as you run the server twice in a row:

Exception in thread "main" java.net.BindException: Address already in use
	at java.base/sun.nio.ch.Net.bind0(Native Method)
	at java.base/java.net.ServerSocket.bind(ServerSocket.java:395)
	at java.base/java.net.ServerSocket.<init>(ServerSocket.java:262)

It means exactly what it says: something already has that port reserved. There are two very different causes.

Cause 1: another process has it

The most common one is an earlier run of your own server that you did not shut down. It is diagnosed in a second:

ss -tlnp | grep 9090
# LISTEN 0 50 0.0.0.0:9090 0.0.0.0:* users:(("java",pid=4711,fd=7))

kill 4711          # and if it resists, kill -9 4711

On macOS and on systems without ss:

lsof -i :9090
netstat -an | grep 9090

Cause 2: TIME_WAIT

This one is subtler and more frustrating. When a TCP connection closes, the end that closes first leaves the connection in the TIME_WAIT state for a couple of minutes. It is deliberate: it stops delayed packets from that connection confusing a new connection with the same four-tuple.

The effect is that you can stop your server and be unable to start it again for two minutes, even though ss shows no process listening.

The solution is SO_REUSEADDR:

// The correct form: create UNBOUND, configure, and then bind.
ServerSocket server = new ServerSocket();
server.setReuseAddress(true);                         // BEFORE the bind
server.bind(new InetSocketAddress(9090), 100);

setReuseAddress(true) must be set before the bind, and that is why the no-argument constructor has to be used. If you use new ServerSocket(9090) the bind has already happened and calling setReuseAddress afterwards achieves nothing.

Warning. SO_REUSEADDR lets you bind a port that is in TIME_WAIT, but it does not let you bind a port another process is actively listening on. If you still get BindException with setReuseAddress(true), the cause is number 1 and it is time for ss -tlnp.

In BiblioTech, the ServerSocket will always be created with this three-step pattern. It is one of those things that cost nothing and save a lot of frustration.

  1. The sequential server and the demonstration of its limit

Let us start with the simplest thing that works, to see why it is not enough.

package com.nexussoftware.bibliotech.network;

import java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
import java.io.PrintWriter;
import java.net.InetSocketAddress;
import java.net.ServerSocket;
import java.net.Socket;
import java.nio.charset.StandardCharsets;
import java.util.logging.Logger;

/**
 * SEQUENTIAL echo server: it serves one client at a time.
 * It exists to demonstrate the limit of the model, not to be used.
 */
public class SequentialServer {

    private static final Logger LOG = Logger.getLogger(SequentialServer.class.getName());

    public static void main(String[] args) throws IOException {
        ServerSocket server = new ServerSocket();
        server.setReuseAddress(true);
        server.bind(new InetSocketAddress(9090), 50);

        LOG.info("Sequential server listening on 9090");

        try (server) {
            while (true) {
                // ONE connection is accepted...
                Socket connection = server.accept();
                LOG.info(() -> "Connection from " + connection.getRemoteSocketAddress());

                // ...and served IN FULL before returning to accept().
                // While this conversation lasts, nobody else is served.
                try (connection) {
                    connection.setSoTimeout(60_000);
                    serve(connection);
                } catch (IOException e) {
                    LOG.warning("Failure serving the connection: " + e.getMessage());
                }

                LOG.info("Connection finished; back to accept()");
            }
        }
    }

    private static void serve(Socket connection) throws IOException {
        BufferedReader reader = new BufferedReader(
                new InputStreamReader(connection.getInputStream(), StandardCharsets.UTF_8));
        PrintWriter writer = new PrintWriter(
                new java.io.BufferedWriter(
                        new java.io.OutputStreamWriter(connection.getOutputStream(),
                                StandardCharsets.UTF_8)),
                true);

        writer.println("200 ECHO SERVER. Type QUIT to finish.");

        String line;
        while ((line = reader.readLine()) != null) {
            if (line.equalsIgnoreCase("QUIT")) {
                writer.println("221 BYE");
                return;
            }
            writer.println("ECHO: " + line);
        }
        // readLine() has returned null: the client closed without saying goodbye.
    }
}

The demonstration

Start the server and open three terminals with telnet localhost 9090:

TERMINAL 1 (first client)            TERMINAL 2 (second client)
------------------------------       ------------------------------
$ telnet localhost 9090              $ telnet localhost 9090
Trying 127.0.0.1...                  Trying 127.0.0.1...
Connected to localhost.              Connected to localhost.
200 ECHO SERVER. Type QUIT...        _
hello                                (nothing. No greeting.
ECHO: hello                           The TCP connection is made,
                                      but the server has not called
                                      accept() and reads nothing)

The second client connects, because the operating system completes the three-way handshake and leaves it in the backlog. But it receives neither the protocol greeting nor a reply to anything, because the server's single thread is still tied up in the first conversation.

Type QUIT in terminal 1 and watch terminal 2:

TERMINAL 1                           TERMINAL 2
------------------------------       ------------------------------
QUIT                                 200 ECHO SERVER. Type QUIT...
221 BYE                              (now it does! With the thread freed,
Connection closed.                    the accept() picks up this connection)

This is the exact demonstration of the problem. The server log confirms it:

INFO: Sequential server listening on 9090
INFO: Connection from /127.0.0.1:52310
INFO: Connection finished; back to accept()
INFO: Connection from /127.0.0.1:52311

The two connections were established almost at the same time, but the server processed them one after the other.

Why it is unacceptable

Problem Consequence
One slow client blocks everybody If Marta goes to lunch with the telnet open, nobody else can query
The backlog fills up With more than 50 waiting, the rest get ConnectException
A malicious client kills the service Connecting and sending nothing is enough to leave the server useless until the setSoTimeout
It does not use the machine A 16-core server using one

And notice the third point, because it is the most serious one: a single client that connects and goes quiet leaves the service down for 60 seconds. Without setSoTimeout, forever. A sequential server is a server anybody can take down from a terminal.

  1. The concurrent server: one thread per connection

The obvious solution, and the first correct step: for every accepted connection, launch a thread to serve it, and return immediately to accept().

while (true) {
    Socket connection = server.accept();

    // A new thread per connection. The loop returns to accept() instantly.
    Thread thread = new Thread(() -> {
        try (connection) {
            serve(connection);
        } catch (IOException e) {
            LOG.warning("Failure while serving: " + e.getMessage());
        }
    }, "client-" + connection.getPort());
    thread.start();
}

It works. With this, the three telnet terminals are served at the same time. But it has a serious scalability problem, and it is the same one you already saw in 08-05 when we explained why pools exist:

Problem with "one thread per connection" Detail
Memory cost Each platform thread reserves its own stack: 512 KB - 1 MB. A thousand clients are 1 GB in stacks alone
Creation cost Creating a thread costs tens of microseconds; with thousands of short connections, it shows
Scheduling cost With many more threads than cores, the system spends time context-switching instead of working
No limit The worst one. Nothing stops ten thousand connections creating ten thousand threads and bringing the JVM down with OutOfMemoryError: unable to create new native thread

The last point turns this into a denial-of-service vulnerability: an attacker opens connections in a loop and your server kills itself creating threads. It is not theoretical; it is the cheapest attack there is against a server like that.

  1. The correct solution: ExecutorService

This is where module 8 really pays off. A bounded pool solves the four problems at a stroke: the threads are reused (no creation cost), they are a fixed number (bounded memory and reasonable scheduling) and, above all, there is a limit: if more connections arrive than fit, they are queued or refused, but the JVM does not fall over.

package com.nexussoftware.bibliotech.network;

import java.util.concurrent.ExecutorService;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

/**
 * Thread factory with readable names. A thread called "bibliotech-client-3"
 * in a thread dump or in a log line is worth a thousand times more than
 * "pool-1-thread-3". This is 08-02 and 08-05 applied.
 */
public class ServerThreadFactory implements ThreadFactory {

    private final String prefix;
    private final AtomicInteger counter = new AtomicInteger(1);

    public ServerThreadFactory(String prefix) {
        this.prefix = prefix;
    }

    @Override
    public Thread newThread(Runnable r) {
        Thread thread = new Thread(r, prefix + "-" + counter.getAndIncrement());
        // NOT daemon: we want the shutdown to be explicit and orderly,
        // not the JVM killing half-finished conversations when main ends.
        thread.setDaemon(false);
        return thread;
    }
}

ThreadFactory is a functional interface with a single method, newThread, so it can also be written as a lambda (04-05); here it goes as a named class because it needs the counter as state.

Creating the pool, with every parameter explicit:

// Bounded pool with ThreadPoolExecutor: total control over the behaviour.
ThreadPoolExecutor pool = new ThreadPoolExecutor(
        8,                                     // core threads: always alive
        8,                                     // maximum: bounded, no surprises
        60L, TimeUnit.SECONDS,                 // idle time before dying
        new LinkedBlockingQueue<>(100),        // BOUNDED queue of waiting tasks
        new ServerThreadFactory("bibliotech-client"),
        new ThreadPoolExecutor.AbortPolicy()); // what to do if it does not fit

The two decisions that matter most:

The queue must be bounded. Executors.newFixedThreadPool(8) uses an unbounded LinkedBlockingQueue internally. That means that if a hundred thousand connections arrive, all hundred thousand are queued and you run out of memory — you have swapped an OutOfMemoryError from threads for an OutOfMemoryError from queued tasks. A queue of 100 says "I can have 8 conversations in progress and 100 waiting; beyond that, I refuse".

The rejection policy must be a conscious choice. When the pool is full and the queue is too:

Policy What it does When to use it
AbortPolicy (default) Throws RejectedExecutionException When you want to refuse explicitly and tell the client
CallerRunsPolicy The thread that called submit runs it Backpressure: it throttles the acceptor. Dangerous in a server: the accept thread starts serving a client and stops accepting
DiscardPolicy Drops it silently Almost never: silent loss
DiscardOldestPolicy Drops the oldest in the queue Almost never in servers

For a server, AbortPolicy is the right one: you catch the exception, reply to the client with something like 503 SERVER OVERLOADED and close politely. It is infinitely better than accepting and not replying.

try {
    pool.execute(() -> serveConnection(connection));
} catch (RejectedExecutionException e) {
    // The server is overloaded. We tell the client and close.
    // Refusing fast is BETTER than accepting and not replying.
    refusePolitely(connection, "503 SERVER OVERLOADED");
}

The diagram of the concurrent server

graph TD
    A["Main thread<br/>accept() loop"] -->|Socket 1| B["ExecutorService<br/>pool of 8 threads"]
    A -->|Socket 2| B
    A -->|Socket 3| B
    A -->|Socket N| B
    B --> C["client-1<br/>BTCP conversation"]
    B --> D["client-2<br/>BTCP conversation"]
    B --> E["client-3<br/>BTCP conversation"]
    C --> F["ConcurrentCatalog<br/>ConcurrentHashMap"]
    D --> F
    E --> F
    B -->|pool and queue full| G["RejectedExecutionException<br/>503 SERVER OVERLOADED"]

The main thread does nothing but accept. It picks up the connection, hands it to the pool and returns to accept() in microseconds. The whole conversation —which may last minutes— happens in a pool thread. That separation is what makes a server scale.

How many threads

The rule from 08-05 still holds, and here it is especially favourable. A thread serving a connection spends the vast majority of its time blocked waiting for the client to write something. It consumes no CPU. That is why a network server tolerates many more threads than cores:

  • For CPU-intensive work: cores or cores + 1.
  • For work with a lot of I/O waiting (this): between cores × 4 and cores × 20, depending on how much waiting there is.
  • With BTCP/1 and a local network, a few dozen are enough for hundreds of occasional clients.

BiblioTech will use a configurable pool, with a default of 16 threads, read from bibliotech.properties.

  1. Closing the ServerSocket to unblock the accept

Here is a detail that surprises everybody the first time:

accept() does not respond to Thread.interrupt().

A thread blocked in accept() ignores the interruption completely. The interrupt flag is set, but the thread goes on waiting. The cooperative cancellation protocol of 08-02, which works perfectly with sleep, wait and BlockingQueue.take, is of no use here.

The only way to unblock an accept() is to close the ServerSocket from another thread. When you do, the accept() immediately throws a SocketException:

private volatile boolean running = true;
private ServerSocket server;

/** Main loop. Runs in the acceptor thread. */
public void run() {
    while (running) {
        try {
            Socket connection = server.accept();
            pool.execute(() -> serve(connection));

        } catch (SocketException e) {
            // It may be a real failure OR our own close.
            // The 'running' flag tells the two cases apart:
            // without it, a normal shutdown appears in the log as a serious error.
            if (!running) {
                LOG.info("ServerSocket closed: orderly end of the accept loop");
                break;
            }
            LOG.log(Level.SEVERE, "Unexpected failure in accept()", e);
            break;

        } catch (IOException e) {
            // A failure accepting ONE connection must not bring the server down:
            // it is logged and we carry on accepting.
            LOG.log(Level.WARNING, "Error accepting a connection", e);
        }
    }
}

/** Called from another thread (menu, shutdown hook, signal). */
public void stop() {
    running = false;                 // 1. volatile flag: visible from the other thread
    closeQuietly(server);            // 2. this unblocks the accept() instantly
    // 3. and then, the two-phase shutdown of the pool
}

The three elements are all necessary and none is superfluous:

  1. The volatile flag tells an intentional close from a real failure. Without it, every clean shutdown leaves a SEVERE with a stack trace in the log, and you end up ignoring SEVERE entries.
  2. Closing the ServerSocket is what unblocks. The flag alone would do nothing, because the thread is inside accept() and never looks at it again.
  3. The order matters: the flag first, the close afterwards. The other way round there is a window in which accept() throws the exception and running is still true, and the clean shutdown is logged as an error.

The two-phase shutdown of the pool

Closing the ServerSocket prevents new connections, but it does not touch the conversations in progress. For those, the two-phase shutdown of 08-05:

public void stop() {
    running = false;
    closeQuietly(server);

    // PHASE 1: no new tasks are accepted, but the ones in progress finish.
    pool.shutdown();
    try {
        // We give a reasonable deadline for the conversations to end by themselves.
        if (!pool.awaitTermination(10, TimeUnit.SECONDS)) {
            LOG.warning("Conversations still active after 10 s: forcing the close");

            // PHASE 2: interrupt whatever threads are left.
            pool.shutdownNow();

            if (!pool.awaitTermination(5, TimeUnit.SECONDS)) {
                LOG.severe("The pool has not finished; there are stuck threads");
            }
        }
    } catch (InterruptedException e) {
        // We were interrupted while waiting: force it and RESTORE the flag.
        pool.shutdownNow();
        Thread.currentThread().interrupt();    // 08-02: never swallow the interruption
    }
}

An important and often forgotten detail. shutdownNow() interrupts the pool threads, but a thread blocked in socket.read() does not respond to the interruption either — just like accept(). That is why the server sets setSoTimeout on every client connection: it turns the eternal block into a periodic SocketTimeoutException, which is a point at which the thread can check its interrupt flag and leave. Without setSoTimeout, shutdownNow() cannot stop a silent client and the awaitTermination runs out. This is the real, practical reason why the read time limit is not optional in a server.

  1. Shared state between connections

With a pool of 16 threads serving connections, 16 threads access the catalogue at once. This is not an exceptional case: it is the server's normal operation.

And here comes the reward for two modules of work. ConcurrentCatalog, which you wrote in 08-06 with ConcurrentHashMap, CopyOnWriteArrayList and LongAdder, is already prepared for this. Not a single line has to be touched:

// Inside the thread serving Marta, and at the same time inside the one serving Diego,
// and at the same time inside the one serving Nuria:
Material material = catalog.findByIsbn(isbn);          // safe: ConcurrentHashMap
catalog.registerQuery();                               // safe: LongAdder
BiblioTech shared state Structure Why it is safe
Material catalogue ConcurrentHashMap in ConcurrentCatalog Lock-free reads, atomic writes per segment
Query counters LongAdder in BiblioTechStatistics Atomic increments with minimal contention
Loan registry ReadWriteLock in SafeLoanRegistry Invariant across two maps under a single lock
Listener list CopyOnWriteArrayList Many reads, rare writes
Configuration Immutable object loaded at startup Immutable = safe by construction

What you do have to watch out for:

  • Compound operations are still not atomic. if (catalog.isAvailable(isbn)) catalog.markAsLent(isbn) has a race between the two calls: two clients can pass the if at the same time and both lend the same book. The solution, already known from 08-06, is to expose the complete operation as an atomic method of the service (SafeLoanRegistry.lendIfAvailable(...)), not to chain two safe operations.
  • Each connection must have its own local state. The BufferedReader, PrintWriter and conversation variables belong to that conversation. Never instance fields of the server shared between threads. This mistake —using a server field for "the" client's writer— produces replies arriving at the wrong client, and it is a spectacular bug to diagnose.
  • The active-connection counters must be atomic: AtomicInteger.

  1. BiblioTech: the complete catalogue server

Now for real. Everything together.

package com.nexussoftware.bibliotech.network;

import com.nexussoftware.bibliotech.domain.Material;
import com.nexussoftware.bibliotech.service.ConcurrentCatalog;
import com.nexussoftware.bibliotech.service.SafeLoanRegistry;

import java.io.BufferedReader;
import java.io.BufferedWriter;
import java.io.IOException;
import java.io.InputStreamReader;
import java.io.OutputStreamWriter;
import java.io.PrintWriter;
import java.net.InetSocketAddress;
import java.net.ServerSocket;
import java.net.Socket;
import java.net.SocketException;
import java.net.SocketTimeoutException;
import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.logging.Level;
import java.util.logging.Logger;

/**
 * Server for BiblioTech's BTCP/1 protocol.
 *
 * - One acceptor thread that only accepts and delegates.
 * - A BOUNDED pool of threads that serve the conversations (module 8).
 * - Shared state in ConcurrentCatalog (ConcurrentHashMap): it was already safe.
 * - Strict validation of everything arriving over the network.
 * - setSoTimeout per connection to evict idle clients.
 * - Two-phase graceful shutdown.
 */
public class CatalogServer implements AutoCloseable {

    private static final Logger LOG = Logger.getLogger(CatalogServer.class.getName());

    // --- Defensive limits: every value coming from the network is bounded ---
    /** No legitimate BTCP/1 request exceeds this. Prevents memory exhaustion. */
    private static final int MAX_LINE_LENGTH = 512;
    /** Consecutive requests per connection: stops a client monopolising a thread. */
    private static final int MAX_REQUESTS = 1_000;
    /** A BiblioTech ISBN has this shape and nothing else. */
    private static final int MAX_ISBN = 20;
    private static final int MAX_EMPLOYEE = 60;

    private final int port;
    private final String bindAddress;
    private final int threads;
    private final int waitQueue;
    private final int idleTimeoutMs;

    private final ConcurrentCatalog catalog;
    private final SafeLoanRegistry loans;

    private ServerSocket server;
    private ThreadPoolExecutor pool;

    /** volatile: written by the stopping thread, read by the acceptor thread (08-04). */
    private volatile boolean running = false;

    private final AtomicInteger activeConnections = new AtomicInteger();
    private final AtomicLong totalConnections = new AtomicLong();
    private final AtomicLong servedRequests = new AtomicLong();
    private final AtomicLong refusedConnections = new AtomicLong();

    public CatalogServer(int port, String bindAddress, int threads,
                         int waitQueue, int idleTimeoutMs,
                         ConcurrentCatalog catalog,
                         SafeLoanRegistry loans) {
        this.port = port;
        this.bindAddress = bindAddress;
        this.threads = threads;
        this.waitQueue = waitQueue;
        this.idleTimeoutMs = idleTimeoutMs;
        this.catalog = catalog;
        this.loans = loans;
    }

    // =================================================================
    // Startup
    // =================================================================

    /** Prepares the socket and the pool. Does not block. */
    public void start() throws IOException {
        // Three-step pattern: create unbound, configure, bind.
        // setReuseAddress MUST come before the bind, or it achieves nothing.
        server = new ServerSocket();
        server.setReuseAddress(true);
        server.bind(new InetSocketAddress(bindAddress, port), waitQueue);

        ThreadFactory factory = new ThreadFactory() {
            private final AtomicInteger n = new AtomicInteger(1);

            @Override
            public Thread newThread(Runnable r) {
                // Readable names: in a thread dump and in every log line
                // you will know exactly who is doing what (08-02).
                Thread t = new Thread(r, "bibliotech-client-" + n.getAndIncrement());
                t.setDaemon(false);
                return t;
            }
        };

        pool = new ThreadPoolExecutor(
                threads, threads,
                60L, TimeUnit.SECONDS,
                // BOUNDED queue: newFixedThreadPool uses an unbounded one and that
                // turns an overload into an OutOfMemoryError.
                new LinkedBlockingQueue<>(waitQueue),
                factory,
                // Refusing fast and telling the client is better than
                // accepting and not replying.
                new ThreadPoolExecutor.AbortPolicy());

        running = true;
        LOG.info(() -> "CatalogServer BTCP/1 listening on "
                + bindAddress + ":" + port
                + " (pool=" + threads + ", queue=" + waitQueue + ")");
    }

    /** Accept loop. Blocks until close() is called. */
    public void run() {
        while (running) {
            Socket connection = null;
            try {
                connection = server.accept();       // BLOCKS here

                // IDLE time limit per connection. It is what evicts the
                // clients that connect and go quiet, and also what lets
                // shutdownNow() actually stop the threads.
                connection.setSoTimeout(idleTimeoutMs);
                connection.setTcpNoDelay(true);     // interactive protocol: no Nagle

                final Socket accepted = connection;
                totalConnections.incrementAndGet();

                // The acceptor thread does NOT serve: it delegates and returns to accept().
                pool.execute(() -> serveConnection(accepted));

            } catch (RejectedExecutionException e) {
                // Pool and queue full: the server is overloaded.
                refusedConnections.incrementAndGet();
                LOG.warning("Server overloaded; connection refused");
                refuse(connection, "503 SERVER OVERLOADED");

            } catch (SocketException e) {
                // Our own close() OR a real failure. The flag decides.
                if (!running) {
                    LOG.info("Accept loop finished in an orderly way");
                    return;
                }
                LOG.log(Level.SEVERE, "Unexpected failure in accept()", e);
                return;

            } catch (IOException e) {
                // A failure accepting ONE connection does not bring the server down.
                LOG.log(Level.WARNING, "Error accepting a connection", e);
                closeQuietly(connection);
            }
        }
    }

    // =================================================================
    // Serving one connection (runs in a pool thread)
    // =================================================================

    private void serveConnection(Socket connection) {
        String client = String.valueOf(connection.getRemoteSocketAddress());
        int active = activeConnections.incrementAndGet();
        LOG.info(() -> "[" + client + "] connected (" + active + " active)");

        // try-with-resources: the socket is closed whatever happens.
        try (connection) {
            BufferedReader reader = new BufferedReader(
                    new InputStreamReader(connection.getInputStream(),
                            StandardCharsets.UTF_8));
            PrintWriter writer = new PrintWriter(
                    new BufferedWriter(
                            new OutputStreamWriter(connection.getOutputStream(),
                                    StandardCharsets.UTF_8)),
                    false);     // no autoFlush: we flush by hand after each reply

            // The server speaks first, as the protocol specifies.
            respond(writer, "200 BIBLIOTECH BTCP/1");

            int requests = 0;
            while (!Thread.currentThread().isInterrupted()) {

                if (requests++ >= MAX_REQUESTS) {
                    // A client cannot keep a thread indefinitely.
                    respond(writer, "421 TOO MANY REQUESTS");
                    LOG.warning("[" + client + "] exceeded the request limit");
                    return;
                }

                String line = readBoundedLine(reader);
                if (line == null) {
                    LOG.info(() -> "[" + client + "] closed the connection");
                    return;
                }

                servedRequests.incrementAndGet();
                boolean carryOn = process(line, writer, client);
                if (!carryOn) {
                    return;
                }
            }

        } catch (SocketTimeoutException e) {
            // The client connected and went quiet. It is evicted: the thread returns to the pool.
            LOG.info(() -> "[" + client + "] evicted for inactivity ("
                    + idleTimeoutMs + " ms)");

        } catch (LineTooLongException e) {
            LOG.warning("[" + client + "] sent a line longer than "
                    + MAX_LINE_LENGTH + " characters; connection cut");

        } catch (SocketException e) {
            // "Connection reset": the client went down. It is normal, not a serious error.
            LOG.info(() -> "[" + client + "] connection lost: " + e.getMessage());

        } catch (IOException e) {
            LOG.log(Level.WARNING, "[" + client + "] I/O failure", e);

        } catch (RuntimeException e) {
            // A programming failure while handling ONE client cannot
            // kill the pool thread or affect the others. This catch
            // is the error boundary of 06-07 applied to each connection.
            LOG.log(Level.SEVERE, "[" + client + "] unexpected failure", e);

        } finally {
            int left = activeConnections.decrementAndGet();
            LOG.info(() -> "[" + client + "] disconnected (" + left + " active)");
        }
    }

    // =================================================================
    // Processing one request
    // =================================================================

    /** Returns false if the connection must be closed after this request. */
    private boolean process(String line, PrintWriter writer, String client) {
        // The command is separated from the rest at the FIRST space: the
        // employee name may contain spaces and we do not want to split it.
        String command;
        String rest;
        int space = line.indexOf(' ');
        if (space < 0) {
            command = line;
            rest = "";
        } else {
            command = line.substring(0, space);
            rest = line.substring(space + 1).strip();
        }

        // The commands are compared in upper case to be tolerant,
        // but the protocol documents them in upper case.
        switch (command.toUpperCase()) {

            case "QUERY" -> {
                if (!validIsbn(rest)) {
                    respond(writer, "400 BAD REQUEST invalid isbn");
                    return true;
                }
                Material material = catalog.findByIsbn(rest);
                if (material == null) {
                    respond(writer, "404 NOT FOUND " + rest);
                } else {
                    respond(writer, "200 OK " + format(material));
                }
                return true;
            }

            case "LIST" -> {
                if (!rest.isEmpty()) {
                    respond(writer, "400 BAD REQUEST LIST takes no arguments");
                    return true;
                }
                List<Material> materials = catalog.all();
                respond(writer, "201 LIST " + materials.size());
                for (Material m : materials) {
                    respond(writer, format(m));
                }
                respond(writer, ".");
                return true;
            }

            case "LEND" -> {
                int sep = rest.indexOf(' ');
                if (sep < 0) {
                    respond(writer, "400 BAD REQUEST missing arguments");
                    return true;
                }
                String isbn = rest.substring(0, sep);
                String employee = rest.substring(sep + 1).strip();

                if (!validIsbn(isbn) || !validEmployee(employee)) {
                    respond(writer, "400 BAD REQUEST invalid arguments");
                    return true;
                }
                if (catalog.findByIsbn(isbn) == null) {
                    respond(writer, "404 NOT FOUND " + isbn);
                    return true;
                }

                // ATOMIC COMPOUND OPERATION: checking availability and lending
                // in a single call to the service. Chaining two safe
                // operations does NOT give a safe operation (08-06).
                boolean lent = loans.lendIfAvailable(isbn, employee);
                if (lent) {
                    LOG.info(() -> "[" + client + "] loan of " + isbn
                            + " to " + employee);
                    respond(writer, "200 OK loan registered");
                } else {
                    respond(writer, "409 NOT AVAILABLE " + isbn);
                }
                return true;
            }

            case "QUIT" -> {
                respond(writer, "221 BYE");
                return false;       // close the connection
            }

            default -> {
                // We do NOT return the command to the client unsanitised:
                // it could contain terminal escape sequences.
                respond(writer, "400 BAD REQUEST unknown command: "
                        + sanitise(command));
                return true;
            }
        }
    }

    private String format(Material m) {
        return m.getIsbn() + "|" + m.getTitle() + "|"
                + m.getType() + "|" + m.isAvailable();
    }

    // =================================================================
    // Low-level protocol I/O
    // =================================================================

    /** Writes one protocol line and FLUSHES. Without flush, the client waits forever. */
    private void respond(PrintWriter writer, String line) {
        writer.print(line + "\n");          // explicit \n: the protocol demands it
        writer.flush();                     // MANDATORY
        if (writer.checkError()) {
            // PrintWriter swallows IOException; checkError is the only clue.
            LOG.fine("Failed to write the reply; the client probably closed");
        }
    }

    /** Internal exception for a line that exceeds the limit. */
    private static class LineTooLongException extends IOException {
        LineTooLongException(String m) {
            super(m);
        }
    }

    /**
     * Reads a BOUNDED line. BufferedReader.readLine() has no limit:
     * a malicious client can send 100 MB with not one line break
     * and exhaust the server's memory. This is a mandatory defence
     * in any exposed server.
     */
    private String readBoundedLine(BufferedReader reader) throws IOException {
        StringBuilder sb = new StringBuilder();
        int c;
        while ((c = reader.read()) != -1) {
            if (c == '\n') {
                // We also accept \r\n by dropping the trailing carriage return.
                int end = sb.length();
                if (end > 0 && sb.charAt(end - 1) == '\r') {
                    sb.setLength(end - 1);
                }
                return sb.toString();
            }
            if (sb.length() >= MAX_LINE_LENGTH) {
                throw new LineTooLongException(
                        "Line longer than " + MAX_LINE_LENGTH + " characters");
            }
            sb.append((char) c);
        }
        // End of stream. Anything half-written is discarded: it is not a valid line.
        return null;
    }

    // =================================================================
    // Validation: NEVER trust what arrives over the network
    // =================================================================

    private boolean validIsbn(String isbn) {
        if (isbn == null || isbn.isBlank() || isbn.length() > MAX_ISBN) {
            return false;
        }
        // ALLOW list of characters: digits and hyphens only. Everything else out.
        // An allow list is always safer than a deny list.
        for (int i = 0; i < isbn.length(); i++) {
            char c = isbn.charAt(i);
            if (!Character.isDigit(c) && c != '-') {
                return false;
            }
        }
        return true;
    }

    private boolean validEmployee(String employee) {
        if (employee == null || employee.isBlank() || employee.length() > MAX_EMPLOYEE) {
            return false;
        }
        for (int i = 0; i < employee.length(); i++) {
            char c = employee.charAt(i);
            // Letters, digits, spaces, underscores and hyphens. No
            // control characters, which would break the protocol or the log.
            if (!Character.isLetterOrDigit(c) && c != ' ' && c != '_' && c != '-') {
                return false;
            }
        }
        return true;
    }

    /**
     * Sanitises a text coming from the network before returning or logging it.
     * Without this, a client can inject \n (a fake reply in the
     * protocol) or ANSI escape sequences (which manipulate the terminal
     * of whoever reads the log). This is LOG INJECTION, a real problem.
     */
    private String sanitise(String text) {
        StringBuilder sb = new StringBuilder();
        int limit = Math.min(text.length(), 40);
        for (int i = 0; i < limit; i++) {
            char c = text.charAt(i);
            sb.append(Character.isISOControl(c) ? '?' : c);
        }
        return sb.toString();
    }

    // =================================================================
    // Graceful shutdown
    // =================================================================

    @Override
    public void close() {
        if (!running) {
            return;
        }
        LOG.info("Shutting down CatalogServer...");

        // 1. Flag BEFORE the close: it tells a clean shutdown from a failure.
        running = false;

        // 2. Closing the ServerSocket is the ONLY thing that unblocks accept().
        //    Thread.interrupt() cannot do it.
        closeQuietly(server);

        // 3. Two-phase shutdown of the pool (08-05).
        if (pool != null) {
            pool.shutdown();        // no new tasks; the current ones carry on
            try {
                if (!pool.awaitTermination(10, TimeUnit.SECONDS)) {
                    LOG.warning("Conversations active after 10 s: forcing the close");
                    pool.shutdownNow();     // interrupts the threads
                    if (!pool.awaitTermination(5, TimeUnit.SECONDS)) {
                        LOG.severe("There are threads left unfinished");
                    }
                }
            } catch (InterruptedException e) {
                pool.shutdownNow();
                Thread.currentThread().interrupt();   // 08-02: restore the flag
            }
        }

        LOG.info(() -> String.format(
                "CatalogServer stopped. Connections: %d total, %d refused. "
                + "Requests served: %d",
                totalConnections.get(), refusedConnections.get(),
                servedRequests.get()));
    }

    private void refuse(Socket connection, String message) {
        if (connection == null) {
            return;
        }
        // Courtesy: telling the client why it is being closed, instead of
        // cutting it off with no explanation. It costs two lines and saves support.
        try (connection) {
            PrintWriter writer = new PrintWriter(
                    new OutputStreamWriter(connection.getOutputStream(),
                            StandardCharsets.UTF_8));
            writer.print(message + "\n");
            writer.flush();
        } catch (IOException e) {
            LOG.fine("Could not notify the refusal: " + e.getMessage());
        }
    }

    private void closeQuietly(java.io.Closeable resource) {
        if (resource != null) {
            try {
                resource.close();
            } catch (IOException ignored) {
                // While closing there is nothing left to save.
            }
        }
    }

    // --- Metrics for the BiblioTech menu ---

    public int activeConnections() {
        return activeConnections.get();
    }

    public long totalConnections() {
        return totalConnections.get();
    }

    public long servedRequests() {
        return servedRequests.get();
    }
}

The startup class

package com.nexussoftware.bibliotech.presentation;

import com.nexussoftware.bibliotech.persistence.Configuration;
import com.nexussoftware.bibliotech.network.CatalogServer;
import com.nexussoftware.bibliotech.service.ConcurrentCatalog;
import com.nexussoftware.bibliotech.service.SafeLoanRegistry;
import com.nexussoftware.bibliotech.infra.LogConfiguration;

import java.io.IOException;
import java.util.logging.Logger;

/** Startup of BiblioTech's catalogue server. */
public class BiblioTechServerApp {

    private static final Logger LOG = Logger.getLogger(BiblioTechServerApp.class.getName());

    public static void main(String[] args) throws IOException {
        LogConfiguration.initialise();      // the module 6 logger, already configured

        // The configuration comes from bibliotech.properties (07-07).
        Configuration config = Configuration.load();
        int port = config.integer("network.port", 9090);
        String bind = config.text("network.bind", "0.0.0.0");
        int threads = config.integer("network.threads", 16);
        int queue = config.integer("network.queue", 100);
        int idle = config.integer("network.idle.ms", 60_000);

        ConcurrentCatalog catalog = ConcurrentCatalog.loadFromFile();
        SafeLoanRegistry loans = new SafeLoanRegistry(catalog);

        CatalogServer server = new CatalogServer(
                port, bind, threads, queue, idle, catalog, loans);
        server.start();

        // Shutdown hook: it runs on Ctrl+C or on SIGTERM (which is what
        // Docker, systemd or Kubernetes send when stopping a container).
        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            LOG.info("Shutdown signal received");
            server.close();
        }, "bibliotech-shutdown"));

        // Blocks here until the ServerSocket is closed.
        server.run();

        LOG.info("Server finished");
    }
}

  1. Input validation: never trust the network

It is worth pausing on this, because it is the difference between an exercise and a server that can go into production.

Everything arriving over a socket is untrusted input. It does not matter that it is Nexus Software's internal network: that network includes an intern's laptop, a visitor's phone and any compromised machine. A client may be a program other than yours, an old version, or somebody with nc and curiosity.

These are the defences CatalogServer carries and the attack each one stops:

Defence Attack it prevents
MAX_LINE_LENGTH with a bounded read A client sends 100 MB with no \n and exhausts the server's memory
MAX_REQUESTS per connection A client monopolises a pool thread indefinitely
setSoTimeout per connection Clients that connect and go quiet exhaust the pool
Bounded pool and queue Thousands of simultaneous connections bring the JVM down
validIsbn with a character allow list Injection of odd data into the catalogue, the log or the reply
sanitise() before returning or logging Injection of line breaks (fake replies) and of ANSI escapes into the log
Atomic compound operation (lendIfAvailable) Two clients lend the same book simultaneously
catch (RuntimeException) per connection A failure with one client kills the thread and affects the others

Two general principles:

Allow list, never deny list. validIsbn accepts only digits and hyphens. The alternative —rejecting dangerous characters— always falls short, because the list of dangerous characters grows with every new context. Define what you accept and reject the rest.

Sanitise what goes out too. If you return to the client the command it sent, and that command contains \n200 OK fake, you have let it inject a reply into your protocol. And if you write it into the log unsanitised, an attacker can inject fake log lines or escape sequences that manipulate the terminal of the administrator reading it. It is the same kind of problem as SQL injection, applied to another context.

About TLS. BTCP/1 travels in the clear: anybody with access to the network segment can read the queries and the loans. To encrypt it, Java offers SSLServerSocket and SSLSocket (package javax.net.ssl), which are used exactly like ServerSocket and Socket but encrypt the traffic with TLS; the real work is in the certificates and the key stores. Network security is covered in depth in 12-07; here it is enough to know that it exists and that a service with sensitive data must not run in the clear.

  1. Testing it with telnet and with nc

Terminal 1 — start the server:

javac -d classes $(find src -name "*.java")
java -cp classes com.nexussoftware.bibliotech.presentation.BiblioTechServerApp
INFO: CatalogServer BTCP/1 listening on 0.0.0.0:9090 (pool=16, queue=100)

Terminal 2 — connect with telnet and type the requests:

telnet localhost 9090

A complete real session (what you type is marked):

Trying 127.0.0.1...
Connected to localhost.
Escape character is '^]'.
200 BIBLIOTECH BTCP/1
QUERY 978-0000000001                                     <-- YOU TYPE
200 OK 978-0000000001|Effective Java|BOOK|true
QUERY 999                                                <-- YOU TYPE
404 NOT FOUND 999
QUERY ../../etc/passwd                                   <-- YOU TYPE
400 BAD REQUEST invalid isbn
LIST                                                     <-- YOU TYPE
201 LIST 3
978-0000000001|Effective Java|BOOK|true
978-0000000002|Design Patterns|BOOK|false
978-0000000003|Refactoring|BOOK|true
.
LEND 978-0000000003 Nuria Vidal                          <-- YOU TYPE
200 OK loan registered
LEND 978-0000000003 Diego Alonso                         <-- YOU TYPE
409 NOT AVAILABLE 978-0000000003
MAKE ME A COFFEE                                         <-- YOU TYPE
400 BAD REQUEST unknown command: MAKE
QUIT                                                     <-- YOU TYPE
221 BYE
Connection closed by foreign host.

Terminal 3 — at the same time as terminal 2, with nc:

printf 'LIST\nQUIT\n' | nc localhost 9090
200 BIBLIOTECH BTCP/1
201 LIST 3
978-0000000001|Effective Java|BOOK|true
978-0000000002|Design Patterns|BOOK|false
978-0000000003|Refactoring|BOOK|true
.
221 BYE

And the server log with the three simultaneous connections:

INFO: [/127.0.0.1:52440] connected (1 active)
INFO: [/127.0.0.1:52441] connected (2 active)
INFO: [/127.0.0.1:52441] loan of 978-0000000003 to Nuria Vidal
INFO: [/127.0.0.1:52442] connected (3 active)
INFO: [/127.0.0.1:52442] closed the connection
INFO: [/127.0.0.1:52442] disconnected (2 active)
INFO: [/127.0.0.1:52441] disconnected (1 active)
INFO: [/127.0.0.1:52440] disconnected (0 active)

Compare this log with the sequential server's. There the connections were processed one after another; here they overlap. That is the whole difference, and it comes from one line: pool.execute(...) instead of serving in the acceptor thread itself.

Tests worth running

  1. Connect with telnet and type nothing for 60 seconds. The server evicts you: [client] evicted for inactivity (60000 ms). Without setSoTimeout, that telnet would hold a pool thread forever.
  2. Test the 09-02 client against this server. Now that both ends exist, CatalogClientTest should work from start to finish. It is the first time the code from the two lessons meets.
  3. Overload the server. With threads=2 and queue=2 in the configuration, open five telnet sessions at once: the last ones get 503 SERVER OVERLOADED and are closed. Refusing fast and with an explanation is far better than accepting and not replying.
  4. Send a giant line: python3 -c "print('A'*100000)" | nc localhost 9090. The server cuts the connection and logs the attempt, instead of accumulating a hundred thousand characters in memory.
  5. Stop the server with Ctrl+C while there are clients connected. The shutdown hook fires, the accept unblocks, and the log shows the two-phase shutdown and the final statistics.

  1. Graceful shutdown with a shutdown hook

A server is not stopped with System.exit() or by killing the process. A real server receives a signal —SIGTERM from systemd, from Docker or from Kubernetes; SIGINT from a Ctrl+C— and must react by closing properly.

Runtime.getRuntime().addShutdownHook(new Thread(() -> {
    LOG.info("Shutdown signal received");
    server.close();
}, "bibliotech-shutdown"));

A shutdown hook is a thread the JVM starts when it is about to finish. Rules to know:

Rule Detail
It runs on SIGTERM and SIGINT Ctrl+C, kill, docker stop, systemctl stop
It does not run on SIGKILL (kill -9) Nothing can intercept kill -9. That is why important state is persisted, not entrusted to the hook
It must be fast The system usually gives a deadline (10 s in Docker by default) before moving to SIGKILL
It must not call System.exit() It causes a deadlock: the JVM is already exiting
Several hooks run in parallel With no guaranteed order between them
Exceptions it throws are ignored Log them yourself, or they disappear

In BiblioTech, the hook calls GracefulShutdown —the infrastructure class from module 8— which now also includes the server:

Runtime.getRuntime().addShutdownHook(new Thread(() -> {
    LOG.info("Shutting down BiblioTech...");
    server.close();                   // 1. stop accepting and close conversations
    GracefulShutdown.stopServices();  // 2. notice, maintenance and reservation pools
    LoanStore.dump();                 // 3. persist what is pending (module 7)
    LOG.info("BiblioTech stopped correctly");
}, "bibliotech-shutdown"));

The order is deliberate and not interchangeable: first you stop accepting new work, then you finish the work in progress, and finally you persist. The other way round you would persist a state that is still changing.

  1. Limits of the model and what comes next

The model you have built —one platform thread per connection, with a bounded pool— is solid, easy to understand and correct for the vast majority of internal services. But it has a ceiling, and it is worth knowing where.

The C10K problem

With a pool of 16 threads you can serve 16 simultaneous conversations. If the connections are long and mostly idle —think of a chat, or of notifications, where each client keeps the connection open for hours and speaks twice— the model breaks: 10,000 idle clients would require 10,000 threads, that is, several gigabytes in stacks alone.

The classic answer: non-blocking NIO

Java has had a non-blocking I/O API in java.nio.channels since version 1.4:

// Non-blocking I/O: ONE thread serves thousands of connections.
Selector selector = Selector.open();
ServerSocketChannel channel = ServerSocketChannel.open();
channel.bind(new InetSocketAddress(9090));
channel.configureBlocking(false);                        // key point: non-blocking
channel.register(selector, SelectionKey.OP_ACCEPT);

while (true) {
    selector.select();          // blocks until SOME channel has activity
    for (SelectionKey key : selector.selectedKeys()) {
        if (key.isAcceptable()) { /* accept */ }
        if (key.isReadable())   { /* read whatever is there, without blocking */ }
    }
}

A Selector lets a single thread watch thousands of channels and work only on the ones that have data. It is what Netty, Tomcat NIO and practically every high-performance Java server use.

The price is complexity: the code becomes a state machine, each connection needs its own partial buffer because reads return "whatever is there", and debugging it is markedly harder. Do not use it unless you have measured that you need it.

The modern answer: virtual threads

Java 21 changed the calculation completely. Virtual threads are threads managed by the JVM, not by the operating system, that cost a few hundred bytes instead of a megabyte and can be created by the million. When one blocks on I/O, the JVM unmounts it from the system thread and mounts another.

// Java 21+: one VIRTUAL thread per connection. No pool, no practical limit.
try (ExecutorService executor = Executors.newVirtualThreadPerTaskExecutor()) {
    while (running) {
        Socket connection = server.accept();
        executor.submit(() -> serveConnection(connection));
    }
}

That code is the same mental model of one thread per connection, with the usual blocking, readable code, but scaling to hundreds of thousands of connections. It is Java's answer to the C10K problem without paying NIO's price.

Virtual threads are covered in 10-06. For now, keep the idea: the model you have learned does not become obsolete, it becomes even more suitable, because the argument against "one thread per connection" was the cost of the thread, and that cost has disappeared.

Model Connections it supports Complexity When
One thread per connection, no pool Hundreds Low Never in production: no limit
Bounded pool of platform threads Hundreds Low Internal services. What you have done
NIO with a Selector Hundreds of thousands High High performance, if it has been measured
Virtual threads (Java 21+) Hundreds of thousands Low The new default (10-06)

Common Mistakes and Tips

Serving the connection in the accept thread. That is the sequential server: a slow client blocks everybody, and a malicious one takes the service down by connecting and going quiet. The acceptor thread only accepts and delegates.

Using Executors.newFixedThreadPool() without thinking. It carries an unbounded queue: under load, you swap one failure for another. Use ThreadPoolExecutor with a bounded LinkedBlockingQueue and an explicit rejection policy.

Expecting Thread.interrupt() to unblock an accept(). It does not. Only closing the ServerSocket unblocks it. And the same goes for socket.read(): that is why setSoTimeout is not optional in a server.

Calling setReuseAddress after the bind. It has no effect. The ServerSocket must be created with no arguments, configured and bound afterwards.

Not setting setSoTimeout on the accepted connections. Every client that connects and goes quiet keeps a pool thread. With a pool of 16, sixteen idle clients —or sixteen forgotten telnet sessions— leave the server with no capacity, and no error appears in the log.

Using readLine() with no length limit in an exposed server. A hundred megabytes with no line break exhaust memory. Read with a bound.

Sharing the PrintWriter or the BufferedReader between connections. Each conversation has its own, always local variables of the method serving it. An instance field of the server produces replies crossed between clients, and it is a memorable bug to diagnose.

Chaining two safe operations believing the result is safe. if (available) lend() has a race even if both calls are atomic separately. The compound operation must be a method of the service.

Letting an exception kill the pool thread. A catch (RuntimeException) around the serving of each connection is mandatory. Without it, a failure with one particular client takes a pool thread with it and, with enough cases, the server runs out of threads.

Returning or logging what arrived over the network without sanitising it. A line break injects a fake reply into the protocol; an ANSI escape sequence manipulates the terminal of whoever reads the log.

Forgetting the flush() in the server's replies. It is the number-one bug of 09-02, now from the other side: the client waits for a reply that is still in your buffer.

Not closing the client socket when finished. Each socket is a file descriptor, and there is a per-process limit (often 1024 by default). A server that leaks sockets ends up with Too many open files, an error that appears hours after starting and baffles everybody. The try (connection) solves it.

A diagnostic tip. ss -tnp | grep 9090 shows you all the connections live with their state. Many in CLOSE_WAIT means you are not closing your sockets; many in TIME_WAIT is normal after many short connections. It is the network equivalent of watching memory: it tells you whether you are leaking.

Exercises

Exercise 1: BiblioTech statistics server

Add to the server a second port (9092) that serves a status report in plain text, meant for Nexus Software's systems team to query with nc or from a dashboard.

Requirements:

  • Its own ServerSocket on 9092, with its own acceptor thread, sharing the pool with the main server.
  • On connect, with no need to send anything, the server writes the report and closes: active, total and refused connections, requests served, active pool threads, queued tasks, uptime in seconds and memory used.
  • Listen only on 127.0.0.1: the metrics must not be reachable from the network. Explain why in a comment.
  • It must work with nc localhost 9092 and with curl -s telnet://localhost:9092.

Exercise 2: Per-client connection limiter

Add to CatalogServer a limit of simultaneous connections per IP address, so that a single machine cannot take up the whole pool.

Requirements:

  • A ConcurrentHashMap<String, AtomicInteger> counting active connections per IP.
  • A configurable maximum (3 by default) per IP; 127.0.0.1 exempt, so as not to block yourself during testing.
  • When exceeded, reply 429 TOO MANY CONNECTIONS and close, incrementing a counter of limit-based refusals.
  • The counter must always be decremented when the connection ends, even if there was an exception. Think carefully about where that decrement goes.
  • The map entries must be removed when they reach zero, or the map grows indefinitely. Use an atomic operation for the "decrement and perhaps remove" pair, and explain why doing it in two steps would have a race.

Test it by opening four telnet sessions from the same machine with the exemption disabled.

Exercise 3: Server load test

Write a class LoadTest that measures the behaviour of CatalogServer under load using the CatalogClient of 09-02.

Requirements:

  • Parameters: number of concurrent clients, requests per client and host/port.
  • Each client runs as a task in an ExecutorService, connects, makes N queries for random ISBNs from a list, and closes.
  • Collect with atomic structures: successful requests, errors, and latencies in order to compute minimum, mean, median, 95th percentile and maximum.
  • Use a CountDownLatch so that all clients start at once (that way you measure real load, not a ramp) and to wait for them to finish.
  • A final report with requests per second and the latency table.
  • A two-phase shutdown of the test's pool.

Run it with 4, 16 and 64 clients against a 16-thread server and comment on what you observe.

Solutions

Solution 1

package com.nexussoftware.bibliotech.network;

import java.io.BufferedWriter;
import java.io.IOException;
import java.io.OutputStreamWriter;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.net.ServerSocket;
import java.net.Socket;
import java.net.SocketException;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.logging.Level;
import java.util.logging.Logger;

/**
 * BiblioTech statistics server.
 *
 * SECURITY: it listens ONLY on 127.0.0.1. A service's metrics reveal
 * information useful to an attacker (load, saturation, memory, windows of
 * low activity) and are usually the first place something leaks by accident.
 * Exposing them to the network would be a design flaw; if they must be seen
 * from outside, it is done over an SSH tunnel or behind an authenticated proxy (12-07).
 */
public class StatisticsServer implements AutoCloseable {

    private static final Logger LOG = Logger.getLogger(StatisticsServer.class.getName());

    private final int port;
    private final CatalogServer main;
    private final ThreadPoolExecutor mainPool;
    private final ExecutorService pool;
    private final long startInstant = System.currentTimeMillis();

    private ServerSocket server;
    private volatile boolean running = false;

    public StatisticsServer(int port, CatalogServer main,
                            ThreadPoolExecutor mainPool, ExecutorService pool) {
        this.port = port;
        this.main = main;
        this.mainPool = mainPool;
        this.pool = pool;       // the pool is shared: no threads of its own needed
    }

    public void start() throws IOException {
        server = new ServerSocket();
        server.setReuseAddress(true);
        // The third argument of bind TIES the listening to loopback:
        // not even a firewall failure would expose this port.
        server.bind(new InetSocketAddress(
                InetAddress.getLoopbackAddress(), port), 10);
        running = true;

        // Its own acceptor thread, because accept() blocks.
        Thread acceptor = new Thread(this::run, "bibliotech-stats-acceptor");
        acceptor.setDaemon(true);       // this one can be a daemon: it has no state
        acceptor.start();

        LOG.info(() -> "Statistics server on 127.0.0.1:" + port);
    }

    private void run() {
        while (running) {
            try {
                Socket connection = server.accept();
                connection.setSoTimeout(5_000);
                // Delegated to the shared pool: we do not serve in the acceptor.
                pool.execute(() -> serveReport(connection));

            } catch (SocketException e) {
                if (!running) {
                    LOG.info("Statistics server stopped");
                    return;
                }
                LOG.log(Level.SEVERE, "Failure in the statistics accept()", e);
                return;

            } catch (IOException e) {
                LOG.log(Level.WARNING, "Error accepting a statistics connection", e);
            }
        }
    }

    /**
     * This service has no protocol: on connect, it writes and closes.
     * That is what makes it queryable with a simple "nc localhost 9092".
     */
    private void serveReport(Socket connection) {
        try (connection) {
            BufferedWriter writer = new BufferedWriter(
                    new OutputStreamWriter(connection.getOutputStream(),
                            StandardCharsets.UTF_8));
            writer.write(buildReport());
            writer.flush();         // the usual flush
        } catch (IOException e) {
            LOG.fine("Failure serving the report: " + e.getMessage());
        }
    }

    private String buildReport() {
        Runtime rt = Runtime.getRuntime();
        long used = (rt.totalMemory() - rt.freeMemory()) / (1024 * 1024);
        long maximum = rt.maxMemory() / (1024 * 1024);
        long seconds = (System.currentTimeMillis() - startInstant) / 1000;

        StringBuilder sb = new StringBuilder();
        sb.append("=== BIBLIOTECH - SERVER STATUS ===\n");
        sb.append(String.format("%-28s %s%n", "Up for",
                formatDuration(seconds)));
        sb.append('\n');
        sb.append("--- CONNECTIONS ---\n");
        sb.append(String.format("%-28s %d%n", "Active", main.activeConnections()));
        sb.append(String.format("%-28s %d%n", "Total", main.totalConnections()));
        sb.append(String.format("%-28s %d%n", "Requests served",
                main.servedRequests()));
        sb.append('\n');
        sb.append("--- THREAD POOL ---\n");
        sb.append(String.format("%-28s %d%n", "Active threads",
                mainPool.getActiveCount()));
        sb.append(String.format("%-28s %d%n", "Pool size",
                mainPool.getPoolSize()));
        sb.append(String.format("%-28s %d%n", "Queued tasks",
                mainPool.getQueue().size()));
        sb.append(String.format("%-28s %d%n", "Completed tasks",
                mainPool.getCompletedTaskCount()));
        sb.append('\n');
        sb.append("--- MEMORY ---\n");
        sb.append(String.format("%-28s %d MB of %d MB (%.1f%%)%n",
                "Used", used, maximum, 100.0 * used / maximum));
        sb.append(String.format("%-28s %d%n", "Available cores",
                rt.availableProcessors()));
        return sb.toString();
    }

    private String formatDuration(long seconds) {
        // Without java.time (that is 10-05): simple integer arithmetic.
        long h = seconds / 3600;
        long m = (seconds % 3600) / 60;
        long s = seconds % 60;
        return String.format("%dh %02dm %02ds", h, m, s);
    }

    @Override
    public void close() {
        running = false;
        try {
            if (server != null) {
                server.close();         // unblocks the accept()
            }
        } catch (IOException ignored) {
            // While closing there is nothing left to save.
        }
    }
}

Test:

nc localhost 9092
=== BIBLIOTECH - SERVER STATUS ===
Up for                       0h 04m 17s

--- CONNECTIONS ---
Active                       2
Total                        47
Requests served              183

--- THREAD POOL ---
Active threads               2
Pool size                    16
Queued tasks                 0
Completed tasks              45

--- MEMORY ---
Used                         28 MB of 4096 MB (0.7%)
Available cores              8

Comments. Three design decisions. The first and most important: tying the listening to InetAddress.getLoopbackAddress() in the bind. It is not the same as trusting a firewall: the socket does not even exist on the other interfaces, so a network configuration error cannot expose it. The second: the main server's pool is shared instead of creating another; the statistics service handles instantaneous, occasional connections, and giving it threads of its own would waste memory. The third: this service has no protocol, it writes and closes, and that is precisely what makes it queryable with nc and nothing else. Notice too that the statistics acceptor thread is a daemon, unlike the main pool's: it has no state to lose, so there is no problem with the JVM killing it on exit.

Solution 2

package com.nexussoftware.bibliotech.network;

import java.net.Socket;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.logging.Logger;

/**
 * Limiter of simultaneous connections per IP address.
 *
 * Stops a single machine (by mistake or on purpose) taking up the whole
 * thread pool of the server and leaving the other workstations with no service.
 */
public class IpLimiter {

    private static final Logger LOG = Logger.getLogger(IpLimiter.class.getName());

    private final int maxPerIp;
    private final boolean exemptLoopback;

    /** Counts active connections per IP. Concurrent: every thread touches it. */
    private final ConcurrentHashMap<String, AtomicInteger> perIp = new ConcurrentHashMap<>();
    private final AtomicLong refusedByLimit = new AtomicLong();

    public IpLimiter(int maxPerIp, boolean exemptLoopback) {
        this.maxPerIp = maxPerIp;
        this.exemptLoopback = exemptLoopback;
    }

    /**
     * Tries to reserve a slot for this IP.
     * Returns true if admitted; false if the limit has been reached.
     */
    public boolean admit(Socket connection) {
        String ip = ipOf(connection);

        if (exemptLoopback && isLoopback(ip)) {
            // Without the exemption, testing the server with several telnet
            // sessions from the machine itself would be impossible.
            return true;
        }

        // computeIfAbsent is atomic: if two threads arrive at once with the
        // same IP, only one creates the AtomicInteger and both use the same one.
        AtomicInteger counter = perIp.computeIfAbsent(ip, key -> new AtomicInteger());

        // Compare-and-swap loop: the atomic way of saying
        // "increment ONLY IF the current value is lower than the maximum".
        // An "if (get() < max) incrementAndGet()" would have a race:
        // two threads could pass the if at once and exceed the limit.
        while (true) {
            int current = counter.get();
            if (current >= maxPerIp) {
                refusedByLimit.incrementAndGet();
                LOG.warning("Limit of " + maxPerIp
                        + " connections reached for " + ip);
                return false;
            }
            if (counter.compareAndSet(current, current + 1)) {
                return true;    // slot secured
            }
            // compareAndSet failed: another thread changed the value in between.
            // We retry with the new value. This is 08-06.
        }
    }

    /**
     * Releases the slot. It MUST always be called, from a finally.
     */
    public void release(Socket connection) {
        String ip = ipOf(connection);

        if (exemptLoopback && isLoopback(ip)) {
            return;     // nothing was reserved, nothing is released
        }

        // computeIfPresent does the "decrement and, if it reaches zero,
        // remove the entry" pair ATOMICALLY. Returning null from the
        // function removes the key from the map.
        //
        // Doing it in two steps would have a classic race:
        //     if (counter.decrementAndGet() == 0) perIp.remove(ip);
        // Between the decrement and the remove, another thread may increment
        // that same counter to 1 (a new connection from the same IP), and
        // the remove would delete a counter that is already 1. The new
        // connection would be counted in an orphaned object and the limit
        // would stop applying to that IP.
        perIp.computeIfPresent(ip, (key, counter) -> {
            int left = counter.decrementAndGet();
            return left <= 0 ? null : counter;       // null removes the entry
        });
    }

    private String ipOf(Socket connection) {
        return connection.getInetAddress().getHostAddress();
    }

    private boolean isLoopback(String ip) {
        return ip.equals("127.0.0.1") || ip.equals("0:0:0:0:0:0:0:1") || ip.equals("::1");
    }

    public long refusedByLimit() {
        return refusedByLimit.get();
    }

    public int connectedIps() {
        return perIp.size();
    }
}

Integration in the accept loop of CatalogServer:

// In run(), after accepting:
Socket accepted = server.accept();
accepted.setSoTimeout(idleTimeoutMs);
accepted.setTcpNoDelay(true);

if (!limiter.admit(accepted)) {
    refuse(accepted, "429 TOO MANY CONNECTIONS");
    continue;       // it never got delegated to the pool: nothing to release
}

try {
    pool.execute(() -> serveConnection(accepted));
} catch (RejectedExecutionException e) {
    // CAREFUL: the slot was ALREADY reserved. If it is not released here,
    // one slot leaks on every pool refusal and the IP ends up
    // blocked forever.
    limiter.release(accepted);
    refuse(accepted, "503 SERVER OVERLOADED");
}

And in serveConnection, the decrement goes in the finally that already existed:

} finally {
    limiter.release(connection);        // ALWAYS, whatever happens
    int left = activeConnections.decrementAndGet();
    LOG.info(() -> "[" + client + "] disconnected (" + left + " active)");
}

Comments. This exercise concentrates three concurrency lessons from module 8 applied to a real problem.

The compare-and-swap loop is the part most often got wrong. The intuitive code, if (counter.get() < max) counter.incrementAndGet(), has an obvious race: with the limit at 3 and the counter at 2, two threads can read 2, both pass the if and both increment, leaving it at 4. The loop with compareAndSet increments only if the value has not changed since we read it, and retries if it did. It is exactly the pattern of 08-06.

The computeIfPresent to decrement and remove solves a subtler race, explained in the code comment: the pair "decrement" and "remove if zero" must be a single atomic operation, or a new connection from the same IP can slip in between the two and end up counted in a counter that has just been taken out of the map. And removing the entry is not a cosmetic detail: without it, perIp accumulates one entry for every IP that has ever connected, and that is a slow but certain memory leak.

The release in the catch for RejectedExecutionException is the failure that is hardest to find. The slot is reserved before delegating to the pool; if the pool refuses, that slot stays reserved forever because the finally of serveConnection will never run. Three overloads are enough for a workstation to be permanently blocked, and the symptom —"from Marta's computer you cannot get in, but from mine you can"— is one of those that take a whole day of debugging.

Solution 3

package com.nexussoftware.bibliotech.network;

import com.nexussoftware.bibliotech.exception.BiblioTechException;

import java.util.Arrays;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;

/**
 * Load test of CatalogServer using the CatalogClient of 09-02.
 * Measures latencies, error rate and requests per second.
 */
public class LoadTest {

    private static final String[] ISBNS = {
            "978-0000000001",
            "978-0000000002",
            "978-0000000003",
            "000-0000000000"    // non-existent on purpose: must give 404, not an error
    };

    private final String host;
    private final int port;

    private final AtomicInteger successful = new AtomicInteger();
    private final AtomicInteger errors = new AtomicInteger();
    private final AtomicInteger failedConnections = new AtomicInteger();
    private final AtomicLong latencyIndex = new AtomicLong();

    /** Latencies in microseconds. Pre-allocated array: no synchronisation on write. */
    private long[] latencies;

    public LoadTest(String host, int port) {
        this.host = host;
        this.port = port;
    }

    public void run(int clients, int requestsPerClient) throws InterruptedException {
        latencies = new long[clients * requestsPerClient];

        // Pool for the test clients, with named threads.
        ExecutorService pool = Executors.newFixedThreadPool(clients, r -> {
            Thread t = new Thread(r);
            t.setName("load-" + t.getId());
            return t;
        });

        // Two latches: one so that ALL of them start at once (real load,
        // not a ramp), and another to wait for ALL of them to finish (08-05).
        CountDownLatch start = new CountDownLatch(1);
        CountDownLatch finish = new CountDownLatch(clients);

        for (int i = 0; i < clients; i++) {
            final int id = i;
            pool.execute(() -> {
                try {
                    start.await();      // everybody waits here for the starting gun
                    work(id, requestsPerClient);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                } finally {
                    finish.countDown();
                }
            });
        }

        System.out.printf("Launching %d clients x %d requests = %d requests%n%n",
                clients, requestsPerClient, clients * requestsPerClient);

        long t0 = System.nanoTime();
        start.countDown();                      // go!
        finish.await();                         // we wait for everybody
        long durationNs = System.nanoTime() - t0;

        // Two-phase shutdown of the test pool (08-05).
        pool.shutdown();
        if (!pool.awaitTermination(10, TimeUnit.SECONDS)) {
            pool.shutdownNow();
        }

        report(durationNs);
    }

    private void work(int id, int requests) {
        // Each client opens its own connection: CatalogClient is NOT
        // thread-safe, and this way we also measure real connections.
        try (CatalogClient client = new CatalogClient(host, port)) {
            client.connect();

            for (int i = 0; i < requests; i++) {
                String isbn = ISBNS[(id + i) % ISBNS.length];

                long t0 = System.nanoTime();
                try {
                    client.query(isbn);         // null (404) is a success too
                    long us = (System.nanoTime() - t0) / 1000;

                    // Each thread writes into its own array slot:
                    // no collisions, no synchronisation. The index is
                    // handed out by an atomic counter.
                    int pos = (int) latencyIndex.getAndIncrement();
                    if (pos < latencies.length) {
                        latencies[pos] = us;
                    }
                    successful.incrementAndGet();

                } catch (BiblioTechException e) {
                    errors.incrementAndGet();
                }
            }
        } catch (BiblioTechException e) {
            // We could not even connect: counted separately, because it means
            // something different (overload) from an error in one request.
            failedConnections.incrementAndGet();
        }
    }

    private void report(long durationNs) {
        int done = (int) Math.min(latencyIndex.get(), latencies.length);
        if (done == 0) {
            System.out.println("No request was completed.");
            System.out.println("Failed connections: " + failedConnections.get());
            return;
        }

        long[] sorted = Arrays.copyOf(latencies, done);
        Arrays.sort(sorted);        // 05-09: sorting for the percentiles

        long sum = 0;
        for (long v : sorted) {
            sum += v;
        }

        double seconds = durationNs / 1_000_000_000.0;
        double perSecond = successful.get() / seconds;

        System.out.println("=== LOAD TEST RESULT ===");
        System.out.printf("%-26s %.2f s%n", "Total duration", seconds);
        System.out.printf("%-26s %d%n", "Successful requests", successful.get());
        System.out.printf("%-26s %d%n", "Failed requests", errors.get());
        System.out.printf("%-26s %d%n", "Failed connections", failedConnections.get());
        System.out.printf("%-26s %.0f%n", "Requests per second", perSecond);
        System.out.println();
        System.out.println("--- LATENCIES (ms) ---");
        System.out.printf("%-26s %.3f%n", "Minimum", sorted[0] / 1000.0);
        System.out.printf("%-26s %.3f%n", "Mean", (sum / (double) done) / 1000.0);
        System.out.printf("%-26s %.3f%n", "Median (p50)", percentile(sorted, 50) / 1000.0);
        System.out.printf("%-26s %.3f%n", "p95", percentile(sorted, 95) / 1000.0);
        System.out.printf("%-26s %.3f%n", "p99", percentile(sorted, 99) / 1000.0);
        System.out.printf("%-26s %.3f%n", "Maximum",
                sorted[sorted.length - 1] / 1000.0);
    }

    private long percentile(long[] sorted, int p) {
        // Index of the percentile over the already-sorted array.
        int index = (int) Math.ceil(p / 100.0 * sorted.length) - 1;
        return sorted[Math.max(0, Math.min(index, sorted.length - 1))];
    }

    public static void main(String[] args) throws InterruptedException {
        String host = args.length > 0 ? args[0] : "localhost";
        int port = args.length > 1 ? Integer.parseInt(args[1]) : 9090;

        for (int clients : new int[]{4, 16, 64}) {
            System.out.println("############ " + clients + " CLIENTS ############");
            new LoadTest(host, port).run(clients, 200);
            System.out.println();
            Thread.sleep(1000);     // room for the server to recover
        }
    }
}

Output against a server with 16 threads and a queue of 100:

############ 4 CLIENTS ############
Launching 4 clients x 200 requests = 800 requests

=== LOAD TEST RESULT ===
Total duration             0.42 s
Successful requests        800
Failed requests            0
Failed connections         0
Requests per second        1897

--- LATENCIES (ms) ---
Minimum                    0.118
Mean                       0.204
Median (p50)               0.171
p95                        0.392
p99                        0.884
Maximum                    4.113

############ 16 CLIENTS ############
Requests per second        6104
Median (p50)               0.224
p95                        0.668
p99                        1.942
Maximum                    11.208

############ 64 CLIENTS ############
Failed connections         0
Requests per second        6438
Median (p50)               1.102
p95                        4.873
p99                        12.441
Maximum                    58.306

Comments. The numbers tell a very clear story, and it is the same one you will see in any real service.

From 4 to 16 clients, throughput multiplies by 3.2 (1897 → 6104 requests/s) and the median latency barely moves. The server had capacity to spare and is now using it: the pool of 16 threads is being exploited.

From 16 to 64 clients, throughput stalls (6104 → 6438, 5 % more) but the median latency multiplies by five and the p99 by six. This is the canonical behaviour of a saturated system: past the saturation point, more load does not produce more work done, only more queue. The 64 clients share the same 16 threads, so each one waits its turn.

Look at the p99 against the median. With 64 clients, half the requests take less than 1.1 ms, but one in a hundred takes more than 12 ms, and the worst 58 ms. The mean lies and the median is not enough: in a real service, what users perceive as "it is slow" is the p95 and the p99, not the mean. Measuring only the mean is the classic mistake of performance testing.

And the maximum of 4.1 ms with only 4 clients is the JVM warming up: the first requests load classes and run still-interpreted code, before the JIT compiler does its work. Any serious performance test discards an initial warm-up period; otherwise those values contaminate the maximum and the p99.

Conclusion

You have built the other end, and with it BiblioTech has genuinely stopped being alone: there is a server on a Nexus Software machine to which Marta, Diego and Nuria connect at the same time from their desks.

You can tell the two classes involved apart and why they are different: ServerSocket waits, Socket converses. There is one of the first per service and one of the second per connected client, and on the ServerSocket you neither read nor write: its only job is to manufacture connections. You know the three decisions of its creation —the port, the backlog that absorbs bursts but does not fix a slow server, and the listen address, which is a security decision: 0.0.0.0 exposes it to the network and 127.0.0.1 only to the machine itself—. And you know that port 0 makes the system assign you a free one, the standard technique in automated tests.

You understand accept(): it blocks, it returns an already-connected Socket —the three-way handshake was completed by the operating system beforehand, which is why a client can be connected and waiting in the queue even though your code is busy—, it accepts setSoTimeout and, above all, it does not respond to Thread.interrupt(). Closing the ServerSocket is the only thing that unblocks it, and that is the piece that governs the whole shutdown.

You have solved the BindException: Address already in use that you will certainly meet: if it is another process, ss -tlnp exposes it; if it is TIME_WAIT, the solution is setReuseAddress(true) before the bind, which forces the three-step pattern: create unbound, configure, bind.

You have seen live, with three telnet terminals, why a sequential server is unacceptable: the second client connects but does not even get the greeting until the first finishes; and worse, a single client that connects and goes quiet leaves the service down. You have also ruled out "one thread per connection with no limit" for what it really is: a denial-of-service vulnerability, because ten thousand connections create ten thousand threads and bring the JVM down.

And you have applied module 8's solution to the problem it was invented for: a bounded ThreadPoolExecutor, with a ThreadFactory giving readable names to the threads so that logs and dumps are worth something, a bounded queue —because newFixedThreadPool carries an unbounded one and that only swaps one OutOfMemoryError for another— and an explicit rejection policy, with AbortPolicy and a courteous 503 SERVER OVERLOADED, because refusing fast and with an explanation is always better than accepting and not replying. The acceptor thread only accepts and delegates; the whole conversation happens in the pool.

You have mastered the graceful shutdown and its three pieces, none dispensable: the volatile flag that tells an intentional close from a real failure —without it, every clean shutdown leaves a SEVERE in the log—, the close of the ServerSocket that unblocks the accept, and the pool's two-phase shutdown with shutdown, awaitTermination, shutdownNow and restoring the interrupt flag. With the detail that ties it all together: shutdownNow() does not unblock a socket.read() either, and that is why the setSoTimeout on each connection is not a defensive ornament but the mechanism that makes the shutdown work.

You have confirmed that the shared state was already solved: ConcurrentCatalog with its ConcurrentHashMap, BiblioTechStatistics with its LongAdder, SafeLoanRegistry with its ReadWriteLock. Not a line had to be touched. The only thing that still demands care is the usual: compound operations are not atomic just because they chain atomic operations, and that is why lending is a single method of the service and not an if followed by a call.

And you know that everything arriving over a socket is untrusted input, even from the internal network. Your server treats it as such: a bounded line read instead of an unlimited readLine(), a cap on requests per connection, a setSoTimeout that evicts the idle, a bounded pool and queue, validation by allow list —never a deny list—, sanitising of everything going back to the client or to the log to prevent injection of line breaks and terminal escapes, and a catch (RuntimeException) per connection so that a failure with one client does not take a pool thread with it.

BiblioTech gains in this lesson CatalogServer —the complete BTCP/1 server, multi-client, validated, measured and with a graceful shutdown—, BiblioTechServerApp with its shutdown hook, and the exercise classes: StatisticsServer tied to loopback, IpLimiter with its compare-and-swap loop, and LoadTest, which has shown you with numbers the behaviour of a saturated system: past the saturation point, more load does not produce more work done, only more queue, and the mean lies while the p99 tells the truth.

You also know where this model's ceiling is and what lies beyond: NIO with a Selector for the C10K problem, powerful and markedly more complex; and Java 21's virtual threads, which do not invalidate what you have learned but make it even more suitable, because the only argument against "one thread per connection" was the cost of the thread and that cost has disappeared (10-06). And you know that BTCP/1 travels in the clear, that SSLServerSocket is the answer, and that network security is covered thoroughly in 12-07.

In the next lesson, DatagramSocket and DatagramPacket, the transport changes. You are going to let go of every TCP guarantee: no connection, no assured delivery, no ordering, no flow control — and you are going to discover why that is sometimes exactly what you want. You will see the DatagramPacket as an envelope with data, length, address and port, and the DatagramSocket as a letterbox; the two classic mistakes everybody makes with them; what the MTU and fragmentation are, and why 512 bytes is the prudent size; how to tolerate losses, duplicates and disorder, and how to build reliability by hand when it is needed; and broadcast and multicast, which TCP simply cannot do. With two applications for BiblioTech that TCP does not allow: a discovery service in which the workstations shout "WHERE IS BIBLIOTECH?" to the whole local network and the server replies with its address —so that the IP no longer has to be configured on every desk— and a telemetry publisher that sends the number of active loans every few seconds without blocking and without caring whether a packet gets lost on the way.

Java Programming Course

Module 1: Introduction to Java

Module 2: Control Flow

Module 3: Object-Oriented Programming

Module 4: Advanced Object-Oriented Programming

Module 5: Data Structures and Collections

Module 6: Exception Handling

Module 7: File Input/Output

Module 8: Multithreading and Concurrency

Module 9: Networking

Module 10: Advanced Topics

Module 11: Java Frameworks and Libraries

Module 12: Building Real-World Applications

© Copyright 2026. All rights reserved