mirror of
https://github.com/zkldi/Tachi.git
synced 2026-10-09 15:19:50 +03:00
fix: ir not queuing scores up correctly
This commit is contained in:
@@ -1,3 +1,4 @@
|
||||
import type { ParserArguments } from "#lib/score-import/worker/types";
|
||||
import type {
|
||||
ImportDocument,
|
||||
ImportTypes,
|
||||
@@ -7,11 +8,12 @@ import type {
|
||||
} from "tachi-common";
|
||||
|
||||
import { log } from "#lib/log/log";
|
||||
import { AwaitScoreImportWorkerResult } from "#lib/score-import/worker/await-worker-result";
|
||||
import { EnqueueScoreImportJob } from "#lib/score-import/worker/enqueue-pg";
|
||||
import { ServerConfig } from "#lib/setup/config";
|
||||
import { Random20Hex } from "#utils/misc";
|
||||
import { ExpectedErr } from "bliss";
|
||||
|
||||
import type { ParserArguments } from "../worker/types";
|
||||
|
||||
import { MakeScoreImport } from "./score-import";
|
||||
import ScoreImportFatalError from "./score-importing/score-import-error";
|
||||
|
||||
@@ -36,13 +38,21 @@ export async function ExpressWrappedScoreImportMain<I extends ImportTypes>(
|
||||
log.debug("Received import request.");
|
||||
|
||||
try {
|
||||
const res = await MakeScoreImport({
|
||||
const jobData = {
|
||||
importID,
|
||||
importType,
|
||||
userIntent,
|
||||
userID,
|
||||
parserArguments,
|
||||
});
|
||||
};
|
||||
|
||||
let res: ImportDocument;
|
||||
if (ServerConfig.USE_EXTERNAL_SCORE_IMPORT_WORKER) {
|
||||
await EnqueueScoreImportJob(jobData);
|
||||
res = await AwaitScoreImportWorkerResult(importID);
|
||||
} else {
|
||||
res = await MakeScoreImport(jobData);
|
||||
}
|
||||
|
||||
return {
|
||||
statusCode: 200,
|
||||
@@ -54,7 +64,7 @@ export async function ExpressWrappedScoreImportMain<I extends ImportTypes>(
|
||||
};
|
||||
} catch (err) {
|
||||
// `ACTION_ScoreImport` throws `ExpectedErr` (mapped from `ScoreImportFatalError`); the
|
||||
// external-worker guard in `MakeScoreImport` still throws `ScoreImportFatalError`.
|
||||
// await helper throws `ScoreImportFatalError` when the worker records a failed import.
|
||||
if (ExpectedErr.is(err) || err instanceof ScoreImportFatalError) {
|
||||
const description = ExpectedErr.is(err) ? err.reason : err.message;
|
||||
const statusCode = ExpectedErr.is(err) ? err.code : err.statusCode;
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
import type { ImportDocument } from "tachi-common";
|
||||
|
||||
import {
|
||||
GetImportTrackerByImportId,
|
||||
LoadImportDocumentById,
|
||||
} from "#lib/db-formats/import-document";
|
||||
import ScoreImportFatalError from "#lib/score-import/framework/score-importing/score-import-error";
|
||||
import { Sleep } from "#utils/misc";
|
||||
|
||||
const POLL_MS = 100;
|
||||
const MAX_WAIT_MS = 600_000;
|
||||
|
||||
/**
|
||||
* After {@link EnqueueScoreImportJob}, blocks until the worker finishes the import
|
||||
* and a completed {@link ImportDocument} is available, the tracker records failure,
|
||||
* or the wait times out.
|
||||
*
|
||||
* Used by {@link ExpressWrappedScoreImportMain} so IR and similar callers keep
|
||||
* synchronous response semantics while processing still runs on the job-queue worker.
|
||||
*/
|
||||
export async function AwaitScoreImportWorkerResult(importID: string): Promise<ImportDocument> {
|
||||
const deadline = Date.now() + MAX_WAIT_MS;
|
||||
|
||||
while (Date.now() < deadline) {
|
||||
const doc = await LoadImportDocumentById(importID);
|
||||
if (doc) {
|
||||
return doc;
|
||||
}
|
||||
|
||||
const tracker = await GetImportTrackerByImportId(importID);
|
||||
if (tracker?.type === "FAILED") {
|
||||
const code = tracker.error.statusCode ?? 500;
|
||||
throw new ScoreImportFatalError(code, tracker.error.message);
|
||||
}
|
||||
|
||||
await Sleep(POLL_MS);
|
||||
}
|
||||
|
||||
throw new ScoreImportFatalError(
|
||||
504,
|
||||
"Score import timed out waiting for the import worker to finish.",
|
||||
);
|
||||
}
|
||||
Reference in New Issue
Block a user