import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, rmSync, writeFileSync } from 'node:fs' import { join } from 'node:path' export interface PendingJob { jobId: string promptId: string clientId: string ownerKey: string folderId: string hideThumbnail: boolean folderLocked?: boolean name?: string prompt: string aspect: string width: number height: number steps: number turbo: boolean seed: number startedAt: number imageName?: string imageSubfolder?: string extendTmpDir?: string extendPart1Path?: string } function pendingRoot() { const config = useRuntimeConfig() return join((config.libraryDir || process.env.LIBRARY_DIR || '/data/library').replace(/\/$/, ''), 'pending') } function pendingPath(jobId: string) { return join(pendingRoot(), `${jobId}.json`) } function ensurePending() { mkdirSync(pendingRoot(), { recursive: true }) } export function writePendingJob(job: PendingJob) { ensurePending() const tmp = pendingPath(job.jobId) + '.tmp' writeFileSync(tmp, JSON.stringify(job, null, 2)) renameSync(tmp, pendingPath(job.jobId)) } export function readPendingJob(jobId: string): PendingJob | null { const path = pendingPath(jobId) if (!existsSync(path)) return null try { return JSON.parse(readFileSync(path, 'utf8')) as PendingJob } catch { return null } } export function listPendingJobs(): PendingJob[] { ensurePending() return readdirSync(pendingRoot()) .filter(name => name.endsWith('.json')) .map(name => readPendingJob(name.replace(/\.json$/, ''))) .filter((job): job is PendingJob => Boolean(job)) } export function deletePendingJob(jobId: string) { rmSync(pendingPath(jobId), { force: true }) } export async function completePendingIfReady(pending: PendingJob) { const history = await fetchHistory(pending.promptId) const video = extractVideo(history, pending.promptId) if (!video) return null let buffer = await downloadComfyVideo(video) if (pending.extendPart1Path && pending.extendTmpDir) { if (!existsSync(pending.extendPart1Path)) { removeExtendTemp(pending.extendTmpDir) throw createError({ statusCode: 500, statusMessage: 'Extension source clip was missing during stitch' }) } try { buffer = await stitchExtension({ part1Path: pending.extendPart1Path, part2: buffer, tmpDir: pending.extendTmpDir }) } catch (error) { removeExtendTemp(pending.extendTmpDir) throw error } removeExtendTemp(pending.extendTmpDir) } const clip = await saveClip({ ownerKey: pending.ownerKey, folderId: pending.folderId, name: pending.name, prompt: pending.prompt, aspect: pending.aspect, width: pending.width, height: pending.height, steps: pending.steps, turbo: pending.turbo, seed: pending.seed, hideThumbnail: pending.hideThumbnail, video: buffer, thumb: null, comfyFilename: video.filename }) await purgeComfyArtifacts({ video, imageName: pending.imageName, imageSubfolder: pending.imageSubfolder, promptId: pending.promptId }) deletePendingJob(pending.jobId) return { type: 'complete' as const, status: 'complete' as const, message: pending.folderLocked ? 'Saved to the locked folder. Unlock it to view.' : (pending.extendPart1Path ? 'Extended video ready' : 'Video ready'), progress: 100, jobId: pending.jobId, clipId: clip.id, filename: pending.folderLocked ? undefined : video.filename, subfolder: pending.folderLocked ? undefined : video.subfolder, mediaType: pending.folderLocked ? undefined : video.type, hideThumbnail: pending.hideThumbnail, folderLocked: pending.folderLocked } } export async function resolvePendingJob(pending: PendingJob) { const done = await completePendingIfReady(pending) if (done) return done if (await isComfyPromptDropped(pending.promptId)) { deletePendingJob(pending.jobId) const message = 'ComfyUI dropped this job. Reloading the Comfy interface clears the queue. Generate again.' return { type: 'error' as const, status: 'error' as const, jobId: pending.jobId, error: message, message, progress: 0 } } return null }