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