From fde6a4ccb99f3d675a2d796e7e157c522d435028 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 16 Feb 2026 00:54:27 +0100 Subject: [PATCH] Restructured client --- src/client.ts | 116 +++++++++++++++++++++++++++++++++++++++++--------- 1 file changed, 96 insertions(+), 20 deletions(-) diff --git a/src/client.ts b/src/client.ts index 234f29c..41db992 100644 --- a/src/client.ts +++ b/src/client.ts @@ -2,16 +2,28 @@ import { WebSocket } from "ws"; import { LiveRunMove } from "./types"; import { randomUUID } from "crypto"; +/* ────────────────────────── + State +────────────────────────── */ + let socket: WebSocket | null = null; +let connecting: Promise | null = null; +let currentOrigin: string | null = null; + let cookies: Record | undefined; type PendingResolver = { resolve: (value: any) => void; reject: (err: any) => void; + timeout: NodeJS.Timeout; }; const pending = new Map(); +/* ────────────────────────── + Cookies +────────────────────────── */ + export function getCookies(): Record | undefined { return cookies; } @@ -20,74 +32,134 @@ export function setCookies(c: Record) { cookies = c; } +/* ────────────────────────── + Utils +────────────────────────── */ + function getWsOrigin(useLocalhost: boolean) { return useLocalhost ? "ws://localhost:3000/api/run" : "wss://dev-run.selimaj.dev/api/run"; } +function rejectAllPending(err: Error) { + for (const { reject, timeout } of pending.values()) { + clearTimeout(timeout); + reject(err); + } + pending.clear(); +} + +/* ────────────────────────── + Socket management +────────────────────────── */ + async function ensureSocket(useLocalhost: boolean): Promise { + const origin = getWsOrigin(useLocalhost); + if ( socket && - socket.url === getWsOrigin(useLocalhost) && - socket.readyState === WebSocket.OPEN + socket.readyState === WebSocket.OPEN && + currentOrigin === origin ) { return socket; } - if (socket && socket.url !== getWsOrigin(useLocalhost)) { - socket.close(); + if (connecting) { + return connecting; } - return new Promise((resolve, reject) => { - const ws = new WebSocket(getWsOrigin(useLocalhost)); + if (socket && currentOrigin !== origin) { + socket.close(); + socket = null; + } + + connecting = new Promise((resolve, reject) => { + const ws = new WebSocket(origin); ws.onopen = async () => { - socket = ws; + try { + const authCookies = getCookies(); + if (!authCookies) { + throw new Error("Cookies not set"); + } - console.log("Authenticating with cookies:", getCookies()); + socket = ws; + currentOrigin = origin; + connecting = null; - await sendUnsafe(useLocalhost, { cookies: getCookies() }); + await sendUnsafe(useLocalhost, { + type: "auth", + cookies: authCookies, + }); - resolve(ws); - }; - - ws.onerror = (err) => { - reject(err); + resolve(ws); + } catch (err) { + ws.close(); + reject(err); + } }; ws.onmessage = (e) => { - const msg = JSON.parse(e.data.toString()); - const { requestId, ok, data, error } = msg; + let msg: any; + try { + msg = JSON.parse(e.data.toString()); + } catch { + console.warn("Invalid WS message:", e.data); + return; + } + const { requestId, ok, data, error } = msg; if (!requestId) return; const pendingReq = pending.get(requestId); if (!pendingReq) return; pending.delete(requestId); + clearTimeout(pendingReq.timeout); ok ? pendingReq.resolve(data) : pendingReq.reject(error); }; + ws.onerror = (err) => { + rejectAllPending(new Error("WebSocket error")); + connecting = null; + reject(err); + }; + ws.onclose = () => { + rejectAllPending(new Error("WebSocket closed")); socket = null; + currentOrigin = null; + connecting = null; }; }); + + return connecting; } +/* ────────────────────────── + RPC helpers +────────────────────────── */ + async function sendUnsafe( useLocalhost: boolean, payload: Record, ): Promise { - if (!socket) { - throw new Error("WebSocket is not connected"); + if (!socket || socket.readyState !== WebSocket.OPEN) { + throw new Error("WebSocket is not open"); } const requestId = randomUUID(); return new Promise((resolve, reject) => { - pending.set(requestId, { resolve, reject }); + const timeout = setTimeout(() => { + pending.delete(requestId); + reject(new Error("WebSocket request timed out")); + }, 10_000); + + pending.set(requestId, { resolve, reject, timeout }); + socket?.send(JSON.stringify({ ...payload, requestId })); }); } @@ -97,9 +169,13 @@ async function send( payload: Record, ): Promise { await ensureSocket(useLocalhost); - return await sendUnsafe(useLocalhost, payload); + return sendUnsafe(useLocalhost, payload); } +/* ────────────────────────── + Public API +────────────────────────── */ + export async function addRun( useLocalhost: boolean, problem: string,