diff --git a/packages/chess/src/net/client.test.ts b/packages/chess/src/net/client.test.ts new file mode 100644 index 0000000..46935db --- /dev/null +++ b/packages/chess/src/net/client.test.ts @@ -0,0 +1,535 @@ +// Tests for GameClient. We inject a mock WebSocket constructor + a mock +// scheduler so every assertion is deterministic (no real sockets, no real +// timers). The mock is defined inline — no external libraries. + +import { describe, it, expect, beforeEach } from "vitest"; +import { GameClient } from "./client.js"; +import type { + GameDeltaPayload, + GameStatePayload, + RoomJoinedPayload, +} from "./types.js"; + +// --------------------------------------------------------------------------- +// Mock WebSocket +// --------------------------------------------------------------------------- +// +// A WebSocket-shaped double. Tests drive lifecycle explicitly via simulate* +// methods so timing is 100% deterministic. We don't extend the global +// WebSocket because we want to construct the mock with `new` from the client. + +interface SentFrame { + raw: string; + parsed: Record; +} + +class MockWebSocket { + // Constants matched to the lib.dom.d.ts shape (numeric states). + static CONNECTING = 0 as const; + static OPEN = 1 as const; + static CLOSING = 2 as const; + static CLOSED = 3 as const; + + readyState: number = MockWebSocket.CONNECTING; + url: string; + + // Event handlers — the client assigns these after construction. + onopen: ((ev: Event) => void) | null = null; + onmessage: ((ev: MessageEvent) => void) | null = null; + onerror: ((ev: Event) => void) | null = null; + onclose: ((ev: CloseEvent) => void) | null = null; + + // Recorded traffic (outgoing frames the client sent). + readonly sent: SentFrame[] = []; + // Intercepted close calls (so tests can assert the client closed voluntarily). + closeCalled = 0; + + constructor(url: string) { + this.url = url; + // Track every constructed instance so tests can reach the latest one. + MockWebSocket.instances.push(this); + } + + static instances: MockWebSocket[] = []; + static reset(): void { + MockWebSocket.instances = []; + } + static latest(): MockWebSocket { + const inst = MockWebSocket.instances.at(-1); + if (!inst) throw new Error("MockWebSocket: no instances yet"); + return inst; + } + + // --- Outbound (client → server) --- + send(data: string): void { + let parsed: Record = {}; + try { + const raw = JSON.parse(data); + if (raw !== null && typeof raw === "object" && !Array.isArray(raw)) { + parsed = raw as Record; + } + } catch { + // Leave parsed as an empty object; tests can still inspect `raw`. + } + this.sent.push({ raw: data, parsed }); + } + + close(): void { + this.closeCalled += 1; + // Mirror real WebSocket: close() puts us into CLOSED and fires onclose. + if (this.readyState !== MockWebSocket.CLOSED) { + this.readyState = MockWebSocket.CLOSED; + this.onclose?.(new CloseEvent("close")); + } + } + + // --- Inbound lifecycle helpers (driven by tests) --- + simulateOpen(): void { + this.readyState = MockWebSocket.OPEN; + this.onopen?.(new Event("open")); + } + simulateMessage(msg: unknown): void { + const data = JSON.stringify(msg); + this.onmessage?.(new MessageEvent("message", { data })); + } + simulateRawMessage(raw: string): void { + this.onmessage?.(new MessageEvent("message", { data: raw })); + } + simulateError(): void { + this.onerror?.(new Event("error")); + } + simulateClose(): void { + this.readyState = MockWebSocket.CLOSED; + this.onclose?.(new CloseEvent("close")); + } +} + +// --------------------------------------------------------------------------- +// Mock scheduler +// --------------------------------------------------------------------------- +// +// A deterministic replacement for setTimeout. `scheduled` holds pending +// callbacks; tests advance the clock by calling `flush()` or `runNext()`. + +interface Scheduled { + id: number; + cb: () => void; + delay: number; + cancelled: boolean; +} + +class MockScheduler { + private nextId = 1; + readonly scheduled: Scheduled[] = []; + + schedule = (cb: () => void, ms: number): number => { + const entry: Scheduled = { id: this.nextId++, cb, delay: ms, cancelled: false }; + this.scheduled.push(entry); + return entry.id; + }; + + cancel = (handle: unknown): void => { + if (typeof handle !== "number") return; + const hit = this.scheduled.find((s) => s.id === handle); + if (hit) hit.cancelled = true; + }; + + /** Run every pending (non-cancelled) callback, in insertion order. */ + flush(): void { + // Snapshot and clear so callbacks that schedule more work don't loop forever. + const pending = this.scheduled.splice(0, this.scheduled.length); + for (const entry of pending) { + if (!entry.cancelled) entry.cb(); + } + } +} + +// --------------------------------------------------------------------------- +// Fixture factory +// --------------------------------------------------------------------------- + +function makeClient(options?: { + maxReconnectAttempts?: number; +}): { client: GameClient; scheduler: MockScheduler } { + const scheduler = new MockScheduler(); + const client = new GameClient("ws://test.local/ws", { + WebSocketCtor: MockWebSocket as unknown as typeof WebSocket, + setTimeoutFn: scheduler.schedule, + clearTimeoutFn: scheduler.cancel, + baseReconnectDelayMs: 1000, + maxReconnectDelayMs: 30_000, + maxReconnectAttempts: options?.maxReconnectAttempts ?? 10, + }); + return { client, scheduler }; +} + +beforeEach(() => { + MockWebSocket.reset(); +}); + +// --------------------------------------------------------------------------- +// Tests +// --------------------------------------------------------------------------- + +describe("GameClient — connect()", () => { + it("opens a WebSocket to the configured URL", async () => { + const { client } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + const socket = MockWebSocket.latest(); + expect(socket.url).toBe("ws://test.local/ws"); + socket.simulateOpen(); + await expect(pending).resolves.toBeUndefined(); + }); + + it("sends room.join with token after open", async () => { + const { client } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + const socket = MockWebSocket.latest(); + socket.simulateOpen(); + await pending; + expect(socket.sent).toHaveLength(1); + const frame = socket.sent[0]!.parsed; + expect(frame["type"]).toBe("room.join"); + expect(frame["token"]).toBe("tok-1"); + expect(frame["v"]).toBe(1); + expect(frame["seq"]).toBe(1); + expect(typeof frame["ts"]).toBe("number"); + expect(frame["payload"]).toEqual({ code: "ABC123" }); + }); + + it('emits "connected" on open', async () => { + const { client } = makeClient(); + const seen: string[] = []; + client.on("connected", () => seen.push("connected")); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + expect(seen).toEqual(["connected"]); + }); + + it("rejects the connect promise if the socket errors before open", async () => { + const { client } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateError(); + await expect(pending).rejects.toThrow(/WebSocket/i); + }); +}); + +describe("GameClient — inbound messages", () => { + it('fires "game.state" with the payload', async () => { + const { client } = makeClient(); + const received: GameStatePayload[] = []; + client.on("game.state", (ev) => received.push(ev.payload)); + const pending = client.connect("ABC123", "tok-1"); + const socket = MockWebSocket.latest(); + socket.simulateOpen(); + await pending; + + const payload: GameStatePayload = { + facts: [{ id: 1, attr: "PieceType", value: "pawn" }], + turn: "white", + lastSeq: 5, + moveHistory: ["e2-e4"], + activeRules: [], + fen: "startpos", + }; + socket.simulateMessage({ + v: 1, seq: 5, ts: Date.now(), type: "game.state", payload, + }); + expect(received).toHaveLength(1); + expect(received[0]).toEqual(payload); + expect(client.currentSeq).toBe(5); + }); + + it('fires "game.delta" with the payload', async () => { + const { client } = makeClient(); + const received: GameDeltaPayload[] = []; + client.on("game.delta", (ev) => received.push(ev.payload)); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + + const payload: GameDeltaPayload = { + inserted: [{ id: 1, attr: "Position", value: 28 }], + retracted: [{ id: 1, attr: "Position", value: 12 }], + moveNotation: "e2e4", + turn: "black", + gameOver: null, + }; + MockWebSocket.latest().simulateMessage({ + v: 1, seq: 6, ts: Date.now(), type: "game.delta", payload, + }); + expect(received).toHaveLength(1); + expect(received[0]).toEqual(payload); + expect(client.currentSeq).toBe(6); + }); + + it("captures token on room.joined and uses it for subsequent sends", async () => { + const { client } = makeClient(); + const pending = client.connect("ABC123", "initial-tok"); + const socket = MockWebSocket.latest(); + socket.simulateOpen(); + await pending; + + const joinedPayload: RoomJoinedPayload = { + code: "ABC123", + token: "server-issued-tok", + color: "black", + activeRules: [], + }; + socket.simulateMessage({ + v: 1, seq: 1, ts: Date.now(), type: "room.joined", payload: joinedPayload, + }); + expect(client.currentToken).toBe("server-issued-tok"); + + // Next send should use the captured token. + client.sendMove("e2", "e4"); + const move = socket.sent.at(-1)!.parsed; + expect(move["token"]).toBe("server-issued-tok"); + expect(move["type"]).toBe("game.move"); + expect(move["payload"]).toEqual({ from: "e2", to: "e4" }); + }); + + it("ignores malformed JSON frames", async () => { + const { client } = makeClient(); + const seen: unknown[] = []; + client.on("game.state", (ev) => seen.push(ev.payload)); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + expect(() => + MockWebSocket.latest().simulateRawMessage("{not json"), + ).not.toThrow(); + expect(seen).toEqual([]); + }); + + it("ignores unknown message types without throwing", async () => { + const { client } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + expect(() => + MockWebSocket.latest().simulateMessage({ + v: 1, seq: 1, ts: Date.now(), type: "some.future.type", payload: {}, + }), + ).not.toThrow(); + }); +}); + +describe("GameClient — sendMove()", () => { + it("sends a well-formed game.move message", async () => { + const { client } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + const socket = MockWebSocket.latest(); + socket.simulateOpen(); + await pending; + + client.sendMove("e2", "e4"); + const move = socket.sent.at(-1)!.parsed; + expect(move["v"]).toBe(1); + expect(move["type"]).toBe("game.move"); + expect(move["token"]).toBe("tok-1"); + expect(move["payload"]).toEqual({ from: "e2", to: "e4" }); + // Seq is monotonic — room.join was 1, this one must be 2. + expect(move["seq"]).toBe(2); + }); + + it("includes promoteTo only when provided", async () => { + const { client } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + + client.sendMove("e7", "e8", "queen"); + const move = MockWebSocket.latest().sent.at(-1)!.parsed; + expect(move["payload"]).toEqual({ from: "e7", to: "e8", promoteTo: "queen" }); + }); + + it("drops sends when the socket is not open", async () => { + const { client } = makeClient(); + // Never simulateOpen(); readyState stays CONNECTING. + void client.connect("ABC123", "tok-1"); + client.sendMove("e2", "e4"); + expect(MockWebSocket.latest().sent).toHaveLength(0); + }); +}); + +describe("GameClient — disconnect & reconnect", () => { + it('fires "disconnected" with willReconnect: true on unexpected close', async () => { + const { client } = makeClient(); + const events: Array<{ willReconnect: boolean }> = []; + client.on("disconnected", (ev) => events.push({ willReconnect: ev.willReconnect })); + const pending = client.connect("ABC123", "tok-1"); + const socket = MockWebSocket.latest(); + socket.simulateOpen(); + await pending; + + socket.simulateClose(); + expect(events).toEqual([{ willReconnect: true }]); + }); + + it("schedules a reconnect attempt with exponential backoff", async () => { + const { client, scheduler } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + + MockWebSocket.latest().simulateClose(); + // One timer should be scheduled with the base delay (1000ms). + expect(scheduler.scheduled).toHaveLength(1); + expect(scheduler.scheduled[0]!.delay).toBe(1000); + + // Flushing creates a new mock socket. + const before = MockWebSocket.instances.length; + scheduler.flush(); + expect(MockWebSocket.instances.length).toBe(before + 1); + + // Second close → second scheduled delay should be 2000ms (2^1). + MockWebSocket.latest().simulateOpen(); + MockWebSocket.latest().simulateClose(); + // After reconnecting successfully, reconnectAttempts was reset to 0 on open. + // The fresh close starts the counter over at 1000ms. + expect(scheduler.scheduled).toHaveLength(1); + expect(scheduler.scheduled[0]!.delay).toBe(1000); + }); + + it("applies exponential backoff across consecutive failures", async () => { + const { client, scheduler } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + + // Close without re-opening → each retry attempt grows the delay. + MockWebSocket.latest().simulateClose(); + expect(scheduler.scheduled.at(-1)!.delay).toBe(1000); + scheduler.flush(); + + // New socket never opens → close it immediately. + MockWebSocket.latest().simulateClose(); + expect(scheduler.scheduled.at(-1)!.delay).toBe(2000); + scheduler.flush(); + + MockWebSocket.latest().simulateClose(); + expect(scheduler.scheduled.at(-1)!.delay).toBe(4000); + }); + + it("caps the backoff at maxReconnectDelayMs", async () => { + const { client, scheduler } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + + // Burn through 6 failed reconnects: 1s, 2s, 4s, 8s, 16s, 30s (capped). + for (let i = 0; i < 6; i++) { + MockWebSocket.latest().simulateClose(); + scheduler.flush(); + } + // The 6th scheduled delay is at the cap (2^5=32s → capped to 30s). + MockWebSocket.latest().simulateClose(); + const latest = scheduler.scheduled.at(-1)!; + expect(latest.delay).toBe(30_000); + }); + + it('emits willReconnect: false after max attempts', async () => { + const { client, scheduler } = makeClient({ maxReconnectAttempts: 2 }); + const events: Array<{ willReconnect: boolean }> = []; + client.on("disconnected", (ev) => events.push({ willReconnect: ev.willReconnect })); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + + // attempt 1 + MockWebSocket.latest().simulateClose(); + scheduler.flush(); + // attempt 2 + MockWebSocket.latest().simulateClose(); + scheduler.flush(); + // attempt 3 — this one exceeds the cap, so willReconnect is false. + MockWebSocket.latest().simulateClose(); + + expect(events.map((e) => e.willReconnect)).toEqual([true, true, false]); + // And no further timer should be scheduled. + expect(scheduler.scheduled).toHaveLength(0); + }); + + it("close() prevents reconnect", async () => { + const { client, scheduler } = makeClient(); + const events: Array<{ willReconnect: boolean }> = []; + client.on("disconnected", (ev) => events.push({ willReconnect: ev.willReconnect })); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + + client.close(); + expect(events).toEqual([{ willReconnect: false }]); + expect(scheduler.scheduled.filter((s) => !s.cancelled)).toHaveLength(0); + }); + + it("re-sends room.join on reconnect attempt", async () => { + const { client, scheduler } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + + const firstSocket = MockWebSocket.latest(); + firstSocket.simulateClose(); + scheduler.flush(); + + const secondSocket = MockWebSocket.latest(); + expect(secondSocket).not.toBe(firstSocket); + secondSocket.simulateOpen(); + // After the reconnected socket opens it should also send room.join. + const joinFrame = secondSocket.sent[0]!.parsed; + expect(joinFrame["type"]).toBe("room.join"); + expect(joinFrame["token"]).toBe("tok-1"); + expect(joinFrame["payload"]).toEqual({ code: "ABC123" }); + }); +}); + +describe("GameClient — sequence tracking", () => { + it("tracks the highest server seq seen", async () => { + const { client } = makeClient(); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + const socket = MockWebSocket.latest(); + + socket.simulateMessage({ + v: 1, seq: 5, ts: Date.now(), type: "game.state", + payload: { facts: [], turn: "white", lastSeq: 5, moveHistory: [], activeRules: [], fen: "" }, + }); + expect(client.currentSeq).toBe(5); + + socket.simulateMessage({ + v: 1, seq: 8, ts: Date.now(), type: "game.delta", + payload: { inserted: [], retracted: [], moveNotation: "e2e4", turn: "black", gameOver: null }, + }); + expect(client.currentSeq).toBe(8); + + // Out-of-order older message should NOT lower the seq. + socket.simulateMessage({ + v: 1, seq: 6, ts: Date.now(), type: "game.delta", + payload: { inserted: [], retracted: [], moveNotation: "e7e5", turn: "white", gameOver: null }, + }); + expect(client.currentSeq).toBe(8); + }); +}); + +describe("GameClient — listener management", () => { + it("off() removes a listener", async () => { + const { client } = makeClient(); + const seen: number[] = []; + const listener = () => seen.push(1); + client.on("game.state", listener); + const pending = client.connect("ABC123", "tok-1"); + MockWebSocket.latest().simulateOpen(); + await pending; + + client.off("game.state", listener); + MockWebSocket.latest().simulateMessage({ + v: 1, seq: 1, ts: Date.now(), type: "game.state", + payload: { facts: [], turn: "white", lastSeq: 1, moveHistory: [], activeRules: [], fen: "" }, + }); + expect(seen).toEqual([]); + }); +}); diff --git a/packages/chess/src/net/client.ts b/packages/chess/src/net/client.ts new file mode 100644 index 0000000..2c7322c --- /dev/null +++ b/packages/chess/src/net/client.ts @@ -0,0 +1,446 @@ +// WebSocket client library for the chess server. +// +// Responsibilities: +// * Open/close a WebSocket connection to the game server +// * Authenticate via room.create / room.join (token managed internally) +// * Track server `seq` numbers for reconnect reconciliation +// * Emit typed events for every server message + connection lifecycle +// * Reconnect on unexpected close with exponential backoff +// +// The client is intentionally dumb: it does NOT keep game state, it does NOT +// predict moves (that's P4.10). It is a thin, typed transport. + +import type { + ClientMessage, + ErrorPayload, + GameDeltaPayload, + GameEndPayload, + GameMovePayload, + GameStatePayload, + PromotionPiece, + RoomCreatedPayload, + RoomJoinedPayload, +} from "./types.js"; + +// --------------------------------------------------------------------------- +// Event surface +// --------------------------------------------------------------------------- + +export type GameClientEvent = + | { type: "game.state"; payload: GameStatePayload } + | { type: "game.delta"; payload: GameDeltaPayload } + | { type: "game.end"; payload: GameEndPayload } + | { type: "room.created"; payload: RoomCreatedPayload } + | { type: "room.joined"; payload: RoomJoinedPayload } + | { type: "error"; payload: ErrorPayload } + | { type: "connected" } + | { type: "disconnected"; willReconnect: boolean }; + +export type GameClientEventType = GameClientEvent["type"]; + +type EventOfType = Extract< + GameClientEvent, + { type: T } +>; + +type Listener = (event: EventOfType) => void; + +// Heterogeneous internal listener map — narrowed via the public `on()` API. +// Using `unknown` avoids `any` while still permitting one map for all types. +type AnyListener = (event: GameClientEvent) => void; + +// Shape emitted to listeners of a lifecycle-only event. We accept these two +// shapes when callers invoke `emit()` so the compiler stays honest about the +// discriminated union without resorting to casts. +interface LifecycleConnected { + type: "connected"; +} +interface LifecycleDisconnected { + type: "disconnected"; + willReconnect: boolean; +} + +// --------------------------------------------------------------------------- +// Configuration +// --------------------------------------------------------------------------- + +export interface GameClientOptions { + /** Max number of consecutive reconnect attempts before giving up. */ + maxReconnectAttempts?: number; + /** Base delay in ms (doubled each attempt). */ + baseReconnectDelayMs?: number; + /** Upper bound on any single backoff delay. */ + maxReconnectDelayMs?: number; + /** WebSocket constructor to use. Defaults to globalThis.WebSocket. Used for tests. */ + WebSocketCtor?: typeof WebSocket; + /** Scheduler for delayed reconnect. Defaults to setTimeout. Used for tests. */ + setTimeoutFn?: (cb: () => void, ms: number) => unknown; + /** Cancellation for `setTimeoutFn`. Defaults to clearTimeout. Used for tests. */ + clearTimeoutFn?: (handle: unknown) => void; +} + +// --------------------------------------------------------------------------- +// Envelope factory +// --------------------------------------------------------------------------- + +const PROTOCOL_VERSION = 1 as const; + +interface OutgoingEnvelope { + v: 1; + seq: number; + ts: number; + type: string; + token?: string; + payload: unknown; +} + +// Minimal server-message shape we actually read at this layer. Full validation +// is not re-done on the client — the server is authoritative and schema- +// checked; the client trusts `type`/`payload`/`seq` presence and forwards. +interface IncomingEnvelope { + type?: unknown; + seq?: unknown; + payload?: unknown; +} + +// --------------------------------------------------------------------------- +// GameClient +// --------------------------------------------------------------------------- + +export class GameClient { + // Configuration (immutable after construction). + private readonly url: string; + private readonly maxReconnectAttempts: number; + private readonly baseReconnectDelayMs: number; + private readonly maxReconnectDelayMs: number; + private readonly WSCtor: typeof WebSocket; + private readonly scheduleTimer: (cb: () => void, ms: number) => unknown; + private readonly cancelTimer: (handle: unknown) => void; + + // Session state. + private ws: WebSocket | null = null; + private code: string | null = null; + private token: string | null = null; + private lastSeq = 0; + private clientSeq = 0; + private reconnectAttempts = 0; + private reconnectTimer: unknown = null; + // `closed` is set when close() is called: it suppresses auto-reconnect. + private closed = false; + + // Listeners keyed by event type. + private readonly listeners = new Map(); + + constructor(url: string, options: GameClientOptions = {}) { + this.url = url; + this.maxReconnectAttempts = options.maxReconnectAttempts ?? 10; + this.baseReconnectDelayMs = options.baseReconnectDelayMs ?? 1000; + this.maxReconnectDelayMs = options.maxReconnectDelayMs ?? 30_000; + + // Resolve WebSocket / timer primitives from options with sensible global + // fallbacks. This is the injection seam used by the test harness. + const wsCtor = options.WebSocketCtor ?? globalThis.WebSocket; + if (typeof wsCtor !== "function") { + throw new Error( + "GameClient: no WebSocket constructor available (pass WebSocketCtor in options)", + ); + } + this.WSCtor = wsCtor; + this.scheduleTimer = + options.setTimeoutFn ?? + ((cb, ms) => globalThis.setTimeout(cb, ms) as unknown); + this.cancelTimer = + options.clearTimeoutFn ?? + ((handle) => { + globalThis.clearTimeout(handle as ReturnType); + }); + } + + // ------------------------------------------------------------------------- + // Event API + // ------------------------------------------------------------------------- + + on(type: T, listener: Listener): void { + const arr = this.listeners.get(type); + // We up-cast to AnyListener here because the map is heterogeneous; the + // public `on` signature guarantees each listener only ever receives its + // own discriminated variant, which we enforce at `emit()` sites. + const cast = listener as unknown as AnyListener; + if (arr) arr.push(cast); + else this.listeners.set(type, [cast]); + } + + off(type: T, listener: Listener): void { + const arr = this.listeners.get(type); + if (!arr) return; + const cast = listener as unknown as AnyListener; + const idx = arr.indexOf(cast); + if (idx >= 0) arr.splice(idx, 1); + } + + // ------------------------------------------------------------------------- + // Connection API + // ------------------------------------------------------------------------- + + /** + * Open a WebSocket, send `room.join` with the given token, and resolve once + * the socket is open. The promise rejects only if the FIRST connection + * attempt fails (subsequent drops trigger the reconnect state machine). + * + * Callers who don't yet have a token (creating a fresh room) can call + * `connectAndCreate()` instead. + */ + async connect(code: string, token: string): Promise { + this.code = code; + this.token = token; + this.closed = false; + this.reconnectAttempts = 0; + return this.openConnection({ autoJoin: true }); + } + + /** + * Open a WebSocket and send `room.create` without a token. The resulting + * `room.created` event will carry the server-issued token; consumers should + * listen for that to capture it (or use `connect()` once they have one). + */ + async connectAndCreate(rulesetIds?: string[]): Promise { + this.code = null; + this.token = null; + this.closed = false; + this.reconnectAttempts = 0; + // Use a one-shot bootstrap that fires room.create instead of room.join + // after the socket opens. We keep the bootstrap local to avoid leaking + // create-only state onto the instance. + return this.openConnection({ autoCreate: rulesetIds }); + } + + /** + * Close the socket permanently. Suppresses auto-reconnect. + */ + close(): void { + this.closed = true; + if (this.reconnectTimer !== null) { + this.cancelTimer(this.reconnectTimer); + this.reconnectTimer = null; + } + this.ws?.close(); + this.ws = null; + } + + // ------------------------------------------------------------------------- + // Send API + // ------------------------------------------------------------------------- + + /** + * Send a raw client message. The envelope (v/seq/ts) is added automatically. + * Messages sent while the socket is not OPEN are silently dropped; callers + * that need durability should queue at a higher layer. + */ + send(partial: { type: string; payload: unknown }, token?: string): void { + const socket = this.ws; + if (!socket || socket.readyState !== this.WSCtor.OPEN) return; + const effectiveToken = token ?? this.token ?? undefined; + const envelope: OutgoingEnvelope = { + v: PROTOCOL_VERSION, + seq: ++this.clientSeq, + ts: Date.now(), + type: partial.type, + payload: partial.payload, + ...(effectiveToken !== undefined ? { token: effectiveToken } : {}), + }; + socket.send(JSON.stringify(envelope)); + } + + /** Convenience: send a `game.move` with the stored token. */ + sendMove(from: string, to: string, promoteTo?: PromotionPiece): void { + const payload: GameMovePayload = promoteTo !== undefined + ? { from, to, promoteTo } + : { from, to }; + this.send({ type: "game.move", payload }); + } + + // ------------------------------------------------------------------------- + // Accessors (primarily for tests & reconnect logic) + // ------------------------------------------------------------------------- + + get currentSeq(): number { + return this.lastSeq; + } + + get isConnected(): boolean { + return this.ws !== null && this.ws.readyState === this.WSCtor.OPEN; + } + + get currentToken(): string | null { + return this.token; + } + + get currentCode(): string | null { + return this.code; + } + + // ------------------------------------------------------------------------- + // Internals + // ------------------------------------------------------------------------- + + private openConnection( + bootstrap: { + autoJoin?: boolean; + autoCreate?: string[] | undefined; + } = {}, + ): Promise { + return new Promise((resolve, reject) => { + let settled = false; + const ws = new this.WSCtor(this.url); + this.ws = ws; + + ws.onopen = () => { + this.reconnectAttempts = 0; + this.emit({ type: "connected" }); + if (!settled) { + settled = true; + resolve(); + } + // Fire the bootstrap message AFTER resolving so callers that await + // `connect()` see the open state when their awaiter runs. + if (bootstrap.autoJoin && this.code !== null) { + this.send( + { type: "room.join", payload: { code: this.code } }, + this.token ?? undefined, + ); + } else if (bootstrap.autoCreate !== undefined) { + const payload = + bootstrap.autoCreate.length > 0 + ? { rulesetIds: bootstrap.autoCreate } + : {}; + this.send({ type: "room.create", payload }); + } + }; + + ws.onmessage = (event: MessageEvent) => { + // `event.data` is typed as `any` by lib.dom — narrow defensively. + const data = event.data; + if (typeof data === "string") this.handleIncoming(data); + }; + + ws.onerror = () => { + if (!settled) { + settled = true; + reject(new Error("WebSocket connection failed")); + } + // Non-initial errors are handled via onclose (browsers fire close + // after error for any failed connection). No extra work here. + }; + + ws.onclose = () => { + const willReconnect = + !this.closed && this.reconnectAttempts < this.maxReconnectAttempts; + this.emit({ type: "disconnected", willReconnect }); + // Release the socket BEFORE scheduling the next attempt so + // `isConnected` reports false in listener callbacks. + this.ws = null; + if (willReconnect) this.scheduleReconnect(); + }; + }); + } + + private scheduleReconnect(): void { + if (this.closed) return; + if (this.reconnectAttempts >= this.maxReconnectAttempts) return; + const attempt = this.reconnectAttempts; + this.reconnectAttempts += 1; + const delay = Math.min( + this.baseReconnectDelayMs * Math.pow(2, attempt), + this.maxReconnectDelayMs, + ); + this.reconnectTimer = this.scheduleTimer(() => { + this.reconnectTimer = null; + if (this.closed) return; + // Reconnects always re-send room.join with the stored token — a + // create-only client that never captured a token cannot meaningfully + // reconnect, so we require `code` to be set. + if (this.code === null) return; + this.openConnection({ autoJoin: true }).catch(() => { + // Failure here will itself fire onclose → reschedule. Swallowing + // avoids an unhandled rejection on every retry. + }); + }, delay); + } + + private handleIncoming(raw: string): void { + let parsed: unknown; + try { + parsed = JSON.parse(raw); + } catch { + return; // Malformed frames are silently ignored at this layer. + } + if (typeof parsed !== "object" || parsed === null) return; + const msg = parsed as IncomingEnvelope; + + if (typeof msg.seq === "number" && msg.seq > this.lastSeq) { + this.lastSeq = msg.seq; + } + + if (typeof msg.type !== "string") return; + // Capture the server-issued token on room.created / room.joined so that + // subsequent sends and reconnects are authenticated without the caller + // having to wire it up manually. + this.captureAuth(msg); + + this.dispatchServerMessage(msg.type, msg.payload); + } + + private captureAuth(msg: IncomingEnvelope): void { + if (msg.type !== "room.created" && msg.type !== "room.joined") return; + const p = msg.payload; + if (typeof p !== "object" || p === null) return; + const rec = p as { code?: unknown; token?: unknown }; + if (typeof rec.code === "string") this.code = rec.code; + if (typeof rec.token === "string") this.token = rec.token; + } + + private dispatchServerMessage(type: string, payload: unknown): void { + // The client trusts the server's schema; we don't revalidate here. + // Each branch narrows `payload` with a cast that's specific to the + // event type we're constructing, keeping the union sound at emit sites. + switch (type) { + case "game.state": + this.emit({ type, payload: payload as GameStatePayload }); + return; + case "game.delta": + this.emit({ type, payload: payload as GameDeltaPayload }); + return; + case "game.end": + this.emit({ type, payload: payload as GameEndPayload }); + return; + case "room.created": + this.emit({ type, payload: payload as RoomCreatedPayload }); + return; + case "room.joined": + this.emit({ type, payload: payload as RoomJoinedPayload }); + return; + case "error": + this.emit({ type, payload: payload as ErrorPayload }); + return; + default: + // Unknown message types are ignored (forward-compatible with servers + // that add new message types — they're a no-op on old clients). + return; + } + } + + private emit(event: GameClientEvent): void; + private emit(event: LifecycleConnected): void; + private emit(event: LifecycleDisconnected): void; + private emit(event: GameClientEvent): void { + const arr = this.listeners.get(event.type); + if (!arr) return; + // Iterate over a copy so listeners that call `off()` mid-dispatch don't + // skip subsequent listeners. + for (const listener of arr.slice()) listener(event); + } +} + +// Re-export the ClientMessage type for convenience (unused re-export guarded +// by verbatimModuleSyntax via `export type`). +export type { ClientMessage }; diff --git a/packages/chess/src/net/types.ts b/packages/chess/src/net/types.ts new file mode 100644 index 0000000..9ed212d --- /dev/null +++ b/packages/chess/src/net/types.ts @@ -0,0 +1,132 @@ +// Client-side protocol types. Derived from packages/server/PROTOCOL.md, NOT +// imported from @paratype/chess-server (wrong dependency direction). +// +// These are kept structurally compatible with the server's Zod-inferred types +// but live independently so the chess package does not depend on the server. + +export type Color = "white" | "black"; +export type Winner = "white" | "black" | "draw"; +export type PromotionPiece = "queen" | "rook" | "bishop" | "knight"; + +export type GameEndReason = + | "checkmate" + | "stalemate" + | "50-move" + | "threefold" + | "insufficient" + | "player_left"; + +export type ErrorCode = + | "ILLEGAL_MOVE" + | "NOT_YOUR_TURN" + | "GAME_OVER" + | "ROOM_NOT_FOUND" + | "ROOM_FULL" + | "SERVER_FULL" + | "VERSION_MISMATCH" + | "RATE_LIMIT" + | "MSG_TOO_LARGE" + | "BAD_TOKEN" + | "INVALID_MESSAGE"; + +export interface Fact { + id: number; + attr: string; + value: unknown; +} + +// --------------------------------------------------------------------------- +// Server → Client payloads +// --------------------------------------------------------------------------- + +export interface GameStatePayload { + facts: Fact[]; + turn: Color; + lastSeq: number; + moveHistory: string[]; + activeRules: string[]; + fen: string; +} + +export interface GameOver { + winner: Winner; + reason: GameEndReason; +} + +export interface GameDeltaPayload { + inserted: Fact[]; + retracted: Fact[]; + moveNotation: string; + turn: Color; + gameOver: GameOver | null; +} + +export interface GameEndPayload { + winner: Winner; + reason: string; + finalFen: string; +} + +export interface RoomCreatedPayload { + code: string; + token: string; + color: Color; +} + +export interface RoomJoinedPayload { + code: string; + token: string; + color: Color; + activeRules: string[]; +} + +export interface ErrorPayload { + code: string; + message: string; + fatal: boolean; +} + +// --------------------------------------------------------------------------- +// Client → Server payloads (shape for send()) +// --------------------------------------------------------------------------- + +export interface RoomCreatePayload { + rulesetIds?: string[]; +} + +export interface RoomJoinPayload { + code: string; +} + +export interface GameMovePayload { + from: string; + to: string; + promoteTo?: PromotionPiece; +} + +// --------------------------------------------------------------------------- +// Envelope types (for tests and serialization) +// --------------------------------------------------------------------------- + +export interface MessageEnvelope { + v: 1; + seq: number; + ts: number; + type: TType; + token?: string; + payload: TPayload; +} + +export type ServerMessage = + | MessageEnvelope<"game.state", GameStatePayload> + | MessageEnvelope<"game.delta", GameDeltaPayload> + | MessageEnvelope<"game.end", GameEndPayload> + | MessageEnvelope<"room.created", RoomCreatedPayload> + | MessageEnvelope<"room.joined", RoomJoinedPayload> + | MessageEnvelope<"error", ErrorPayload>; + +export type ClientMessage = + | MessageEnvelope<"room.create", RoomCreatePayload> + | MessageEnvelope<"room.join", RoomJoinPayload> + | MessageEnvelope<"room.leave", Record> + | MessageEnvelope<"game.move", GameMovePayload>;