| 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 | } |