From 3ce3323dc379431a4aa79b9f23bd9aa38a5688e9 Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 18:32:04 -0400 Subject: [PATCH] bug: pending match imports can no longer get stuck --- .../enums/SteamMatchHistoryQueues.ts | 1 + .../jobs/ReconcilePendingMatchImports.spec.ts | 93 +++++++++++++++++++ .../jobs/ReconcilePendingMatchImports.ts | 85 +++++++++++++++++ .../jobs/ResolveMatchMetadata.spec.ts | 56 +++++++++++ .../jobs/ResolveMatchMetadata.ts | 16 +++- .../steam-match-history.module.ts | 21 +++++ 6 files changed, 271 insertions(+), 1 deletion(-) create mode 100644 src/steam-match-history/jobs/ReconcilePendingMatchImports.spec.ts create mode 100644 src/steam-match-history/jobs/ReconcilePendingMatchImports.ts create mode 100644 src/steam-match-history/jobs/ResolveMatchMetadata.spec.ts diff --git a/src/steam-match-history/enums/SteamMatchHistoryQueues.ts b/src/steam-match-history/enums/SteamMatchHistoryQueues.ts index 29334180b..23153aeae 100644 --- a/src/steam-match-history/enums/SteamMatchHistoryQueues.ts +++ b/src/steam-match-history/enums/SteamMatchHistoryQueues.ts @@ -2,6 +2,7 @@ export enum SteamMatchHistoryQueues { PollAllSteamMatchHistory = "PollAllSteamMatchHistory", ResolveMatchMetadata = "ResolveMatchMetadata", ParseImportedDemo = "ParseImportedDemo", + ReconcilePendingMatchImports = "ReconcilePendingMatchImports", ProcessUploadedDemo = "ProcessUploadedDemo", CheckSteamBansForMatch = "CheckSteamBansForMatch", CheckSteamBans = "CheckSteamBans", diff --git a/src/steam-match-history/jobs/ReconcilePendingMatchImports.spec.ts b/src/steam-match-history/jobs/ReconcilePendingMatchImports.spec.ts new file mode 100644 index 000000000..89c787f91 --- /dev/null +++ b/src/steam-match-history/jobs/ReconcilePendingMatchImports.spec.ts @@ -0,0 +1,93 @@ +import { ReconcilePendingMatchImports } from "./ReconcilePendingMatchImports"; + +type FakeJob = { state: string; failedReason?: string }; + +const build = ( + stale: string[], + jobs: Record = {}, +) => { + const updates: unknown[][] = []; + const postgres = { + query: jest.fn(async (sql: string, params: unknown[] = []) => { + if (sql.includes("SELECT valve_match_id")) { + return stale.map((valve_match_id) => ({ valve_match_id })); + } + if (sql.includes("SET status = 'Failed'")) { + updates.push(params); + return [{ valve_match_id: params[0] }]; + } + return []; + }), + }; + const queue = { + getJob: jest.fn(async (id: string) => { + const job = jobs[id]; + return job + ? { failedReason: job.failedReason, getState: async () => job.state } + : undefined; + }), + }; + const logger = { warn: jest.fn() }; + + const reconcile = new ReconcilePendingMatchImports( + logger as never, + postgres as never, + queue as never, + queue as never, + ); + + return { reconcile, updates }; +}; + +describe("ReconcilePendingMatchImports", () => { + it("fails a stranded row whose jobs are gone", async () => { + const { reconcile, updates } = build(["111"]); + + await reconcile.process(); + + expect(updates).toEqual([ + ["111", "import job ended without finishing the import", "15 minutes"], + ]); + }); + + it("carries the failed job's reason onto the row", async () => { + const { reconcile, updates } = build(["222"], { + "resolve-222": { + state: "failed", + failedReason: 'column "parties" does not exist', + }, + }); + + await reconcile.process(); + + expect(updates[0]?.[1]).toBe('column "parties" does not exist'); + }); + + it("fails a row whose parse job stalled out after resolve completed", async () => { + const { reconcile, updates } = build(["333"], { + "resolve-333": { state: "completed" }, + "parse-333": { + state: "failed", + failedReason: "job stalled more than allowable limit", + }, + }); + + await reconcile.process(); + + expect(updates[0]?.[1]).toBe("job stalled more than allowable limit"); + }); + + it.each(["waiting", "active", "delayed", "prioritized", "waiting-children"])( + "leaves a row alone while a job is %s", + async (state) => { + const { reconcile, updates } = build(["444"], { + "resolve-444": { state: "completed" }, + "parse-444": { state }, + }); + + await reconcile.process(); + + expect(updates).toHaveLength(0); + }, + ); +}); diff --git a/src/steam-match-history/jobs/ReconcilePendingMatchImports.ts b/src/steam-match-history/jobs/ReconcilePendingMatchImports.ts new file mode 100644 index 000000000..5c140adfc --- /dev/null +++ b/src/steam-match-history/jobs/ReconcilePendingMatchImports.ts @@ -0,0 +1,85 @@ +import { Logger } from "@nestjs/common"; +import { InjectQueue, WorkerHost } from "@nestjs/bullmq"; +import { Job, Queue } from "bullmq"; +import { UseQueue } from "src/utilities/QueueProcessors"; +import { PostgresService } from "../../postgres/postgres.service"; +import { SteamMatchHistoryQueues } from "../enums/SteamMatchHistoryQueues"; + +// A pending import only leaves Queued/Parsing from inside its own job, so a job +// that stalls out, is lost, or finishes without deciding strands the row, and +// only Failed rows can be retried. This fails any row whose jobs are all done. +@UseQueue( + "SteamMatchHistory", + SteamMatchHistoryQueues.ReconcilePendingMatchImports, +) +export class ReconcilePendingMatchImports extends WorkerHost { + // The row is written before its job is added, so a fresh row with no job + // yet is not stranded. + private static readonly GRACE = "15 minutes"; + + private static readonly FINISHED_STATES = new Set([ + "completed", + "failed", + "unknown", + ]); + + constructor( + private readonly logger: Logger, + private readonly postgres: PostgresService, + @InjectQueue(SteamMatchHistoryQueues.ResolveMatchMetadata) + private readonly resolveQueue: Queue, + @InjectQueue(SteamMatchHistoryQueues.ParseImportedDemo) + private readonly parseQueue: Queue, + ) { + super(); + } + + async process(): Promise { + const rows = await this.postgres.query>( + `SELECT valve_match_id::text AS valve_match_id + FROM public.pending_match_imports + WHERE status IN ('Queued', 'Parsing') + AND updated_at < now() - $1::interval`, + [ReconcilePendingMatchImports.GRACE], + ); + + for (const { valve_match_id } of rows) { + const jobs = ( + await Promise.all([ + this.resolveQueue.getJob(`resolve-${valve_match_id}`), + this.parseQueue.getJob(`parse-${valve_match_id}`), + ]) + ).filter((job): job is Job => !!job); + + const states = await Promise.all(jobs.map((job) => job.getState())); + if ( + states.some( + (state) => !ReconcilePendingMatchImports.FINISHED_STATES.has(state), + ) + ) { + continue; + } + + const reason = + jobs.map((job) => job.failedReason).find(Boolean) ?? + "import job ended without finishing the import"; + + const failed = await this.postgres.query< + Array<{ valve_match_id: string }> + >( + `UPDATE public.pending_match_imports + SET status = 'Failed', error = $2 + WHERE valve_match_id = $1::numeric + AND status IN ('Queued', 'Parsing') + AND updated_at < now() - $3::interval + RETURNING valve_match_id`, + [valve_match_id, reason, ReconcilePendingMatchImports.GRACE], + ); + if (failed.length > 0) { + this.logger.warn( + `reconcile-pending-match-imports failed stranded valve_match_id=${valve_match_id}: ${reason}`, + ); + } + } + } +} diff --git a/src/steam-match-history/jobs/ResolveMatchMetadata.spec.ts b/src/steam-match-history/jobs/ResolveMatchMetadata.spec.ts new file mode 100644 index 000000000..85c4cbae8 --- /dev/null +++ b/src/steam-match-history/jobs/ResolveMatchMetadata.spec.ts @@ -0,0 +1,56 @@ +import { ResolveMatchMetadata } from "./ResolveMatchMetadata"; + +const VALVE_MATCH_ID = "1174469974867590708"; + +const build = () => { + const failures: unknown[][] = []; + const postgres = { + query: jest.fn(async (sql: string, params: unknown[] = []) => { + if (sql.includes("SELECT share_code")) { + throw new Error('column "parties" does not exist'); + } + if (sql.includes("SET status = 'Failed'")) { + failures.push(params); + } + return []; + }), + }; + const logger = { log: jest.fn(), warn: jest.fn() }; + + const job = new ResolveMatchMetadata( + logger as never, + postgres as never, + {} as never, + {} as never, + {} as never, + ); + + const run = (attemptsMade: number) => + job.process({ + data: { valve_match_id: VALVE_MATCH_ID }, + attemptsMade, + opts: { attempts: 5 }, + } as never); + + return { run, failures }; +}; + +describe("ResolveMatchMetadata", () => { + it("leaves the row for the retry while attempts remain", async () => { + const { run, failures } = build(); + + await expect(run(0)).rejects.toThrow("parties"); + + expect(failures).toHaveLength(0); + }); + + it("marks the row Failed when the last attempt throws", async () => { + const { run, failures } = build(); + + await expect(run(4)).rejects.toThrow("parties"); + + expect(failures).toEqual([ + [VALVE_MATCH_ID, 'column "parties" does not exist'], + ]); + }); +}); diff --git a/src/steam-match-history/jobs/ResolveMatchMetadata.ts b/src/steam-match-history/jobs/ResolveMatchMetadata.ts index 916b84fac..5e58ff7e8 100644 --- a/src/steam-match-history/jobs/ResolveMatchMetadata.ts +++ b/src/steam-match-history/jobs/ResolveMatchMetadata.ts @@ -30,8 +30,22 @@ export class ResolveMatchMetadata extends WorkerHost { } async process(job: Job): Promise { - const { valve_match_id } = job.data; + try { + await this.resolve(job.data.valve_match_id); + } catch (error) { + const lastAttempt = + (job.attemptsMade ?? 0) >= (job.opts.attempts ?? 1) - 1; + if (lastAttempt) { + await this.markFailed( + job.data.valve_match_id, + (error as Error)?.message ?? String(error), + ); + } + throw error; + } + } + private async resolve(valve_match_id: string): Promise { const rows = await this.postgres.query< Array<{ share_code: string; diff --git a/src/steam-match-history/steam-match-history.module.ts b/src/steam-match-history/steam-match-history.module.ts index aa2738e0e..ffddcd50f 100644 --- a/src/steam-match-history/steam-match-history.module.ts +++ b/src/steam-match-history/steam-match-history.module.ts @@ -26,6 +26,7 @@ import { DrainSteamBans } from "./jobs/DrainSteamBans"; import { PollSteamMatchHistoryForUser } from "./jobs/PollSteamMatchHistoryForUser"; import { ResolveMatchMetadata } from "./jobs/ResolveMatchMetadata"; import { ParseImportedDemo } from "./jobs/ParseImportedDemo"; +import { ReconcilePendingMatchImports } from "./jobs/ReconcilePendingMatchImports"; import { ProcessUploadedDemo } from "./jobs/ProcessUploadedDemo"; @Module({ @@ -39,6 +40,9 @@ import { ProcessUploadedDemo } from "./jobs/ProcessUploadedDemo"; BullModule.registerQueue({ name: SteamMatchHistoryQueues.ParseImportedDemo, }), + BullModule.registerQueue({ + name: SteamMatchHistoryQueues.ReconcilePendingMatchImports, + }), BullModule.registerQueue({ name: SteamMatchHistoryQueues.ProcessUploadedDemo, }), @@ -75,6 +79,10 @@ import { ProcessUploadedDemo } from "./jobs/ProcessUploadedDemo"; name: SteamMatchHistoryQueues.ParseImportedDemo, adapter: BullMQAdapter, }), + BullBoardModule.forFeature({ + name: SteamMatchHistoryQueues.ReconcilePendingMatchImports, + adapter: BullMQAdapter, + }), BullBoardModule.forFeature({ name: SteamMatchHistoryQueues.ProcessUploadedDemo, adapter: BullMQAdapter, @@ -100,6 +108,7 @@ import { ProcessUploadedDemo } from "./jobs/ProcessUploadedDemo"; PollSteamMatchHistoryForUser, ResolveMatchMetadata, ParseImportedDemo, + ReconcilePendingMatchImports, ProcessUploadedDemo, ...getQueuesProcessors("SteamMatchHistory"), loggerFactory(), @@ -111,6 +120,8 @@ export class SteamMatchHistoryModule { constructor( @InjectQueue(SteamMatchHistoryQueues.PollAllSteamMatchHistory) queue: Queue, + @InjectQueue(SteamMatchHistoryQueues.ReconcilePendingMatchImports) + reconcileQueue: Queue, ) { if (process.env.RUN_MIGRATIONS) { return; @@ -125,5 +136,15 @@ export class SteamMatchHistoryModule { }, }, ); + + void reconcileQueue.add( + ReconcilePendingMatchImports.name, + {}, + { + repeat: { + pattern: "*/15 * * * *", + }, + }, + ); } }