Skip to content
512 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

The docs service answers every artifacts call: docs can be made, listed, shared, moved, trashed, restored, searched, versioned and edited live in their own rooms, agents read, write and recall them only where their person and everyone in the conversation can, and folio events go out on the bus, while Docs' pages keep working as before.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}