Skip to content
516 linesCodeBlameRaw
1/**
2 * One page's live room, as a Durable Object named by the page id.
3 *
4 * It owns the page's Yjs document: every editor's socket syncs with it
5 * (the y-protocols sync and awareness messages, as binary frames), and
6 * every change made on the server (an agent's edit, an accepted
7 * suggestion, a restore, a comment) is applied here, so all of them merge
8 * as one CRDT. It holds sockets with the WebSocket Hibernation API, so a
9 * page people have open but aren't typing in costs nothing between
10 * keystrokes; the document is read back from storage when it wakes.
11 *
12 * Storage (the object's own SQLite): the document as a snapshot plus the
13 * updates since, compacted every so often. A few seconds after a burst of
14 * edits (an alarm), the room saves the Markdown rendition, the search
15 * index, backlinks and history to D1 (src/persist.ts). Half a minute after
16 * the Markdown first changes, it brings the page's passages in the
17 * semantic index up to date (src/indexer.ts), on the same alarm, so typing
18 * never waits on embedding.
19 *
20 * The room authorizes nothing about who may open the page: the Worker
21 * checks the viewer's role before forwarding a socket (src/index.ts,
22 * `live`) and puts it in ROOM_MEMBER_HEADER. The room enforces that role:
23 * a socket that may only view or comment never changes the document.
24 */
25import { DurableObject } from "cloudflare:workers";
26
27import type { DocEditTarget, DocRole, DocThreadAction, DocVersionKind, DocsLiveEvent, MemberProfile, ServiceBinding } from "@g1t/contracts";
28import * as decoding from "lib0/decoding";
29import * as encoding from "lib0/encoding";
30import * as syncProtocol from "y-protocols/sync";
31import * as Y from "yjs";
32
33import { atLeast } from "./access.ts";
34import { seed } from "./blocks.ts";
35import { anchorThread, applyEdit, findTarget, rangeIds, rangeMarkdown, restoreFrom, unanchorThread } from "./edits.ts";
36import { bodyCitations } from "./citations.ts";
37import { citationNodes, documentMarkdown, mentionedIds, outline, type Outline } from "./markdown.ts";
38import { indexPage, type DocsJob } from "./indexer.ts";
39import { save } from "./persist.ts";
40import { applyThreadAction, listThreads, setQuote, type ThreadResult } from "./threads.ts";
41
42/** Header the Worker sets on a socket it forwards: who it is, as JSON (`RoomMember`). */
43export const ROOM_MEMBER_HEADER = "x-g1t-docs-member";
44
45export type RoomMember = {
46 page_id: string;
47 workspace_slug: string;
48 /** `user:<id>`. */
49 key: string;
50 member: MemberProfile;
51 role: DocRole;
52};
53
54/** Who made a change on the server, for history. */
55export type Origin = { key: string; kind: DocVersionKind; note: string | null; authors?: string[] };
56
57type Attachment = RoomMember & { clients: number[] };
58
59type Env = { DB: D1Database; NOTIFY?: ServiceBinding; EVENTS?: ServiceBinding; AI?: Ai; VECTORS?: Vectorize; JOBS?: Queue<DocsJob> };
60
61const FRAGMENT = "document-store";
62const MESSAGE_SYNC = 0;
63const MESSAGE_AWARENESS = 1;
64const MESSAGE_QUERY_AWARENESS = 3;
65/** Save this long after the last change. */
66const SAVE_AFTER_MS = 4_000;
67/** Index the page's passages this long after its Markdown first changed: at most twice a minute while someone types. */
68const INDEX_AFTER_MS = 30_000;
69/** Compact the stored updates into one snapshot past this many. */
70const COMPACT_AT = 300;
71
72export class PageRoom extends DurableObject<Env> {
73 private doc: Y.Doc | null = null;
74 /** The last awareness update each client sent, so newcomers see everyone at once. Lost on hibernation; clients resend every 15 s. */
75 private awareness = new Map<number, Uint8Array>();
76 /** Each client's last awareness clock, to mark it gone with the next one. */
77 private clocks = new Map<number, number>();
78
79 constructor(ctx: DurableObjectState, env: Env) {
80 super(ctx, env);
81 ctx.setWebSocketAutoResponse(new WebSocketRequestResponsePair("ping", "pong"));
82 ctx.blockConcurrencyWhile(async () => {
83 this.ctx.storage.sql.exec("CREATE TABLE IF NOT EXISTS updates (seq INTEGER PRIMARY KEY AUTOINCREMENT, data BLOB NOT NULL)");
84 this.ctx.storage.sql.exec("CREATE TABLE IF NOT EXISTS snapshot (id INTEGER PRIMARY KEY CHECK (id = 1), data BLOB NOT NULL)");
85 this.ctx.storage.sql.exec("CREATE TABLE IF NOT EXISTS meta (key TEXT PRIMARY KEY, value TEXT NOT NULL)");
86 });
87 }
88
89 // ── State ──────────────────────────────────────────────────────────────
90
91 private meta<T>(key: string, fallback: T): T {
92 const row = this.ctx.storage.sql.exec<{ value: string }>("SELECT value FROM meta WHERE key = ?", key).toArray()[0];
93 if (!row) return fallback;
94 try {
95 return JSON.parse(row.value) as T;
96 } catch {
97 return fallback;
98 }
99 }
100
101 private setMeta(key: string, value: unknown): void {
102 this.ctx.storage.sql.exec("INSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT (key) DO UPDATE SET value = excluded.value", key, JSON.stringify(value));
103 }
104
105 /** The document, read from storage the first time it is needed after waking. */
106 private load(): Y.Doc {
107 if (this.doc) return this.doc;
108 const doc = new Y.Doc({ gc: true });
109 const snapshot = this.ctx.storage.sql.exec<{ data: ArrayBuffer }>("SELECT data FROM snapshot WHERE id = 1").toArray()[0];
110 if (snapshot) Y.applyUpdate(doc, new Uint8Array(snapshot.data));
111 for (const row of this.ctx.storage.sql.exec<{ data: ArrayBuffer }>("SELECT data FROM updates ORDER BY seq")) {
112 Y.applyUpdate(doc, new Uint8Array(row.data));
113 }
114 doc.on("update", (update: Uint8Array, origin: unknown) => this.onUpdate(update, origin));
115 this.doc = doc;
116 return doc;
117 }
118
119 private fragment(): Y.XmlFragment {
120 return this.load().getXmlFragment(FRAGMENT);
121 }
122
123 private onUpdate(update: Uint8Array, origin: unknown): void {
124 this.ctx.storage.sql.exec("INSERT INTO updates (data) VALUES (?)", update);
125 const count = this.ctx.storage.sql.exec<{ n: number }>("SELECT COUNT(*) AS n FROM updates").one().n;
126 if (count >= COMPACT_AT) this.compact();
127 // Who changed it: a socket's member, or a change made here.
128 const member = origin instanceof WebSocket ? (origin.deserializeAttachment() as Attachment | null) : null;
129 const key = member ? member.key : (origin as Origin | null)?.key;
130 if (member?.member.kind === "user") {
131 const names = this.meta<string[]>("editor_names", []);
132 const name = member.member.name.toLowerCase();
133 if (!names.includes(name)) this.setMeta("editor_names", [...names, name]);
134 }
135 if (key) {
136 const editors = this.meta<string[]>("editors", []).filter((k) => k !== key);
137 editors.push(key);
138 this.setMeta("editors", editors);
139 const pending = this.meta<string[]>("pending_authors", []);
140 if (!pending.includes(key)) this.setMeta("pending_authors", [...pending, key]);
141 }
142 // Everyone else sees it at once.
143 const encoder = encoding.createEncoder();
144 encoding.writeVarUint(encoder, MESSAGE_SYNC);
145 syncProtocol.writeUpdate(encoder, update);
146 this.send(encoding.toUint8Array(encoder), origin instanceof WebSocket ? origin : null);
147 this.scheduleSave();
148 }
149
150 private compact(): void {
151 const doc = this.load();
152 this.ctx.storage.transactionSync(() => {
153 this.ctx.storage.sql.exec("INSERT INTO snapshot (id, data) VALUES (1, ?) ON CONFLICT (id) DO UPDATE SET data = excluded.data", Y.encodeStateAsUpdate(doc));
154 this.ctx.storage.sql.exec("DELETE FROM updates");
155 });
156 }
157
158 private scheduleSave(): void {
159 this.alarmBy(Date.now() + SAVE_AFTER_MS);
160 }
161
162 /** Makes sure the alarm rings by `when`: the save's, or the index's, whichever is first. */
163 private alarmBy(when: number): void {
164 void this.ctx.storage.getAlarm().then((at) => {
165 if (at == null || at > when) return this.ctx.storage.setAlarm(when);
166 });
167 }
168
169 private send(message: Uint8Array | string, except: WebSocket | null = null): void {
170 for (const socket of this.ctx.getWebSockets()) {
171 if (socket === except) continue;
172 try {
173 socket.send(message);
174 } catch {
175 // Closing already.
176 }
177 }
178 }
179
180 /** Saves to D1 now: the Markdown, search, links, mentions, and a version if one is due or asked for. */
181 private async persist(version: Origin | null = null): Promise<string | null> {
182 const pageId = this.meta<string | null>("page_id", null);
183 if (!pageId) return null;
184 const fragment = this.fragment();
185 const editors = this.meta<string[]>("editors", []);
186 const pending = this.meta<string[]>("pending_authors", []);
187 const markdown = documentMarkdown(fragment);
188 const result = await save(this.env, {
189 page_id: pageId,
190 markdown,
191 citations: bodyCitations(citationNodes(fragment), markdown),
192 editors,
193 mentioned: mentionedIds(fragment).users,
194 editor_names: this.meta<string[]>("editor_names", []),
195 state: Y.encodeStateAsUpdate(this.load()),
196 version: version ? { kind: version.kind, note: version.note, authors: version.authors ?? [version.key] } : null,
197 pending_authors: pending,
198 last_version_at: this.meta<number>("last_version_at", 0),
199 workspace_slug: this.meta<string | null>("workspace_slug", null),
200 });
201 this.setMeta("editors", []);
202 this.setMeta("editor_names", []);
203 if (result.changed && this.meta<number | null>("index_at", null) == null) {
204 const when = Date.now() + INDEX_AFTER_MS;
205 this.setMeta("index_at", when);
206 this.alarmBy(when);
207 }
208 if (result.version_id) {
209 this.setMeta("last_version_at", Date.now());
210 this.setMeta("pending_authors", []);
211 }
212 return result.version_id;
213 }
214
215 override async alarm(): Promise<void> {
216 await this.persist();
217 const due = this.meta<number | null>("index_at", null);
218 if (due == null) return;
219 if (Date.now() < due) return this.alarmBy(due);
220 this.setMeta("index_at", null);
221 const pageId = this.meta<string | null>("page_id", null);
222 // Never throws; a failure is retried on the next change.
223 if (pageId) await indexPage(this.env, pageId);
224 }
225
226 // ── Calls from the Worker ──────────────────────────────────────────────
227
228 /**
229 * Names the page this room is for, and fills an empty document from
230 * Markdown (a template, an agent's new page, or a page made before its
231 * room existed) or from a Yjs state (a duplicate).
232 */
233 async ensure(init: { page_id: string; workspace_slug: string; markdown?: string | null; state?: Uint8Array | null }): Promise<void> {
234 this.setMeta("page_id", init.page_id);
235 if (init.workspace_slug) this.setMeta("workspace_slug", init.workspace_slug);
236 const doc = this.load();
237 const fragment = this.fragment();
238 if (fragment.length > 0) return;
239 if (init.state) Y.applyUpdate(doc, init.state, { key: "system", kind: "created", note: null } satisfies Origin);
240 else seed(doc, fragment, init.markdown ?? "");
241 }
242
243 /** The document as Markdown, and its top-level blocks. */
244 async read(): Promise<{ markdown: string; blocks: Outline[] }> {
245 const fragment = this.fragment();
246 return { markdown: documentMarkdown(fragment), blocks: outline(fragment) };
247 }
248
249 /** The whole document's state, for a duplicate. */
250 async state(): Promise<Uint8Array> {
251 return Y.encodeStateAsUpdate(this.load());
252 }
253
254 /** A target's current Markdown and blocks, or null when it is gone. */
255 async target(target: DocEditTarget): Promise<{ markdown: string; block_ids: string[] } | null> {
256 const fragment = this.fragment();
257 const range = findTarget(fragment, target);
258 if (!range) return null;
259 return { markdown: rangeMarkdown(fragment, range), block_ids: rangeIds(fragment, range) };
260 }
261
262 /** Where each target is now, for marking open suggestions in the editor. */
263 async targets(targets: DocEditTarget[]): Promise<(string[] | null)[]> {
264 const fragment = this.fragment();
265 return targets.map((t) => {
266 const range = findTarget(fragment, t);
267 return range ? rangeIds(fragment, range) : null;
268 });
269 }
270
271 /** Applies an edit and records a version for it. False when the target is gone. */
272 async edit(target: DocEditTarget, markdown: string, origin: Origin): Promise<{ applied: boolean; version_id: string | null }> {
273 const doc = this.load();
274 const applied = applyEdit(doc, this.fragment(), target, markdown, origin);
275 if (!applied) return { applied: false, version_id: null };
276 const version_id = await this.persist(origin);
277 return { applied, version_id };
278 }
279
280 /** Makes the document what a version's was, as a new version. */
281 async restore(input: { state: Uint8Array | null; markdown: string }, origin: Origin): Promise<string | null> {
282 const doc = this.load();
283 if (input.state) {
284 const old = new Y.Doc();
285 Y.applyUpdate(old, input.state);
286 restoreFrom(doc, this.fragment(), old.getXmlFragment(FRAGMENT), origin);
287 } else applyEdit(doc, this.fragment(), { kind: "document" }, input.markdown, origin);
288 return this.persist(origin);
289 }
290
291 /** A comment operation from `actor` with `role`. */
292 async thread(actor: string, role: DocRole, action: DocThreadAction): Promise<ThreadResult> {
293 const doc = this.load();
294 if (action.op === "anchor") {
295 if (!atLeast(role, "comment")) return { ok: false, code: "forbidden", message: "You can read this page but not comment on it." };
296 const thread = doc.getMap<Y.Map<unknown>>("threads").get(action.thread_id);
297 if (!thread) return { ok: false, code: "not_found", message: "That thread is gone." };
298 const quote = anchorThread(doc, this.fragment(), action.anchor, action.head, action.thread_id);
299 if (quote) setQuote(doc, action.thread_id, quote);
300 return { ok: true, value: { quote } };
301 }
302 const result = applyThreadAction(doc, actor, role, action);
303 if (result.ok && action.op === "delete_thread") unanchorThread(doc, this.fragment(), action.thread_id);
304 return result;
305 }
306
307 /** The page's threads, authors as member keys. */
308 async threads(): Promise<ReturnType<typeof listThreads>> {
309 return listThreads(this.load());
310 }
311
312 /**
313 * Shows an agent on the page for a little while (its face in everyone's
314 * presence row) when it edits or suggests: an awareness entry of its
315 * own, which editors drop after 30 s without renewal, as for anyone.
316 */
317 async announce(key: string, name: string): Promise<void> {
318 let id = 0;
319 for (let i = 0; i < key.length; i++) id = (Math.imul(id, 31) + key.charCodeAt(i)) | 0;
320 const client = (id & 0x3fffffff) + 1;
321 const clock = Math.floor(Date.now() / 1000);
322 const state = JSON.stringify({ user: { name, color: "#b8a6ff", key, kind: "agent", avatar: "" } });
323 const update = encoding.createEncoder();
324 encoding.writeVarUint(update, 1);
325 encoding.writeVarUint(update, client);
326 encoding.writeVarUint(update, clock);
327 encoding.writeVarString(update, state);
328 const bytes = encoding.toUint8Array(update);
329 this.awareness.set(client, bytes);
330 this.clocks.set(client, clock);
331 const message = encoding.createEncoder();
332 encoding.writeVarUint(message, MESSAGE_AWARENESS);
333 encoding.writeVarUint8Array(message, bytes);
334 this.send(encoding.toUint8Array(message));
335 }
336
337 /** Tells everyone with the page open. */
338 async notice(event: DocsLiveEvent): Promise<void> {
339 this.send(JSON.stringify(event));
340 }
341
342 /** A member's role changed (null: they can no longer read it): their sockets follow. */
343 async setRole(key: string, role: DocRole | null): Promise<void> {
344 for (const socket of this.ctx.getWebSockets(key)) {
345 const who = socket.deserializeAttachment() as Attachment | null;
346 if (!who) continue;
347 if (!role) {
348 try {
349 socket.close(4403, "No longer allowed");
350 } catch {
351 // Already closed.
352 }
353 continue;
354 }
355 socket.serializeAttachment({ ...who, role });
356 try {
357 socket.send(JSON.stringify({ type: "access", role } satisfies DocsLiveEvent));
358 } catch {
359 // Closing.
360 }
361 }
362 }
363
364 /** Saves now, as before the page is archived or exported. */
365 async flush(): Promise<void> {
366 if (this.doc) await this.persist();
367 }
368
369 /** Closes every socket: the page went to the trash. */
370 async closeAll(reason: string): Promise<void> {
371 for (const socket of this.ctx.getWebSockets()) {
372 try {
373 socket.send(JSON.stringify({ type: "page.archived", page_id: this.meta<string>("page_id", "") } satisfies DocsLiveEvent));
374 socket.close(4410, reason);
375 } catch {
376 // Already closed.
377 }
378 }
379 }
380
381 // ── Sockets ────────────────────────────────────────────────────────────
382
383 override async fetch(request: Request): Promise<Response> {
384 if (request.headers.get("upgrade")?.toLowerCase() !== "websocket") {
385 return new Response("Expected a WebSocket upgrade\n", { status: 426 });
386 }
387 let who: RoomMember;
388 try {
389 who = JSON.parse(request.headers.get(ROOM_MEMBER_HEADER) ?? "") as RoomMember;
390 } catch {
391 return new Response("Missing member\n", { status: 400 });
392 }
393 if (who.workspace_slug) this.setMeta("workspace_slug", who.workspace_slug);
394 if (who.page_id) this.setMeta("page_id", who.page_id);
395 const pair = new WebSocketPair();
396 const [client, server] = Object.values(pair) as [WebSocket, WebSocket];
397 this.ctx.acceptWebSocket(server, [who.key]);
398 server.serializeAttachment({ ...who, clients: [] } satisfies Attachment);
399 // Start syncing: our state vector, and everyone's presence so far.
400 const doc = this.load();
401 const encoder = encoding.createEncoder();
402 encoding.writeVarUint(encoder, MESSAGE_SYNC);
403 syncProtocol.writeSyncStep1(encoder, doc);
404 server.send(encoding.toUint8Array(encoder));
405 for (const update of this.awareness.values()) {
406 const e = encoding.createEncoder();
407 encoding.writeVarUint(e, MESSAGE_AWARENESS);
408 encoding.writeVarUint8Array(e, update);
409 server.send(encoding.toUint8Array(e));
410 }
411 return new Response(null, { status: 101, webSocket: client });
412 }
413
414 override async webSocketMessage(socket: WebSocket, message: string | ArrayBuffer): Promise<void> {
415 // Text frames are for the Worker's notices only; clients speak binary.
416 if (typeof message === "string") return;
417 const who = socket.deserializeAttachment() as Attachment | null;
418 if (!who) return;
419 const data = new Uint8Array(message);
420 const decoder = decoding.createDecoder(data);
421 const type = decoding.readVarUint(decoder);
422 if (type === MESSAGE_SYNC) {
423 // Peek at the sync message: viewers and commenters may ask for the
424 // document (step 1) but never change it (step 2, update).
425 const peek = decoding.createDecoder(data);
426 decoding.readVarUint(peek);
427 const step = decoding.readVarUint(peek);
428 if (step !== syncProtocol.messageYjsSyncStep1 && !atLeast(who.role, "edit")) return;
429 const doc = this.load();
430 const encoder = encoding.createEncoder();
431 encoding.writeVarUint(encoder, MESSAGE_SYNC);
432 syncProtocol.readSyncMessage(decoder, encoder, doc, socket);
433 if (encoding.length(encoder) > 1) socket.send(encoding.toUint8Array(encoder));
434 return;
435 }
436 if (type === MESSAGE_AWARENESS) {
437 const update = decoding.readVarUint8Array(decoder);
438 // Remember which clients this socket speaks for, to clear them when it closes.
439 const entries = awarenessEntries(update);
440 const clients = entries.map((e) => e.id);
441 for (const e of entries) {
442 this.clocks.set(e.id, e.clock);
443 if (e.gone) this.awareness.delete(e.id);
444 else this.awareness.set(e.id, update);
445 }
446 const known = new Set(who.clients);
447 if (clients.some((id) => !known.has(id))) socket.serializeAttachment({ ...who, clients: [...new Set([...who.clients, ...clients])] });
448 const encoder = encoding.createEncoder();
449 encoding.writeVarUint(encoder, MESSAGE_AWARENESS);
450 encoding.writeVarUint8Array(encoder, update);
451 this.send(encoding.toUint8Array(encoder), socket);
452 return;
453 }
454 if (type === MESSAGE_QUERY_AWARENESS) {
455 for (const update of this.awareness.values()) {
456 const encoder = encoding.createEncoder();
457 encoding.writeVarUint(encoder, MESSAGE_AWARENESS);
458 encoding.writeVarUint8Array(encoder, update);
459 socket.send(encoding.toUint8Array(encoder));
460 }
461 }
462 }
463
464 override async webSocketClose(socket: WebSocket, code: number, reason: string): Promise<void> {
465 this.leave(socket);
466 try {
467 socket.close(code, reason);
468 } catch {
469 // Already closed.
470 }
471 }
472
473 override async webSocketError(socket: WebSocket): Promise<void> {
474 this.leave(socket);
475 }
476
477 /** Tells everyone a closed socket's people left: their cursors go. */
478 private leave(socket: WebSocket): void {
479 const who = socket.deserializeAttachment() as Attachment | null;
480 if (!who?.clients.length) return;
481 for (const id of who.clients) this.awareness.delete(id);
482 // Marks each client gone (state null, the next clock). One whose clock
483 // was lost to hibernation is left to time out in the others' editors
484 // (30 s), as y-protocols does for any client that stops renewing.
485 const gone = who.clients.filter((id) => this.clocks.has(id));
486 if (!gone.length) return;
487 const encoder = encoding.createEncoder();
488 encoding.writeVarUint(encoder, gone.length);
489 for (const id of gone) {
490 encoding.writeVarUint(encoder, id);
491 encoding.writeVarUint(encoder, this.clocks.get(id)! + 1);
492 encoding.writeVarString(encoder, "null");
493 this.clocks.delete(id);
494 }
495 const update = encoding.toUint8Array(encoder);
496 const message = encoding.createEncoder();
497 encoding.writeVarUint(message, MESSAGE_AWARENESS);
498 encoding.writeVarUint8Array(message, update);
499 this.send(encoding.toUint8Array(message), socket);
500 }
501}
502
503/** The clients in an awareness update (its format: count, then id, clock, JSON state). */
504export function awarenessEntries(update: Uint8Array): { id: number; clock: number; gone: boolean }[] {
505 const decoder = decoding.createDecoder(update);
506 const count = decoding.readVarUint(decoder);
507 const out: { id: number; clock: number; gone: boolean }[] = [];
508 for (let i = 0; i < count; i++) {
509 const id = decoding.readVarUint(decoder);
510 const clock = decoding.readVarUint(decoder);
511 const state = decoding.readVarString(decoder);
512 out.push({ id, clock, gone: state === "null" });
513 }
514 return out;
515}
516