concurrent parse limit

This commit is contained in:
Sergey Yarkov 2026-03-05 08:36:06 +03:00
parent 9e7daddf51
commit cd0a702270
3 changed files with 95 additions and 78 deletions

View File

@ -6,7 +6,7 @@ import { clients } from "#root/index.js";
import { serializeError } from "serialize-error"; import { serializeError } from "serialize-error";
const sheetName = "Отчеты"; const sheetName = "Отчеты";
const limit = pLimit(10); const MAX_CONCURRENT_PARSE_CLIENT = Number.parseInt(process.env.REPORT_CONCURRENT_PARSE_CLIENTS || "3", 10);
/** @type {xlsx.ParsingOptions} */ /** @type {xlsx.ParsingOptions} */
const xlsxParseOptions = { const xlsxParseOptions = {
@ -27,102 +27,120 @@ const xlsxParseOptions = {
*/ */
export default async (ctx) => { export default async (ctx) => {
const { res, params } = ctx; const { res, params } = ctx;
const clientId = params.get("clientId");
if (!ctx.data) throw new BadRequestError("Cannot proceed uploaded file");
if (!clientId) throw new BadRequestError('"clientId" parameter is required');
const [, files] = ctx.data; const [, files] = ctx.data;
const limit = pLimit(10);
const clientId = params.get("clientId");
const excelFilePath = files?.report?.at(0)?.filepath; const excelFilePath = files?.report?.at(0)?.filepath;
const sseClientRes = clients.get(clientId);
if (!sseClientRes) throw new ServiceUnavailableError("Cannot proceed request due to unknown client");
let isSseClientClosed = false;
sseClientRes.res.on("close", () => (isSseClientClosed = true));
/** @type {ReportModel[]} */
const reports = []; const reports = [];
let isSseClientClosed = false;
let parseResultCount = { equipments: 0, applicants: 0, executors: 0, reportsFullfilled: 0 }; let parseResultCount = { equipments: 0, applicants: 0, executors: 0, reportsFullfilled: 0 };
if (!ctx.data) {
throw new BadRequestError("Cannot proceed uploaded file");
}
if (!clientId) {
throw new BadRequestError('"clientId" parameter is required');
}
const sseClientCtx = clients.get(clientId);
const parseProcessClients = clients
.values()
.filter((sseClientCtx) => sseClientCtx.local?.isParsing)
.toArray();
if (parseProcessClients.length >= MAX_CONCURRENT_PARSE_CLIENT) {
throw new ServiceUnavailableError(`Maximum concurrent parse clients reached limit`);
}
if (!sseClientCtx) {
throw new ServiceUnavailableError("Cannot proceed request due to unknown client");
}
if (sseClientCtx.local.isParsing) {
throw new ServiceUnavailableError("Your request on parsing is processing");
}
sseClientCtx.res.res.on("close", () => (isSseClientClosed = true));
if (excelFilePath) { if (excelFilePath) {
const workbook = xlsx.readFile(excelFilePath, xlsxParseOptions); try {
const sheet = workbook.Sheets[sheetName]; const workbook = xlsx.readFile(excelFilePath, xlsxParseOptions);
const sheet = workbook.Sheets[sheetName];
if (!sheet) throw new Error("Invalid reports workbook.", { cause: { statusCode: 400 } }); if (!sheet) throw new Error("Invalid reports workbook.", { cause: { statusCode: 400 } });
res.sendJson({ message: "OK" }); const data = xlsx.utils.sheet_to_json(sheet, { raw: false, blankrows: true });
const data = xlsx.utils.sheet_to_json(sheet, { raw: false, blankrows: true }); res.sendJson({ message: "OK" });
sseClientCtx.local.isParsing = true;
for (let i = 0; i < data.length; i++) { for (let i = 0; i < data.length; i++) {
const cols = Object.keys(data[i]); const cols = Object.keys(data[i]);
if (i === 0) continue; if (i === 0) continue;
/** Несинхронизированные помеченные отчеты */ /** Несинхронизированные помеченные отчеты */
if ("__EMPTY_3" in data[i] && data[i]["__EMPTY_3"] === "r") { if ("__EMPTY_3" in data[i] && data[i]["__EMPTY_3"] === "r") {
reports.push({
date: data[i][cols[0]],
equipment: data[i][cols[1]],
reason_call: data[i][cols[3]].split("\r\r\n")[1] || "",
job_description: data[i][cols[4]].split("\r\r\n")[1] || "",
root_cause: data[i]["__EMPTY_2"] || "",
applicantName: data[i][cols[3]].split("\r\r\n")[0] || "",
executorNames: data[i][cols[4]].split("\r\r\n")[0] || "",
isMarked: true
});
continue;
}
/** Остальные отчеты */
reports.push({ reports.push({
date: data[i][cols[0]], date: data[i][cols[0]],
equipment: data[i][cols[1]], equipment: data[i][cols[1]],
reason_call: data[i][cols[3]].split("\r\r\n")[1] || "", reason_call: parseReasonCallAndJobDesc(data[i][cols[3]]).trim(),
job_description: data[i][cols[4]].split("\r\r\n")[1] || "", job_description: parseReasonCallAndJobDesc(data[i][cols[4]]).trim(),
root_cause: data[i]["__EMPTY_2"] || "", root_cause: data[i]["__EMPTY_2"] || "",
applicantName: data[i][cols[3]].split("\r\r\n")[0] || "", applicantName: parseNames(data[i][cols[3]]).trim(),
executorNames: data[i][cols[4]].split("\r\r\n")[0] || "", executorNames: parseNames(data[i][cols[4]]).trim()
isMarked: true
}); });
continue;
} }
} catch (error) {
/** Остальные отчеты */ sseClientCtx.local.isParsing = false;
reports.push({ throw error;
date: data[i][cols[0]],
equipment: data[i][cols[1]],
reason_call: parseReasonCallAndJobDesc(data[i][cols[3]]).trim(),
job_description: parseReasonCallAndJobDesc(data[i][cols[4]]).trim(),
root_cause: data[i]["__EMPTY_2"] || "",
applicantName: parseNames(data[i][cols[3]]).trim(),
executorNames: parseNames(data[i][cols[4]]).trim()
});
} }
const requests = reports.map((report) => try {
limit(() => { const requests = reports.map((report) =>
return checkReportAssigments( limit(() => {
report, return checkReportAssigments(
(checkedReport) => { report,
const isEquipmentLinked = checkedReport.equipment instanceof Object; (checkedReport) => {
const isApplicantLinked = checkedReport.applicant instanceof Object; const isEquipmentLinked = checkedReport.equipment instanceof Object;
const isExecutorsLinked = checkedReport.executors && checkedReport.executors.length !== 0; const isApplicantLinked = checkedReport.applicant instanceof Object;
const isExecutorsLinked = checkedReport.executors && checkedReport.executors.length !== 0;
if (isEquipmentLinked) parseResultCount.equipments++; if (isEquipmentLinked) parseResultCount.equipments++;
if (isApplicantLinked) parseResultCount.applicants++; if (isApplicantLinked) parseResultCount.applicants++;
if (isExecutorsLinked) parseResultCount.executors += checkedReport.executors?.length || 0; if (isExecutorsLinked) parseResultCount.executors += checkedReport.executors?.length || 0;
if (isEquipmentLinked && isApplicantLinked && isExecutorsLinked) parseResultCount.reportsFullfilled++; if (isEquipmentLinked && isApplicantLinked && isExecutorsLinked) parseResultCount.reportsFullfilled++;
sseClientRes.sendSSEJson(checkedReport, "progress"); sseClientCtx.res.sendSSEJson(checkedReport, "progress");
}, },
isSseClientClosed isSseClientClosed
); );
}) })
); );
await Promise.all(requests);
await Promise.all(requests) sseClientCtx.res.sendSSEJson(parseResultCount, "done");
.then(() => { } catch (error) {
console.log(parseResultCount); console.error(error);
sseClientRes.sendSSEJson(parseResultCount, "done"); sseClientCtx.res.sendSSEJson({ error: serializeError(error) }, "error");
sseClientRes.res.end(); } finally {
}) sseClientCtx.local.isParsing = false;
.catch((error) => { }
console.error(error);
sseClientRes.sendSSEJson({ error: serializeError(error) }, "error");
sseClientRes.res.end();
})
.finally(() => res.res.end());
} else { } else {
throw new BadRequestError(); throw new BadRequestError();
} }

View File

@ -12,7 +12,7 @@ export default (ctx) => {
const { res } = ctx; const { res } = ctx;
const clientId = crypto.randomUUID(); const clientId = crypto.randomUUID();
clients.set(clientId, ctx.res); clients.set(clientId, ctx);
res.res.setHeaders( res.res.setHeaders(
new Headers({ new Headers({
@ -26,7 +26,6 @@ export default (ctx) => {
res.res.on("close", () => { res.res.on("close", () => {
clients.delete(clientId); clients.delete(clientId);
console.log(clients);
res.res.end(); res.res.end();
}); });
}; };

View File

@ -6,7 +6,7 @@ import db from "./db.js";
const SERVER_PORT = 3000; const SERVER_PORT = 3000;
/** @type {Map<string, import("./http-server/types/http-server.js").HttpResponse>} */ /** @type {Map<string, import("./http-server/types/http-server.js").HttpContext>} */
export const clients = new Map(); export const clients = new Map();
/** /**