diff --git a/src/matches/events/CaptainEvent.spec.ts b/src/matches/events/CaptainEvent.spec.ts new file mode 100644 index 00000000..4df5f946 --- /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 d546c917..c266e8a9 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 b675b84f..ac2fecca 100644 --- a/src/matches/match-events.gateway.spec.ts +++ b/src/matches/match-events.gateway.spec.ts @@ -1,16 +1,99 @@ -// Isolate the gateway from its heavy DI imports; we only exercise the -// dedup/processing ordering logic in handleMatchEvent. -jest.mock("./events", () => ({ MatchEvents: { testEvent: class {} } })); +// Isolate the gateway from its heavy DI imports. +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 {} })); -import { MatchEventsGateway } from "./match-events.gateway"; +import { + FiveStackGameServerWebSocketClient, + MatchEventsGateway, +} from "./match-events.gateway"; + +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 }) { +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 +108,786 @@ 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; +} + +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: "matchId" in overrides ? overrides.matchId : MATCH_1, + messageId: overrides.messageId ?? "msg-1", + data: { + event: overrides.name ?? "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 dedupPut = cache.put.mock.calls.findIndex(([key]) => + key.endsWith(":msg-1"), + ); + expect(dedupPut).toBeGreaterThanOrEqual(0); expect(processor.process.mock.invocationCallOrder[0]).toBeLessThan( - cache.put.mock.invocationCallOrder[0], + cache.put.mock.invocationCallOrder[dedupPut], ); 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(); + }); +}); + +describe("MatchEventsGateway.handleConnection", () => { + it("authenticates a server with the right password", async () => { + const { gateway } = makeGateway(); + const client = socket(); + + await connect(gateway, 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 connect(gateway, client, request); + + 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 connect(gateway, client, basic("not-a-uuid", "x")); + + expect(client.terminate).toHaveBeenCalledTimes(1); + expect(client.authenticated).toBeFalsy(); + }); +}); + +describe("MatchEventsGateway unauthenticated clients", () => { + function pendingAuth() { + let resolveServer: (value: unknown) => void; + 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(); + + gateway.handleConnection(client, basic(SERVER_A, "wrong")); + + const raced = gateway.handleMatchEvent(client, event()); + expect(processor.process).not.toHaveBeenCalled(); + + finishAuth(); + + await expect(raced).resolves.toBeUndefined(); + 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("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("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(); + + 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.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({ name: "techTimeout", data: { [key]: MAP_2 } }), + ); + + expect(processor.process).not.toHaveBeenCalled(); + }, + ); + + it("accepts a techTimeout for its own map", async () => { + const { gateway, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ name: "techTimeout", data: { map_id: MAP_1 } }), + ); + + 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(); + + 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(); - expect(cache.put).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 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); + + 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); + }); + + 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(); + 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 mapStatus, chat and disconnects after the match ends (%s)", + async (status) => { + const { gateway, matches, processor } = makeGateway(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + endMatch(matches, status); + + const lateHost = authedSocket(SERVER_A); + await gateway.handleMatchEvent( + lateHost, + event({ + messageId: "m2", + name: "mapStatus", + data: { status: "Finished" }, + }), + ); + 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(["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) => { + 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 }), + ); + + 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(); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m1" }), + ); + + endMatch(matches); + + await gateway.handleMatchEvent( + authedSocket(SERVER_B), + event({ messageId: "m2", name: "chat" }), + ); + + 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({ name: "chat" }), + ); + + 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); + + endMatch(matches); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m3", name: "chat" }), + ); + expect(processor.process).toHaveBeenCalledTimes(2); + + await gateway.handleMatchEvent( + authedSocket(SERVER_B), + 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(); + + 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" }), + ); + + endMatch(matches, "Canceled"); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + event({ messageId: "m3", name: "chat" }), + ); + + 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" }), + ); + + endMatch(matches, "WaitingForServer"); + + await gateway.handleMatchEvent( + authedSocket(SERVER_A), + 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 d2e03475..f557a86b 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,64 @@ 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; + ended: boolean; +}; export type FiveStackGameServerWebSocketClient = WebSocket.WebSocket & { - id: string; - matchId: string; + authentication?: Promise; + 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 = 10 * 60; + + private static readonly TERMINAL_STATUSES: readonly e_match_status_enum[] = [ + "Finished", + "Canceled", + "Forfeit", + "Tie", + "Surrendered", + ]; + + // 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[] = [ + "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[] = [ + "match_map_id", + "map_id", + ]; + constructor( private readonly logger: Logger, private readonly moduleRef: ModuleRef, @@ -30,8 +79,95 @@ export class MatchEventsGateway { private readonly cache: CacheService, ) {} - async handleConnection( - @ConnectedSocket() client: WebSocket.WebSocket, + // 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, + ) { + client.authentication = this.authenticate(client, request); + } + + @SubscribeMessage("events") + async handleMatchEvent( + @ConnectedSocket() client: FiveStackGameServerWebSocketClient, + @MessageBody() + message: { + mapId?: string; + matchId: string; + messageId: string; + data: { + event: string; + data: Record; + }; + }, + ) { + await client.authentication; + + if (!client.authenticated || !client.serverId) { + return; + } + + const { matchId, mapId, messageId } = message; + const { data, event } = message.data; + + if (!(await this.isHostedBy(client, matchId, event, data))) { + 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}` + : `match-events:${matchId}:${messageId}`; + + if (await this.cache.has(cacheKey)) { + return messageId; + } + + const Processor = MatchEvents[event as keyof typeof MatchEvents]; + + if (!Processor) { + this.logger.warn("unable to find event handler", event); + return messageId; + } + + const processor = + await this.moduleRef.resolve>(Processor); + + processor.setData(matchId, data); + + try { + await processor.process(); + } catch (error) { + // Do NOT write the dedup entry on failure: leave the key absent so the + // game server's redelivery of this messageId is reprocessed instead of + // being silently swallowed by the dedup short-circuit for its TTL. + this.logger.error( + `[${matchId}] error processing game event ${event} (messageId=${messageId}): ${ + (error as Error)?.message + }`, + (error as Error)?.stack, + ); + throw error; + } + + // Mark processed only after success. + await this.cache.put(cacheKey, true, 10); + + return messageId; + } + + private async authenticate( + client: FiveStackGameServerWebSocketClient, request: Request, ) { try { @@ -41,7 +177,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 +186,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,90 +199,179 @@ export class MatchEventsGateway { ip: request.headers["cf-connecting-ip"], }, ); - client.close(); + 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, + 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 (!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( - @MessageBody() - message: { - mapId?: string; - matchId: string; - messageId: string; - data: { - event: string; - data: Record; - }; - }, - ) { - const { matchId, mapId, messageId } = message; + // 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); + } + } - const cacheKey = mapId - ? `match-events:${matchId}:${mapId}:${messageId}` - : `match-events:${matchId}:${messageId}`; + private async isHostedBy( + client: FiveStackGameServerWebSocketClient, + matchId: unknown, + event: string, + data: Record | undefined, + ): Promise { + if (typeof matchId !== "string" || !validate(matchId)) { + return false; + } - if (await this.cache.has(cacheKey)) { - return messageId; + 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 { data, event } = message.data; + const bindingKey = matchId.toLowerCase(); - const Processor = MatchEvents[event as keyof typeof MatchEvents]; + let binding = client.matchBindings?.get(bindingKey); + if ( + !binding || + binding.expiresAt <= Date.now() || + !payloadMapIds.every((id) => binding.mapIds.has(id)) + ) { + binding = await this.lookupBinding(client.serverId, bindingKey); - if (!Processor) { - this.logger.warn("unable to find event handler", event); - return messageId; + client.matchBindings ??= new Map(); + + if (!binding) { + client.matchBindings.delete(bindingKey); + return false; + } + + client.matchBindings.set(bindingKey, binding); } - const processor = - await this.moduleRef.resolve>(Processor); + if ( + binding.ended && + !MatchEventsGateway.isAllowedAfterMatchEnd(event, data) + ) { + return false; + } - processor.setData(matchId, data); + return payloadMapIds.every((id) => binding.mapIds.has(id)); + } - try { - await processor.process(); - } catch (error) { - // Do NOT write the dedup entry on failure: leave the key absent so the - // game server's redelivery of this messageId is reprocessed instead of - // being silently swallowed by the dedup short-circuit for its TTL. - this.logger.error( - `[${matchId}] error processing game event ${event} (messageId=${messageId}): ${ - (error as Error)?.message - }`, - (error as Error)?.stack, + // 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 what isAllowedAfterMatchEnd lets through only. + private async lookupBinding( + serverId: string, + matchId: string, + ): Promise { + 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); + + const ended = MatchEventsGateway.TERMINAL_STATUSES.includes(match.status); + + if (match.server_id) { + await this.cache.put( + lastHostKey, + match.server_id, + MatchEventsGateway.LAST_HOST_TTL_SECONDS, ); - throw error; + + if (match.server_id !== serverId) { + return null; + } + } else if (!ended || (await this.cache.get(lastHostKey)) !== serverId) { + return null; } - // Mark processed only after success. - await this.cache.put(cacheKey, true, 10); + return { + expiresAt: Date.now() + MatchEventsGateway.BINDING_TTL_MS, + mapIds: new Set(match.match_maps.map(({ id }) => id)), + ended, + }; + } - return messageId; + 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}`; } }