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 | */ | |
| 24 | import { DurableObject } from "cloudflare:workers"; | |
| 25 | ||
| 26 | import type { DocEditTarget, DocRole, DocThreadAction, FolioAgentEdit, FolioKind, FoliosLiveEvent, MemberProfile, ServiceBinding } from "@g1t/contracts"; | |
| 27 | import * as decoding from "lib0/decoding"; | |
| 28 | import * as encoding from "lib0/encoding"; | |
| 29 | import * as syncProtocol from "y-protocols/sync"; | |
| 30 | import * as Y from "yjs"; | |
| 31 | ||
| 32 | import { atLeast } from "../access.ts"; | |
| 33 | import { anchorThread, unanchorThread } from "../edits.ts"; | |
| 34 | import type { FileStoreEnv } from "../files.ts"; | |
| 35 | import { indexFolio, type IndexEnv } from "../indexer.ts"; | |
| 36 | import { docFragment, docTarget, docTargets } from "../kinds/doc/index.ts"; | |
| 37 | import { kindModel } from "../kinds/index.ts"; | |
| 38 | import type { AgentForm, FolioOrigin, KindModel } from "../kinds/types.ts"; | |
| 39 | import { saveFolio } from "../persist.ts"; | |
| 40 | import { ROOM_MEMBER_HEADER, awarenessEntries } from "../room.ts"; | |
| 41 | import { applyThreadAction, listThreads, setQuote, type ThreadResult } from "../threads.ts"; | |
| 42 | import { notifyFolioMentions } from "./service.ts"; | |
| 43 | ||
| 44 | export { ROOM_MEMBER_HEADER }; | |
| 45 | ||
| 46 | export type FolioRoomMember = { | |
| 47 | folio_id: string; | |
| 48 | workspace_slug: string; | |
| 49 | /** `user:<id>`. */ | |
| 50 | key: string; | |
| 51 | member: MemberProfile; | |
| 52 | role: DocRole; | |
| 53 | }; | |
| 54 | ||
| 55 | type Attachment = FolioRoomMember & { clients: number[] }; | |
| 56 | ||
| 57 | export 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 | ||
| 69 | const MESSAGE_SYNC = 0; | |
| 70 | const MESSAGE_AWARENESS = 1; | |
| 71 | const MESSAGE_QUERY_AWARENESS = 3; | |
| 72 | /** Save this long after the last change. */ | |
| 73 | const SAVE_AFTER_MS = 4_000; | |
| 74 | /** Index this long after the text first changed: at most twice a minute while someone types. */ | |
| 75 | const INDEX_AFTER_MS = 30_000; | |
| 76 | /** Compact the stored updates into one snapshot past this many. */ | |
| 77 | const COMPACT_AT = 300; | |
| 78 | ||
| 79 | export 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 | } |