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