import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, rmSync, writeFileSync } from 'node:fs' import { join } from 'node:path' import type { Job } from '~/server/utils/jobs' import { extractFirstFrameFromBuffer } from '~/server/utils/ffmpeg' import { mergePermanenceRefs } from '~/utils/globalLocks' 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 familyId?: string parentClipId?: string chainIndex?: number chainStep?: number chainTotal?: number chainLabel?: string stillId?: string stillFilename?: string sound?: boolean referenceImageNames?: string[] referenceStillIds?: Array useIdentityRefs?: boolean globalLocks?: string permanenceRefs?: import('~/utils/globalLocks').PermanenceRef[] shotPermanenceRefs?: import('~/utils/globalLocks').PermanenceRef[][] workflow?: import('~/utils/videoModels').VideoWorkflowId duration?: number cfg?: number fps?: number samplerName?: string scheduler?: string hideInput?: boolean extensions?: { prompt: string; duration: number; permanenceRefs?: import('~/utils/globalLocks').PermanenceRef[] }[] remainingExtensions?: { prompt: string; duration: number; permanenceRefs?: import('~/utils/globalLocks').PermanenceRef[] }[] currentClipId?: string queueId?: string queueAutoRun?: boolean queueBudget?: number stopAfterCurrent?: boolean } 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 remainingAfterCurrentShot(job: { chainIndex?: number extensions?: { prompt: string; duration: number }[] }) { const extensions = job.extensions || [] const index = job.chainIndex || 0 return index === 0 ? extensions : extensions.slice(index) } export function isLastChainShot(job: { chainIndex?: number extensions?: { prompt: string; duration: number }[] remainingExtensions?: { prompt: string; duration: number }[] }) { if (job.remainingExtensions) return job.remainingExtensions.length === 0 return remainingAfterCurrentShot(job).length === 0 } export function pendingFromJob(job: Job, overrides: Partial = {}): PendingJob { const library = job.library if (!library) { throw new Error('Cannot persist a job without library metadata') } const remaining = remainingAfterCurrentShot(library) return { jobId: job.id, promptId: job.promptId || '', clientId: job.clientId, ownerKey: library.ownerKey, folderId: library.folderId, hideThumbnail: library.hideThumbnail, folderLocked: library.folderLocked, name: library.name, prompt: library.prompt, aspect: library.aspect, width: library.width, height: library.height, steps: library.steps, turbo: library.turbo, seed: library.seed, startedAt: job.startedAt, imageName: library.imageName, imageSubfolder: library.imageSubfolder, extendTmpDir: library.extendTmpDir, extendPart1Path: library.extendPart1Path, familyId: library.familyId, parentClipId: library.parentClipId, chainIndex: library.chainIndex, chainStep: library.chainStep, chainTotal: library.chainTotal, chainLabel: library.chainLabel, stillId: library.stillId, stillFilename: library.stillFilename, sound: library.sound, referenceImageNames: library.referenceImageNames, referenceStillIds: library.referenceStillIds, useIdentityRefs: library.useIdentityRefs, workflow: library.workflow, duration: library.duration, cfg: library.cfg, fps: library.fps, samplerName: library.samplerName, scheduler: library.scheduler, hideInput: library.hideInput, extensions: library.extensions, remainingExtensions: remaining, currentClipId: job.clipId, globalLocks: library.globalLocks, permanenceRefs: library.permanenceRefs, shotPermanenceRefs: library.shotPermanenceRefs, queueId: library.queueId, queueAutoRun: library.queueAutoRun, queueBudget: library.queueBudget, stopAfterCurrent: library.stopAfterCurrent, ...overrides } } export function libraryFromPending(pending: PendingJob): NonNullable { return { ownerKey: pending.ownerKey, folderId: pending.folderId, hideThumbnail: pending.hideThumbnail, hideInput: pending.hideInput, folderLocked: pending.folderLocked, name: pending.name, prompt: pending.prompt, aspect: pending.aspect, width: pending.width, height: pending.height, steps: pending.steps, turbo: pending.turbo, seed: pending.seed, cfg: pending.cfg, fps: pending.fps, samplerName: pending.samplerName, scheduler: pending.scheduler, imageName: pending.imageName, imageSubfolder: pending.imageSubfolder, extendTmpDir: pending.extendTmpDir, extendPart1Path: pending.extendPart1Path, familyId: pending.familyId, parentClipId: pending.parentClipId, chainIndex: pending.chainIndex, chainStep: pending.chainStep, chainTotal: pending.chainTotal, chainLabel: pending.chainLabel, stillId: pending.stillId, stillFilename: pending.stillFilename, sound: pending.sound, referenceImageNames: pending.referenceImageNames, referenceStillIds: pending.referenceStillIds, useIdentityRefs: pending.useIdentityRefs, workflow: pending.workflow, duration: pending.duration, extensions: pending.extensions || pending.remainingExtensions, queueId: pending.queueId, queueAutoRun: pending.queueAutoRun, queueBudget: pending.queueBudget, stopAfterCurrent: pending.stopAfterCurrent, globalLocks: pending.globalLocks, permanenceRefs: pending.permanenceRefs, shotPermanenceRefs: pending.shotPermanenceRefs } } 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 patchPendingJob(jobId: string, patch: Partial) { const existing = readPendingJob(jobId) if (!existing) return null const next = { ...existing, ...patch } writePendingJob(next) return next } 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) { if (!isLastChainShot(pending)) return null if (!pending.promptId) return null const alreadySaved = () => { const live = getJob(pending.jobId) if (live?.savedPromptId === 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: live.clipId, hideThumbnail: pending.hideThumbnail, folderLocked: pending.folderLocked } } return null } const inFlight = () => { const current = getJob(pending.jobId) return Boolean(current && current.promptId === pending.promptId && current.saving) } if (inFlight()) return null const saved = alreadySaved() if (saved) return saved const history = await fetchHistory(pending.promptId) if (inFlight()) return null const afterHistory = alreadySaved() if (afterHistory) return afterHistory const video = extractVideo(history, pending.promptId) if (!video) return null let buffer = await downloadComfyVideo(video) if (inFlight()) return null const afterDownload = alreadySaved() if (afterDownload) return afterDownload let segmentFirstFrame: Buffer | undefined 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 { segmentFirstFrame = await extractFirstFrameFromBuffer(buffer, pending.extendTmpDir) } catch { // ensureClipFirstFrame will seek to the join if this frame is missing } try { buffer = await stitchExtension({ part1Path: pending.extendPart1Path, part2: buffer, tmpDir: pending.extendTmpDir }) } catch (error) { removeExtendTemp(pending.extendTmpDir) throw error } removeExtendTemp(pending.extendTmpDir) } const afterStitch = alreadySaved() if (afterStitch) return afterStitch 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, familyId: pending.familyId, parentClipId: pending.parentClipId, chainIndex: pending.chainIndex, stillId: pending.stillId, sound: pending.sound, globalLocks: pending.globalLocks, permanenceRefs: mergePermanenceRefs(pending.permanenceRefs, pending.shotPermanenceRefs?.[pending.chainIndex || 0]), segmentFirstFrame }) await purgeComfyArtifacts({ video, imageName: pending.imageName, imageSubfolder: pending.imageSubfolder, extraImageNames: pending.referenceImageNames, 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 }