Interview Setup
Interview Prompt
Design a real-time messaging application like WhatsApp or Slack. Users send one-to-one and group messages with sub-second delivery. The system must support offline message delivery, read receipts, online presence tracking, media sharing, and optional end-to-end encryption.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| What is the expected group size distribution: private small chats or large broadcast channels? | A bounded private group of up to 256 members uses write fan-out, whereas large channels with thousands of members require fan-out on read to prevent quadratic write amplification. |
| Do we require strict message ordering across conversations or only within a single conversation? | Per-conversation total order requires a single sequence authority per conversation ID, while global ordering across disparate chats is unnecessary. |
| Is message delivery online-only or must offline messages be durably buffered? | Offline delivery mandates durable persistence before acknowledging message receipt to prevent data loss during network disconnects. |
| Is end-to-end encryption required, or can the server parse message plaintext? | End-to-end encryption restricts the server to storing opaque ciphertext, shifting full-text search, moderation parsing, and read receipt verification to client devices. |
Scope
In scope
- One-to-one and group messaging with sub-second delivery
- Horizontal scaling of stateful WebSocket connection gateways
- Message persistence, ordering, and offline sync retrieval
- Delivery receipts (sent, delivered, read) and user presence
- High-level end-to-end encryption architecture using the Signal Protocol
Out of scope (state explicitly)
- Live voice and video streaming infrastructure (WebRTC mesh and SFU topologies)
- Deep cryptographic implementation of zero-knowledge proofs
- Enterprise single sign-on and billing integrations
Functional Requirements
Real-time chat architectures balance low-latency full-duplex socket transport with persistent history and offline notification pipelines.
- One-to-One Direct Messaging: Send and receive real-time text and multimedia messages between two users.
- Group Messaging: Support private group conversations with up to 256 participants as well as large public channels.
- Online and Offline Presence: Track real-time connection status (online, offline, and last seen timestamps).
- Delivery and Read Receipts: Track message lifecycle states including sent (durably accepted), delivered to recipient device, and read by recipient.
- Media Sharing: Securely upload, preview, and share images, videos, audio clips, and documents.
- Message History Synchronization: Durably store messages and support paginated historical retrieval across multiple devices using conversation sequence cursors.
- Offline Push Notifications: Alert disconnected recipients through mobile push notifications when new messages arrive without leaking private message content.
- End-to-End Encryption: Optional client-side encryption ensuring the server acts only as an opaque ciphertext relay.
Non-Functional Requirements
Chat systems prioritize message delivery latency, deterministic per-conversation ordering, and continuous availability across transient mobile disconnects.
- Ultra-Low Latency: Deliver messages between online users in under 100 ms at p99.
- High Availability: Maintain 99.99% service availability with graceful degradation during gateway node failures.
- Deterministic Per-Conversation Ordering: Guarantee that messages within any given conversation render in exact sequence order defined by the conversation sequencing authority.
- High Durability (Write-Before-ACK): Persist all incoming messages to disk in Cassandra with local quorum consistency before returning an ingress acknowledgment to the sender.
- Massive Scalability: Support 50 million daily active users and 10 million concurrent WebSocket connections.
- At-Least-Once Delivery with Deduplication: Ensure reliable message delivery while client applications deduplicate via unique message IDs and servers enforce send idempotency via clientMessageId.
Capacity Estimations
Evaluating socket connection density, message throughput, and storage growth dictates gateway fleet sizing and database partitioning.
| Metric | Calculation | Value |
|---|---|---|
| Daily active users (DAU) | Given product scale | 50M |
| Peak concurrent connections | Estimated 20% active peak concurrency | ~10M concurrent WebSockets |
| Daily messages processed | 50M DAU x 10 messages per user | 500M messages / day |
| Average message throughput | 500M messages ÷ 86,400 seconds | ~5,787 messages / sec |
| Peak message throughput (10x) | Peak surge traffic multiplier | ~58,000 messages / sec |
| Average text message payload | Text, metadata, and encryption wrappers | ~200 bytes |
| Daily text storage volume | 500M x 200 bytes | ~100 GB / day (36.5 TB / year) |
| Daily media storage volume (10%) | 50M media messages x 500 KB average | ~25 TB / day (Amazon S3) |
Throughput and Storage Sizing Details
- Connection Gateway Sizing: 10 million concurrent connections handled across gateway servers supporting 100,000 connections each requires roughly 100 stateful chat server instances. Each connection allocates 64 KB of kernel and user-space socket buffers (~6.4 GB RAM per instance).
- Message Throughput: 500 million messages daily averages roughly 5,787 messages per second, with peak traffic reaching approximately 58,000 messages per second during major events.
- Text Storage Growth: 500 million messages at 200 bytes each generates 100 GB of text storage daily, totaling approximately 36.5 TB per year before replication.
- Media Storage Growth: If 10% of messages include media attachments averaging 500 KB, media storage grows by 25 TB daily, managed via Amazon S3 with automated lifecycle transitions to cold storage.
Architecture Diagram
Interview strategy: Separate stateful WebSocket connection termination from stateless business services and the Message Service persistence layer. Explain how the Session Service in Redis coordinates cross-pod message routing when sender and receiver devices terminate on different gateway instances.
The system separates stateful WebSocket termination from message sequencing, business logic, and persistent storage. Each connected device maintains an open WebSocket with an assigned Chat Gateway Server.
When a message arrives, the Chat Gateway delegates persistence to the Message Service, which assigns a monotonic sequence number within that conversation and commits the record to Cassandra with local quorum consistency. Once persisted, the gateway returns a server ACK ({status: "sent", server_message_id, conversation_seq}) to the sender and queries the Session Service in Redis to locate active device sockets for the recipient. If the recipient is connected to another gateway instance, the message routes via the Redis Pub/Sub fast path directly to the recipient's gateway instance. If the delivery signal fails or the socket disconnects, the message is already durable in Cassandra and will be retrieved during client sync. If the recipient has no active sessions, the gateway publishes a wake-up event to the Push Notification Service via Kafka.
Component Deep Dives
Comprehensive breakdown of WebSocket connection lifecycle, multi-device routing, sequence allocation, bounded time-series storage in Cassandra, group fan-out strategies, presence tracking, and end-to-end encryption.
WebSocket Connection Gateway Layer
The gateway fleet manages persistent full-duplex TCP connections, handling socket keep-alives, connection heartbeats, and client reconnection storms.
- Transport Protocol: WebSocket provides full-duplex framing with tiny headers (2 to 14 bytes) compared to HTTP request-response overhead.
- Connection Density: Each instance runs an event-driven networking runtime (Netty or Go epoll/kqueue) managing approximately 100,000 concurrent sockets with 64 KB buffers.
- Load Balancing: L4 load balancers distribute long-lived TCP connections across the gateway fleet, while TLS termination is handled by a TLS-aware edge proxy or gateway ingress, with the Session Service in Redis acting as the authoritative registry for user device locations.
- Session Heartbeats: Clients dispatch lightweight ping frames every 30 seconds. If no heartbeat arrives within 90 seconds, the gateway terminates the connection and removes the device session.
Session Registry Service (Redis)
The Session Service maintains an in-memory directory mapping user IDs to their active device sessions, server identifiers, and socket descriptors.
- Multi-Device Data Model: Stored as
sessions:{user_id}containing a hash of active device records, where each field maps adevice_idto its hostserver_id, connection timestamp, and last heartbeat. - Automatic Expiration: Each device session carries a 5-minute time-to-live (TTL) refreshed periodically by client heartbeats, ensuring disconnected devices clean up automatically without stale routing entries.
Message Router (Cross-Server Communication)
The Message Router bridges isolated gateway instances so that devices connected to different physical servers can communicate with minimal latency.
- Fast Path (Redis Pub/Sub): For online-to-online direct messages, the gateway publishes the payload to a channel named after the recipient's
server_id, achieving sub-5 ms cross-pod routing. Because the message is already safely persisted in Cassandra before routing, any dropped Pub/Sub notification or network glitch is safely recovered by the client via Cassandra sync upon reconnect. - Durable Asynchronous Streams (Kafka): For offline wake-up alerts, large group fan-out, and search indexing pipelines, messages route through partitioned Kafka topics to provide backpressure buffering and replay capabilities. Kafka is strictly used for asynchronous background jobs and not as an ordinary mailbox for direct messages.
Message Persistence & Sequencing Service
The Message Service coordinates message idempotency, assigns deterministic sequence identifiers, and commits records to Cassandra before acknowledging delivery to the sender.
- Identity, Idempotency, and Ordering: The client attaches a
clientMessageIdfor retry idempotency, enabling safe resends without duplicate processing if an ACK is lost. The Message Service assigns a globally unique 64-bit Snowflakemessage_idfor internal server-side identity and correlation. Within the specific conversation, it assigns a monotonically increasingconversation_seqto establish an authoritative, deterministic total order. - Sequence Allocation Mechanism: To prevent race conditions and duplicate sequence numbers across concurrent senders, the Message Service coordinates sequencing per conversation. Writes for a given
conversation_idroute through a conversation partition sequencer (such as an assigned consistent-hash worker node or atomic sequence counter lease) that increments and assigns the next sequential integer atomically before writing to Cassandra. - Write-Before-ACK Guarantee: Messages are written to Cassandra with local quorum consistency before the server returns an acknowledgment frame ({status: "sent", server_message_id, conversation_seq}) to the sender, ensuring durability against server node crashes before confirming receipt.
Cassandra Message Store Architecture & Bucket-Aware Sync
Cassandra serves as the primary time-series message repository, optimized for high-throughput write workloads and chronological range scans within bounded partitions.
- Bounded Partition Key:
((conversation_id, bucket), conversation_seq)groups messages by conversation and sequence range bucket (such asbucket = conversation_seq / 100000). This bounds partition size to under 100 MB even for active groups with millions of messages, preventing hot partition degradation. - Clustering Key:
conversation_seq DESCstores rows physically sorted by descending conversation sequence number, enabling single-seek retrieval of recent chat history. - Bucket-Aware Sync Progression: When a client issues a sync request with its sync cursor (
last_synced_conversation_seq), the Message Service calculates the starting bucket (last_synced_conversation_seq / 100000) and queries Cassandra for rows whereconversation_seq > last_synced_conversation_seqwithin that bucket. If the returned batch is smaller than the requested limit and newer sequence buckets exist, the service queries subsequent buckets until the page size is fulfilled or the latest sequence is reached. - LSM-Tree Storage Engine: Append-only sequential writes handle over 58,000 peak writes per second without random disk seek bottlenecks.
Real-Time User Presence and Heartbeats
The presence subsystem tracks online status and last-seen timestamps without generating broadcast storms across large contact lists.
- Heartbeat Mechanism: Active device sessions maintain
presence:{user_id} = "online"in Redis with a 60-second TTL, refreshed every 30 seconds as long as at least one device session is active. - Throttled Presence Fan-Out: Rather than broadcasting presence changes to all 10,000 contacts of a user, the system broadcasts updates only to mutual contacts in active open chat windows, batching updates every 5 seconds.
Group Message Fan-Out Service
The Group Service manages membership rosters, admin permissions, and distribution strategies for multi-party conversations.
- Private Small Groups (≤ 256 members): Uses fan-out on write (push model). The gateway synchronously queries active device session records in Redis and delivers frames directly to connected participant gateway instances via Redis Pub/Sub fast path.
- Large Public Channels (1,000+ members): For large enterprise channels (such as Slack or Discord), fan-out on write creates quadratic write amplification. The system persists the message once in the channel timeline, and clients fetch updates via fan-out on read (pull model) or subscribe to channel event rooms.
- Group End-to-End Encryption (Sender Keys): Group messaging under E2EE uses Signal-style Sender Keys. Each sender maintains a ratchet-based sender-key state for the group and distributes the symmetric key material to other members through pairwise 1:1 encrypted channels (using Double Ratchet). When group membership changes (a member leaves or is removed), the sender rotates the sender key to preserve backward and forward secrecy.
Media Ingestion and CDN Distribution
Multimedia assets (images, videos, documents) bypass the stateful WebSocket connection to keep socket buffers lightweight and responsive.
- Client-Side Encryption: For E2EE chats, the client generates an ephemeral symmetric key, encrypts the media locally, and uploads the encrypted blob directly to Amazon S3 via a signed HTTP PUT URL.
- Encrypted Thumbnails: The client generates a low-resolution thumbnail preview, encrypts it locally, and embeds the ciphertext thumbnail along with the encrypted media key inside the WebSocket message payload. Recipients download the encrypted blob from the CDN and decrypt it client-side.
Push Notification Integration
When a recipient has no active device sessions, the gateway delegates message alerting to native mobile push notification networks.
- Trigger Condition: If the Session Service reports no active sessions in Redis, the gateway publishes a wake-up event to the
push-notificationsKafka topic. - Privacy and Collapsing: Push notifications convey generic alerts (such as "Alice sent 5 messages") without exposing private message ciphertext or sensitive plaintext to Apple APNs or Google FCM. The durable message remains in Cassandra, ready to be fetched when the recipient launches the app.
End-to-End Encryption (Signal Protocol Architecture)
For privacy-first messaging (such as WhatsApp), the server operates solely as an untrusted ciphertext router and public key directory.
- Key Agreement (X3DH): Users publish pre-key bundles (identity key, signed pre-key, one-time pre-keys) to the E2EE Key Service. Initiators fetch the bundle and derive shared session secrets without server knowledge.
- Ratchet Forward Secrecy & Healing: The Double Ratchet algorithm derives fresh ephemeral encryption keys for every message. This ensures forward secrecy (protecting past messages if current ephemeral keys are compromised) and post-compromise security (healing future session security once a new Diffie-Hellman ratchet exchange occurs).
Kafka Event Streaming Topology
Detailed topic layout, partition keys, and consumer group responsibilities:
Topic: group-fanout
Partitions: 64
Partition key: group_id
Retention: 24 hours
Producer: Group Service when coordinating asynchronous fan-out for groups with many participants
Event Schema: { message_id, conversation_id, conversation_seq, sender_id, payload, timestamp }
Topic: push-notifications
Partitions: 64
Partition key: user_id
Retention: 7 days
Producer: Chat Gateway Server when recipient has no active sessions
Event Schema: { user_id, conversation_id, sender_name, unread_count, collapse_key }
Consumer Groups:
1. group-fanout-workers: Resolves group membership rosters and dispatches to member chat servers via Redis Pub/Sub
2. push-workers: Dispatches batched push notifications to Apple APNs and Google FCM
Dead Letter Queue: push-notifications-dlq after 3 retries
Synchronous Ingress Flow:
Client WebSocket frame -> Chat Gateway validates auth & payload -> Message Service allocates monotonically increasing conversation_seq and persists to Cassandra with local quorum -> Server ACK returned to sender -> Chat Gateway queries Session Service in Redis for active recipient device sockets
Cross-Pod Delivery & Offline Catch-Up Flow:
If recipient is ONLINE on another server -> Publish to recipient server Redis Pub/Sub channel (fast path)
If delivery fails, socket disconnects, or session is stale -> No message is lost; Cassandra holds durable copy
If recipient has no active sessions -> Publish wake-up event to push-notifications Kafka topic
Authoritative recovery -> Offline and reconnecting clients sync missing messages directly from Cassandra by querying where conversation_seq > last_synced_conversation_seq across sequence bucketsAPI Design
WebSocket framing protocol for real-time messaging combined with RESTful APIs for chat history synchronization and group management.
Client WebSocket Type Definitions
TypeScript interfaces defining message frames, server acknowledgments, receipts, and presence events:
export type MessageContentType = "text" | "image" | "video" | "audio" | "document";
export type DeliveryReceiptStatus = "sent" | "delivered" | "read";
export interface SendMessagePayload {
type: "send_message";
clientMessageId: string; // Client-generated UUID for retry idempotency
conversationId: string;
contentType: MessageContentType;
content: string; // Plaintext or E2EE encrypted ciphertext
mediaUrl?: string; // Encrypted media blob URL (if media attachment)
encryptedMediaKey?: string; // Client-encrypted symmetric media key
encryptedThumbnail?: string; // Client-encrypted base64 thumbnail preview
clientTimestamp: number;
}
export interface ServerAckPayload {
type: "ack";
clientMessageId: string;
serverMessageId: string; // Globally unique 64-bit Snowflake ID (server message identity & correlation)
conversationId: string;
conversationSeq: number; // Monotonically increasing sequence within this conversation
status: "sent"; // Confirms durable acceptance in Cassandra (not recipient delivery)
serverTimestamp: number;
}
export interface IncomingMessagePush {
type: "new_message";
messageId: string;
conversationId: string;
conversationSeq: number; // Deterministic sequence number for client ordering
senderId: string;
contentType: MessageContentType;
content: string;
mediaUrl?: string;
encryptedMediaKey?: string;
encryptedThumbnail?: string;
timestamp: number;
}
export interface DeliveryStatusUpdate {
type: "message_status";
conversationId: string;
messageId: string;
conversationSeq: number;
recipientId: string;
deviceId?: string; // Device acknowledging delivery
status: DeliveryReceiptStatus;
timestamp: number;
}
export interface ReadCursorUpdate {
type: "read_cursor";
conversationId: string;
readUpToSeq: number; // Highest conversation sequence number read by the user
timestamp: number;
}
export interface UserPresenceEvent {
type: "presence_change";
userId: string;
status: "online" | "offline";
lastActiveTimestamp?: number;
}
export interface RealtimeChatClient {
sendMessage(payload: SendMessagePayload): Promise<ServerAckPayload>;
sendReceipt(status: DeliveryStatusUpdate): Promise<void>;
updateReadCursor(update: ReadCursorUpdate): Promise<void>;
syncConversation(conversationId: string, afterSeq: number, limit?: number): Promise<IncomingMessagePush[]>;
}Send Message Frame
Dispatched by the sender client over an established WebSocket connection:
{
"type": "send_message",
"client_message_id": "c9a1b8e4-7d2f-4a31-90ef-112233445566",
"conversation_id": "conv_8839102830",
"content_type": "text",
"content": "Hey, are we still meeting at 3 PM?",
"client_timestamp": 1710320000000
}Server Acknowledgment Frame
Returned immediately to the sender once the message is durably committed to Cassandra:
{
"type": "ack",
"client_message_id": "c9a1b8e4-7d2f-4a31-90ef-112233445566",
"server_message_id": "1541815603606036480",
"conversation_id": "conv_8839102830",
"conversation_seq": 48201,
"status": "sent",
"server_timestamp": 1710320000045
}Incoming Message Delivery Frame
Pushed by the recipient's gateway server over their active WebSocket socket:
{
"type": "new_message",
"message_id": "1541815603606036480",
"conversation_id": "conv_8839102830",
"conversation_seq": 48201,
"sender_id": "usr_99201481",
"content_type": "text",
"content": "Hey, are we still meeting at 3 PM?",
"timestamp": 1710320000045
}Delivery and Read Status Updates
Emitted when a recipient device receives or renders a message:
{
"type": "message_status",
"conversation_id": "conv_8839102830",
"message_id": "1541815603606036480",
"conversation_seq": 48201,
"recipient_id": "usr_10293847",
"status": "read",
"timestamp": 1710320001200
}REST Endpoints (History and Group Operations)
Used for catching up on missed messages during reconnect and managing group rosters:
GET /api/v1/conversations/{conversation_id}/messages?after_seq={conversation_seq}&limit=50
Authorization: Bearer <user_token>
HTTP/1.1 200 OK
Content-Type: application/json
{
"messages": [
{
"message_id": "1541815603606036480",
"conversation_id": "conv_8839102830",
"conversation_seq": 48201,
"sender_id": "usr_99201481",
"content": "Hey, are we still meeting at 3 PM?",
"content_type": "text",
"created_at": "2026-03-13T10:00:00.045Z"
}
],
"next_cursor_seq": 48201,
"has_more": false
}
POST /api/v1/groups
Authorization: Bearer <user_token>
Content-Type: application/json
{
"name": "Architecture Working Group",
"member_ids": ["usr_99201481", "usr_10293847", "usr_55667788"]
}
HTTP/1.1 201 Created
Content-Type: application/json
{
"group_id": "grp_3344556677",
"created_at": "2026-03-13T10:05:00.000Z"
}Common Error Responses
400 Bad Request: invalid input, missing required fields, or malformed JSON payload 401 Unauthorized: missing or invalid authentication token or API key 403 Forbidden: authenticated caller lacks required permissions for this resource 404 Not Found: requested resource ID does not exist 409 Conflict: duplicate write or version conflict, retry with a unique idempotency key 422 Unprocessable Entity: syntactically valid request failed semantic business validation 429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header 500 Internal Error: unexpected server failure, retry safely with an idempotency key 503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff 440 Login Timeout: WebSocket connection session expired, client reconnect is required
Data Model
Polyglot data architecture pairing Cassandra for bounded time-series message streams with Redis for transient multi-device session state and MySQL for relational user and group metadata.
Cassandra: Messages Table
Stores immutable message records clustered physically by descending sequence number within bounded partitions:
CREATE TABLE messages (
conversation_id UUID,
bucket INT, -- Sequence bucket (conversation_seq / 100000) to bound partition size
conversation_seq BIGINT, -- Monotonically increasing sequence within this conversation
message_id BIGINT, -- Globally unique 64-bit Snowflake ID (server message identity & lookup)
sender_id UUID,
content TEXT, -- Encrypted ciphertext or plaintext
content_type VARCHAR(32), -- 'text', 'image', 'video', 'voice', 'document'
media_url TEXT, -- Object storage reference for encrypted media blob
encrypted_media_key TEXT, -- Client-encrypted symmetric key for media decryption
encrypted_thumbnail TEXT, -- Client-encrypted preview thumbnail
created_at TIMESTAMP,
PRIMARY KEY ((conversation_id, bucket), conversation_seq)
) WITH CLUSTERING ORDER BY (conversation_seq DESC);Cassandra: User Conversation Index
Maintains the user's active inbox view sorted by latest message activity:
CREATE TABLE user_conversations (
user_id UUID,
last_message_at TIMESTAMP,
conversation_id UUID,
conversation_type VARCHAR(16), -- '1:1' or 'group'
last_read_seq BIGINT, -- User's local read watermark in this conversation
last_message_preview TEXT,
unread_count INT,
PRIMARY KEY (user_id, last_message_at)
) WITH CLUSTERING ORDER BY (last_message_at DESC);Redis: Multi-Device Session and Presence Schema
High-speed in-memory hash keys supporting multiple concurrent devices per user with automatic TTL expiration:
# Active Multi-Device WebSocket Session Registry
Key: sessions:{user_id}
Type: Hash {
"dev_ios_901": "{ server_id: 'chat-gw-042', connected_at: 1710320000000, last_heartbeat: 1710320030000 }",
"dev_web_412": "{ server_id: 'chat-gw-019', connected_at: 1710319500000, last_heartbeat: 1710320025000 }"
}
TTL: 300 seconds (refreshed per device on heartbeat)
# User Online Presence Key (cached aggregated state)
Key: presence:{user_id}
Type: String "online"
TTL: 60 seconds (refreshed while at least one device session is active)
# Last Seen Timestamp Store
Key: last_seen:{user_id}
Type: String "1710320001200"MySQL: User and Group Metadata Tables
Relational schemas managing user identity, group rosters, and public encryption keys:
CREATE TABLE users (
user_id UUID PRIMARY KEY,
phone VARCHAR(20) UNIQUE,
name VARCHAR(128),
avatar_url TEXT,
public_key BLOB, -- Identity public key for E2EE
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);CREATE TABLE groups (
group_id UUID PRIMARY KEY,
name VARCHAR(256),
avatar_url TEXT,
created_by UUID,
max_members INT DEFAULT 256,
created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE group_members (
group_id UUID,
user_id UUID,
role VARCHAR(16) DEFAULT 'member', -- 'admin', 'member'
joined_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (group_id, user_id)
);Kafka Event Streaming Layout
Message routing queues and push notification topologies:
Topic: group-fanout
Partitions: 64
Partition key: group_id
Retention: 24 hours
Producer: Group Service when coordinating asynchronous fan-out for groups with many participants
Event Schema: { message_id, conversation_id, conversation_seq, sender_id, payload, timestamp }
Topic: push-notifications
Partitions: 64
Partition key: user_id
Retention: 7 days
Producer: Chat Gateway Server when recipient has no active sessions
Event Schema: { user_id, conversation_id, sender_name, unread_count, collapse_key }
Consumer Groups:
1. group-fanout-workers: Resolves group membership rosters and dispatches to member chat servers via Redis Pub/Sub
2. push-workers: Dispatches batched push notifications to Apple APNs and Google FCM
Dead Letter Queue: push-notifications-dlq after 3 retries
Synchronous Ingress Flow:
Client WebSocket frame -> Chat Gateway validates auth & payload -> Message Service allocates monotonically increasing conversation_seq and persists to Cassandra with local quorum -> Server ACK returned to sender -> Chat Gateway queries Session Service in Redis for active recipient device sockets
Cross-Pod Delivery & Offline Catch-Up Flow:
If recipient is ONLINE on another server -> Publish to recipient server Redis Pub/Sub channel (fast path)
If delivery fails, socket disconnects, or session is stale -> No message is lost; Cassandra holds durable copy
If recipient has no active sessions -> Publish wake-up event to push-notifications Kafka topic
Authoritative recovery -> Offline and reconnecting clients sync missing messages directly from Cassandra by querying where conversation_seq > last_synced_conversation_seq across sequence bucketsFault Tolerance
Resilience patterns for gateway failover, transient network partitions, ordering recovery, and offline synchronization.
Fault Tolerance Strategies
| Failure Domain | Resilience Mechanism |
|---|---|
| Write-Before-ACK Persistence | Messages are committed to Cassandra with local quorum replication before returning an ingress confirmation to the sender, ensuring durable persistence against node crashes. |
| Client Retry Idempotency (clientMessageId) | Clients attach a unique clientMessageId to outgoing packets, allowing the server to detect duplicate resends and return existing results safely without creating duplicate database records. |
| Multi-Datacenter Cassandra Replication | Cross-region asynchronous replication with local quorum reads and writes protects against complete datacenter outages. |
| Redis Pub/Sub Fast-Path with Cassandra Catch-Up | Cross-pod delivery routes via Redis Pub/Sub for sub-millisecond delivery to online users. If a delivery signal drops or the socket disconnects, the message is already durable in Cassandra, and the recipient catches up via normal sync upon reconnect. |
Problem-Specific Failure Scenarios
1. Chat Gateway Server Crashes
- When a gateway instance fails, all active WebSockets on that server drop simultaneously.
- Client applications detect socket closure and reconnect with exponential backoff to another available gateway pod.
- The new gateway registers the updated
server_idin the Redis Session Service and initiates a sync query to fetch messages received during the disconnect window.
2. Out-of-Order Message Arrival and Client Retries
- Network latency variations on mobile connections can cause rapid messages to arrive out of sequence or trigger duplicate client retries.
- If network interruptions trigger client retries, the server checks the
clientMessageIdto return the existing ACK without creating duplicate rows. The Message Service assigns an authoritativeconversation_sequpon initial ingestion, allowing clients to buffer incoming packets and render items strictly in order ofconversation_seq.
3. Offline Recipient Catch-Up Synchronization
- When a user is offline, incoming messages are durably stored in Cassandra and unread counters increment.
- Upon reconnecting, the client issues a sync request providing its
last_synced_conversation_seq. The Message Service calculates the starting sequence bucket and fetches all newer messages from Cassandra across bounded partitions in paginated 50-item batches.
4. Group Fan-Out Queue Congestion
- High-velocity group conversations can saturate single-worker queues during viral discussions.
- The Group Service separates large groups into dedicated Kafka topic partitions, enabling independent consumer scaling without impacting one-to-one message traffic.
5. Multi-Region Network Partitions and Replication
- Messages are immutable records identified by globally unique
message_ids. Within a conversation, sequence allocation is assigned by the designated conversation partition sequencer. - Normal message writes append distinct immutable records rather than competing to overwrite existing rows. Cross-region multi-master configurations route sequencing through a partition owner or allocate partitioned sequence ranges per region to prevent sequence collisions, while Cassandra asynchronous replication synchronizes committed records across datacenters.
Additional Considerations
Advanced topics covering typing indicators, full-text message search, multi-device synchronization, and data lifecycle management.
Ephemeral Typing Indicators
Typing events are lightweight ephemeral signals that are never persisted to disk. When User A types, the client sends a typing frame over the WebSocket. The gateway looks up active session IDs for conversation participants in Redis and forwards the frame directly, automatically expiring the indicator after 5 seconds of inactivity.
Full-Text Message Search
For non-E2EE deployments (such as Slack), messages are ingested asynchronously into Elasticsearch clusters indexed by conversation_id and message content. For end-to-end encrypted chats (such as WhatsApp), the server stores only ciphertext, while full-text search runs locally on client devices using embedded SQLite databases.
Multi-Device Synchronization
When a user operates multiple concurrent devices (mobile, tablet, desktop web), the Session Service stores multiple active sessions under sessions:{user_id}. Incoming messages and read receipt cursor updates are duplicated across all registered device sockets. At the product level, the message status transitions to "delivered" once at least one registered recipient device acknowledges receipt, and transitions to "read" once the user advances their monotonic read cursor (read_up_to_seq).
Rate Limiting and Abuse Prevention
Token bucket rate limiters in the gateway enforce a cap of 100 messages per minute per user to block automated spam bots, with file attachments restricted to 100 MB per upload.
Related Problems and Concepts
WebSocket scaling and socket registry patterns connect directly to Design a Real-Time User Presence Service. Offline push delivery pipelines integrate with Design a Distributed Notification System, while group fan-out mathematics mirror principles from Design a News Feed System. Deepen your foundational networking and messaging knowledge in Network Protocols (HTTP, gRPC, WebSocket, DNS), Message Queues Fundamentals, Caching Patterns and Invalidation, and System Design Interview Patterns.
Interview Walkthrough
- 25-Minute Interview Strategy
Focus on the stateful connection lifecycle and cross-pod routing before discussing Cassandra time-series models and E2EE key exchanges.
- Requirements and Scope: 1:1 vs Group Chat (3 min)
- WebSocket Gateway Fleet and Redis Session Registry (6 min)
- Cross-Pod Routing via Redis Pub/Sub and Push Dispatch (6 min)
- Cassandra Bounded Storage Model and Sequence Ordering (5 min)
- Offline Push Notifications and Reconnect Sync (5 min)
- Establish the communication scope early, differentiating the low-latency direct push model for 1:1 chats from fan-out architectures for large groups.
- Illustrate the stateful connection lifecycle clearly, showing how TCP load balancing with TLS termination, in-memory socket maps, and Redis session registries coordinate cross-server delivery.
- Emphasize per-conversation ordering over global ordering, explaining how server-assigned conversation sequence numbers guarantee deterministic client rendering.
- Highlight write-before-ACK persistence in Cassandra to prove that acknowledged messages are committed before returning ingress confirmation.
- Describe how read receipts operate via monotonic cursors rather than quadratic per-message broadcast updates.
- Outline the Signal Protocol (X3DH and Double Ratchet) when E2EE is in scope, clarifying that encrypted messaging transforms the server into an untrusted blind relay.
Engineering Trade-offs
Key architectural trade-offs across transport protocols, database storage engines, encryption standards, sequence ordering mechanisms, and group distribution models.
Transport Protocol: WebSocket vs Long Polling vs Server-Sent Events
Evaluating client-server transport protocols for real-time bidirectional messaging latency and connection efficiency.
| Protocol | Operational Mechanism | Latency | Bidirectional | Connection Overhead | Target Use Case |
|---|---|---|---|---|---|
| WebSocket ⭐ | Persistent full-duplex TCP socket with lightweight framing headers (2 to 14 bytes) | < 100 ms | True full-duplex | 1 persistent TCP socket per client | Real-time chat, collaborative workspaces, multiplayer gaming |
| HTTP Long Polling | Client sends HTTP request held open by server until new message data arrives | 100 to 500 ms | Simulated via separate POST | New HTTP handshake and header overhead per message batch | Fallback transport behind restrictive enterprise corporate firewalls |
| Server-Sent Events (SSE) | Persistent HTTP connection streaming server-to-client events over standard HTTP/2 | < 100 ms | Server to client only | 1 HTTP/2 stream per client | Unidirectional event feeds, stock tickers, notification streams |
Message Storage Engine: Cassandra vs Relational SQL vs Document Stores
Analyzing database engines against append-only time-series message patterns.
| Database Engine | Architectural Characteristics | Trade-off Assessment |
|---|---|---|
| Cassandra (NoSQL) ⭐ | Partition by (conversation_id, bucket) with conversation_seq clustering order and LSM-tree writes | Handles over 58,000 writes per second effortlessly, scales horizontally linearly, and natively clusters chronological messages for single-seek reads while bucketing bounds partition sizes. |
| MySQL / PostgreSQL (Sharded) | Relational B-tree tables sharded horizontally by conversation_id | Strong ACID consistency for metadata, but B-tree page splits and index locking create severe write bottlenecks at scale. |
| MongoDB (Document Store) | Document-per-conversation or bucketed message arrays | Flexible schema, but document size limits and WiredTiger memory overhead make append-heavy time-series chat less efficient than Cassandra. |
End-to-End Encryption: Signal Protocol vs Transport-Only TLS
Balancing user privacy and security guarantees against server-side feature capabilities.
| Security Architecture | Cryptographic Properties | Trade-off Profile |
|---|---|---|
| Transport Layer Security (TLS Only) | Encrypted in transit over HTTPS and WSS while the server stores plaintext | Enables server-side full-text search, server-side content moderation, and rich link parsing, but server compromise exposes all user communications. |
| Signal Protocol (E2EE) ⭐ | X3DH key exchange with Double Ratchet per-message ephemeral keys | Provides forward secrecy and post-compromise healing so the server cannot read content, but shifts search and moderation to client devices. |
Cross-Server Message Routing: Redis Pub/Sub vs Distributed Kafka Queue
Comparing low-latency memory brokers with durable partitioned event streams.
| Routing Mechanism | Execution Model | Trade-off Assessment |
|---|---|---|
| Redis Pub/Sub (Fast Path) ⭐ | Publish directly to recipient gateway server channel | Sub-millisecond routing latency for active online users. Because messages are committed to Cassandra before routing, any dropped signals recover via client sync upon reconnect. |
| Kafka Broker (Durable Path) ⭐ | Partitioned topic per group ID or background worker pool | Provides durable buffering, replay capabilities, and consumer backpressure, making it ideal for push notification batches and large group fan-out workflows. |
Group Chat Delivery: Push Fan-Out vs Pull on Read
Managing delivery overhead between small private chats and massive public channels.
| Group Scale | Delivery Paradigm | Architectural Rationale |
|---|---|---|
| Small Groups (< 256 members) | Fan-out on write (push model) | Iterates through member session list and pushes directly to active sockets with small, bounded write amplification. |
| Large Channels (1,000+ members) | Fan-out on read (pull model) |
|
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.