System Design Problem

Design a Real-Time Chat Application (WhatsApp / Slack)

Commonly Asked By:MetaSlackMicrosoftDiscord

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)

QuestionWhy 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.

MetricCalculationValue
Daily active users (DAU)Given product scale50M
Peak concurrent connectionsEstimated 20% active peak concurrency~10M concurrent WebSockets
Daily messages processed50M DAU x 10 messages per user500M messages / day
Average message throughput500M messages ÷ 86,400 seconds~5,787 messages / sec
Peak message throughput (10x)Peak surge traffic multiplier~58,000 messages / sec
Average text message payloadText, metadata, and encryption wrappers~200 bytes
Daily text storage volume500M 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.

Loading...

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 a device_id to its host server_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 clientMessageId for retry idempotency, enabling safe resends without duplicate processing if an ACK is lost. The Message Service assigns a globally unique 64-bit Snowflake message_id for internal server-side identity and correlation. Within the specific conversation, it assigns a monotonically increasing conversation_seq to 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_id route 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 as bucket = 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 DESC stores 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 where conversation_seq > last_synced_conversation_seq within 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-notifications Kafka 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 buckets

API 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:

TYPESCRIPT
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:

JSON
{
  "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:

JSON
{
  "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:

JSON
{
  "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:

JSON
{
  "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:

HTTP
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:

SQL
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:

SQL
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:

SQL
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
);
SQL
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:

YAML
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 buckets

Fault Tolerance

Resilience patterns for gateway failover, transient network partitions, ordering recovery, and offline synchronization.

Fault Tolerance Strategies

Failure DomainResilience Mechanism
Write-Before-ACK PersistenceMessages 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 ReplicationCross-region asynchronous replication with local quorum reads and writes protects against complete datacenter outages.
Redis Pub/Sub Fast-Path with Cassandra Catch-UpCross-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_id in 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 clientMessageId to return the existing ACK without creating duplicate rows. The Message Service assigns an authoritative conversation_seq upon initial ingestion, allowing clients to buffer incoming packets and render items strictly in order of conversation_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.

ProtocolOperational MechanismLatencyBidirectionalConnection OverheadTarget Use Case
WebSocket ⭐Persistent full-duplex TCP socket with lightweight framing headers (2 to 14 bytes)< 100 msTrue full-duplex1 persistent TCP socket per clientReal-time chat, collaborative workspaces, multiplayer gaming
HTTP Long PollingClient sends HTTP request held open by server until new message data arrives100 to 500 msSimulated via separate POSTNew HTTP handshake and header overhead per message batchFallback transport behind restrictive enterprise corporate firewalls
Server-Sent Events (SSE)Persistent HTTP connection streaming server-to-client events over standard HTTP/2< 100 msServer to client only1 HTTP/2 stream per clientUnidirectional 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 EngineArchitectural CharacteristicsTrade-off Assessment
Cassandra (NoSQL) ⭐Partition by (conversation_id, bucket) with conversation_seq clustering order and LSM-tree writesHandles 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_idStrong 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 arraysFlexible 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 ArchitectureCryptographic PropertiesTrade-off Profile
Transport Layer Security (TLS Only)Encrypted in transit over HTTPS and WSS while the server stores plaintextEnables 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 keysProvides 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 MechanismExecution ModelTrade-off Assessment
Redis Pub/Sub (Fast Path) ⭐Publish directly to recipient gateway server channelSub-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 poolProvides 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 ScaleDelivery ParadigmArchitectural 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)
  • Messages are stored once in the channel timeline
  • clients fetch new entries on demand or subscribe to channel rooms, avoiding massive write amplification.

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...