You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

149 lines
3.7 KiB

/**
* Thin WebSocket client for the dating server's `/api/realtime`
* endpoint. Authenticates via the same `dating_session` cookie the
* HTTP API uses (the cookie is sent with the upgrade request because
* the WebSocket URL shares hostname with the page — same-site, so
* `SameSite=Lax` lets the cookie travel).
*
* The connection auto-reconnects with exponential backoff so the user
* doesn't have to refresh after a server restart or temporary
* disconnect. `dispose()` from the returned handle cancels every
* pending reconnect and closes the socket.
*/
import { resolveDatingApiBase, DATING_API_DEFAULT_PORT } from './api.ts';
import type { DatingMessage } from './types.ts';
export type DatingRealtimeEvent =
| { type: 'hello'; userId: string }
| { type: 'dating.message.created'; message: DatingMessage };
export type DatingRealtimeListener = (event: DatingRealtimeEvent) => void;
export interface DatingRealtimeOptions {
readonly url?: string;
/** Cap on the reconnect backoff. Default: 30 s. */
readonly maxReconnectMs?: number;
}
export interface DatingRealtimeHandle {
/** Subscribe to incoming events. Returns an unsubscribe. */
subscribe(listener: DatingRealtimeListener): () => void;
/** True once the server has accepted the connection. */
connected(): boolean;
dispose(): void;
}
function defaultRealtimeUrl(): string {
const base = resolveDatingApiBase(DATING_API_DEFAULT_PORT);
const wsBase = base.replace(/^http/, 'ws');
return `${wsBase}/api/realtime`;
}
export function connectDatingRealtime(
options: DatingRealtimeOptions = {}
): DatingRealtimeHandle {
const url = options.url ?? defaultRealtimeUrl();
const maxReconnectMs = options.maxReconnectMs ?? 30_000;
const listeners = new Set<DatingRealtimeListener>();
let socket: WebSocket | null = null;
let isConnected = false;
let disposed = false;
let attempt = 0;
let reconnectHandle: ReturnType<typeof setTimeout> | null = null;
function nextDelay(): number {
// 1s, 2s, 4s, 8s, … capped at `maxReconnectMs`.
return Math.min(maxReconnectMs, 1000 * Math.pow(2, Math.min(attempt, 6)));
}
function dispatch(event: DatingRealtimeEvent) {
for (const listener of [...listeners]) {
try {
listener(event);
} catch {
// listeners must not block one another
}
}
}
function connect() {
if (disposed) return;
if (typeof WebSocket === 'undefined') return; // SSR guard
try {
socket = new WebSocket(url);
} catch {
scheduleReconnect();
return;
}
socket.addEventListener('open', () => {
attempt = 0;
isConnected = true;
});
socket.addEventListener('message', (event) => {
let payload: unknown;
try {
payload = JSON.parse(typeof event.data === 'string' ? event.data : '');
} catch {
return;
}
if (
typeof payload !== 'object' ||
payload === null ||
typeof (payload as { type?: unknown }).type !== 'string'
) {
return;
}
dispatch(payload as DatingRealtimeEvent);
});
const onClose = () => {
isConnected = false;
socket = null;
scheduleReconnect();
};
socket.addEventListener('close', onClose);
socket.addEventListener('error', onClose);
}
function scheduleReconnect() {
if (disposed) return;
const delay = nextDelay();
attempt += 1;
reconnectHandle = setTimeout(connect, delay);
}
connect();
return {
subscribe(listener) {
listeners.add(listener);
return () => listeners.delete(listener);
},
connected() {
return isConnected;
},
dispose() {
if (disposed) return;
disposed = true;
listeners.clear();
if (reconnectHandle !== null) {
clearTimeout(reconnectHandle);
reconnectHandle = null;
}
if (socket !== null) {
try {
socket.close();
} catch {
// best-effort
}
socket = null;
}
}
};
}

Powered by TurnKey Linux.