Functional Requirements:
ID Functional Requirement
FR1 Users can send and receive text messages in real time.
FR2 Users can access and reply to conversations from web and mobile applications.
FR3 Users can send images, videos, files, emojis, and stickers.
FR4 Users can send voice messages.
FR5 Users can react to messages.
FR6 Users can see another user's online/offline presence and last seen.
FR7 Users can hide their online status and/or last seen through privacy settings.
FR8 The system maintains conversation history, including messages received while the user was offline.
FR9 The messages are encrypted from end to end
FR10 1:1 messaging
fr11 n:n messages, IN CAN CREATE GROUP CHATS
Non-Functional Requirements:
ID Non-Functional Requirement
NFR1 Low Latency / Real-Time Sync: Messages should be delivered to connected users in near real time.
NFR2 High Availability: Messaging should remain available despite individual server/instance/AZ failures.
NFR3 Scalability: The system must horizontally scale to support hundreds of millions of users and millions of messages per second.
NFR4 Security: User data and messages must be protected in transit and at rest, with authentication and authorization.
NFR5 Reliable Delivery: Use at-least-once delivery so transient failures do not silently lose messages.
NFR6 Idempotency / Deduplication: Duplicate message deliveries must not result in duplicate messages being stored/displayed.
NFR7 Message Ordering: Messages within a conversation should maintain a consistent logical order.
Estimate the scale of the system. Consider daily active users, read/write ratio, storage requirements, bandwidth, and any relevant QPS calculations...
| Contract Protocol Purpose |
POST /conversations | HTTPS | Create 1:1/group conversation |
GET /conversations | HTTPS | Get user's conversations |
GET /conversations/{conversationId}/messages?cursor=...&limit=50 | HTTPS | Load/catch up message history |
POST /conversations/{conversationId}/members | HTTPS | Add group member |
DELETE /conversations/{conversationId}/members/{userId} | HTTPS | Remove/leave member |
POST /media/uploads | HTTPS | Request media_id + presigned S3 upload URL |
WSS /ws | WebSocket | Establish persistent real-time channel |
MESSAGE_SEND | WebSocket event | Send message |
MESSAGE_ACK | WebSocket event | Server accepted/persisted message |
MESSAGE_RECEIVED | WebSocket event | Deliver message to recipient |
MESSAGE_DELIVERED | WebSocket event | Recipient device received message |
MESSAGE_READ | WebSocket event | Update read position |
REACTION_ADD | WebSocket event | Add reaction |
REACTION_REMOVE | WebSocket event | Remove reaction |
TYPING_START/STOP | WebSocket event | Ephemeral typing state |
HEARTBEAT | WebSocket event | Keep connection/presence alive |
Describe the overall system architecture. Identify the main components needed to solve the problem end-to-end. Use the diagramming tool to create a block diagram.
| Area Architecture Decision Why / Rationale How it works in our design |
| Client Architecture | Web + Mobile clients use the same backend | Avoid separate backend architectures per client while maintaining consistent messaging behavior | Both clients consume the same REST APIs and establish WebSocket connections to the messaging backend |
| Real-Time Communication | Persistent WebSocket connections | Messaging requires low-latency bidirectional communication; polling would create excessive requests and latency | Client establishes a persistent WSS connection through the load balancer to a messaging instance |
| Transport Security | TLS/WSS | Protect traffic between clients and backend from interception/tampering | REST uses HTTPS and WebSocket traffic uses WSS |
| End-to-End Encryption | Encrypt/decrypt message content on endpoints | Backend should transport/store ciphertext rather than plaintext message content | Sender encrypts → backend transports/stores ciphertext → recipient decrypts |
| Application Architecture | Horizontally scaled stateless messaging nodes | We may have millions of concurrent users and millions of messages/sec; one server cannot handle the workload | Many application instances handle REST requests and persistent WebSocket connections |
| Compute Model | Long-running containerized compute | WebSockets are long-lived connections and messaging traffic is continuously active | Containers remain alive and maintain client WebSocket connections while scaling horizontally |
| Traffic Distribution | Layer-7 load balancing | Incoming HTTP/WSS traffic must be distributed across healthy messaging instances | New connections/requests are routed across the available application instances |
| Application State | Application instances remain stateless with durable state externalized | If an instance crashes or scales in, messages must not disappear with it | Persistent state is stored externally; clients reconnect to another healthy instance |
| Connection Recovery | Clients reconnect and synchronize after connection loss | WebSocket connections are inherently transient and can disappear when a node fails | Client reconnects through the load balancer and retrieves messages/state it missed |
| Message Delivery Semantics | At-least-once + idempotency | Retrying is safer than silently losing a message, but retries can generate duplicates | Every message has a unique message_id; duplicate processing is detected/ignored |
| Message Ordering | Ordering scoped to a conversation, not globally | Global ordering across hundreds of millions of users is unnecessary and expensive | Messages/events from the same conversation follow an ordered sequence/partition |
| Persistent Message Storage | Distributed NoSQL message store | Workload is enormous, write-heavy and primarily queried by conversation rather than relational joins | Messages are stored using conversation-oriented access patterns |
| Message Partitioning | Partition by conversation_id | Dominant read pattern is “give me messages belonging to conversation X” | Messages from the same conversation are colocated logically for efficient history queries |
| Sort Strategy | Sequence/time + message ID | Conversation history needs deterministic chronological retrieval | Messages within a conversation partition are sorted by their ordering key |
| Hot Conversation Handling | Bucket/shard very hot conversations | A massive group could concentrate excessive traffic on one logical partition | conversation_id can be extended with a bucket/shard identifier |
| Data Replication | Replicated durable data | Failure of a physical node/AZ must not destroy conversation history | Database infrastructure maintains redundant copies of persistent data |
| Message Event Pipeline | Distributed asynchronous event stream | Sender processing should not synchronously wait for every recipient delivery, especially for groups | After persistence, a message event enters the messaging stream and delivery processing happens asynchronously |
| Event Partitioning | Partition event stream by conversation_id | Helps preserve ordering for events belonging to the same conversation | Events for a conversation consistently map to the same logical event partition |
| Delivery Workers | Separate asynchronous delivery consumers | Message ingestion and recipient delivery have different scaling characteristics | Workers consume message events, resolve recipients/connections and perform fan-out/delivery work |
| Fan-Out Strategy | Async fan-out with potential hybrid strategy | One group message can generate hundreds/thousands of delivery operations | Small conversations can favor fan-out-on-write; very large groups can shift more work toward read/catch-up |
| Connection Registry | Distributed ephemeral mapping of users/devices → active connections | Delivery processing must know whether and where a recipient is currently connected | Connection information is registered when WebSocket sessions are established |
| Presence | Heartbeat + TTL | Explicit logout/disconnect events aren't reliable because devices can lose network/power suddenly | Clients send heartbeats; stale entries expire automatically and user becomes offline |
| Multi-Device Support | One user can have multiple sessions/connections | Same account may simultaneously be active on phone, laptop, browser, etc. | Connection registry maintains multiple active connections for the same user_id |
| Multi-Device Delivery | Deliver events to active devices without duplicating persistent messages | Each device should synchronize while only one logical message exists | One message is persisted; delivery layer distributes its event to the user's active connections |
| Offline Messaging | Durable history + catch-up synchronization | Keeping messages indefinitely in queues for offline users would be inefficient | Messages remain in persistent storage; reconnecting clients request everything after their last known position |
| Read State | Store last_read_message_id / cursor per user+conversation | A boolean record for every user × message becomes enormous | Advancing one read cursor implicitly marks previous messages as read |
| Cache | Cache only hot/ephemeral state | Caching entire conversation history would be expensive and largely wasteful | Presence, connections and frequently accessed metadata are cached; persistent DB remains source of truth |
| Cache Strategy | Cache-aside | Cache must not become the authoritative persistent data store | Application checks cache → on miss reads persistent store → optionally populates cache |
| Media Storage | Separate object storage from message database | Images/videos/audio/files are large binary objects and shouldn't consume message database capacity | Message DB stores metadata/reference; binary content lives in object storage |
| Media Upload | Client uploads directly to object storage | Routing a 200 MB video through application containers wastes bandwidth, memory and compute | Backend authorizes upload; client transfers binary directly to object storage |
| Media Metadata | Store references rather than binary payloads | Keeps message records small and efficiently queryable | Message contains media_id, type, metadata/reference rather than the actual video/image |
| Async Media Processing | Event-driven processing after upload when needed | Transcoding, thumbnails, scanning, etc. shouldn't block message APIs | Successful media upload can trigger asynchronous processing pipelines |
| Conversation Model | Unified conversation_id for 1:1 and groups | Avoid designing separate messaging models for direct and group conversations | Both direct and group messages reference a conversation; membership determines recipients |
| Group Membership | Membership stored independently from individual messages | Group membership changes over time and shouldn't be duplicated inside every message | Conversation membership determines who participates/receives events |
| Read Scalability | Query conversation-oriented partitions + optional hot cache | History requests should avoid scanning the entire global message dataset | Client queries specific conversation partitions with cursor-based pagination |
| Write Scalability | Horizontal partitioning + asynchronous delivery | Millions of message writes/sec cannot depend on a single writer/node | Writes distribute across partitions while downstream delivery scales independently |
| Connection Scalability | Horizontally scale WebSocket nodes independently | Concurrent connection count is a different scaling dimension from message throughput | More messaging nodes can be added as concurrent WebSocket connections increase |
| Autoscaling | Scale from workload signals, not only CPU | A WebSocket server can exhaust connection capacity while CPU remains relatively low | Connection count, throughput, CPU/memory and related operational metrics can drive scaling |
| High Availability | Multiple application instances across failure domains | Loss of one instance/AZ shouldn't take Messenger offline | Load balancing sends traffic to healthy instances and clients reconnect after failures |
| Backpressure | Buffer asynchronous work through the event pipeline | Sudden message spikes shouldn't immediately overwhelm delivery workers | Event stream absorbs temporary producer/consumer throughput differences |
| Observability | Centralized logs, metrics and alarms | Distributed asynchronous systems are difficult to troubleshoot without visibility | Track latency, failures, active connections, throughput, consumer lag, retries, etc. |
| Service Authorization | Least-privilege workload identities | Application components should access only the infrastructure/resources they require | Runtime workloads receive scoped permissions rather than static credentials |
| Failure Philosophy | Assume compute and connections will fail | At this scale, failures are routine rather than exceptional | Durable state + retries + idempotency + reconnection + async processing allow recovery |
Define the data model. Identify the main entities, their attributes, and relationships. Consider the choice of database type (SQL vs NoSQL) and justify your decision based on access patterns...
Deep dive into 2-3 key components. Explain how they work, how they scale, discuss tradeoffs, capacity, and any relevant algorithms or data structures.