A Multi-Client Chat Server with Boost.Asio: Sessions, Write Queues and Strands [#31-1]

What we are building, and what goes wrong without care

This article builds a small but complete TCP chat server with Boost.Asio: clients connect, send NICK alice, and every line they type is broadcast to everyone else in the room. New joiners get the last 100 messages replayed. The whole server is one file of about 220 lines, and the full listing at the end compiles as-is with C++17.

A chat server is the classic Asio exercise because it forces you to deal with the three problems that make asynchronous networking hard:

  1. Concurrent writes to one socket. A broadcast can target a client that is still busy receiving the previous message.
  2. Shared state touched from many handlers. The member list changes on join and leave while a broadcast is iterating over it.
  3. Object lifetime with no obvious owner. A session is alive as long as its socket has pending operations, and must die cleanly when the peer disappears.

The design below solves each one with a specific tool: a per-session write queue, a strand per object, and shared_from_this() captured in handlers. The companion post #50-1 covers what happens after this works on one machine: horizontal scaling, pub/sub fan-out, presence and reconnection.

Requirements: C++17, Boost 1.70 or newer (the listing only uses APIs available since 1.70: make_strand, async_accept with an executor argument, expires_after).


The shape of the program

There are three classes:

  • Server owns the acceptor and a signal handler for Ctrl+C / SIGTERM.
  • ChatRoom owns the member set and the history. All of its state is touched only on its own strand.
  • Session owns one socket, a read buffer, a write queue and an idle timer. All of its state is touched only on its socket’s strand.

The key idea is that nothing is shared between strands except by posting a message. The room never touches a session’s queue directly; it calls session->deliver(msg), which posts to the session’s strand. A session never touches the member set; it calls room.join(self), which posts to the room’s strand. That gives us thread safety with zero mutexes, and the io_context can run on as many threads as there are cores.


Framing: TCP is a byte stream, not a message stream

The first bug people hit is assuming one async_read_some equals one message. TCP gives you no such guarantee: one send("hello\n") from the client may arrive as hel then lo\n, and two quick sends may arrive glued together as hello\nworld\n. You need an explicit framing rule.

For a line-based chat the rule is “a message ends at \n”, and async_read_until(socket, streambuf, '\n', ...) implements it. Two details matter:

  • async_read_until may read past the delimiter. The extra bytes stay in the streambuf and are returned immediately by the next async_read_until call without touching the socket. That is why the buffer must be a member that survives between reads, and why you should consume exactly one line from it (std::getline on an istream over the streambuf does this) instead of calling buffer.consume(buffer.size()).
  • A streambuf grows without limit by default. A client that sends gigabytes without a newline will happily exhaust your memory. Construct it with a maximum size, asio::streambuf read_buf_(4096), and async_read_until fails with asio::error::not_found once the limit is hit without a delimiter. Treat that as a protocol violation and disconnect.

Telnet and Windows clients send \r\n, so strip a trailing \r too. If you later need binary payloads, switch to length-prefix framing: read exactly 4 bytes with asio::async_read, decode the size, check it against a maximum, then read exactly that many bytes. The delimiter approach breaks as soon as the payload can contain the delimiter.

std::string take_line() {
    std::istream is(&read_buf_);
    std::string line;
    std::getline(is, line);                // consumes up to and including '\n'
    if (!line.empty() && line.back() == '\r') line.pop_back();
    return line;
}

Session lifetime: who owns a connection?

After async_accept hands us a socket, we do this:

std::make_shared<Session>(std::move(socket), room_)->start();

The temporary shared_ptr dies at the semicolon, yet the session lives on. The reason is that start() launches an async read whose handler captures self = shared_from_this(). The pending operation now owns the session. Every handler that continues the conversation (next read, next write) launches another operation with another copy of self. When a handler exits without starting new work, typically after an error, the reference count drops to zero and the session destructs, closing the socket.

Two rules follow:

  • Never call shared_from_this() in the constructor. The object is not yet owned by a shared_ptr at that point; since C++17 the call throws std::bad_weak_ptr (before C++17 it was undefined behavior). That is why the work starts in a separate start().
  • Capture self, not just this. Capturing [this] alone compiles and even works in quick tests, until a peer disconnects while a write is still pending; then the completion handler runs on freed memory. With self in the capture list, the object is guaranteed to exist until the handler finishes.

The room holds shared_ptr<Session> in its member set too, which adds a second owner. That is fine as long as leave always runs: every error path in the session calls shutdown(), and shutdown() posts room_.leave(self).


The write queue: one async_write at a time

This is the most important part of the article. asio::async_write is a composed operation: internally it calls async_write_some repeatedly until the whole buffer is sent. If a broadcast calls async_write for message B while message A is still in flight, the two operations’ partial writes can interleave on the socket, and the client receives something like alihello bob: hice: .... Asio’s documentation states that the program must ensure no other write operation is performed on the stream until the composed operation completes.

The fix is a queue plus a simple invariant: a write is in flight if and only if the queue is non-empty.

void enqueue(Message msg) {
    if (closed_) return;
    if (write_q_.size() >= kMaxQueue) return shutdown();   // slow consumer
    bool idle = write_q_.empty();
    write_q_.push_back(std::move(msg));
    if (idle) write_next();                // start the chain only if nothing is in flight
}

void write_next() {
    asio::async_write(socket_, asio::buffer(*write_q_.front()),
        [self = shared_from_this()](error_code ec, std::size_t) {
            if (ec) return self->on_error(ec);
            self->write_q_.pop_front();    // pop only after the write finished
            if (!self->write_q_.empty()) self->write_next();
        });
}

Details that are easy to get wrong:

  • Pop after completion, not before. asio::buffer(...) does not copy; it points at the string. The string must stay alive and unmodified until the handler runs. Popping it first leaves Asio sending freed memory.
  • Use std::deque, not std::vector. deque::push_back invalidates iterators but keeps references to existing elements valid; vector::push_back can reallocate and move the string the in-flight write points to. Here Message is shared_ptr<const std::string>, so the string lives on the heap regardless, but the rule is worth knowing if you store strings by value.
  • Share one buffer across all recipients. A broadcast to 500 members with std::string by value makes 500 copies. shared_ptr<const std::string> makes one allocation that every session queue references, and const guarantees nobody mutates it while other writes are pending.
  • Bound the queue. A client on a bad mobile link, or a malicious one that never reads, will make its queue grow forever while the rest of the room keeps talking. When the queue reaches kMaxQueue, the server disconnects that client. Dropping messages is the alternative, but in chat a silent gap is usually worse than a reconnect that replays history.

I have debugged the interleaved-write bug in more than one codebase, and the frustrating part is that it never shows up locally. On loopback, every async_write_some completes in one go and the messages never overlap. The corruption only appears under real network latency with a busy room, which is exactly when it is hardest to reproduce. If you see “garbled messages, but only in production”, check for a second async_write before blaming the client.


Strands instead of mutexes

We run io.run() from one thread per core, which means two handlers can execute simultaneously. A strand is an executor that guarantees handlers posted through it never run concurrently and run in the order they were posted.

  • The acceptor is given asio::make_strand(io_) as the executor for new sockets. Every session’s socket therefore lives on its own strand, and because a completion handler runs on the I/O object’s executor by default, all of the session’s read, write and timer handlers are serialized automatically.
  • ChatRoom has its own strand. join, leave, broadcast and close_all do nothing but asio::post(strand_, ...).

Why not a mutex around the member set? You can, but then broadcast would hold the lock while calling deliver on every member, and if deliver ever did anything that tried to take the same lock (for example, a session disconnecting inline and calling leave), you deadlock. With strands every cross-object call is a posted message, so there is no lock to hold. Asio deadlock debugging (#49-3) digs into exactly those lock-in-callback failures.

A subtle point: use asio::post, not asio::dispatch, for these cross-object calls. dispatch may run the function inline if you are already on the target executor. In close_all, the room iterates over members_ and calls m->close() on each. If that chain ever executed leave inline, it would erase from members_ in the middle of the loop and invalidate the iterator. post always queues, so the loop finishes first.

The alternative to strands is one io_context per thread, with each connection pinned to one of them. That removes strand overhead and gives good cache locality, but a room with members on different threads still needs a cross-thread hand-off, and a hot room cannot spread over cores. For a chat server, where the expensive part is fan-out across connections, one shared io_context with many strands is the simpler default.


Disconnects, errors and the idle timeout

When a client closes normally, the pending read completes with asio::error::eof. When it vanishes abruptly (killed process, NAT timeout, RST), you see connection_reset or, on the write side, broken_pipe. When we close the socket ourselves, every pending operation completes with operation_aborted. All of these lead to the same place, shutdown(), but only unexpected ones are worth logging, otherwise normal traffic floods the log:

void on_error(error_code ec) {
    if (ec != asio::error::eof && ec != asio::error::connection_reset &&
        ec != asio::error::operation_aborted)
        std::cerr << "session: " << ec.message() << "\n";
    shutdown();
}

void shutdown() {
    if (closed_) return;                   // read and write errors can both arrive
    closed_ = true;
    error_code ignored;
    idle_.cancel();
    socket_.shutdown(tcp::socket::shutdown_both, ignored);
    socket_.close(ignored);                // pending ops complete with operation_aborted
    room_.leave(shared_from_this());
}

The closed_ flag matters because a dead peer usually makes both the read and the write fail, and each handler calls shutdown(). Without the flag, you would announce “bob left” twice.

TCP does not tell you that a peer has silently disappeared, for example when a laptop lid closes. The read just waits forever. An idle timer solves this: re-arm a steady_timer on every read, and close the session when it fires. Two traps:

  • The timer handler captures self, so a pending timer keeps the session alive. shutdown() must cancel it, or dead sessions linger for the full timeout.
  • expires_after() cancels a pending wait, but if the timer has already expired and its handler is queued, that handler still runs with a success code. Checking idle_.expiry() > now inside the handler detects “I was re-armed after being queued” and avoids killing a session that just sent a message.
void arm_idle() {
    idle_.expires_after(kIdle);            // cancels any pending wait
    idle_.async_wait([self = shared_from_this()](error_code ec) {
        if (ec) return;                    // operation_aborted: re-armed or closed
        if (self->idle_.expiry() > std::chrono::steady_clock::now())
            return;                        // re-armed after this handler was queued
        self->shutdown();
    });
}

My first version of this timer forgot the cancel in shutdown(). Everything looked fine functionally, but memory kept climbing under a reconnect-heavy load test, because every disconnected session sat in memory for five minutes waiting for its timer. Printing a line from the Session destructor is the fastest way to catch this class of bug: if destructors do not appear when clients leave, something is still holding a shared_ptr.


The room: history replay and ordering

void ChatRoom::join(std::shared_ptr<Session> s) {
    asio::post(strand_, [this, s = std::move(s)] {
        for (const auto& m : history_) s->deliver(m);   // replay before announcing
        members_.insert(s);
        auto note = std::make_shared<const std::string>("* " + s->nick() + " joined\n");
        remember(note);
        for (const auto& m : members_)
            if (m != s) m->deliver(note);
    });
}

Because the replay and the insert happen inside one strand handler, no broadcast can sneak in between them. The new member sees history, then live traffic, with no gap and no duplicates. Messages from one sender reach every recipient in the order they were sent, since each session’s reads are serialized and each post to the room strand is ordered. Messages from different senders that arrive at the same moment are ordered once, on the room strand, so every recipient sees the same order.

nick_ is read on the room strand but written on the session strand. That is safe here because it is written exactly once, before room_.join() is posted, and post establishes the happens-before relationship. If you allow renaming, route the change through the room strand as well.


Graceful shutdown

On SIGINT or SIGTERM, the server closes the acceptor, then posts close_all() to the room. Each session shuts down its socket, pending operations fail with operation_aborted, handlers exit without starting new work, and io.run() returns on every thread once no work remains. There is no io.stop() call: stopping the context abandons queued handlers, which is how “left” messages and final log lines get lost. Letting the work drain is what makes the shutdown graceful.

The acceptor and the signal set share a strand, because the signal handler closes the acceptor while an accept handler may be running on another thread.


Full listing

#include <boost/asio.hpp>
#include <algorithm>
#include <chrono>
#include <csignal>
#include <deque>
#include <iostream>
#include <memory>
#include <set>
#include <string>
#include <thread>
#include <vector>

namespace asio = boost::asio;
using asio::ip::tcp;
using boost::system::error_code;

class Session;
using Message = std::shared_ptr<const std::string>;

class ChatRoom {
public:
    explicit ChatRoom(asio::io_context& io) : strand_(asio::make_strand(io)) {}

    void join(std::shared_ptr<Session> s);
    void leave(std::shared_ptr<Session> s);
    void broadcast(std::string text, std::shared_ptr<Session> sender);
    void close_all();

private:
    void remember(const Message& m) {
        history_.push_back(m);
        if (history_.size() > kMaxHistory) history_.pop_front();
    }

    static constexpr std::size_t kMaxHistory = 100;
    asio::strand<asio::io_context::executor_type> strand_;
    std::set<std::shared_ptr<Session>> members_;   // touched only on strand_
    std::deque<Message> history_;                  // touched only on strand_
};

class Session : public std::enable_shared_from_this<Session> {
public:
    Session(tcp::socket socket, ChatRoom& room)
        : socket_(std::move(socket)),
          idle_(socket_.get_executor()),
          room_(room),
          read_buf_(kMaxLine) {}

    void start() { read_nick(); }

    // Thread-safe: may be called from the room's strand or anywhere else.
    void deliver(Message msg) {
        asio::post(socket_.get_executor(),
                   [self = shared_from_this(), msg = std::move(msg)]() mutable {
                       self->enqueue(std::move(msg));
                   });
    }

    void close() {
        asio::post(socket_.get_executor(),
                   [self = shared_from_this()] { self->shutdown(); });
    }

    // Written once before join() is posted, read-only afterwards.
    const std::string& nick() const { return nick_; }

private:
    static constexpr std::size_t kMaxLine = 4096;
    static constexpr std::size_t kMaxQueue = 256;
    static constexpr std::chrono::minutes kIdle{5};

    std::string take_line() {
        std::istream is(&read_buf_);
        std::string line;
        std::getline(is, line);
        if (!line.empty() && line.back() == '\r') line.pop_back();
        return line;
    }

    void read_nick() {
        arm_idle();
        asio::async_read_until(socket_, read_buf_, '\n',
            [self = shared_from_this()](error_code ec, std::size_t) {
                if (ec) return self->on_error(ec);
                std::string line = self->take_line();
                self->nick_ = (line.rfind("NICK ", 0) == 0 && line.size() > 5)
                                  ? line.substr(5, 32)
                                  : "anon";
                self->room_.join(self);
                self->read_line();
            });
    }

    void read_line() {
        arm_idle();
        asio::async_read_until(socket_, read_buf_, '\n',
            [self = shared_from_this()](error_code ec, std::size_t) {
                if (ec) return self->on_error(ec);
                std::string line = self->take_line();
                if (!line.empty())
                    self->room_.broadcast(self->nick_ + ": " + line + "\n", self);
                self->read_line();
            });
    }

    void enqueue(Message msg) {
        if (closed_) return;
        if (write_q_.size() >= kMaxQueue) return shutdown();
        bool idle = write_q_.empty();
        write_q_.push_back(std::move(msg));
        if (idle) write_next();
    }

    void write_next() {
        asio::async_write(socket_, asio::buffer(*write_q_.front()),
            [self = shared_from_this()](error_code ec, std::size_t) {
                if (ec) return self->on_error(ec);
                self->write_q_.pop_front();
                if (!self->write_q_.empty()) self->write_next();
            });
    }

    void arm_idle() {
        idle_.expires_after(kIdle);
        idle_.async_wait([self = shared_from_this()](error_code ec) {
            if (ec) return;
            if (self->idle_.expiry() > std::chrono::steady_clock::now()) return;
            self->shutdown();
        });
    }

    void on_error(error_code ec) {
        if (ec != asio::error::eof && ec != asio::error::connection_reset &&
            ec != asio::error::operation_aborted)
            std::cerr << "session: " << ec.message() << "\n";
        shutdown();
    }

    void shutdown() {
        if (closed_) return;
        closed_ = true;
        error_code ignored;
        idle_.cancel();
        socket_.shutdown(tcp::socket::shutdown_both, ignored);
        socket_.close(ignored);
        room_.leave(shared_from_this());
    }

    tcp::socket socket_;
    asio::steady_timer idle_;
    ChatRoom& room_;
    asio::streambuf read_buf_;
    std::deque<Message> write_q_;
    std::string nick_;
    bool closed_ = false;
};

void ChatRoom::join(std::shared_ptr<Session> s) {
    asio::post(strand_, [this, s = std::move(s)] {
        for (const auto& m : history_) s->deliver(m);
        members_.insert(s);
        auto note = std::make_shared<const std::string>("* " + s->nick() + " joined\n");
        remember(note);
        for (const auto& m : members_)
            if (m != s) m->deliver(note);
    });
}

void ChatRoom::leave(std::shared_ptr<Session> s) {
    asio::post(strand_, [this, s = std::move(s)] {
        if (members_.erase(s) == 0) return;             // never joined, or already gone
        auto note = std::make_shared<const std::string>("* " + s->nick() + " left\n");
        remember(note);
        for (const auto& m : members_) m->deliver(note);
    });
}

void ChatRoom::broadcast(std::string text, std::shared_ptr<Session> sender) {
    auto msg = std::make_shared<const std::string>(std::move(text));
    asio::post(strand_, [this, msg, sender = std::move(sender)] {
        remember(msg);
        for (const auto& m : members_)
            if (m != sender) m->deliver(msg);
    });
}

void ChatRoom::close_all() {
    asio::post(strand_, [this] {
        for (const auto& m : members_) m->close();      // close() posts; never erases inline
    });
}

class Server {
public:
    Server(asio::io_context& io, unsigned short port)
        : io_(io),
          strand_(asio::make_strand(io)),
          acceptor_(strand_, tcp::endpoint(tcp::v4(), port)),
          signals_(strand_, SIGINT, SIGTERM),
          room_(io) {
        signals_.async_wait([this](error_code, int) {
            acceptor_.close();
            room_.close_all();
        });
        accept();
    }

private:
    void accept() {
        acceptor_.async_accept(asio::make_strand(io_),
            [this](error_code ec, tcp::socket socket) {
                if (ec == asio::error::operation_aborted) return;
                if (!ec) std::make_shared<Session>(std::move(socket), room_)->start();
                accept();
            });
    }

    asio::io_context& io_;
    asio::strand<asio::io_context::executor_type> strand_;
    tcp::acceptor acceptor_;
    asio::signal_set signals_;
    ChatRoom room_;
};

int main(int argc, char* argv[]) {
    unsigned short port = argc > 1 ? static_cast<unsigned short>(std::stoi(argv[1])) : 5555;
    asio::io_context io;
    Server server(io, port);

    unsigned n = std::max(1u, std::thread::hardware_concurrency());
    std::vector<std::thread> pool;
    for (unsigned i = 1; i < n; ++i) pool.emplace_back([&io] { io.run(); });
    io.run();
    for (auto& t : pool) t.join();
}

Build and try it with two terminals:

g++ -std=c++17 -O2 chat_server.cpp -o chat_server -pthread   # add -lboost_system on Boost < 1.69
./chat_server 5555

# terminal 2 and 3
nc localhost 5555
NICK alice
hello

Where to go from here

  • TLS. Replace tcp::socket with asio::ssl::stream<tcp::socket> and add an async_handshake before read_nick(). Everything else, including the write queue, stays the same. See SSL/TLS with Asio (#30-2).
  • Authentication and rate limiting. Both belong in the session before room_.join(): validate a token in the first line instead of NICK, and keep a token bucket per session that closes the connection when exceeded.
  • Testing. Build with -fsanitize=thread and hammer the server with a script that connects, sends, and disconnects in a loop. ThreadSanitizer will flag any state you touch outside its strand.
  • More than one machine. A single process with one room strand serializes all broadcasts for that room on one core at a time. Splitting rooms across processes, fanning out through pub/sub and tracking presence is the subject of chat server architecture (#50-1).

Next: REST API server #31-2 Previous: SSL/TLS #30-2