diff --git a/package.json b/package.json index f6c12de62..1be0a064c 100644 --- a/package.json +++ b/package.json @@ -93,6 +93,7 @@ "@nestjs/schematics": "^11.0.5", "@nestjs/testing": "^11.1.3", "@testcontainers/postgresql": "^12.0.4", + "testcontainers": "^12.0.4", "@types/archiver": "^6.0.2", "@types/express": "^5.0.0", "@types/express-session": "^1.18.0", diff --git a/src/configs/redis.ts b/src/configs/redis.ts index 22345f5a6..31deecbb9 100644 --- a/src/configs/redis.ts +++ b/src/configs/redis.ts @@ -13,6 +13,20 @@ export default (): { : undefined, password: process.env.REDIS_PASSWORD, }, + // Playcast relay traffic: its own socket so fragment payloads never queue + // behind session lookups, and failing fast so an outage answers viewers + // and game servers with a 503 instead of holding their requests open. + relay: { + db: 1, + host: process.env.REDIS_HOST || "redis", + port: process.env.REDIS_SERVICE_PORT + ? parseInt(process.env.REDIS_SERVICE_PORT) + : undefined, + password: process.env.REDIS_PASSWORD, + enableOfflineQueue: false, + maxRetriesPerRequest: 1, + commandTimeout: 2000, + }, sub: { db: 1, host: process.env.REDIS_HOST || "redis", diff --git a/src/configs/types/RedisConfig.ts b/src/configs/types/RedisConfig.ts index 2e418f804..ca02b5d78 100644 --- a/src/configs/types/RedisConfig.ts +++ b/src/configs/types/RedisConfig.ts @@ -6,6 +6,9 @@ export type RedisConfig = { host: string; port: number; password: string; + enableOfflineQueue?: boolean; + maxRetriesPerRequest?: number | null; + commandTimeout?: number; } >; }; diff --git a/src/matches/jobs/StopMatchBroadcast.ts b/src/matches/jobs/StopMatchBroadcast.ts index c0571479a..f3e7e7225 100644 --- a/src/matches/jobs/StopMatchBroadcast.ts +++ b/src/matches/jobs/StopMatchBroadcast.ts @@ -30,7 +30,7 @@ export class StopMatchBroadcast extends WorkerHost { return; } - this.matchRelay.removeBroadcast(matchId); + await this.matchRelay.removeBroadcast(matchId); if (!(await this.gameStreamer.stopLiveIfRunning(matchId))) { return; diff --git a/src/matches/match-relay/match-relay.controller.ts b/src/matches/match-relay/match-relay.controller.ts index feb681b37..d9d4faa3c 100644 --- a/src/matches/match-relay/match-relay.controller.ts +++ b/src/matches/match-relay/match-relay.controller.ts @@ -1,92 +1,84 @@ -import { Controller, Get, Post, Req, Res, Param } from "@nestjs/common"; +import { Controller, Get, Post, Req, Res, Param, Logger } from "@nestjs/common"; import { Request, Response } from "express"; import { MatchRelayService } from "./match-relay.service"; +import { FragmentField } from "./types/fragment.types"; @Controller("match-relay/:id") export class MatchRelayController { - constructor(private readonly matchRelayService: MatchRelayService) {} + constructor( + private readonly logger: Logger, + private readonly matchRelayService: MatchRelayService, + ) {} @Get("sync") - public handleSyncGet( + public async handleSyncGet( @Param("id") matchId: string, @Req() request: Request, @Res() response: Response, ) { - this.matchRelayService.getSyncInfo(request, response, matchId); + await this.relay(matchId, response, () => + this.matchRelayService.getSyncInfo(request, response, matchId), + ); } @Get(":fragment/start") - public handleGetStart( + public async handleGetStart( @Param("id") matchId: string, @Param("fragment") fragment: string, @Res() response: Response, ) { - this.matchRelayService.getStart(response, matchId, parseInt(fragment)); + await this.relay(matchId, response, () => + this.matchRelayService.getStart(response, matchId, parseInt(fragment)), + ); } @Get(":fragment/full") - public handleGetFull( + public async handleGetFull( @Param("id") matchId: string, @Param("fragment") fragment: string, @Res() response: Response, ) { - this.matchRelayService.getFragment( - response, - matchId, - parseInt(fragment), - "full", - ); + await this.getFragment(response, matchId, fragment, "full"); } @Get(":fragment/delta") - public handleGetDelta( + public async handleGetDelta( @Param("id") matchId: string, @Param("fragment") fragment: string, @Res() response: Response, ) { - this.matchRelayService.getFragment( - response, - matchId, - parseInt(fragment), - "delta", - ); + await this.getFragment(response, matchId, fragment, "delta"); } @Get(":token/:fragment/start") - public handleGetStartWithToken( + public async handleGetStartWithToken( @Param("id") matchId: string, @Param("fragment") fragment: string, @Res() response: Response, ) { - this.matchRelayService.getStart(response, matchId, parseInt(fragment)); + await this.relay(matchId, response, () => + this.matchRelayService.getStart(response, matchId, parseInt(fragment)), + ); } @Get(":token/:fragment/full") - public handleGetFullWithToken( + public async handleGetFullWithToken( @Param("id") matchId: string, + @Param("token") token: string, @Param("fragment") fragment: string, @Res() response: Response, ) { - this.matchRelayService.getFragment( - response, - matchId, - parseInt(fragment), - "full", - ); + await this.getFragment(response, matchId, fragment, "full", token); } @Get(":token/:fragment/delta") - public handleGetDeltaWithToken( + public async handleGetDeltaWithToken( @Param("id") matchId: string, + @Param("token") token: string, @Param("fragment") fragment: string, @Res() response: Response, ) { - this.matchRelayService.getFragment( - response, - matchId, - parseInt(fragment), - "delta", - ); + await this.getFragment(response, matchId, fragment, "delta", token); } @Post(":token/:fragment/start") @@ -97,14 +89,7 @@ export class MatchRelayController { @Req() request: Request, @Res() response: Response, ) { - this.matchRelayService.postField( - request, - response, - token, - "start", - matchId, - parseInt(fragment), - ); + await this.postField(request, response, token, "start", matchId, fragment); } @Post(":token/:fragment/full") @@ -115,14 +100,7 @@ export class MatchRelayController { @Req() request: Request, @Res() response: Response, ) { - this.matchRelayService.postField( - request, - response, - token, - "full", - matchId, - parseInt(fragment), - ); + await this.postField(request, response, token, "full", matchId, fragment); } @Post(":token/:fragment/delta") @@ -133,13 +111,62 @@ export class MatchRelayController { @Req() request: Request, @Res() response: Response, ) { - this.matchRelayService.postField( - request, - response, - token, - "delta", - matchId, - parseInt(fragment), + await this.postField(request, response, token, "delta", matchId, fragment); + } + + private getFragment( + response: Response, + matchId: string, + fragment: string, + field: FragmentField, + token?: string, + ) { + return this.relay(matchId, response, () => + this.matchRelayService.getFragment( + response, + matchId, + parseInt(fragment), + field, + token, + ), + ); + } + + private postField( + request: Request, + response: Response, + token: string, + field: FragmentField, + matchId: string, + fragment: string, + ) { + return this.relay(matchId, response, () => + this.matchRelayService.postField( + request, + response, + token, + field, + matchId, + parseInt(fragment), + ), ); } + + private async relay( + matchId: string, + response: Response, + handle: () => Promise, + ) { + try { + await handle(); + } catch (error) { + this.logger.error( + `[${matchId}] relay request failed: ${(error as Error)?.message}`, + ); + if (!response.headersSent) { + response.writeHead(503, { "Cache-Control": "no-store" }); + } + response.end(); + } + } } diff --git a/src/matches/match-relay/match-relay.service.spec.ts b/src/matches/match-relay/match-relay.service.spec.ts deleted file mode 100644 index 2cfa5c268..000000000 --- a/src/matches/match-relay/match-relay.service.spec.ts +++ /dev/null @@ -1,205 +0,0 @@ -import { EventEmitter } from "events"; -import { MatchRelayService } from "./match-relay.service"; - -const fakeResponse = () => { - let resolveEnded: () => void; - const ended = new Promise((resolve) => { - resolveEnded = resolve; - }); - - const response = { - statusCode: undefined as number | undefined, - headers: {} as Record, - body: undefined as unknown, - ended, - writeHead(code: number, headers?: unknown) { - response.statusCode = code; - if (headers && typeof headers === "object") { - Object.assign(response.headers, headers); - } - return response; - }, - setHeader(name: string, value: unknown) { - response.headers[name] = value; - }, - end(body?: unknown) { - response.body = body; - resolveEnded(); - return response; - }, - }; - - return response; -}; - -describe("MatchRelayService", () => { - const matchId = "match-1"; - const token = "s845489096165654t8799308478907"; - - let service: MatchRelayService; - - beforeEach(() => { - service = new MatchRelayService({ - log: jest.fn(), - warn: jest.fn(), - error: jest.fn(), - } as any); - }); - - const openPost = ( - field: "start" | "full" | "delta", - fragment: number, - query: Record = {}, - ) => { - const request = Object.assign(new EventEmitter(), { query }); - const response = fakeResponse(); - - service.postField( - request as any, - response as any, - token, - field, - matchId, - fragment, - ); - - return { - response, - finish: async () => { - request.emit("data", Buffer.from(`${field}-${fragment}`)); - request.emit("end"); - await response.ended; - return response; - }, - }; - }; - - const post = ( - field: "start" | "full" | "delta", - fragment: number, - query: Record = {}, - ) => openPost(field, fragment, query).finish(); - - const sync = (query: Record = {}) => { - const response = fakeResponse(); - service.getSyncInfo({ query } as any, response as any, matchId); - return response; - }; - - const getStart = (fragment: number) => { - const response = fakeResponse(); - service.getStart(response as any, matchId, fragment); - return response; - }; - - const startBroadcastAt = async (fragment: number) => { - await post("start", fragment, { - tick: "100", - tps: "64", - map: "de_inferno", - keyframe_interval: "3", - protocol: "5", - }); - await post("full", fragment, { tick: "100" }); - await post("delta", fragment, { endtick: "292" }); - }; - - it("reports the fragment the broadcast signed up at, with numeric fields", async () => { - await startBroadcastAt(42); - - const response = sync({ fragment: "0" }); - - expect(response.statusCode).toBe(200); - expect(JSON.parse(response.body as string)).toEqual( - expect.objectContaining({ - fragment: 42, - signup_fragment: 42, - tick: 100, - endtick: 292, - maxtick: 292, - tps: 64, - keyframe_interval: 3, - map: "de_inferno", - protocol: 5, - }), - ); - }); - - // CS2 sends tps as a decimal (64.0), not an integer, so the coercion has to - // accept a fractional part or clients get tps back as a string. - it("reports how long clients keep playing once the server stops posting", async () => { - await startBroadcastAt(42); - - // Clients sit 7 fragments behind the newest and still have that one to play. - expect(service.playoutSeconds(matchId)).toBe(8 * 3); - }); - - it("has nothing to play out for a broadcast it does not hold", () => { - expect(service.playoutSeconds(matchId)).toBe(0); - }); - - it("reports a fractional tps as a number", async () => { - await post("start", 42, { - tick: "100", - tps: "64.0", - map: "de_inferno", - }); - await post("full", 42, { tick: "100" }); - await post("delta", 42, { endtick: "292" }); - - const response = sync({ fragment: "0" }); - - expect(JSON.parse(response.body as string)).toEqual( - expect.objectContaining({ tps: 64 }), - ); - }); - - it("keeps fields that are not numeric in /sync as the strings the server sent", async () => { - await post("start", 42, { - tick: "100", - tps: "64", - map: "3070284539", - protocol: "5", - }); - await post("full", 42, { tick: "100" }); - await post("delta", 42, { endtick: "292" }); - - const body = JSON.parse(sync({ fragment: "0" }).body as string); - - expect(body.map).toBe("3070284539"); - expect(body.tick).toBe(100); - }); - - it("serves start only at the fragment the broadcast signed up at", async () => { - await startBroadcastAt(42); - - expect(getStart(42).statusCode).toBe(200); - expect(getStart(0).statusCode).toBe(404); - }); - - it("moves the signup fragment when the game server re-sends start", async () => { - await startBroadcastAt(42); - await post("start", 50, { tick: "900", tps: "64", map: "de_inferno" }); - - expect(getStart(50).statusCode).toBe(200); - expect(getStart(42).statusCode).toBe(404); - }); - - it("asks for start again when a fragment arrives before any start", async () => { - const response = await post("full", 7, { tick: "100" }); - - expect(response.statusCode).toBe(205); - }); - - it("asks for start again when a fragment arrives before the start data has", async () => { - const start = openPost("start", 42, { tick: "100", tps: "64" }); - - const early = await post("full", 42, { tick: "100" }); - expect(early.statusCode).toBe(205); - - await start.finish(); - - const late = await post("full", 43, { tick: "292" }); - expect(late.statusCode).toBe(200); - }); -}); diff --git a/src/matches/match-relay/match-relay.service.ts b/src/matches/match-relay/match-relay.service.ts index 5b488f3bd..66220bc82 100644 --- a/src/matches/match-relay/match-relay.service.ts +++ b/src/matches/match-relay/match-relay.service.ts @@ -2,13 +2,18 @@ import zlib from "zlib"; import { promisify } from "util"; import { Request, Response } from "express"; import { Injectable, Logger } from "@nestjs/common"; -import { - Fragment, - StartFieldData, - FullFieldData, - DeltaFieldData, -} from "./types/fragment.types"; - +import { Redis } from "ioredis"; +import { RedisManagerService } from "../../redis/redis-manager/redis-manager.service"; +import { FieldMeta, FragmentField } from "./types/fragment.types"; + +// Broadcast state lives in redis rather than in the process, so an api restart +// or rollout does not drop every live broadcast and more than one replica can +// serve the relay. +// +// Everything a broadcast stores is keyed by its token as well as the match. A +// new token (a new map, or the server restarting the broadcast) starts over at +// low fragment numbers, so keying by match alone would let a half-finished +// cleanup serve one map's fragment as another's. @Injectable() export class MatchRelayService { private static readonly NUMERIC_QUERY_FIELDS = [ @@ -23,29 +28,44 @@ export class MatchRelayService { // behind the newest one, and the client then plays in real time. private static readonly SYNC_LAG_FRAGMENTS = 7; + private static readonly FRAGMENT_TTL_SECONDS = 60; + + // Every post pushes this back, so it only runs out once the game server has + // been silent this long. + private static readonly BROADCAST_TTL_SECONDS = 60 * 60; + + // Fragment data expires after FRAGMENT_TTL_SECONDS; this only bounds the + // index that /sync walks. + private static readonly INDEX_WINDOW = 64; + + private static readonly IMMUTABLE = "public, max-age=31536000, immutable"; + private readonly gzip = promisify(zlib.gzip); - private readonly broadcasts: { - [key: string]: { - steamId: string; - masterCookie: string; - fragments: Map; - }; - } = {}; + private readonly redis: Redis; - constructor(private readonly logger: Logger) {} + constructor( + private readonly logger: Logger, + redisManager: RedisManagerService, + ) { + this.redis = redisManager.getConnection("relay"); + } - public removeBroadcast(matchId: string) { - delete this.broadcasts[matchId]; + public async removeBroadcast(matchId: string) { + const token = await this.currentToken(matchId); + if (token) { + await this.clearBroadcast(matchId, token); + } + await this.redis.del(MatchRelayService.tokenKey(matchId)); } // How long a client keeps playing after the game server stops posting: it is // SYNC_LAG_FRAGMENTS behind the last fragment, plus that fragment itself. // keyframe_interval is the fragment length the game server announced in start. - public playoutSeconds(matchId: string): number { - const keyframeInterval = Number( - this.broadcasts[matchId]?.fragments.get(0)?.start?.keyframe_interval, - ); + public async playoutSeconds(matchId: string): Promise { + const token = await this.currentToken(matchId); + const start = token ? await this.readStartMeta(matchId, token) : null; + const keyframeInterval = Number(start?.keyframe_interval); if (!(keyframeInterval > 0)) { return 0; @@ -54,14 +74,23 @@ export class MatchRelayService { return (MatchRelayService.SYNC_LAG_FRAGMENTS + 1) * keyframeInterval; } - public getStart(response: Response, matchId: string, fragmentIndex: number) { - const broadcast = this.broadcasts[matchId]; - const startFragment = broadcast?.fragments.get(0); - - if ( - startFragment == null || - startFragment.start?.signup_fragment != fragmentIndex - ) { + public async getStart( + response: Response, + matchId: string, + fragmentIndex: number, + ) { + const token = await this.currentToken(matchId); + const [meta, data] = token + ? await Promise.all([ + this.readStartMeta(matchId, token), + this.redis.hgetBuffer( + MatchRelayService.startKey(matchId, token), + "data", + ), + ]) + : [null, null]; + + if (meta == null || meta.signup_fragment != fragmentIndex) { return this.relayError( response, 404, @@ -69,61 +98,101 @@ export class MatchRelayService { ); } - this.serveBlob(response, startFragment, "start"); + this.serveBlob(response, data, meta); } - public getFragment( + // A token in the url pins the request to one broadcast, which is what makes + // the fragment safe to cache: indexes start over when a new map starts a new + // broadcast, so the same url without a token can name different data. + public async getFragment( response: Response, matchId: string, fragmentIndex: number, - field: "start" | "full" | "delta", + field: FragmentField, + token?: string, ) { - const broadcast = this.broadcasts[matchId]; - if (!broadcast) { - this.relayError(response, 404, `broadcast not found`); + const currentToken = await this.currentToken(matchId); + + if (!currentToken) { + this.relayError(response, 404, `broadcast not found`, token); return; } - const fragment = broadcast.fragments.get(fragmentIndex); - if (!fragment) { - response.writeHead(404, "fragment not found"); - response.end(); + if (token !== undefined && token !== currentToken) { + this.relayError( + response, + 404, + `broadcast has moved on, please re-sync`, + token, + ); return; } - this.serveBlob(response, fragment, field); + const fragmentKey = MatchRelayService.fragmentKey( + matchId, + currentToken, + fragmentIndex, + ); + const [metaJson, data] = await Promise.all([ + this.redis.hget(fragmentKey, `${field}_meta`), + this.redis.hgetBuffer(fragmentKey, field), + ]); + + if (!metaJson || !data) { + this.relayError(response, 404, "fragment not found", token); + return; + } + + this.serveBlob( + response, + data, + JSON.parse(metaJson), + token ? MatchRelayService.IMMUTABLE : undefined, + ); } - public getSyncInfo( + public async getSyncInfo( request: Request, response: Response, matchId: string, - ): void { + ): Promise { const nowMs = Date.now(); response.setHeader("Cache-Control", "public, max-age=3"); response.setHeader("Expires", new Date(nowMs + 3000).toUTCString()); - const broadcast = this.broadcasts[matchId]; - if (!broadcast) { + const token = await this.currentToken(matchId); + + if (!token) { this.relayError(response, 404, `broadcast not found`); return; } - const startFragment = broadcast.fragments.get(0); + const [start, hasStartData, indexes] = await Promise.all([ + this.readStartMeta(matchId, token), + this.redis.hexists(MatchRelayService.startKey(matchId, token), "data"), + this.redis.zrange(MatchRelayService.indexKey(matchId, token), 0, -1), + ]); - if (startFragment == null || startFragment.start?.data == null) { + if (start == null || !hasStartData) { this.relayError(response, 404, `broadcast has not started yet`); return; } - let fragmentIndex: number | null = null; - const fragmentParam = request.query.fragment as string | undefined; - let fragment: Fragment | null = null; + response.setHeader("X-Broadcast-Token", token); - const maxIndex = - broadcast.fragments.size > 0 - ? Math.max(...Array.from(broadcast.fragments.keys())) - : 0; + const fragments = await this.readFragmentMetas( + matchId, + token, + indexes.map(Number), + ); + const signupFragment = start.signup_fragment || 0; + const maxIndex = fragments.length + ? fragments[fragments.length - 1].index + : 0; + + let fragmentIndex: number; + let fragment: (typeof fragments)[number] | undefined; + const fragmentParam = request.query.fragment as string | undefined; if (fragmentParam == null) { fragmentIndex = Math.max( @@ -131,30 +200,21 @@ export class MatchRelayService { maxIndex - MatchRelayService.SYNC_LAG_FRAGMENTS, ); - if ( - fragmentIndex >= 0 && - fragmentIndex >= (startFragment.start.signup_fragment || 0) - ) { - const _fragment = broadcast.fragments.get(fragmentIndex); - if (this.isSyncReady(_fragment)) { - fragment = _fragment; - } + if (fragmentIndex >= signupFragment) { + fragment = fragments.find( + (candidate) => + candidate.index === fragmentIndex && + MatchRelayService.isSyncReady(candidate), + ); } } else { - fragmentIndex = parseInt(fragmentParam); - - if (fragmentIndex < (startFragment.start?.signup_fragment || 0)) { - fragmentIndex = startFragment.start?.signup_fragment || 0; - } - - for (let i = fragmentIndex; i <= maxIndex; i++) { - const _fragment = broadcast.fragments.get(i); - if (this.isSyncReady(_fragment)) { - fragment = _fragment; - fragmentIndex = i; - break; - } - } + fragmentIndex = Math.max(parseInt(fragmentParam), signupFragment); + fragment = fragments.find( + (candidate) => + candidate.index >= fragmentIndex && + MatchRelayService.isSyncReady(candidate), + ); + fragmentIndex = fragment?.index ?? fragmentIndex; } if (!fragment) { @@ -166,128 +226,235 @@ export class MatchRelayService { return; } - response.writeHead(200, { "Content-Type": "application/json" }); - if (startFragment.start?.protocol == null) { - if (!startFragment.start) { - startFragment.start = {}; - } - startFragment.start.protocol = 5; - } - - const fragTick = fragment.full?.tick; - const fragEndtick = fragment.delta?.endtick; - const fragTimestamp = fragment.delta?.timestamp; + const lastFragment = fragments[fragments.length - 1]; + const endTick = [...fragments] + .reverse() + .find((candidate) => candidate.delta?.endtick != null)?.delta?.endtick; + response.writeHead(200, { "Content-Type": "application/json" }); response.end( JSON.stringify({ - tick: fragTick, - endtick: fragEndtick, - maxtick: this.getMatchBroadcastEndTick(broadcast.fragments), - rtdelay: (nowMs - (fragTimestamp || nowMs)) / 1000, - rcvage: - (nowMs - - (this.getLastFragment(broadcast.fragments)?.delta?.timestamp || - nowMs)) / - 1000, + tick: fragment.full?.tick, + endtick: fragment.delta?.endtick, + maxtick: endTick ?? 0, + rtdelay: (nowMs - (fragment.delta?.timestamp || nowMs)) / 1000, + rcvage: (nowMs - (lastFragment?.delta?.timestamp || nowMs)) / 1000, fragment: fragmentIndex, - signup_fragment: startFragment.start?.signup_fragment, - tps: startFragment.start?.tps, - keyframe_interval: startFragment.start?.keyframe_interval, - map: startFragment.start?.map, - protocol: startFragment.start?.protocol, + signup_fragment: start.signup_fragment, + tps: start.tps, + keyframe_interval: start.keyframe_interval, + map: start.map, + protocol: start.protocol ?? 5, }), ); } - public postField( + // Answers only once the field is stored: a 200 tells the game server it can + // move on, and a failed write has to come back as an error instead. + public async postField( request: Request, response: Response, token: string, - field: "start" | "full" | "delta", + field: FragmentField, matchId: string, fragmentIndex: number, - ): void { - const [steamId, masterCookie] = token.split("t"); - - if (!this.broadcasts[matchId]) { - this.broadcasts[matchId] = { - steamId, - masterCookie, - fragments: new Map(), - }; + ): Promise { + await this.claimBroadcast(matchId, token); + + const startKey = MatchRelayService.startKey(matchId, token); + + // 205 makes the server re-send start, so it also covers a start whose body hasn't landed. + if (field != "start" && !(await this.redis.hexists(startKey, "data"))) { + response.writeHead(205); + response.end(); + return; } - const broadcast = this.broadcasts[matchId]; + const meta: FieldMeta = {}; + Object.entries(request.query).forEach(([key, value]) => { + meta[key] = MatchRelayService.parseQueryValue(key, value); + }); + + const body = await MatchRelayService.readBody(request); - if ( - broadcast.steamId !== steamId || - broadcast.masterCookie !== masterCookie - ) { - broadcast.steamId = steamId; - broadcast.masterCookie = masterCookie; - broadcast.fragments.clear(); + let data = body; + try { + data = await this.gzip(body); + meta.gipped = true; + } catch (error) { + this.logger.error(`cannot gzip: ${error}`); + meta.gipped = false; } + meta.timestamp = Date.now(); - const signupFragment = fragmentIndex; + const write = this.redis.multi(); - if (field == "start") { - fragmentIndex = 0; + if (field === "start") { + meta.signup_fragment = fragmentIndex; + write + .del(startKey) + .hset(startKey, { meta: JSON.stringify(meta), data }) + .expire(startKey, MatchRelayService.BROADCAST_TTL_SECONDS); + } else { + const fragmentKey = MatchRelayService.fragmentKey( + matchId, + token, + fragmentIndex, + ); + const indexKey = MatchRelayService.indexKey(matchId, token); + write + .hset(fragmentKey, { + [`${field}_meta`]: JSON.stringify(meta), + [field]: data, + }) + .expire(fragmentKey, MatchRelayService.FRAGMENT_TTL_SECONDS) + .zadd(indexKey, fragmentIndex, String(fragmentIndex)) + .zremrangebyscore( + indexKey, + "-inf", + `(${fragmentIndex - MatchRelayService.INDEX_WINDOW}`, + ) + .expire(indexKey, MatchRelayService.BROADCAST_TTL_SECONDS) + .expire(startKey, MatchRelayService.BROADCAST_TTL_SECONDS); } - // 205 makes the server re-send start, so it also covers a start whose body hasn't landed. - if (field != "start" && broadcast.fragments.get(0)?.start?.data == null) { - response.writeHead(205); - response.end(); - return; + const results = await write.exec(); + if (!results) { + throw new Error("relay write was aborted"); + } + for (const [error] of results) { + if (error) { + throw error; + } } response.writeHead(200); - if (!broadcast.fragments.has(fragmentIndex)) { - broadcast.fragments.set(fragmentIndex, {}); + response.end(); + } + + private currentToken(matchId: string) { + return this.redis.get(MatchRelayService.tokenKey(matchId)); + } + + // A left-over key from the previous broadcast can only ever be read under its + // own token, so a clear that fails leaves nothing wrong behind: the keys + // simply expire. + private async claimBroadcast(matchId: string, token: string) { + const previous = await this.redis.set( + MatchRelayService.tokenKey(matchId), + token, + "EX", + MatchRelayService.BROADCAST_TTL_SECONDS, + "GET", + ); + + if (previous !== null && previous !== token) { + await this.clearBroadcast(matchId, previous).catch((error) => { + this.logger.warn( + `[${matchId}] could not clear the previous broadcast: ${ + (error as Error)?.message + }`, + ); + }); } + } + + private async clearBroadcast(matchId: string, token: string) { + const indexKey = MatchRelayService.indexKey(matchId, token); + const indexes = await this.redis.zrange(indexKey, 0, -1); - const fragment = broadcast.fragments.get(fragmentIndex)!; + await this.redis.del( + MatchRelayService.startKey(matchId, token), + indexKey, + ...indexes.map((index) => + MatchRelayService.fragmentKey(matchId, token, Number(index)), + ), + ); + } - if (fragment[field] == null) { - fragment[field] = {}; + private async readStartMeta( + matchId: string, + token: string, + ): Promise { + const meta = await this.redis.hget( + MatchRelayService.startKey(matchId, token), + "meta", + ); + return meta ? JSON.parse(meta) : null; + } + + // Only the metadata: whether a fragment is sync-ready is decided from it, and + // it is written in the same step as the data it describes. + private async readFragmentMetas( + matchId: string, + token: string, + indexes: Array, + ) { + if (indexes.length === 0) { + return []; } - if (field === "start") { - fragment.start!.signup_fragment = signupFragment; + const pipeline = this.redis.pipeline(); + for (const index of indexes) { + pipeline.hmget( + MatchRelayService.fragmentKey(matchId, token, index), + "full_meta", + "delta_meta", + ); } - Object.entries(request.query).forEach(([key, value]) => { - fragment[field]![key] = MatchRelayService.parseQueryValue(key, value); - }); + const results = (await pipeline.exec()) ?? []; + + return indexes + .map((index, position) => { + const [fullMeta, deltaMeta] = (results[position]?.[1] ?? []) as [ + string | null, + string | null, + ]; + return { + index, + full: fullMeta ? (JSON.parse(fullMeta) as FieldMeta) : undefined, + delta: deltaMeta ? (JSON.parse(deltaMeta) as FieldMeta) : undefined, + }; + }) + .filter((fragment) => fragment.full || fragment.delta); + } - const body: Buffer[] = []; - request.on("data", function (data: Buffer) { - body.push(data); - }); + private static isSyncReady(fragment: { + full?: FieldMeta; + delta?: FieldMeta; + }): boolean { + return ( + fragment.full != null && + fragment.delta != null && + (fragment.full.tick != null || fragment.delta.tick != null) && + fragment.delta.endtick != null && + fragment.delta.timestamp != null + ); + } - request.on("end", () => { - const totalBuffer = Buffer.concat(body); + private static async readBody(request: Request): Promise { + const chunks: Buffer[] = []; + for await (const chunk of request) { + chunks.push(chunk as Buffer); + } + return Buffer.concat(chunks); + } - if (fragment[field] == null) { - fragment[field] = {}; - } + private static tokenKey(matchId: string) { + return `match-relay:${matchId}:token`; + } - this.gzip(totalBuffer) - .then((compressedBlob: Buffer) => { - fragment[field]!.gipped = true; - fragment[field]!.data = compressedBlob; - }) - .catch((error: Error) => { - this.logger.error(`cannot gzip: ${error}`); - fragment[field]!.gipped = false; - fragment[field]!.data = totalBuffer; - }) - .finally(() => { - response.end(); - fragment[field]!.timestamp = Date.now(); - this.cleanupOldFragments(matchId); - }); - }); + private static startKey(matchId: string, token: string) { + return `match-relay:${matchId}:${token}:start`; + } + + private static indexKey(matchId: string, token: string) { + return `match-relay:${matchId}:${token}:fragments`; + } + + private static fragmentKey(matchId: string, token: string, index: number) { + return `match-relay:${matchId}:${token}:fragment:${index}`; } // Clients read /sync's tick, tps, etc. as JSON numbers, as Valve's reference relay sends them. @@ -303,85 +470,28 @@ export class MatchRelayService { return Number(value); } + // A token-scoped miss must never be cached: the fragment may simply not have + // arrived yet. private relayError( response: Response, code: number, explanation: string, + token?: string, ): void { - response.writeHead(code, { "X-Reason": explanation }); - response.end(); - } - - private isSyncReady(fragment: Fragment | undefined): boolean { - return ( - fragment != null && - fragment.full?.data != null && - fragment.delta?.data != null && - (fragment.full?.tick != null || fragment.delta?.tick != null) && - fragment.delta?.endtick != null && - fragment.delta?.timestamp != null - ); - } - - private cleanupOldFragments(matchId: string): void { - const broadcast = this.broadcasts[matchId]; - if (!broadcast) { - return; + const headers: Record = { "X-Reason": explanation }; + if (token) { + headers["Cache-Control"] = "no-store"; } - - const now = Date.now(); - const indicesToDelete: number[] = []; - - for (const [index, fragment] of broadcast.fragments.entries()) { - if (index === 0) { - continue; - } - if (fragment?.delta?.timestamp != null) { - const timeDiff = now - fragment.delta.timestamp; - if (timeDiff > 60000) { - indicesToDelete.push(index); - } - } - } - - for (const index of indicesToDelete) { - broadcast.fragments.delete(index); - } - } - - private getMatchBroadcastEndTick(broadcast: Map): number { - const sortedIndices = Array.from(broadcast.keys()).sort((a, b) => b - a); - for (const index of sortedIndices) { - const fragment = broadcast.get(index); - if (fragment?.delta?.endtick != null) { - return fragment.delta.endtick; - } - } - return 0; - } - - private getLastFragment( - broadcast: Map, - ): Fragment | undefined { - if (broadcast.size === 0) { - return undefined; - } - const maxIndex = Math.max(...Array.from(broadcast.keys())); - return broadcast.get(maxIndex); + response.writeHead(code, headers); + response.end(); } private serveBlob( response: Response, - fragmentRec: Fragment | undefined, - field: string, + blob: Buffer | null, + meta: FieldMeta, + cacheControl?: string, ): void { - const fieldData = fragmentRec?.[field] as - | StartFieldData - | FullFieldData - | DeltaFieldData - | undefined; - const blob = fieldData?.data; - if (!blob) { response.writeHead(404, "Field not found"); response.end(); @@ -391,9 +501,12 @@ export class MatchRelayService { const headers: { [key: string]: string } = { "Content-Type": "application/octet-stream", }; - if (fieldData.gipped) { + if (meta.gipped) { headers["Content-Encoding"] = "gzip"; } + if (cacheControl) { + headers["Cache-Control"] = cacheControl; + } response.writeHead(200, headers); response.end(blob); } diff --git a/src/matches/match-relay/types/fragment.types.ts b/src/matches/match-relay/types/fragment.types.ts index 7a34c08f1..b18c4c5ec 100644 --- a/src/matches/match-relay/types/fragment.types.ts +++ b/src/matches/match-relay/types/fragment.types.ts @@ -1,35 +1,17 @@ -export type StartFieldData = { - data?: Buffer; +export type FragmentField = "start" | "full" | "delta"; + +// What the game server sent in the query string alongside a field's body, plus +// what the relay records about it. Numeric protocol fields are stored as +// numbers so /sync can return them as JSON numbers. +export type FieldMeta = { gipped?: boolean; + timestamp?: number; signup_fragment?: number; tick?: number; + endtick?: number; tps?: number; map?: string; keyframe_interval?: number; protocol?: number; - [key: string]: any; -}; - -export type FullFieldData = { - data?: Buffer; - gipped?: boolean; - tick?: number; - [key: string]: any; + [key: string]: unknown; }; - -export type DeltaFieldData = { - data?: Buffer; - gipped?: boolean; - timestamp?: number; - endtick?: number; - [key: string]: any; -}; - -export type Fragment = { - start?: StartFieldData; - full?: FullFieldData; - delta?: DeltaFieldData; - [key: string]: any; -}; - -export type Broadcast = Fragment[]; diff --git a/src/matches/matches.controller.match-events.spec.ts b/src/matches/matches.controller.match-events.spec.ts index 071cd7ea1..b2e9a1332 100644 --- a/src/matches/matches.controller.match-events.spec.ts +++ b/src/matches/matches.controller.match-events.spec.ts @@ -127,8 +127,8 @@ describe("MatchesController — match_events on-demand servers", () => { })), }; matchRelay = { - removeBroadcast: jest.fn(), - playoutSeconds: jest.fn(() => 0), + removeBroadcast: jest.fn(async (): Promise => undefined), + playoutSeconds: jest.fn(async () => 0), }; inPlayMaps = []; @@ -449,7 +449,7 @@ describe("MatchesController — match_events on-demand servers", () => { }); it("lets relay viewers play out what they had buffered", async () => { - matchRelay.playoutSeconds.mockReturnValue(24); + matchRelay.playoutSeconds.mockResolvedValue(24); await finish(); @@ -469,7 +469,7 @@ describe("MatchesController — match_events on-demand servers", () => { it("adds the relay play-out on top of tv_delay mid-map", async () => { inPlayMaps = [{ id: "map-1" }]; - matchRelay.playoutSeconds.mockReturnValue(24); + matchRelay.playoutSeconds.mockResolvedValue(24); await finish(); diff --git a/src/matches/matches.controller.ts b/src/matches/matches.controller.ts index 9ea548dd7..57bed278f 100644 --- a/src/matches/matches.controller.ts +++ b/src/matches/matches.controller.ts @@ -991,7 +991,13 @@ export class MatchesController { } if (!scheduled) { - this.matchRelayService.removeBroadcast(matchId); + await this.matchRelayService.removeBroadcast(matchId).catch((error) => { + this.logger.error( + `[${matchId}] failed to remove the relay broadcast: ${ + (error as Error)?.message + }`, + ); + }); } } @@ -1038,7 +1044,11 @@ export class MatchesController { ); } - return feedEndsIn + this.matchRelayService.playoutSeconds(matchId); + const playout = await this.matchRelayService + .playoutSeconds(matchId) + .catch(() => 0); + + return feedEndsIn + playout; } private static stopOnDemandServerJobOptions(delaySeconds = 0) { diff --git a/test/match-relay.spec.ts b/test/match-relay.spec.ts new file mode 100644 index 000000000..bc813964f --- /dev/null +++ b/test/match-relay.spec.ts @@ -0,0 +1,406 @@ +import { PassThrough } from "stream"; +import { gunzipSync } from "zlib"; +import { Logger } from "@nestjs/common"; +import IORedis, { Redis } from "ioredis"; +import { GenericContainer, StartedTestContainer } from "testcontainers"; +import { MatchRelayService } from "../src/matches/match-relay/match-relay.service"; + +const fakeResponse = () => { + let resolveEnded: () => void; + const ended = new Promise((resolve) => { + resolveEnded = resolve; + }); + + const response = { + statusCode: undefined as number | undefined, + headers: {} as Record, + body: undefined as unknown, + headersSent: false, + ended, + writeHead(code: number, headers?: unknown) { + response.statusCode = code; + response.headersSent = true; + if (headers && typeof headers === "object") { + Object.assign(response.headers, headers); + } + return response; + }, + setHeader(name: string, value: unknown) { + response.headers[name] = value; + }, + end(body?: unknown) { + response.body = body; + resolveEnded(); + return response; + }, + }; + + return response; +}; + +// Runs against a real redis, like production: the relay's correctness is in +// what survives in redis between requests, which a mock cannot show. +describe("MatchRelayService", () => { + const matchId = "match-1"; + const token = "s845489096165654t8799308478907"; + + let container: StartedTestContainer; + let redis: Redis; + let service: MatchRelayService; + + const newService = () => + new MatchRelayService(new Logger("MatchRelayTest"), { + getConnection: () => redis, + } as any); + + beforeAll(async () => { + container = await new GenericContainer("redis:8.8-alpine") + .withExposedPorts(6379) + .start(); + redis = new IORedis({ + host: container.getHost(), + port: container.getMappedPort(6379), + }); + }, 120_000); + + afterAll(async () => { + redis?.disconnect(); + await container?.stop(); + }); + + beforeEach(async () => { + await redis.flushall(); + service = newService(); + }); + + const openPost = ( + field: "start" | "full" | "delta", + fragment: number, + query: Record = {}, + postToken = token, + relay = service, + ) => { + const request = Object.assign(new PassThrough(), { query }); + const response = fakeResponse(); + + const done = relay.postField( + request as any, + response as any, + postToken, + field, + matchId, + fragment, + ); + + return { + finish: async (body = `${field}-${fragment}`) => { + request.end(Buffer.from(body)); + await done; + await response.ended; + return response; + }, + }; + }; + + const post = ( + field: "start" | "full" | "delta", + fragment: number, + query: Record = {}, + postToken = token, + relay = service, + ) => openPost(field, fragment, query, postToken, relay).finish(); + + const sync = async (query: Record = {}, relay = service) => { + const response = fakeResponse(); + await relay.getSyncInfo({ query } as any, response as any, matchId); + return response; + }; + + const getStart = async (fragment: number) => { + const response = fakeResponse(); + await service.getStart(response as any, matchId, fragment); + return response; + }; + + const getFragment = async ( + fragment: number, + field: "full" | "delta", + fragmentToken?: string, + relay = service, + ) => { + const response = fakeResponse(); + await relay.getFragment( + response as any, + matchId, + fragment, + field, + fragmentToken, + ); + return response; + }; + + const startBroadcastAt = async (fragment: number, postToken = token) => { + await post( + "start", + fragment, + { + tick: "100", + tps: "64", + map: "de_inferno", + keyframe_interval: "3", + protocol: "5", + }, + postToken, + ); + await post("full", fragment, { tick: "100" }, postToken); + await post("delta", fragment, { endtick: "292" }, postToken); + }; + + const postFragments = async (from: number, to: number) => { + for (let fragment = from; fragment <= to; fragment++) { + await post("full", fragment, { tick: String(fragment * 192) }); + await post("delta", fragment, { endtick: String(fragment * 192 + 192) }); + } + }; + + it("reports the fragment the broadcast signed up at, with numeric fields", async () => { + await startBroadcastAt(42); + + const response = await sync({ fragment: "0" }); + + expect(response.statusCode).toBe(200); + expect(JSON.parse(response.body as string)).toEqual( + expect.objectContaining({ + fragment: 42, + signup_fragment: 42, + tick: 100, + endtick: 292, + maxtick: 292, + tps: 64, + keyframe_interval: 3, + map: "de_inferno", + protocol: 5, + }), + ); + }); + + it("starts a new client 7 fragments behind the newest", async () => { + await startBroadcastAt(42); + await postFragments(43, 52); + + const body = JSON.parse((await sync()).body as string); + + expect(body.fragment).toBe(45); + }); + + it("tells a proxy which broadcast the sync belongs to", async () => { + await startBroadcastAt(42); + + expect((await sync()).headers["X-Broadcast-Token"]).toBe(token); + }); + + it("reports how long clients keep playing once the server stops posting", async () => { + await startBroadcastAt(42); + + // Clients sit 7 fragments behind the newest and still have that one to play. + expect(await service.playoutSeconds(matchId)).toBe(8 * 3); + }); + + it("has nothing to play out for a broadcast it does not hold", async () => { + expect(await service.playoutSeconds(matchId)).toBe(0); + }); + + // CS2 sends tps as a decimal (64.0), not an integer, so the coercion has to + // accept a fractional part or clients get tps back as a string. + it("reports a fractional tps as a number", async () => { + await post("start", 42, { tick: "100", tps: "64.0", map: "de_inferno" }); + await post("full", 42, { tick: "100" }); + await post("delta", 42, { endtick: "292" }); + + const response = await sync({ fragment: "0" }); + + expect(JSON.parse(response.body as string)).toEqual( + expect.objectContaining({ tps: 64 }), + ); + }); + + it("keeps fields that are not numeric in /sync as the strings the server sent", async () => { + await post("start", 42, { + tick: "100", + tps: "64", + map: "3070284539", + protocol: "5", + }); + await post("full", 42, { tick: "100" }); + await post("delta", 42, { endtick: "292" }); + + const body = JSON.parse((await sync({ fragment: "0" })).body as string); + + expect(body.map).toBe("3070284539"); + expect(body.tick).toBe(100); + }); + + it("serves start only at the fragment the broadcast signed up at", async () => { + await startBroadcastAt(42); + + expect((await getStart(42)).statusCode).toBe(200); + expect((await getStart(0)).statusCode).toBe(404); + }); + + it("moves the signup fragment when the game server re-sends start", async () => { + await startBroadcastAt(42); + await post("start", 50, { tick: "900", tps: "64", map: "de_inferno" }); + + expect((await getStart(50)).statusCode).toBe(200); + expect((await getStart(42)).statusCode).toBe(404); + }); + + it("asks for start again when a fragment arrives before any start", async () => { + const response = await post("full", 7, { tick: "100" }); + + expect(response.statusCode).toBe(205); + }); + + it("asks for start again when a fragment arrives before the start data has", async () => { + const start = openPost("start", 42, { tick: "100", tps: "64" }); + + const early = await post("full", 42, { tick: "100" }); + expect(early.statusCode).toBe(205); + + await start.finish(); + + const late = await post("full", 43, { tick: "292" }); + expect(late.statusCode).toBe(200); + }); + + it("serves the posted fragment back, gzipped", async () => { + await startBroadcastAt(42); + + const response = await getFragment(42, "full"); + + expect(response.statusCode).toBe(200); + expect(response.headers["Content-Encoding"]).toBe("gzip"); + expect(gunzipSync(response.body as Buffer).toString()).toBe("full-42"); + }); + + it("keeps a broadcast going across an api restart", async () => { + await startBroadcastAt(42); + await postFragments(43, 52); + + const restarted = newService(); + + expect((await sync({}, restarted)).statusCode).toBe(200); + expect( + (await getFragment(50, "delta", undefined, restarted)).statusCode, + ).toBe(200); + expect( + (await post("full", 53, { tick: "1" }, token, restarted)).statusCode, + ).toBe(200); + }); + + it("drops the old broadcast's fragments when a new broadcast starts", async () => { + await startBroadcastAt(42); + await postFragments(43, 45); + + const nextToken = "s845489096165654t1111111111111"; + await startBroadcastAt(3, nextToken); + + expect((await getFragment(44, "full")).statusCode).toBe(404); + expect((await getFragment(3, "full")).statusCode).toBe(200); + expect((await sync()).headers["X-Broadcast-Token"]).toBe(nextToken); + }); + + it("lets a fragment under the current token be cached for good", async () => { + await startBroadcastAt(42); + + const scoped = await getFragment(42, "full", token); + expect(scoped.statusCode).toBe(200); + expect(scoped.headers["Cache-Control"]).toContain("immutable"); + + const unscoped = await getFragment(42, "full"); + expect(unscoped.headers["Cache-Control"]).toBeUndefined(); + }); + + it("refuses a fragment under an old token, uncached", async () => { + await startBroadcastAt(42); + + const stale = await getFragment(42, "full", "s1t2"); + + expect(stale.statusCode).toBe(404); + expect(stale.headers["Cache-Control"]).toBe("no-store"); + }); + + it("does not cache a fragment that has not arrived yet", async () => { + await startBroadcastAt(42); + + const missing = await getFragment(43, "full", token); + + expect(missing.statusCode).toBe(404); + expect(missing.headers["Cache-Control"]).toBe("no-store"); + }); + + it("forgets everything about a removed broadcast", async () => { + await startBroadcastAt(42); + await postFragments(43, 45); + + await service.removeBroadcast(matchId); + + expect((await sync()).statusCode).toBe(404); + expect((await getFragment(44, "full")).statusCode).toBe(404); + expect(await service.playoutSeconds(matchId)).toBe(0); + expect(await redis.keys(`match-relay:${matchId}:*`)).toEqual([]); + }); + + it("lets fragments expire a minute after they arrive", async () => { + await startBroadcastAt(42); + + const ttl = await redis.ttl(`match-relay:${matchId}:${token}:fragment:42`); + + expect(ttl).toBeGreaterThan(0); + expect(ttl).toBeLessThanOrEqual(60); + }); + + it("keeps the start of a broadcast alive for as long as fragments arrive", async () => { + await startBroadcastAt(42); + const startKey = `match-relay:${matchId}:${token}:start`; + await redis.expire(startKey, 5); + + await postFragments(43, 43); + + expect(await redis.ttl(startKey)).toBeGreaterThan(3000); + }); + + it("does not tell the game server a fragment was stored when it was not", async () => { + await startBroadcastAt(42); + await redis.set( + `match-relay:${matchId}:${token}:fragment:43`, + "not a hash", + ); + + const request = Object.assign(new PassThrough(), { query: { tick: "1" } }); + const response = fakeResponse(); + const done = service.postField( + request as any, + response as any, + token, + "full", + matchId, + 43, + ); + request.end(Buffer.from("full-43")); + + await expect(done).rejects.toThrow(); + expect(response.statusCode).toBeUndefined(); + }); + + it("never reads an old broadcast's data under the new one's token", async () => { + await startBroadcastAt(42); + await postFragments(43, 45); + + const nextToken = "s845489096165654t1111111111111"; + await startBroadcastAt(3, nextToken); + + expect(await redis.keys(`match-relay:${matchId}:${token}:*`)).toEqual([]); + expect((await getFragment(44, "full", nextToken)).statusCode).toBe(404); + }); +});