System Design

0.0(0)
Studied by 2 people
call kaiCall Kai
learnLearn
examPractice Test
spaced repetitionSpaced Repetition
heart puzzleMatch
flashcardsFlashcards
GameKnowt Play
Card Sorting

1/55

encourage image

There's no tags or description

Looks like no tags are added yet.

Last updated 4:54 AM on 10/1/26
Name
Mastery
Learn
Test
Matching
Spaced
Call with Kai
Chat

No analytics yet

Send a link to your students to track their progress

56 Terms

1
New cards

options for scaling writes

  1. vertical scaling and Database Choices

  2. Sharding and Partitioning

  3. Handling Bursts with Queues and Load Shedding

  4. Batching and Hierarchical Aggregation


2
New cards

Delivery Framework

Requirements, Core Entities, API or Interface, Data Flow, High Level Design, Deep Dives

3
New cards

Non Functional Requirement List

  1. CAP Theorem: Should your system prioritize consistency or availability? Note, partition tolerance is a given in distributed systems.

  2. Environment Constraints: Are there any constraints on the environment in which your system will run? For example, are you running on a mobile device with limited battery life? Running on devices with limited memory or limited bandwidth (e.g. streaming video on 3G)?

  3. Scalability: All systems need to scale, but does this system have unique scaling requirements? For example, does it have bursty traffic at a specific time of day? Are there events, like holidays, that will cause a significant increase in traffic? Also consider the read vs write ratio here. Does your system need to scale reads or writes more?

  4. Latency: How quickly does the system need to respond to user requests? Specifically consider any requests that require meaningful computation. For example, low latency search when designing Yelp.

  5. Durability: How important is it that the data in your system is not lost? For example, a social network might be able to tolerate some data loss, but a banking system cannot.

  6. Security: How secure does the system need to be? Consider data protection, access control, and compliance with regulations.

  7. Fault Tolerance: How well does the system need to handle failures? Consider redundancy, failover, and recovery mechanisms.

  8. Compliance: Are there legal or regulatory requirements the system needs to meet? Consider industry standards, data protection laws, and other regulations.


4
New cards

Three most important layers in OSI model

layer 7, Application Layer (HTTP, REST, etc)

layer 4, Transport Layer (TCP/UDP)

layer 3, Network Layer (IP Addresses)

5
New cards

What happens when you type something into browser

DNS converts to IP, set up a TCP connection over IP, send our HTTP request, get a response, and tear down the connection

6
New cards

Transport Layer Protocols and their tradeoffs

UDP: (video streaming, phone calls, etc)

1. Connectionless: No handshake or connection setup

  1. No guarantee of delivery: Packets may be lost without notification

  2. No ordering: Packets may arrive in a different order than sent

  3. Lower latency: Less overhead means faster transmission

TCP: (everything else)

  1. Connection-oriented: Establishes a dedicated connection before data transfer

  2. Reliable delivery: Guarantees that data arrives in order and without errors

  3. Flow control: Prevents overwhelming receivers with too much data

  4. Congestion control: Adapts to network congestion to prevent collapse


7
New cards

Application Layer Protocols

HTTP, REST, GraphQL, gRPC, Server Sent Events, Websockets, WebRTC,

8
New cards

HTTP: definition, common methods & error messages

request-response protocol

Common Methods:

  • GET: Request data from the server. GET requests should be idempotent and don't have a body.

  • POST: Send data to the server.

  • PUT: Update data on the server.

  • PATCH: Update a resource partially.

  • DELETE: Delete data from the server. DELETE requests should be idempotent.

Common Status Codes:

  • 2xx: success

  • 3xx: moved

  • 4xx: Client Error

  • 5xx: Server error


9
New cards

When to use GraphQL

overfetching, underfetching, or need flexible data (e.g., browser vs app needs different layouts and data)

10
New cards

gRPC when to use

gRPC shines in microservices architectures where services need to communicate efficiently. Its strong typing helps catch errors at compile time rather than runtime, and its binary protocol is more efficient than JSON over HTTP

Consider gRPC for internal service-to-service communication, when performance is critical, not for front-facing APIs

11
New cards

What is SSE and when to use

Server sent events, a nice hack on top of HTTP that allows a server to stream many messages, over time, in a single response from the server.

with SSE, the server can push many messages as "chunks" in a single response from the server rather than one big JSON blob

You'll find SSE useful in system design interviews in situations where you want clients to get notifications or events as soon as they happen.

12
New cards

Web sockets: what is it and when to use

WebSockets provide a persistent, TCP-style connection between client and server, allowing for real-time, bidirectional communication with broad support

WebSockets come up in system design interviews when you need high-frequency, persistent, bi-directional communication between client and server. Think real-time applications, games, and other use-cases where you need to send and receive messages as soon as they happen.

13
New cards

WebRTC: What it is and when to use

WebRTC enables direct peer-to-peer communication between browsers without requiring an intermediary server for the data exchange. WebRTC can be perfect for collaborative applications like document editors (if needing to scale for many clients) and is especially useful for video/audio calling and conferencing applications.

14
New cards

Cassandra: Overview + Benefits

distributed NoSQL database. It implements a partitioned wide-column storage model with eventually consistent semantics

15
New cards

Cassandra: Partitioning Description

Cassandra achieves horizontal scalability by partitioning data across many nodes in its cluster. In order to partition data successfully, Cassandra makes use of consistent hashing (nodes in a ring rather than hashing with modulo for easier rehashing) with virtual nodes (to even spread)

16
New cards

Cassandra: Replication

Virtual nodes are mapped to physical nodes; data written to a virtual node will also be replicated to other virtual nodes mapping to different physical nodes

17
New cards

Cassandra: Consistency

Cassandra allows you to choose from a list of "consistency levels" for reads and writes (tuned consistency), which are required node response numbers for a write or a read to succeed. These enforce different consistency vs. availability behavior depending on the combination used. These range from ONE, where a single replica needs to respond, to QUORUM (majority), to ALL, where all replicas must respond.

18
New cards

Cassandra’s Storage Model

Cassandra leverages a data structure called a Log Structured Merge Tree (LSM tree) index to achieve speed in write throughput

Every create / update / delete is a new entry (with some exceptions). Cassandra uses the ordering of these updates to determine the "state" of a row. For example, if a row is created and then it is updated later, Cassandra will understand the state of the row by looking at the creation and then the update vs. looking at just a single row.

19
New cards

Cassandra: Gossip

Cassandra nodes communicate information throughout the cluster via "gossip", which is a peer-to-peer scheme for distributing information between nodes. Universal knowledge of the cluster makes every node aware and able to participate in all operations of the database, eliminating any single points of failure and allowing Cassandra to be a very reliable database for availability-skewing system designs.

20
New cards

Cassandra: Fault Tolerance

When a node gossips with a node that doesn't respond, Cassandra's failure detection logic "convicts" that node and stops routing writes to it. The convicted node can re-enter the cluster when it starts heartbeating again. Cassandra treats a convicted node as unavailable but keeps it in the ring until an administrator explicitly changes the topology. This is done to prevent intermittent communication failures / node restarts from causing the cluster to re-balance data.

21
New cards

Dedicated Load Balancer Types

Layer 4 load balancers have some key characteristics, they ...

  • Maintain persistent TCP connections between client and server.

  • Are fast and efficient due to minimal packet inspection.

  • Cannot make routing decisions based on application data.

  • Are typically used when raw performance is the priority

  • WEBSOCKETS

Layer 7 load balancers have some key characteristics, they ...

  • Terminate incoming connections and create new ones to backend servers.

  • Can route based on request content (URL, headers, cookies, etc.).

  • More CPU-intensive due to packet inspection.

  • Provide more flexibility and features.

  • Better suited for HTTP-based traffic.

  • Everything else but Websockets


22
New cards

Failover with load balancers

Load balancers offer health checks, as they are essential for high availability.

A Layer 7 health check might make an HTTP request to the server and make sure the response is success (e.g. a 200 status code vs a 500 indicating internal failures or no response indicating a crash).

23
New cards

Regionalization

Data locality - try to keep databases near servers, but if users are spread out, hard to deal with. Options: CDNs and Regional Partitioning (If we have a lot of users in a single region, we can partition our data by region so that each region only has data relevant to it.)

24
New cards

How to handle failures in system

retry with exponential backoff + jitter

make sure requests are idempotent with idempotency key. The idempotency key is a unique identifier for a request that we can use to make sure the same request is idempotent (Idempotent APIs are APIs that can be called multiple times and they produce the same result every time.)

thundering herd issue/cascading failures: (If your database has gone down cold and you need to boot it up one instance at a time, having a firehose of retries and angry users might pin down an instance from ever getting started.)

  1. The circuit breaker monitors for failures when calling external services

  2. When failures exceed a threshold, the circuit "trips" to an open state

  3. While open, requests immediately fail without attempting the actual call

  4. After a timeout period, the circuit transitions to a "half-open" state

  5. A test request determines whether to close the circuit or keep it open


25
New cards

Common API Patterns

Pagination, Idempotency keys, filtering and sorting, consistent error responses, versioning

26
New cards

Pagination

Offset-based pagination is the simplest approach and used by most websites. You specify how many records to skip and how many to return: /events?offset=20&limit=10 gets records 21-30
Cursor-based pagination solves this by using a pointer to a specific record instead of counting from the beginning. The cursor is typically an encoded reference to a specific record (like an ID or timestamp).

27
New cards

API Security Principles

Authentication verifies identity - proving the user is who they claim to be. Authorizationverifies permissions - checking if that authenticated user is allowed to perform the specific action they're requesting.

API keys are long, randomly generated strings that act like passwords for applications rather than humans. When a client makes a request, they include their API key in the Authorization header, and your server looks up that key to identify which application is making the request. (Not for user-facing stuff)

JWT tokens (JSON web tokens), on the other hand, encode user information directly into the token itself rather than storing session state on your server. When a user logs in successfully, your server creates a JWT containing their user ID, permissions, and an expiration time, then signs the entire token with a secret key. Conveniently, when that JWT comes back with future requests, you can verify it's authentic by checking the signature, and you can read the user information directly from the token without any database lookups.

Role-Based Access Control: Real systems have different types of users with different permissions. RBAC assigns roles to users and permissions to roles.

Rate Limiting and Throttling: You typically implement rate limiting at the API gateway level or using middleware in your application, restricting how many requests a client can make in a given time period.

28
New cards

Ad Aggregator Problem: Users can click on ads and be redirected to the target

The user clicks on the ad, which will then send a request to our server. Our server can then track the click and respond with a redirect to the advertiser's website via a 302 (redirect) status code.

29
New cards

Ad Aggregator Problem: Advertisers can query ad click metrics over time at 1 minute intervals

  1. The click processor service writes the event to a stream like Kafka or Kinesis.

  2. A stream processor like Flink or Spark Streaming reads the events from the stream and aggregates them in real-time.

  3. The aggregated data is stored in our OLAP database (any database) for querying.

  4. Advertisers can query the OLAP database to get metrics on their ads in near real-time.

Note we may want to use Cassandra bc write-heavy, but Cassandra not good with range queries - which is why we need a diff solution.

30
New cards

Ad Aggregator Problem: How to Scale to Support 10k clicks per second? + how to deal with very popular ads

  1. Click Processor Service: We can easily scale this service horizontally by adding more instances; will need load balancer now

  2. Stream: Need to shard Kinesis/Kafka. Sharding by AdId is a natural choice, this way, the stream processor can read from multiple shards in parallel since they will be independent of each other (all events for a given AdId will be in the same shard).

  3. Stream Processor: The stream processor, like Flink, can also be scaled horizontally by adding more tasks or jobs. We'll have a separate Flink job reading from each shard doing the aggregation for the AdIds in that shard.

  4. OLAP Database/Database Sharding: You can shard by AdvertiserId so all data for a given advertiser lives on the same node, making queries for that advertiser's ads faster.

To solve the hot shard problem, we need a way of further partitioning the data. One popular approach is to update the partition key by appending a random number to the AdId. We could do this only for the popular ads as determined by ad spend or previous click volume. This way, the partition key becomes AdId:0-N where Nis the number of additional partitions for that AdId.

31
New cards

Ad Aggregator Problem: How can we ensure that we don't lose any click data?

Kafka replicates data across multiple brokers within a cluster, and Kinesis replicates across multiple availability zones within a region. Even if a node goes down, the data is not lost. Importantly for our system, they also allow us to enable persistent storage, so even if the data is consumed by the stream processor, it is still stored in the stream for a certain period of time.
Stream processors like Flink also have a feature called checkpointing. This is where the processor periodically writes its state to a persistent storage like S3. Aggregation windows are short though (a minute), so doesn’t make sense necessarily.

32
New cards

Ad Aggregator: How can we prevent abuse from users clicking on ads multiple times?

Rather than asking user to log in to view our ad so we can generate a JWT/session token, the better approach is to have the Ad Placement Service generate a unique impression ID for each ad instance shown to each user. This impression ID would be sent to the browser along with the ad and will serve as an idempotency key. When the user clicks on the ad, the browser sends the impression ID along with the click data. Check if the impression ID is signed with a secret key as well by the ad service, to make sure it is legitimate.

33
New cards

Ad Aggregator: How can we ensure that advertisers can query metrics at low latency?

The stream solution with Kafka —> Flink —> (OLAP) database works well.

For querying over very large time windows, like days, weeks, or years, we can pre-aggregate the data in the OLAP database. This can be done by creating a new table that stores the aggregated data at a higher level of granularity, like daily or weekly. This can be via a nightly cron job that runs a query to aggregate the data and store it in the new table.

34
New cards

FB Live Comments: Viewers can post comments on a Live video feed

Users will initiate a POST request to the POST /comments/:liveVideoId endpoint with the comment message. The server will then validate the request and store the comment in the database.

35
New cards

FB Live Comments: Viewers can see comments made before they joined the live feed

For the history of comments, users should be able to scroll up to view progressively older comments - this UI pattern is called "infinite scrolling" and is commonly used in chat applications. Cursor pagination is a technique that uses a cursor to specify the starting point for fetching a set of results. The cursor is a unique identifier that points to a specific item in the list of results. Initially, the cursor is set to the most recent comment, and it is updated each time the user scrolls to load more.

36
New cards

FB Live Comments: How can we ensure comments are broadcasted to viewers in real-time?

Polling as an initial solution is good but not real time; instead, we consider Websockets or Server-Sent Events. SSE is a persistent connection that operates over standard HTTP, making it simpler to implement than WebSockets. The server can push data to the client in real-time, while client-to-server communication happens through regular HTTP requests.

This is a better solution given our read/write ratio imbalance. The infrequent comment creation can use standard HTTP POST requests, while the frequent reads benefit from SSE's efficient one-way streaming.

37
New cards

FB Live Comments: How will the system scale to support millions of concurrent viewers?

When we add more servers to handle the load, viewers watching the same live video may end up connected to different servers, and when a comment is posted by a viewer on one server, it might not be sent to a viewer on another server.

The first thing we need to do is separate out the write and read traffic by creating Realtime Messaging Servers that are responsible for sending comments to viewers. We separate this out because the write traffic is much lower than the read traffic and we need to scale them independently.

When a new comment is created, how does each server learn about it? The most common solution is pub/sub. The comment management service publishes a message to a channel whenever a new comment is created, and all Realtime Messaging Servers subscribe to this channel. They then send the comment to all viewers watching that live video.

We can improve efficiency by partitioning the comment stream into different channels. Each Realtime Messaging Server subscribes only to the channels it needs, determined by the viewers connected to it.

Since having a channel per live video would consume a lot resources (and may be infeasible with systems like Kafka), we use a hashing function to distribute across N channels: hash(liveVideoId) % N. This ensures a reasonable upper bound on channel count.

38
New cards

FB Live Comments: What happens when the World Cup final goes live and suddenly hundreds of millions of viewers are watching the same stream, with comments pouring in at thousands per second?

An improved approach to streaming over every comment is to recognize that users don't need to see every comment. They just need to see enough comments to feel the energy of the crowd. We can sample the comment stream and only deliver a subset to each viewer.

39
New cards

FB Live Comments: How do we handle client disconnections and ensure viewers don't miss comments?

SSE has a built-in mechanism for handling reconnections through the Last-Event-IDheader. Every comment we send includes a unique event ID (the comment ID works well). When the browser detects a disconnection, it automatically reconnects and includes the last event ID it received. On the server side, when we see this header, we replay missed comments before resuming normal streaming.

40
New cards

What approach to use when we have some allowable latency (e.g., 5 seconds) and need semi-real-time updates from the system?

Simple polling: The client makes a request to the server at a regular interval and the server responds with the current state of the world.

41
New cards

What approach to use when we need semi-real time updates from the servers to our clients and we do not want to set up and tear down a TCP connection constantly?

The idea is also simple: the client makes a request to the server and the server holds the request open until new data is available. It's as if the server is just taking really long to process the request. The server then responds with the data, finalizes the HTTP requests, and the client immediately makes a new HTTP request. This repeats until the server has new data to send. If no data has come through in a long while, we might even return an empty response to the client so that they can make another request. This may introduce some extra latency, but is generally more efficient than short polling solution.

Long Polling is a great solution for applications where a long async process is running but you want to know when it finishes, as soon as it finishes - like is often the case in payment processing. We'll long-poll for the payment status before showing the user a success page. Standard length of poll is 15-30 seconds.

42
New cards

What approach to use when servers needs to push data to clients in real-time constantly, but clients don’t need to send data back frequently?

Server-Sent Events: Normally HTTP responses have a header like Content-Length which tells the client how much data to expect. SSE instead uses a special header Transfer-Encoding: chunked which tells the client that the response is a series of chunks - we don't know how many there are or how big they are until we send them. This allows us to move from a single, atomic request/response to a more granular "stream" of data.

A very popular use-case for SSE today is AI chat apps which frequently involve the need to stream new tokens (words) to the user as they are generated to keep the UI responsive. We also used it for our FB Live Comments solution.

43
New cards

What solution to use when you have high frequency writes and reads and need updates in real time?

Websockets build on HTTP through an "upgrade" protocol, which allows an existing TCP connection to change L7 protocols.

  1. Client initiates WebSocket handshake over HTTP

  2. Connection upgrades to WebSocket protocol

  3. Both client and server can send messages

  4. Connection stays open until explicitly closed

L4 load balancers will support websockets natively since the same TCP connection is used for each request. Note that if you do not need high-frequency communication, a very common pattern is to have SSE subscriptions for updates and do writes over simple HTTP POST/PUT whenever they occur.

There are many issues associated with the required stateful connections that WebSockets bring, so many architectures will terminate WebSockets into a WebSocket service early in their infrastructure. This service can then handle the connection management and scaling concerns and allows the rest of the system to remain stateless.


44
New cards

What solution to use when you have audio/video conferencing calls, multiplayer games, or potentially document editors?

WebRTC enables direct peer-to-peer communication between browsers. Clients talk to a central "signaling server" which keeps track of which peers are available together with their connection information. Once a client has the connection information for another peer, they can try to establish a direct connection without going through any intermediary servers. It reduces server load.

The WebRTC standard includes two methods to work around restrictions to inbound connections:

  • STUN: A protocol and a set of techniques like "hole punching" which allows peers to establish publically routable addresses and ports.

  • TURN: A way to bounce requests through a central server which can then be routed to the appropriate peer. (Basically workaround, exactly what we are trying to avoid: connection with server.)


45
New cards

What is a way to propagate updates to server that are not quite real-time?

With Simple Polling, we're using a pull-based model. Our client is constantly asking the server for updates and the server needs to maintain the state necessary to respond to those requests. The most common way to do this is to have a database where we can store the updates (e.g. all of the messages in the chat room), and from this database our clients can pull the updates they need when they can. For our chat app, we'd basically be polling for "what messages have been sent to the room with a timestamp larger than the last message I received?".

46
New cards

What methods can we use to send data in real time from our database to our server(s)?

Pushing via consistent hashes (with Zookeeper) or pub/sub

47
New cards

Describe consistent hashing as a method to send data in real time from our database to our servers

We always have 1 server who "owns" the connections for that user. To do this, we'll have a central service that knows how many servers there are N and can assign them each a number 0 through N-1. This is frequently Apache ZooKeeper or Etcd which allows us to manage this metadata and allows the servers to keep in sync as it updates, though in practice there are many alternatives.

With simple modulo hashing, changing the number of servers would require almost all users to disconnect and reconnect to different servers - an expensive operation that disrupts service.

Consistent hashing solves this by minimizing the number of connections that need to move when scaling. It maps both servers and users onto a hash ring, and each user connects to the next server they encounter when moving clockwise around the ring.

Consistent hashing is ideal when you need to maintain persistent connections (WebSocket/SSE) and your system needs to scale dynamically. It's particularly valuable when each connection requires significant server-side state that would be expensive to rebuild.

48
New cards

Describe Pub/Sub as a method to send data in real time from our database to our servers

In this model, we have a single service that is responsible for collecting updates from the source and then broadcasting them to all interested clients. Popular options here include Kafka and Redis.

The pub/sub service becomes the biggest source of state for our realtime updates. Our persistent connections are now made to lightweight servers which simply subscribe to the relevant topics and forward the updates to the appropriate clients. When clients connect, we don't need them to connect to a specific endpoint server (like we did with consistent hashing) and instead can connect to any of them. Once connected, the endpoint server will register the client with the pub/sub server so that any updates can be sent to them.

If you're using a pub/sub model, you'll probably need to talk about the single point of failure and bottleneck of the pub/sub service. Redis cluster is a popular way to scale pub/sub service which involves sharding the subscriptions by their key across multiple hosts. This scales up the number of subscriptions you can support and the throughput. For inbound connections to the endpoint servers, you'll probably want to use a load balancer with a "least connections" strategy. This will help ensure that you're distributing the load across the servers in the cluster.

49
New cards

How to handle connection failures and reconnection in a real-time system?

For recovery, you need to track what messages or updates a client has received. When they reconnect, they should get everything they missed. This often means maintaining a per-user message queue or implementing sequence numbers that clients can reference during reconnection. Using Redis streams for this is a popular option.

50
New cards

What happens when a single user has millions of followers who all need the same update in a real-time system?

The solution involves strategic caching and hierarchical distribution. Instead of writing the update to millions of individual user feeds, cache the update once and distribute through multiple layers. Regional servers can pull the update and push to their local clients, reducing the load on any single component.

51
New cards

How do you maintain message ordering across distributed servers in a real-time system?

Vector clocks or logical timestamps help establish ordering relationships between messages. Each server maintains its own clock, and messages include timestamp information that helps recipients determine the correct order.

Most commonly in user-facing Apps, for critical ordering requirements, you might need to funnel all related messages through a single server or partition. This trades some scalability for consistency guarantees, but simplifies the ordering problem significantly.

52
New cards

Description of B-Trees

B-tree indexes are the most common type of database index, providing an efficient way to organize data for fast searches and updates. They achieve this by maintaining a balanced tree structure that minimizes the number of disk reads needed to find any piece of data.

  1. They maintain sorted order, making range queries and ORDER BY operations efficient

  2. They're self-balancing, ensuring predictable performance even as data grows

  3. They minimize disk I/O by matching their structure to how databases store data

  4. They handle both equality searches (email = 'x') and range searches (age > 25) equally well

  5. They remain balanced even with random inserts and deletes, avoiding the performance cliffs you might see with simpler tree structures

Postgres uses B-Trees

53
New cards

LSM trees description

With B-trees, each write means finding the right leaf page, reading it into memory, updating it, and writing it back to disk. For a few thousand writes per second, this works fine. But when you're processing 100,000 writes per second, those random disk seeks become a bottleneck. Instead of updating data in place like B-trees, LSM trees use an append-only approach that's built for write-heavy workloads. While LSM trees excel at writes, they make reads more complex.

Cassandra and DynamoDB use LSM trees.

54
New cards

Design a Messaging System: How can users be able to start group chats with multiple participants?

We'll start with a simple service behind an L4 load balancer (we're using Websockets) which can write Chatmetadata to a database.

  1. User connects to the service and sends a createChat message.

  2. The service creates a Chat record in the database along with a ChatParticipant record for each user in the chat. For small chats this can be done in a single DynamoDB transaction (up to 100 items), but for chats near the 100-participant limit we may need to batch the writes.

  3. The service returns the chatId to the user.

For the ChatParticipant table, we'll want to be able to (1) look up all participants for a given chat and (2) look up all chats for a given user. We could do this with two composite keys.


55
New cards

Design a messaging system: How can users send/receive messages.

To send a message:

  1. User sends a sendMessage message to the Chat Server.

  2. The Chat Server looks up all participants in the chat via the ChatParticipant table.

  3. The Chat Server looks up the websocket connection for each participant in its internal hash table and sends the message via each connection.

Note: We're assuming all users are online, connected to the same Chat Server, and that we have a websocket connection for each of them. This is not scalable.

56
New cards

Design a Messaging system: Users should be able to receive messages sent while they are not online (up to 30 days).

We're going to need to start storing messages in our database so that we can deliver them to users even when they're offline.

Let's keep an "Inbox" for each user which will contain all undelivered messages. When messages are sent, we'll write them to the inbox of each recipient user. If they're already online, we can go ahead and try to deliver the message immediately. If they're not online, we'll store the message and wait for them to come back later.