Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions src/steam-match-history/enums/SteamMatchHistoryQueues.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ export enum SteamMatchHistoryQueues {
PollAllSteamMatchHistory = "PollAllSteamMatchHistory",
ResolveMatchMetadata = "ResolveMatchMetadata",
ParseImportedDemo = "ParseImportedDemo",
ReconcilePendingMatchImports = "ReconcilePendingMatchImports",
ProcessUploadedDemo = "ProcessUploadedDemo",
CheckSteamBansForMatch = "CheckSteamBansForMatch",
CheckSteamBans = "CheckSteamBans",
Expand Down
93 changes: 93 additions & 0 deletions src/steam-match-history/jobs/ReconcilePendingMatchImports.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,93 @@
import { ReconcilePendingMatchImports } from "./ReconcilePendingMatchImports";

type FakeJob = { state: string; failedReason?: string };

const build = (
stale: string[],
jobs: Record<string, FakeJob | undefined> = {},
) => {
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);
},
);
});
85 changes: 85 additions & 0 deletions src/steam-match-history/jobs/ReconcilePendingMatchImports.ts
Original file line number Diff line number Diff line change
@@ -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<void> {
const rows = await this.postgres.query<Array<{ valve_match_id: string }>>(
`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}`,
);
}
}
}
}
56 changes: 56 additions & 0 deletions src/steam-match-history/jobs/ResolveMatchMetadata.spec.ts
Original file line number Diff line number Diff line change
@@ -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'],
]);
});
});
16 changes: 15 additions & 1 deletion src/steam-match-history/jobs/ResolveMatchMetadata.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,22 @@ export class ResolveMatchMetadata extends WorkerHost {
}

async process(job: Job<ResolveMatchMetadataPayload>): Promise<void> {
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<void> {
const rows = await this.postgres.query<
Array<{
share_code: string;
Expand Down
21 changes: 21 additions & 0 deletions src/steam-match-history/steam-match-history.module.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand All @@ -39,6 +40,9 @@ import { ProcessUploadedDemo } from "./jobs/ProcessUploadedDemo";
BullModule.registerQueue({
name: SteamMatchHistoryQueues.ParseImportedDemo,
}),
BullModule.registerQueue({
name: SteamMatchHistoryQueues.ReconcilePendingMatchImports,
}),
BullModule.registerQueue({
name: SteamMatchHistoryQueues.ProcessUploadedDemo,
}),
Expand Down Expand Up @@ -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,
Expand All @@ -100,6 +108,7 @@ import { ProcessUploadedDemo } from "./jobs/ProcessUploadedDemo";
PollSteamMatchHistoryForUser,
ResolveMatchMetadata,
ParseImportedDemo,
ReconcilePendingMatchImports,
ProcessUploadedDemo,
...getQueuesProcessors("SteamMatchHistory"),
loggerFactory(),
Expand All @@ -111,6 +120,8 @@ export class SteamMatchHistoryModule {
constructor(
@InjectQueue(SteamMatchHistoryQueues.PollAllSteamMatchHistory)
queue: Queue,
@InjectQueue(SteamMatchHistoryQueues.ReconcilePendingMatchImports)
reconcileQueue: Queue,
) {
if (process.env.RUN_MIGRATIONS) {
return;
Expand All @@ -125,5 +136,15 @@ export class SteamMatchHistoryModule {
},
},
);

void reconcileQueue.add(
ReconcilePendingMatchImports.name,
{},
{
repeat: {
pattern: "*/15 * * * *",
},
},
);
}
}
Loading