// server/utils/imageArchiveQueue.ts - 图床归档队列 worker:从数据库领取任务、限流上传并按退避策略重试。 import { randomUUID } from "node:crypto"; import { clearTimeout, setTimeout } from "node:timers"; import { consola } from "consola"; import { claimRunnableImageArchiveTasks, finishImageGenerationArchiveFailed, finishImageGenerationArchiveSuccess, getUserArchiveIdentity, scheduleImageGenerationArchiveRetry, type IImageArchiveTask } from "~~/server/utils/imageGenerationRecords"; import { uploadImageFromUrl } from "~~/server/utils/lsky"; interface IImageArchiveWorkerConfig { concurrency: number; maxAttempts: number; lockTtlMs: number; pollMs: number; } interface IImageArchiveWorkerState { activeCount: number; draining: boolean; started: boolean; timer: ReturnType | null; workerId: string; } const DEFAULT_CONCURRENCY = 2; const DEFAULT_MAX_ATTEMPTS = 3; const DEFAULT_LOCK_TTL_MS = 10 * 60 * 1000; const DEFAULT_POLL_MS = 15 * 1000; const RETRY_DELAYS_MS = [30 * 1000, 2 * 60 * 1000, 10 * 60 * 1000]; const globalForArchiveWorker = globalThis as unknown as { imageArchiveWorkerState?: IImageArchiveWorkerState; }; const workerState = globalForArchiveWorker.imageArchiveWorkerState ?? (globalForArchiveWorker.imageArchiveWorkerState = { activeCount: 0, draining: false, started: false, timer: null, workerId: `archive-${randomUUID()}` }); /** 启动应用内图床归档 worker;重复调用只会唤醒同一个 worker */ export const startImageArchiveWorker = () => { if (!workerState.started) { workerState.started = true; consola.info("[imageArchiveQueue] worker 启动", { workerId: workerState.workerId, concurrency: getImageArchiveWorkerConfig().concurrency }); } scheduleDrain(0); }; /** 生图成功入队后调用,用于尽快扫描新任务 */ export const wakeImageArchiveWorker = () => { if (!workerState.started) return; scheduleDrain(0); }; const drainImageArchiveQueue = async () => { if (workerState.draining) return; workerState.draining = true; let nextDelayMs = DEFAULT_POLL_MS; try { const config = getImageArchiveWorkerConfig(); const availableSlots = config.concurrency - workerState.activeCount; if (availableSlots <= 0) { nextDelayMs = config.pollMs; return; } const tasks = await claimRunnableImageArchiveTasks({ limit: availableSlots, workerId: workerState.workerId, lockTtlMs: config.lockTtlMs }); if (tasks.length === 0) { nextDelayMs = config.pollMs; return; } for (const task of tasks) { runImageArchiveTask(task, config); } nextDelayMs = workerState.activeCount < config.concurrency ? 0 : config.pollMs; } catch (error) { consola.error("[imageArchiveQueue] 扫描归档任务失败", { workerId: workerState.workerId, error: toArchiveLogError(error) }); } finally { workerState.draining = false; if (workerState.started) { scheduleDrain(nextDelayMs); } } }; const runImageArchiveTask = ( task: IImageArchiveTask, config: IImageArchiveWorkerConfig ) => { workerState.activeCount += 1; void (async () => { const startedAt = Date.now(); consola.info("[imageArchiveQueue] 归档开始", { workerId: workerState.workerId, recordId: task.id.toString(), attempt: task.archiveAttempts }); try { const identity = await getUserArchiveIdentity(task.userId); const uploaded = await uploadImageFromUrl({ imageUrl: task.imageUrl, userId: task.userId, username: identity.username, recordId: task.id, createdAt: task.startedAt }); await finishImageGenerationArchiveSuccess(task.id, { hostedImageUrl: uploaded.publicUrl, imageMimeType: uploaded.mimetype, hostedResponse: uploaded.response }); consola.info("[imageArchiveQueue] 归档成功", { workerId: workerState.workerId, recordId: task.id.toString(), mimeType: uploaded.mimetype, byteLength: uploaded.byteLength, elapsedMs: Date.now() - startedAt }); } catch (error) { await handleImageArchiveFailure(task, config, error); } finally { workerState.activeCount = Math.max(0, workerState.activeCount - 1); wakeImageArchiveWorker(); } })(); }; const handleImageArchiveFailure = async ( task: IImageArchiveTask, config: IImageArchiveWorkerConfig, error: unknown ) => { const safeError = toArchiveLogError(error); const errorMessage = getArchiveErrorMessage(error); const shouldStop = task.archiveAttempts >= config.maxAttempts; consola.error("[imageArchiveQueue] 归档失败", { workerId: workerState.workerId, recordId: task.id.toString(), attempt: task.archiveAttempts, terminal: shouldStop, error: safeError }); if (shouldStop) { await finishImageGenerationArchiveFailed(task.id, errorMessage); return; } await scheduleImageGenerationArchiveRetry({ recordId: task.id, nextRunAt: new Date(Date.now() + getRetryDelayMs(task.archiveAttempts)), errorMessage }); }; const scheduleDrain = (delayMs: number) => { if (workerState.timer) { clearTimeout(workerState.timer); } workerState.timer = setTimeout(() => { workerState.timer = null; void drainImageArchiveQueue(); }, delayMs); workerState.timer.unref?.(); }; const getImageArchiveWorkerConfig = (): IImageArchiveWorkerConfig => { return { concurrency: getPositiveIntegerEnv( "IMAGE_ARCHIVE_CONCURRENCY", DEFAULT_CONCURRENCY ), maxAttempts: getPositiveIntegerEnv( "IMAGE_ARCHIVE_MAX_ATTEMPTS", DEFAULT_MAX_ATTEMPTS ), lockTtlMs: DEFAULT_LOCK_TTL_MS, pollMs: DEFAULT_POLL_MS }; }; const getRetryDelayMs = (attempt: number) => { const lastDelayMs = RETRY_DELAYS_MS[RETRY_DELAYS_MS.length - 1] ?? 0; return ( RETRY_DELAYS_MS[Math.min(attempt - 1, RETRY_DELAYS_MS.length - 1)] ?? lastDelayMs ); }; const getPositiveIntegerEnv = (name: string, fallback: number) => { const value = Number.parseInt(process.env[name] ?? "", 10); return Number.isInteger(value) && value > 0 ? value : fallback; }; const getArchiveErrorMessage = (error: unknown) => { if (error instanceof Error && error.message) return error.message; if (typeof error === "string" && error) return error; return "图片归档失败"; }; const toArchiveLogError = (error: unknown) => { const maybeError = error as { message?: string; name?: string; response?: { status?: number; }; status?: number; statusCode?: number; statusMessage?: string; }; return { name: maybeError.name, message: maybeError.message ?? maybeError.statusMessage, status: maybeError.status ?? maybeError.statusCode ?? maybeError.response?.status }; };