From 8ce160b8aef25fa5e11c0bbb44ef5b865a0714c4 Mon Sep 17 00:00:00 2001 From: zk Date: Sat, 13 Jun 2026 11:22:35 +0100 Subject: [PATCH] feat: maybe stabilise pb cleaning (#1662) --- ...0_pb_dirty_committed_and_calculated_at.sql | 28 + .../db/src/generated/public/GameProfile.ts | 2 + typescript/db/src/generated/public/Pb.ts | 2 + typescript/db/src/generated/public/Session.ts | 2 + .../src/lib/dirty-queues/calculation-run.ts | 15 + .../src/lib/dirty-queues/claim-dirty-queue.ts | 99 ++++ .../lib/dirty-queues/dirty-queue-race.test.ts | 496 ++++++++++++++++++ .../server/src/lib/jobs/drain-dirty-queues.ts | 90 ++-- .../framework/pb/pb-dirty.test.ts | 216 ++++++++ .../score-import/framework/pb/process-pbs.ts | 13 +- .../score-import/framework/pb/upsert-pb-pg.ts | 30 +- .../framework/ugpt-stats/update-ugpt-stats.ts | 14 +- .../src/utils/calculations/recalc-scores.ts | 7 +- 13 files changed, 954 insertions(+), 60 deletions(-) create mode 100644 db/migrations/20260612120000_pb_dirty_committed_and_calculated_at.sql create mode 100644 typescript/server/src/lib/dirty-queues/calculation-run.ts create mode 100644 typescript/server/src/lib/dirty-queues/claim-dirty-queue.ts create mode 100644 typescript/server/src/lib/dirty-queues/dirty-queue-race.test.ts create mode 100644 typescript/server/src/lib/score-import/framework/pb/pb-dirty.test.ts diff --git a/db/migrations/20260612120000_pb_dirty_committed_and_calculated_at.sql b/db/migrations/20260612120000_pb_dirty_committed_and_calculated_at.sql new file mode 100644 index 000000000..06c752dfb --- /dev/null +++ b/db/migrations/20260612120000_pb_dirty_committed_and_calculated_at.sql @@ -0,0 +1,28 @@ +-- Only committed scores should mark PBs dirty on insert/update (staging scores are invisible +-- to async drains). Deletes always enqueue: a removed score may have been merged into the +-- current PB (including staging scores written during a failed import). +CREATE OR REPLACE FUNCTION enqueue_pb_dirty() RETURNS trigger AS $$ +BEGIN + IF TG_OP = 'DELETE' THEN + INSERT INTO pb_dirty (user_id, chart_id) + SELECT DISTINCT score_old.user_id, score_old.chart_id + FROM score_old + ORDER BY score_old.user_id, score_old.chart_id + ON CONFLICT DO NOTHING; + ELSE + INSERT INTO pb_dirty (user_id, chart_id) + SELECT DISTINCT score_new.user_id, score_new.chart_id + FROM score_new + WHERE score_new.committed + ORDER BY score_new.user_id, score_new.chart_id + ON CONFLICT DO NOTHING; + END IF; + RETURN NULL; +END; +$$ LANGUAGE plpgsql; + +-- When a calculation run finishes, last_clean_started_at records when that run *started*. +-- Writers only apply if their runStartedAt is >= the row's last_clean_started_at (stale runs lose). +ALTER TABLE pb ADD COLUMN last_clean_started_at TIMESTAMPTZ NOT NULL DEFAULT now(); +ALTER TABLE session ADD COLUMN last_clean_started_at TIMESTAMPTZ NOT NULL DEFAULT now(); +ALTER TABLE game_profile ADD COLUMN last_clean_started_at TIMESTAMPTZ NOT NULL DEFAULT now(); diff --git a/typescript/db/src/generated/public/GameProfile.ts b/typescript/db/src/generated/public/GameProfile.ts index 7409b12c8..13b8e3722 100644 --- a/typescript/db/src/generated/public/GameProfile.ts +++ b/typescript/db/src/generated/public/GameProfile.ts @@ -33,6 +33,8 @@ export default interface GameProfileTable { data: ColumnType; showcase: ColumnType; + + last_clean_started_at: ColumnType; } export type GameProfile = Selectable; diff --git a/typescript/db/src/generated/public/Pb.ts b/typescript/db/src/generated/public/Pb.ts index b2c935041..17125042e 100644 --- a/typescript/db/src/generated/public/Pb.ts +++ b/typescript/db/src/generated/public/Pb.ts @@ -41,6 +41,8 @@ export default interface PbTable { highlight: ColumnType; time_achieved: ColumnType; + + last_clean_started_at: ColumnType; } export type Pb = Selectable; diff --git a/typescript/db/src/generated/public/Session.ts b/typescript/db/src/generated/public/Session.ts index e3269ec5b..02de70153 100644 --- a/typescript/db/src/generated/public/Session.ts +++ b/typescript/db/src/generated/public/Session.ts @@ -31,6 +31,8 @@ export default interface SessionTable { highlight: ColumnType; textsearch: ColumnType; + + last_clean_started_at: ColumnType; } export type Session = Selectable; diff --git a/typescript/server/src/lib/dirty-queues/calculation-run.ts b/typescript/server/src/lib/dirty-queues/calculation-run.ts new file mode 100644 index 000000000..710948222 --- /dev/null +++ b/typescript/server/src/lib/dirty-queues/calculation-run.ts @@ -0,0 +1,15 @@ +import DB from "#services/pg/db"; +import { sql } from "kysely"; + +/** + * ISO timestamp captured when a stats calculation run begins. Upserts only apply when + * {@link runStartedAt} is >= the row's existing `last_clean_started_at`. + * + * Uses Postgres `now()` so the guard compares against the same clock as `last_clean_started_at`. + */ +export type CalculationRunStartedAt = string; + +export async function newCalculationRunStartedAt(): Promise { + const result = await sql<{ ts: CalculationRunStartedAt }>`select now() as ts`.execute(DB); + return result.rows[0]!.ts; +} diff --git a/typescript/server/src/lib/dirty-queues/claim-dirty-queue.ts b/typescript/server/src/lib/dirty-queues/claim-dirty-queue.ts new file mode 100644 index 000000000..96c99b585 --- /dev/null +++ b/typescript/server/src/lib/dirty-queues/claim-dirty-queue.ts @@ -0,0 +1,99 @@ +import type { Game } from "tachi-db"; + +import DB from "#services/pg/db"; +import { sql } from "kysely"; + +export interface ClaimedPbDirtyRow { + user_id: number; + chart_id: string; + chart_game: Game; +} + +export interface ClaimedSessionDirtyRow { + session_id: string; +} + +export interface ClaimedGameProfileDirtyRow { + user_id: number; + game: Game; +} + +/** + * Atomically claim `pb_dirty` rows for processing. Uses `FOR UPDATE SKIP LOCKED` so + * concurrent drain workers never process the same pair. + */ +export function claimPbDirtyRows(limit: number): Promise> { + return DB.transaction().execute(async (trx) => { + const result = await sql` + WITH to_claim AS ( + SELECT d.user_id, d.chart_id, c.game AS chart_game + FROM pb_dirty d + INNER JOIN chart c ON c.id = d.chart_id + ORDER BY d.enqueued_at ASC + LIMIT ${limit} + FOR UPDATE OF d SKIP LOCKED + ), + deleted AS ( + DELETE FROM pb_dirty d + USING to_claim t + WHERE d.user_id = t.user_id + AND d.chart_id = t.chart_id + RETURNING d.user_id, d.chart_id, t.chart_game + ) + SELECT user_id, chart_id, chart_game FROM deleted + `.execute(trx); + + return result.rows; + }); +} + +/** Atomically claim `session_dirty` rows (`FOR UPDATE SKIP LOCKED`). */ +export function claimSessionDirtyRows(limit: number): Promise> { + return DB.transaction().execute(async (trx) => { + const result = await sql` + WITH to_claim AS ( + SELECT d.session_id + FROM session_dirty d + ORDER BY d.enqueued_at ASC + LIMIT ${limit} + FOR UPDATE OF d SKIP LOCKED + ), + deleted AS ( + DELETE FROM session_dirty d + USING to_claim t + WHERE d.session_id = t.session_id + RETURNING d.session_id + ) + SELECT session_id FROM deleted + `.execute(trx); + + return result.rows; + }); +} + +/** Atomically claim `game_profile_dirty` rows (`FOR UPDATE SKIP LOCKED`). */ +export function claimGameProfileDirtyRows( + limit: number, +): Promise> { + return DB.transaction().execute(async (trx) => { + const result = await sql` + WITH to_claim AS ( + SELECT d.user_id, d.game + FROM game_profile_dirty d + ORDER BY d.enqueued_at ASC + LIMIT ${limit} + FOR UPDATE OF d SKIP LOCKED + ), + deleted AS ( + DELETE FROM game_profile_dirty d + USING to_claim t + WHERE d.user_id = t.user_id + AND d.game = t.game + RETURNING d.user_id, d.game + ) + SELECT user_id, game FROM deleted + `.execute(trx); + + return result.rows; + }); +} diff --git a/typescript/server/src/lib/dirty-queues/dirty-queue-race.test.ts b/typescript/server/src/lib/dirty-queues/dirty-queue-race.test.ts new file mode 100644 index 000000000..776c2b314 --- /dev/null +++ b/typescript/server/src/lib/dirty-queues/dirty-queue-race.test.ts @@ -0,0 +1,496 @@ +import { ComputeChartStabilityChecksum } from "#game-implementations/utils/derivation-checksum"; +import { newCalculationRunStartedAt } from "#lib/dirty-queues/calculation-run"; +import { + claimGameProfileDirtyRows, + claimPbDirtyRows, + claimSessionDirtyRows, +} from "#lib/dirty-queues/claim-dirty-queue"; +import { clearPbDirtyForUser, drainPbDirty } from "#lib/jobs/drain-dirty-queues"; +import { log } from "#lib/log/log"; +import { CreatePBDoc } from "#lib/score-import/framework/pb/create-pb-doc"; +import { ProcessPBs } from "#lib/score-import/framework/pb/process-pbs"; +import { upsertPbFromMongoDoc } from "#lib/score-import/framework/pb/upsert-pb-pg"; +import { mongoScoreDataToPg, pgScoreDataToAPI } from "#lib/v3/migration-tools"; +import DB from "#services/pg/db"; +import { mkFakeScoreIIDXSP } from "#test-utils/misc"; +import { seedUser } from "#test-utils/pg-fixtures"; +import { Testing511Song, Testing511SPA, TestingIIDXSPScore } from "#test-utils/test-data"; +import { GetChartForIDGuaranteed } from "#utils/db"; +import { UnixMillisecondsToISO8601 } from "#utils/time"; +import { type PgScoreData } from "tachi-common"; +import { describe, expect, it } from "vitest"; + +let raceSeedCounter = 0; + +function nextId(prefix: string) { + raceSeedCounter += 1; + return `${prefix}-${raceSeedCounter}`; +} + +async function readPbPercent(userId: number, chartId: string): Promise { + const row = await DB.selectFrom("pb") + .select(["pb.data", "pb.derived_data", "pb.judgements"]) + .where("pb.user_id", "=", userId) + .where("pb.chart_id", "=", chartId) + .where("pb.lens", "is", null) + .executeTakeFirst(); + + if (!row) { + throw new Error(`missing pb for user ${userId} chart ${chartId}`); + } + + const scoreData = pgScoreDataToAPI("iidx-sp", { + data: row.data, + derived: row.derived_data, + judgements: row.judgements, + } as PgScoreData<"iidx-sp">); + + return scoreData.percent; +} + +async function seedIidxChart(chartId: string, songId: string) { + await DB.insertInto("song") + .values({ + id: songId, + legacy_id: 95_000 + raceSeedCounter, + game_group: "iidx", + title: Testing511Song.title, + artist: Testing511Song.artist, + search_terms: Testing511Song.searchTerms, + alt_titles: Testing511Song.altTitles, + data: Testing511Song.data, + fts_document: "", + }) + .onConflict((oc) => oc.doNothing()) + .execute(); + + const chartDoc = { + ...Testing511SPA, + chartID: chartId, + song: { ...Testing511Song, id: songId }, + }; + + await DB.insertInto("chart") + .values({ + id: chartId, + legacy_id: `legacy-race-${chartId}`, + game: "iidx-sp", + song_id: songId, + difficulty: chartDoc.difficulty, + level: chartDoc.level, + level_num: chartDoc.levelNum, + is_primary: chartDoc.isPrimary, + versions: chartDoc.versions, + data: chartDoc.data, + derivation_checksum: ComputeChartStabilityChecksum("iidx-sp", chartDoc), + }) + .onConflict((oc) => oc.doNothing()) + .execute(); +} + +async function insertCommittedScore(opts: { + chartId: string; + percent: number; + scoreId: string; + userId: number; +}) { + const scoreData = { + ...TestingIIDXSPScore.scoreData, + percent: opts.percent, + }; + const doc = mkFakeScoreIIDXSP({ + userID: opts.userId, + scoreID: opts.scoreId, + chartID: opts.chartId, + scoreData, + calculatedData: TestingIIDXSPScore.calculatedData, + timeAchieved: 1_700_000_000_000, + timeAdded: 1_700_000_000_000, + }); + const { data, derived, judgements } = mongoScoreDataToPg("iidx-sp", doc.scoreData); + const ts = UnixMillisecondsToISO8601(1_700_000_000_000); + + await DB.insertInto("score") + .values({ + id: opts.scoreId, + user_id: opts.userId, + chart_id: opts.chartId, + game: "iidx-sp", + session_id: null, + import_id: null, + data: JSON.stringify(data), + derived_data: JSON.stringify(derived), + judgements: JSON.stringify(judgements), + calculated_data: JSON.stringify(doc.calculatedData), + meta: JSON.stringify(doc.scoreMeta), + time_achieved: ts, + time_added: ts, + highlight: false, + comment: null, + committed: true, + }) + .execute(); +} + +async function insertUncommittedScore(opts: { + chartId: string; + importId: string; + percent: number; + scoreId: string; + userId: number; +}) { + const scoreData = { + ...TestingIIDXSPScore.scoreData, + percent: opts.percent, + }; + const doc = mkFakeScoreIIDXSP({ + userID: opts.userId, + scoreID: opts.scoreId, + chartID: opts.chartId, + scoreData, + calculatedData: TestingIIDXSPScore.calculatedData, + timeAchieved: 1_700_000_100_000, + timeAdded: 1_700_000_100_000, + }); + const { data, derived, judgements } = mongoScoreDataToPg("iidx-sp", doc.scoreData); + const ts = UnixMillisecondsToISO8601(1_700_000_100_000); + + await DB.insertInto("score") + .values({ + id: opts.scoreId, + user_id: opts.userId, + chart_id: opts.chartId, + game: "iidx-sp", + session_id: null, + import_id: opts.importId, + data: JSON.stringify(data), + derived_data: JSON.stringify(derived), + judgements: JSON.stringify(judgements), + calculated_data: JSON.stringify(doc.calculatedData), + meta: JSON.stringify(doc.scoreMeta), + time_achieved: ts, + time_added: ts, + highlight: false, + comment: null, + committed: false, + }) + .execute(); +} + +/** Run `workerCount` async workers until `work()` returns zero (queue drained). */ +async function runParallelUntilIdle( + workerCount: number, + work: () => Promise, +): Promise { + let total = 0; + + while (true) { + const batchTotals = await Promise.all(Array.from({ length: workerCount }, () => work())); + const moved = batchTotals.reduce((sum, n) => sum + n, 0); + + if (moved === 0) { + break; + } + + total += moved; + } + + return total; +} + +describe("dirty queue race harness", () => { + it("parallel pb_dirty claims are disjoint (SKIP LOCKED)", async () => { + const pairs = await Promise.all( + Array.from({ length: 12 }, async () => { + const { id: user_id } = await seedUser({ username: nextId("claim-user") }); + const chartId = nextId("chart-claim"); + const songId = nextId("song-claim"); + await seedIidxChart(chartId, songId); + return { user_id, chart_id: chartId }; + }), + ); + + await DB.insertInto("pb_dirty").values(pairs).execute(); + + const workerCount = 8; + const claims = await Promise.all( + Array.from({ length: workerCount }, () => claimPbDirtyRows(5)), + ); + + const claimedPairs = claims.flat().map((r) => `${r.user_id}:${r.chart_id}`); + expect(new Set(claimedPairs).size).toBe(claimedPairs.length); + expect(claimedPairs).toHaveLength(pairs.length); + + const remaining = await DB.selectFrom("pb_dirty").selectAll().execute(); + expect(remaining).toHaveLength(0); + }); + + it("parallel session_dirty and game_profile_dirty claims are disjoint", async () => { + const { id: userId } = await seedUser({ username: nextId("sess-gp-claim") }); + const sessionIds = [nextId("sess-a"), nextId("sess-b"), nextId("sess-c")]; + const now = new Date().toISOString(); + + for (const sessionId of sessionIds) { + await DB.insertInto("session") + .values({ + id: sessionId, + user_id: userId, + game: "iidx-sp", + name: sessionId, + description: null, + time_inserted: now, + time_started: now, + time_ended: now, + calculated_data: JSON.stringify({}), + highlight: false, + }) + .execute(); + } + + await DB.insertInto("session_dirty") + .values(sessionIds.map((session_id) => ({ session_id }))) + .execute(); + + const sessionClaims = await Promise.all([ + claimSessionDirtyRows(2), + claimSessionDirtyRows(2), + claimSessionDirtyRows(2), + ]); + const claimedSessions = sessionClaims.flat().map((r) => r.session_id); + expect(new Set(claimedSessions).size).toBe(sessionIds.length); + + const users = await Promise.all([ + seedUser({ username: nextId("gp-u1") }), + seedUser({ username: nextId("gp-u2") }), + seedUser({ username: nextId("gp-u3") }), + ]); + + await DB.insertInto("game_profile_dirty") + .values( + users.map((u) => ({ + user_id: u.id, + game: "iidx-sp" as const, + })), + ) + .execute(); + + const profileClaims = await Promise.all([ + claimGameProfileDirtyRows(2), + claimGameProfileDirtyRows(2), + claimGameProfileDirtyRows(2), + ]); + const claimedProfiles = profileClaims.flat().map((r) => `${r.user_id}:${r.game}`); + expect(new Set(claimedProfiles).size).toBe(users.length); + }); + + it("concurrent pb upserts: newer-started run wins when stale finishes last", async () => { + const { id: userId } = await seedUser({ username: nextId("pb-upsert-race") }); + const chartId = nextId("chart-upsert-race"); + const songId = nextId("song-upsert-race"); + await seedIidxChart(chartId, songId); + + await insertCommittedScore({ + chartId, + scoreId: nextId("score-upsert-race"), + userId, + percent: 88, + }); + + const chart = await GetChartForIDGuaranteed(chartId); + const pbDoc = await CreatePBDoc("iidx-sp", userId, chart, log); + expect(pbDoc).toBeDefined(); + + const newerRun = await newCalculationRunStartedAt(); + const olderRun = new Date(Date.parse(newerRun) - 120_000).toISOString(); + + let releaseStale!: () => void; + const staleMayFinish = new Promise((resolve) => { + releaseStale = resolve; + }); + + const staleDoc = { + ...pbDoc!, + scoreData: { + ...pbDoc!.scoreData, + percent: 1, + }, + }; + + const [staleApplied, newerApplied] = await Promise.all([ + (async () => { + await staleMayFinish; + return DB.transaction().execute((trx) => + upsertPbFromMongoDoc(trx, staleDoc, olderRun), + ); + })(), + (async () => { + const applied = await DB.transaction().execute((trx) => + upsertPbFromMongoDoc(trx, pbDoc!, newerRun), + ); + releaseStale(); + return applied; + })(), + ]); + + expect(newerApplied).toBe(true); + expect(staleApplied).toBe(false); + expect(await readPbPercent(userId, chartId)).toBe(88); + + const row = await DB.selectFrom("pb") + .select(["pb.last_clean_started_at"]) + .where("pb.user_id", "=", userId) + .where("pb.chart_id", "=", chartId) + .where("pb.lens", "is", null) + .executeTakeFirstOrThrow(); + + expect(row.last_clean_started_at).toBe(newerRun); + }); + + it("concurrent session updates: newer-started run wins when stale finishes last", async () => { + const { id: userId } = await seedUser({ username: nextId("sess-race") }); + const sessionId = nextId("sess-race"); + const now = new Date().toISOString(); + + await DB.insertInto("session") + .values({ + id: sessionId, + user_id: userId, + game: "iidx-sp", + name: "race session", + description: null, + time_inserted: now, + time_started: now, + time_ended: now, + calculated_data: JSON.stringify({ marker: "initial" }), + highlight: false, + }) + .execute(); + + const newerRun = await newCalculationRunStartedAt(); + const olderRun = new Date(Date.parse(newerRun) - 120_000).toISOString(); + + let releaseStale!: () => void; + const staleMayFinish = new Promise((resolve) => { + releaseStale = resolve; + }); + + await Promise.all([ + (async () => { + await staleMayFinish; + await DB.updateTable("session") + .set({ + calculated_data: JSON.stringify({ marker: "stale" }), + last_clean_started_at: olderRun, + }) + .where("session.id", "=", sessionId) + .where("session.last_clean_started_at", "<=", olderRun) + .execute(); + })(), + (async () => { + await DB.updateTable("session") + .set({ + calculated_data: JSON.stringify({ marker: "newer" }), + last_clean_started_at: newerRun, + }) + .where("session.id", "=", sessionId) + .where("session.last_clean_started_at", "<=", newerRun) + .execute(); + releaseStale(); + })(), + ]); + + const row = await DB.selectFrom("session") + .select(["session.calculated_data", "session.last_clean_started_at"]) + .where("session.id", "=", sessionId) + .executeTakeFirstOrThrow(); + + expect((row.calculated_data as { marker: string }).marker).toBe("newer"); + expect(row.last_clean_started_at).toBe(newerRun); + }); + + it("parallel drainPbDirty workers fully drain without double-processing", async () => { + const { id: userId } = await seedUser({ username: nextId("drain-race") }); + const chartIds = [nextId("c1"), nextId("c2"), nextId("c3")]; + + for (const chartId of chartIds) { + const songId = nextId(`song-${chartId}`); + await seedIidxChart(chartId, songId); + await insertCommittedScore({ + chartId, + scoreId: nextId(`score-${chartId}`), + userId, + percent: 70 + chartIds.indexOf(chartId), + }); + } + + const workerCount = 4; + const processed = await runParallelUntilIdle(workerCount, drainPbDirty); + + expect(processed).toBe(chartIds.length); + + const remaining = await DB.selectFrom("pb_dirty") + .selectAll() + .where("pb_dirty.user_id", "=", userId) + .execute(); + expect(remaining).toHaveLength(0); + + for (const chartId of chartIds) { + expect(await readPbPercent(userId, chartId)).toBe(70 + chartIds.indexOf(chartId)); + } + }); + + it("drain ignores uncommitted scores even when pb_dirty row exists (staging race)", async () => { + const { id: userId } = await seedUser({ username: nextId("staging-race") }); + const chartId = nextId("chart-staging"); + const songId = nextId("song-staging"); + const importId = nextId("import-staging"); + await seedIidxChart(chartId, songId); + + await insertCommittedScore({ + chartId, + scoreId: nextId("score-committed"), + userId, + percent: 50, + }); + + await ProcessPBs("iidx-sp", userId, new Set([chartId]), log); + expect(await readPbPercent(userId, chartId)).toBe(50); + await clearPbDirtyForUser(userId, [chartId]); + + await insertUncommittedScore({ + chartId, + importId, + scoreId: nextId("score-staging"), + userId, + percent: 99, + }); + + const triggerDirty = await DB.selectFrom("pb_dirty") + .selectAll() + .where("pb_dirty.user_id", "=", userId) + .where("pb_dirty.chart_id", "=", chartId) + .execute(); + expect(triggerDirty).toHaveLength(0); + + // Simulate the old bug: a dirty row exists while staging is still in flight. + await DB.insertInto("pb_dirty").values({ user_id: userId, chart_id: chartId }).execute(); + + await drainPbDirty(); + expect(await readPbPercent(userId, chartId)).toBe(50); + + await DB.updateTable("score") + .set({ committed: true }) + .where("score.import_id", "=", importId) + .execute(); + + const afterCommitDirty = await DB.selectFrom("pb_dirty") + .selectAll() + .where("pb_dirty.user_id", "=", userId) + .where("pb_dirty.chart_id", "=", chartId) + .execute(); + expect(afterCommitDirty).toHaveLength(1); + + await drainPbDirty(); + expect(await readPbPercent(userId, chartId)).toBe(99); + }); +}); diff --git a/typescript/server/src/lib/jobs/drain-dirty-queues.ts b/typescript/server/src/lib/jobs/drain-dirty-queues.ts index ffc36c1af..51ff9fdfd 100644 --- a/typescript/server/src/lib/jobs/drain-dirty-queues.ts +++ b/typescript/server/src/lib/jobs/drain-dirty-queues.ts @@ -4,6 +4,12 @@ import { ToScoreDocument, } from "#lib/db-formats/score"; import { SELECT_SESSION_DOCUMENT } from "#lib/db-formats/session"; +import { newCalculationRunStartedAt } from "#lib/dirty-queues/calculation-run"; +import { + claimGameProfileDirtyRows, + claimPbDirtyRows, + claimSessionDirtyRows, +} from "#lib/dirty-queues/claim-dirty-queue"; import { log } from "#lib/log/log"; import { CreateSessionCalcData } from "#lib/score-import/framework/calculated-data/session"; import { ProcessPBs } from "#lib/score-import/framework/pb/process-pbs"; @@ -36,16 +42,11 @@ const SESSION_DIRTY_CAP = 50_000; const GAME_PROFILE_DIRTY_CAP = 5_000; /** - * Drain the `pb_dirty` queue: group entries by (game, playtype, user_id), - * call ProcessPBs per group, and delete processed rows. + * Drain the `pb_dirty` queue: claim rows with SKIP LOCKED, group by (game, playtype, user_id), + * call ProcessPBs per group. Claiming deletes rows atomically so concurrent workers are safe. */ export async function drainPbDirty(): Promise { - const rows = await DB.selectFrom("pb_dirty") - .innerJoin("chart", "chart.id", "pb_dirty.chart_id") - .select(["pb_dirty.user_id", "pb_dirty.chart_id", "chart.game as chart_game"]) - .orderBy("pb_dirty.enqueued_at", "asc") - .limit(PB_DIRTY_BATCH) - .execute(); + const rows = await claimPbDirtyRows(PB_DIRTY_BATCH); if (rows.length === 0) { return 0; @@ -53,7 +54,13 @@ export async function drainPbDirty(): Promise { const groups = new Map< string, - { chartIDs: Set; game: GameGroup; playtype: LEGACY_Playtype; userID: integer } + { + chartIDs: Set; + game: GameGroup; + playtype: LEGACY_Playtype; + runStartedAt: Awaited>; + userID: integer; + } >(); for (const row of rows) { @@ -63,7 +70,13 @@ export async function drainPbDirty(): Promise { let group = groups.get(key); if (!group) { - group = { game, playtype, userID: row.user_id, chartIDs: new Set() }; + group = { + game, + playtype, + userID: row.user_id, + chartIDs: new Set(), + runStartedAt: await newCalculationRunStartedAt(), + }; groups.set(key, group); } @@ -73,17 +86,9 @@ export async function drainPbDirty(): Promise { for (const group of groups.values()) { const v3Game = LEGACY_GameGroupPTToGame(group.game, group.playtype); // eslint-disable-next-line no-await-in-loop - await ProcessPBs(v3Game, group.userID, group.chartIDs, log); - } - - const processedPairs = rows.map((r) => [r.user_id, r.chart_id] as const); - - for (const [userId, chartId] of processedPairs) { - // eslint-disable-next-line no-await-in-loop - await DB.deleteFrom("pb_dirty") - .where("pb_dirty.user_id", "=", userId) - .where("pb_dirty.chart_id", "=", chartId) - .execute(); + await ProcessPBs(v3Game, group.userID, group.chartIDs, log, { + runStartedAt: group.runStartedAt, + }); } log.info(`Drained ${rows.length} pb_dirty entries across ${groups.size} user/game groups.`); @@ -141,14 +146,10 @@ export async function drainScoreRederive(): Promise { } /** - * Drain `session_dirty`: recompute `session.calculated_data` from visible scores in that session. + * Drain `session_dirty`: claim rows with SKIP LOCKED, recompute `session.calculated_data`. */ export async function drainSessionDirty(): Promise { - const rows = await DB.selectFrom("session_dirty") - .select(["session_dirty.session_id"]) - .orderBy("session_dirty.enqueued_at", "asc") - .limit(SESSION_DIRTY_BATCH) - .execute(); + const rows = await claimSessionDirtyRows(SESSION_DIRTY_BATCH); if (rows.length === 0) { return 0; @@ -156,6 +157,7 @@ export async function drainSessionDirty(): Promise { for (const row of rows) { const sessionId = row.session_id; + const runStartedAt = await newCalculationRunStartedAt(); // eslint-disable-next-line no-await-in-loop const scoreRows = await DB.selectFrom("score") @@ -168,10 +170,6 @@ export async function drainSessionDirty(): Promise { .execute(); if (scoreRows.length === 0) { - // eslint-disable-next-line no-await-in-loop - await DB.deleteFrom("session_dirty") - .where("session_dirty.session_id", "=", sessionId) - .execute(); continue; } @@ -182,10 +180,6 @@ export async function drainSessionDirty(): Promise { .executeTakeFirst(); if (!sessionRow) { - // eslint-disable-next-line no-await-in-loop - await DB.deleteFrom("session_dirty") - .where("session_dirty.session_id", "=", sessionId) - .execute(); continue; } @@ -196,13 +190,10 @@ export async function drainSessionDirty(): Promise { await DB.updateTable("session") .set({ calculated_data: JSON.stringify(calculatedData), + last_clean_started_at: runStartedAt, }) .where("session.id", "=", sessionId) - .execute(); - - // eslint-disable-next-line no-await-in-loop - await DB.deleteFrom("session_dirty") - .where("session_dirty.session_id", "=", sessionId) + .where("session.last_clean_started_at", "<=", runStartedAt) .execute(); } @@ -212,14 +203,10 @@ export async function drainSessionDirty(): Promise { } /** - * Drain `game_profile_dirty`: recompute `game_profile` ratings/classes for each (user, playtype). + * Drain `game_profile_dirty`: claim rows with SKIP LOCKED, recompute ratings/classes. */ export async function drainGameProfileDirty(): Promise { - const rows = await DB.selectFrom("game_profile_dirty") - .select(["game_profile_dirty.user_id", "game_profile_dirty.game"]) - .orderBy("game_profile_dirty.enqueued_at", "asc") - .limit(GAME_PROFILE_DIRTY_BATCH) - .execute(); + const rows = await claimGameProfileDirtyRows(GAME_PROFILE_DIRTY_BATCH); if (rows.length === 0) { return 0; @@ -227,15 +214,12 @@ export async function drainGameProfileDirty(): Promise { for (const row of rows) { const userId = row.user_id; + const runStartedAt = await newCalculationRunStartedAt(); // eslint-disable-next-line no-await-in-loop - await UpdateUsersGamePlaytypeStats(row.game as V3Game, userId, null, log); - - // eslint-disable-next-line no-await-in-loop - await DB.deleteFrom("game_profile_dirty") - .where("game_profile_dirty.user_id", "=", userId) - .where("game_profile_dirty.game", "=", row.game) - .execute(); + await UpdateUsersGamePlaytypeStats(row.game as V3Game, userId, null, log, { + runStartedAt, + }); } log.info(`Drained ${rows.length} game_profile_dirty entries.`); diff --git a/typescript/server/src/lib/score-import/framework/pb/pb-dirty.test.ts b/typescript/server/src/lib/score-import/framework/pb/pb-dirty.test.ts new file mode 100644 index 000000000..b164033e2 --- /dev/null +++ b/typescript/server/src/lib/score-import/framework/pb/pb-dirty.test.ts @@ -0,0 +1,216 @@ +import { newCalculationRunStartedAt } from "#lib/dirty-queues/calculation-run"; +import { log } from "#lib/log/log"; +import { mongoScoreDataToPg, pgScoreDataToAPI } from "#lib/v3/migration-tools"; +import DB from "#services/pg/db"; +import { mkFakeScoreIIDXSP } from "#test-utils/misc"; +import { seedUser } from "#test-utils/pg-fixtures"; +import { Testing511Song, Testing511SPA, TestingIIDXSPScore } from "#test-utils/test-data"; +import { GetChartForIDGuaranteed } from "#utils/db"; +import { UnixMillisecondsToISO8601 } from "#utils/time"; +import { type PgScoreData } from "tachi-common"; +import { describe, expect, it } from "vitest"; + +import { CreatePBDoc } from "./create-pb-doc"; +import { upsertPbFromMongoDoc } from "./upsert-pb-pg"; + +async function seedIidx511Chart() { + await DB.insertInto("song") + .values({ + id: Testing511Song.id, + legacy_id: 1, + game_group: "iidx", + title: Testing511Song.title, + artist: Testing511Song.artist, + search_terms: Testing511Song.searchTerms, + alt_titles: Testing511Song.altTitles, + data: Testing511Song.data, + fts_document: "", + }) + .onConflict((oc) => oc.doNothing()) + .execute(); + + await DB.insertInto("chart") + .values({ + id: Testing511SPA.chartID, + legacy_id: Testing511SPA.chartID, + game: "iidx-sp", + song_id: Testing511Song.id, + difficulty: Testing511SPA.difficulty, + level: Testing511SPA.level, + level_num: Testing511SPA.levelNum, + is_primary: Testing511SPA.isPrimary, + versions: Testing511SPA.versions, + data: Testing511SPA.data, + }) + .onConflict((oc) => oc.doNothing()) + .execute(); +} + +async function insertScoreRow(opts: { + chartId: string; + committed: boolean; + importId?: string | null; + percent?: number; + scoreId: string; + userId: number; +}) { + const scoreData = { + ...TestingIIDXSPScore.scoreData, + percent: opts.percent ?? TestingIIDXSPScore.scoreData.percent, + }; + const doc = mkFakeScoreIIDXSP({ + userID: opts.userId, + scoreID: opts.scoreId, + chartID: opts.chartId, + scoreData, + calculatedData: TestingIIDXSPScore.calculatedData, + timeAchieved: 1_700_000_000_000, + timeAdded: 1_700_000_000_000, + }); + const { data, derived, judgements } = mongoScoreDataToPg("iidx-sp", doc.scoreData); + const ts = UnixMillisecondsToISO8601(1_700_000_000_000); + + await DB.insertInto("score") + .values({ + id: opts.scoreId, + user_id: opts.userId, + chart_id: opts.chartId, + game: "iidx-sp", + session_id: null, + import_id: opts.importId ?? null, + data: JSON.stringify(data), + derived_data: JSON.stringify(derived), + judgements: JSON.stringify(judgements), + calculated_data: JSON.stringify(doc.calculatedData), + meta: JSON.stringify(doc.scoreMeta), + time_achieved: ts, + time_added: ts, + highlight: false, + comment: null, + committed: opts.committed, + }) + .execute(); +} + +describe("pb_dirty trigger and last_clean_started_at guards", () => { + it("does not enqueue pb_dirty for uncommitted score insert", async () => { + const { id: userId } = await seedUser(); + await seedIidx511Chart(); + + await insertScoreRow({ + chartId: Testing511SPA.chartID, + scoreId: "score_uncommitted", + userId, + committed: false, + importId: "import-staging-1", + }); + + const dirty = await DB.selectFrom("pb_dirty") + .selectAll() + .where("pb_dirty.user_id", "=", userId) + .where("pb_dirty.chart_id", "=", Testing511SPA.chartID) + .execute(); + + expect(dirty).toHaveLength(0); + }); + + it("enqueues pb_dirty for committed score insert", async () => { + const { id: userId } = await seedUser(); + await seedIidx511Chart(); + + await insertScoreRow({ + chartId: Testing511SPA.chartID, + scoreId: "score_committed", + userId, + committed: true, + }); + + const dirty = await DB.selectFrom("pb_dirty") + .selectAll() + .where("pb_dirty.user_id", "=", userId) + .where("pb_dirty.chart_id", "=", Testing511SPA.chartID) + .execute(); + + expect(dirty).toHaveLength(1); + }); + + it("enqueues pb_dirty when a staged score is committed", async () => { + const { id: userId } = await seedUser(); + await seedIidx511Chart(); + const importId = "import-commit-1"; + + await insertScoreRow({ + chartId: Testing511SPA.chartID, + scoreId: "score_staged", + userId, + committed: false, + importId, + }); + + await DB.updateTable("score") + .set({ committed: true }) + .where("score.id", "=", "score_staged") + .execute(); + + const dirty = await DB.selectFrom("pb_dirty") + .selectAll() + .where("pb_dirty.user_id", "=", userId) + .where("pb_dirty.chart_id", "=", Testing511SPA.chartID) + .execute(); + + expect(dirty).toHaveLength(1); + }); + + it("skips pb upsert when a newer calculation run already wrote the row", async () => { + const { id: userId } = await seedUser(); + await seedIidx511Chart(); + + await insertScoreRow({ + chartId: Testing511SPA.chartID, + scoreId: "score_calc_at", + userId, + committed: true, + percent: 99, + }); + + const chart = await GetChartForIDGuaranteed(Testing511SPA.chartID); + const newerRun = await newCalculationRunStartedAt(); + const pbDoc = await CreatePBDoc("iidx-sp", userId, chart, log); + expect(pbDoc).toBeDefined(); + + await DB.transaction().execute(async (trx) => { + await upsertPbFromMongoDoc(trx, pbDoc!, newerRun); + }); + + const staleRun = new Date(Date.parse(newerRun) - 60_000).toISOString(); + const staleDoc = { + ...pbDoc!, + scoreData: { + ...pbDoc!.scoreData, + percent: 1, + }, + }; + + const applied = await DB.transaction().execute((trx) => + upsertPbFromMongoDoc(trx, staleDoc, staleRun), + ); + + expect(applied).toBe(false); + + const row = await DB.selectFrom("pb") + .select(["pb.data", "pb.derived_data", "pb.judgements", "pb.last_clean_started_at"]) + .where("pb.user_id", "=", userId) + .where("pb.chart_id", "=", Testing511SPA.chartID) + .where("pb.lens", "is", null) + .executeTakeFirstOrThrow(); + + const scoreData = pgScoreDataToAPI("iidx-sp", { + data: row.data, + derived: row.derived_data, + judgements: row.judgements, + } as PgScoreData<"iidx-sp">); + + expect(scoreData.percent).toBe(99); + expect(row.last_clean_started_at).toBe(newerRun); + }); +}); diff --git a/typescript/server/src/lib/score-import/framework/pb/process-pbs.ts b/typescript/server/src/lib/score-import/framework/pb/process-pbs.ts index 01984c2a6..5dea55835 100644 --- a/typescript/server/src/lib/score-import/framework/pb/process-pbs.ts +++ b/typescript/server/src/lib/score-import/framework/pb/process-pbs.ts @@ -1,5 +1,9 @@ import type { KtLogger } from "#lib/log/log"; +import { + type CalculationRunStartedAt, + newCalculationRunStartedAt, +} from "#lib/dirty-queues/calculation-run"; import DB from "#services/pg/db"; import { GetChartForIDGuaranteed } from "#utils/db"; import { type integer, type V3Game } from "tachi-common"; @@ -7,6 +11,10 @@ import { type integer, type V3Game } from "tachi-common"; import { CreatePBDoc, type PBScoreDocumentNoRank, UpdateChartRanking } from "./create-pb-doc"; import { upsertPbFromMongoDoc } from "./upsert-pb-pg"; +export interface ProcessPBsOptions { + runStartedAt?: CalculationRunStartedAt; +} + /** * Process, recalculate and update a users PBs for this set of chartIDs. * @@ -18,11 +26,14 @@ export async function ProcessPBs( userID: integer, chartIDs: Set, log: KtLogger, + options?: ProcessPBsOptions, ): Promise { if (chartIDs.size === 0) { return; } + const runStartedAt = options?.runStartedAt ?? (await newCalculationRunStartedAt()); + const chartIDsArray = [...chartIDs]; const promises = chartIDsArray.map((chartID) => GetChartForIDGuaranteed(chartID).then((chart) => CreatePBDoc(game, userID, chart, log)), @@ -69,7 +80,7 @@ export async function ProcessPBs( await DB.transaction().execute(async (trx) => { // TODO(zk): parallelize? for (const doc of pbDocs) { - await upsertPbFromMongoDoc(trx, doc); + await upsertPbFromMongoDoc(trx, doc, runStartedAt); } }); diff --git a/typescript/server/src/lib/score-import/framework/pb/upsert-pb-pg.ts b/typescript/server/src/lib/score-import/framework/pb/upsert-pb-pg.ts index 06c53a459..31dc3fd2a 100644 --- a/typescript/server/src/lib/score-import/framework/pb/upsert-pb-pg.ts +++ b/typescript/server/src/lib/score-import/framework/pb/upsert-pb-pg.ts @@ -1,3 +1,4 @@ +import type { CalculationRunStartedAt } from "#lib/dirty-queues/calculation-run"; import type { Kysely } from "kysely"; import type { Database } from "tachi-db"; @@ -13,11 +14,14 @@ export type PBScoreDocumentNoRank = Omit< /** * Inserts or updates a `pb` row plus `pb_composed_from` for one user/chart (lens = null). + * + * Returns whether the row was written. Stale runs (see `last_clean_started_at`) are skipped. */ export async function upsertPbFromMongoDoc( db: Kysely, pbDoc: PBScoreDocumentNoRank, -): Promise { + runStartedAt: CalculationRunStartedAt, +): Promise { const game = pbDoc.game; const { data, derived, judgements } = mongoScoreDataToPg(game, pbDoc.scoreData); const judgementsJson = JSON.stringify(judgements); @@ -40,11 +44,12 @@ export async function upsertPbFromMongoDoc( .executeTakeFirst(); let pbId: string; + let applied = false; if (existing) { pbId = existing.row_id; - await db + const updated = await db .updateTable("pb") .set({ data: JSON.stringify(data), @@ -62,9 +67,14 @@ export async function upsertPbFromMongoDoc( pbDoc.timeAchieved !== null && pbDoc.timeAchieved !== undefined ? UnixMillisecondsToISO8601(pbDoc.timeAchieved) : null, + last_clean_started_at: runStartedAt, }) .where("row_id", "=", pbId) - .execute(); + .where("last_clean_started_at", "<=", runStartedAt) + .returning("row_id") + .executeTakeFirst(); + + applied = updated !== undefined; } else { const inserted = await db .insertInto("pb") @@ -87,11 +97,21 @@ export async function upsertPbFromMongoDoc( pbDoc.timeAchieved !== null && pbDoc.timeAchieved !== undefined ? UnixMillisecondsToISO8601(pbDoc.timeAchieved) : null, + last_clean_started_at: runStartedAt, }) .returning("row_id") - .executeTakeFirstOrThrow(); + .executeTakeFirst(); + + if (!inserted) { + return false; + } pbId = inserted.row_id; + applied = true; + } + + if (!applied) { + return false; } await db.deleteFrom("pb_composed_from").where("pb_id", "=", pbId).execute(); @@ -108,4 +128,6 @@ export async function upsertPbFromMongoDoc( ) .execute(); } + + return true; } diff --git a/typescript/server/src/lib/score-import/framework/ugpt-stats/update-ugpt-stats.ts b/typescript/server/src/lib/score-import/framework/ugpt-stats/update-ugpt-stats.ts index d27474147..6f0ac8d34 100644 --- a/typescript/server/src/lib/score-import/framework/ugpt-stats/update-ugpt-stats.ts +++ b/typescript/server/src/lib/score-import/framework/ugpt-stats/update-ugpt-stats.ts @@ -1,5 +1,9 @@ import type { KtLogger } from "#lib/log/log"; +import { + type CalculationRunStartedAt, + newCalculationRunStartedAt, +} from "#lib/dirty-queues/calculation-run"; import { newGameProfilePreferenceColumns } from "#lib/game-settings/create-game-settings"; import DB from "#services/pg/db"; import { loadUserGameStats } from "#utils/class"; @@ -11,13 +15,18 @@ import type { ClassProcessOptions } from "../profile-calculated-data/class-proce import { CalculateProfileRatings } from "../calculated-data/profile"; import { CalculateUGPTClasses, ProcessClassDeltas } from "../profile-calculated-data/classes"; +export type UpdateUsersGamePlaytypeStatsOptions = { + runStartedAt?: CalculationRunStartedAt; +} & ClassProcessOptions; + export async function UpdateUsersGamePlaytypeStats( game: V3Game, userID: integer, classProvider: ClassProvider | null, log: KtLogger, - options?: ClassProcessOptions, + options?: UpdateUsersGamePlaytypeStatsOptions, ): Promise> { + const runStartedAt = options?.runStartedAt ?? (await newCalculationRunStartedAt()); log.debug(`Calculating Ratings...`); const ratings = await CalculateProfileRatings(game, userID); @@ -53,9 +62,11 @@ export async function UpdateUsersGamePlaytypeStats( .set({ ratings: JSON.stringify(ratings), classes: JSON.stringify(nextClasses), + last_clean_started_at: runStartedAt, }) .where("user_id", "=", userID) .where("game", "=", game) + .where("last_clean_started_at", "<=", runStartedAt) .execute(); } else { const hasAnyScores = await DB.selectFrom("score") @@ -82,6 +93,7 @@ export async function UpdateUsersGamePlaytypeStats( game, ratings: JSON.stringify(ratings), classes: JSON.stringify(classes), + last_clean_started_at: runStartedAt, ...newGameProfilePreferenceColumns(game), }) .execute(); diff --git a/typescript/server/src/utils/calculations/recalc-scores.ts b/typescript/server/src/utils/calculations/recalc-scores.ts index a0fd31ae9..79b437dfe 100644 --- a/typescript/server/src/utils/calculations/recalc-scores.ts +++ b/typescript/server/src/utils/calculations/recalc-scores.ts @@ -18,7 +18,12 @@ export async function RecalcAllScores(): Promise { */ export async function UpdateAllPBs(): Promise { await DB.insertInto("pb_dirty") - .expression(DB.selectFrom("score").select(["score.user_id", "score.chart_id"]).distinct()) + .expression( + DB.selectFrom("score") + .select(["score.user_id", "score.chart_id"]) + .where("score.committed", "=", true) + .distinct(), + ) .onConflict((oc) => oc.doNothing()) .execute(); }