From 78c589fdd232c83e093da7ba3e06fcf09b41a631 Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 18:20:33 -0400 Subject: [PATCH 1/3] bug: bind game-server sockets to their own match and drop events after failed auth A failed login on /ws/matches only called close(), which waits up to 30s for the closing handshake while ws keeps dispatching messages, and events sent while the password lookup was still running were never checked at all. Failed auth now terminates the socket, and the event handler refuses any socket that has not authenticated. An authenticated server could also post events for any match. Events are now accepted only from the match's server_id, and a match_map_id in the payload has to belong to that match. The lookup is cached per connection for 5s. Once a match ends its server_id is cleared, so the last server seen hosting it can still flush late events for an hour. Refused events are acknowledged so the plugin stops resending them. --- src/matches/match-events.gateway.spec.ts | 574 ++++++++++++++++++++++- src/matches/match-events.gateway.ts | 171 ++++++- 2 files changed, 708 insertions(+), 37 deletions(-) diff --git a/src/matches/match-events.gateway.spec.ts b/src/matches/match-events.gateway.spec.ts index b675b84f6..aa40c5a04 100644 --- a/src/matches/match-events.gateway.spec.ts +++ b/src/matches/match-events.gateway.spec.ts @@ -1,16 +1,82 @@ -// Isolate the gateway from its heavy DI imports; we only exercise the -// dedup/processing ordering logic in handleMatchEvent. +// Isolate the gateway from its heavy DI imports. jest.mock("./events", () => ({ MatchEvents: { testEvent: class {} } })); jest.mock("src/hasura/hasura.service", () => ({ HasuraService: class {} })); jest.mock("src/cache/cache.service", () => ({ CacheService: class {} })); -import { MatchEventsGateway } from "./match-events.gateway"; +import { + FiveStackGameServerWebSocketClient, + MatchEventsGateway, +} from "./match-events.gateway"; -function makeGateway(opts: { cacheHit?: boolean; processImpl?: () => any }) { +const SERVER_A = "a0000000-0000-4000-8000-00000000000a"; +const SERVER_B = "b0000000-0000-4000-8000-00000000000b"; +const MATCH_1 = "10000000-0000-4000-8000-000000000001"; +const MATCH_2 = "20000000-0000-4000-8000-000000000002"; +const MAP_1 = "11000000-0000-4000-8000-000000000011"; +const MAP_2 = "22000000-0000-4000-8000-000000000022"; + +type MatchRow = { + server_id: string | null; + status: string; + match_maps: Array<{ id: string }>; +}; + +function makeGateway( + opts: { + cacheHit?: boolean; + processImpl?: () => any; + matches?: Record; + servers?: Record; + serverQuery?: () => Promise; + } = {}, +) { + const store = new Map(); const cache = { - has: jest.fn().mockResolvedValue(opts.cacheHit ?? false), - put: jest.fn().mockResolvedValue(undefined), + has: jest.fn( + async (key: string) => (opts.cacheHit ?? false) || store.has(key), + ), + get: jest.fn(async (key: string) => store.get(key)), + put: jest.fn(async (key: string, value: unknown) => { + store.set(key, value); + return true; + }), + }; + + const matches: Record = opts.matches ?? { + [MATCH_1]: { + server_id: SERVER_A, + status: "Live", + match_maps: [{ id: MAP_1 }], + }, + [MATCH_2]: { + server_id: SERVER_B, + status: "Live", + match_maps: [{ id: MAP_2 }], + }, + }; + const servers = opts.servers ?? { + [SERVER_A]: { id: SERVER_A, api_password: "password-a" }, + [SERVER_B]: { id: SERVER_B, api_password: "password-b" }, }; + + const hasura = { + query: jest.fn(async (query: any) => { + if (query.servers_by_pk) { + if (opts.serverQuery) { + return opts.serverQuery(); + } + return { + servers_by_pk: servers[query.servers_by_pk.__args.id] ?? null, + }; + } + if (query.matches_by_pk) { + const row = matches[query.matches_by_pk.__args.id]; + return { matches_by_pk: row ? { ...row } : null }; + } + throw new Error("unexpected query"); + }), + }; + const processor = { setData: jest.fn(), process: jest.fn().mockImplementation(opts.processImpl ?? (async () => {})), @@ -25,52 +91,508 @@ function makeGateway(opts: { cacheHit?: boolean; processImpl?: () => any }) { const gateway = new MatchEventsGateway( logger as any, moduleRef as any, - {} as any, + hasura as any, cache as any, ); - return { gateway, cache, processor, logger }; + + const matchLookups = () => + hasura.query.mock.calls.filter(([query]) => query.matches_by_pk).length; + + return { + gateway, + cache, + store, + hasura, + matches, + processor, + moduleRef, + logger, + matchLookups, + }; } -const message = { - matchId: "match-1", - messageId: "msg-1", - data: { event: "testEvent", data: {} }, -}; +function socket() { + return { + terminate: jest.fn(), + close: jest.fn(), + } as unknown as FiveStackGameServerWebSocketClient & { + terminate: jest.Mock; + close: jest.Mock; + }; +} + +function authedSocket(serverId = SERVER_A) { + const client = socket(); + client.authenticated = true; + client.serverId = serverId; + return client; +} + +function basic(serverId: string, password: string) { + return { + headers: { + authorization: `Basic ${Buffer.from(`${serverId}:${password}`).toString("base64")}`, + }, + } as any; +} + +function event( + overrides: { + matchId?: unknown; + messageId?: string; + data?: Record; + } = {}, +) { + return { + matchId: MATCH_1, + messageId: "msg-1", + data: { event: "testEvent", data: {} }, + ...overrides, + ...(overrides.data + ? { data: { event: "testEvent", data: overrides.data } } + : {}), + } as any; +} describe("MatchEventsGateway.handleMatchEvent dedup", () => { it("marks the event processed only after process() succeeds", async () => { - const { gateway, cache, processor } = makeGateway({}); - const result = await gateway.handleMatchEvent(message as any); + const { gateway, cache, processor } = makeGateway(); + const result = await gateway.handleMatchEvent(authedSocket(), event()); expect(processor.process).toHaveBeenCalledTimes(1); - expect(cache.put).toHaveBeenCalledTimes(1); - // Ordering: process resolved before the dedup entry was written. + const dedupPuts = cache.put.mock.calls.filter(([key]) => + key.endsWith(":msg-1"), + ); + expect(dedupPuts).toHaveLength(1); expect(processor.process.mock.invocationCallOrder[0]).toBeLessThan( - cache.put.mock.invocationCallOrder[0], + cache.put.mock.invocationCallOrder[ + cache.put.mock.calls.findIndex(([key]) => key.endsWith(":msg-1")) + ], ); expect(result).toBe("msg-1"); }); it("does NOT write the dedup entry when process() throws (so redelivery retries)", async () => { - const { gateway, cache, processor } = makeGateway({ + const { gateway, store, processor } = makeGateway({ processImpl: async () => { throw new Error("hasura down"); }, }); - await expect(gateway.handleMatchEvent(message as any)).rejects.toThrow( - "hasura down", - ); + await expect( + gateway.handleMatchEvent(authedSocket(), event()), + ).rejects.toThrow("hasura down"); expect(processor.process).toHaveBeenCalledTimes(1); - expect(cache.put).not.toHaveBeenCalled(); + expect(store.has(`match-events:${MATCH_1}:msg-1`)).toBe(false); }); it("short-circuits on a dedup hit without processing", async () => { - const { gateway, cache, processor } = makeGateway({ cacheHit: true }); - const result = await gateway.handleMatchEvent(message as any); + const { gateway, processor } = makeGateway({ cacheHit: true }); + const result = await gateway.handleMatchEvent(authedSocket(), event()); expect(result).toBe("msg-1"); expect(processor.process).not.toHaveBeenCalled(); - expect(cache.put).not.toHaveBeenCalled(); + }); +}); + +describe("MatchEventsGateway.handleConnection", () => { + it("authenticates a server with the right password", async () => { + const { gateway } = makeGateway(); + const client = socket(); + + await gateway.handleConnection(client, basic(SERVER_A, "password-a")); + + expect(client.authenticated).toBe(true); + expect(client.serverId).toBe(SERVER_A); + expect(client.terminate).not.toHaveBeenCalled(); + }); + + it.each([ + ["no header", { headers: {} }], + ["a non-basic header", { headers: { authorization: "Bearer x" } }], + ["an empty basic header", { headers: { authorization: "Basic " } }], + [ + "credentials without a colon", + { + headers: { + authorization: `Basic ${Buffer.from("nocolon").toString("base64")}`, + }, + }, + ], + ["the wrong password", basic(SERVER_A, "password-b")], + ["an unknown server", basic(MATCH_1, "password-a")], + ])("terminates on %s", async (_, request) => { + const { gateway } = makeGateway(); + const client = socket(); + + await gateway.handleConnection(client, request as any); + + expect(client.terminate).toHaveBeenCalledTimes(1); + expect(client.close).not.toHaveBeenCalled(); + expect(client.authenticated).toBeFalsy(); + expect(client.serverId).toBeUndefined(); + }); + + it("terminates when the server lookup throws", async () => { + const { gateway } = makeGateway({ + serverQuery: async () => { + throw new Error("invalid uuid"); + }, + }); + const client = socket(); + + await gateway.handleConnection(client, basic("not-a-uuid", "x")); + + expect(client.terminate).toHaveBeenCalledTimes(1); + expect(client.authenticated).toBeFalsy(); + }); +}); + +describe("MatchEventsGateway unauthenticated clients", () => { + it("ignores events that arrive while auth is still pending and after it fails", async () => { + let resolveServer: (value: unknown) => void; + const { gateway, processor, moduleRef, matchLookups } = makeGateway({ + serverQuery: () => + new Promise((resolve) => { + resolveServer = resolve; + }), + }); + const client = socket(); + + const connecting = gateway.handleConnection( + client, + basic(SERVER_A, "wrong"), + ); + + await expect( + gateway.handleMatchEvent(client, event()), + ).resolves.toBeUndefined(); + + resolveServer!({ + servers_by_pk: { id: SERVER_A, api_password: "password-a" }, + }); + await connecting; + + expect(client.terminate).toHaveBeenCalledTimes(1); + + await expect( + gateway.handleMatchEvent(client, event({ messageId: "msg-2" })), + ).resolves.toBeUndefined(); + + expect(matchLookups()).toBe(0); + expect(moduleRef.resolve).not.toHaveBeenCalled(); + expect(processor.process).not.toHaveBeenCalled(); + }); + + it("ignores events on a socket that never authenticated", async () => { + const { gateway, processor } = makeGateway(); + + await expect( + gateway.handleMatchEvent(socket(), event()), + ).resolves.toBeUndefined(); + expect(processor.process).not.toHaveBeenCalled(); + }); +}); + +describe("MatchEventsGateway match binding", () => { + afterEach(() => { + jest.restoreAllMocks(); + }); + + it("accepts an event for the server's own match", async () => { + const { gateway, processor } = makeGateway(); + + const result = await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ data: { match_map_id: MAP_1 } }), + ); + + expect(result).toBe("msg-1"); + expect(processor.setData).toHaveBeenCalledWith(MATCH_1, { + match_map_id: MAP_1, + }); + expect(processor.process).toHaveBeenCalledTimes(1); + }); + + it("refuses, acknowledges and warns on an event for another server's match", async () => { + const { gateway, processor, logger, moduleRef } = makeGateway(); + + const result = await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ matchId: MATCH_2 }), + ); + + expect(result).toBe("msg-1"); + expect(moduleRef.resolve).not.toHaveBeenCalled(); + expect(processor.process).not.toHaveBeenCalled(); + expect(logger.warn).toHaveBeenCalledWith( + "game server event refused: match is not hosted by this server", + expect.objectContaining({ serverId: SERVER_A, matchId: MATCH_2 }), + ); + }); + + it("refuses its own match when the payload names another match's map", async () => { + const { gateway, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ data: { match_map_id: MAP_2 } }), + ); + + expect(processor.process).not.toHaveBeenCalled(); + }); + + it("refuses a match_map_id that is not a string", async () => { + const { gateway, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ data: { match_map_id: { _neq: MAP_1 } } }), + ); + + expect(processor.process).not.toHaveBeenCalled(); + }); + + it("passes a null match_map_id through, since it cannot name another match", async () => { + const { gateway, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ data: { match_map_id: null } }), + ); + + expect(processor.process).toHaveBeenCalledTimes(1); + }); + + it("refuses a match that does not exist", async () => { + const { gateway, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ matchId: "30000000-0000-4000-8000-000000000003" }), + ); + + expect(processor.process).not.toHaveBeenCalled(); + }); + + it.each([["not-a-uuid"], [42], [undefined], [{ id: MATCH_1 }]])( + "refuses a malformed match id (%p) without querying", + async (matchId) => { + const { gateway, processor, matchLookups } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ matchId }), + ); + + expect(matchLookups()).toBe(0); + expect(processor.process).not.toHaveBeenCalled(); + }, + ); + + it("looks the binding up once per connection within its ttl", async () => { + const now = jest.spyOn(Date, "now").mockReturnValue(1_000_000); + const { gateway, processor, matchLookups } = makeGateway(); + const client = authedSocket(SERVER_A); + + await gateway.handleMatchEvent(client, event({ messageId: "m1" })); + await gateway.handleMatchEvent(client, event({ messageId: "m2" })); + now.mockReturnValue(1_004_999); + await gateway.handleMatchEvent(client, event({ messageId: "m3" })); + + expect(matchLookups()).toBe(1); + expect(processor.process).toHaveBeenCalledTimes(3); + + now.mockReturnValue(1_005_001); + await gateway.handleMatchEvent(client, event({ messageId: "m4" })); + + expect(matchLookups()).toBe(2); + }); + + it("does not cache a refusal", async () => { + const { gateway, matchLookups, matches, processor } = makeGateway(); + const client = authedSocket(SERVER_A); + + await gateway.handleMatchEvent(client, event({ matchId: MATCH_2 })); + matches[MATCH_2].server_id = SERVER_A; + await gateway.handleMatchEvent( + client, + event({ matchId: MATCH_2, messageId: "msg-2" }), + ); + + expect(matchLookups()).toBe(2); + expect(processor.process).toHaveBeenCalledTimes(1); + }); + + describe("legitimate edge flows", () => { + it("moves with the match: the old server is refused and the new one accepted once the binding expires", async () => { + const now = jest.spyOn(Date, "now").mockReturnValue(1_000_000); + const { gateway, matches, processor } = makeGateway(); + const oldServer = authedSocket(SERVER_A); + + await gateway.handleMatchEvent(oldServer, event({ messageId: "m1" })); + expect(processor.process).toHaveBeenCalledTimes(1); + + matches[MATCH_1].server_id = SERVER_B; + now.mockReturnValue(1_010_000); + + await gateway.handleMatchEvent(oldServer, event({ messageId: "m2" })); + expect(processor.process).toHaveBeenCalledTimes(1); + + await gateway.handleMatchEvent( + authedSocket(SERVER_B), + event({ messageId: "m3" }), + ); + expect(processor.process).toHaveBeenCalledTimes(2); + }); + + it("re-checks on a fresh connection after a reconnect", async () => { + const { gateway, matches, processor, matchLookups } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + matches[MATCH_1].server_id = SERVER_B; + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m2" }), + ); + + expect(matchLookups()).toBe(2); + expect(processor.process).toHaveBeenCalledTimes(1); + }); + + it.each(["Finished", "Surrendered", "Canceled", "Forfeit", "Tie"])( + "lets the last host flush late events after the match ends (%s) and server_id is cleared", + async (status) => { + const { gateway, matches, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + matches[MATCH_1].server_id = null; + matches[MATCH_1].status = status; + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m2", data: { match_map_id: MAP_1 } }), + ); + + expect(processor.process).toHaveBeenCalledTimes(2); + }, + ); + + it("keeps an ended match closed to every other server", async () => { + const { gateway, matches, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + matches[MATCH_1].server_id = null; + matches[MATCH_1].status = "Finished"; + + await gateway.handleMatchEvent( + authedSocket(SERVER_B), + event({ messageId: "m2" }), + ); + + expect(processor.process).toHaveBeenCalledTimes(1); + }); + + it("refuses an ended match nobody was seen hosting", async () => { + const { gateway, processor } = makeGateway({ + matches: { + [MATCH_1]: { + server_id: null, + status: "Finished", + match_maps: [{ id: MAP_1 }], + }, + }, + }); + + await gateway.handleMatchEvent(authedSocket(SERVER_A), event()); + + expect(processor.process).not.toHaveBeenCalled(); + }); + + it("refuses the previous host once the match moved and then ended", async () => { + const { gateway, matches, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + matches[MATCH_1].server_id = SERVER_B; + await gateway.handleMatchEvent( + authedSocket(SERVER_B), + event({ messageId: "m2" }), + ); + expect(processor.process).toHaveBeenCalledTimes(2); + + matches[MATCH_1].server_id = null; + matches[MATCH_1].status = "Finished"; + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m3" }), + ); + expect(processor.process).toHaveBeenCalledTimes(2); + + await gateway.handleMatchEvent( + authedSocket(SERVER_B), + event({ messageId: "m4" }), + ); + expect(processor.process).toHaveBeenCalledTimes(3); + }); + + it("learns the new host from a refused lookup even if it never posted before the match ended", async () => { + const { gateway, matches, processor } = makeGateway(); + const oldServer = authedSocket(SERVER_A); + + await gateway.handleMatchEvent(oldServer, event({ messageId: "m1" })); + + matches[MATCH_1].server_id = SERVER_B; + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m2" }), + ); + + matches[MATCH_1].server_id = null; + matches[MATCH_1].status = "Canceled"; + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m3" }), + ); + + expect(processor.process).toHaveBeenCalledTimes(1); + }); + + it("refuses the previous host while the match waits for a new server", async () => { + const { gateway, matches, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + matches[MATCH_1].server_id = null; + matches[MATCH_1].status = "WaitingForServer"; + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m2" }), + ); + + expect(processor.process).toHaveBeenCalledTimes(1); + }); }); }); diff --git a/src/matches/match-events.gateway.ts b/src/matches/match-events.gateway.ts index d2e034758..ec05ba445 100644 --- a/src/matches/match-events.gateway.ts +++ b/src/matches/match-events.gateway.ts @@ -5,6 +5,7 @@ import { WebSocketGateway, } from "@nestjs/websockets"; import WebSocket from "ws"; +import { validate } from "uuid"; import { Request } from "express"; import { ModuleRef } from "@nestjs/core"; import { MatchEvents } from "./events"; @@ -13,16 +14,38 @@ import { Logger } from "@nestjs/common"; import { HasuraService } from "src/hasura/hasura.service"; import { CacheService } from "src/cache/cache.service"; import { timingSafeStringEqual } from "src/utilities/timingSafeStringEqual"; +import type { e_match_status_enum } from "../../generated"; + +type MatchBinding = { + expiresAt: number; + mapIds: Set; +}; export type FiveStackGameServerWebSocketClient = WebSocket.WebSocket & { - id: string; - matchId: string; + authenticated?: boolean; + serverId?: string; + matchBindings?: Map; }; @WebSocketGateway({ path: "/ws/matches", }) export class MatchEventsGateway { + // A server moved off a match is still accepted for it until this runs out. + // Checking every event instead would add a Hasura round trip to each damage + // and kill, which arrive many times a second during a round. + private static readonly BINDING_TTL_MS = 5 * 1000; + + private static readonly LAST_HOST_TTL_SECONDS = 60 * 60; + + private static readonly TERMINAL_STATUSES: readonly e_match_status_enum[] = [ + "Finished", + "Canceled", + "Forfeit", + "Tie", + "Surrendered", + ]; + constructor( private readonly logger: Logger, private readonly moduleRef: ModuleRef, @@ -31,7 +54,7 @@ export class MatchEventsGateway { ) {} async handleConnection( - @ConnectedSocket() client: WebSocket.WebSocket, + @ConnectedSocket() client: FiveStackGameServerWebSocketClient, request: Request, ) { try { @@ -41,7 +64,7 @@ export class MatchEventsGateway { this.logger.warn("game server connection rejected: missing auth", { ip: request.headers["cf-connecting-ip"], }); - client.close(); + client.terminate(); return; } @@ -50,7 +73,7 @@ export class MatchEventsGateway { this.logger.warn("game server connection rejected: malformed auth", { ip: request.headers["cf-connecting-ip"], }); - client.close(); + client.terminate(); return; } @@ -63,7 +86,7 @@ export class MatchEventsGateway { ip: request.headers["cf-connecting-ip"], }, ); - client.close(); + client.terminate(); return; } @@ -80,20 +103,28 @@ export class MatchEventsGateway { }, }); - if (!timingSafeStringEqual(servers_by_pk?.api_password, apiPassword)) { - client.close(); + if ( + !servers_by_pk?.id || + !timingSafeStringEqual(servers_by_pk.api_password, apiPassword) + ) { + client.terminate(); this.logger.warn("game server auth failure", { serverId, ip: request.headers["cf-connecting-ip"], }); + return; } + + client.serverId = servers_by_pk.id; + client.authenticated = true; } catch { - client.close(); + client.terminate(); } } @SubscribeMessage("events") async handleMatchEvent( + @ConnectedSocket() client: FiveStackGameServerWebSocketClient, @MessageBody() message: { mapId?: string; @@ -105,7 +136,26 @@ export class MatchEventsGateway { }; }, ) { + if (!client.authenticated || !client.serverId) { + return; + } + const { matchId, mapId, messageId } = message; + const { data, event } = message.data; + + if (!(await this.isHostedBy(client, matchId, data?.match_map_id))) { + this.logger.warn( + "game server event refused: match is not hosted by this server", + { + serverId: client.serverId, + matchId, + event, + }, + ); + // The plugin resends an unacknowledged message every few seconds for as + // long as it runs, so a refusal is acknowledged to make it drop it. + return messageId; + } const cacheKey = mapId ? `match-events:${matchId}:${mapId}:${messageId}` @@ -115,8 +165,6 @@ export class MatchEventsGateway { return messageId; } - const { data, event } = message.data; - const Processor = MatchEvents[event as keyof typeof MatchEvents]; if (!Processor) { @@ -149,4 +197,105 @@ export class MatchEventsGateway { return messageId; } + + // Stats and round events name their match map in the payload, so the map has + // to belong to the match too or a server could write into another match + // through its own. + private async isHostedBy( + client: FiveStackGameServerWebSocketClient, + matchId: unknown, + matchMapId: unknown, + ): Promise { + if (typeof matchId !== "string" || !validate(matchId)) { + return false; + } + + if ( + matchMapId !== undefined && + matchMapId !== null && + typeof matchMapId !== "string" + ) { + return false; + } + + const bindingKey = matchId.toLowerCase(); + const mapId = + typeof matchMapId === "string" ? matchMapId.toLowerCase() : undefined; + + const cached = client.matchBindings?.get(bindingKey); + if ( + cached && + cached.expiresAt > Date.now() && + (!mapId || cached.mapIds.has(mapId)) + ) { + return true; + } + + const mapIds = await this.hostedMatchMapIds(client.serverId, bindingKey); + + client.matchBindings ??= new Map(); + + if (!mapIds) { + client.matchBindings.delete(bindingKey); + return false; + } + + client.matchBindings.set(bindingKey, { + expiresAt: Date.now() + MatchEventsGateway.BINDING_TTL_MS, + mapIds, + }); + + return !mapId || mapIds.has(mapId); + } + + // Ending a match clears its server_id while the server is still flushing + // late events (the map Finished that follows a surrender, chat, disconnects, + // retries), so an ended match stays open to the last server seen hosting it. + private async hostedMatchMapIds( + serverId: string, + matchId: string, + ): Promise | null> { + const { matches_by_pk: match } = await this.hasura.query({ + matches_by_pk: { + __args: { + id: matchId, + }, + server_id: true, + status: true, + match_maps: { + id: true, + }, + }, + }); + + if (!match) { + return null; + } + + const lastHostKey = MatchEventsGateway.lastHostKey(matchId); + + let hosted: boolean; + if (match.server_id) { + await this.cache.put( + lastHostKey, + match.server_id, + MatchEventsGateway.LAST_HOST_TTL_SECONDS, + ); + hosted = match.server_id === serverId; + } else { + hosted = + MatchEventsGateway.TERMINAL_STATUSES.includes(match.status) && + (await this.cache.get(lastHostKey)) === serverId; + } + + if (!hosted) { + return null; + } + + return new Set(match.match_maps.map(({ id }) => id)); + } + + private static lastHostKey(matchId: string) { + return `match-events:last-host:${matchId}`; + } } From 1d55926cb016b500a4d169b2c4714a044f6a9395 Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 18:49:35 -0400 Subject: [PATCH 2/3] bug: close the review gaps in the game-server match binding - techTimeout names its map as map_id, which the binding never checked, so a server could still zero another match's timeouts. map_id is now held to the same rule as match_map_id. - captain updated the player's row in every lineup they were in, across all matches. It is now scoped to this match's two lineups. - After a match ends, its last host may only send mapStatus, chat and player-disconnected, and only for 10 minutes. Before, it could send any event, including surrender, for an hour. - Events sent while the password lookup is still running now wait for it instead of being dropped and resent 10s later out of order. --- src/matches/events/CaptainEvent.spec.ts | 65 +++++ src/matches/events/CaptainEvent.ts | 3 + src/matches/match-events.gateway.spec.ts | 320 ++++++++++++++++++----- src/matches/match-events.gateway.ts | 253 ++++++++++-------- 4 files changed, 473 insertions(+), 168 deletions(-) create mode 100644 src/matches/events/CaptainEvent.spec.ts diff --git a/src/matches/events/CaptainEvent.spec.ts b/src/matches/events/CaptainEvent.spec.ts new file mode 100644 index 000000000..4df5f946f --- /dev/null +++ b/src/matches/events/CaptainEvent.spec.ts @@ -0,0 +1,65 @@ +import { Logger } from "@nestjs/common"; +import CaptainEvent from "./CaptainEvent"; + +describe("CaptainEvent", () => { + let processor: CaptainEvent; + let hasura: { mutation: jest.Mock }; + let matchAssistant: { getMatchLineups: jest.Mock }; + + beforeEach(() => { + hasura = { mutation: jest.fn(async () => ({})) }; + matchAssistant = { + getMatchLineups: jest.fn(async () => ({ + lineup_1_id: "lineup-1", + lineup_2_id: "lineup-2", + lineup_players: [ + { steam_id: "76561198000000001", discord_id: null }, + { steam_id: null, discord_id: "discord-2", placeholder_name: "bob" }, + ] as Array>, + })), + }; + + processor = new CaptainEvent( + new Logger("CaptainEventTest"), + hasura as any, + matchAssistant as any, + {} as any, + {} as any, + ); + }); + + function updateWhere() { + return hasura.mutation.mock.calls[0][0].update_match_lineup_players.__args + .where; + } + + it("only changes the captain in this match's lineups", async () => { + processor.setData("11111111-1111-1111-1111-111111111111", { + claim: true, + steam_id: "76561198000000001", + player_name: "alice", + }); + + await processor.process(); + + expect(updateWhere()).toEqual({ + steam_id: { _eq: "76561198000000001" }, + match_lineup_id: { _in: ["lineup-1", "lineup-2"] }, + }); + }); + + it("scopes a placeholder player's claim to this match too", async () => { + processor.setData("11111111-1111-1111-1111-111111111111", { + claim: true, + steam_id: "0", + player_name: "bob", + }); + + await processor.process(); + + expect(updateWhere()).toEqual({ + discord_id: { _eq: "discord-2" }, + match_lineup_id: { _in: ["lineup-1", "lineup-2"] }, + }); + }); +}); diff --git a/src/matches/events/CaptainEvent.ts b/src/matches/events/CaptainEvent.ts index d546c9172..c266e8a91 100644 --- a/src/matches/events/CaptainEvent.ts +++ b/src/matches/events/CaptainEvent.ts @@ -35,6 +35,9 @@ export default class CaptainEvent extends MatchEventProcessor<{ [lineup_player.steam_id ? "steam_id" : "discord_id"]: { _eq: lineup_player.steam_id || lineup_player.discord_id, }, + match_lineup_id: { + _in: [match.lineup_1_id, match.lineup_2_id], + }, }, _set: { captain: this.data.claim, diff --git a/src/matches/match-events.gateway.spec.ts b/src/matches/match-events.gateway.spec.ts index aa40c5a04..09e58a57d 100644 --- a/src/matches/match-events.gateway.spec.ts +++ b/src/matches/match-events.gateway.spec.ts @@ -1,5 +1,22 @@ // Isolate the gateway from its heavy DI imports. -jest.mock("./events", () => ({ MatchEvents: { testEvent: class {} } })); +jest.mock("./events", () => { + class TestEvent {} + return { + MatchEvents: Object.fromEntries( + [ + "testEvent", + "mapStatus", + "chat", + "player-disconnected", + "surrender", + "score", + "restoreRound", + "techTimeout", + "kill", + ].map((name) => [name, TestEvent]), + ), + }; +}); jest.mock("src/hasura/hasura.service", () => ({ HasuraService: class {} })); jest.mock("src/cache/cache.service", () => ({ CacheService: class {} })); @@ -136,21 +153,30 @@ function basic(serverId: string, password: string) { } as any; } +async function connect( + gateway: MatchEventsGateway, + client: FiveStackGameServerWebSocketClient, + request: any, +) { + gateway.handleConnection(client, request); + await client.authentication; +} + function event( overrides: { matchId?: unknown; messageId?: string; + name?: string; data?: Record; } = {}, ) { return { - matchId: MATCH_1, - messageId: "msg-1", - data: { event: "testEvent", data: {} }, - ...overrides, - ...(overrides.data - ? { data: { event: "testEvent", data: overrides.data } } - : {}), + matchId: "matchId" in overrides ? overrides.matchId : MATCH_1, + messageId: overrides.messageId ?? "msg-1", + data: { + event: overrides.name ?? "testEvent", + data: overrides.data ?? {}, + }, } as any; } @@ -160,14 +186,12 @@ describe("MatchEventsGateway.handleMatchEvent dedup", () => { const result = await gateway.handleMatchEvent(authedSocket(), event()); expect(processor.process).toHaveBeenCalledTimes(1); - const dedupPuts = cache.put.mock.calls.filter(([key]) => + const dedupPut = cache.put.mock.calls.findIndex(([key]) => key.endsWith(":msg-1"), ); - expect(dedupPuts).toHaveLength(1); + expect(dedupPut).toBeGreaterThanOrEqual(0); expect(processor.process.mock.invocationCallOrder[0]).toBeLessThan( - cache.put.mock.invocationCallOrder[ - cache.put.mock.calls.findIndex(([key]) => key.endsWith(":msg-1")) - ], + cache.put.mock.invocationCallOrder[dedupPut], ); expect(result).toBe("msg-1"); }); @@ -200,7 +224,7 @@ describe("MatchEventsGateway.handleConnection", () => { const { gateway } = makeGateway(); const client = socket(); - await gateway.handleConnection(client, basic(SERVER_A, "password-a")); + await connect(gateway, client, basic(SERVER_A, "password-a")); expect(client.authenticated).toBe(true); expect(client.serverId).toBe(SERVER_A); @@ -225,7 +249,7 @@ describe("MatchEventsGateway.handleConnection", () => { const { gateway } = makeGateway(); const client = socket(); - await gateway.handleConnection(client, request as any); + await connect(gateway, client, request); expect(client.terminate).toHaveBeenCalledTimes(1); expect(client.close).not.toHaveBeenCalled(); @@ -241,7 +265,7 @@ describe("MatchEventsGateway.handleConnection", () => { }); const client = socket(); - await gateway.handleConnection(client, basic("not-a-uuid", "x")); + await connect(gateway, client, basic("not-a-uuid", "x")); expect(client.terminate).toHaveBeenCalledTimes(1); expect(client.authenticated).toBeFalsy(); @@ -249,30 +273,36 @@ describe("MatchEventsGateway.handleConnection", () => { }); describe("MatchEventsGateway unauthenticated clients", () => { - it("ignores events that arrive while auth is still pending and after it fails", async () => { + function pendingAuth() { let resolveServer: (value: unknown) => void; - const { gateway, processor, moduleRef, matchLookups } = makeGateway({ + const harness = makeGateway({ serverQuery: () => new Promise((resolve) => { resolveServer = resolve; }), }); + return { + ...harness, + finishAuth: () => + resolveServer({ + servers_by_pk: { id: SERVER_A, api_password: "password-a" }, + }), + }; + } + + it("never processes events sent while auth is pending or after it failed", async () => { + const { gateway, processor, moduleRef, matchLookups, finishAuth } = + pendingAuth(); const client = socket(); - const connecting = gateway.handleConnection( - client, - basic(SERVER_A, "wrong"), - ); + gateway.handleConnection(client, basic(SERVER_A, "wrong")); - await expect( - gateway.handleMatchEvent(client, event()), - ).resolves.toBeUndefined(); + const raced = gateway.handleMatchEvent(client, event()); + expect(processor.process).not.toHaveBeenCalled(); - resolveServer!({ - servers_by_pk: { id: SERVER_A, api_password: "password-a" }, - }); - await connecting; + finishAuth(); + await expect(raced).resolves.toBeUndefined(); expect(client.terminate).toHaveBeenCalledTimes(1); await expect( @@ -284,6 +314,22 @@ describe("MatchEventsGateway unauthenticated clients", () => { expect(processor.process).not.toHaveBeenCalled(); }); + it("holds events sent while auth is pending and processes them once it succeeds", async () => { + const { gateway, processor, finishAuth } = pendingAuth(); + const client = socket(); + + gateway.handleConnection(client, basic(SERVER_A, "password-a")); + + const early = gateway.handleMatchEvent(client, event()); + await Promise.resolve(); + expect(processor.process).not.toHaveBeenCalled(); + + finishAuth(); + + await expect(early).resolves.toBe("msg-1"); + expect(processor.process).toHaveBeenCalledTimes(1); + }); + it("ignores events on a socket that never authenticated", async () => { const { gateway, processor } = makeGateway(); @@ -331,28 +377,45 @@ describe("MatchEventsGateway match binding", () => { ); }); - it("refuses its own match when the payload names another match's map", async () => { - const { gateway, processor } = makeGateway(); + it.each(["match_map_id", "map_id"])( + "refuses its own match when %s names another match's map", + async (key) => { + const { gateway, processor } = makeGateway(); - await gateway.handleMatchEvent( - authedSocket(SERVER_A), - event({ data: { match_map_id: MAP_2 } }), - ); + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ name: "techTimeout", data: { [key]: MAP_2 } }), + ); - expect(processor.process).not.toHaveBeenCalled(); - }); + expect(processor.process).not.toHaveBeenCalled(); + }, + ); - it("refuses a match_map_id that is not a string", async () => { + it("accepts a techTimeout for its own map", async () => { const { gateway, processor } = makeGateway(); await gateway.handleMatchEvent( authedSocket(SERVER_A), - event({ data: { match_map_id: { _neq: MAP_1 } } }), + event({ name: "techTimeout", data: { map_id: MAP_1 } }), ); - expect(processor.process).not.toHaveBeenCalled(); + expect(processor.process).toHaveBeenCalledTimes(1); }); + it.each(["match_map_id", "map_id"])( + "refuses a %s that is not a string", + async (key) => { + const { gateway, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ data: { [key]: { _neq: MAP_1 } } }), + ); + + expect(processor.process).not.toHaveBeenCalled(); + }, + ); + it("passes a null match_map_id through, since it cannot name another match", async () => { const { gateway, processor } = makeGateway(); @@ -409,6 +472,70 @@ describe("MatchEventsGateway match binding", () => { expect(matchLookups()).toBe(2); }); + it("does not let a warm binding for one match admit another match", async () => { + const { gateway, processor } = makeGateway(); + const client = authedSocket(SERVER_A); + + await gateway.handleMatchEvent(client, event({ messageId: "m1" })); + await gateway.handleMatchEvent( + client, + event({ matchId: MATCH_2, messageId: "m2" }), + ); + + expect(processor.process).toHaveBeenCalledTimes(1); + }); + + it("does not let a warm binding admit another match's map", async () => { + const { gateway, processor } = makeGateway(); + const client = authedSocket(SERVER_A); + + await gateway.handleMatchEvent( + client, + event({ messageId: "m1", data: { match_map_id: MAP_1 } }), + ); + await gateway.handleMatchEvent( + client, + event({ messageId: "m2", data: { match_map_id: MAP_2 } }), + ); + await gateway.handleMatchEvent( + client, + event({ messageId: "m3", name: "techTimeout", data: { map_id: MAP_2 } }), + ); + + expect(processor.process).toHaveBeenCalledTimes(1); + }); + + it("re-reads a warm binding when the payload names a map it has not seen", async () => { + const { gateway, matches, processor, matchLookups } = makeGateway(); + const client = authedSocket(SERVER_A); + const MAP_3 = "33000000-0000-4000-8000-000000000033"; + + await gateway.handleMatchEvent(client, event({ messageId: "m1" })); + matches[MATCH_1].match_maps.push({ id: MAP_3 }); + await gateway.handleMatchEvent( + client, + event({ messageId: "m2", data: { match_map_id: MAP_3 } }), + ); + + expect(matchLookups()).toBe(2); + expect(processor.process).toHaveBeenCalledTimes(2); + }); + + it("drops a warm binding as soon as a re-read finds the match moved", async () => { + const { gateway, matches, processor } = makeGateway(); + const client = authedSocket(SERVER_A); + + await gateway.handleMatchEvent(client, event({ messageId: "m1" })); + matches[MATCH_1].server_id = SERVER_B; + await gateway.handleMatchEvent( + client, + event({ messageId: "m2", data: { match_map_id: MAP_2 } }), + ); + await gateway.handleMatchEvent(client, event({ messageId: "m3" })); + + expect(processor.process).toHaveBeenCalledTimes(1); + }); + it("does not cache a refusal", async () => { const { gateway, matchLookups, matches, processor } = makeGateway(); const client = authedSocket(SERVER_A); @@ -424,7 +551,24 @@ describe("MatchEventsGateway match binding", () => { expect(processor.process).toHaveBeenCalledTimes(1); }); + it("remembers the host for ten minutes", async () => { + const { gateway, cache } = makeGateway(); + + await gateway.handleMatchEvent(authedSocket(SERVER_A), event()); + + expect(cache.put).toHaveBeenCalledWith( + `match-events:last-host:${MATCH_1}`, + SERVER_A, + 600, + ); + }); + describe("legitimate edge flows", () => { + function endMatch(matches: Record, status = "Finished") { + matches[MATCH_1].server_id = null; + matches[MATCH_1].status = status; + } + it("moves with the match: the old server is refused and the new one accepted once the binding expires", async () => { const now = jest.spyOn(Date, "now").mockReturnValue(1_000_000); const { gateway, matches, processor } = makeGateway(); @@ -466,7 +610,7 @@ describe("MatchEventsGateway match binding", () => { }); it.each(["Finished", "Surrendered", "Canceled", "Forfeit", "Tie"])( - "lets the last host flush late events after the match ends (%s) and server_id is cleared", + "lets the last host flush mapStatus, chat and disconnects after the match ends (%s)", async (status) => { const { gateway, matches, processor } = makeGateway(); @@ -475,18 +619,71 @@ describe("MatchEventsGateway match binding", () => { event({ messageId: "m1" }), ); - matches[MATCH_1].server_id = null; - matches[MATCH_1].status = status; + endMatch(matches, status); + + const lateHost = authedSocket(SERVER_A); + await gateway.handleMatchEvent( + lateHost, + event({ messageId: "m2", name: "mapStatus" }), + ); + await gateway.handleMatchEvent( + lateHost, + event({ messageId: "m3", name: "chat" }), + ); + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m4", name: "player-disconnected" }), + ); + + expect(processor.process).toHaveBeenCalledTimes(4); + }, + ); + + it.each(["surrender", "score", "restoreRound", "techTimeout", "kill"])( + "refuses %s from the last host once the match has ended", + async (name) => { + const { gateway, matches, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + endMatch(matches); await gateway.handleMatchEvent( authedSocket(SERVER_A), - event({ messageId: "m2", data: { match_map_id: MAP_1 } }), + event({ messageId: "m2", name }), ); - expect(processor.process).toHaveBeenCalledTimes(2); + expect(processor.process).toHaveBeenCalledTimes(1); }, ); + it("keeps refusing result changes on a connection whose ended binding is cached", async () => { + const { gateway, matches, processor, matchLookups } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + endMatch(matches); + + const lateHost = authedSocket(SERVER_A); + await gateway.handleMatchEvent( + lateHost, + event({ messageId: "m2", name: "chat" }), + ); + await gateway.handleMatchEvent( + lateHost, + event({ messageId: "m3", name: "surrender" }), + ); + + expect(matchLookups()).toBe(2); + expect(processor.process).toHaveBeenCalledTimes(2); + }); + it("keeps an ended match closed to every other server", async () => { const { gateway, matches, processor } = makeGateway(); @@ -495,12 +692,11 @@ describe("MatchEventsGateway match binding", () => { event({ messageId: "m1" }), ); - matches[MATCH_1].server_id = null; - matches[MATCH_1].status = "Finished"; + endMatch(matches); await gateway.handleMatchEvent( authedSocket(SERVER_B), - event({ messageId: "m2" }), + event({ messageId: "m2", name: "chat" }), ); expect(processor.process).toHaveBeenCalledTimes(1); @@ -517,7 +713,10 @@ describe("MatchEventsGateway match binding", () => { }, }); - await gateway.handleMatchEvent(authedSocket(SERVER_A), event()); + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ name: "chat" }), + ); expect(processor.process).not.toHaveBeenCalled(); }); @@ -537,27 +736,28 @@ describe("MatchEventsGateway match binding", () => { ); expect(processor.process).toHaveBeenCalledTimes(2); - matches[MATCH_1].server_id = null; - matches[MATCH_1].status = "Finished"; + endMatch(matches); await gateway.handleMatchEvent( authedSocket(SERVER_A), - event({ messageId: "m3" }), + event({ messageId: "m3", name: "chat" }), ); expect(processor.process).toHaveBeenCalledTimes(2); await gateway.handleMatchEvent( authedSocket(SERVER_B), - event({ messageId: "m4" }), + event({ messageId: "m4", name: "chat" }), ); expect(processor.process).toHaveBeenCalledTimes(3); }); it("learns the new host from a refused lookup even if it never posted before the match ended", async () => { const { gateway, matches, processor } = makeGateway(); - const oldServer = authedSocket(SERVER_A); - await gateway.handleMatchEvent(oldServer, event({ messageId: "m1" })); + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); matches[MATCH_1].server_id = SERVER_B; await gateway.handleMatchEvent( @@ -565,12 +765,11 @@ describe("MatchEventsGateway match binding", () => { event({ messageId: "m2" }), ); - matches[MATCH_1].server_id = null; - matches[MATCH_1].status = "Canceled"; + endMatch(matches, "Canceled"); await gateway.handleMatchEvent( authedSocket(SERVER_A), - event({ messageId: "m3" }), + event({ messageId: "m3", name: "chat" }), ); expect(processor.process).toHaveBeenCalledTimes(1); @@ -584,12 +783,11 @@ describe("MatchEventsGateway match binding", () => { event({ messageId: "m1" }), ); - matches[MATCH_1].server_id = null; - matches[MATCH_1].status = "WaitingForServer"; + endMatch(matches, "WaitingForServer"); await gateway.handleMatchEvent( authedSocket(SERVER_A), - event({ messageId: "m2" }), + event({ messageId: "m2", name: "chat" }), ); expect(processor.process).toHaveBeenCalledTimes(1); diff --git a/src/matches/match-events.gateway.ts b/src/matches/match-events.gateway.ts index ec05ba445..70cdf0fff 100644 --- a/src/matches/match-events.gateway.ts +++ b/src/matches/match-events.gateway.ts @@ -19,9 +19,11 @@ import type { e_match_status_enum } from "../../generated"; type MatchBinding = { expiresAt: number; mapIds: Set; + ended: boolean; }; export type FiveStackGameServerWebSocketClient = WebSocket.WebSocket & { + authentication?: Promise; authenticated?: boolean; serverId?: string; matchBindings?: Map; @@ -36,7 +38,7 @@ export class MatchEventsGateway { // and kill, which arrive many times a second during a round. private static readonly BINDING_TTL_MS = 5 * 1000; - private static readonly LAST_HOST_TTL_SECONDS = 60 * 60; + private static readonly LAST_HOST_TTL_SECONDS = 10 * 60; private static readonly TERMINAL_STATUSES: readonly e_match_status_enum[] = [ "Finished", @@ -46,6 +48,22 @@ export class MatchEventsGateway { "Surrendered", ]; + // What a server still sends once its match is over: the map Finished that + // follows a surrender, post-match chat and players leaving. Anything that + // could rewrite the result is refused. + private static readonly EVENTS_AFTER_MATCH_END: readonly string[] = [ + "mapStatus", + "chat", + "player-disconnected", + ]; + + // Processors write straight to the match map a payload names, so each of + // these has to be a map of the match the event is for. + private static readonly PAYLOAD_MAP_ID_KEYS: readonly string[] = [ + "match_map_id", + "map_id", + ]; + constructor( private readonly logger: Logger, private readonly moduleRef: ModuleRef, @@ -53,73 +71,13 @@ export class MatchEventsGateway { private readonly cache: CacheService, ) {} - async handleConnection( + // Nest binds the message handlers without waiting for this, so handlers wait + // on the promise instead of seeing a half-authenticated socket. + handleConnection( @ConnectedSocket() client: FiveStackGameServerWebSocketClient, request: Request, ) { - try { - const authHeader = request.headers.authorization; - - if (!authHeader || !authHeader.startsWith("Basic ")) { - this.logger.warn("game server connection rejected: missing auth", { - ip: request.headers["cf-connecting-ip"], - }); - client.terminate(); - return; - } - - const base64Credentials = authHeader.split(" ").at(1); - if (!base64Credentials) { - this.logger.warn("game server connection rejected: malformed auth", { - ip: request.headers["cf-connecting-ip"], - }); - client.terminate(); - return; - } - - const decoded = Buffer.from(base64Credentials, "base64").toString(); - const colonIndex = decoded.indexOf(":"); - if (colonIndex === -1) { - this.logger.warn( - "game server connection rejected: invalid credentials format", - { - ip: request.headers["cf-connecting-ip"], - }, - ); - client.terminate(); - return; - } - - const serverId = decoded.substring(0, colonIndex); - const apiPassword = decoded.substring(colonIndex + 1); - - const { servers_by_pk } = await this.hasura.query({ - servers_by_pk: { - __args: { - id: serverId, - }, - id: true, - api_password: true, - }, - }); - - if ( - !servers_by_pk?.id || - !timingSafeStringEqual(servers_by_pk.api_password, apiPassword) - ) { - client.terminate(); - this.logger.warn("game server auth failure", { - serverId, - ip: request.headers["cf-connecting-ip"], - }); - return; - } - - client.serverId = servers_by_pk.id; - client.authenticated = true; - } catch { - client.terminate(); - } + client.authentication = this.authenticate(client, request); } @SubscribeMessage("events") @@ -136,6 +94,8 @@ export class MatchEventsGateway { }; }, ) { + await client.authentication; + if (!client.authenticated || !client.serverId) { return; } @@ -143,7 +103,7 @@ export class MatchEventsGateway { const { matchId, mapId, messageId } = message; const { data, event } = message.data; - if (!(await this.isHostedBy(client, matchId, data?.match_map_id))) { + if (!(await this.isHostedBy(client, matchId, event, data))) { this.logger.warn( "game server event refused: match is not hosted by this server", { @@ -198,63 +158,134 @@ export class MatchEventsGateway { return messageId; } - // Stats and round events name their match map in the payload, so the map has - // to belong to the match too or a server could write into another match - // through its own. + private async authenticate( + client: FiveStackGameServerWebSocketClient, + request: Request, + ) { + try { + const authHeader = request.headers.authorization; + + if (!authHeader || !authHeader.startsWith("Basic ")) { + this.logger.warn("game server connection rejected: missing auth", { + ip: request.headers["cf-connecting-ip"], + }); + client.terminate(); + return; + } + + const base64Credentials = authHeader.split(" ").at(1); + if (!base64Credentials) { + this.logger.warn("game server connection rejected: malformed auth", { + ip: request.headers["cf-connecting-ip"], + }); + client.terminate(); + return; + } + + const decoded = Buffer.from(base64Credentials, "base64").toString(); + const colonIndex = decoded.indexOf(":"); + if (colonIndex === -1) { + this.logger.warn( + "game server connection rejected: invalid credentials format", + { + ip: request.headers["cf-connecting-ip"], + }, + ); + client.terminate(); + return; + } + + const serverId = decoded.substring(0, colonIndex); + const apiPassword = decoded.substring(colonIndex + 1); + + const { servers_by_pk } = await this.hasura.query({ + servers_by_pk: { + __args: { + id: serverId, + }, + id: true, + api_password: true, + }, + }); + + if ( + !servers_by_pk?.id || + !timingSafeStringEqual(servers_by_pk.api_password, apiPassword) + ) { + client.terminate(); + this.logger.warn("game server auth failure", { + serverId, + ip: request.headers["cf-connecting-ip"], + }); + return; + } + + client.serverId = servers_by_pk.id; + client.authenticated = true; + } catch { + client.terminate(); + } + } + private async isHostedBy( client: FiveStackGameServerWebSocketClient, matchId: unknown, - matchMapId: unknown, + event: string, + data: Record | undefined, ): Promise { if (typeof matchId !== "string" || !validate(matchId)) { return false; } - if ( - matchMapId !== undefined && - matchMapId !== null && - typeof matchMapId !== "string" - ) { - return false; + const payloadMapIds: string[] = []; + for (const key of MatchEventsGateway.PAYLOAD_MAP_ID_KEYS) { + const value = data?.[key]; + if (value === undefined || value === null) { + continue; + } + if (typeof value !== "string") { + return false; + } + payloadMapIds.push(value.toLowerCase()); } const bindingKey = matchId.toLowerCase(); - const mapId = - typeof matchMapId === "string" ? matchMapId.toLowerCase() : undefined; - const cached = client.matchBindings?.get(bindingKey); + let binding = client.matchBindings?.get(bindingKey); if ( - cached && - cached.expiresAt > Date.now() && - (!mapId || cached.mapIds.has(mapId)) + !binding || + binding.expiresAt <= Date.now() || + !payloadMapIds.every((id) => binding.mapIds.has(id)) ) { - return true; - } + binding = await this.lookupBinding(client.serverId, bindingKey); - const mapIds = await this.hostedMatchMapIds(client.serverId, bindingKey); + client.matchBindings ??= new Map(); - client.matchBindings ??= new Map(); + if (!binding) { + client.matchBindings.delete(bindingKey); + return false; + } - if (!mapIds) { - client.matchBindings.delete(bindingKey); - return false; + client.matchBindings.set(bindingKey, binding); } - client.matchBindings.set(bindingKey, { - expiresAt: Date.now() + MatchEventsGateway.BINDING_TTL_MS, - mapIds, - }); + if ( + binding.ended && + !MatchEventsGateway.EVENTS_AFTER_MATCH_END.includes(event) + ) { + return false; + } - return !mapId || mapIds.has(mapId); + return payloadMapIds.every((id) => binding.mapIds.has(id)); } // Ending a match clears its server_id while the server is still flushing - // late events (the map Finished that follows a surrender, chat, disconnects, - // retries), so an ended match stays open to the last server seen hosting it. - private async hostedMatchMapIds( + // late events, so an ended match stays open to the last server seen hosting + // it, for EVENTS_AFTER_MATCH_END only. + private async lookupBinding( serverId: string, matchId: string, - ): Promise | null> { + ): Promise { const { matches_by_pk: match } = await this.hasura.query({ matches_by_pk: { __args: { @@ -274,25 +305,33 @@ export class MatchEventsGateway { const lastHostKey = MatchEventsGateway.lastHostKey(matchId); - let hosted: boolean; + let ended = false; if (match.server_id) { await this.cache.put( lastHostKey, match.server_id, MatchEventsGateway.LAST_HOST_TTL_SECONDS, ); - hosted = match.server_id === serverId; + + if (match.server_id !== serverId) { + return null; + } } else { - hosted = - MatchEventsGateway.TERMINAL_STATUSES.includes(match.status) && - (await this.cache.get(lastHostKey)) === serverId; - } + if ( + !MatchEventsGateway.TERMINAL_STATUSES.includes(match.status) || + (await this.cache.get(lastHostKey)) !== serverId + ) { + return null; + } - if (!hosted) { - return null; + ended = true; } - return new Set(match.match_maps.map(({ id }) => id)); + return { + expiresAt: Date.now() + MatchEventsGateway.BINDING_TTL_MS, + mapIds: new Set(match.match_maps.map(({ id }) => id)), + ended, + }; } private static lastHostKey(matchId: string) { From 0afdb164f10df2233a916ee7863a4f98908f102a Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 19:10:38 -0400 Subject: [PATCH 3/3] bug: hold ended matches to the post-match rules before server_id clears - A match now counts as ended from its status, not from server_id being cleared. The controller clears server_id only after other side effects, and until then the host could still send surrender or score into a Finished match. - After a match ends, mapStatus is accepted only for WaitingForTV, UploadingDemo and Finished. Any other status could reopen or re-award a map, such as the unplayed third map of a 2-0 Bo3. - The auth lookup gives up after 10s and terminates the socket, so a stalled Hasura cannot make an unauthenticated socket hold its messages in memory. --- src/matches/match-events.gateway.spec.ts | 99 +++++++++++++++++++++++- src/matches/match-events.gateway.ts | 85 ++++++++++++++------ 2 files changed, 159 insertions(+), 25 deletions(-) diff --git a/src/matches/match-events.gateway.spec.ts b/src/matches/match-events.gateway.spec.ts index 09e58a57d..ac2feccaa 100644 --- a/src/matches/match-events.gateway.spec.ts +++ b/src/matches/match-events.gateway.spec.ts @@ -330,6 +330,31 @@ describe("MatchEventsGateway unauthenticated clients", () => { expect(processor.process).toHaveBeenCalledTimes(1); }); + it("gives up on an auth lookup that stalls and drops the held events", async () => { + jest.useFakeTimers(); + try { + const { gateway, processor } = makeGateway({ + serverQuery: () => new Promise(() => {}), + }); + const client = socket(); + + gateway.handleConnection(client, basic(SERVER_A, "password-a")); + const held = gateway.handleMatchEvent(client, event()); + + await jest.advanceTimersByTimeAsync(9_999); + expect(client.terminate).not.toHaveBeenCalled(); + + await jest.advanceTimersByTimeAsync(1); + + await expect(held).resolves.toBeUndefined(); + expect(client.terminate).toHaveBeenCalledTimes(1); + expect(client.authenticated).toBeFalsy(); + expect(processor.process).not.toHaveBeenCalled(); + } finally { + jest.useRealTimers(); + } + }); + it("ignores events on a socket that never authenticated", async () => { const { gateway, processor } = makeGateway(); @@ -624,7 +649,11 @@ describe("MatchEventsGateway match binding", () => { const lateHost = authedSocket(SERVER_A); await gateway.handleMatchEvent( lateHost, - event({ messageId: "m2", name: "mapStatus" }), + event({ + messageId: "m2", + name: "mapStatus", + data: { status: "Finished" }, + }), ); await gateway.handleMatchEvent( lateHost, @@ -639,6 +668,74 @@ describe("MatchEventsGateway match binding", () => { }, ); + it.each(["WaitingForTV", "UploadingDemo", "Finished"])( + "lets the last host report the end-of-map status %s after the match ends", + async (mapStatus) => { + const { gateway, matches, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + endMatch(matches); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ + messageId: "m2", + name: "mapStatus", + data: { status: mapStatus }, + }), + ); + + expect(processor.process).toHaveBeenCalledTimes(2); + }, + ); + + it.each([["Surrendered"], ["Live"], ["Paused"], ["Knife"], [undefined]])( + "refuses a map status of %p from the last host once the match has ended", + async (mapStatus) => { + const { gateway, matches, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + endMatch(matches); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ + messageId: "m2", + name: "mapStatus", + data: { status: mapStatus, winning_lineup_id: "lineup-2" }, + }), + ); + + expect(processor.process).toHaveBeenCalledTimes(1); + }, + ); + + it("treats a match as ended by its status even before server_id is cleared", async () => { + const { gateway, matches, processor } = makeGateway(); + + matches[MATCH_1].status = "Finished"; + + const host = authedSocket(SERVER_A); + await gateway.handleMatchEvent( + host, + event({ messageId: "m1", name: "surrender" }), + ); + await gateway.handleMatchEvent( + host, + event({ messageId: "m2", name: "chat" }), + ); + + expect(processor.process).toHaveBeenCalledTimes(1); + }); + it.each(["surrender", "score", "restoreRound", "techTimeout", "kill"])( "refuses %s from the last host once the match has ended", async (name) => { diff --git a/src/matches/match-events.gateway.ts b/src/matches/match-events.gateway.ts index 70cdf0fff..f557a86bc 100644 --- a/src/matches/match-events.gateway.ts +++ b/src/matches/match-events.gateway.ts @@ -48,15 +48,23 @@ export class MatchEventsGateway { "Surrendered", ]; - // What a server still sends once its match is over: the map Finished that - // follows a surrender, post-match chat and players leaving. Anything that - // could rewrite the result is refused. + // What a server still sends once its match is over: post-match chat, players + // leaving and the end-of-map statuses (the map Finished that follows a + // surrender). Anything that could rewrite the result is refused, including a + // map status that would reopen or re-award a map. private static readonly EVENTS_AFTER_MATCH_END: readonly string[] = [ - "mapStatus", "chat", "player-disconnected", ]; + private static readonly MAP_STATUSES_AFTER_MATCH_END: readonly string[] = [ + "WaitingForTV", + "UploadingDemo", + "Finished", + ]; + + private static readonly AUTH_TIMEOUT_MS = 10 * 1000; + // Processors write straight to the match map a payload names, so each of // these has to be a map of the match the event is for. private static readonly PAYLOAD_MAP_ID_KEYS: readonly string[] = [ @@ -198,15 +206,17 @@ export class MatchEventsGateway { const serverId = decoded.substring(0, colonIndex); const apiPassword = decoded.substring(colonIndex + 1); - const { servers_by_pk } = await this.hasura.query({ - servers_by_pk: { - __args: { - id: serverId, + const { servers_by_pk } = await this.withAuthTimeout( + this.hasura.query({ + servers_by_pk: { + __args: { + id: serverId, + }, + id: true, + api_password: true, }, - id: true, - api_password: true, - }, - }); + }), + ); if ( !servers_by_pk?.id || @@ -227,6 +237,25 @@ export class MatchEventsGateway { } } + // Events wait on the auth lookup, so a stalled one would hold every message + // an unauthenticated socket sends in memory. + private async withAuthTimeout(lookup: Promise): Promise { + let timer: NodeJS.Timeout; + try { + return await Promise.race([ + lookup, + new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error("game server auth timed out")), + MatchEventsGateway.AUTH_TIMEOUT_MS, + ); + }), + ]); + } finally { + clearTimeout(timer); + } + } + private async isHostedBy( client: FiveStackGameServerWebSocketClient, matchId: unknown, @@ -271,7 +300,7 @@ export class MatchEventsGateway { if ( binding.ended && - !MatchEventsGateway.EVENTS_AFTER_MATCH_END.includes(event) + !MatchEventsGateway.isAllowedAfterMatchEnd(event, data) ) { return false; } @@ -281,7 +310,7 @@ export class MatchEventsGateway { // Ending a match clears its server_id while the server is still flushing // late events, so an ended match stays open to the last server seen hosting - // it, for EVENTS_AFTER_MATCH_END only. + // it, for what isAllowedAfterMatchEnd lets through only. private async lookupBinding( serverId: string, matchId: string, @@ -305,7 +334,8 @@ export class MatchEventsGateway { const lastHostKey = MatchEventsGateway.lastHostKey(matchId); - let ended = false; + const ended = MatchEventsGateway.TERMINAL_STATUSES.includes(match.status); + if (match.server_id) { await this.cache.put( lastHostKey, @@ -316,15 +346,8 @@ export class MatchEventsGateway { if (match.server_id !== serverId) { return null; } - } else { - if ( - !MatchEventsGateway.TERMINAL_STATUSES.includes(match.status) || - (await this.cache.get(lastHostKey)) !== serverId - ) { - return null; - } - - ended = true; + } else if (!ended || (await this.cache.get(lastHostKey)) !== serverId) { + return null; } return { @@ -334,6 +357,20 @@ export class MatchEventsGateway { }; } + private static isAllowedAfterMatchEnd( + event: string, + data: Record | undefined, + ) { + if (event === "mapStatus") { + return ( + typeof data?.status === "string" && + MatchEventsGateway.MAP_STATUSES_AFTER_MATCH_END.includes(data.status) + ); + } + + return MatchEventsGateway.EVENTS_AFTER_MATCH_END.includes(event); + } + private static lastHostKey(matchId: string) { return `match-events:last-host:${matchId}`; }