Skip to content
189 linesCodeBlameRaw
1/**
2 * A page's live connection: the Yjs document and everyone's presence,
3 * synced over `wss://<site>/<workspace>/-/docs/live?page=<id>` with the
4 * page's room (services/docs src/room.ts). Binary frames speak the
5 * y-protocols sync and awareness messages; text frames are the service's
6 * own notices (`DocsLiveEvent`): a rename, a new suggestion, a version,
7 * a change of access.
8 *
9 * It reconnects with backoff when the socket drops, and says where it is
10 * (`status`) so the page can show "Offline, changes will sync".
11 * Browser-only.
12 */
13import type { DocsLiveEvent } from "@g1t/contracts";
14import * as decoding from "lib0/decoding";
15import * as encoding from "lib0/encoding";
16import * as awarenessProtocol from "y-protocols/awareness";
17import * as syncProtocol from "y-protocols/sync";
18import * as Y from "yjs";
19
20const MESSAGE_SYNC = 0;
21const MESSAGE_AWARENESS = 1;
22const MESSAGE_QUERY_AWARENESS = 3;
23
24export type LiveStatus = "connecting" | "synced" | "offline" | "closed";
25
26export class DocsProvider {
27 readonly doc: Y.Doc;
28 readonly awareness: awarenessProtocol.Awareness;
29 status: LiveStatus = "connecting";
30 private socket: WebSocket | null = null;
31 private attempts = 0;
32 private timer: ReturnType<typeof setTimeout> | null = null;
33 private keepalive: ReturnType<typeof setInterval> | null = null;
34 private stopped = false;
35 private readonly statusListeners = new Set<(status: LiveStatus) => void>();
36 private readonly eventListeners = new Set<(event: DocsLiveEvent) => void>();
37
38 constructor(
39 private readonly url: string,
40 doc?: Y.Doc,
41 ) {
42 this.doc = doc ?? new Y.Doc();
43 this.awareness = new awarenessProtocol.Awareness(this.doc);
44 this.doc.on("update", this.onDocUpdate);
45 this.awareness.on("update", this.onAwarenessUpdate);
46 if (typeof window !== "undefined") {
47 window.addEventListener("beforeunload", this.onUnload);
48 window.addEventListener("online", this.onOnline);
49 }
50 this.connect();
51 }
52
53 onStatus(listener: (status: LiveStatus) => void): () => void {
54 this.statusListeners.add(listener);
55 listener(this.status);
56 return () => this.statusListeners.delete(listener);
57 }
58
59 onEvent(listener: (event: DocsLiveEvent) => void): () => void {
60 this.eventListeners.add(listener);
61 return () => this.eventListeners.delete(listener);
62 }
63
64 private setStatus(status: LiveStatus) {
65 if (this.status === status) return;
66 this.status = status;
67 for (const l of this.statusListeners) l(status);
68 }
69
70 private connect() {
71 if (this.stopped) return;
72 this.setStatus(this.attempts ? "offline" : "connecting");
73 const socket = new WebSocket(this.url);
74 socket.binaryType = "arraybuffer";
75 this.socket = socket;
76 socket.onopen = () => {
77 this.attempts = 0;
78 // Our state vector: the room answers with what we lack, and asks for what it lacks.
79 const encoder = encoding.createEncoder();
80 encoding.writeVarUint(encoder, MESSAGE_SYNC);
81 syncProtocol.writeSyncStep1(encoder, this.doc);
82 socket.send(encoding.toUint8Array(encoder));
83 if (this.awareness.getLocalState() !== null) {
84 const a = encoding.createEncoder();
85 encoding.writeVarUint(a, MESSAGE_AWARENESS);
86 encoding.writeVarUint8Array(a, awarenessProtocol.encodeAwarenessUpdate(this.awareness, [this.doc.clientID]));
87 socket.send(encoding.toUint8Array(a));
88 }
89 const q = encoding.createEncoder();
90 encoding.writeVarUint(q, MESSAGE_QUERY_AWARENESS);
91 socket.send(encoding.toUint8Array(q));
92 if (this.keepalive) clearInterval(this.keepalive);
93 // Answered at the edge without waking the room.
94 this.keepalive = setInterval(() => socket.readyState === WebSocket.OPEN && socket.send("ping"), 25_000);
95 };
96 socket.onmessage = (event) => {
97 if (typeof event.data === "string") {
98 if (event.data === "pong") return;
99 try {
100 const parsed = JSON.parse(event.data) as DocsLiveEvent;
101 for (const l of this.eventListeners) l(parsed);
102 } catch {
103 // Not ours.
104 }
105 return;
106 }
107 this.receive(new Uint8Array(event.data as ArrayBuffer));
108 };
109 socket.onclose = (event) => {
110 if (this.keepalive) clearInterval(this.keepalive);
111 this.socket = null;
112 // Others stop seeing our cursor; we stop seeing theirs.
113 awarenessProtocol.removeAwarenessStates(
114 this.awareness,
115 [...this.awareness.getStates().keys()].filter((id) => id !== this.doc.clientID),
116 this,
117 );
118 // 4403: no longer allowed; 4410: in the trash. Neither comes back by retrying.
119 if (event.code === 4403 || event.code === 4410 || this.stopped) {
120 this.setStatus("closed");
121 return;
122 }
123 this.setStatus("offline");
124 const delay = Math.min(30_000, 500 * 2 ** this.attempts) + Math.random() * 500;
125 this.attempts++;
126 this.timer = setTimeout(() => this.connect(), delay);
127 };
128 }
129
130 private receive(data: Uint8Array) {
131 const decoder = decoding.createDecoder(data);
132 const type = decoding.readVarUint(decoder);
133 if (type === MESSAGE_SYNC) {
134 const encoder = encoding.createEncoder();
135 encoding.writeVarUint(encoder, MESSAGE_SYNC);
136 const step = syncProtocol.readSyncMessage(decoder, encoder, this.doc, this);
137 if (encoding.length(encoder) > 1) this.socket?.send(encoding.toUint8Array(encoder));
138 if (step === syncProtocol.messageYjsSyncStep2) this.setStatus("synced");
139 return;
140 }
141 if (type === MESSAGE_AWARENESS) {
142 awarenessProtocol.applyAwarenessUpdate(this.awareness, decoding.readVarUint8Array(decoder), this);
143 }
144 }
145
146 private onDocUpdate = (update: Uint8Array, origin: unknown) => {
147 if (origin === this) return;
148 const encoder = encoding.createEncoder();
149 encoding.writeVarUint(encoder, MESSAGE_SYNC);
150 syncProtocol.writeUpdate(encoder, update);
151 if (this.socket?.readyState === WebSocket.OPEN) this.socket.send(encoding.toUint8Array(encoder));
152 // While offline, the update stays in the document and goes in the next sync.
153 };
154
155 private onAwarenessUpdate = ({ added, updated, removed }: { added: number[]; updated: number[]; removed: number[] }, origin: unknown) => {
156 if (origin === this) return;
157 const changed = [...added, ...updated, ...removed];
158 const encoder = encoding.createEncoder();
159 encoding.writeVarUint(encoder, MESSAGE_AWARENESS);
160 encoding.writeVarUint8Array(encoder, awarenessProtocol.encodeAwarenessUpdate(this.awareness, changed));
161 if (this.socket?.readyState === WebSocket.OPEN) this.socket.send(encoding.toUint8Array(encoder));
162 };
163
164 private onUnload = () => {
165 awarenessProtocol.removeAwarenessStates(this.awareness, [this.doc.clientID], "unload");
166 };
167
168 private onOnline = () => {
169 if (this.socket || this.stopped) return;
170 if (this.timer) clearTimeout(this.timer);
171 this.attempts = 0;
172 this.connect();
173 };
174
175 destroy() {
176 this.stopped = true;
177 if (this.timer) clearTimeout(this.timer);
178 if (this.keepalive) clearInterval(this.keepalive);
179 awarenessProtocol.removeAwarenessStates(this.awareness, [this.doc.clientID], "destroy");
180 this.doc.off("update", this.onDocUpdate);
181 this.awareness.off("update", this.onAwarenessUpdate);
182 if (typeof window !== "undefined") {
183 window.removeEventListener("beforeunload", this.onUnload);
184 window.removeEventListener("online", this.onOnline);
185 }
186 this.socket?.close();
187 this.awareness.destroy();
188 }
189}