97 lines
2.7 KiB
JavaScript
97 lines
2.7 KiB
JavaScript
// Task runner with retry, timeout, and structured logging
|
|
import { createLogger } from "./logger.mjs";
|
|
import { PrismaClient } from "@prisma/client";
|
|
import { PrismaMariaDb } from "@prisma/adapter-mariadb";
|
|
|
|
const logger = createLogger("task-runner");
|
|
|
|
function createPrisma() {
|
|
const base = process.env.DATABASE_URL.replace("mysql://", "mariadb://");
|
|
const sep = base.includes("?") ? "&" : "?";
|
|
const connectionString = `${base}${sep}connection_limit=3&pool_timeout=30`;
|
|
const adapter = new PrismaMariaDb(connectionString);
|
|
return new PrismaClient({ adapter });
|
|
}
|
|
|
|
async function runTask(taskKey, taskName, fn, { timeoutMs = 30 * 60 * 1000, maxRetries = 2 } = {}) {
|
|
const prisma = createPrisma();
|
|
const startedAt = new Date();
|
|
let attempt = 0;
|
|
let lastError;
|
|
|
|
logger.info(`Starting task: ${taskName}`, { taskKey });
|
|
|
|
while (attempt <= maxRetries) {
|
|
if (attempt > 0) {
|
|
const delay = Math.min(5000 * Math.pow(2, attempt - 1), 30000);
|
|
logger.warn(`Retry ${attempt}/${maxRetries} for ${taskName} after ${delay}ms`, { taskKey });
|
|
await new Promise((r) => setTimeout(r, delay));
|
|
}
|
|
|
|
try {
|
|
// Timeout wrapper
|
|
const timeoutPromise = new Promise((_, reject) =>
|
|
setTimeout(() => reject(new Error(`Task timed out after ${timeoutMs / 60000} minutes`)), timeoutMs)
|
|
);
|
|
|
|
const result = await Promise.race([fn(prisma), timeoutPromise]);
|
|
|
|
// Log success
|
|
await prisma.taskLog.create({
|
|
data: {
|
|
taskKey,
|
|
taskName,
|
|
status: "success",
|
|
startedAt,
|
|
finishedAt: new Date(),
|
|
duration: Date.now() - startedAt.getTime(),
|
|
triggerBy: "cron",
|
|
result: result ? JSON.parse(JSON.stringify(result).slice(0, 2000)) : undefined,
|
|
},
|
|
}).catch(() => {});
|
|
|
|
logger.info(`Task completed: ${taskName}`, {
|
|
taskKey,
|
|
durationMs: Date.now() - startedAt.getTime(),
|
|
attempt,
|
|
});
|
|
|
|
await prisma.$disconnect();
|
|
return result;
|
|
} catch (error) {
|
|
lastError = error;
|
|
attempt++;
|
|
logger.error(`Task attempt failed: ${taskName}`, {
|
|
taskKey,
|
|
attempt,
|
|
error: error.message,
|
|
});
|
|
}
|
|
}
|
|
|
|
// All retries exhausted
|
|
const errorMsg = (lastError?.message || String(lastError)).slice(0, 1000);
|
|
await prisma.taskLog.create({
|
|
data: {
|
|
taskKey,
|
|
taskName,
|
|
status: "failed",
|
|
startedAt,
|
|
finishedAt: new Date(),
|
|
duration: Date.now() - startedAt.getTime(),
|
|
triggerBy: "cron",
|
|
error: errorMsg,
|
|
},
|
|
}).catch(() => {});
|
|
|
|
logger.error(`Task failed after ${maxRetries + 1} attempts: ${taskName}`, {
|
|
taskKey,
|
|
error: errorMsg,
|
|
});
|
|
|
|
await prisma.$disconnect();
|
|
throw lastError;
|
|
}
|
|
|
|
export { runTask, createLogger };
|