Chat rooms
An end-to-end chat room: one shared protocol and SQL layer, the same room built twice (Socket.IO on Bun, and a Cloudflare Durable Object with hibernation), and a React hook. Background: Realtime fundamentals, Socket.IO, Durable Objects, Realtime in React & Next.js and the socket itself in WebSockets.
Architecture
| Piece | Job | Runs in |
|---|---|---|
protocol.ts | Zod union for client → server, TS union for server → client | both servers and the browser |
token.ts | sign and verify a short-lived HMAC ticket | Next.js route, both servers |
store.ts | the SQL: insert, page, edit, delete, read pointers, rate counts | both servers (both are SQLite) |
text.ts | sanitize text; rate and typing constants | both servers, client throttle |
| Server A | Socket.IO rooms, bun:sqlite file | one Bun process |
| Server B | Worker (auth) + one ChatRoom object per room | Cloudflare |
use-chat.ts | reducer + socket + optimistic sends | React |
chat/shared/ # imported by every other folderprotocol.tstoken.tsstore.tstext.tssio/server.ts # Server Ado/ # Server Bsrc/worker.tssrc/room.tswrangler.jsoncreact/ # client| Pick | When |
|---|---|
| Socket.IO on Bun | you already run a Node/Bun server; want long-polling fallback, acks, namespaces; one region is fine |
| Durable Object | many independent rooms, global users, no servers to run; idle rooms should cost nothing |
Both servers speak the same JSON messages, so the React client only changes transport.
Protocol
Clients never send a user id, message id or timestamp: the server derives the user from the token and assigns
the rest. clientId (a UUID made by the sender) is the idempotency key for resends.
import { z } from "zod";
export const MAX_TEXT = 2000;
export const PAGE = 50;
const Room = z.string().regex(/^[a-z0-9-]{1,64}$/);
const Id = z.number().int().positive();
const Text = z.string().trim().min(1).max(MAX_TEXT);
// Client -> server. Never carries a user id.
export const ClientMsg = z.discriminatedUnion("type", [
z.object({ type: z.literal("join"), room: Room }),
z.object({ type: z.literal("leave"), room: Room }),
z.object({
type: z.literal("send"),
room: Room,
clientId: z.uuid(), // idempotency key
text: Text,
}),
z.object({
type: z.literal("edit"),
room: Room,
id: Id,
text: Text,
}),
z.object({
type: z.literal("delete"),
room: Room,
id: Id,
}),
z.object({ type: z.literal("typing"), room: Room }),
z.object({
type: z.literal("read"),
room: Room,
upTo: Id,
}),
z.object({
type: z.literal("history"),
room: Room,
before: Id,
limit: z.number().int().min(1).max(100).default(PAGE),
}),
]);
export type ClientMsg = z.infer<typeof ClientMsg>;
export type ClientMsgIn = z.input<typeof ClientMsg>;
export type User = { id: string; name: string };
export type ChatMessage = {
id: number; // server-assigned, increasing per room
room: string;
userId: string;
name: string;
text: string; // plain text, never HTML
ts: number; // server clock, ms
editedAt: number | null;
deleted: boolean;
clientId: string;
};
// Server -> client.
export type ServerMsg =
| {
type: "joined";
room: string;
me: User;
history: ChatMessage[];
hasMore: boolean;
members: User[];
reads: Record<string, number>;
}
| { type: "presence"; room: string; members: User[] }
| { type: "message"; msg: ChatMessage }
| { type: "edited"; msg: ChatMessage }
| { type: "deleted"; room: string; id: number }
| {
type: "typing";
room: string;
user: User;
until: number;
}
| {
type: "read";
room: string;
userId: string;
upTo: number;
}
| {
type: "history";
room: string;
msgs: ChatMessage[];
hasMore: boolean;
}
| { type: "error"; code: ErrorCode; retryIn?: number };
export type ErrorCode =
| "bad_message"
| "forbidden"
| "rate_limited"
| "not_found";
export function parseClientMsg(
raw: unknown,
): ClientMsg | null {
if (typeof raw !== "string") return null;
try {
const r = ClientMsg.safeParse(JSON.parse(raw));
return r.success ? r.data : null;
} catch {
return null; // not JSON
}
}| Client sends | Server answers | Who receives it |
|---|---|---|
join (DO: implicit on connect) | joined: me, last 50 messages, hasMore, members, read pointers | the joiner |
| (join, or last tab closes) | presence: current members | whole room |
send | message with server id, ts and the clientId | whole room, sender included (confirms it) |
edit / delete | edited / deleted | whole room |
typing | typing with until | room minus sender |
read | read: user's pointer | whole room |
history with before | history page, oldest first | the requester |
| anything invalid | error: bad_message, forbidden, rate_limited (with retryIn) | the sender |
- Validate on the server with Zod; on the client a TS type is enough (you trust your server).
- Envelopes, versioning and sequencing in general: fundamentals: Message envelopes.
Server rules
| Rule | How | Why |
|---|---|---|
| Identity from the handshake | verified token → socket.data.user or the socket attachment | a client-sent userId is forgeable |
| Server ids | INTEGER PRIMARY KEY AUTOINCREMENT | one order per room; ids increase but may skip |
| Server time | Date.now() on insert | client clocks drift and lie |
| Idempotent sends | UNIQUE (user_id, client_id) + ON CONFLICT DO NOTHING | a resend after reconnect returns the original |
| Membership | Socket.IO: socket.rooms.has(room); DO: message room must equal the socket's room | no posting into rooms you never joined |
| Size | frame cap, Zod max(2000) after trim() | the DO runtime accepts 32 MiB frames |
| Author-only edit/delete | WHERE user_id = ? in the UPDATE | no ownership check in app code to forget |
| Rate limit | count of the user's messages in the last 10 s, from SQL | survives restarts and hibernation |
More on trust boundaries: fundamentals: Authority & trust.
Auth
Browsers can't set headers on a WebSocket, so the page fetches a short-lived ticket from the app (session cookie), then presents it in the handshake. Background: WebSockets: Authentication.
import { z } from "zod";
import type { User } from "./protocol";
// Minimal HMAC-signed ticket that works in Bun, Node and
// Workers. In production use your auth provider's JWTs.
const Claims = z.object({
sub: z.string().min(1).max(64),
name: z.string().min(1).max(40),
exp: z.number(), // ms since epoch
});
type Claims = z.infer<typeof Claims>;
const enc = new TextEncoder();
const b64u = (b: Uint8Array) =>
btoa(String.fromCharCode(...b))
.replace(/\+/g, "-").replace(/\//g, "_")
.replace(/=+$/, "");
const unb64u = (s: string) =>
Uint8Array.from(
atob(s.replace(/-/g, "+").replace(/_/g, "/")),
(c) => c.charCodeAt(0),
);
const key = (secret: string) =>
crypto.subtle.importKey(
"raw", enc.encode(secret),
{ name: "HMAC", hash: "SHA-256" },
false, ["sign", "verify"],
);
export async function signToken(
c: Claims, secret: string,
): Promise<string> {
const body = b64u(enc.encode(JSON.stringify(c)));
const sig = await crypto.subtle.sign(
"HMAC", await key(secret), enc.encode(body),
);
return `${body}.${b64u(new Uint8Array(sig))}`;
}
export async function verifyToken(
token: string, secret: string,
): Promise<User | null> {
const [body, sig, extra] = token.split(".");
if (!body || !sig || extra !== undefined) return null;
try {
const ok = await crypto.subtle.verify(
"HMAC", await key(secret), unb64u(sig),
enc.encode(body),
);
if (!ok) return null;
const json = new TextDecoder().decode(unb64u(body));
const c = Claims.parse(JSON.parse(json));
if (c.exp < Date.now()) return null;
return { id: c.sub, name: c.name };
} catch {
return null; // malformed base64, JSON or claims
}
}import { currentUser } from "@/lib/auth";
import { signToken } from "@/lib/chat/token";
const TTL_MS = 60_000; // only has to outlive the handshake
// POST, same-origin, session cookie: returns a ticket
// for the WebSocket URL or the Socket.IO auth field.
export async function POST(): Promise<Response> {
const secret = process.env.CHAT_SECRET;
if (!secret) throw new Error("CHAT_SECRET not set");
const user = await currentUser();
if (!user) return new Response(null, { status: 401 });
const exp = Date.now() + TTL_MS;
const token = await signToken(
{ sub: user.id, name: user.name, exp },
secret,
);
return Response.json(
{ token },
{ headers: { "Cache-Control": "no-store" } },
);
}| Transport | Ticket goes in | Checked by |
|---|---|---|
| Socket.IO | io(url, { auth: { token } }) → socket.handshake.auth.token | io.use() middleware, once per connection |
| WebSocket to a DO | ?token= in the URL | the Worker, before the (billed) object request |
- Short TTL (60 s): query strings end up in logs. Fetch a new ticket for every reconnect.
- Check
Originon WebSocket upgrades: CORS doesn't cover them. Socket.IO: the engine'sallowRequest(itscorsoption covers long-polling only); Worker: compare toAPP_ORIGIN. - The Worker forwards the verified user in an
X-Userheader it overwrites, so clients can't inject one. The object is reachable only through the Worker.
Storage & history
Both servers use SQLite (bun:sqlite and the object's database), so one module holds all the SQL behind a
two-method adapter.
import { PAGE } from "./protocol";
import type { ChatMessage, User } from "./protocol";
// The two drivers (bun:sqlite, Durable Object SQL) behind
// one tiny interface, so both servers share these queries.
export type SqlArg = string | number | null;
export interface Sql {
exec(script: string): void;
all<T>(query: string, ...args: SqlArg[]): T[];
}
const MAX_ID = Number.MAX_SAFE_INTEGER;
export const SCHEMA = `
CREATE TABLE IF NOT EXISTS messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
room TEXT NOT NULL,
user_id TEXT NOT NULL,
name TEXT NOT NULL,
text TEXT NOT NULL,
ts INTEGER NOT NULL,
edited_at INTEGER,
deleted INTEGER NOT NULL DEFAULT 0,
client_id TEXT NOT NULL,
UNIQUE (user_id, client_id)
);
CREATE INDEX IF NOT EXISTS msg_room
ON messages (room, id);
CREATE INDEX IF NOT EXISTS msg_user
ON messages (user_id, ts);
CREATE TABLE IF NOT EXISTS reads (
room TEXT NOT NULL,
user_id TEXT NOT NULL,
up_to INTEGER NOT NULL,
PRIMARY KEY (room, user_id)
);`;
type Row = {
id: number;
room: string;
user_id: string;
name: string;
text: string;
ts: number;
edited_at: number | null;
deleted: number;
client_id: string;
};
const toMsg = (r: Row): ChatMessage => ({
id: r.id,
room: r.room,
userId: r.user_id,
name: r.name,
text: r.text,
ts: r.ts,
editedAt: r.edited_at,
deleted: r.deleted === 1,
clientId: r.client_id,
});
type NewMsg = {
room: string;
user: User;
text: string;
clientId: string;
};
export function createStore(sql: Sql) {
sql.exec(SCHEMA);
const n = (q: string, ...a: SqlArg[]) =>
sql.all<{ n: number }>(q, ...a)[0]?.n ?? 0;
return {
// Idempotent: a resend with the same clientId
// returns the row stored the first time.
insert(m: NewMsg, now = Date.now()): ChatMessage {
sql.all(
`INSERT INTO messages
(room, user_id, name, text, ts, client_id)
VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT (user_id, client_id) DO NOTHING`,
m.room, m.user.id, m.user.name, m.text, now,
m.clientId,
);
const [row] = sql.all<Row>(
`SELECT * FROM messages
WHERE user_id = ? AND client_id = ?`,
m.user.id, m.clientId,
);
if (!row) throw new Error("insert failed");
return toMsg(row);
},
// Keyset pagination: newest first in SQL, oldest
// first in the result. `before` is an exclusive id.
page(room: string, before = MAX_ID, limit = PAGE) {
const rows = sql.all<Row>(
`SELECT * FROM messages
WHERE room = ? AND id < ?
ORDER BY id DESC LIMIT ?`,
room, before, limit + 1,
);
return {
msgs: rows.slice(0, limit).reverse().map(toMsg),
hasMore: rows.length > limit,
};
},
// Only the author may edit; returns null otherwise.
edit(room: string, userId: string, id: number,
text: string, now = Date.now()): ChatMessage | null {
const [row] = sql.all<Row>(
`UPDATE messages SET text = ?, edited_at = ?
WHERE id = ? AND room = ? AND user_id = ?
AND deleted = 0
RETURNING *`,
text, now, id, room, userId,
);
return row ? toMsg(row) : null;
},
// Soft delete: keeps the id so pages stay stable.
remove(room: string, userId: string, id: number) {
return sql.all(
`UPDATE messages SET text = '', deleted = 1
WHERE id = ? AND room = ? AND user_id = ?
AND deleted = 0
RETURNING id`,
id, room, userId,
).length === 1;
},
// Read pointers only move forward.
markRead(room: string, userId: string, upTo: number) {
return n(
`INSERT INTO reads (room, user_id, up_to)
VALUES (?, ?, ?)
ON CONFLICT (room, user_id) DO UPDATE
SET up_to = MAX(up_to, excluded.up_to)
RETURNING up_to AS n`,
room, userId, upTo,
);
},
reads(room: string): Record<string, number> {
type R = { user_id: string; up_to: number };
const rows = sql.all<R>(
"SELECT user_id, up_to FROM reads WHERE room = ?",
room,
);
return Object.fromEntries(
rows.map((r) => [r.user_id, r.up_to]),
);
},
unread(room: string, userId: string) {
return n(
`SELECT COUNT(*) AS n FROM messages
WHERE room = ? AND user_id != ? AND deleted = 0
AND id > COALESCE((SELECT up_to FROM reads
WHERE room = ? AND user_id = ?), 0)`,
room, userId, room, userId,
);
},
// For slow mode: this user's newest message here.
lastSent(room: string, userId: string): number | null {
return sql.all<{ t: number | null }>(
`SELECT MAX(ts) AS t FROM messages
WHERE room = ? AND user_id = ?`,
room, userId,
)[0]?.t ?? null;
},
// For rate limits: messages by a user since `since`.
sentSince(userId: string, since: number) {
return n(
`SELECT COUNT(*) AS n FROM messages
WHERE user_id = ? AND ts > ?`,
userId, since,
);
},
};
}
export type Store = ReturnType<typeof createStore>;| Concern | Choice |
|---|---|
| History on join | last 50 (page(room)), plus hasMore |
| Older pages | keyset: id < before ORDER BY id DESC LIMIT n + 1; the extra row means hasMore |
Why not OFFSET | new messages shift offsets (duplicates or gaps) and OFFSET scans skipped rows |
| Deletes | soft: blank text, deleted = 1; ids stay stable and moderators keep an audit trail |
| Read receipts | one pointer per user per room, only moving forward (MAX) |
| Indexes | (room, id) for pages, (user_id, ts) for rate limits |
| Retention | prune with an alarm or cron: DELETE … WHERE ts < ? |
Server A: Socket.IO on Bun
@socket.io/bun-engine puts Socket.IO on Bun.serve. One event name, msg, carries the shared union so both
servers share the protocol (named events work too: Socket.IO: Typed events).
import { Database } from "bun:sqlite";
import { Server as Engine } from "@socket.io/bun-engine";
import { Server, type Socket } from "socket.io";
import {
ClientMsg,
type ErrorCode,
type ServerMsg,
type User,
} from "../shared/protocol";
import { createStore, type Sql } from "../shared/store";
import {
cleanText,
RATE,
TYPING_GAP,
TYPING_TTL,
} from "../shared/text";
import { verifyToken } from "../shared/token";
const SECRET = Bun.env.CHAT_SECRET;
if (!SECRET) throw new Error("CHAT_SECRET not set");
const db = new Database(Bun.env.CHAT_DB ?? "chat.sqlite");
db.run("PRAGMA journal_mode = WAL");
const sql: Sql = {
exec: (script) => void db.run(script),
all: (q, ...args) => db.query(q).all(...args) as never,
};
const store = createStore(sql);
// One event each way, carrying the shared Zod union.
type C2S = { msg: (m: unknown) => void };
type S2C = { msg: (m: ServerMsg) => void };
type Data = { user: User; typedAt: Map<string, number> };
type ChatSocket = Socket<C2S, S2C, {}, Data>;
const ORIGIN = Bun.env.APP_ORIGIN ?? "http://localhost:5173";
const io = new Server<C2S, S2C, {}, Data>();
const engine = new Engine({
path: "/socket.io/",
maxHttpBufferSize: 16 * 1024, // bytes per message
cors: { origin: ORIGIN }, // HTTP long-polling only
// WebSocket upgrades skip CORS: check Origin here.
allowRequest: async (req) => {
const origin = req.headers.get("origin");
if (origin && origin !== ORIGIN) throw "bad origin";
},
});
io.bind(engine);
const key = (room: string) => `room:${room}`;
const emit = (s: ChatSocket, m: ServerMsg) =>
void s.emit("msg", m);
const fail = (s: ChatSocket, code: ErrorCode) =>
emit(s, { type: "error", code });
// Auth once, in the handshake. The user comes from the
// verified token, never from anything the client sends.
io.use(async (socket, next) => {
const token: unknown = socket.handshake.auth.token;
const user = typeof token === "string"
? await verifyToken(token, SECRET)
: null;
if (!user) return next(new Error("unauthorized"));
socket.data.user = user;
socket.data.typedAt = new Map();
next();
});
async function members(room: string, skip?: string) {
const all = await io.in(key(room)).fetchSockets();
const byId = new Map<string, User>();
for (const s of all) {
if (s.id !== skip) byId.set(s.data.user.id, s.data.user);
}
return [...byId.values()]; // one entry per user, not tab
}
async function presence(room: string, skip?: string) {
io.to(key(room)).emit("msg", {
type: "presence",
room,
members: await members(room, skip),
});
}
async function onMsg(socket: ChatSocket, m: ClientMsg) {
const { user } = socket.data;
const room = key(m.room);
const now = Date.now();
if (m.type === "join") {
await socket.join(room);
const { msgs, hasMore } = store.page(m.room);
emit(socket, {
type: "joined",
room: m.room,
me: user,
history: msgs,
hasMore,
members: await members(m.room),
reads: store.reads(m.room),
});
return presence(m.room);
}
if (!socket.rooms.has(room)) {
return fail(socket, "forbidden");
}
switch (m.type) {
case "leave":
await socket.leave(room);
return presence(m.room);
case "send": {
const since = now - RATE.windowMs;
if (store.sentSince(user.id, since) >= RATE.max) {
return emit(socket, {
type: "error",
code: "rate_limited",
retryIn: RATE.windowMs,
});
}
const text = cleanText(m.text);
if (!text) {
return fail(socket, "bad_message");
}
const msg = store.insert({
room: m.room, user, text, clientId: m.clientId,
});
io.to(room).emit("msg", { type: "message", msg });
return;
}
case "edit": {
const text = cleanText(m.text);
const msg = text
? store.edit(m.room, user.id, m.id, text)
: null;
if (!msg) {
return fail(socket, "forbidden");
}
io.to(room).emit("msg", { type: "edited", msg });
return;
}
case "delete":
if (!store.remove(m.room, user.id, m.id)) {
return fail(socket, "forbidden");
}
io.to(room).emit("msg", {
type: "deleted", room: m.room, id: m.id,
});
return;
case "typing": {
const last = socket.data.typedAt.get(room) ?? 0;
if (now - last < TYPING_GAP) return; // throttle
socket.data.typedAt.set(room, now);
socket.to(room).emit("msg", { // everyone but me
type: "typing",
room: m.room,
user,
until: now + TYPING_TTL,
});
return;
}
case "read": {
const upTo = store.markRead(m.room, user.id, m.upTo);
io.to(room).emit("msg", {
type: "read", room: m.room, userId: user.id, upTo,
});
return;
}
case "history": {
const page = store.page(m.room, m.before, m.limit);
return emit(socket, {
type: "history", room: m.room, ...page,
});
}
default:
m satisfies never;
}
}
io.on("connection", (socket) => {
socket.on("msg", (raw) => {
const parsed = ClientMsg.safeParse(raw);
if (!parsed.success) {
return fail(socket, "bad_message");
}
onMsg(socket, parsed.data).catch((err) => {
console.error("chat handler failed", err);
fail(socket, "bad_message");
});
});
// Still in its rooms here; gone by "disconnect".
socket.on("disconnecting", () => {
for (const room of socket.rooms) {
if (!room.startsWith("room:")) continue;
void presence(room.slice(5), socket.id);
}
});
});
export default {
port: Number(Bun.env.PORT ?? 3000),
...engine.handler(),
};bun add socket.io @socket.io/bun-engine zod
CHAT_SECRET=dev-secret bun sio/server.ts- Client:
io(url, { auth: { token }, transports: ["websocket"] }), thensocket.emit("msg", …)andsocket.on("msg", …). - Presence comes from
fetchSockets(), which also works across nodes with an adapter.bun:sqlitedoes not: for several nodes move the store to Postgres (Socket.IO: Scaling with adapters). disconnectingstill sees the socket's rooms; bydisconnectthey're gone.
Server B: Durable Object
The Worker authenticates, then routes /rooms/:room/ws to the object named room. The object keeps no state
in memory: presence comes from ctx.getWebSockets(), per-socket state from attachments, the rest from SQLite.
import { verifyToken } from "../../shared/token";
import { ChatRoom } from "./room";
export { ChatRoom }; // the runtime needs the class export
export interface Env {
CHAT: DurableObjectNamespace<ChatRoom>;
CHAT_SECRET: string; // .dev.vars locally, secret in prod
APP_ORIGIN: string;
}
const ROUTE = /^\/rooms\/([a-z0-9-]{1,64})\/ws$/;
const deny = (status: number, text: string) =>
new Response(text, { status });
export default {
async fetch(req, env): Promise<Response> {
const url = new URL(req.url);
const room = ROUTE.exec(url.pathname)?.[1];
if (!room) return deny(404, "Not found");
if (req.headers.get("Upgrade") !== "websocket") {
return deny(426, "Expected WebSocket");
}
// Reject cross-site pages (browsers always send Origin).
const origin = req.headers.get("Origin");
if (origin && origin !== env.APP_ORIGIN) {
return deny(403, "Bad origin");
}
// Browsers can't set headers on a WebSocket: short-lived
// ticket in the query string, verified here, before
// the (billed) Durable Object request.
const token = url.searchParams.get("token") ?? "";
const user = await verifyToken(token, env.CHAT_SECRET);
if (!user) return deny(401, "Unauthorized");
const headers = new Headers(req.headers);
headers.set("X-User", JSON.stringify(user)); // overwrite
const stub = env.CHAT.getByName(room); // one per room
return stub.fetch(new Request(req, { headers }));
},
} satisfies ExportedHandler<Env>;import { DurableObject } from "cloudflare:workers";
import {
parseClientMsg,
type ClientMsg,
type ErrorCode,
type ServerMsg,
type User,
} from "../../shared/protocol";
import { createStore, type Store } from "../../shared/store";
import {
cleanText,
RATE,
TYPING_GAP,
TYPING_TTL,
} from "../../shared/text";
import type { Env } from "./worker";
// Per-socket state that must survive hibernation
// (serialized, max 16 KiB).
type Att = { user: User; room: string; typedAt: number };
const MAX_FRAME = 16 * 1024; // chars
export class ChatRoom extends DurableObject<Env> {
store: Store;
// Runs on first use AND after every hibernation wake.
constructor(ctx: DurableObjectState, env: Env) {
super(ctx, env);
const sql = ctx.storage.sql;
this.store = createStore({
exec: (script) => void sql.exec(script),
all: (q, ...args) =>
sql.exec(q, ...args).toArray() as never,
});
// Answered at the edge without waking the object.
ctx.setWebSocketAutoResponse(
new WebSocketRequestResponsePair("ping", "pong"),
);
}
// The Worker has already authenticated the user.
async fetch(req: Request): Promise<Response> {
const user = JSON.parse(
req.headers.get("X-User") ?? "null",
) as User | null;
const room = this.ctx.id.name; // set by getByName()
if (!user || !room) {
return new Response("Forbidden", { status: 403 });
}
const pair = new WebSocketPair();
const [client, server] = [pair[0], pair[1]];
this.ctx.acceptWebSocket(server, [user.id]); // tag
const att: Att = { user, room, typedAt: 0 };
server.serializeAttachment(att);
this.send(server, {
type: "joined", room, me: user,
...this.page(room), members: this.members(),
reads: this.store.reads(room),
});
// Announce only the user's first tab.
if (this.ctx.getWebSockets(user.id).length === 1) {
this.presence(room);
}
return new Response(null, {
status: 101,
webSocket: client,
});
}
async webSocketMessage(
ws: WebSocket,
raw: string | ArrayBuffer,
) {
const att = ws.deserializeAttachment() as Att;
// The runtime accepts up to 32 MiB: cap it first.
if (typeof raw === "string" && raw.length > MAX_FRAME) {
return this.fail(ws, "bad_message");
}
const m = parseClientMsg(raw);
if (!m || m.room !== att.room) {
return this.fail(ws, "bad_message");
}
this.handle(ws, att, m);
}
async webSocketClose(ws: WebSocket, code: number) {
const { user, room } = ws.deserializeAttachment() as Att;
ws.close(code, "bye"); // automatic from 2026-04-07
const others = this.ctx.getWebSockets(user.id);
if (others.every((s) => s === ws)) {
this.presence(room, ws); // last tab closed
}
}
async webSocketError(ws: WebSocket, err: unknown) {
console.error("socket error", err);
}
handle(ws: WebSocket, att: Att, m: ClientMsg) {
const { user, room } = att;
const now = Date.now();
switch (m.type) {
case "join": // resync after a client-side hiccup
return this.send(ws, {
type: "joined", room, me: user,
...this.page(room), members: this.members(),
reads: this.store.reads(room),
});
case "leave":
return ws.close(1000, "left");
case "send": {
const since = now - RATE.windowMs;
const sent = this.store.sentSince(user.id, since);
if (sent >= RATE.max) {
return this.send(ws, {
type: "error",
code: "rate_limited",
retryIn: RATE.windowMs,
});
}
const wait = this.slowModeWait(room, user.id, now);
if (wait > 0) {
return this.send(ws, {
type: "error",
code: "rate_limited",
retryIn: wait,
});
}
const text = cleanText(m.text);
if (!text) return this.fail(ws, "bad_message");
const msg = this.store.insert({
room, user, text, clientId: m.clientId,
});
return this.broadcast({ type: "message", msg });
}
case "edit": {
const text = cleanText(m.text);
const msg = text
? this.store.edit(room, user.id, m.id, text)
: null;
if (!msg) return this.fail(ws, "forbidden");
return this.broadcast({ type: "edited", msg });
}
case "delete":
if (!this.store.remove(room, user.id, m.id)) {
return this.fail(ws, "forbidden");
}
return this.broadcast({
type: "deleted", room, id: m.id,
});
case "typing":
if (now - att.typedAt < TYPING_GAP) return;
ws.serializeAttachment({ ...att, typedAt: now });
return this.broadcast({
type: "typing", room, user,
until: now + TYPING_TTL,
}, ws);
case "read":
return this.broadcast({
type: "read", room, userId: user.id,
upTo: this.store.markRead(room, user.id, m.upTo),
});
case "history":
return this.send(ws, {
type: "history", room,
...this.store.page(room, m.before, m.limit),
});
default:
m satisfies never;
}
}
// RPC for moderators (any Worker with the binding):
// await env.CHAT.getByName("lobby").setSlowMode(30_000)
async setSlowMode(ms: number): Promise<void> {
this.ctx.storage.kv.put("slowMs", Math.max(0, ms));
}
slowModeWait(room: string, userId: string, now: number) {
const kv = this.ctx.storage.kv;
const gap = kv.get<number>("slowMs") ?? 0;
if (gap === 0) return 0;
const last = this.store.lastSent(room, userId) ?? 0;
return Math.max(0, last + gap - now);
}
page(room: string) {
const { msgs, hasMore } = this.store.page(room);
return { history: msgs, hasMore };
}
// Online users, one entry per user (not per tab).
members(skip?: WebSocket): User[] {
const byId = new Map<string, User>();
for (const ws of this.ctx.getWebSockets()) {
if (ws === skip) continue;
const { user } = ws.deserializeAttachment() as Att;
byId.set(user.id, user);
}
return [...byId.values()];
}
presence(room: string, skip?: WebSocket) {
this.broadcast({
type: "presence", room, members: this.members(skip),
});
}
send(ws: WebSocket, m: ServerMsg) {
try {
ws.send(JSON.stringify(m));
} catch {
// closing or closed; webSocketClose cleans up
}
}
fail(ws: WebSocket, code: ErrorCode) {
this.send(ws, { type: "error", code });
}
broadcast(m: ServerMsg, skip?: WebSocket) {
const data = JSON.stringify(m);
for (const ws of this.ctx.getWebSockets()) {
if (ws === skip) continue;
try {
ws.send(data);
} catch {
// ignore sockets that are closing
}
}
}
}{
"$schema": "./node_modules/wrangler/config-schema.json",
"name": "chat",
"main": "src/worker.ts",
"compatibility_date": "2026-09-01",
"vars": { "APP_ORIGIN": "http://localhost:5173" },
"durable_objects": {
"bindings": [
{ "name": "CHAT", "class_name": "ChatRoom" }
]
},
"exports": {
"ChatRoom": {
"type": "durable-object",
"storage": "sqlite"
}
}
}bun add zod
echo 'CHAT_SECRET="dev-secret"' > do/.dev.vars
# serves ws://localhost:8787/rooms/:room/ws
cd do && bunx wrangler dev
bunx wrangler secret put CHAT_SECRET && bunx wrangler deploy| Socket.IO on Bun | Durable Object | |
|---|---|---|
| Room | socket.join("room:lobby") | one object per room name |
| Connection state | socket.data (in memory) | tags + serializeAttachment (survives hibernation) |
| Presence | io.in(room).fetchSockets() | ctx.getWebSockets() + attachments |
| Fan-out | io.to(room).emit / socket.to(room) | loop over getWebSockets() |
| Join | explicit join message | the connection is the join |
| Heartbeat | Engine.IO ping/pong | protocol pings plus "ping" → "pong" auto-response, both without waking |
| History | one bun:sqlite file for all rooms | each room's own database |
| Scale | vertical, then adapter + shared DB | automatic per room; one room ≈ one thread |
| Idle cost | the process keeps running | none while hibernated |
Typing, presence & receipts
| Signal | Stored | Sent to | Limits | Clears when |
|---|---|---|---|---|
| Presence | no: derived from open sockets | whole room | one entry per user, not per tab | last tab closes |
| Typing | no | room minus sender | client ≤ 1 per 2 s; server drops closer ones | until passes (5 s) or that user's message arrives |
| Read receipt | yes: pointer per user | whole room | moves forward only | never; unread = later messages by others |
| Message | yes | whole room | 5 per 10 s per user | edit, soft delete |
- Typing and presence are best-effort: losing one is harmless, so they skip storage and acks.
- Show receipts as "seen by" on the newest message each user's pointer has reached.
- Presence at scale (heartbeats, sweeps, diffs): fundamentals: Presence.
Moderation
// Control characters (except tab and newline) and bidi
// overrides, which can spoof how a line reads.
const CTRL = /(?![\n\t])\p{Cc}/gu;
const BIDI = /[\u202A-\u202E\u2066-\u2069]/g;
export function cleanText(input: string): string {
return input
.normalize("NFC")
.replace(CTRL, "")
.replace(BIDI, "")
.replace(/\n{3,}/g, "\n\n") // at most one blank line
.trim();
}
export const RATE = { max: 5, windowMs: 10_000 };
export const TYPING_GAP = 2_000; // min ms between relays
export const TYPING_TTL = 5_000; // client hides after this| Control | Here | Notes |
|---|---|---|
| Max length | Zod max(2000), frame cap, maxHttpBufferSize | the client maxLength is a courtesy only |
| Sanitize | NFC, strip control and bidi-override characters, collapse blank lines | keep < and >: they're harmless as text |
| Render as text | React text nodes, white-space: pre-wrap | never dangerouslySetInnerHTML; Markdown only with raw HTML off and http/https/mailto links |
| Rate limit | 5 per 10 s per user | reply rate_limited with retryIn; the client disables Send |
| Slow mode | per-room gap | see Recipes |
| Edit/delete | author only; soft delete | add a role claim for moderators |
| Mute or ban | not shown | a bans table checked at handshake and on send; close the user's sockets (DO: getWebSockets(userId)) |
| Images and links | not shown | upload to storage, send a URL; unfurl link previews server-side |
React client
A reducer holds the room; the hook owns the socket, reconnects with backoff and jitter, and resends unconfirmed
messages after each joined (safe: clientId dedupes). Sharing one socket across components with
useSyncExternalStore is in Realtime in React & Next.js.
import type {
ChatMessage,
ServerMsg,
User,
} from "../shared/protocol";
export type Pending = { clientId: string; text: string };
export type ChatState = {
online: boolean;
me: User | null;
msgs: ChatMessage[]; // ascending id
pending: Pending[]; // optimistic, not yet echoed
hasMore: boolean;
members: User[];
typing: Record<string, { name: string; until: number }>;
reads: Record<string, number>; // userId -> last read id
};
export type Action =
| { type: "online"; online: boolean }
| { type: "pending"; p: Pending }
| { type: "server"; m: ServerMsg };
export const initial: ChatState = {
online: false, me: null, msgs: [], pending: [],
hasMore: false, members: [], typing: {}, reads: {},
};
const upsert = (list: ChatMessage[], m: ChatMessage) =>
[...list.filter((x) => x.id !== m.id), m]
.sort((a, b) => a.id - b.id);
export function reducer(s: ChatState, a: Action): ChatState {
if (a.type === "online") return { ...s, online: a.online };
if (a.type === "pending") {
return { ...s, pending: [...s.pending, a.p] };
}
const m = a.m;
switch (m.type) {
case "joined":
return {
...s, me: m.me, msgs: m.history, hasMore: m.hasMore,
members: m.members, reads: m.reads,
};
case "presence":
return { ...s, members: m.members };
case "message": {
// A message ends that user's typing indicator.
const { [m.msg.userId]: _, ...typing } = s.typing;
return {
...s, typing,
msgs: upsert(s.msgs, m.msg), // dedupes resends
pending: s.pending.filter(
(p) => p.clientId !== m.msg.clientId,
),
};
}
case "edited":
return { ...s, msgs: upsert(s.msgs, m.msg) };
case "deleted":
return {
...s,
msgs: s.msgs.map((x) => x.id === m.id
? { ...x, text: "", deleted: true } : x),
};
case "typing":
return {
...s,
typing: {
...s.typing,
[m.user.id]: { name: m.user.name, until: m.until },
},
};
case "read":
return {
...s, reads: { ...s.reads, [m.userId]: m.upTo },
};
case "history": {
const have = new Set(s.msgs.map((x) => x.id));
const older = m.msgs.filter((x) => !have.has(x.id));
return {
...s,
hasMore: m.hasMore,
msgs: [...older, ...s.msgs],
};
}
case "error":
return s; // show a toast; rate_limited has retryIn
}
}import {
useCallback, useEffect, useReducer, useRef,
} from "react";
import type {
ClientMsgIn, ServerMsg,
} from "../shared/protocol";
import { initial, reducer } from "./chat-state";
// getToken must be stable (useCallback) or the socket
// reconnects on every render.
export function useChat(
room: string,
getToken: () => Promise<string>,
base = "", // e.g. "wss://chat.example.com"; "" = same host
) {
const [state, dispatch] = useReducer(reducer, initial);
const ws = useRef<WebSocket | null>(null);
const pending = useRef(state.pending);
pending.current = state.pending;
const post = useCallback((m: ClientMsgIn) => {
const s = ws.current;
if (s?.readyState === WebSocket.OPEN) {
s.send(JSON.stringify(m));
}
}, []);
useEffect(() => {
let stopped = false;
let tries = 0;
let timer: ReturnType<typeof setTimeout> | undefined;
const retry = () => {
if (stopped) return;
const ms = Math.min(30_000, 500 * 2 ** tries++);
timer = setTimeout(connect, ms * Math.random());
};
async function connect() {
let token: string;
try {
token = await getToken(); // fresh, short-lived
} catch {
return retry();
}
if (stopped) return;
const q = `?token=${encodeURIComponent(token)}`;
const url = `${base}/rooms/${room}/ws${q}`;
const s = new WebSocket(url);
ws.current = s;
s.onopen = () => {
tries = 0;
dispatch({ type: "online", online: true });
};
s.onmessage = (e: MessageEvent<string>) => {
if (e.data === "pong") return; // auto-response
const m = JSON.parse(e.data) as ServerMsg;
dispatch({ type: "server", m });
if (m.type !== "joined") return;
for (const p of pending.current) {
// Safe to resend: the server dedupes clientId.
post({ type: "send", room, ...p });
}
};
s.onclose = () => {
dispatch({ type: "online", online: false });
retry();
};
}
void connect();
return () => {
stopped = true;
clearTimeout(timer);
ws.current?.close(1000);
};
}, [room, getToken, post, base]);
const send = useCallback((text: string) => {
const p = { clientId: crypto.randomUUID(), text };
dispatch({ type: "pending", p });
post({ type: "send", room, ...p });
}, [room, post]);
const typing = useCallback(
() => post({ type: "typing", room }), [room, post]);
const markRead = useCallback((upTo: number) =>
post({ type: "read", room, upTo }), [room, post]);
const asked = useRef(0); // one request per page
const oldest = state.msgs[0]?.id;
const loadOlder = useCallback(() => {
if (!oldest || asked.current === oldest) return;
asked.current = oldest;
post({ type: "history", room, before: oldest });
}, [room, post, oldest]);
return { state, send, typing, markRead, loadOlder };
}import type { ChatMessage } from "../shared/protocol";
// Text only: React escapes {m.text}. Never pass chat text
// to dangerouslySetInnerHTML or a Markdown renderer
// that allows raw HTML.
export function Message({ m }: { m: ChatMessage }) {
if (m.deleted) return <li className="gone">deleted</li>;
return (
<li style={{ whiteSpace: "pre-wrap" }}>
<b>{m.name}</b> {m.text}
{m.editedAt !== null && <small> (edited)</small>}
</li>
);
}getToken:useCallback(() => fetch("/api/chat-token", { method: "POST" }).then(…), []).- Pending messages render from
state.pending(grayed out) until theirmessageecho arrives. - For the Socket.IO server, swap the
WebSocketforio()andsend/onmessageforemit/on("msg").
Recipes
Typing indicator
Throttle on the client, throttle again on the server, expire on the client.
import { useRef, useState } from "react";
import { MAX_TEXT } from "../shared/protocol";
import { TYPING_GAP } from "../shared/text";
type Props = {
send: (text: string) => void;
typing: () => void;
};
// Throttle "typing" on the client too: at most one
// event per TYPING_GAP while keys are pressed.
export function Composer({ send, typing }: Props) {
const [text, setText] = useState("");
const last = useRef(0);
return (
<form
onSubmit={(e) => {
e.preventDefault();
if (!text.trim()) return;
send(text);
setText("");
last.current = 0;
}}
>
<textarea
value={text}
maxLength={MAX_TEXT}
onChange={(e) => {
setText(e.target.value);
const now = Date.now();
if (now - last.current > TYPING_GAP) {
last.current = now;
typing();
}
}}
/>
<button type="submit">Send</button>
</form>
);
}import { useEffect, useState } from "react";
import type { ChatState } from "./chat-state";
// Re-render once a second so stale entries disappear
// even when no new event arrives.
export function TypingLine({ s }: { s: ChatState }) {
const [now, setNow] = useState(Date.now);
useEffect(() => {
const t = setInterval(() => setNow(Date.now()), 1000);
return () => clearInterval(t);
}, []);
const names = Object.entries(s.typing)
.filter(([id, t]) => id !== s.me?.id && t.until > now)
.map(([, t]) => t.name);
if (names.length === 0) return null;
const who = names.length > 2
? "Several people are"
: `${names.join(" and ")} ${
names.length > 1 ? "are" : "is"
}`;
return <p aria-live="polite">{who} typing…</p>;
}Unread count
Server side it's store.unread(room, userId) (for badges and notifications); client side, count past your pointer
and move the pointer when the newest message is on screen.
import { useEffect, type RefObject } from "react";
// Mark the newest message read once it is on screen and
// the tab is visible. The server only moves pointers up.
export function useReadMarker(
last: RefObject<HTMLElement | null>,
lastId: number | undefined,
markRead: (upTo: number) => void,
) {
useEffect(() => {
const el = last.current;
if (!el || lastId === undefined) return;
const io = new IntersectionObserver(([entry]) => {
if (!entry?.isIntersecting) return;
if (document.visibilityState !== "visible") return;
markRead(lastId);
});
io.observe(el);
return () => io.disconnect();
}, [last, lastId, markRead]);
}
export function unreadCount(
msgs: { id: number; userId: string; deleted: boolean }[],
meId: string,
readUpTo = 0,
) {
return msgs.filter((m) =>
m.id > readUpTo && m.userId !== meId && !m.deleted,
).length;
}History pagination
Load the previous page near the top, and keep the viewport still while it's prepended.
import {
useLayoutEffect, useRef, type RefObject,
} from "react";
// Prepending older messages must not move what the
// reader is looking at: shift scrollTop by the added
// height. First render lands at the bottom.
export function useKeepScroll(
box: RefObject<HTMLElement | null>,
firstId: number | undefined,
) {
const height = useRef(0);
const first = useRef<number | undefined>(undefined);
useLayoutEffect(() => {
const el = box.current;
if (!el) return;
if (firstId !== first.current) {
el.scrollTop += el.scrollHeight - height.current;
}
first.current = firstId;
height.current = el.scrollHeight;
}); // every render, before paint
}import { useRef } from "react";
import type { ChatState } from "./chat-state";
import { Message } from "./message";
import { useKeepScroll } from "./use-keep-scroll";
type Props = { s: ChatState; loadOlder: () => void };
export function MessageList({ s, loadOlder }: Props) {
const box = useRef<HTMLUListElement>(null);
useKeepScroll(box, s.msgs[0]?.id);
return (
<ul
ref={box}
style={{ overflowY: "auto", height: 400 }}
onScroll={(e) => {
const top = e.currentTarget.scrollTop;
if (top < 80 && s.hasMore) loadOlder();
}}
>
{s.msgs.map((m) => <Message key={m.id} m={m} />)}
</ul>
);
}Slow mode
A per-room minimum gap between one user's messages, set by a moderator over RPC and computed from stored messages,
so it survives restarts and hibernation. In ChatRoom (already in room.ts above):
// RPC for moderators (any Worker with the binding):
// await env.CHAT.getByName("lobby").setSlowMode(30_000)
async setSlowMode(ms: number): Promise<void> {
this.ctx.storage.kv.put("slowMs", Math.max(0, ms));
}
slowModeWait(room: string, userId: string, now: number) {
const kv = this.ctx.storage.kv;
const gap = kv.get<number>("slowMs") ?? 0;
if (gap === 0) return 0;
const last = this.store.lastSent(room, userId) ?? 0;
return Math.max(0, last + gap - now);
}const wait = this.slowModeWait(room, user.id, now);
if (wait > 0) {
return this.send(ws, {
type: "error",
code: "rate_limited",
retryIn: wait,
});
}store.lastSent is SELECT MAX(ts) … WHERE room = ? AND user_id = ?, served by the (user_id, ts) index.
Notify offline members
Debounce with the room's alarm: a burst of messages becomes one push per member. A queue consumer sends the Web Push (Notifications & Push).
import type { Store } from "../../shared/store";
export type PushJob = {
userId: string;
room: string;
unread: number;
};
const DELAY_MS = 30_000; // one push per burst
// Call after each insert. setAlarm replaces any alarm, so
// only arm when none is pending (a debounce).
export async function armNotify(ctx: DurableObjectState) {
if ((await ctx.storage.getAlarm()) === null) {
await ctx.storage.setAlarm(Date.now() + DELAY_MS);
}
}
// Call from alarm(): members with unread messages and no
// open socket. Anyone with a read pointer is a member.
export async function notifyOffline(
ctx: DurableObjectState,
store: Store,
room: string,
queue: Queue<PushJob>,
) {
const jobs = Object.keys(store.reads(room)).flatMap(
(userId) => {
if (ctx.getWebSockets(userId).length > 0) return [];
const unread = store.unread(room, userId);
return unread > 0 ? [{ userId, room, unread }] : [];
},
);
if (jobs.length > 0) {
await queue.sendBatch(jobs.map((body) => ({ body })));
}
}Wire-up in ChatRoom: await armNotify(this.ctx) after each insert; an alarm() method that calls
notifyOffline with this.ctx, this.store, the room name and this.env.PUSH (a queue binding in
wrangler.jsonc). Store a "notified up to" id per user so a member isn't pushed twice for the same messages.
References
- MDN: WebSocket (opens in a new tab), IntersectionObserver (opens in a new tab), Push API (opens in a new tab)
- Socket.IO v4 docs (opens in a new tab): rooms (opens in a new tab), middlewares (opens in a new tab), Bun engine (opens in a new tab)
- Cloudflare: Durable Objects WebSockets (opens in a new tab), SQLite storage API (opens in a new tab), Alarms (opens in a new tab)
- Cloudflare Edge Chat Demo (opens in a new tab): the original DO chat example
- Bun: SQLite (opens in a new tab)
- Zod (opens in a new tab): schemas and
discriminatedUnion - OWASP: Cross-Site Scripting Prevention (opens in a new tab)
- Trojan Source (CVE-2021-42574) (opens in a new tab): why bidi overrides are stripped
- use-the-index-luke: Paging through results (opens in a new tab): keyset pagination