Chat Server Architecture in C++: Scaling Out with Pub/Sub, Presence, Backpressure and Reconnects [#50-1]
From one process to many
The previous part (#31-1) built a single-process Boost.Asio chat server: one session per socket, a write queue per session, and a room that fans messages out on its own strand. That design is correct, and on one machine it goes a long way. This article is about the moment one machine is no longer enough, or no longer acceptable because a single process is a single point of failure.
As soon as there are two server processes, the in-memory room from #31-1 stops being the truth. Alice is connected to gateway A, Bob to gateway B, and both are in #general. A message Alice sends arrives at A, but A has no socket for Bob. Everything in this article follows from that one fact:
- Fan-out must cross process boundaries.
- Ordering can no longer be decided by one strand.
- Presence (“is Bob online, and where?”) becomes distributed state.
- Failure is partial: one gateway can die while the others keep running, and thousands of clients reconnect at once.
The C++ code in each section is intentionally small and library-agnostic. The hard part of a chat system is choosing where state lives, not writing the socket loop.
Split the server into two tiers
The most common layout separates the process that holds connections from the logic that decides who receives what.
Gateways accept TCP or WebSocket connections, authenticate them, run the per-connection read loop and write queue, and keep a local map of which rooms their connected users are in. They hold no durable state; if a gateway dies, its clients reconnect somewhere else and nothing is lost except the sockets.
The message path assigns each message an order, persists it, and publishes it to every gateway that has a member of that room. This can be a separate “room service”, or simply a shared broker plus a database that the gateways talk to directly.
This split is what lets you scale the two dimensions independently. Connection count is limited by file descriptors, memory for buffers and TLS state, and the per-connection CPU cost of heartbeats; message throughput is limited by the broker and storage. They rarely grow at the same rate.
Put a layer-4 load balancer (or DNS with several addresses) in front of the gateways. Sticky sessions are not needed for long-lived connections, since a connection is pinned to one gateway by nature. What does matter is that the load balancer’s idle timeout is longer than your heartbeat interval, otherwise it will silently cut quiet connections.
Fan-out through a pub/sub broker
The simplest cross-gateway fan-out is a pub/sub channel per room. When Alice sends to #general, gateway A publishes to room:general. Every gateway that has at least one local member of #general is subscribed to that channel, receives the message once, and delivers it to its local sockets.
Two properties make this efficient. Each gateway receives a room message once, regardless of how many local members it has, and inside the gateway the payload is one shared buffer handed to every session’s queue, exactly as in #31-1. The broker’s work scales with the number of gateways per room, not the number of users.
The gateway needs reference-counted subscriptions: subscribe when the first local member joins a room, unsubscribe when the last one leaves.
class PubSub { // thin wrapper over Redis, NATS, ...
public:
virtual ~PubSub() = default;
virtual void publish(const std::string& channel, std::string payload) = 0;
virtual void subscribe(const std::string& channel) = 0;
virtual void unsubscribe(const std::string& channel) = 0;
};
// All methods run on one strand (or one thread) of the gateway.
class LocalRooms {
public:
explicit LocalRooms(PubSub& bus) : bus_(bus) {}
void join(const std::string& room, const std::shared_ptr<Session>& s) {
auto& members = rooms_[room];
if (members.empty()) bus_.subscribe("room:" + room); // first local member
members.push_back(s);
}
void leave(const std::string& room, const Session* s) {
auto it = rooms_.find(room);
if (it == rooms_.end()) return;
auto& v = it->second;
v.erase(std::remove_if(v.begin(), v.end(),
[s](const std::weak_ptr<Session>& w) {
auto p = w.lock();
return !p || p.get() == s;
}),
v.end());
if (v.empty()) { // last local member
bus_.unsubscribe("room:" + room);
rooms_.erase(it);
}
}
// Called for every message the bus delivers to this gateway.
void on_bus_message(const std::string& room, const Message& msg) {
auto it = rooms_.find(room);
if (it == rooms_.end()) return; // late message after unsubscribe
for (const auto& w : it->second)
if (auto s = w.lock()) s->deliver(msg); // one buffer shared by all
}
private:
PubSub& bus_;
std::unordered_map<std::string, std::vector<std::weak_ptr<Session>>> rooms_;
};
LocalRooms holds weak_ptr<Session> on purpose. The session’s lifetime is owned by its pending I/O operations; the room registry should never be the thing that keeps a dead connection alive. A stale entry is harmless: lock() fails and the next leave sweeps it out.
What pub/sub does not give you
Redis Pub/Sub, and core NATS without JetStream, deliver at most once. There is no replay. If a gateway’s connection to the broker blips for two seconds, messages published in those two seconds are gone for that gateway. Redis also protects itself from slow subscribers: its default client-output-buffer-limit for pub/sub clients disconnects a subscriber whose unread output exceeds 32 MB, or stays above 8 MB for 60 seconds. A gateway that stalls under load gets kicked off, and loses everything published until it resubscribes.
There is a quieter version of the same problem. SUBSCRIBE is asynchronous. Between a user joining a room and the broker acknowledging the subscription, messages published to that room do not reach this gateway. On a busy room, the first message or two after join simply never arrive.
The first time I ran a pub/sub-based design with several gateways, both of these showed up as “sometimes a message is missing for one user”, which is about the least actionable bug report there is. The fix is not a better broker setting. It is to stop treating pub/sub as the source of truth: persist every message with a sequence number, use pub/sub only as a low-latency notification, and let clients repair gaps from storage.
Sequence numbers make loss detectable
Give every message a per-room, monotonically increasing sequence number, assigned in exactly one place. Options, from simplest:
- A counter in Redis:
INCR room:general:seq, then store the message and publish it. Cheap, but the increment, store and publish are three steps. If the gateway crashes between them, a sequence number exists that no message ever carries. - An append-only log that assigns IDs itself, such as Redis Streams (
XADDreturns the entry ID) or a Kafka partition per room group. Assignment and persistence become one step, and consumers can resume from an ID. The trade-off is operational weight and, for Kafka, higher end-to-end latency than raw pub/sub. - A room owner process chosen by consistent hashing of the room ID. It orders messages in memory on a strand, like the #31-1 room, and writes them to storage. This keeps the single-writer simplicity, at the cost of handling ownership moves when nodes join or leave.
With sequence numbers, the client (or the gateway on the client’s behalf) can classify every incoming message:
class RoomCursor {
public:
enum class Action { Deliver, Duplicate, Gap };
Action on_message(std::uint64_t seq) {
if (seq <= last_seq_) return Action::Duplicate; // replayed or redelivered
if (seq != last_seq_ + 1) return Action::Gap; // fetch last_seq_+1 .. seq-1 first
last_seq_ = seq;
return Action::Deliver;
}
std::uint64_t last_seq() const { return last_seq_; } // sent on reconnect
private:
std::uint64_t last_seq_ = 0;
};
Duplicates become harmless, which means every other part of the system is allowed to deliver at least once and retry freely. Gaps trigger a history fetch (“give me room general from seq 1042”). The same fetch runs on join and on reconnect, which closes the subscribe race described above: subscribe first, then load history up to the present, then discard any live message the cursor reports as a duplicate.
One detail: a gap caused by a crashed writer will never be filled. Give the gap repair a timeout, and after it, skip forward rather than blocking the room forever.
Persist first or publish first?
If the message path persists before publishing, a message that anyone has seen is guaranteed to be in history, and reconnecting clients never lose it. The cost is that delivery latency includes a storage write.
If it publishes first and persists asynchronously, delivery is faster, but a crash can produce a message that some users saw and that is missing from history, which confuses everyone later.
For chat, persist first is almost always the right default. A storage write on an append-optimized store is short compared to the network latency to users’ devices, and “the message I saw is gone” is a worse bug than an extra few milliseconds. Typing indicators, read receipts and presence updates are the exception: they are ephemeral, can be lost without harm, and should skip storage and go straight to pub/sub.
Presence: eventually consistent by design
Presence answers two questions: is a user online, and which gateway holds their connection (needed for direct messages and for “kick user” admin actions).
The standard pattern is a key per connection with a TTL, refreshed by the gateway:
- On connect:
SET presence:{user}:{conn_id} {gateway_id} EX 60. - Every 20 seconds: the gateway refreshes all its keys in one pipelined batch, not one round trip per user.
- On clean disconnect: delete the key immediately.
- On gateway crash: nothing deletes the keys, so they expire. The user appears online for up to 60 seconds after their gateway died.
That last point is a deliberate trade-off. A short TTL means faster detection but more refresh traffic; a long TTL means stale presence. Keying by connection, not by user, handles multiple devices naturally: the user is online if any presence:{user}:* key exists (in practice, keep a per-user set of connection IDs rather than scanning keys).
Presence changes are more expensive than presence state. If every online and offline transition is broadcast to every member of every room the user is in, a large room generates traffic proportional to members squared. Common mitigations are to send presence only for friends or the currently visible room, to coalesce changes into periodic batches, and to debounce flapping connections so a phone switching from Wi-Fi to mobile data does not announce offline and online within a second.
Backpressure at every hop
A chat system has at least three queues, and each needs a limit and a policy for what happens when it is full.
Per connection (gateway to client). This is the write queue from #31-1. When a client cannot keep up, you either drop low-priority traffic (typing indicators first), or disconnect it and let it recover through history on reconnect. What you must not do is let the queue grow without bound, because one slow phone then consumes memory proportional to the whole room’s traffic.
Per gateway (broker to gateway). The gateway must read from the broker faster than messages arrive, or the broker disconnects it as described above. Keep the broker callback trivial: parse, find the local room, hand off the shared buffer, return. Anything slow, such as permission checks against a database, belongs before publish, not after.
Into the message path (clients to storage). A per-user rate limit (a token bucket at the gateway) keeps one misbehaving client from flooding a room. Reject over-limit messages with an explicit error so honest clients can show “slow down” rather than silently dropping text.
The general rule is to push the rejection as close to the source as possible, and to make every drop visible in metrics. Silent drops in the middle of the pipeline are the ones that turn into week-long investigations.
Heartbeats and dead connections
TCP will not tell you that a client has vanished. A phone that loses coverage, or a laptop that sleeps, leaves the server with a socket that looks open indefinitely. TCP keepalive exists, but on Linux the default idle time before the first probe is two hours, and it only proves the remote kernel is reachable.
Use an application-level heartbeat: the client sends a ping every 20 to 30 seconds (or the server does, and expects a pong), and the gateway closes any connection that has been silent for two or three intervals. The idle timer from #31-1 is exactly this mechanism. Keep the interval below the idle timeout of every load balancer and NAT on the path, since those drop quiet flows without telling either side. WebSocket has ping and pong frames built in, which the WebSocket article (#30-1) covers.
Reconnection without a stampede
Gateways restart: deploys, crashes, kernel updates. Every client on that gateway reconnects, and if all of them retry at fixed intervals they arrive in synchronized waves, each wave re-running authentication, presence writes and history fetches. The result can take down the gateways that were healthy.
Clients should back off exponentially with random jitter. Randomizing over the whole interval, the variant AWS’s architecture blog calls “full jitter”, spreads retries most evenly:
std::chrono::milliseconds reconnect_delay(int attempt, std::mt19937& rng) {
using std::chrono::milliseconds;
const milliseconds base{500};
const milliseconds cap{30'000};
milliseconds ceiling = std::min(cap, base * (1LL << std::min(attempt, 6)));
std::uniform_int_distribution<long long> pick(0, ceiling.count());
return milliseconds{pick(rng)}; // "full jitter"
}
On reconnect the client sends its session token and, per room, RoomCursor::last_seq(). The gateway resubscribes, fetches everything after that sequence, and the user sees no gap.
On the server side, planned restarts should drain: stop accepting new connections, tell connected clients to reconnect (a small control message), and close connections gradually over a window of tens of seconds, not all at once on SIGTERM. Kubernetes and systemd both allow a stop timeout long enough for this; production deployment (#50-5) covers the process side.
Decisions to make up front
Before writing more code, it helps to write down an answer to each of these, because each one changes the design:
| Decision | Common choice | Why |
|---|---|---|
| Delivery guarantee | At least once, with client dedupe by sequence | Retries become safe everywhere |
| Ordering scope | Per room | Global order is expensive and rarely needed |
| Fan-out | Pub/sub channel per room, subscribed by gateways with local members | Broker load scales with gateways, not users |
| Durability | Persist before publish | Messages users saw are always in history |
| Presence | Per-connection keys with a TTL, refreshed in batches | Survives gateway crashes without cleanup code |
| Slow client | Bounded queue, then disconnect | Recovery happens through history on reconnect |
| Reconnect | Exponential backoff with full jitter, resume from last sequence | Avoids reconnect storms and gaps |
None of these require exotic technology. The single-process server from #31-1, a broker, and a store with an append-friendly schema get you a long way, as long as the boundaries between them are explicit.