Skip to content
249 linesCodeBlameRaw
1/**
2 * An artifact's live connection, whatever its kind: the Yjs document and
3 * everyone's presence, synced over
4 * `wss://<site>/<workspace>/-/artifacts/live?folio=<id>` with the folio's
5 * room (services/artifacts src/folios/room.ts). Binary frames speak the
6 * y-protocols sync and awareness messages; text frames are the service's
7 * own notices (`FoliosLiveEvent`): a rename, a new suggestion, a version,
8 * a change of access.
9 *
10 * It reconnects with backoff when the socket drops, and says where it is
11 * (`status`) so the page can show "Offline, changes will sync".
12 *
13 * Given the document as the page's data carries it (the state the room
14 * last saved), it starts from that: the editor opens on it at once, and
15 * the room's answer to our state vector is only what changed since. Any
16 * edit made before the room answers stays in the document and goes out
17 * in that same exchange (sync step 2), so nothing typed early is lost.
18 * Browser-only.
19 */
20import type { FoliosLiveEvent } from "@g1t/contracts";
21import * as decoding from "lib0/decoding";
22import * as encoding from "lib0/encoding";
23import * as awarenessProtocol from "y-protocols/awareness";
24import * as syncProtocol from "y-protocols/sync";
25import * as Y from "yjs";
26
27// With extensions, so provider.test.ts runs under Node as the other tests do.
28import { openLive } from "../../lib/live-socket.ts";
29import { heldOpen } from "../../lib/notify-store.ts";
30
31const MESSAGE_SYNC = 0;
32const MESSAGE_AWARENESS = 1;
33const MESSAGE_QUERY_AWARENESS = 3;
34
35export type LiveStatus = "connecting" | "synced" | "offline" | "closed";
36
37/** The document a page carries (base64 of a Yjs update), as bytes; null when it carries none or it is not base64. */
38export function pageStateBytes(state: string | null | undefined): Uint8Array | null {
39 if (!state) return null;
40 try {
41 return Uint8Array.from(atob(state), (c) => c.charCodeAt(0));
42 } catch {
43 return null;
44 }
45}
46
47export class FolioProvider {
48 readonly doc: Y.Doc;
49 readonly awareness: awarenessProtocol.Awareness;
50 status: LiveStatus = "connecting";
51 /** Whether the document started from the page's saved state (so the editor need not wait for the room). */
52 readonly seeded: boolean;
53 private socket: WebSocket | null = null;
54 private attempts = 0;
55 /** When the current connection opened; the backoff starts over only once one holds. */
56 private openedAt: number | null = null;
57 private timer: ReturnType<typeof setTimeout> | null = null;
58 private keepalive: ReturnType<typeof setInterval> | null = null;
59 private stopped = false;
60 /** Between asking for a socket ticket and opening the socket. */
61 private opening = false;
62 private readonly statusListeners = new Set<(status: LiveStatus) => void>();
63 private readonly eventListeners = new Set<(event: FoliosLiveEvent) => void>();
64
65 private readonly url: string;
66
67 constructor(url: string, options: { doc?: Y.Doc; state?: Uint8Array | null } = {}) {
68 this.url = url;
69 this.doc = options.doc ?? new Y.Doc();
70 let seeded = false;
71 if (options.state?.byteLength) {
72 try {
73 // No origin of our own: the seed is the room's work, not an edit to send it.
74 Y.applyUpdate(this.doc, options.state, "seed");
75 seeded = true;
76 } catch (error) {
77 console.error("artifacts: the saved document could not be read; waiting for the room", error);
78 }
79 }
80 this.seeded = seeded;
81 this.awareness = new awarenessProtocol.Awareness(this.doc);
82 this.doc.on("update", this.onDocUpdate);
83 this.awareness.on("update", this.onAwarenessUpdate);
84 if (typeof window !== "undefined") {
85 window.addEventListener("beforeunload", this.onUnload);
86 window.addEventListener("online", this.onOnline);
87 }
88 this.connect();
89 }
90
91 onStatus(listener: (status: LiveStatus) => void): () => void {
92 this.statusListeners.add(listener);
93 listener(this.status);
94 return () => this.statusListeners.delete(listener);
95 }
96
97 onEvent(listener: (event: FoliosLiveEvent) => void): () => void {
98 this.eventListeners.add(listener);
99 return () => this.eventListeners.delete(listener);
100 }
101
102 private setStatus(status: LiveStatus) {
103 if (this.status === status) return;
104 this.status = status;
105 for (const l of this.statusListeners) l(status);
106 }
107
108 private connect() {
109 if (this.stopped || this.opening) return;
110 this.setStatus(this.attempts ? "offline" : "connecting");
111 // A page opened with an access token adds a socket ticket first (lib/live-socket.ts).
112 const target = new URL(this.url);
113 this.opening = true;
114 openLive(
115 target.pathname,
116 () => Object.fromEntries(target.searchParams),
117 (address) => {
118 this.opening = false;
119 this.open(address);
120 },
121 () => {
122 if (!this.stopped) return false;
123 this.opening = false;
124 return true;
125 },
126 );
127 }
128
129 private open(address: string) {
130 const socket = new WebSocket(address);
131 socket.binaryType = "arraybuffer";
132 this.socket = socket;
133 socket.onopen = () => {
134 this.openedAt = Date.now();
135 // Our state vector: the room answers with what we lack, and asks for what it lacks.
136 const encoder = encoding.createEncoder();
137 encoding.writeVarUint(encoder, MESSAGE_SYNC);
138 syncProtocol.writeSyncStep1(encoder, this.doc);
139 socket.send(encoding.toUint8Array(encoder));
140 if (this.awareness.getLocalState() !== null) {
141 const a = encoding.createEncoder();
142 encoding.writeVarUint(a, MESSAGE_AWARENESS);
143 encoding.writeVarUint8Array(a, awarenessProtocol.encodeAwarenessUpdate(this.awareness, [this.doc.clientID]));
144 socket.send(encoding.toUint8Array(a));
145 }
146 const q = encoding.createEncoder();
147 encoding.writeVarUint(q, MESSAGE_QUERY_AWARENESS);
148 socket.send(encoding.toUint8Array(q));
149 if (this.keepalive) clearInterval(this.keepalive);
150 // Answered at the edge without waking the room.
151 this.keepalive = setInterval(() => socket.readyState === WebSocket.OPEN && socket.send("ping"), 25_000);
152 };
153 socket.onmessage = (event) => {
154 if (typeof event.data === "string") {
155 if (event.data === "pong") return;
156 try {
157 const parsed = JSON.parse(event.data) as FoliosLiveEvent;
158 for (const l of this.eventListeners) l(parsed);
159 } catch {
160 // Not ours.
161 }
162 return;
163 }
164 this.receive(new Uint8Array(event.data as ArrayBuffer));
165 };
166 socket.onclose = (event) => {
167 if (this.keepalive) clearInterval(this.keepalive);
168 this.socket = null;
169 // Others stop seeing our cursor; we stop seeing theirs.
170 awarenessProtocol.removeAwarenessStates(
171 this.awareness,
172 [...this.awareness.getStates().keys()].filter((id) => id !== this.doc.clientID),
173 this,
174 );
175 // 4403: no longer allowed; 4410: in the trash. Neither comes back by retrying.
176 if (event.code === 4403 || event.code === 4410 || this.stopped) {
177 this.setStatus("closed");
178 return;
179 }
180 this.setStatus("offline");
181 // Only a connection that held starts the backoff over.
182 if (heldOpen(this.openedAt)) this.attempts = 0;
183 this.openedAt = null;
184 const delay = Math.min(30_000, 500 * 2 ** this.attempts) + Math.random() * 500;
185 this.attempts++;
186 this.timer = setTimeout(() => this.connect(), delay);
187 };
188 }
189
190 private receive(data: Uint8Array) {
191 const decoder = decoding.createDecoder(data);
192 const type = decoding.readVarUint(decoder);
193 if (type === MESSAGE_SYNC) {
194 const encoder = encoding.createEncoder();
195 encoding.writeVarUint(encoder, MESSAGE_SYNC);
196 const step = syncProtocol.readSyncMessage(decoder, encoder, this.doc, this);
197 if (encoding.length(encoder) > 1) this.socket?.send(encoding.toUint8Array(encoder));
198 if (step === syncProtocol.messageYjsSyncStep2) this.setStatus("synced");
199 return;
200 }
201 if (type === MESSAGE_AWARENESS) {
202 awarenessProtocol.applyAwarenessUpdate(this.awareness, decoding.readVarUint8Array(decoder), this);
203 }
204 }
205
206 private onDocUpdate = (update: Uint8Array, origin: unknown) => {
207 if (origin === this) return;
208 const encoder = encoding.createEncoder();
209 encoding.writeVarUint(encoder, MESSAGE_SYNC);
210 syncProtocol.writeUpdate(encoder, update);
211 if (this.socket?.readyState === WebSocket.OPEN) this.socket.send(encoding.toUint8Array(encoder));
212 // While offline, the update stays in the document and goes in the next sync.
213 };
214
215 private onAwarenessUpdate = ({ added, updated, removed }: { added: number[]; updated: number[]; removed: number[] }, origin: unknown) => {
216 if (origin === this) return;
217 const changed = [...added, ...updated, ...removed];
218 const encoder = encoding.createEncoder();
219 encoding.writeVarUint(encoder, MESSAGE_AWARENESS);
220 encoding.writeVarUint8Array(encoder, awarenessProtocol.encodeAwarenessUpdate(this.awareness, changed));
221 if (this.socket?.readyState === WebSocket.OPEN) this.socket.send(encoding.toUint8Array(encoder));
222 };
223
224 private onUnload = () => {
225 awarenessProtocol.removeAwarenessStates(this.awareness, [this.doc.clientID], "unload");
226 };
227
228 private onOnline = () => {
229 if (this.socket || this.opening || this.stopped) return;
230 if (this.timer) clearTimeout(this.timer);
231 this.attempts = 0;
232 this.connect();
233 };
234
235 destroy() {
236 this.stopped = true;
237 if (this.timer) clearTimeout(this.timer);
238 if (this.keepalive) clearInterval(this.keepalive);
239 awarenessProtocol.removeAwarenessStates(this.awareness, [this.doc.clientID], "destroy");
240 this.doc.off("update", this.onDocUpdate);
241 this.awareness.off("update", this.onAwarenessUpdate);
242 if (typeof window !== "undefined") {
243 window.removeEventListener("beforeunload", this.onUnload);
244 window.removeEventListener("online", this.onOnline);
245 }
246 this.socket?.close();
247 this.awareness.destroy();
248 }
249}