import { describe, expect, it, mock } from "bun:test"; import { createRealtimeConnectionController, type RealtimeConnectionAttempt, type RealtimeSocket, } from "./realtime-connection.ts"; class FakeRealtimeSocket implements RealtimeSocket { static instances: FakeRealtimeSocket[] = []; onclose: ((event: { type: string }) => void) | null = null; onerror: ((event: { type: string }) => void) | null = null; onmessage: ((event: { data: unknown }) => void) | null = null; onopen: ((event: { type: string }) => void) | null = null; sent: string[] = []; closeCount = 0; constructor(readonly url: string) { FakeRealtimeSocket.instances.push(this); } close(): void { this.closeCount += 1; this.onclose?.({ type: "close" }); } open(): void { this.onopen?.({ type: "open" }); } send(data: string): void { this.sent.push(data); } } const createTimerHarness = () => { const timers = new Map void; delay: number }>(); let nextTimer = 1; return { clearTimeout: mock((timer: ReturnType) => { timers.delete(Number(timer)); }), runNext(): number | null { const entry = timers.entries().next().value as | [number, { callback: () => void; delay: number }] | undefined; if (!entry) { return null; } const [timer, value] = entry; timers.delete(timer); value.callback(); return value.delay; }, setTimeout: mock((callback: () => void, delay: number) => { const timer = nextTimer; nextTimer += 1; timers.set(timer, { callback, delay }); return timer as unknown as ReturnType; }), timers, }; }; const defaultAttempt: RealtimeConnectionAttempt = { registrationMessage: "register", url: "wss://realtime.example.test/v1", }; describe("createRealtimeConnectionController", () => { it("registers, refreshes registration, forwards messages, and stops", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const onMessage = mock(() => {}); const onRegistrationSent = mock(() => {}); const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage, onRegistrationSent, registrationRefreshMs: 100, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: FakeRealtimeSocket, }); connection.start(); const socket = FakeRealtimeSocket.instances[0]; expect(onRegistrationSent).not.toHaveBeenCalled(); socket?.open(); expect(socket?.sent).toEqual(["register"]); expect(onRegistrationSent).toHaveBeenCalledTimes(1); expect(onRegistrationSent.mock.calls[0]).toEqual(["register"]); socket?.onmessage?.({ data: "sync" }); expect(onMessage).toHaveBeenCalledWith("sync"); expect(timers.runNext()).toBe(100); expect(socket?.sent).toEqual(["register", "register"]); expect(onRegistrationSent).toHaveBeenCalledTimes(2); expect(onRegistrationSent.mock.calls[1]).toEqual(["register"]); connection.stop(); expect(timers.timers.size).toBe(0); expect(socket?.onclose).toBeDefined(); }); it("reconnects after the configured delay", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const onRegistrationSent = mock(() => {}); const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onRegistrationSent, reconnectDelayMs: 10, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: FakeRealtimeSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.open(); FakeRealtimeSocket.instances[0]?.close(); expect(timers.runNext()).toBe(10); expect(FakeRealtimeSocket.instances).toHaveLength(2); FakeRealtimeSocket.instances[1]?.open(); expect(onRegistrationSent.mock.calls).toEqual([["register"], ["register"]]); }); it("settles transport state before reporting a registration send", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); let connection: | ReturnType | undefined; const onRegistrationSent = mock(() => { connection?.stop(); }); connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onRegistrationSent, registrationRefreshMs: 100, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: FakeRealtimeSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.open(); expect(onRegistrationSent.mock.calls).toEqual([["register"]]); expect(FakeRealtimeSocket.instances[0]?.closeCount).toBe(1); expect(timers.timers.size).toBe(0); }); it("reports automatic reconnect setup failures and stops retrying", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const onAsyncTransportError = mock(() => {}); const onReconnectSetupError = mock(() => {}); let constructionCount = 0; class ReconnectFailureSocket extends FakeRealtimeSocket { constructor(url: string) { constructionCount += 1; if (constructionCount === 2) { throw new Error("reconnect constructor failed"); } super(url); } } const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onAsyncTransportError, onReconnectSetupError, reconnectDelayMs: 10, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: ReconnectFailureSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.close(); expect(() => timers.runNext()).not.toThrow(); expect(onAsyncTransportError).toHaveBeenCalledTimes(1); expect(onAsyncTransportError.mock.calls[0]?.[0]).toEqual( new Error("reconnect constructor failed"), ); expect(onReconnectSetupError).toHaveBeenCalledWith( new Error("reconnect constructor failed"), ); expect(timers.timers.size).toBe(0); connection.start(); expect(FakeRealtimeSocket.instances).toHaveLength(2); }); it("reports reconnect policy failures from the timer boundary", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const onAsyncTransportError = mock(() => {}); let policyCheckCount = 0; const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onAsyncTransportError, reconnectDelayMs: 10, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, shouldAttemptConnection: () => { policyCheckCount += 1; if (policyCheckCount === 3) { throw new Error("reconnect policy failed"); } return true; }, socketConstructor: FakeRealtimeSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.close(); expect(() => timers.runNext()).not.toThrow(); expect(onAsyncTransportError).toHaveBeenCalledWith( new Error("reconnect policy failed"), ); expect(timers.timers.size).toBe(0); }); it("contains asynchronous transport error observer failures", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); let constructionCount = 0; class ReconnectFailureSocket extends FakeRealtimeSocket { constructor(url: string) { constructionCount += 1; if (constructionCount === 2) { throw new Error("reconnect constructor failed"); } super(url); } } const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onAsyncTransportError: () => { throw new Error("observer failed"); }, reconnectDelayMs: 10, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: ReconnectFailureSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.close(); expect(() => timers.runNext()).not.toThrow(); expect(timers.timers.size).toBe(0); }); it("closes partially initialized reconnect sockets when handler setup throws", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const onAsyncTransportError = mock(() => {}); class PartialSetupFailureSocket extends FakeRealtimeSocket { constructor(url: string) { super(url); if (FakeRealtimeSocket.instances.length === 2) { Object.defineProperty(this, "onmessage", { configurable: true, set() { throw new Error("handler assignment failed"); }, }); } } } const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onAsyncTransportError, reconnectDelayMs: 10, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: PartialSetupFailureSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.close(); expect(() => timers.runNext()).not.toThrow(); expect(onAsyncTransportError).toHaveBeenCalledWith( new Error("handler assignment failed"), ); expect(FakeRealtimeSocket.instances[1]?.closeCount).toBe(1); expect(timers.timers.size).toBe(0); connection.stop(); expect(FakeRealtimeSocket.instances[1]?.closeCount).toBe(1); }); it("closes partially initialized reconnect sockets when late handler setup throws", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const onAsyncTransportError = mock(() => {}); class LateSetupFailureSocket extends FakeRealtimeSocket { constructor(url: string) { super(url); if (FakeRealtimeSocket.instances.length === 2) { Object.defineProperty(this, "onerror", { configurable: true, set() { throw new Error("late handler assignment failed"); }, }); } } } const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onAsyncTransportError, reconnectDelayMs: 10, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: LateSetupFailureSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.close(); expect(() => timers.runNext()).not.toThrow(); expect(onAsyncTransportError).toHaveBeenCalledWith( new Error("late handler assignment failed"), ); expect(FakeRealtimeSocket.instances[1]?.closeCount).toBe(1); expect(timers.timers.size).toBe(0); }); it("preserves direct setup errors when partial socket cleanup throws", () => { FakeRealtimeSocket.instances = []; class SetupAndCleanupFailureSocket extends FakeRealtimeSocket { constructor(url: string) { super(url); Object.defineProperty(this, "onmessage", { configurable: true, set() { throw new Error("handler assignment failed"); }, }); } override close(): void { this.closeCount += 1; throw new Error("socket cleanup failed"); } } const connection = createRealtimeConnectionController({ onMessage: () => {}, resolveConnectionAttempt: () => defaultAttempt, socketConstructor: SetupAndCleanupFailureSocket, }); expect(() => connection.start()).toThrow("handler assignment failed"); expect(FakeRealtimeSocket.instances[0]?.closeCount).toBe(1); expect(() => connection.stop()).not.toThrow(); }); it("reports reconnect setup errors when partial socket cleanup throws", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const onAsyncTransportError = mock(() => {}); class ReconnectSetupAndCleanupFailureSocket extends FakeRealtimeSocket { constructor(url: string) { super(url); if (FakeRealtimeSocket.instances.length === 2) { Object.defineProperty(this, "onerror", { configurable: true, set() { throw new Error("late handler assignment failed"); }, }); } } override close(): void { this.closeCount += 1; if (FakeRealtimeSocket.instances.length === 2) { throw new Error("socket cleanup failed"); } this.onclose?.({ type: "close" }); } } const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onAsyncTransportError, reconnectDelayMs: 10, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: ReconnectSetupAndCleanupFailureSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.close(); expect(() => timers.runNext()).not.toThrow(); expect(onAsyncTransportError).toHaveBeenCalledWith( new Error("late handler assignment failed"), ); expect(FakeRealtimeSocket.instances[1]?.closeCount).toBe(1); expect(timers.timers.size).toBe(0); }); it("start preempts a pending reconnect", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: FakeRealtimeSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.close(); connection.start(); expect(FakeRealtimeSocket.instances).toHaveLength(2); expect(timers.timers.size).toBe(0); }); it("ignores stale reconnect callbacks when timer handles are reused", () => { FakeRealtimeSocket.instances = []; let currentCallback: (() => void) | null = null; const connection = createRealtimeConnectionController({ clearTimeout: mock(() => { currentCallback = null; }), onMessage: () => {}, resolveConnectionAttempt: () => defaultAttempt, setTimeout: mock((callback: () => void) => { currentCallback = callback; return 1 as unknown as ReturnType; }), socketConstructor: FakeRealtimeSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.close(); const staleReconnect = currentCallback; connection.start(); FakeRealtimeSocket.instances[1]?.close(); const currentReconnect = currentCallback; staleReconnect?.(); expect(FakeRealtimeSocket.instances).toHaveLength(2); currentReconnect?.(); expect(FakeRealtimeSocket.instances).toHaveLength(3); }); it("keeps URL and registration data atomic for each socket attempt", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); let attempt = defaultAttempt; const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, registrationRefreshMs: 100, resolveConnectionAttempt: () => attempt, setTimeout: timers.setTimeout, socketConstructor: FakeRealtimeSocket, }); connection.start(); attempt = { registrationMessage: "register-next", url: "wss://next-realtime.example.test/v1", }; FakeRealtimeSocket.instances[0]?.open(); expect(timers.runNext()).toBe(100); expect(FakeRealtimeSocket.instances[0]).toMatchObject({ sent: ["register", "register"], url: defaultAttempt.url, }); FakeRealtimeSocket.instances[0]?.close(); expect(timers.runNext()).toBe(1000); FakeRealtimeSocket.instances[1]?.open(); expect(FakeRealtimeSocket.instances[1]).toMatchObject({ sent: ["register-next"], url: attempt.url, }); }); it("surfaces synchronous setup failures and recovers registration send failures", () => { class ConstructorFailureSocket extends FakeRealtimeSocket { constructor(url: string) { super(url); throw new Error("constructor failed"); } } const constructorFailure = createRealtimeConnectionController({ onMessage: () => {}, resolveConnectionAttempt: () => defaultAttempt, socketConstructor: ConstructorFailureSocket, }); expect(() => constructorFailure.start()).toThrow("constructor failed"); const resolverFailure = createRealtimeConnectionController({ onMessage: () => {}, resolveConnectionAttempt: () => { throw new Error("resolver failed"); }, socketConstructor: FakeRealtimeSocket, }); expect(() => resolverFailure.start()).toThrow("resolver failed"); FakeRealtimeSocket.instances = []; class SendFailureSocket extends FakeRealtimeSocket { override send(): void { throw new Error("send failed"); } } const timers = createTimerHarness(); const onRegistrationSent = mock(() => {}); const onAsyncTransportError = mock(() => {}); const sendFailure = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onAsyncTransportError, onRegistrationSent, reconnectDelayMs: 10, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: SendFailureSocket, }); sendFailure.start(); expect(() => FakeRealtimeSocket.instances[0]?.open()).not.toThrow(); expect(onAsyncTransportError).toHaveBeenCalledWith( new Error("send failed"), ); expect(onRegistrationSent).not.toHaveBeenCalled(); expect(FakeRealtimeSocket.instances[0]?.closeCount).toBe(1); expect(timers.runNext()).toBe(10); expect(FakeRealtimeSocket.instances).toHaveLength(2); }); it("reports registration refresh send failures and reconnects", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const onRegistrationSent = mock(() => {}); const onAsyncTransportError = mock(() => {}); class RefreshFailureSocket extends FakeRealtimeSocket { override send(data: string): void { if (this.sent.length > 0) { throw new Error("refresh failed"); } super.send(data); } } const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onAsyncTransportError, onRegistrationSent, reconnectDelayMs: 10, registrationRefreshMs: 100, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: RefreshFailureSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.open(); expect(() => timers.runNext()).not.toThrow(); expect(onAsyncTransportError).toHaveBeenCalledWith( new Error("refresh failed"), ); expect(onRegistrationSent.mock.calls).toEqual([["register"]]); expect(FakeRealtimeSocket.instances[0]?.closeCount).toBe(1); expect(timers.runNext()).toBe(10); expect(FakeRealtimeSocket.instances).toHaveLength(2); FakeRealtimeSocket.instances[1]?.open(); expect(onRegistrationSent.mock.calls).toEqual([["register"], ["register"]]); }); it("lets the transport error observer cancel send-failure recovery", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); class SendFailureSocket extends FakeRealtimeSocket { override send(): void { throw new Error("send failed"); } } let connection: | ReturnType | undefined; connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onAsyncTransportError: () => { connection?.stop(); }, onMessage: () => {}, reconnectDelayMs: 10, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: SendFailureSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.open(); expect(timers.timers.size).toBe(0); expect(FakeRealtimeSocket.instances).toHaveLength(1); }); it("contains registration observer failures without retiring the socket", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); const onAsyncTransportError = mock(() => {}); const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, onAsyncTransportError, onRegistrationSent: () => { throw new Error("observer failed"); }, registrationRefreshMs: 100, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, socketConstructor: FakeRealtimeSocket, }); connection.start(); expect(() => FakeRealtimeSocket.instances[0]?.open()).not.toThrow(); expect(FakeRealtimeSocket.instances[0]?.closeCount).toBe(0); expect(timers.timers.size).toBe(1); expect(() => timers.runNext()).not.toThrow(); expect(FakeRealtimeSocket.instances[0]?.closeCount).toBe(0); expect(timers.timers.size).toBe(1); expect(onAsyncTransportError).not.toHaveBeenCalled(); }); it("retains the next refresh when timer handles are reused", () => { FakeRealtimeSocket.instances = []; let currentCallback: (() => void) | null = null; const connection = createRealtimeConnectionController({ clearTimeout: mock(() => { currentCallback = null; }), onMessage: () => {}, registrationRefreshMs: 100, resolveConnectionAttempt: () => defaultAttempt, setTimeout: mock((callback: () => void) => { currentCallback = callback; return 1 as unknown as ReturnType; }), socketConstructor: FakeRealtimeSocket, }); connection.start(); const socket = FakeRealtimeSocket.instances[0]; socket?.open(); const firstRefresh = currentCallback; currentCallback = null; firstRefresh?.(); expect(currentCallback).not.toBeNull(); const secondRefresh = currentCallback; secondRefresh?.(); expect(socket?.sent).toEqual(["register", "register", "register"]); }); it("stops when no connection attempt is available", () => { FakeRealtimeSocket.instances = []; const connection = createRealtimeConnectionController({ onMessage: () => {}, resolveConnectionAttempt: () => null, socketConstructor: FakeRealtimeSocket, }); connection.start(); expect(FakeRealtimeSocket.instances).toHaveLength(0); }); it("rechecks attempt policy before reconnecting", () => { FakeRealtimeSocket.instances = []; const timers = createTimerHarness(); let reconnectAllowed = true; const connection = createRealtimeConnectionController({ clearTimeout: timers.clearTimeout, onMessage: () => {}, resolveConnectionAttempt: () => defaultAttempt, setTimeout: timers.setTimeout, shouldAttemptConnection: () => reconnectAllowed, socketConstructor: FakeRealtimeSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.close(); reconnectAllowed = false; expect(timers.runNext()).toBe(1000); expect(FakeRealtimeSocket.instances).toHaveLength(1); }); it("reports current socket errors", () => { FakeRealtimeSocket.instances = []; const onSocketError = mock(() => {}); const connection = createRealtimeConnectionController({ onMessage: () => {}, onSocketError, resolveConnectionAttempt: () => defaultAttempt, socketConstructor: FakeRealtimeSocket, }); connection.start(); FakeRealtimeSocket.instances[0]?.onerror?.({ type: "error" }); expect(onSocketError).toHaveBeenCalledWith({ type: "error" }); }); it("rejects invalid timing policies", () => { for (const options of [ { reconnectDelayMs: -1 }, { reconnectDelayMs: Number.NaN }, { registrationRefreshMs: 0 }, { registrationRefreshMs: Number.POSITIVE_INFINITY }, ]) { expect(() => createRealtimeConnectionController({ onMessage: () => {}, resolveConnectionAttempt: () => defaultAttempt, socketConstructor: FakeRealtimeSocket, ...options, }), ).toThrow(); } }); });