feat(thressgame-coverage): Wave 8 (WS protocol v2 + suspended execution + request-choice)

- T43: WS protocol v2 schema; protocolVersion field; RequestChoice/SubmitChoice/ProtocolVersionMismatch messages; v1 backward-compat
- T44: server-side request-choice broadcast on push; submit-choice validation (kind/forPlayer/value-type); ordered LIFO matching
- T45: PendingChoices stack on GAME_ENTITY; pushPendingChoice/popPendingChoice/peekPendingChoice helpers; serializePendingChoice (Map<->Array roundtrip); MAX_CHOICE_DEPTH=8 enforced
- T46: submitChoiceAndResume(engine, choiceId, value); descriptor-by-id lookup; bindings restored; remaining primitives executed via runPrimitives from primitiveIndex+1
- T47: request-choice primitive; SuspendedExecution exception mechanism; dispatcher catches and stops sibling iteration; deterministic choiceId via session counter
- T48: AutoChoiceResolver test transport (answersByKind / answersById); drainPendingChoices LIFO walk
- T49: server-side choice timeout enforcement; auto-resolve to first-option-per-kind; disconnect handler (forfeit / pause)
- T50: ChoiceTimeoutPolicy on GAME_ENTITY (timeout-with-default | no-timeout); CreateGameRequest extended; default 60s

Tests: 2533 -> 2658 (+125). bun run check exit 0.
This commit is contained in:
Joey Yakimowich-Payne 2026-04-26 11:54:24 -06:00
commit d4931a50ee
No known key found for this signature in database
37 changed files with 6841 additions and 295 deletions

View file

@ -28,7 +28,10 @@ import {
import { RateLimiter } from "./middleware.js";
import {
PROTOCOL_VERSION,
validateMessageString,
SUPPORTED_PROTOCOL_VERSIONS,
negotiateVersion,
shouldSkipV2Broadcast,
validateAnyMessageString,
type CustomModifierRegisterPayload,
type ErrorCode,
type Fact as WireFact,
@ -39,19 +42,33 @@ import {
type ModifierProfileRejectReason,
type ModifierProfileUpdatePayload,
type PresetActivation,
type RequestChoice,
type RoomCreatePayload,
type RoomJoinPayload,
type RoomSetPresetsPayload,
type ServerMessage,
type SubmitChoice,
type SupportedProtocolVersion,
type V2Message,
} from "./protocol.js";
import { DEFAULT_GRACE_MS, reconnectManager } from "./reconnect.js";
import { RoomRegistry } from "./rooms.js";
import { resolveLayoutRequest, toResolvedLayout } from "./layouts.js";
import {
choiceTimeoutManager,
decideDisconnectAction,
firstDefaultForKind,
getChoiceTimeoutPolicy,
hasPendingChoice,
} from "./choice-timeout.js";
import {
peekPendingChoice,
popPendingChoice,
validateProfile,
type ActionResult,
type ModifierProfile,
type ModifierValidationErrorCode,
type PendingChoice,
type PlayerAction,
} from "@paratype/chess";
@ -83,6 +100,16 @@ export interface ClientData {
token?: string;
/** Lazy-initialised on first message; one bucket per connection. */
rateLimiter?: RateLimiter;
/**
* T43: client capability version negotiated on the first frame
* via `negotiateVersion`. Pinned for the lifetime of the
* connection — re-negotiation isn't supported. Defaults to v1
* (the legacy behaviour) until the first `room.create` /
* `room.join` frame declares otherwise. The broadcast layer
* consults `shouldSkipV2Broadcast(protocolVersion)` before
* sending v2-only frames (request-choice).
*/
protocolVersion?: SupportedProtocolVersion;
}
// ---------------------------------------------------------------------------
@ -130,6 +157,67 @@ export function unregisterConnection(ws: ServerWebSocket<ClientData>): void {
roomRegistry.markDisconnected(roomCode, token);
// T49 — choice-timeout disconnect handler. If the leaver had an
// unresolved choice on the engine's stack, the active
// `ChoiceTimeoutPolicy` (T50) determines the outcome:
// - timeout-with-default → forfeit immediately (the leaver loses;
// opponent gets game.end; session is torn down). Skips the
// standard reconnect-grace window because the choice flow has
// already pinned the leaver to a synchronous decision.
// - no-timeout → enter a paused state. We mark the
// room paused but DO NOT start the grace timer — `no-timeout`
// means "wait indefinitely", so a 60s grace expiring into a
// game.end would defeat the policy. The leaver's reconnect
// path clears the pause flag and resumes the prompt.
// - none (no pending) → fall through to the standard grace.
const session = sessionRegistry.get(roomCode);
const action =
session !== undefined
? decideDisconnectAction(
getChoiceTimeoutPolicy(session.getEngine()),
hasPendingChoice(session.getEngine()),
)
: "none";
if (action === "forfeit") {
// Cancel any pending auto-default timer on this room — the
// forfeit decision supersedes the auto-resolve path.
choiceTimeoutManager.cancelAll(roomCode);
broadcastGameEnd(roomCode, token, "player_left");
roomRegistry.leaveRoom(roomCode, token);
if (roomRegistry.getRoom(roomCode) === undefined) {
sessionRegistry.delete(roomCode);
}
setActiveRooms(roomRegistry.getRoomCount());
logger
.child({ roomCode })
.info(
{ token },
"T49: forfeit on disconnect (timeout-with-default + pending choice)",
);
return;
}
if (action === "pause") {
// Mark the room paused. Subsequent move attempts from the
// remaining player are gated downstream (out of T49 scope —
// T49 owns the *transition into* paused, not the move-gate
// policy). No grace timer is started so the room persists
// until the leaver reconnects or the opponent explicitly
// leaves.
const room = roomRegistry.getRoom(roomCode);
if (room !== undefined) {
room.pausedByChoiceDisconnect = { byToken: token, since: Date.now() };
}
logger
.child({ roomCode })
.info(
{ token },
"T49: paused on disconnect (no-timeout + pending choice)",
);
return;
}
// Capture roomCode + token into the closure so onExpire doesn't need
// any ambient `this`. The closure runs up to DEFAULT_GRACE_MS later
// on the event loop — by then ws.data may already be GC'd.
@ -184,6 +272,263 @@ function broadcastToRoom(code: string, msg: ServerMessage): void {
}
}
// ---------------------------------------------------------------------------
// T44 — request-choice broadcast / submit-choice validation
// ---------------------------------------------------------------------------
//
// When the engine pushes a PendingChoice (T45 stack) — typically via the
// request-choice primitive (T47) firing inside a trigger cascade — the WS
// layer is responsible for surfacing the prompt to the appropriate
// client(s). The flow is:
//
// 1. Engine pushes a frame onto `GAME_ENTITY.PendingChoices`.
// 2. broadcast layer calls `broadcastTopChoiceIfNew(roomCode, session)`
// which peeks the top, confirms it hasn't already been broadcast for
// this session, and emits a v2 `request-choice` frame to every
// connected v2 client whose color matches `forPlayer`. v1 clients
// are skipped (T43 negotiation; no fall back to auto-resolve here —
// that's T49's job).
// 3. Client(s) submit a `submit-choice` frame; the v2 dispatch in
// `handleV2Frame` validates choiceId-vs-top and value-vs-kind, then
// pops the frame (T44 owns validation; the engine resume mechanism
// lands in T46 — until then the popped value is dropped on the
// floor with a logger note).
//
// Per-session bookkeeping: we track the set of choiceIds already
// broadcast for each room so a re-entry into `broadcastTopChoiceIfNew`
// after a downstream applyMove doesn't double-emit the same prompt.
// The set is cleared whenever a choiceId is popped (resolved) so the
// memory footprint stays O(stack-depth) ≤ MAX_CHOICE_DEPTH = 8.
const broadcastedChoiceIds = new Map<string, Set<string>>();
function getBroadcastedSet(roomCode: string): Set<string> {
let set = broadcastedChoiceIds.get(roomCode);
if (set === undefined) {
set = new Set();
broadcastedChoiceIds.set(roomCode, set);
}
return set;
}
/**
* Build the v2 request-choice frame from a PendingChoice. Bindings are
* deliberately NOT serialised onto the wire — they're internal engine
* resume state, not client-relevant. `options` is left unset for now;
* future work (T47) populates it from the request-choice primitive's
* params when the legal value-space is known up front.
*/
function buildRequestChoice(choice: PendingChoice): RequestChoice {
const frame: RequestChoice = {
kind: "request-choice",
protocolVersion: 2,
choiceId: choice.choiceId,
descriptorId: choice.descriptorId,
prompt: choice.prompt,
choiceKind: choice.kind,
forPlayer: choice.forPlayer,
};
// Only attach optional fields when present so the wire shape stays
// compact and v2 clients see the same shape they assert against.
if (choice.timeout !== undefined) {
return { ...frame, timeout: choice.timeout };
}
if (choice.expiresAtTimestamp !== undefined) {
return { ...frame, expiresAtTimestamp: choice.expiresAtTimestamp };
}
return frame;
}
/**
* Send a v2 request-choice frame to every connected client in `roomCode`
* whose color matches `forPlayer` and who negotiated protocolVersion >= 2.
* v1 clients are SKIPPED (T43): the v1 envelope has no `request-choice`
* carrier and a v1 client wouldn't know what to do with one. The trigger
* suspension primitive falls back to auto-resolve via T49 in that case.
*/
function sendRequestChoice(
roomCode: string,
choice: PendingChoice,
): void {
const room = roomRegistry.getRoom(roomCode);
if (!room) return;
const frame = buildRequestChoice(choice);
const json = JSON.stringify(frame);
for (const ws of getConnectionsInRoom(roomCode)) {
// v1 client → skip; the suspension semantics fall back to
// auto-resolve elsewhere in the pipeline.
const negotiated = ws.data.protocolVersion ?? 1;
if (shouldSkipV2Broadcast(negotiated)) continue;
// forPlayer gate — `"both"` means everyone, otherwise the
// single-color value restricts the prompt to that player. We
// resolve the socket's color from the room's player record
// because `ws.data` doesn't carry it.
if (choice.forPlayer !== "both") {
const player =
ws.data.token !== undefined
? room.players.get(ws.data.token)
: undefined;
if (!player || player.color !== choice.forPlayer) continue;
}
ws.send(json);
}
}
/**
* Inspect the top of the engine's PendingChoices stack and broadcast it
* to the appropriate clients, IF it hasn't already been broadcast for
* this room. Idempotent — safe to call after every state-mutating
* operation; only newly-pushed frames generate wire traffic.
*
* Exported for tests + the choice-flow integration points (`applyMove`
* post-tick, `performAction` post-tick, manual push for testing). The
* helper deliberately doesn't pop — the resolve path (submit-choice in
* `handleV2Frame`) owns popping.
*/
export function broadcastTopChoiceIfNew(
roomCode: string,
session: GameSession,
): void {
const top = peekPendingChoice(session.getEngine());
if (top === undefined) return;
const set = getBroadcastedSet(roomCode);
if (set.has(top.choiceId)) return;
set.add(top.choiceId);
sendRequestChoice(roomCode, top);
// T49 — arm the timeout timer at the same moment the prompt becomes
// visible to the client. We consult the engine's policy fact (T50)
// every call so a future "policy changes mid-game" pathway lands
// here automatically. The handler is captured into the closure so
// the timer fires against the SAME (roomCode, top) pair even if the
// stack churns underneath us before expiry.
armChoiceTimeoutFor(roomCode, session, top);
}
/**
* T49 — install a per-(roomCode, choiceId) timer that auto-resolves
* the prompt when the policy is `timeout-with-default` and the
* configured number of seconds elapses without a real
* `submit-choice`. No-op when:
* - the policy is `no-timeout` (no timer is ever armed); OR
* - a timer for this choiceId is already armed (re-broadcast
* idempotency — the underlying ChoiceTimeoutManager.arm is
* itself defensive about replacing an existing entry, but
* we shortcut here so we don't churn the timer's internal
* setTimeout handle on every `broadcastTopChoiceIfNew` call).
*
* The expiry callback simulates a synthetic submit-choice with
* `firstDefaultForKind(top.kind)`: it pops the top frame and
* clears the broadcast bookkeeping so a subsequent inner choice
* can be re-broadcast cleanly. T46's resume mechanism will replace
* the bare pop here with the real "thread the value back into the
* engine resume context" call — same parity as the human submit
* path in `handleSubmitChoice`.
*/
function armChoiceTimeoutFor(
roomCode: string,
session: GameSession,
top: PendingChoice,
): void {
const policy = getChoiceTimeoutPolicy(session.getEngine());
if (policy.mode !== "timeout-with-default") return;
if (choiceTimeoutManager.isArmed(roomCode, top.choiceId)) return;
const durationMs = policy.seconds * 1000;
const choiceId = top.choiceId;
const kind = top.kind;
choiceTimeoutManager.arm(roomCode, choiceId, durationMs, () => {
// Defensive: the session/room may have been torn down between
// arming and firing. A late-arriving timer must never crash the
// server or operate on stale state.
const liveSession = sessionRegistry.get(roomCode);
if (liveSession === undefined) return;
const stillTop = peekPendingChoice(liveSession.getEngine());
// Race guard: a real submit-choice that landed milliseconds
// before expiry has already popped the frame. The
// ChoiceTimeoutManager.cancel call in handleSubmitChoice
// should suppress this callback in that race, but the macrotask
// ordering between `clearTimeout` and an already-queued timeout
// callback is implementation-defined — re-checking the top
// before any side-effect is the only correct guard.
if (stillTop === undefined || stillTop.choiceId !== choiceId) return;
// Auto-resolve with the first option for this kind. The value
// bypasses the wire-side `isValidChoiceValue` check (it never
// touches the wire); the engine resume path (T46) will do its
// own kind-specific legality gate.
const defaultValue = firstDefaultForKind(kind);
void defaultValue; // T46: thread into engine resume context.
const popped = popPendingChoice(liveSession.getEngine());
if (popped !== undefined) {
forgetBroadcastedChoiceId(roomCode, popped.choiceId);
}
logger
.child({ roomCode })
.info(
{ choiceId, kind },
"T49: choice timeout expired; auto-resolved with first option",
);
});
}
/**
* Drop bookkeeping for a choiceId once it has been resolved. Called
* from the submit-choice handler after `popPendingChoice` succeeds.
*/
function forgetBroadcastedChoiceId(roomCode: string, choiceId: string): void {
const set = broadcastedChoiceIds.get(roomCode);
if (!set) return;
set.delete(choiceId);
if (set.size === 0) broadcastedChoiceIds.delete(roomCode);
}
/**
* Validate that `value` is a legal choice for the requested kind.
* Mirror the spec locked at T44:
* - "rps" → "rock" | "paper" | "scissors"
* - "piece" → number (entityId; non-negative integer)
* - "square" → number 0..63
* - "column" → number 0..7
* - "row" → number 0..7
* Returns `true` when `value` matches the kind's expected value-space.
*
* Engines / presets can still reject the value semantically downstream
* (e.g. picking an opponent's piece in a self-only choice), but the
* structural check here gates obvious garbage at the wire boundary so
* we don't feed arbitrary types into the resume mechanism (T46).
*/
function isValidChoiceValue(
kind: PendingChoice["kind"],
value: unknown,
): boolean {
switch (kind) {
case "rps":
return value === "rock" || value === "paper" || value === "scissors";
case "piece":
return (
typeof value === "number" &&
Number.isInteger(value) &&
value >= 0
);
case "square":
return (
typeof value === "number" &&
Number.isInteger(value) &&
value >= 0 &&
value <= 63
);
case "column":
case "row":
return (
typeof value === "number" &&
Number.isInteger(value) &&
value >= 0 &&
value <= 7
);
}
}
// ---------------------------------------------------------------------------
// Envelope builders
// ---------------------------------------------------------------------------
@ -217,6 +562,242 @@ function errorMessage(
// Message dispatch
// ---------------------------------------------------------------------------
/**
* Send a v2 top-level frame (e.g. `protocol-version-mismatch`).
* v2 frames are NOT wrapped in the v1 envelope — they travel as
* standalone JSON objects with a `kind` discriminator. See T43
* design notes in protocol.ts.
*/
function sendV2(ws: ServerWebSocket<ClientData>, msg: V2Message): void {
ws.send(JSON.stringify(msg));
}
/**
* T43: per-frame capability negotiation. Inspects the envelope's
* optional `protocolVersion` field and either pins it on
* `ws.data.protocolVersion` (first declaration) or rejects the
* frame on mismatch (unknown version) / re-negotiation attempt
* (a different version than the one already pinned).
*
* Returns `true` when the caller should keep processing the
* frame; `false` when the negotiation failed and the socket has
* already been sent a `protocol-version-mismatch` (and closed).
*/
function negotiateOrReject(
ws: ServerWebSocket<ClientData>,
declared: number | undefined,
): boolean {
// Treat `undefined` as "no declaration on THIS frame" — that is a
// no-op (the pinned value, if any, stands; the default of v1 if
// not yet pinned). Avoids forcing every subsequent frame to repeat
// the declaration.
if (declared === undefined) {
if (ws.data.protocolVersion === undefined) ws.data.protocolVersion = 1;
return true;
}
const negotiated = negotiateVersion(declared);
if (negotiated === "mismatch") {
sendV2(ws, {
kind: "protocol-version-mismatch",
supported: [...SUPPORTED_PROTOCOL_VERSIONS],
});
ws.close();
return false;
}
// First-time pin OR matching re-declaration both succeed; only a
// contradicting subsequent declaration fails. We treat that as a
// mismatch to keep the invariant "one version per connection" —
// a misbehaving client that flips versions mid-stream is not a
// case we want to silently accept.
if (
ws.data.protocolVersion !== undefined &&
ws.data.protocolVersion !== negotiated
) {
sendV2(ws, {
kind: "protocol-version-mismatch",
supported: [...SUPPORTED_PROTOCOL_VERSIONS],
});
ws.close();
return false;
}
ws.data.protocolVersion = negotiated;
return true;
}
/**
* T43/T44: route a v2 top-level frame after `validateMessage` narrowed
* it. T43 covers schema + negotiation; T44 wires the actual semantics
* for `submit-choice` (peek-top-of-stack, validate choiceId match,
* validate value-vs-kind, pop). `request-choice` and
* `protocol-version-mismatch` are server-only — a client emitting one
* is misuse and surfaces as a non-fatal INVALID_MESSAGE.
*/
function handleV2Frame(
ws: ServerWebSocket<ClientData>,
frame: V2Message,
): void {
switch (frame.kind) {
case "submit-choice":
handleSubmitChoice(ws, frame);
return;
case "request-choice":
case "protocol-version-mismatch":
// Server-only frames; clients have no business sending these.
sendTo(
ws,
errorMessage(
"INVALID_MESSAGE",
`server-only frame "${frame.kind}" received from client`,
false,
),
);
return;
}
}
/**
* T44 — validate and apply an inbound `submit-choice` frame.
*
* Validation pipeline (rejections all surface as non-fatal
* INVALID_MESSAGE — wire-shape was already proven by `V2MessageSchema`,
* so failures here are semantic):
*
* 1. The submitter must be authenticated into a room (`BAD_TOKEN`
* otherwise — same wire code as game.move).
* 2. The room must have a live session.
* 3. The PendingChoices stack must be non-empty.
* 4. `frame.choiceId` MUST equal the top frame's `choiceId`.
* Out-of-order submissions (resolving an outer frame before its
* inner) are rejected verbatim — LIFO is non-negotiable.
* 5. The submitting player's color must be authorised by the top
* frame's `forPlayer` ("both" admits either side; a single-color
* value rejects the opposite side as `BAD_TOKEN`).
* 6. `frame.value` must structurally match the top frame's `kind`
* (see `isValidChoiceValue`). Failures emit a `protocol.invalid-
* choice-value` message.
*
* On full success: the frame is popped (`popPendingChoice`),
* bookkeeping for the broadcasted-set is cleared, and the popped
* value is currently DROPPED — T46 will replace this with the real
* resume mechanism (`submitChoiceAndResume`). The pop itself is
* preserved so the LIFO contract holds even before T46 lands.
*/
function handleSubmitChoice(
ws: ServerWebSocket<ClientData>,
frame: SubmitChoice,
): void {
const { roomCode, token } = ws.data;
if (roomCode === undefined || token === undefined) {
sendTo(
ws,
errorMessage("BAD_TOKEN", "not authenticated into a room", false),
);
return;
}
const player = roomRegistry.getPlayerByToken(roomCode, token);
if (!player) {
sendTo(ws, errorMessage("BAD_TOKEN", "unknown token for room", false));
return;
}
const session = sessionRegistry.get(roomCode);
if (!session) {
sendTo(
ws,
errorMessage(
"INVALID_MESSAGE",
"internal error: missing game session",
true,
),
);
ws.close();
return;
}
const top = peekPendingChoice(session.getEngine());
if (top === undefined) {
sendTo(
ws,
errorMessage(
"INVALID_MESSAGE",
"submit-choice received but no pending choice on the stack",
false,
),
);
return;
}
// LIFO guard. The client must resolve the innermost (top) frame —
// submitting a different choiceId is either a stale retry from a
// prior frame or a misordered nested-choice resolution.
if (frame.choiceId !== top.choiceId) {
sendTo(
ws,
errorMessage(
"INVALID_MESSAGE",
`submit-choice choiceId mismatch: expected top "${top.choiceId}", got "${frame.choiceId}"`,
false,
),
);
return;
}
// forPlayer gate. "both" admits either color; otherwise the
// submitter's color must match. The opposite-color path is BAD_TOKEN
// because the misuse is "this player has no authority to resolve
// this prompt" — same family as a turn-order violation.
if (top.forPlayer !== "both" && top.forPlayer !== player.color) {
sendTo(
ws,
errorMessage(
"BAD_TOKEN",
`submit-choice not authorised: prompt is for ${top.forPlayer}, submitter is ${player.color}`,
false,
),
);
return;
}
// Structural value check against the prompt's kind. Failures land
// under `protocol.invalid-choice-value` per the T44 spec.
if (!isValidChoiceValue(top.kind, frame.value)) {
sendTo(
ws,
errorMessage(
"INVALID_MESSAGE",
`protocol.invalid-choice-value: value ${JSON.stringify(frame.value)} is not a legal "${top.kind}" choice`,
false,
),
);
return;
}
// T49 — cancel the auto-default timer first so an in-flight macrotask
// for THIS choiceId can't race with the human submission. The
// expiry callback's race-guard would catch it anyway, but cancelling
// up front keeps the timer registry size bounded by the actual
// number of unresolved choices and avoids an audible "auto-resolve
// fired after submit" log line in the latency window.
choiceTimeoutManager.cancel(roomCode, frame.choiceId);
// Pop the top frame and drop the bookkeeping entry. T46 will replace
// the bare pop with the real resume mechanism (which threads the
// value back into the engine's resume context); for now we satisfy
// the LIFO contract and clear our broadcast tracking so the next
// pending choice on this socket can be re-broadcast cleanly.
const popped = popPendingChoice(session.getEngine());
if (popped !== undefined) {
forgetBroadcastedChoiceId(roomCode, popped.choiceId);
}
logger
.child({ clientId: ws.data.clientId, roomCode })
.info(
{ choiceId: frame.choiceId, kind: top.kind },
"submit-choice (pre-T46: value accepted, resume not yet wired)",
);
}
/**
* Entry point for every inbound WS frame. Order of checks mirrors
* PROTOCOL.md §Error Handling: framing → size → parse → dispatch.
@ -230,7 +811,7 @@ export function handleMessage(
incMessages();
const str = typeof raw === "string" ? raw : raw.toString("utf8");
const result = validateMessageString(str);
const result = validateAnyMessageString(str);
if (!result.ok) {
// VERSION_MISMATCH is fatal per PROTOCOL.md; other parse failures are
// still fatal in v1 because we have no way to resync on a malformed
@ -244,7 +825,24 @@ export function handleMessage(
return;
}
const msg = result.data;
// T43: route v2 top-level frames separately. Inbound v2 frames are
// currently `submit-choice` (T44 wires the handler) — for now we
// surface a non-fatal INVALID_MESSAGE since no choice flow has been
// implemented yet. `protocol-version-mismatch` from a client is
// illegal (server-only); fail closed.
if (result.data.wire === "v2") {
handleV2Frame(ws, result.data.message);
return;
}
const msg = result.data.message;
// T43: every v1 frame can carry an optional `protocolVersion`
// envelope field (clients announce their capability on the first
// frame, typically `room.create` / `room.join`). Negotiate before
// dispatch so subsequent v2-only outbound traffic (e.g.
// request-choice) consults the pinned value.
if (!negotiateOrReject(ws, msg.protocolVersion)) return;
// Server-originated messages arriving from a client are protocol errors
// — we never expect to see them inbound. The union includes them so the
// single Schema can round-trip; here we gate them out.
@ -403,7 +1001,18 @@ function handleRoomCreate(
profile,
payload.preferredColor,
);
sessionRegistry.create(code, rulesetIds, resolvedLayout, profile);
// T50 — thread the wire-supplied choice-timeout policy through to
// the engine. `payload.choiceTimeout` is optional on the wire; when
// omitted the engine falls back to its DEFAULT_CHOICE_TIMEOUT_POLICY
// so old clients (and clients that explicitly omit the field)
// produce identical engine state to clients that supply the default.
sessionRegistry.create(
code,
rulesetIds,
resolvedLayout,
profile,
payload.choiceTimeout,
);
ws.data.roomCode = code;
ws.data.token = token;
setActiveRooms(roomRegistry.getRoomCount());
@ -457,7 +1066,19 @@ function handleRoomJoin(
payload.code,
envelopeToken,
);
if (existing && reconnectManager.isPending(envelopeToken)) {
// T49 — a `no-timeout` choice-disconnect parks the room in
// `pausedByChoiceDisconnect` WITHOUT starting a grace window
// (`isPending` would be false). Treat the paused-and-mine case
// as an alternative reconnect signal so the leaver can still
// resume their pending choice.
const pausedRoom = roomRegistry.getRoom(payload.code);
const isPausedReconnect =
existing !== undefined &&
pausedRoom?.pausedByChoiceDisconnect?.byToken === envelopeToken;
if (
existing &&
(reconnectManager.isPending(envelopeToken) || isPausedReconnect)
) {
handleReconnect(ws, payload.code, existing.token, existing.color);
return;
}
@ -550,24 +1171,37 @@ function handleReconnect(
token: string,
color: "white" | "black",
): void {
// cancelGrace MUST return deltas (isPending was true above), but we
// defend against a race where the timer fires between the isPending
// check and here. If the window expired we fall back to treating this
// as a failed reconnect — the caller already emitted game.end.
const missed = reconnectManager.cancelGrace(token);
if (missed === undefined) {
sendTo(
ws,
errorMessage(
"ROOM_NOT_FOUND",
`reconnect grace expired for room ${code}`,
false,
),
);
return;
// T49 — a paused-by-choice-disconnect reconnect path skips the
// ReconnectManager entirely (we never armed a grace timer for it).
// Detect that case BEFORE the cancelGrace call so an undefined
// return value isn't misread as an expired grace. `missed` is the
// empty array in the paused case — there were no deltas to buffer
// because the game was paused, not running.
const room = roomRegistry.getRoom(code);
const isPausedReconnect =
room?.pausedByChoiceDisconnect?.byToken === token;
let missed: ReturnType<typeof reconnectManager.cancelGrace> | undefined;
if (isPausedReconnect) {
missed = [];
} else {
// cancelGrace MUST return deltas (isPending was true above), but we
// defend against a race where the timer fires between the isPending
// check and here. If the window expired we fall back to treating this
// as a failed reconnect — the caller already emitted game.end.
missed = reconnectManager.cancelGrace(token);
if (missed === undefined) {
sendTo(
ws,
errorMessage(
"ROOM_NOT_FOUND",
`reconnect grace expired for room ${code}`,
false,
),
);
return;
}
}
const session = sessionRegistry.get(code);
const room = roomRegistry.getRoom(code);
if (!session || !room) {
sendTo(
ws,
@ -580,6 +1214,19 @@ function handleReconnect(
return;
}
// T49 — clear the paused flag so subsequent move/action handlers
// (downstream of T49 scope) un-gate. The pending-choice frame on
// the engine stack stays put and is re-broadcast below via the
// request-choice idempotence helper.
if (room.pausedByChoiceDisconnect !== undefined) {
delete room.pausedByChoiceDisconnect;
// Drop the broadcast bookkeeping for this room's pending
// choices so the prompt is RE-broadcast to the returning
// client (the original frame went out before they
// disconnected and was never delivered to a re-bound socket).
broadcastedChoiceIds.delete(code);
}
roomRegistry.markConnected(code, token);
ws.data.roomCode = code;
ws.data.token = token;
@ -633,6 +1280,15 @@ function handleReconnect(
for (const delta of missed) {
sendTo(ws, envelope("game.delta", delta.payload));
}
// T49 — if the engine still has a pending choice (typical of the
// paused-reconnect path, but also benign if a non-pause reconnect
// happens to land mid-choice), re-broadcast the top frame to the
// re-bound socket. The bookkeeping was cleared above so the
// idempotence guard inside `broadcastTopChoiceIfNew` permits the
// re-emit. No timer is re-armed for `no-timeout` policies; the
// existing `armChoiceTimeoutFor` no-ops in that mode.
broadcastTopChoiceIfNew(code, session);
}
function handleRoomLeave(ws: ServerWebSocket<ClientData>): void {
@ -765,6 +1421,12 @@ function handleGameMove(
// this — the pending slot is shared across the room.
applyPendingProfileIfAny(roomCode, session);
// T44 — if the move's tick pushed a request-choice onto the
// PendingChoices stack (request-choice primitive fired during a
// trigger arm), surface it to the appropriate client(s) NOW. v1
// clients are skipped; the broadcast is idempotent across re-entry.
broadcastTopChoiceIfNew(roomCode, session);
// If any preset durations expired during this move's tick, push the
// new set so clients stop rendering those rules. We skip the broadcast
// when the set is byte-identical to pre-move — the common case —
@ -957,6 +1619,11 @@ function handleGameAction(
// the existing reconciliation path for free.
broadcastGameStateSnapshot(roomCode, session);
// T44 — same hook as handleGameMove: if performAction's pipeline
// pushed a request-choice frame, surface it to the appropriate
// client(s) immediately after the snapshot.
broadcastTopChoiceIfNew(roomCode, session);
// Terminal-state guard — mirrors handleGameMove. If the action
// happened to end the game (rare in v1 but possible once future
// presets ship), emit game.end so the UI can transition out of

View file

@ -0,0 +1,569 @@
// T49 — choice-timeout enforcement + disconnect handler tests.
//
// Two layers of coverage live here:
//
// 1. *Pure-helper* unit tests against the policy + default tables in
// `choice-timeout.ts`. These need no engine and no timers — they
// pin the locked T49 contract (default value per kind, disconnect
// action mapping) so a future drift produces an immediate test
// failure rather than a silent semantic regression.
// 2. *Manager* tests against `ChoiceTimeoutManager` driven by
// vitest fake timers. Same structural pattern as
// `reconnect.test.ts` — the timer registry is storage-only so the
// tests stay free of WS / engine wiring.
// 3. *Wired* tests against `broadcast.ts` integration via the
// mock-WS pattern from `ws.request-choice.test.ts`. These cover
// the end-to-end paths the task brief calls out:
// - timer fires → top frame is auto-popped
// - real submit-choice cancels the pending timer
// - disconnect mid-choice under `timeout-with-default` →
// forfeit (game.end to opponent + session torn down)
// - disconnect mid-choice under `no-timeout` →
// paused state (no game.end, no grace timer, room intact)
//
// The task spec mandates fake timers (MUST NOT use Date.now). All
// timer-sensitive assertions advance the clock explicitly via
// `vi.advanceTimersByTime`.
import type { ServerWebSocket } from "bun";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import {
pushPendingChoice,
type ChoiceTimeoutPolicyValue,
type PendingChoice,
} from "@paratype/chess";
import {
broadcastTopChoiceIfNew,
handleMessage,
registerConnection,
roomRegistry,
sessionRegistry,
unregisterConnection,
type ClientData,
} from "./broadcast.js";
import {
ChoiceTimeoutManager,
choiceTimeoutManager,
decideDisconnectAction,
firstDefaultForKind,
getChoiceTimeoutPolicy,
hasPendingChoice,
} from "./choice-timeout.js";
import { PROTOCOL_VERSION, type ClientMessage } from "./protocol.js";
// ---------------------------------------------------------------------------
// Mock ServerWebSocket — same shape as ws.request-choice.test.ts
// ---------------------------------------------------------------------------
interface MockWs extends ServerWebSocket<ClientData> {
readonly sent: unknown[];
readonly closed: boolean;
}
function makeMockWs(clientId: string): MockWs {
const sent: unknown[] = [];
const closedFlag = { value: false };
const ws = {
data: { clientId } as ClientData,
sent,
get closed(): boolean {
return closedFlag.value;
},
send(msg: string | Buffer): number {
const str = typeof msg === "string" ? msg : msg.toString("utf8");
sent.push(JSON.parse(str));
return str.length;
},
close(): void {
closedFlag.value = true;
},
} as unknown as MockWs;
return ws;
}
function findMsgOfType(ws: MockWs, type: string): unknown | undefined {
return ws.sent.find(
(m) =>
typeof m === "object" &&
m !== null &&
((m as { type?: unknown }).type === type ||
(m as { kind?: unknown }).kind === type),
);
}
function nextMsgOfType(
ws: MockWs,
type: string,
): { payload?: Record<string, unknown>; [k: string]: unknown } {
const idx = ws.sent.findIndex(
(m) =>
typeof m === "object" &&
m !== null &&
((m as { type?: unknown }).type === type ||
(m as { kind?: unknown }).kind === type),
);
if (idx < 0) {
const tags = ws.sent.map(
(m) =>
(m as { type?: string; kind?: string }).type ??
(m as { kind?: string }).kind,
);
throw new Error(
`no message of type/kind "${type}" in inbox (got ${JSON.stringify(tags)})`,
);
}
const msg = ws.sent[idx] as { payload?: Record<string, unknown> };
ws.sent.splice(idx, 1);
return msg;
}
function sendClient(
ws: MockWs,
type: ClientMessage["type"],
payload: unknown,
opts: { seq?: number; protocolVersion?: number; token?: string } = {},
): void {
const env: Record<string, unknown> = {
v: PROTOCOL_VERSION,
seq: opts.seq ?? 1,
ts: Date.now(),
type,
payload,
};
if (opts.protocolVersion !== undefined) {
env["protocolVersion"] = opts.protocolVersion;
}
if (opts.token !== undefined) {
env["token"] = opts.token;
}
handleMessage(ws, JSON.stringify(env));
}
function sendV2(ws: MockWs, frame: Record<string, unknown>): void {
handleMessage(ws, JSON.stringify(frame));
}
interface RoomCtx {
white: MockWs;
black: MockWs;
code: string;
whiteToken: string;
blackToken: string;
}
function setupRoom(opts?: {
choiceTimeout?: ChoiceTimeoutPolicyValue;
}): RoomCtx {
const white = makeMockWs(`white-${Math.random().toString(36).slice(2, 8)}`);
const black = makeMockWs(`black-${Math.random().toString(36).slice(2, 8)}`);
registerConnection(white);
registerConnection(black);
const createPayload: Record<string, unknown> = { rulesetIds: [] };
if (opts?.choiceTimeout !== undefined) {
createPayload["choiceTimeout"] = opts.choiceTimeout;
}
sendClient(white, "room.create", createPayload, { protocolVersion: 2 });
const created = nextMsgOfType(white, "room.created");
const code = created["payload"]!["code"] as string;
const whiteToken = created["payload"]!["token"] as string;
sendClient(
black,
"room.join",
{ code },
{ protocolVersion: 2 },
);
const joined = nextMsgOfType(black, "room.joined");
const blackToken = joined["payload"]!["token"] as string;
// Drain initial game.state frames.
nextMsgOfType(white, "game.state");
nextMsgOfType(black, "game.state");
return { white, black, code, whiteToken, blackToken };
}
function buildPendingChoice(overrides: Partial<PendingChoice>): PendingChoice {
return {
choiceId: "choice-1",
descriptorId: "test-descriptor",
triggerPath: [0],
primitiveIndex: 0,
bindings: new Map(),
kind: "rps",
prompt: "rock-paper-scissors?",
forPlayer: "both",
...overrides,
};
}
function teardownRoom(ctx: RoomCtx): void {
// Force-close both sockets without going through the disconnect
// handler so the singleton state from one test doesn't bleed into
// the next. We then explicitly tear down the room via a
// roomRegistry.leaveRoom call so room codes don't accumulate.
unregisterConnection(ctx.white);
unregisterConnection(ctx.black);
// Last-resort cleanup for any timer or paused flag that survived
// the disconnect path (e.g. forfeit branches that already cleaned up).
choiceTimeoutManager.cancelAll(ctx.code);
}
// ---------------------------------------------------------------------------
// Pure-helper tests
// ---------------------------------------------------------------------------
describe("T49 — firstDefaultForKind (locked default table)", () => {
it("rps → 'rock'", () => {
expect(firstDefaultForKind("rps")).toBe("rock");
});
it("piece → -1 (sentinel; T46 resume layer interprets)", () => {
expect(firstDefaultForKind("piece")).toBe(-1);
});
it("square → 0 (a1)", () => {
expect(firstDefaultForKind("square")).toBe(0);
});
it("column → 0 (file a)", () => {
expect(firstDefaultForKind("column")).toBe(0);
});
it("row → 0 (rank 1)", () => {
expect(firstDefaultForKind("row")).toBe(0);
});
});
describe("T49 — decideDisconnectAction (locked policy table)", () => {
it("no pending choice → 'none' regardless of policy", () => {
expect(
decideDisconnectAction(
{ mode: "timeout-with-default", seconds: 60 },
false,
),
).toBe("none");
expect(
decideDisconnectAction({ mode: "no-timeout" }, false),
).toBe("none");
});
it("timeout-with-default + pending → 'forfeit'", () => {
expect(
decideDisconnectAction(
{ mode: "timeout-with-default", seconds: 30 },
true,
),
).toBe("forfeit");
});
it("no-timeout + pending → 'pause'", () => {
expect(
decideDisconnectAction({ mode: "no-timeout" }, true),
).toBe("pause");
});
});
// ---------------------------------------------------------------------------
// ChoiceTimeoutManager (in-isolation; mirrors reconnect.test.ts)
// ---------------------------------------------------------------------------
describe("ChoiceTimeoutManager", () => {
let mgr: ChoiceTimeoutManager;
beforeEach(() => {
mgr = new ChoiceTimeoutManager();
vi.useFakeTimers();
});
afterEach(() => {
mgr.clearAll();
vi.useRealTimers();
});
it("arm + isArmed reflects the registry state", () => {
mgr.arm("ROOM-A", "choice-1", 1_000, () => {});
expect(mgr.isArmed("ROOM-A", "choice-1")).toBe(true);
expect(mgr.isArmed("ROOM-A", "choice-2")).toBe(false);
expect(mgr.isArmed("ROOM-B", "choice-1")).toBe(false);
});
it("timer fires onExpire after duration and removes the entry", () => {
const onExpire = vi.fn();
mgr.arm("ROOM-A", "choice-1", 100, onExpire);
expect(mgr.isArmed("ROOM-A", "choice-1")).toBe(true);
vi.advanceTimersByTime(150);
expect(onExpire).toHaveBeenCalledTimes(1);
expect(mgr.isArmed("ROOM-A", "choice-1")).toBe(false);
});
it("cancel within the window suppresses the callback", () => {
const onExpire = vi.fn();
mgr.arm("ROOM-A", "choice-1", 1_000, onExpire);
expect(mgr.cancel("ROOM-A", "choice-1")).toBe(true);
vi.advanceTimersByTime(2_000);
expect(onExpire).not.toHaveBeenCalled();
expect(mgr.cancel("ROOM-A", "choice-1")).toBe(false); // idempotent
});
it("cancelAll(roomCode) drops every timer in that room only", () => {
const a = vi.fn();
const b = vi.fn();
const c = vi.fn();
mgr.arm("ROOM-A", "c-1", 100, a);
mgr.arm("ROOM-A", "c-2", 100, b);
mgr.arm("ROOM-B", "c-3", 100, c);
mgr.cancelAll("ROOM-A");
vi.advanceTimersByTime(200);
expect(a).not.toHaveBeenCalled();
expect(b).not.toHaveBeenCalled();
expect(c).toHaveBeenCalledTimes(1);
});
it("re-arming the same key replaces the prior timer (no leak)", () => {
const first = vi.fn();
const second = vi.fn();
mgr.arm("ROOM-A", "c-1", 100, first);
mgr.arm("ROOM-A", "c-1", 100, second);
vi.advanceTimersByTime(200);
expect(first).not.toHaveBeenCalled();
expect(second).toHaveBeenCalledTimes(1);
expect(mgr.size()).toBe(0);
});
});
// ---------------------------------------------------------------------------
// End-to-end (broadcast.ts) integration
// ---------------------------------------------------------------------------
describe("T49 — choice timeout enforcement (e2e via broadcast)", () => {
beforeEach(() => {
vi.useFakeTimers();
});
afterEach(() => {
vi.useRealTimers();
});
it("auto-resolves the top frame after the policy deadline elapses (timeout-with-default)", () => {
const ctx = setupRoom({
choiceTimeout: { mode: "timeout-with-default", seconds: 5 },
});
const session = sessionRegistry.get(ctx.code)!;
// Confirm the engine seeded the policy fact correctly.
expect(getChoiceTimeoutPolicy(session.getEngine())).toEqual({
mode: "timeout-with-default",
seconds: 5,
});
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "rps-timeout",
kind: "rps",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(ctx.code, session);
// Both clients receive the prompt; the timer is armed.
nextMsgOfType(ctx.white, "request-choice");
nextMsgOfType(ctx.black, "request-choice");
expect(choiceTimeoutManager.isArmed(ctx.code, "rps-timeout")).toBe(true);
expect(hasPendingChoice(session.getEngine())).toBe(true);
// Just before the deadline: nothing has happened.
vi.advanceTimersByTime(4_999);
expect(hasPendingChoice(session.getEngine())).toBe(true);
// Cross the deadline: the auto-resolve fires, the frame is popped,
// and the registry forgets the entry.
vi.advanceTimersByTime(2);
expect(hasPendingChoice(session.getEngine())).toBe(false);
expect(choiceTimeoutManager.isArmed(ctx.code, "rps-timeout")).toBe(false);
teardownRoom(ctx);
});
it("does NOT arm a timer under no-timeout policy; the prompt waits indefinitely", () => {
const ctx = setupRoom({ choiceTimeout: { mode: "no-timeout" } });
const session = sessionRegistry.get(ctx.code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "no-timeout-prompt",
kind: "rps",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(ctx.code, session);
nextMsgOfType(ctx.white, "request-choice");
nextMsgOfType(ctx.black, "request-choice");
expect(
choiceTimeoutManager.isArmed(ctx.code, "no-timeout-prompt"),
).toBe(false);
// Advance an hour: the engine still holds the pending frame.
vi.advanceTimersByTime(60 * 60 * 1_000);
expect(hasPendingChoice(session.getEngine())).toBe(true);
teardownRoom(ctx);
});
it("a real submit-choice cancels the pending auto-default timer", () => {
const ctx = setupRoom({
choiceTimeout: { mode: "timeout-with-default", seconds: 5 },
});
const session = sessionRegistry.get(ctx.code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "submit-cancels",
kind: "rps",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(ctx.code, session);
nextMsgOfType(ctx.white, "request-choice");
nextMsgOfType(ctx.black, "request-choice");
expect(
choiceTimeoutManager.isArmed(ctx.code, "submit-cancels"),
).toBe(true);
sendV2(ctx.white, {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "submit-cancels",
value: "paper",
});
expect(findMsgOfType(ctx.white, "error")).toBeUndefined();
expect(
choiceTimeoutManager.isArmed(ctx.code, "submit-cancels"),
).toBe(false);
// Advance well past the original deadline — no second pop, no
// crash; the auto-resolve callback was cancelled.
vi.advanceTimersByTime(60_000);
expect(hasPendingChoice(session.getEngine())).toBe(false);
teardownRoom(ctx);
});
});
describe("T49 — disconnect handler", () => {
beforeEach(() => {
vi.useFakeTimers();
});
afterEach(() => {
vi.useRealTimers();
});
it("timeout-with-default + disconnect mid-choice → forfeit (game.end + room torn down)", () => {
const ctx = setupRoom({
choiceTimeout: { mode: "timeout-with-default", seconds: 30 },
});
const session = sessionRegistry.get(ctx.code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "mid-flight",
kind: "rps",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(ctx.code, session);
nextMsgOfType(ctx.white, "request-choice");
nextMsgOfType(ctx.black, "request-choice");
// White disconnects with an unresolved choice on the stack.
unregisterConnection(ctx.white);
// Black receives game.end (forfeit; opponent wins) immediately —
// the choice-timeout-disconnect path supersedes the standard
// 60s grace window, which would otherwise have suppressed the
// game.end until expiry.
const end = nextMsgOfType(ctx.black, "game.end");
expect(end["payload"]!["winner"]).toBe("black");
expect(end["payload"]!["reason"]).toBe("player_left");
// White's slot is gone (leaveRoom called inline); black's slot
// survives so the room itself stays alive until they leave.
expect(roomRegistry.getPlayerByToken(ctx.code, ctx.whiteToken)).toBeUndefined();
// No pending auto-default timer survived the forfeit.
expect(
choiceTimeoutManager.isArmed(ctx.code, "mid-flight"),
).toBe(false);
// The room is NOT in the paused-by-choice state (forfeit chose
// a hard tear-down for the leaver, not a pause).
const room = roomRegistry.getRoom(ctx.code);
expect(room?.pausedByChoiceDisconnect).toBeUndefined();
unregisterConnection(ctx.black);
});
it("no-timeout + disconnect mid-choice → paused; no game.end, no grace timer", () => {
const ctx = setupRoom({ choiceTimeout: { mode: "no-timeout" } });
const session = sessionRegistry.get(ctx.code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "paused-prompt",
kind: "rps",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(ctx.code, session);
nextMsgOfType(ctx.white, "request-choice");
nextMsgOfType(ctx.black, "request-choice");
unregisterConnection(ctx.white);
// Black received NO game.end — the game is paused, not over.
expect(findMsgOfType(ctx.black, "game.end")).toBeUndefined();
// The room survives, with the paused flag set.
const room = roomRegistry.getRoom(ctx.code);
expect(room).toBeDefined();
expect(room!.pausedByChoiceDisconnect).toMatchObject({
byToken: ctx.whiteToken,
});
// The session is still alive and still holds the pending frame.
const aliveSession = sessionRegistry.get(ctx.code);
expect(aliveSession).toBeDefined();
expect(hasPendingChoice(aliveSession!.getEngine())).toBe(true);
// Advance an hour; nothing fires, the room stays paused.
vi.advanceTimersByTime(60 * 60 * 1_000);
expect(roomRegistry.getRoom(ctx.code)).toBeDefined();
expect(hasPendingChoice(aliveSession!.getEngine())).toBe(true);
unregisterConnection(ctx.black);
// Force-clear any leftover state from the second disconnect.
if (roomRegistry.getRoom(ctx.code) !== undefined) {
roomRegistry.leaveRoom(ctx.code, ctx.blackToken);
sessionRegistry.delete(ctx.code);
}
choiceTimeoutManager.cancelAll(ctx.code);
});
it("no-pending-choice disconnect falls through to the standard reconnect-grace path (unchanged)", () => {
const ctx = setupRoom({
choiceTimeout: { mode: "timeout-with-default", seconds: 30 },
});
// No pushPendingChoice — engine stack is empty.
unregisterConnection(ctx.white);
// Standard grace path: the room survives the disconnect (the
// 60s grace timer is armed), the slot is marked disconnected,
// and no game.end fires immediately.
const room = roomRegistry.getRoom(ctx.code);
expect(room).toBeDefined();
expect(room!.pausedByChoiceDisconnect).toBeUndefined();
expect(findMsgOfType(ctx.black, "game.end")).toBeUndefined();
teardownRoom(ctx);
});
});

View file

@ -0,0 +1,302 @@
// T49 — choice-timeout enforcement + disconnect handler.
//
// This module owns the *server-side* timing layer that complements the
// engine-owned `ChoiceTimeoutPolicy` fact (T50, on `GAME_ENTITY`). The
// engine itself never schedules timers — it only carries the policy
// value so a single source of truth is bound to the game session
// (chess/src/schema.ts §"T50 — per-game choice-timeout policy"). The
// WS layer ((un-)registerConnection + broadcastTopChoiceIfNew +
// handleSubmitChoice in `broadcast.ts`) consults this module to:
//
// 1. ARM a timer when a `request-choice` is broadcast and the policy
// is `timeout-with-default`. On expiry the top frame is auto-
// resolved with `firstDefaultForKind` and the resume mechanism
// proceeds as if the player had submitted that value (today,
// pre-T46, the popped frame is dropped — same behaviour as the
// real submit-choice handler).
// 2. CANCEL the timer when a real `submit-choice` arrives so the
// auto-default never races a successful submission.
// 3. DECIDE on disconnect: with `timeout-with-default` policy a
// mid-choice disconnect forfeits the leaver; with `no-timeout`
// the game enters a paused state until reconnect.
//
// Locked by `decisions.md` §"Choice Timeout & Disconnect" (notepad
// `thressgame-coverage`); no policy decisions are made here — this
// module is pure mechanism for the policy already chosen at T50 and
// stored on the engine.
//
// ## Why a separate module
//
// Keeping the timer registry out of `broadcast.ts` is deliberate:
// - it stays storage-only (no WS imports), mirroring the
// `ReconnectManager` split — same shape, same testability;
// - the disconnect-policy decision is a pure function of (policy,
// hasPendingChoice) and easy to unit-test without a mock socket;
// - tests can drive the manager with `vi.useFakeTimers()` without
// spinning up the broadcast layer (the task spec mandates fake
// timers — `MUST NOT DO: Use Date.now in tests`).
import {
GAME_ENTITY,
peekPendingChoice,
type ChessEngine,
type ChoiceTimeoutPolicyValue,
type PendingChoice,
} from "@paratype/chess";
// ---------------------------------------------------------------------------
// Per-kind first-option defaults
// ---------------------------------------------------------------------------
/**
* Canonical "first option" value used to auto-resolve a pending
* choice when its timer expires under `timeout-with-default`.
*
* Locked by the T49 task brief ("first option" defaults table):
* - rps → "rock" — first rock-paper-scissors token.
* - piece → -1 — sentinel id; T46's resume layer
* treats it as "no piece chosen"
* (or "first ally piece" depending
* on the descriptor's contract).
* The wire-side structural check
* (`isValidChoiceValue` in
* broadcast.ts) rejects -1, which
* is intentional: a *real* client
* cannot submit -1 — the server
* only constructs it for the auto-
* resolve path, which bypasses the
* wire validator.
* - square → 0 — a1 in 0-indexed square space.
* - column → 0 — file a.
* - row → 0 — rank 1.
*
* The function is exhaustive over `PendingChoice["kind"]`; adding a
* new kind to the schema fails the typecheck here.
*/
export function firstDefaultForKind(
kind: PendingChoice["kind"],
): "rock" | number {
switch (kind) {
case "rps":
return "rock";
case "piece":
return -1;
case "square":
case "column":
case "row":
return 0;
}
}
// ---------------------------------------------------------------------------
// Policy reader
// ---------------------------------------------------------------------------
/**
* Read the active `ChoiceTimeoutPolicy` from `GAME_ENTITY`. The engine
* always seeds this fact at construction time (T50 default
* `{ mode: "timeout-with-default", seconds: 60 }`) so the lookup
* is total — a missing fact would indicate engine corruption and
* the function returns the wire-default to fail open.
*
* Returning the default rather than throwing means a hypothetical
* test that constructs an engine without seeding the policy still
* gets predictable behaviour; production paths always have the fact.
*/
export function getChoiceTimeoutPolicy(
engine: ChessEngine,
): ChoiceTimeoutPolicyValue {
const fact = engine.session.get(GAME_ENTITY, "ChoiceTimeoutPolicy") as
| ChoiceTimeoutPolicyValue
| undefined;
return fact ?? { mode: "timeout-with-default", seconds: 60 };
}
/**
* Convenience wrapper: is there an unresolved choice frame on the
* engine's stack? Used by the disconnect handler to decide whether
* the choice-timeout policy is even relevant — a disconnect with no
* pending choice falls through to the standard reconnect-grace
* pipeline unchanged.
*/
export function hasPendingChoice(engine: ChessEngine): boolean {
return peekPendingChoice(engine) !== undefined;
}
// ---------------------------------------------------------------------------
// Disconnect-policy decision
// ---------------------------------------------------------------------------
/**
* What should the WS layer do when a player disconnects with a
* pending choice on the stack?
*
* - `"forfeit"` — `timeout-with-default` mode. The disconnected
* player loses immediately; the opponent receives `game.end`
* and the session is torn down without waiting for the standard
* reconnect-grace window (the choice-timeout decision is more
* specific than the generic disconnect path).
* - `"pause"` — `no-timeout` mode. The game enters a paused
* state until the leaver reconnects; no `game.end` is broadcast,
* no auto-resolve fires, and the standard grace-window timer
* is suppressed (tearing the room down on a 60s grace expiry
* would defeat the "paused indefinitely" contract).
* - `"none"` — no pending choice. The disconnect is unrelated
* to the choice flow; the standard reconnect-grace path applies
* verbatim.
*
* Pure function of (policy, hasPending) — no I/O, no state. The
* `decisions.md` table (notepad `thressgame-coverage` §"Choice
* Timeout & Disconnect") locks the policy → action mapping; this
* implementation must mirror that table exactly. Any drift requires
* a notepad amendment first.
*/
export type DisconnectAction = "forfeit" | "pause" | "none";
export function decideDisconnectAction(
policy: ChoiceTimeoutPolicyValue,
pending: boolean,
): DisconnectAction {
if (!pending) return "none";
switch (policy.mode) {
case "timeout-with-default":
return "forfeit";
case "no-timeout":
return "pause";
}
}
// ---------------------------------------------------------------------------
// ChoiceTimeoutManager — per-(roomCode, choiceId) timer registry
// ---------------------------------------------------------------------------
/**
* Composite key for a timer entry. (roomCode, choiceId) is the
* minimal disambiguator: a single room may have multiple stacked
* choices in flight (LIFO depth ≤ MAX_CHOICE_DEPTH = 8) and the
* top-of-stack id changes as inner frames resolve.
*
* String concatenation is fine — choiceIds are server-minted UUID-
* like tokens and roomCodes are 6-char alphanumeric; no separator
* collision is reachable.
*/
function timerKey(roomCode: string, choiceId: string): string {
return `${roomCode}\x00${choiceId}`;
}
interface TimerEntry {
handle: ReturnType<typeof setTimeout>;
/** Captured at arm time so cancellation can confirm we're cancelling
* the same logical entry the caller intends. */
roomCode: string;
choiceId: string;
}
/**
* Storage-only timer registry for the choice-timeout flow. Mirrors
* the structural shape of `ReconnectManager` — start/cancel/clearAll
* — so the broadcast layer's mental model is uniform across the two
* timer subsystems.
*
* A single process-global instance is exported as `choiceTimeoutManager`
* (alongside `roomRegistry` / `sessionRegistry` / `reconnectManager`
* in broadcast.ts). Tests construct fresh instances per `describe`
* block to avoid cross-test pollution.
*/
export class ChoiceTimeoutManager {
private readonly timers = new Map<string, TimerEntry>();
/**
* Arm a timer for `(roomCode, choiceId)`. On expiry, `onExpire`
* runs on the event loop and the entry is forgotten from the
* registry BEFORE the callback fires — same semantics as
* `ReconnectManager.startGrace` so an `onExpire` that calls back
* into `isArmed` sees `false`.
*
* Calling `arm` twice for the same key replaces the previous
* timer (defensive; the broadcast layer is supposed to call
* `arm` at most once per choiceId, but a re-broadcast under
* `broadcastTopChoiceIfNew`'s idempotence check shouldn't leak a
* timer on the off chance the bookkeeping diverges).
*/
arm(
roomCode: string,
choiceId: string,
durationMs: number,
onExpire: () => void,
): void {
const key = timerKey(roomCode, choiceId);
const existing = this.timers.get(key);
if (existing) clearTimeout(existing.handle);
const handle = setTimeout(() => {
// Forget the entry FIRST so onExpire's downstream calls
// (`isArmed`, `cancel`) observe a consistent post-fire state.
this.timers.delete(key);
onExpire();
}, durationMs);
// Same `unref` defensive call as ReconnectManager: a pending
// timer must not block process shutdown in tests that forget
// to clearAll() — Bun's setTimeout is Node-compatible.
if (typeof (handle as { unref?: () => void }).unref === "function") {
(handle as { unref: () => void }).unref();
}
this.timers.set(key, { handle, roomCode, choiceId });
}
/**
* Cancel the armed timer for `(roomCode, choiceId)`. Returns
* `true` if a timer was cancelled, `false` if no timer was
* pending (already fired or never armed). Idempotent — safe to
* call from the submit-choice handler unconditionally.
*/
cancel(roomCode: string, choiceId: string): boolean {
const key = timerKey(roomCode, choiceId);
const entry = this.timers.get(key);
if (!entry) return false;
clearTimeout(entry.handle);
this.timers.delete(key);
return true;
}
/**
* Cancel every timer scoped to `roomCode`. Used at room teardown
* (game.end / leaveRoom / forfeit on choice-disconnect) so a
* pending auto-resolve doesn't fire after the session is gone.
*/
cancelAll(roomCode: string): void {
for (const [key, entry] of this.timers.entries()) {
if (entry.roomCode === roomCode) {
clearTimeout(entry.handle);
this.timers.delete(key);
}
}
}
/** True iff an unfired timer is registered for the key. */
isArmed(roomCode: string, choiceId: string): boolean {
return this.timers.has(timerKey(roomCode, choiceId));
}
/**
* Diagnostics-only: clear every timer without firing callbacks.
* Tests use this in `afterEach` to keep the singleton clean.
*/
clearAll(): void {
for (const entry of this.timers.values()) {
clearTimeout(entry.handle);
}
this.timers.clear();
}
/** Number of armed timers (for tests / metrics). */
size(): number {
return this.timers.size;
}
}
/**
* Process-global manager shared by `broadcast.ts`. Singleton matches
* the rest of the server's wiring (roomRegistry, sessionRegistry,
* reconnectManager).
*/
export const choiceTimeoutManager = new ChoiceTimeoutManager();

View file

@ -18,6 +18,7 @@ import {
validateProfile,
type ActionResult,
type ActivationRequest,
type ChoiceTimeoutPolicyValue,
type ModifierProfile,
type ModifierValidationErrorCode,
type PlayerAction,
@ -139,17 +140,31 @@ export class GameSession {
* pass straight through to `EngineOptions.profile`, which seeds
* the modifier facts and auto-activates the
* `__modifier-profile-integration__` preset.
* @param choiceTimeout — optional per-game choice-timeout policy
* (T50). Threaded straight through to `EngineOptions.choiceTimeout`
* which seeds the `ChoiceTimeoutPolicy` fact on `GAME_ENTITY`.
* When omitted the engine seeds its own
* `DEFAULT_CHOICE_TIMEOUT_POLICY` so the server's wire-default
* and the engine's construction-default coincide. The protocol
* schema has already enforced `seconds >= 1` upstream; this
* layer trusts the value.
*/
constructor(
rulesetIds: readonly string[] = [],
layout?: StartingLayout,
profile?: ModifierProfile,
choiceTimeout?: ChoiceTimeoutPolicyValue,
) {
// Build the options bag once, including only the keys that were
// supplied — ChessEngine treats missing keys as "use default".
const opts: { layout?: StartingLayout; profile?: ModifierProfile } = {};
const opts: {
layout?: StartingLayout;
profile?: ModifierProfile;
choiceTimeout?: ChoiceTimeoutPolicyValue;
} = {};
if (layout !== undefined) opts.layout = layout;
if (profile !== undefined) opts.profile = profile;
if (choiceTimeout !== undefined) opts.choiceTimeout = choiceTimeout;
this.engine = Object.keys(opts).length > 0
? new ChessEngine(opts)
: new ChessEngine();
@ -289,6 +304,22 @@ export class GameSession {
return this.engine.getCurrentTurn();
}
/**
* T44 — escape hatch for the WS broadcast layer to introspect the
* underlying engine when wiring suspended-execution flows
* (request-choice / submit-choice). The broadcast layer needs to
* peek/pop the engine's `PendingChoices` stack via the chess-side
* pending-choices helpers; those helpers take a `ChessEngine`
* directly so we expose it here rather than mirroring every helper
* onto GameSession. Keep the surface narrow — production code
* should funnel state mutations through GameSession's typed methods
* (applyMove, performAction, reconcileProfile). This accessor is
* deliberately scoped to the choice-flow integration point.
*/
getEngine(): ChessEngine {
return this.engine;
}
/**
* Is the game terminally over? Returns the terminal descriptor or null.
* Cheap — just reads the sticky flag set in applyMove.
@ -432,20 +463,28 @@ export class GameSessionRegistry {
* `layout` is optional; when provided, the engine opens from it
* instead of the FIDE default. `profile` is optional; when provided
* the engine seeds modifier facts on piece entities and auto-
* activates the integration preset.
* activates the integration preset. `choiceTimeout` (T50) is the
* per-game choice-timeout policy; when omitted the engine falls
* back to `DEFAULT_CHOICE_TIMEOUT_POLICY`.
*/
create(
code: string,
rulesetIds?: readonly string[],
layout?: StartingLayout,
profile?: ModifierProfile,
choiceTimeout?: ChoiceTimeoutPolicyValue,
): GameSession {
if (this.sessions.has(code)) {
throw new Error(
`GameSessionRegistry: session already exists for code "${code}"`,
);
}
const session = new GameSession(rulesetIds ?? [], layout, profile);
const session = new GameSession(
rulesetIds ?? [],
layout,
profile,
choiceTimeout,
);
this.sessions.set(code, session);
return session;
}

View file

@ -2,12 +2,22 @@ import { describe, it, expect } from "vitest";
import {
validateMessage,
validateMessageString,
validateAnyMessage,
validateAnyMessageString,
negotiateVersion,
shouldSkipV2Broadcast,
PROTOCOL_VERSION,
SUPPORTED_PROTOCOL_VERSIONS,
LATEST_PROTOCOL_VERSION,
ClientMessageSchema,
ServerMessageSchema,
ModifierProfileSchema,
ModifierProfileUpdatePayloadSchema,
RoomCreatePayloadSchema,
RequestChoiceSchema,
SubmitChoiceSchema,
ProtocolVersionMismatchSchema,
V2MessageSchema,
MODIFIER_PROFILE_INVALID,
MODIFIER_PROFILE_NO_KING,
MODIFIER_PROFILE_INVULN_KING,
@ -15,6 +25,8 @@ import {
type AnyMessage,
type ClientMessage,
type ServerMessage,
type RequestChoice,
type SubmitChoice,
} from "./protocol.js";
// ---------------------------------------------------------------------------
@ -474,6 +486,85 @@ describe("RoomCreatePayloadSchema preferredColor", () => {
});
});
// ---------------------------------------------------------------------------
// T50 — choiceTimeout on room.create
// ---------------------------------------------------------------------------
describe("RoomCreatePayloadSchema choiceTimeout (T50)", () => {
it("omitted is valid (legacy clients land on engine default)", () => {
const r = RoomCreatePayloadSchema.safeParse({});
expect(r.success).toBe(true);
if (r.success) {
// Field must be absent rather than coerced — the engine, not the
// wire schema, is responsible for falling back to
// DEFAULT_CHOICE_TIMEOUT_POLICY when the option is missing.
expect(r.data.choiceTimeout).toBeUndefined();
}
});
it("accepts explicit timeout-with-default + positive seconds", () => {
const r = RoomCreatePayloadSchema.safeParse({
choiceTimeout: { mode: "timeout-with-default", seconds: 60 },
});
expect(r.success).toBe(true);
if (r.success) {
expect(r.data.choiceTimeout).toEqual({
mode: "timeout-with-default",
seconds: 60,
});
}
});
it("accepts no-timeout (no seconds field)", () => {
const r = RoomCreatePayloadSchema.safeParse({
choiceTimeout: { mode: "no-timeout" },
});
expect(r.success).toBe(true);
if (r.success) {
expect(r.data.choiceTimeout).toEqual({ mode: "no-timeout" });
}
});
it("rejects an unknown mode discriminator", () => {
const r = RoomCreatePayloadSchema.safeParse({
choiceTimeout: { mode: "fast-forward", seconds: 30 },
});
expect(r.success).toBe(false);
});
it("rejects timeout-with-default with non-positive seconds", () => {
// `seconds` must be >= 1 — non-positive timeouts would disable
// the very feature the policy enables. The wire-layer bound keeps
// the engine contract simple (engine trusts the value verbatim).
const r = RoomCreatePayloadSchema.safeParse({
choiceTimeout: { mode: "timeout-with-default", seconds: 0 },
});
expect(r.success).toBe(false);
});
it("rejects timeout-with-default with negative seconds", () => {
const r = RoomCreatePayloadSchema.safeParse({
choiceTimeout: { mode: "timeout-with-default", seconds: -10 },
});
expect(r.success).toBe(false);
});
it("rejects timeout-with-default with non-integer seconds", () => {
const r = RoomCreatePayloadSchema.safeParse({
choiceTimeout: { mode: "timeout-with-default", seconds: 1.5 },
});
expect(r.success).toBe(false);
});
it("rejects timeout-with-default missing seconds", () => {
const r = RoomCreatePayloadSchema.safeParse({
choiceTimeout: { mode: "timeout-with-default" },
});
expect(r.success).toBe(false);
});
it("composes with rulesetIds + preferredColor + choiceTimeout", () => {
const r = RoomCreatePayloadSchema.safeParse({
rulesetIds: ["piece-hp"],
preferredColor: "random",
choiceTimeout: { mode: "no-timeout" },
});
expect(r.success).toBe(true);
});
});
describe("validateMessageString", () => {
it("parses a valid JSON string frame", () => {
const msg: ClientMessage = {
@ -1027,3 +1118,418 @@ describe("game.action schema validation", () => {
expect(mod.KNOWN_MESSAGE_TYPES).toContain("game.action");
});
});
// ---------------------------------------------------------------------------
// T43 — WS protocol v2 schema + version negotiation
// ---------------------------------------------------------------------------
//
// Two flavours of test:
// - Schema parsing for the new v2 frames (request-choice / submit-choice
// / protocol-version-mismatch) and their discriminated union.
// - The pure `negotiateVersion` helper that resolves a client's
// declared `protocolVersion` against the server's supported set.
// Plus a backward-compat smoke for v1 — old envelope frames must still
// validate byte-identically through both `validateMessage` and the new
// unified `validateAnyMessage` entry point.
describe("T43 — supported version constants", () => {
it("exposes both v1 and v2 in the supported list", () => {
// The order matters for the protocol-version-mismatch wire frame:
// the server emits this list verbatim so clients can show a
// human-readable "supports v1 or v2" message. Lock it down.
expect([...SUPPORTED_PROTOCOL_VERSIONS]).toEqual([1, 2]);
});
it("LATEST_PROTOCOL_VERSION pins the highest supported version", () => {
expect(LATEST_PROTOCOL_VERSION).toBe(2);
});
it("envelope PROTOCOL_VERSION (the v1 wire literal) is unchanged", () => {
// v2 added new top-level frames; the v1 envelope itself is
// byte-identical pre/post T43, so this constant MUST stay 1.
expect(PROTOCOL_VERSION).toBe(1);
});
});
describe("T43 — negotiateVersion helper", () => {
it("resolves undefined to v1 (backward compat for pre-T43 clients)", () => {
// The whole point: a client that doesn't know about
// protocolVersion should keep working as a v1 client. The
// server uses this fallback to skip v2-only broadcasts via
// `shouldSkipV2Broadcast`, so request-choice never lands on a
// v1 client unable to render it.
expect(negotiateVersion(undefined)).toBe(1);
});
it("resolves 1 to 1 (explicit v1 declaration)", () => {
expect(negotiateVersion(1)).toBe(1);
});
it("resolves 2 to 2 (v2 client opts in)", () => {
expect(negotiateVersion(2)).toBe(2);
});
it("returns 'mismatch' for an unknown future version", () => {
// A future v3 client must NOT be silently downgraded — that
// would let the server pretend to understand a frame shape it
// doesn't, hiding bugs. Force the client to handle the
// protocol-version-mismatch frame instead.
expect(negotiateVersion(3)).toBe("mismatch");
expect(negotiateVersion(99)).toBe("mismatch");
});
it("returns 'mismatch' for nonsense values (negative, zero)", () => {
// Even though the wire schema rejects non-positive integers
// for the field, callers may receive unparsed inbound JSON
// before validation runs — the helper should refuse all
// non-supported values, not just the positive ones.
expect(negotiateVersion(0)).toBe("mismatch");
expect(negotiateVersion(-1)).toBe("mismatch");
});
});
describe("T43 — shouldSkipV2Broadcast", () => {
it("v1 clients skip v2-only broadcasts", () => {
// This is the wire-level expression of "v1 client receives no
// request-choice broadcasts" — the server-side broadcast
// helpers consult this before pushing a v2 frame.
expect(shouldSkipV2Broadcast(1)).toBe(true);
});
it("v2 clients receive v2-only broadcasts", () => {
expect(shouldSkipV2Broadcast(2)).toBe(false);
});
});
// Reusable v2 fixtures.
const validRequestChoice = {
kind: "request-choice",
protocolVersion: 2,
choiceId: "choice-1",
descriptorId: "modifier:rps-duel",
prompt: "Choose rock, paper, or scissors",
choiceKind: "rps",
forPlayer: "both",
options: ["rock", "paper", "scissors"],
timeout: 30_000,
expiresAtTimestamp: 1_745_000_030_000,
} as const satisfies RequestChoice;
const validSubmitChoice = {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "choice-1",
value: "rock",
} as const satisfies SubmitChoice;
describe("T43 — RequestChoiceSchema", () => {
it("accepts a fully-populated request-choice", () => {
const r = RequestChoiceSchema.safeParse(validRequestChoice);
expect(r.success).toBe(true);
});
it("accepts a minimal request-choice (no options/timeout/expires)", () => {
const r = RequestChoiceSchema.safeParse({
kind: "request-choice",
protocolVersion: 2,
choiceId: "c1",
descriptorId: "modifier:promote",
prompt: "Pick a piece",
choiceKind: "piece",
forPlayer: "white",
});
expect(r.success).toBe(true);
});
it("rejects protocolVersion !== 2 (v1 has no request-choice)", () => {
const r = RequestChoiceSchema.safeParse({
...validRequestChoice,
protocolVersion: 1,
});
expect(r.success).toBe(false);
});
it("rejects an unknown choiceKind", () => {
const r = RequestChoiceSchema.safeParse({
...validRequestChoice,
choiceKind: "coin-flip",
});
expect(r.success).toBe(false);
});
it("rejects empty-string choiceId", () => {
const r = RequestChoiceSchema.safeParse({
...validRequestChoice,
choiceId: "",
});
expect(r.success).toBe(false);
});
it("rejects negative timeout", () => {
const r = RequestChoiceSchema.safeParse({
...validRequestChoice,
timeout: -1,
});
expect(r.success).toBe(false);
});
it("accepts forPlayer = 'white' | 'black' | 'both'", () => {
for (const forPlayer of ["white", "black", "both"] as const) {
const r = RequestChoiceSchema.safeParse({
...validRequestChoice,
forPlayer,
});
expect(r.success).toBe(true);
}
});
it("rejects forPlayer = 'draw' (not a valid choice target)", () => {
const r = RequestChoiceSchema.safeParse({
...validRequestChoice,
forPlayer: "draw",
});
expect(r.success).toBe(false);
});
});
describe("T43 — SubmitChoiceSchema", () => {
it("accepts a well-formed submit-choice", () => {
const r = SubmitChoiceSchema.safeParse(validSubmitChoice);
expect(r.success).toBe(true);
});
it("accepts arbitrary value shapes (kind-dependent legality)", () => {
// The wire schema treats `value` as `unknown` because legal
// shape depends on the underlying `choiceKind` (string vs
// pieceId number vs square name). Engine-side validation
// gates the actual content; the wire just transports it.
const cases: unknown[] = [
"rock",
42,
{ square: "e4" },
["rock", "paper"],
null,
];
for (const value of cases) {
const r = SubmitChoiceSchema.safeParse({
...validSubmitChoice,
value,
});
expect(r.success).toBe(true);
}
});
it("rejects protocolVersion !== 2", () => {
const r = SubmitChoiceSchema.safeParse({
...validSubmitChoice,
protocolVersion: 1,
});
expect(r.success).toBe(false);
});
it("rejects empty choiceId", () => {
const r = SubmitChoiceSchema.safeParse({
...validSubmitChoice,
choiceId: "",
});
expect(r.success).toBe(false);
});
it("rejects wrong literal kind", () => {
const r = SubmitChoiceSchema.safeParse({
...validSubmitChoice,
kind: "request-choice",
});
expect(r.success).toBe(false);
});
});
describe("T43 — ProtocolVersionMismatchSchema", () => {
it("accepts a well-formed mismatch frame", () => {
const r = ProtocolVersionMismatchSchema.safeParse({
kind: "protocol-version-mismatch",
supported: [1, 2],
});
expect(r.success).toBe(true);
});
it("requires at least one supported version", () => {
const r = ProtocolVersionMismatchSchema.safeParse({
kind: "protocol-version-mismatch",
supported: [],
});
expect(r.success).toBe(false);
});
it("rejects non-positive supported versions", () => {
const r = ProtocolVersionMismatchSchema.safeParse({
kind: "protocol-version-mismatch",
supported: [0, 1],
});
expect(r.success).toBe(false);
});
});
describe("T43 — V2MessageSchema (discriminated union)", () => {
it("narrows on `kind` to RequestChoice", () => {
const r = V2MessageSchema.safeParse(validRequestChoice);
expect(r.success).toBe(true);
if (r.success) expect(r.data.kind).toBe("request-choice");
});
it("narrows on `kind` to SubmitChoice", () => {
const r = V2MessageSchema.safeParse(validSubmitChoice);
expect(r.success).toBe(true);
});
it("narrows on `kind` to ProtocolVersionMismatch", () => {
const r = V2MessageSchema.safeParse({
kind: "protocol-version-mismatch",
supported: [1, 2],
});
expect(r.success).toBe(true);
});
it("rejects an unknown v2 kind", () => {
const r = V2MessageSchema.safeParse({
kind: "future-v3-frame",
protocolVersion: 2,
});
expect(r.success).toBe(false);
});
});
describe("T43 — validateAnyMessage entry-point routing", () => {
it("routes a v1 envelope frame through the v1 path", () => {
const r = validateAnyMessage({
...envelope,
type: "room.create",
payload: {},
});
expect(r.ok).toBe(true);
if (r.ok) {
expect(r.data.wire).toBe("v1");
if (r.data.wire === "v1") {
expect(r.data.message.type).toBe("room.create");
}
}
});
it("routes a v2 top-level frame through the v2 path", () => {
const r = validateAnyMessage(validRequestChoice);
expect(r.ok).toBe(true);
if (r.ok) {
expect(r.data.wire).toBe("v2");
if (r.data.wire === "v2") {
expect(r.data.message.kind).toBe("request-choice");
}
}
});
it("routes submit-choice through the v2 path", () => {
const r = validateAnyMessage(validSubmitChoice);
expect(r.ok).toBe(true);
if (r.ok) {
expect(r.data.wire).toBe("v2");
}
});
it("a v2 frame with stray `v` falls through to v1's clearer error", () => {
// Defensive routing: a malformed mix (kind=v2 frame plus a v
// field) should not silently parse as v2 — the v1 path's
// VERSION_MISMATCH gate gives the human a more obvious clue
// about what's wrong (likely a copy-paste of the wrong
// envelope shape).
const r = validateAnyMessage({
v: 99,
kind: "submit-choice",
protocolVersion: 2,
choiceId: "c1",
value: 1,
});
expect(r.ok).toBe(false);
if (!r.ok) expect(r.error).toMatch(/VERSION_MISMATCH/);
});
it("propagates JSON.parse errors via validateAnyMessageString", () => {
const r = validateAnyMessageString("{not json");
expect(r.ok).toBe(false);
if (!r.ok) expect(r.error).toMatch(/INVALID_MESSAGE.*malformed JSON/);
});
it("rejects v2-shaped frames missing required fields", () => {
const r = validateAnyMessage({
kind: "request-choice",
protocolVersion: 2,
// missing choiceId / descriptorId / prompt / choiceKind / forPlayer
});
expect(r.ok).toBe(false);
if (!r.ok) expect(r.error).toMatch(/INVALID_MESSAGE/);
});
});
describe("T43 — backward compat (v1 envelope unchanged)", () => {
it("v1 frame without protocolVersion field still parses (legacy clients)", () => {
// Pre-T43 clients have no idea protocolVersion exists; their
// frames omit it entirely. Server must continue to accept.
const r = validateMessage({
...envelope,
type: "room.create",
payload: {},
});
expect(r.ok).toBe(true);
});
it("v1 frame WITH protocolVersion: 2 still parses (v2 client opt-in)", () => {
// This is the actual handshake path: a v2-aware client adds
// the field to its first envelope frame to declare capability.
// The envelope must accept it without breaking the schema.
const r = validateMessage({
...envelope,
protocolVersion: 2,
type: "room.create",
payload: {},
});
expect(r.ok).toBe(true);
});
it("v1 frame with protocolVersion: 1 also parses (explicit declaration)", () => {
const r = validateMessage({
...envelope,
protocolVersion: 1,
type: "room.join",
payload: { code: "ABC123" },
});
expect(r.ok).toBe(true);
});
it("validateMessage rejects v2 top-level frames (v1-only entry point)", () => {
// The legacy `validateMessage` is v1-only by contract — v2
// frames have no `v` envelope and must take the
// `validateAnyMessage` path instead. This test pins the
// separation so callers don't accidentally route v2 traffic
// into the v1 parser.
const r = validateMessage(validRequestChoice);
expect(r.ok).toBe(false);
if (!r.ok) expect(r.error).toMatch(/VERSION_MISMATCH/);
});
it("validateMessageString remains v1-only (string entry point)", () => {
const r = validateMessageString(JSON.stringify(validSubmitChoice));
expect(r.ok).toBe(false);
if (!r.ok) expect(r.error).toMatch(/VERSION_MISMATCH/);
});
it("envelope rejects negative protocolVersion field value", () => {
// The field is `z.number().int().positive().optional()` —
// negative values fail at the envelope shape gate before
// negotiateVersion ever runs. Belt-and-braces.
const r = validateMessage({
...envelope,
protocolVersion: -1,
type: "room.leave",
payload: {},
});
expect(r.ok).toBe(false);
});
});

View file

@ -1,4 +1,4 @@
// Chess server WebSocket protocol v1 — Zod schemas & validation.
// Chess server WebSocket protocol v1/v2 — Zod schemas & validation.
// See PROTOCOL.md for the full spec.
import { z } from "zod";
import type { ModifierProfile } from "@paratype/chess";
@ -7,8 +7,36 @@ import type { ModifierProfile } from "@paratype/chess";
// Primitives
// ---------------------------------------------------------------------------
/**
* Wire envelope version for v1 messages. A *separate* concept from
* `protocolVersion` (T43): the envelope `v` discriminates the message
* envelope shape (introduced in v1 and unchanged), while
* `protocolVersion` is a CLIENT-DECLARED capability flag negotiated
* at handshake time. v2 introduces new message kinds (request-choice
* / submit-choice) without altering the v1 envelope, so existing
* clients keep parsing their own traffic byte-for-byte unchanged.
*/
export const PROTOCOL_VERSION = 1 as const;
/**
* Highest *client capability* version this server understands (T43).
* Clients announce their capability via `protocolVersion` on the
* first frame (`room.create` / `room.join`); see `negotiateVersion`.
*
* v1 — original protocol; no request-choice / submit-choice support.
* v2 — adds the player-choice flow (T43-T50).
*
* Older v1 clients (or any client that omits `protocolVersion`) are
* treated as v1 and never receive request-choice broadcasts, so
* trigger suspension degrades gracefully (auto-resolve fallback;
* see plan T43 line 1442). Unknown versions are rejected fatally
* with `protocol-version-mismatch`.
*/
export const SUPPORTED_PROTOCOL_VERSIONS = [1, 2] as const;
export type SupportedProtocolVersion =
(typeof SUPPORTED_PROTOCOL_VERSIONS)[number];
export const LATEST_PROTOCOL_VERSION = 2 as const;
export const ColorSchema = z.enum(["white", "black"]);
export type Color = z.infer<typeof ColorSchema>;
@ -111,6 +139,16 @@ const envelopeShape = {
seq: z.number().int().nonnegative(),
ts: z.number().int().positive(),
token: z.string().uuid().optional(),
/**
* T43: client capability declaration. Clients announce on their
* first frame which top-level protocol version they speak; the
* server pins this on `ws.data.protocolVersion` for the duration
* of the connection. Optional for backward compat — absent =
* treated as v1 by `negotiateVersion`. The envelope `v` field
* remains `1` because the envelope SHAPE itself didn't change in
* v2; only new message *kinds* were added at the top level.
*/
protocolVersion: z.number().int().positive().optional(),
} as const;
// A permissive envelope-only parser used to inspect `v` and `type` before
@ -320,6 +358,51 @@ void _modifierProfileKeyCheck;
export const PreferredColorSchema = z.enum(["white", "black", "random"]);
export type PreferredColor = z.infer<typeof PreferredColorSchema>;
/**
* T50 — per-game choice-timeout policy. Threaded through `room.create`
* into `EngineOptions.choiceTimeout`; the chess engine seeds the
* resolved value onto `GAME_ENTITY` under the `ChoiceTimeoutPolicy`
* attr at construction time so the WS-layer timer + disconnect
* handler (T49) has a single authoritative source bound to the
* session.
*
* Two modes (locked verbatim by the chess-side `decisions.md`
* § "Choice Timeout & Disconnect"):
* - `"timeout-with-default"` — `seconds` is the per-choice budget;
* on expiry the server auto-submits the FIRST option of the
* pending request-choice and resumes (T49). `seconds` MUST be
* `>= 1` — non-positive values are nonsensical and would
* disable the very feature the policy enables. The wire schema
* enforces the lower bound here so the engine layer can trust
* the value verbatim. The schema does NOT cap the upper bound;
* UX guidance (per the plan) is to keep values reasonable
* (~30–120s) but extreme values land on the engine as-is.
* - `"no-timeout"` — no timer is armed; pending choices wait
* indefinitely. T49 routes a mid-choice disconnect under this
* mode to a "paused" game state instead of a forfeit.
*/
export const ChoiceTimeoutPolicySchema = z.discriminatedUnion("mode", [
z.object({
mode: z.literal("timeout-with-default"),
seconds: z.number().int().min(1),
}),
z.object({
mode: z.literal("no-timeout"),
}),
]);
export type ChoiceTimeoutPolicy = z.infer<typeof ChoiceTimeoutPolicySchema>;
/**
* T50 — canonical fallback for `room.create.choiceTimeout` when the
* field is omitted on the wire. Mirrors the chess-side
* `DEFAULT_CHOICE_TIMEOUT_POLICY` so legacy clients that never send
* the field land on the same engine state as new clients that do.
*/
export const DEFAULT_CHOICE_TIMEOUT_POLICY: ChoiceTimeoutPolicy = {
mode: "timeout-with-default",
seconds: 60,
};
export const RoomCreatePayloadSchema = z.object({
rulesetIds: z.array(z.string()).optional(),
layout: LayoutRequestSchema.optional(),
@ -337,6 +420,16 @@ export const RoomCreatePayloadSchema = z.object({
* with legacy clients that never sent a preference.
*/
preferredColor: PreferredColorSchema.optional(),
/**
* T50 — per-game choice-timeout policy threaded into the engine at
* construction time. Optional on the wire so legacy clients that
* don't yet send the field continue to work; when omitted the
* server falls back to {@link DEFAULT_CHOICE_TIMEOUT_POLICY}
* (`{ mode: "timeout-with-default", seconds: 60 }`). Validated by
* {@link ChoiceTimeoutPolicySchema}: `seconds >= 1` is enforced
* here so the engine layer can trust the value verbatim.
*/
choiceTimeout: ChoiceTimeoutPolicySchema.optional(),
});
export type RoomCreatePayload = z.infer<typeof RoomCreatePayloadSchema>;
@ -962,6 +1055,179 @@ export const KNOWN_MESSAGE_TYPES = [
] as const;
export type MessageType = (typeof KNOWN_MESSAGE_TYPES)[number];
// ---------------------------------------------------------------------------
// T43 — Protocol v2: request-choice / submit-choice / version negotiation
// ---------------------------------------------------------------------------
//
// v2 introduces the *player-choice flow* (plan T43-T50) used by the
// suspended-execution primitive `request-choice`. The new message
// kinds travel as TOP-LEVEL frames (no v1 envelope wrapping) because
// (a) they predate any room state on the wire and (b) the v1 envelope
// was deliberately left untouched so existing v1 clients keep parsing
// their own traffic byte-for-byte unchanged.
//
// Discriminator: `kind` (string literal). This avoids colliding with
// the v1 envelope's `type` field — a v1 client's envelope parser
// inspects `type`, not `kind`, and `kind` ≠ `type` means an old
// client never accidentally narrows a v2 message into a v1 union.
export const ChoiceKindSchema = z.enum([
"rps",
"piece",
"square",
"column",
"row",
]);
export type ChoiceKind = z.infer<typeof ChoiceKindSchema>;
/**
* Which player(s) the request is targeted at. `"both"` is used for
* RPS-style simultaneous choices where each player submits privately
* and the resolver merges the two values into the binding (plan
* T47 line 1499). Single-color values restrict the prompt to that
* specific player; the opposite player's submit is rejected as an
* unauthorised submission.
*/
export const ChoiceForPlayerSchema = z.enum(["white", "black", "both"]);
export type ChoiceForPlayer = z.infer<typeof ChoiceForPlayerSchema>;
/**
* Server → client: ask a (subset of) players to make a structured
* choice. Sent only to clients that negotiated `protocolVersion >= 2`;
* v1 clients receive nothing for this and the server's
* suspended-execution layer falls back to "auto-resolve with the
* first option" (plan T43 line 1442; wired in T44/T46).
*
* The shape is intentionally FLAT (no v1 envelope) to keep v2 frames
* easy to introspect on the wire and to avoid forcing clients to
* construct a fake `seq`/`ts` envelope around a server-pushed prompt.
*/
export const RequestChoiceSchema = z.object({
kind: z.literal("request-choice"),
protocolVersion: z.literal(2),
/** Server-minted unique id per choice. Clients echo this back on
* `submit-choice` so the server can match the response to the
* top-of-stack PendingChoice (plan T44 line 1457: "validate
* choiceId matches top of stack"). */
choiceId: z.string().min(1),
/** Owning descriptor (modifier id) — for diagnostics + so future
* presets can scope choice resolution per descriptor. */
descriptorId: z.string().min(1),
/** Human-readable prompt copy for the client UI. The server doesn't
* localise — clients are responsible for rendering. */
prompt: z.string(),
/** Discriminates the *value space* the client must choose from.
* Named `choiceKind` rather than `kind` to avoid collision with
* the message-level `kind` discriminator. */
choiceKind: ChoiceKindSchema,
forPlayer: ChoiceForPlayerSchema,
/** Optional concrete option list (e.g. specific squares/pieces).
* Schema-level `unknown` because the legal value-space depends
* on `choiceKind` and the engine validates server-side; the wire
* just transports the array intact. */
options: z.array(z.unknown()).optional(),
/** Optional milliseconds-from-now deadline. */
timeout: z.number().int().positive().optional(),
/** Optional unix-ms wall-clock deadline. Coexists with `timeout`
* so reconnecting clients can compute remaining time accurately
* without trusting their local relative clock. */
expiresAtTimestamp: z.number().int().nonnegative().optional(),
});
export type RequestChoice = z.infer<typeof RequestChoiceSchema>;
/**
* Client → server: resolve a previously-broadcast choice. The
* server's T44 handler validates that `choiceId` matches the top
* of the LIFO stack and that the submitting player is authorised
* (the prompt's `forPlayer` includes their color). `value` is
* unstructured at the wire because legal shapes are
* `choiceKind`-dependent (a square = string, an rps = "rock"|...,
* a piece = pieceId number). Engine-side validation gates on
* `choiceKind` after envelope parsing.
*/
export const SubmitChoiceSchema = z.object({
kind: z.literal("submit-choice"),
protocolVersion: z.literal(2),
choiceId: z.string().min(1),
value: z.unknown(),
});
export type SubmitChoice = z.infer<typeof SubmitChoiceSchema>;
/**
* Server → client: connection-fatal handshake rejection. Emitted
* when a client's declared `protocolVersion` is unrecognised
* (e.g. a future v3 client connecting to a v1/v2-only server, or
* an obvious garbage value). `supported` lets the client display a
* targeted upgrade/downgrade message instead of guessing.
*
* Distinct from the v1 `error.code = "VERSION_MISMATCH"` path
* (which signals a bad envelope `v`). This is the *capability*
* mismatch — the client's REQUESTED version isn't on the menu.
*/
export const ProtocolVersionMismatchSchema = z.object({
kind: z.literal("protocol-version-mismatch"),
/** Versions the server understands. Wire-stable list (mirrors
* `SUPPORTED_PROTOCOL_VERSIONS`) so clients can display a
* human-readable "this server supports v1 or v2" message. */
supported: z.array(z.number().int().positive()).min(1),
});
export type ProtocolVersionMismatch = z.infer<
typeof ProtocolVersionMismatchSchema
>;
/** Discriminated union of every v2 top-level frame. Used by callers
* that need to validate an inbound v2 frame without going through
* the v1 envelope path (e.g. T44 will parse `submit-choice` here). */
export const V2MessageSchema = z.discriminatedUnion("kind", [
RequestChoiceSchema,
SubmitChoiceSchema,
ProtocolVersionMismatchSchema,
]);
export type V2Message = z.infer<typeof V2MessageSchema>;
/**
* Outcome of negotiating a client's declared `protocolVersion`
* against `SUPPORTED_PROTOCOL_VERSIONS`. The server stores the
* resolved version on `ws.data.protocolVersion` and uses it to
* gate v2-only outbound traffic (e.g. request-choice broadcasts).
*
* - `undefined` (field omitted) → resolves to v1 for backward
* compat with pre-T43 clients. Such clients never receive
* v2 broadcasts.
* - `1` or `2` → resolves to that version verbatim.
* - Anything else → returns `"mismatch"`. Caller must emit a
* `protocol-version-mismatch` frame and close the socket
* (fail-fast — the connection cannot recover).
*
* Pure function (no side effects, no I/O); easy to unit-test.
*/
export function negotiateVersion(
clientVersion: number | undefined,
): SupportedProtocolVersion | "mismatch" {
if (clientVersion === undefined) return 1;
if (
(SUPPORTED_PROTOCOL_VERSIONS as readonly number[]).includes(clientVersion)
) {
return clientVersion as SupportedProtocolVersion;
}
return "mismatch";
}
/**
* Should the server skip a v2-only broadcast for a client at the
* given negotiated version? Centralises the "v1 client doesn't
* receive request-choice" rule so the broadcast layer (T44) has
* one explicit place to consult.
*
* Returns `true` when the client is on v1 (pre-T43); `false`
* when on v2 (the only version that understands request-choice).
*/
export function shouldSkipV2Broadcast(
negotiated: SupportedProtocolVersion,
): boolean {
return negotiated < 2;
}
// ---------------------------------------------------------------------------
// validateMessage — Result-style entry point
// ---------------------------------------------------------------------------
@ -972,7 +1238,17 @@ const ok = <T>(data: T): Result<T, never> => ({ ok: true, data });
const err = <E>(error: E): Result<never, E> => ({ ok: false, error });
/**
* Validate a decoded JSON value against the protocol.
* Tagged union returned by `validateAnyMessage` (T43). v1 frames
* narrow to `AnyMessage`; v2 top-level frames (request-choice /
* submit-choice / version-mismatch) narrow to `V2Message`. Callers
* branch on `data.wire` to tell them apart.
*/
export type ValidatedMessage =
| { wire: "v1"; message: AnyMessage }
| { wire: "v2"; message: V2Message };
/**
* Validate a decoded JSON value against the v1 envelope protocol.
*
* The input is `unknown` — callers that start from a raw string MUST
* `JSON.parse` first (and catch its throw) before handing a value here.
@ -981,6 +1257,13 @@ const err = <E>(error: E): Result<never, E> => ({ ok: false, error });
* message. On failure returns `{ ok: false, error }` with a descriptive
* string. Version mismatches are surfaced with a `VERSION_MISMATCH:` prefix
* so callers can disconnect fatally without re-parsing.
*
* NB: This entry point is v1-ONLY. v2 top-level frames (request-choice /
* submit-choice / protocol-version-mismatch) are NOT routed through
* here — they have no `v` envelope field. The broadcast layer uses
* `validateAnyMessage` to handle both wire formats from one entry
* point. Existing v1 callers (and tests) keep their semantics
* byte-for-byte.
*/
export function validateMessage(raw: unknown): Result<AnyMessage, string> {
// 1. Shape-check the envelope first so we can give precise errors about
@ -1019,8 +1302,10 @@ export function validateMessage(raw: unknown): Result<AnyMessage, string> {
}
/**
* Convenience: parse a raw WebSocket string frame. Handles the JSON.parse
* throw and funnels it into the same Result shape as `validateMessage`.
* Convenience: parse a raw WebSocket string frame as a v1 envelope
* message. Handles the JSON.parse throw and funnels it into the
* same Result shape as `validateMessage`. v1 ONLY — see
* `validateAnyMessageString` for the v1+v2 unified entry point.
*/
export function validateMessageString(
raw: string,
@ -1035,6 +1320,59 @@ export function validateMessageString(
return validateMessage(decoded);
}
/**
* Unified v1+v2 entry point (T43). Routing:
* - Object with `v` field → v1 envelope path; equivalent to
* `validateMessage`.
* - Object with `kind` field but no `v` → v2 top-level frame;
* parsed via `V2MessageSchema`.
*
* Used by the WS broadcast layer so a single inbound dispatch can
* accept both v1 and v2 traffic without callers branching on shape.
*/
export function validateAnyMessage(
raw: unknown,
): Result<ValidatedMessage, string> {
if (typeof raw !== "object" || raw === null || Array.isArray(raw)) {
return err("INVALID_MESSAGE: message must be a JSON object");
}
const obj = raw as Record<string, unknown>;
// v2 top-level frame? `kind` is the discriminator AND `v` is
// absent (v2 frames are NOT wrapped in the v1 envelope — they're
// standalone top-level objects). We probe for `kind` first so a
// stray `v` field on a malformed v2 frame falls through to the
// v1 path's clearer error messages instead of being routed here.
if (obj["v"] === undefined && typeof obj["kind"] === "string") {
const parsedV2 = V2MessageSchema.safeParse(raw);
if (!parsedV2.success) {
return err(`INVALID_MESSAGE: ${formatZodError(parsedV2.error)}`);
}
return ok({ wire: "v2", message: parsedV2.data });
}
const v1 = validateMessage(raw);
if (!v1.ok) return err(v1.error);
return ok({ wire: "v1", message: v1.data });
}
/**
* String-frame variant of `validateAnyMessage`. Handles JSON.parse
* + dispatches to v1 envelope or v2 top-level path based on shape.
*/
export function validateAnyMessageString(
raw: string,
): Result<ValidatedMessage, string> {
let decoded: unknown;
try {
decoded = JSON.parse(raw);
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
return err(`INVALID_MESSAGE: malformed JSON (${msg})`);
}
return validateAnyMessage(decoded);
}
function formatZodError(error: z.ZodError): string {
// Collapse issues into a compact single-line description. Keeping this
// deterministic is useful for tests and log greppability.

View file

@ -118,6 +118,25 @@ export interface Room {
* the state machine linear — at any instant a room has AT MOST
* one proposal awaiting consent.
*/
/**
* T49 — set when a player disconnected mid-choice under a
* `no-timeout` `ChoiceTimeoutPolicy` (T50). The room is parked in
* a "paused" state: the standard reconnect-grace timer is
* suppressed (no auto-game-end on disconnect), the engine's
* pending choice frame stays on the stack, and downstream move /
* action handlers are expected to gate on this flag (out of T49
* scope — T49 owns the *transition into* paused, not the
* downstream move-gate policy). Cleared on reconnect by
* `handleReconnect` so the surviving prompt is re-broadcast.
*
* `byToken` is the leaver's player token; `since` is unix-ms of
* the disconnect, useful for diagnostics ("paused 3 minutes ago")
* without forcing a separate audit log.
*/
pausedByChoiceDisconnect?: {
byToken: string;
since: number;
};
proposalState?: {
/** The candidate profile, already validated against the layout
* at receipt time. Stored verbatim; if consent approves, this

View file

@ -0,0 +1,531 @@
// T44 — server emits request-choice to v2 clients on engine push, and
// validates submit-choice against the LIFO PendingChoices stack.
//
// We exercise the broadcast and validation layers directly via
// `handleMessage` + `broadcastTopChoiceIfNew`, mirroring the mock-WS
// pattern used in `broadcast.test.ts`. The chess engine's request-
// choice primitive (T47) is not yet wired into a fireable trigger,
// so we drive the stack with the public `pushPendingChoice` helper —
// this matches how T44's broadcast hook is reached at runtime
// (engine pushes → server peeks/broadcasts) and lets us assert the
// wire-level contract independently of T47's trigger plumbing.
import type { ServerWebSocket } from "bun";
import { describe, it, expect } from "vitest";
import { pushPendingChoice, type PendingChoice } from "@paratype/chess";
import {
broadcastTopChoiceIfNew,
handleMessage,
registerConnection,
sessionRegistry,
unregisterConnection,
type ClientData,
} from "./broadcast.js";
import { PROTOCOL_VERSION, type ClientMessage } from "./protocol.js";
// ---------------------------------------------------------------------------
// Mock ServerWebSocket — same shape as broadcast.test.ts
// ---------------------------------------------------------------------------
interface MockWs extends ServerWebSocket<ClientData> {
readonly sent: unknown[];
readonly closed: boolean;
}
function makeMockWs(clientId: string): MockWs {
const sent: unknown[] = [];
const closedFlag = { value: false };
const ws = {
data: { clientId } as ClientData,
sent,
get closed(): boolean {
return closedFlag.value;
},
send(msg: string | Buffer): number {
const str = typeof msg === "string" ? msg : msg.toString("utf8");
sent.push(JSON.parse(str));
return str.length;
},
close(): void {
closedFlag.value = true;
},
} as unknown as MockWs;
return ws;
}
function nextMsgOfType(
ws: MockWs,
type: string,
): { payload?: Record<string, unknown>; [k: string]: unknown } {
const idx = ws.sent.findIndex(
(m) =>
typeof m === "object" &&
m !== null &&
((m as { type?: unknown }).type === type ||
(m as { kind?: unknown }).kind === type),
);
if (idx < 0) {
const tags = ws.sent.map(
(m) => (m as { type?: string; kind?: string }).type ?? (m as { kind?: string }).kind,
);
throw new Error(
`no message of type/kind "${type}" in inbox (got ${JSON.stringify(tags)})`,
);
}
const msg = ws.sent[idx] as { payload?: Record<string, unknown> };
ws.sent.splice(idx, 1);
return msg;
}
function findMsgOfType(ws: MockWs, type: string): unknown | undefined {
return ws.sent.find(
(m) =>
typeof m === "object" &&
m !== null &&
((m as { type?: unknown }).type === type ||
(m as { kind?: unknown }).kind === type),
);
}
/** Send a v1-envelope message; optionally declare client capability. */
function sendClient(
ws: MockWs,
type: ClientMessage["type"],
payload: unknown,
opts: { seq?: number; protocolVersion?: number; token?: string } = {},
): void {
const envelope: Record<string, unknown> = {
v: PROTOCOL_VERSION,
seq: opts.seq ?? 1,
ts: Date.now(),
type,
payload,
};
if (opts.protocolVersion !== undefined) {
envelope["protocolVersion"] = opts.protocolVersion;
}
if (opts.token !== undefined) {
envelope["token"] = opts.token;
}
handleMessage(ws, JSON.stringify(envelope));
}
/** Send a v2 top-level frame (no envelope wrapping). */
function sendV2(ws: MockWs, frame: Record<string, unknown>): void {
handleMessage(ws, JSON.stringify(frame));
}
/** Set up a fresh room with both players connected as v2 clients. */
function setupRoom(opts?: {
whiteV2?: boolean;
blackV2?: boolean;
}): {
white: MockWs;
black: MockWs;
code: string;
whiteToken: string;
blackToken: string;
} {
const whiteV2 = opts?.whiteV2 ?? true;
const blackV2 = opts?.blackV2 ?? true;
const white = makeMockWs(`white-${Math.random().toString(36).slice(2, 8)}`);
const black = makeMockWs(`black-${Math.random().toString(36).slice(2, 8)}`);
registerConnection(white);
registerConnection(black);
sendClient(
white,
"room.create",
{ rulesetIds: [] },
whiteV2 ? { protocolVersion: 2 } : {},
);
const created = nextMsgOfType(white, "room.created");
const code = created["payload"]!["code"] as string;
const whiteToken = created["payload"]!["token"] as string;
sendClient(
black,
"room.join",
{ code },
blackV2 ? { protocolVersion: 2 } : {},
);
const joined = nextMsgOfType(black, "room.joined");
const blackToken = joined["payload"]!["token"] as string;
// Drain the game.state frames each side gets at start.
nextMsgOfType(white, "game.state");
nextMsgOfType(black, "game.state");
return { white, black, code, whiteToken, blackToken };
}
function buildPendingChoice(overrides: Partial<PendingChoice>): PendingChoice {
return {
choiceId: "choice-1",
descriptorId: "test-descriptor",
triggerPath: [0],
primitiveIndex: 0,
bindings: new Map(),
kind: "rps",
prompt: "rock-paper-scissors?",
forPlayer: "both",
...overrides,
};
}
// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------
describe("T44 — request-choice broadcast", () => {
it("v2 clients matching forPlayer=both both receive the prompt", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "rps-1",
kind: "rps",
forPlayer: "both",
prompt: "pick",
}),
);
broadcastTopChoiceIfNew(code, session);
const wMsg = nextMsgOfType(white, "request-choice");
expect(wMsg["choiceId"]).toBe("rps-1");
expect(wMsg["choiceKind"]).toBe("rps");
expect(wMsg["forPlayer"]).toBe("both");
expect(wMsg["protocolVersion"]).toBe(2);
const bMsg = nextMsgOfType(black, "request-choice");
expect(bMsg["choiceId"]).toBe("rps-1");
unregisterConnection(white);
unregisterConnection(black);
});
it("v2 client with non-matching color does NOT receive a single-color prompt", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "white-only",
kind: "square",
forPlayer: "white",
}),
);
broadcastTopChoiceIfNew(code, session);
nextMsgOfType(white, "request-choice");
expect(findMsgOfType(black, "request-choice")).toBeUndefined();
unregisterConnection(white);
unregisterConnection(black);
});
it("v1 client receives NO request-choice broadcast (T43 negotiation)", () => {
// White is v1, black is v2. Both get a forPlayer="both" prompt;
// only the v2 client should see it.
const { white, black, code } = setupRoom({
whiteV2: false,
blackV2: true,
});
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "v1-skip",
forPlayer: "both",
kind: "rps",
}),
);
broadcastTopChoiceIfNew(code, session);
expect(findMsgOfType(white, "request-choice")).toBeUndefined();
nextMsgOfType(black, "request-choice");
unregisterConnection(white);
unregisterConnection(black);
});
it("broadcastTopChoiceIfNew is idempotent — re-calling does not re-emit", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({ choiceId: "once", forPlayer: "both" }),
);
broadcastTopChoiceIfNew(code, session);
broadcastTopChoiceIfNew(code, session);
// Each socket gets exactly one request-choice frame.
const whiteFrames = white.sent.filter(
(m) => (m as { kind?: string }).kind === "request-choice",
);
const blackFrames = black.sent.filter(
(m) => (m as { kind?: string }).kind === "request-choice",
);
expect(whiteFrames).toHaveLength(1);
expect(blackFrames).toHaveLength(1);
unregisterConnection(white);
unregisterConnection(black);
});
});
describe("T44 — submit-choice validation", () => {
it("rejects submit-choice when no pending choice on the stack", () => {
const { white, black, code, whiteToken } = setupRoom();
void black;
void code;
void whiteToken;
sendV2(white, {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "nonexistent",
value: "rock",
});
const err = nextMsgOfType(white, "error");
expect(err["payload"]!["code"]).toBe("INVALID_MESSAGE");
expect(String(err["payload"]!["message"])).toContain("no pending choice");
unregisterConnection(white);
unregisterConnection(black);
});
it("rejects submit-choice with mismatched choiceId (LIFO violation)", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "outer",
kind: "rps",
forPlayer: "both",
}),
);
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "inner",
kind: "rps",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(code, session);
nextMsgOfType(white, "request-choice");
nextMsgOfType(black, "request-choice");
// White attempts to resolve the OUTER frame while INNER is on top.
sendV2(white, {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "outer",
value: "rock",
});
const err = nextMsgOfType(white, "error");
expect(err["payload"]!["code"]).toBe("INVALID_MESSAGE");
expect(String(err["payload"]!["message"])).toContain("choiceId mismatch");
unregisterConnection(white);
unregisterConnection(black);
});
it("rejects submit-choice whose value does not match the kind", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "sq-1",
kind: "square",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(code, session);
nextMsgOfType(white, "request-choice");
nextMsgOfType(black, "request-choice");
// Square requires a number 0..63; "rock" is rps-shaped garbage.
sendV2(white, {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "sq-1",
value: "rock",
});
const err = nextMsgOfType(white, "error");
expect(err["payload"]!["code"]).toBe("INVALID_MESSAGE");
expect(String(err["payload"]!["message"])).toContain(
"protocol.invalid-choice-value",
);
unregisterConnection(white);
unregisterConnection(black);
});
it("accepts a well-formed submit-choice and pops the top frame", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "rps-ok",
kind: "rps",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(code, session);
nextMsgOfType(white, "request-choice");
nextMsgOfType(black, "request-choice");
sendV2(white, {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "rps-ok",
value: "rock",
});
// No error frame should arrive.
expect(findMsgOfType(white, "error")).toBeUndefined();
// Stack is now empty — peek confirms the pop happened.
// Re-broadcast call must be a no-op.
broadcastTopChoiceIfNew(code, session);
unregisterConnection(white);
unregisterConnection(black);
});
it("rejects submit-choice from a player not authorised by forPlayer", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "white-only",
kind: "rps",
forPlayer: "white",
}),
);
broadcastTopChoiceIfNew(code, session);
nextMsgOfType(white, "request-choice");
// Black tries to resolve a white-only prompt.
sendV2(black, {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "white-only",
value: "rock",
});
const err = nextMsgOfType(black, "error");
expect(err["payload"]!["code"]).toBe("BAD_TOKEN");
unregisterConnection(white);
unregisterConnection(black);
});
it("validates each kind's value-space — square accepts 0..63, rejects 64", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "sq-edge",
kind: "square",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(code, session);
nextMsgOfType(white, "request-choice");
nextMsgOfType(black, "request-choice");
// 64 is out of range.
sendV2(white, {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "sq-edge",
value: 64,
});
const err = nextMsgOfType(white, "error");
expect(String(err["payload"]!["message"])).toContain(
"protocol.invalid-choice-value",
);
unregisterConnection(white);
unregisterConnection(black);
});
it("validates each kind's value-space — column accepts 0..7, rejects 8", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "col-edge",
kind: "column",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(code, session);
nextMsgOfType(white, "request-choice");
nextMsgOfType(black, "request-choice");
sendV2(white, {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "col-edge",
value: 8,
});
const err = nextMsgOfType(white, "error");
expect(String(err["payload"]!["message"])).toContain(
"protocol.invalid-choice-value",
);
unregisterConnection(white);
unregisterConnection(black);
});
it("validates piece kind — accepts non-negative integer, rejects negative", () => {
const { white, black, code } = setupRoom();
const session = sessionRegistry.get(code)!;
pushPendingChoice(
session.getEngine(),
buildPendingChoice({
choiceId: "pc-1",
kind: "piece",
forPlayer: "both",
}),
);
broadcastTopChoiceIfNew(code, session);
nextMsgOfType(white, "request-choice");
nextMsgOfType(black, "request-choice");
sendV2(white, {
kind: "submit-choice",
protocolVersion: 2,
choiceId: "pc-1",
value: -1,
});
const err = nextMsgOfType(white, "error");
expect(String(err["payload"]!["message"])).toContain(
"protocol.invalid-choice-value",
);
unregisterConnection(white);
unregisterConnection(black);
});
});