Merge pull request #649 from TNG-dev/zkldi/issue-600

This commit is contained in:
zkldi
2022-02-01 02:28:48 +00:00
committed by GitHub
3 changed files with 92 additions and 35 deletions
@@ -1,4 +1,5 @@
import db from "external/mongo/db";
import { rootLogger } from "lib/logger/logger";
import { integer } from "tachi-common";
/**
@@ -1,10 +1,15 @@
/* eslint-disable no-await-in-loop */
import { ServerConfig } from "lib/setup/config";
import { ScoreImportJobData } from "../worker/types";
import { GetInputParser } from "./common/get-input-parser";
import ScoreImportMain from "./score-importing/score-import-main";
import { ImportTypes, ImportDocument } from "tachi-common";
import { ImportTypes, ImportDocument, integer } from "tachi-common";
import ScoreImportQueue, { ScoreImportQueueEvents } from "../worker/queue";
import ScoreImportFatalError from "./score-importing/score-import-error";
import { Sleep } from "utils/misc";
import CreateLogCtx from "lib/logger/logger";
const logger = CreateLogCtx(__filename);
/**
* Makes a score import given ScoreImportJobData.
@@ -21,17 +26,53 @@ export async function MakeScoreImport<I extends ImportTypes>(
jobData: ScoreImportJobData<I>
): Promise<ImportDocument> {
if (ServerConfig.USE_EXTERNAL_SCORE_IMPORT_WORKER && process.env.IS_JOB === undefined) {
const job = await ScoreImportQueue.add(`Import ${jobData.importID}`, jobData, {
jobId: jobData.importID,
});
let timesAttempted = 1;
const data = await job.waitUntilFinished(ScoreImportQueueEvents);
// There's no chance this thing goes on 7 times.
// if it does, this import has been trying for the past 6 hours or so.
while (timesAttempted <= 7) {
const job = await ScoreImportQueue.add(
`Import ${jobData.importID}${timesAttempted > 0 ? ` (TRY${timesAttempted})` : ""}`,
jobData,
{
jobId: `${jobData.importID}:TRY${timesAttempted}`,
}
);
if (data.success) {
return data.importDocument;
} else {
throw new ScoreImportFatalError(data.statusCode, data.description);
const data = await job.waitUntilFinished(ScoreImportQueueEvents);
if (data.success) {
return data.importDocument;
} else if (data.statusCode !== 409) {
throw new ScoreImportFatalError(data.statusCode, data.description);
}
const backoff = ExponentialBackoff(timesAttempted - 1);
logger.info(
`User ${jobData.userID} already had an import ongoing. (${
jobData.importID
}) Backing off for ${(backoff / 1_000).toFixed(2)} seconds.`
);
// If we get here, we were 409'd and the user already has an ongoing
// import.
// In the interest of not just throwing scores away, we'll back off a bit
// and then restart the job.
await Sleep(backoff);
timesAttempted++;
}
logger.error(
`User ${jobData.userID} didn't get an import through in around 6 hours. Has their lock gotten stuck?`,
jobData
);
throw new ScoreImportFatalError(
409,
"Couldn't get an import through in the past 6 hours, at all."
);
} else {
const InputParser = GetInputParser(jobData);
@@ -44,3 +85,16 @@ export async function MakeScoreImport<I extends ImportTypes>(
);
}
}
function ExponentialBackoff(exponent: integer) {
// n | backoff
// 0 | 4 Seconds
// 1 | 16 Seconds
// 2 | 64 Seconds
// 3 | 256 Seconds
// 4 | 1024 Seconds
// ...
// ends at 7, which is around 4 hours.
return 1000 * 4 ** exponent;
}
@@ -1,5 +1,5 @@
import db from "external/mongo/db";
import { KtLogger } from "lib/logger/logger";
import { KtLogger, rootLogger } from "lib/logger/logger";
import { ScoreImportJob } from "lib/score-import/worker/types";
import {
Game,
@@ -12,7 +12,7 @@ import {
PublicUserDocument,
GetGameConfig,
} from "tachi-common";
import { GetMillisecondsSince } from "utils/misc";
import { GetMillisecondsSince, Sleep } from "utils/misc";
import { GetUserWithID } from "utils/user";
import { ConverterFunction, ImportInputParser } from "../../import-types/common/types";
import { Converters } from "../../import-types/converters";
@@ -50,32 +50,34 @@ export default async function ScoreImportMain<D, C>(
);
}
let logger;
if (!providedLogger) {
// If they weren't given to us -
// we create an "import logger".
// this holds a reference to the user's name, ID, and type
// of score import for any future debugging.
logger = CreateScoreLogger(user, importID, importType);
logger.debug("Received import request.");
} else {
logger = providedLogger;
}
const hasNoOngoingImport = await CheckAndSetOngoingImportLock(user.id);
if (hasNoOngoingImport) {
logger.info(`User ${userID} made an import while they had one ongoing.`);
// @danger
// Throwing away an import if the user already has one outgoing is *bad*, as in the case
// of degraded performance we might just start throwing scores away.
// Under normal circumstances, there is no scenario where a user would have two ongoing
// imports at the same time - even if they were using single-score imports on a 5 second
// chart, as each score import takes only around ~10-15milliseconds.
throw new ScoreImportFatalError(409, "This user already has an ongoing import.");
}
try {
const hasNoOngoingImport = await CheckAndSetOngoingImportLock(user.id);
if (hasNoOngoingImport) {
// @danger
// Throwing away an import if the user already has one outgoing is *bad*, as in the case
// of degraded performance we might just start throwing scores away.
// Under normal circumstances, there is no scenario where a user would have two ongoing
// imports at the same time - even if they were using single-score imports on a 5 second
// chart, as each score import takes only around ~10-15milliseconds.
throw new ScoreImportFatalError(409, "This user already has an ongoing import.");
}
const timeStarted = Date.now();
let logger;
if (!providedLogger) {
// If they weren't given to us -
// we create an "import logger".
// this holds a reference to the user's name, ID, and type
// of score import for any future debugging.
logger = CreateScoreLogger(user, importID, importType);
logger.debug("Received import request.");
} else {
logger = providedLogger;
}
SetJobProgress(job, "Parsing score data.");