diff --git a/server/src/lib/score-import/framework/score-importing/insert-score.test.ts b/server/src/lib/score-import/framework/score-importing/insert-score.test.ts index 3eee2b36e..0c9feaebb 100644 --- a/server/src/lib/score-import/framework/score-importing/insert-score.test.ts +++ b/server/src/lib/score-import/framework/score-importing/insert-score.test.ts @@ -16,7 +16,10 @@ t.test("#QueueScoreInsert, #InsertQueue", async (t) => { t.test("Single Queue Test", async (t) => { // fake score doc - const res = await QueueScoreInsert({ scoreID: "foo" } as unknown as ScoreDocument); + const res = await QueueScoreInsert({ + scoreID: "foo", + userID: 1, + } as unknown as ScoreDocument); t.equal( res, @@ -25,7 +28,7 @@ t.test("#QueueScoreInsert, #InsertQueue", async (t) => { ); // this is the best way to get the size of the queue - const flushSize = await InsertQueue(); + const flushSize = await InsertQueue(1); t.equal(flushSize, 1, "QueueScoreInsert should append the score to the queue."); @@ -42,17 +45,22 @@ t.test("#QueueScoreInsert, #InsertQueue", async (t) => { t.end(); }); - let r = await InsertQueue(); // flush queue just incase former test fails. + let r = await InsertQueue(1); // flush queue just incase former test fails. t.test("Queue Overflow Test", async (t) => { for (let i = 0; i < 499; i++) { // eslint-disable-next-line no-await-in-loop - await QueueScoreInsert({ scoreID: i, chartID: "test" } as unknown as ScoreDocument); + await QueueScoreInsert({ + scoreID: i, + chartID: "test", + userID: 1, + } as unknown as ScoreDocument); } const overflowRes = await QueueScoreInsert({ scoreID: "foo", chartID: "test", + userID: 1, } as unknown as ScoreDocument); t.equal( @@ -61,7 +69,7 @@ t.test("#QueueScoreInsert, #InsertQueue", async (t) => { "Appending 500 items to the queue should result in them being inserted." ); - const flushRes = await InsertQueue(); + const flushRes = await InsertQueue(1); t.equal(flushRes, 0, "The queue should now be empty."); @@ -79,20 +87,25 @@ t.test("#QueueScoreInsert, #InsertQueue", async (t) => { t.end(); }); - r = await InsertQueue(); // flush queue just incase former test fails. + r = await InsertQueue(1); // flush queue just incase former test fails. t.equal(r, 0, "Queue should be empty after test."); t.test("Queue Dedupe Test", async (t) => { - await QueueScoreInsert({ scoreID: 1, chartID: "foo" } as unknown as ScoreDocument); + await QueueScoreInsert({ + scoreID: 1, + chartID: "foo", + userID: 1, + } as unknown as ScoreDocument); const r2 = await QueueScoreInsert({ scoreID: 1, chartID: "foo", + userID: 1, } as unknown as ScoreDocument); t.equal(r2, null, "Should return null when a duplicate scoreID is submitted"); - const flushRes = await InsertQueue(); + const flushRes = await InsertQueue(1); t.equal(flushRes, 1, "Should flush 1 score document"); @@ -103,5 +116,37 @@ t.test("#QueueScoreInsert, #InsertQueue", async (t) => { t.equal(dbRes.length, 1, "Should only insert one document"); }); + t.test("Should give separate users separate queues.", async (t) => { + await QueueScoreInsert({ + scoreID: "1", + chartID: "foo", + userID: 1, + } as unknown as ScoreDocument); + await QueueScoreInsert({ + scoreID: "2", + chartID: "foo", + userID: 2, + } as unknown as ScoreDocument); + + const r1 = await InsertQueue(1); + t.equal(r1, 1, "Queue for userID 1 should have length 1."); + + const r2 = await InsertQueue(2); + t.equal(r2, 1, "Queue for userID 2 should also have length 1."); + + t.end(); + }); + + t.test("Should not throw if InsertQueue is called on an empty queue.", async (t) => { + try { + await InsertQueue(1); + t.pass("Did not throw when inserting an empty queue."); + } catch (err) { + t.fail(err); + } + + t.end(); + }); + t.end(); }); diff --git a/server/src/lib/score-import/framework/score-importing/insert-score.ts b/server/src/lib/score-import/framework/score-importing/insert-score.ts index 4a21077ec..37738f5ad 100644 --- a/server/src/lib/score-import/framework/score-importing/insert-score.ts +++ b/server/src/lib/score-import/framework/score-importing/insert-score.ts @@ -1,46 +1,68 @@ -import { ScoreDocument } from "tachi-common"; +import { ScoreDocument, integer } from "tachi-common"; import db from "external/mongo/db"; import CreateLogCtx from "lib/logger/logger"; const logger = CreateLogCtx(__filename); -const ScoreQueue: ScoreDocument[] = []; -export let ScoreIDs: Set = new Set(); const MAX_PIPELINE_LENGTH = 500; +interface ScoreQueue { + queue: ScoreDocument[]; + scoreIDSet: Set; +} + +const ScoreQueues: Record = {}; + /** - * Adds a score to a queue to be inserted in batch to the database. - * @param score - The score document to queue. - * @returns True on success, The amount of scores inserted on auto-pipeline-flush, and null if - * the score provided is already loaded. + * Returns this user's score queue. A score queue is a temporary place scores are saved + * so that they can be inserted into the database in bulk. + * + * This massively improves performance on large imports instead of constantly running single imports. + * + * If a score queue does not exist for the user, one is created. */ -export function QueueScoreInsert(score: ScoreDocument) { - if (ScoreIDs.has(score.scoreID)) { - // skip - logger.verbose(`Triggered skip for ID ${score.scoreID}`); - return null; +function GetOrSetScoreQueue(userID: integer) { + const queue = ScoreQueues[userID]; + + if (!queue) { + return SetScoreQueue(userID); } - ScoreQueue.push(score); - ScoreIDs.add(score.scoreID); + return queue; +} - if (ScoreQueue.length >= MAX_PIPELINE_LENGTH) { - logger.verbose(`Triggered pipeline flush with len ${ScoreQueue.length}.`); - return InsertQueue(); - } +export function GetScoreQueueMaybe(userID: integer): ScoreQueue | undefined { + return ScoreQueues[userID]; +} - return true; +function SetScoreQueue(userID: integer) { + const queue: ScoreQueue = { + queue: [], + scoreIDSet: new Set(), + }; + + ScoreQueues[userID] = queue; + + return queue; } /** - * Bulk inserts the entire Queue. - * @warn Be cautious of inducing race conditions when using this function. + * Adds a new score to the given queue. */ -export async function InsertQueue() { - const temp = ScoreQueue.splice(0); - if (temp.length !== 0) { - ScoreIDs = new Set(); +function AddToScoreQueue(scoreQueue: ScoreQueue, score: ScoreDocument) { + scoreQueue.queue.push(score); + scoreQueue.scoreIDSet.add(score.scoreID); +} + +export async function InsertQueue(userID: integer) { + const scoreQueue = GetOrSetScoreQueue(userID); + + const queuedScores = scoreQueue.queue.splice(0); + + if (queuedScores.length !== 0) { + delete ScoreQueues[userID]; + try { - await db.scores.insert(temp); + await db.scores.insert(queuedScores); } catch (err) { logger.warn( `Triggered duplicate key protection. Race condition protected against, but this is not good.` @@ -49,5 +71,32 @@ export async function InsertQueue() { } } - return temp.length; + return queuedScores.length; +} + +/** + * Adds a score to a queue to be inserted in batch to the database. + * @param score - The score document to queue. + * @returns True on success, The amount of scores inserted on auto-pipeline-flush, and null if + * the score provided is already loaded. + */ +export function QueueScoreInsert(score: ScoreDocument) { + const scoreQueue = GetOrSetScoreQueue(score.userID); + + if (scoreQueue.scoreIDSet.has(score.scoreID)) { + // skip + logger.verbose(`Score ID ${score.scoreID} was already queued to be imported.`); + return null; + } + + AddToScoreQueue(scoreQueue, score); + + logger.debug(`ScoreQueue for ${score.userID} is now at ${scoreQueue.queue.length}.`); + + if (scoreQueue.queue.length >= MAX_PIPELINE_LENGTH) { + logger.verbose(`Triggered pipeline flush with len ${scoreQueue.queue.length}.`); + return InsertQueue(score.userID); + } + + return true; } diff --git a/server/src/lib/score-import/framework/score-importing/score-importing.ts b/server/src/lib/score-import/framework/score-importing/score-importing.ts index 99d600d30..8571f2b22 100644 --- a/server/src/lib/score-import/framework/score-importing/score-importing.ts +++ b/server/src/lib/score-import/framework/score-importing/score-importing.ts @@ -25,7 +25,7 @@ import { import { DryScore } from "../common/types"; import { OrphanScore } from "../orphans/orphans"; import { HydrateScore } from "./hydrate-score"; -import { InsertQueue, QueueScoreInsert, ScoreIDs } from "./insert-score"; +import { GetScoreQueueMaybe, InsertQueue, QueueScoreInsert } from "./insert-score"; import { CreateScoreID } from "./score-id"; /** @@ -92,7 +92,7 @@ export async function ImportAllIterableData( // Flush the score queue out after finishing most of the import. This ensures no scores get left in the // queue. - const emptied = await InsertQueue(); + const emptied = await InsertQueue(userID); if (emptied) { logger.verbose(`Emptied ${emptied} documents from score queue.`); @@ -309,7 +309,8 @@ async function HydrateAndInsertScore( return null; } - if (ScoreIDs.has(scoreID)) { + // If this users score queue + if (GetScoreQueueMaybe(userID)?.scoreIDSet.has(scoreID)) { logger.verbose(`Skipped score.`); return null; } @@ -323,7 +324,7 @@ async function HydrateAndInsertScore( res = await QueueScoreInsert(score); } - // emergency state - this is a last resort for avoiding doubled imports + // this is a last resort for avoiding doubled imports if (res === null) { logger.verbose(`Skipped score - Race Condition protection triggered.`); return null;