Skip to content
100 linesCodeBlameRaw
1import type { Client } from "./client";
2import { socketUrl } from "./server";
3
4/**
5 * One of the site's live sockets, kept open: a conversation's
6 * (`/<workspace>/-/chat/live`) or the person's feed (`/-/live`). A socket
7 * cannot carry the `Authorization` header, so each opening asks
8 * `GET /-/live/ticket?path=<path>` for a ticket (good for a minute, for
9 * that socket alone; apps/web/app/lib/socket-ticket.ts) and adds it to
10 * the address. A dropped socket opens again, with a new ticket, waiting
11 * longer each time up to half a minute; a `ping` every 25 s keeps it
12 * from being taken for gone.
13 */
14
15export type SocketStatus = "connecting" | "live" | "offline";
16
17type Handlers = {
18 onMessage: (data: unknown) => void;
19 onStatus?: (status: SocketStatus) => void;
20 /** Opened again after a drop: what was missed meanwhile should be read again. */
21 onReconnect?: () => void;
22 onOpen?: (send: (frame: unknown) => void) => void;
23};
24
25const PING_MS = 25_000;
26
27export function openSocket(client: Client, path: string, query: Record<string, string>, handlers: Handlers) {
28 let socket: WebSocket | null = null;
29 let closed = false;
30 let attempts = 0;
31 let opened = false;
32 let retry: ReturnType<typeof setTimeout> | null = null;
33 let ping: ReturnType<typeof setInterval> | null = null;
34
35 const send = (frame: unknown) => {
36 if (socket?.readyState === WebSocket.OPEN) socket.send(typeof frame === "string" ? frame : JSON.stringify(frame));
37 };
38
39 const again = () => {
40 if (ping) clearInterval(ping);
41 ping = null;
42 if (closed) return;
43 handlers.onStatus?.("offline");
44 attempts += 1;
45 retry = setTimeout(open, Math.min(30_000, 1_000 * 2 ** Math.min(attempts - 1, 5)));
46 };
47
48 async function open() {
49 if (closed) return;
50 handlers.onStatus?.("connecting");
51 let ticket: string;
52 try {
53 const token = await client.token();
54 const response = await fetch(`${client.server.site}/-/live/ticket?path=${encodeURIComponent(path)}`, {
55 headers: { authorization: `Bearer ${token}`, accept: "application/json" },
56 });
57 if (!response.ok) throw new Error(`ticket ${response.status}`);
58 ticket = ((await response.json()) as { ticket: string }).ticket;
59 } catch {
60 return again();
61 }
62 if (closed) return;
63 const address = `${socketUrl(client.server, path)}?${new URLSearchParams({ ...query, ticket })}`;
64 const next = new WebSocket(address);
65 socket = next;
66 next.onopen = () => {
67 const wasOpened = opened;
68 opened = true;
69 attempts = 0;
70 handlers.onStatus?.("live");
71 ping = setInterval(() => send("ping"), PING_MS);
72 handlers.onOpen?.(send);
73 if (wasOpened) handlers.onReconnect?.();
74 };
75 next.onmessage = (event) => {
76 if (event.data === "pong") return;
77 try {
78 handlers.onMessage(JSON.parse(String(event.data)));
79 } catch {
80 // Not JSON: nothing for the app.
81 }
82 };
83 next.onclose = () => {
84 if (socket === next) again();
85 };
86 }
87
88 void open();
89
90 return {
91 send,
92 close() {
93 closed = true;
94 if (retry) clearTimeout(retry);
95 if (ping) clearInterval(ping);
96 socket?.close();
97 socket = null;
98 },
99 };
100}