../

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

PieceJobRuns in
protocol.tsZod union for client → server, TS union for server → clientboth servers and the browser
token.tssign and verify a short-lived HMAC ticketNext.js route, both servers
store.tsthe SQL: insert, page, edit, delete, read pointers, rate countsboth servers (both are SQLite)
text.tssanitize text; rate and typing constantsboth servers, client throttle
Server ASocket.IO rooms, bun:sqlite fileone Bun process
Server BWorker (auth) + one ChatRoom object per roomCloudflare
use-chat.tsreducer + socket + optimistic sendsReact
chat/
chat/shared/        # imported by every other folderprotocol.tstoken.tsstore.tstext.tssio/server.ts  # Server Ado/            # Server Bsrc/worker.tssrc/room.tswrangler.jsoncreact/         # client
PickWhen
Socket.IO on Bunyou already run a Node/Bun server; want long-polling fallback, acks, namespaces; one region is fine
Durable Objectmany 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.

shared/protocol.ts
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 sendsServer answersWho receives it
join (DO: implicit on connect)joined: me, last 50 messages, hasMore, members, read pointersthe joiner
(join, or last tab closes)presence: current memberswhole room
sendmessage with server id, ts and the clientIdwhole room, sender included (confirms it)
edit / deleteedited / deletedwhole room
typingtyping with untilroom minus sender
readread: user's pointerwhole room
history with beforehistory page, oldest firstthe requester
anything invaliderror: 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

RuleHowWhy
Identity from the handshakeverified token → socket.data.user or the socket attachmenta client-sent userId is forgeable
Server idsINTEGER PRIMARY KEY AUTOINCREMENTone order per room; ids increase but may skip
Server timeDate.now() on insertclient clocks drift and lie
Idempotent sendsUNIQUE (user_id, client_id) + ON CONFLICT DO NOTHINGa resend after reconnect returns the original
MembershipSocket.IO: socket.rooms.has(room); DO: message room must equal the socket's roomno posting into rooms you never joined
Sizeframe cap, Zod max(2000) after trim()the DO runtime accepts 32 MiB frames
Author-only edit/deleteWHERE user_id = ? in the UPDATEno ownership check in app code to forget
Rate limitcount of the user's messages in the last 10 s, from SQLsurvives 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.

shared/token.ts
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
  }
}
app/api/chat-token/route.ts
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" } },
  );
}
TransportTicket goes inChecked by
Socket.IOio(url, { auth: { token } }) → socket.handshake.auth.tokenio.use() middleware, once per connection
WebSocket to a DO?token= in the URLthe Worker, before the (billed) object request
  • Short TTL (60 s): query strings end up in logs. Fetch a new ticket for every reconnect.
  • Check Origin on WebSocket upgrades: CORS doesn't cover them. Socket.IO: the engine's allowRequest (its cors option covers long-polling only); Worker: compare to APP_ORIGIN.
  • The Worker forwards the verified user in an X-User header 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.

shared/store.ts
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>;
ConcernChoice
History on joinlast 50 (page(room)), plus hasMore
Older pageskeyset: id < before ORDER BY id DESC LIMIT n + 1; the extra row means hasMore
Why not OFFSETnew messages shift offsets (duplicates or gaps) and OFFSET scans skipped rows
Deletessoft: blank text, deleted = 1; ids stay stable and moderators keep an audit trail
Read receiptsone pointer per user per room, only moving forward (MAX)
Indexes(room, id) for pages, (user_id, ts) for rate limits
Retentionprune 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).

sio/server.ts
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"] }), then socket.emit("msg", …) and socket.on("msg", …).
  • Presence comes from fetchSockets(), which also works across nodes with an adapter. bun:sqlite does not: for several nodes move the store to Postgres (Socket.IO: Scaling with adapters).
  • disconnecting still sees the socket's rooms; by disconnect they'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.

do/src/worker.ts
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>;
do/src/room.ts
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
      }
    }
  }
}
do/wrangler.jsonc
{
  "$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 BunDurable Object
Roomsocket.join("room:lobby")one object per room name
Connection statesocket.data (in memory)tags + serializeAttachment (survives hibernation)
Presenceio.in(room).fetchSockets()ctx.getWebSockets() + attachments
Fan-outio.to(room).emit / socket.to(room)loop over getWebSockets()
Joinexplicit join messagethe connection is the join
HeartbeatEngine.IO ping/pongprotocol pings plus "ping" → "pong" auto-response, both without waking
Historyone bun:sqlite file for all roomseach room's own database
Scalevertical, then adapter + shared DBautomatic per room; one room ≈ one thread
Idle costthe process keeps runningnone while hibernated

Typing, presence & receipts

SignalStoredSent toLimitsClears when
Presenceno: derived from open socketswhole roomone entry per user, not per tablast tab closes
Typingnoroom minus senderclient ≤ 1 per 2 s; server drops closer onesuntil passes (5 s) or that user's message arrives
Read receiptyes: pointer per userwhole roommoves forward onlynever; unread = later messages by others
Messageyeswhole room5 per 10 s per useredit, 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

shared/text.ts
// 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
ControlHereNotes
Max lengthZod max(2000), frame cap, maxHttpBufferSizethe client maxLength is a courtesy only
SanitizeNFC, strip control and bidi-override characters, collapse blank lineskeep < and >: they're harmless as text
Render as textReact text nodes, white-space: pre-wrapnever dangerouslySetInnerHTML; Markdown only with raw HTML off and http/https/mailto links
Rate limit5 per 10 s per userreply rate_limited with retryIn; the client disables Send
Slow modeper-room gapsee Recipes
Edit/deleteauthor only; soft deleteadd a role claim for moderators
Mute or bannot showna bans table checked at handshake and on send; close the user's sockets (DO: getWebSockets(userId))
Images and linksnot shownupload 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.

react/chat-state.ts
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
  }
}
react/use-chat.ts
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 };
}
react/message.tsx
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 their message echo arrives.
  • For the Socket.IO server, swap the WebSocket for io() and send/onmessage for emit/on("msg").

Recipes

Typing indicator

Throttle on the client, throttle again on the server, expire on the client.

react/composer.tsx
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>
  );
}
react/typing-line.tsx
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.

react/use-read-marker.ts
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.

react/use-keep-scroll.ts
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
}
react/message-list.tsx
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):

do/src/room.ts (methods)
// 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);
}
do/src/room.ts (in the send case)
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).

do/src/notify.ts
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