feat: 补充注释
This commit is contained in:
@@ -13,26 +13,41 @@ import {
|
|||||||
import { uploadImageFromUrl } from "~~/server/utils/lsky";
|
import { uploadImageFromUrl } from "~~/server/utils/lsky";
|
||||||
|
|
||||||
interface IImageArchiveWorkerConfig {
|
interface IImageArchiveWorkerConfig {
|
||||||
|
/** 同一 Node 进程内允许同时执行的归档任务数,控制下载/上传带来的内存和网络峰值 */
|
||||||
concurrency: number;
|
concurrency: number;
|
||||||
|
/** 单条归档任务最多尝试次数,超过后记录为最终失败 */
|
||||||
maxAttempts: number;
|
maxAttempts: number;
|
||||||
|
/** RUNNING 任务超过该时间未完成时视为锁失效,可被其他 worker 重新领取 */
|
||||||
lockTtlMs: number;
|
lockTtlMs: number;
|
||||||
|
/** 没有可执行任务时的轮询间隔,避免空转打数据库 */
|
||||||
pollMs: number;
|
pollMs: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
interface IImageArchiveWorkerState {
|
interface IImageArchiveWorkerState {
|
||||||
|
/** 当前进程正在执行的归档任务数量 */
|
||||||
activeCount: number;
|
activeCount: number;
|
||||||
|
/** 防止同一进程内多个扫描循环重叠领取任务 */
|
||||||
draining: boolean;
|
draining: boolean;
|
||||||
|
/** 标记插件是否已经启动过 worker,避免开发热更新重复初始化 */
|
||||||
started: boolean;
|
started: boolean;
|
||||||
|
/** 下一次扫描的定时器句柄,用于新任务入队时提前唤醒 */
|
||||||
timer: ReturnType<typeof setTimeout> | null;
|
timer: ReturnType<typeof setTimeout> | null;
|
||||||
|
/** 写入数据库锁的 worker 标识,用于日志排查和多实例抢占 */
|
||||||
workerId: string;
|
workerId: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** 默认只开 2 个归档槽位,避免多用户同时归档时挤占生图主链路资源 */
|
||||||
const DEFAULT_CONCURRENCY = 2;
|
const DEFAULT_CONCURRENCY = 2;
|
||||||
|
/** 默认最多重试 3 次,覆盖短暂网络抖动,同时避免坏任务长期循环 */
|
||||||
const DEFAULT_MAX_ATTEMPTS = 3;
|
const DEFAULT_MAX_ATTEMPTS = 3;
|
||||||
|
/** 锁超时 10 分钟,进程崩溃或重启后任务可以自动恢复 */
|
||||||
const DEFAULT_LOCK_TTL_MS = 10 * 60 * 1000;
|
const DEFAULT_LOCK_TTL_MS = 10 * 60 * 1000;
|
||||||
|
/** 空队列每 15 秒扫一次;生图成功会主动唤醒,不依赖纯轮询 */
|
||||||
const DEFAULT_POLL_MS = 15 * 1000;
|
const DEFAULT_POLL_MS = 15 * 1000;
|
||||||
|
/** 失败退避节奏:第一次 30 秒,第二次 2 分钟,第三次及以后 10 分钟 */
|
||||||
const RETRY_DELAYS_MS = [30 * 1000, 2 * 60 * 1000, 10 * 60 * 1000];
|
const RETRY_DELAYS_MS = [30 * 1000, 2 * 60 * 1000, 10 * 60 * 1000];
|
||||||
|
|
||||||
|
// Nuxt 开发热更新可能重复加载模块;把状态挂到 globalThis,避免重复 worker 抢同一批任务。
|
||||||
const globalForArchiveWorker = globalThis as unknown as {
|
const globalForArchiveWorker = globalThis as unknown as {
|
||||||
imageArchiveWorkerState?: IImageArchiveWorkerState;
|
imageArchiveWorkerState?: IImageArchiveWorkerState;
|
||||||
};
|
};
|
||||||
@@ -66,7 +81,9 @@ export const wakeImageArchiveWorker = () => {
|
|||||||
scheduleDrain(0);
|
scheduleDrain(0);
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/** 扫描数据库中的可执行归档任务,按当前空闲槽位领取并启动异步执行 */
|
||||||
const drainImageArchiveQueue = async () => {
|
const drainImageArchiveQueue = async () => {
|
||||||
|
// 归档扫描只负责领取任务和启动执行器;真正上传在 runImageArchiveTask 中异步完成。
|
||||||
if (workerState.draining) return;
|
if (workerState.draining) return;
|
||||||
|
|
||||||
workerState.draining = true;
|
workerState.draining = true;
|
||||||
@@ -76,11 +93,13 @@ const drainImageArchiveQueue = async () => {
|
|||||||
const config = getImageArchiveWorkerConfig();
|
const config = getImageArchiveWorkerConfig();
|
||||||
const availableSlots = config.concurrency - workerState.activeCount;
|
const availableSlots = config.concurrency - workerState.activeCount;
|
||||||
|
|
||||||
|
// 槽位满时不再访问数据库,等任务完成后由 finally 主动唤醒下一轮扫描。
|
||||||
if (availableSlots <= 0) {
|
if (availableSlots <= 0) {
|
||||||
nextDelayMs = config.pollMs;
|
nextDelayMs = config.pollMs;
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 数据库层会用条件 update 抢锁,保证多实例部署时同一条记录只会被一个 worker 领取。
|
||||||
const tasks = await claimRunnableImageArchiveTasks({
|
const tasks = await claimRunnableImageArchiveTasks({
|
||||||
limit: availableSlots,
|
limit: availableSlots,
|
||||||
workerId: workerState.workerId,
|
workerId: workerState.workerId,
|
||||||
@@ -96,6 +115,7 @@ const drainImageArchiveQueue = async () => {
|
|||||||
runImageArchiveTask(task, config);
|
runImageArchiveTask(task, config);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 如果本轮没填满并发槽,立刻再扫一次,尽快把可执行任务填到上限。
|
||||||
nextDelayMs =
|
nextDelayMs =
|
||||||
workerState.activeCount < config.concurrency ? 0 : config.pollMs;
|
workerState.activeCount < config.concurrency ? 0 : config.pollMs;
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
@@ -112,12 +132,14 @@ const drainImageArchiveQueue = async () => {
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/** 执行单条归档任务:下载上游图片、上传图床并回写归档结果 */
|
||||||
const runImageArchiveTask = (
|
const runImageArchiveTask = (
|
||||||
task: IImageArchiveTask,
|
task: IImageArchiveTask,
|
||||||
config: IImageArchiveWorkerConfig
|
config: IImageArchiveWorkerConfig
|
||||||
) => {
|
) => {
|
||||||
workerState.activeCount += 1;
|
workerState.activeCount += 1;
|
||||||
|
|
||||||
|
// 不 await 这里的 IIFE:扫描循环要尽快释放,归档任务本身受 activeCount 限流。
|
||||||
void (async () => {
|
void (async () => {
|
||||||
const startedAt = Date.now();
|
const startedAt = Date.now();
|
||||||
|
|
||||||
@@ -128,6 +150,7 @@ const runImageArchiveTask = (
|
|||||||
});
|
});
|
||||||
|
|
||||||
try {
|
try {
|
||||||
|
// 只读取用户展示名用于文件命名,不把用户敏感信息写入日志或响应。
|
||||||
const identity = await getUserArchiveIdentity(task.userId);
|
const identity = await getUserArchiveIdentity(task.userId);
|
||||||
const uploaded = await uploadImageFromUrl({
|
const uploaded = await uploadImageFromUrl({
|
||||||
imageUrl: task.imageUrl,
|
imageUrl: task.imageUrl,
|
||||||
@@ -137,6 +160,7 @@ const runImageArchiveTask = (
|
|||||||
createdAt: task.startedAt
|
createdAt: task.startedAt
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// 成功后补写 hosted 字段;前端仍通过历史接口读取归档结果,不暴露图床完整响应。
|
||||||
await finishImageGenerationArchiveSuccess(task.id, {
|
await finishImageGenerationArchiveSuccess(task.id, {
|
||||||
hostedImageUrl: uploaded.publicUrl,
|
hostedImageUrl: uploaded.publicUrl,
|
||||||
imageMimeType: uploaded.mimetype,
|
imageMimeType: uploaded.mimetype,
|
||||||
@@ -153,12 +177,14 @@ const runImageArchiveTask = (
|
|||||||
} catch (error) {
|
} catch (error) {
|
||||||
await handleImageArchiveFailure(task, config, error);
|
await handleImageArchiveFailure(task, config, error);
|
||||||
} finally {
|
} finally {
|
||||||
|
// 无论成功失败都释放本进程槽位,并立刻唤醒队列处理后续待执行任务。
|
||||||
workerState.activeCount = Math.max(0, workerState.activeCount - 1);
|
workerState.activeCount = Math.max(0, workerState.activeCount - 1);
|
||||||
wakeImageArchiveWorker();
|
wakeImageArchiveWorker();
|
||||||
}
|
}
|
||||||
})();
|
})();
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/** 处理单条归档任务失败:决定最终失败或按退避时间重新入队 */
|
||||||
const handleImageArchiveFailure = async (
|
const handleImageArchiveFailure = async (
|
||||||
task: IImageArchiveTask,
|
task: IImageArchiveTask,
|
||||||
config: IImageArchiveWorkerConfig,
|
config: IImageArchiveWorkerConfig,
|
||||||
@@ -168,6 +194,7 @@ const handleImageArchiveFailure = async (
|
|||||||
const errorMessage = getArchiveErrorMessage(error);
|
const errorMessage = getArchiveErrorMessage(error);
|
||||||
const shouldStop = task.archiveAttempts >= config.maxAttempts;
|
const shouldStop = task.archiveAttempts >= config.maxAttempts;
|
||||||
|
|
||||||
|
// 归档日志只保留排查摘要,不记录 token、图片二进制、完整图床响应或完整上游响应。
|
||||||
consola.error("[imageArchiveQueue] 归档失败", {
|
consola.error("[imageArchiveQueue] 归档失败", {
|
||||||
workerId: workerState.workerId,
|
workerId: workerState.workerId,
|
||||||
recordId: task.id.toString(),
|
recordId: task.id.toString(),
|
||||||
@@ -177,10 +204,12 @@ const handleImageArchiveFailure = async (
|
|||||||
});
|
});
|
||||||
|
|
||||||
if (shouldStop) {
|
if (shouldStop) {
|
||||||
|
// 最终失败只影响归档字段,生图状态保持 SUCCEEDED,前端看到的是本地化“图片归档失败”。
|
||||||
await finishImageGenerationArchiveFailed(task.id, errorMessage);
|
await finishImageGenerationArchiveFailed(task.id, errorMessage);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 可重试失败会清掉本次锁并设置下一次执行时间,期间历史记录仍表现为归档中。
|
||||||
await scheduleImageGenerationArchiveRetry({
|
await scheduleImageGenerationArchiveRetry({
|
||||||
recordId: task.id,
|
recordId: task.id,
|
||||||
nextRunAt: new Date(Date.now() + getRetryDelayMs(task.archiveAttempts)),
|
nextRunAt: new Date(Date.now() + getRetryDelayMs(task.archiveAttempts)),
|
||||||
@@ -188,7 +217,9 @@ const handleImageArchiveFailure = async (
|
|||||||
});
|
});
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/** 安排下一次队列扫描;新唤醒会覆盖旧定时器,确保扫描节奏可控 */
|
||||||
const scheduleDrain = (delayMs: number) => {
|
const scheduleDrain = (delayMs: number) => {
|
||||||
|
// 新任务入队时会用 0ms 重新调度,因此这里先清理旧 timer,避免重复扫描。
|
||||||
if (workerState.timer) {
|
if (workerState.timer) {
|
||||||
clearTimeout(workerState.timer);
|
clearTimeout(workerState.timer);
|
||||||
}
|
}
|
||||||
@@ -197,9 +228,11 @@ const scheduleDrain = (delayMs: number) => {
|
|||||||
workerState.timer = null;
|
workerState.timer = null;
|
||||||
void drainImageArchiveQueue();
|
void drainImageArchiveQueue();
|
||||||
}, delayMs);
|
}, delayMs);
|
||||||
|
// 不让空队列轮询 timer 阻止 Node 进程正常退出。
|
||||||
workerState.timer.unref?.();
|
workerState.timer.unref?.();
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/** 从环境变量读取 worker 配置,非法值统一回退默认值 */
|
||||||
const getImageArchiveWorkerConfig = (): IImageArchiveWorkerConfig => {
|
const getImageArchiveWorkerConfig = (): IImageArchiveWorkerConfig => {
|
||||||
return {
|
return {
|
||||||
concurrency: getPositiveIntegerEnv(
|
concurrency: getPositiveIntegerEnv(
|
||||||
@@ -215,6 +248,7 @@ const getImageArchiveWorkerConfig = (): IImageArchiveWorkerConfig => {
|
|||||||
};
|
};
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/** 根据当前尝试次数选择退避时间,超过表长时沿用最后一个退避值 */
|
||||||
const getRetryDelayMs = (attempt: number) => {
|
const getRetryDelayMs = (attempt: number) => {
|
||||||
const lastDelayMs = RETRY_DELAYS_MS[RETRY_DELAYS_MS.length - 1] ?? 0;
|
const lastDelayMs = RETRY_DELAYS_MS[RETRY_DELAYS_MS.length - 1] ?? 0;
|
||||||
return (
|
return (
|
||||||
@@ -223,17 +257,20 @@ const getRetryDelayMs = (attempt: number) => {
|
|||||||
);
|
);
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/** 读取正整数环境变量,避免 0、负数、NaN 破坏限流和重试策略 */
|
||||||
const getPositiveIntegerEnv = (name: string, fallback: number) => {
|
const getPositiveIntegerEnv = (name: string, fallback: number) => {
|
||||||
const value = Number.parseInt(process.env[name] ?? "", 10);
|
const value = Number.parseInt(process.env[name] ?? "", 10);
|
||||||
return Number.isInteger(value) && value > 0 ? value : fallback;
|
return Number.isInteger(value) && value > 0 ? value : fallback;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/** 保存到数据库的失败摘要只取可读 message,不保存复杂错误对象 */
|
||||||
const getArchiveErrorMessage = (error: unknown) => {
|
const getArchiveErrorMessage = (error: unknown) => {
|
||||||
if (error instanceof Error && error.message) return error.message;
|
if (error instanceof Error && error.message) return error.message;
|
||||||
if (typeof error === "string" && error) return error;
|
if (typeof error === "string" && error) return error;
|
||||||
return "图片归档失败";
|
return "图片归档失败";
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/** 归档日志错误摘要,不包含 stack、请求体、响应体或任何密钥字段 */
|
||||||
const toArchiveLogError = (error: unknown) => {
|
const toArchiveLogError = (error: unknown) => {
|
||||||
const maybeError = error as {
|
const maybeError = error as {
|
||||||
message?: string;
|
message?: string;
|
||||||
|
|||||||
Reference in New Issue
Block a user