feat: maybe stabilise pb cleaning (#1662)

This commit is contained in:
zk
2026-06-13 11:22:35 +01:00
committed by GitHub
parent 99fe1186ff
commit 8ce160b8ae
13 changed files with 954 additions and 60 deletions
@@ -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();
@@ -33,6 +33,8 @@ export default interface GameProfileTable {
data: ColumnType<unknown, unknown | undefined, unknown>;
showcase: ColumnType<unknown, unknown | undefined, unknown>;
last_clean_started_at: ColumnType<string, string | undefined, string>;
}
export type GameProfile = Selectable<GameProfileTable>;
+2
View File
@@ -41,6 +41,8 @@ export default interface PbTable {
highlight: ColumnType<boolean, boolean, boolean>;
time_achieved: ColumnType<string | null, string | null, string | null>;
last_clean_started_at: ColumnType<string, string | undefined, string>;
}
export type Pb = Selectable<PbTable>;
@@ -31,6 +31,8 @@ export default interface SessionTable {
highlight: ColumnType<boolean, boolean, boolean>;
textsearch: ColumnType<string, never, never>;
last_clean_started_at: ColumnType<string, string | undefined, string>;
}
export type Session = Selectable<SessionTable>;
@@ -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<CalculationRunStartedAt> {
const result = await sql<{ ts: CalculationRunStartedAt }>`select now() as ts`.execute(DB);
return result.rows[0]!.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<Array<ClaimedPbDirtyRow>> {
return DB.transaction().execute(async (trx) => {
const result = await sql<ClaimedPbDirtyRow>`
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<Array<ClaimedSessionDirtyRow>> {
return DB.transaction().execute(async (trx) => {
const result = await sql<ClaimedSessionDirtyRow>`
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<Array<ClaimedGameProfileDirtyRow>> {
return DB.transaction().execute(async (trx) => {
const result = await sql<ClaimedGameProfileDirtyRow>`
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;
});
}
@@ -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<number> {
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<number>,
): Promise<number> {
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<void>((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<void>((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);
});
});
@@ -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<number> {
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<number> {
const groups = new Map<
string,
{ chartIDs: Set<string>; game: GameGroup; playtype: LEGACY_Playtype; userID: integer }
{
chartIDs: Set<string>;
game: GameGroup;
playtype: LEGACY_Playtype;
runStartedAt: Awaited<ReturnType<typeof newCalculationRunStartedAt>>;
userID: integer;
}
>();
for (const row of rows) {
@@ -63,7 +70,13 @@ export async function drainPbDirty(): Promise<number> {
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<number> {
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<number> {
}
/**
* 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<number> {
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<number> {
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<number> {
.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<number> {
.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<number> {
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<number> {
}
/**
* 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<number> {
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<number> {
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.`);
@@ -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);
});
});
@@ -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<string>,
log: KtLogger,
options?: ProcessPBsOptions,
): Promise<void> {
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);
}
});
@@ -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<TGame extends V3Game = V3Game> = 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<Database>,
pbDoc: PBScoreDocumentNoRank,
): Promise<void> {
runStartedAt: CalculationRunStartedAt,
): Promise<boolean> {
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;
}
@@ -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<V3Game> | null,
log: KtLogger,
options?: ClassProcessOptions,
options?: UpdateUsersGamePlaytypeStatsOptions,
): Promise<Array<ClassDelta>> {
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();
@@ -18,7 +18,12 @@ export async function RecalcAllScores(): Promise<void> {
*/
export async function UpdateAllPBs(): Promise<void> {
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();
}