feat(chess): add WebSocket client library with reconnect (P4.9)
GameClient provides typed event-driven access to the chess server: - Envelope management: auto v/seq/ts, token captured from room.created/joined - Event emitter: on/off for game.state, game.delta, game.end, room.created, room.joined, error, connected, disconnected(willReconnect) - Exponential backoff reconnect: 1s, 2s, 4s, ... capped at 30s, max 10 attempts - Sequence-ack tracking via currentSeq (highest server seq seen; never regresses) - Dependency injection for WebSocket ctor and timers enables deterministic tests packages/chess/src/net/types.ts mirrors the server wire types without importing from @paratype/chess-server (wrong dependency direction). 21 tests cover connect flow, message dispatch, send semantics, reconnect state machine, seq tracking, and listener management.
This commit is contained in:
parent
2d0b8399f8
commit
39d91b6356
3 changed files with 1113 additions and 0 deletions
535
packages/chess/src/net/client.test.ts
Normal file
535
packages/chess/src/net/client.test.ts
Normal file
|
|
@ -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<string, unknown>;
|
||||
}
|
||||
|
||||
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<string, unknown> = {};
|
||||
try {
|
||||
const raw = JSON.parse(data);
|
||||
if (raw !== null && typeof raw === "object" && !Array.isArray(raw)) {
|
||||
parsed = raw as Record<string, unknown>;
|
||||
}
|
||||
} 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([]);
|
||||
});
|
||||
});
|
||||
446
packages/chess/src/net/client.ts
Normal file
446
packages/chess/src/net/client.ts
Normal file
|
|
@ -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<T extends GameClientEventType> = Extract<
|
||||
GameClientEvent,
|
||||
{ type: T }
|
||||
>;
|
||||
|
||||
type Listener<T extends GameClientEventType> = (event: EventOfType<T>) => 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<GameClientEventType, AnyListener[]>();
|
||||
|
||||
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<typeof setTimeout>);
|
||||
});
|
||||
}
|
||||
|
||||
// -------------------------------------------------------------------------
|
||||
// Event API
|
||||
// -------------------------------------------------------------------------
|
||||
|
||||
on<T extends GameClientEventType>(type: T, listener: Listener<T>): 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<T extends GameClientEventType>(type: T, listener: Listener<T>): 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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
return new Promise<void>((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 };
|
||||
132
packages/chess/src/net/types.ts
Normal file
132
packages/chess/src/net/types.ts
Normal file
|
|
@ -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<TType extends string, TPayload> {
|
||||
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<string, never>>
|
||||
| MessageEnvelope<"game.move", GameMovePayload>;
|
||||
Loading…
Add table
Add a link
Reference in a new issue