| 1 | import type { Client } from "./client"; |
| 2 | import { 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 | |
| 15 | export type SocketStatus = "connecting" | "live" | "offline"; |
| 16 | |
| 17 | type 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 | |
| 25 | const PING_MS = 25_000; |
| 26 | |
| 27 | export 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 | } |