import { existsSync } from 'node:fs' import type { Job, JobEvent } from '~/server/utils/jobs' import { getJob, restoreJob } from '~/server/utils/jobs' import type { PendingJob } from '~/server/utils/pending' import { isLastChainShot, libraryFromPending, pendingFromJob, remainingAfterCurrentShot, writePendingJob, deletePendingJob } from '~/server/utils/pending' import { syncQueueFromJob } from '~/server/utils/shotQueue' function classifyError(message: string) { const lower = message.toLowerCase() if (lower.includes('out of memory') || (lower.includes('cuda') && lower.includes('alloc')) || lower.includes('vram')) { return 'ComfyUI VRAM allocation failed. Try Turbo or a smaller aspect ratio.' } if (lower.includes('econnrefused') || lower.includes('unreachable') || lower.includes('fetch failed')) { return 'ComfyUI host connection dropped. Confirm the desktop instance is running and reachable on the LAN.' } if (lower.includes('timeout')) { return 'Network timeout talking to ComfyUI. The job may still be running on the desktop.' } return message } function sleep(ms: number) { return new Promise(resolve => setTimeout(resolve, ms)) } function chainMeta(job: Job) { const total = Math.max(1, job.library?.chainTotal || 1) const index = job.library?.chainIndex ?? 0 return { total, index, step: index + 1, label: job.library?.chainLabel || (index === 0 ? 'Initial' : `Extension ${index}`), chained: total > 1 } } export function mapChainProgress(job: Job, localPct: number) { const { total, index, chained } = chainMeta(job) const local = Math.max(0, Math.min(100, localPct)) if (!chained) return Math.round(local) const share = 100 / total return Math.min(99, Math.round(index * share + (Math.min(local, 99) / 100) * share)) } export function formatChainMessage(job: Job, message: string, samplePct?: number) { const { chained, step, total, label } = chainMeta(job) if (!chained) return message if (typeof samplePct === 'number') { return `Step ${step}/${total}: Generating ${label} (${samplePct}%)` } if (/waiting 3s buffer/i.test(message)) { return `Step ${step}/${total}: Waiting 3s buffer...` } return `Step ${step}/${total}: ${message}` } export function emitChainJob(job: Job, event: JobEvent, samplePct?: number) { const { chained, step, total, label } = chainMeta(job) const next: JobEvent = { ...event } if (next.message) { next.message = formatChainMessage(job, next.message, samplePct) } if (chained && typeof next.progress === 'number' && next.type !== 'complete' && next.type !== 'checkpoint') { next.overallProgress = mapChainProgress(job, next.progress) } if (chained) { next.chainStep = step next.chainTotal = total next.chainLabel = label } emitJob(job, next) } export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Promise { const persist = options.persist !== false job.socketReady = false job.saving = false const ws = new WebSocket(comfyWsUrl(job.clientId)) let settled = false let finishing = false let pollTimer: ReturnType | null = null let timeout: ReturnType | null = null return new Promise((resolve) => { let localProgress = 0 const startedAt = Date.now() let lastActivity = Date.now() const IDLE_MS = 30 * 60 * 1000 const ABSOLUTE_MS = 3 * 60 * 60 * 1000 const markActivity = () => { lastActivity = Date.now() } const emitLocal = (event: JobEvent, samplePct?: number) => { markActivity() if (typeof event.progress === 'number') localProgress = event.progress emitChainJob(job, event, samplePct) } const cleanup = () => { if (timeout) clearTimeout(timeout) if (pollTimer) clearInterval(pollTimer) timeout = null pollTimer = null try { ws.close() } catch { /* ignore */ } } const armIdle = () => { if (timeout) clearTimeout(timeout) if (settled) return const absLeft = ABSOLUTE_MS - (Date.now() - startedAt) if (absLeft <= 0) { void fail('Timed out waiting for ComfyUI (3 hours).') return } const idleLeft = IDLE_MS - (Date.now() - lastActivity) timeout = setTimeout(() => { if (settled || finishing) return if (Date.now() - lastActivity >= IDLE_MS) { void fail('Timed out waiting for ComfyUI progress (30 minutes with no updates). The desktop job may still be running.') return } armIdle() }, Math.max(1000, Math.min(idleLeft, absLeft))) } const fail = async (error: string) => { if (settled || finishing) return settled = true cleanup() job.status = job.status === 'cancelled' ? 'cancelled' : 'error' job.error = classifyError(error) emitChainJob(job, { type: 'error', error: job.error, message: job.error }) await syncQueueFromJob(job).catch(() => null) if (job.id) deletePendingJob(job.id) removeExtendTemp(job.library?.extendTmpDir) await purgeComfyArtifacts({ imageName: job.library?.imageName, imageSubfolder: job.library?.imageSubfolder, extraImageNames: job.library?.referenceImageNames }) resolve() } const succeed = async () => { if (settled || finishing || !job.promptId) return false finishing = true job.saving = true try { const history = await fetchHistory(job.promptId) const video = extractVideo(history, job.promptId) if (!video) { finishing = false job.saving = false return false } if (settled) return false settled = true cleanup() try { job.video = video const stitching = Boolean(job.library?.extendPart1Path && job.library.extendTmpDir) emitLocal({ type: 'status', message: stitching ? 'Extracting last frame & stitching extension...' : (persist ? 'Saving to library...' : 'Downloading segment...'), progress: stitching ? 97 : 98 }) if (job.library) { let buffer = await downloadComfyVideo(video) if (stitching) { if (!existsSync(job.library.extendPart1Path!)) { throw new Error('Extension source clip was missing during stitch') } buffer = await stitchExtension({ part1Path: job.library.extendPart1Path!, part2: buffer, tmpDir: job.library.extendTmpDir! }) } const clip = await saveClip({ ownerKey: job.library.ownerKey, folderId: job.library.folderId, name: job.library.name, prompt: job.library.prompt, aspect: job.library.aspect, width: job.library.width, height: job.library.height, steps: job.library.steps, turbo: job.library.turbo, seed: job.library.seed, hideThumbnail: job.library.hideThumbnail, hideInput: job.library.hideInput, video: buffer, thumb: (job.library.parentClipId || (job.library.chainIndex || 0) > 0 || /last[\s._-]*frame/i.test(job.library.imageName || '')) ? undefined : job.library.thumb, comfyFilename: video.filename, cfg: job.library.cfg, fps: job.library.fps, samplerName: job.library.samplerName, scheduler: job.library.scheduler, familyId: job.library.familyId || (job.library.familyId = crypto.randomUUID()), parentClipId: job.library.parentClipId, chainIndex: job.library.chainIndex, workflow: job.library.workflow, stillId: job.library.stillId, sound: job.library.sound }) job.clipId = clip.id job.savedPromptId = job.promptId job.hideThumbnail = clip.hideThumbnail job.library.thumb = undefined job.library.name = nextClipPartName(clip.name) job.library.parentClipId = clip.id await syncQueueFromJob(job, clip.id).catch(() => null) await purgeComfyArtifacts({ video, imageName: job.library.imageName, imageSubfolder: job.library.imageSubfolder, extraImageNames: job.library.referenceImageNames, promptId: job.promptId }) if (!persist) { job.segmentBuffer = buffer job.status = 'running' writePendingJob(pendingFromJob(job, { promptId: '', currentClipId: clip.id })) job.saving = false emitLocal({ type: 'checkpoint', message: stitching ? 'Extension checkpoint saved' : 'Initial segment saved', progress: chainMeta(job).chained ? mapChainProgress(job, 99) : 99, clipId: clip.id, hideThumbnail: job.hideThumbnail, folderLocked: job.library.folderLocked }) resolve() return true } removeExtendTemp(job.library.extendTmpDir) job.library.extendTmpDir = undefined } job.status = 'complete' const queuedRemaining = job.library ? remainingAfterCurrentShot(job.library).length : 0 emitJob(job, { type: 'complete', message: job.library?.folderLocked ? 'Saved to the locked folder. Unlock it to view.' : queuedRemaining > 0 ? `Shot ${chainMeta(job).step} of ${chainMeta(job).total} saved. ${queuedRemaining} remaining on Queue.` : (job.library?.extendSourceClipId || (job.library?.chainTotal || 1) > 1 ? 'Extended video ready' : 'Video ready'), progress: 100, filename: job.library?.folderLocked ? undefined : video.filename, subfolder: job.library?.folderLocked ? undefined : video.subfolder, mediaType: job.library?.folderLocked ? undefined : video.type, clipId: job.clipId, hideThumbnail: job.hideThumbnail, folderLocked: job.library?.folderLocked, chainStep: job.library?.chainStep, chainTotal: job.library?.chainTotal, chainLabel: job.library?.chainLabel, queuedRemaining, queueId: job.library?.queueId }) deletePendingJob(job.id) job.saving = false } catch (saveError) { job.saving = false removeExtendTemp(job.library?.extendTmpDir) const message = saveError instanceof Error ? saveError.message : String(saveError) job.status = 'error' job.error = persist ? `Video generated but library save failed: ${message}` : `Segment finished but could not be prepared: ${message}` emitChainJob(job, { type: 'error', error: job.error, message: job.error }) } resolve() return true } catch (error) { finishing = false job.saving = false throw error } } let dropMisses = 0 const pollHistory = async () => { if (settled || finishing || !job.promptId) return try { const inspected = inspectHistory(await fetchHistory(job.promptId), job.promptId) if (inspected.error) { await fail(inspected.error) return } if (inspected.video) { await succeed() return } markActivity() if (await isComfyPromptDropped(job.promptId)) { dropMisses += 1 if (dropMisses >= 20) { await fail('ComfyUI dropped this job. Reloading the Comfy interface clears the queue. Generate again.') } } else { dropMisses = 0 } } catch { // Transient ComfyUI history misses are expected while the graph is still running. } } armIdle() pollTimer = setInterval(() => { void pollHistory() }, 2000) ws.addEventListener('open', () => { job.socketReady = true emitLocal({ type: 'status', message: 'Connected to ComfyUI', progress: Math.max(localProgress, 4) }) }) ws.addEventListener('close', () => { if (settled || finishing || !job.promptId) return void pollHistory() }) ws.addEventListener('message', async (event) => { let payload: { type?: string; data?: Record } try { payload = JSON.parse(String(event.data)) } catch { return } const data = payload.data || {} if (job.promptId && data.prompt_id && data.prompt_id !== job.promptId) return if (payload.type === 'progress') { const value = Number(data.value || 0) const max = Number(data.max || 1) const pct = 12 + Math.round((value / Math.max(max, 1)) * 70) const samplePct = Math.round((value / Math.max(max, 1)) * 100) job.status = 'running' emitLocal({ type: 'progress', step: value, maxStep: max, progress: pct, node: data.node ? String(data.node) : null, message: `Sampling step ${value}/${max} (${samplePct}%)` }, samplePct) } if (payload.type === 'executing') { const node = data.node === null || data.node === undefined ? null : String(data.node) if (node === null) { const saved = await succeed() if (!saved) { emitLocal({ type: 'status', message: 'Waiting for ComfyUI to save the MP4...', progress: 96 }) } return } const label = NODE_LABELS[node] || `Running node ${node}` const encoding = isEncodingNode(node) emitLocal({ type: 'executing', node, progress: encoding ? 90 : Math.max(localProgress, 10), message: encoding ? 'Encoding video...' : `${label}...` }) } if (payload.type === 'execution_error') { const message = String(data.exception_message || data.message || 'ComfyUI node execution failed') await fail(message) } if (payload.type === 'execution_interrupted') { job.status = 'cancelled' await fail('Job interrupted.') } }) }) } export async function waitForComfySocket(job: Job, ms = 4000) { const started = Date.now() while (Date.now() - started < ms) { if (job.socketReady) return await sleep(100) } } const pendingWatches = new Set() export function ensurePendingWatch(pending: PendingJob) { const existing = getJob(pending.jobId) if (existing) return existing const lastShot = isLastChainShot(pending) const job = restoreJob({ id: pending.jobId, clientId: pending.clientId, promptId: pending.promptId, startedAt: pending.startedAt, hideThumbnail: pending.hideThumbnail, library: libraryFromPending(pending) }) if (pending.currentClipId) job.clipId = pending.currentClipId if (!pendingWatches.has(job.id)) { pendingWatches.add(job.id) const resume = async () => { try { if (pending.promptId) { await watchComfyJob(job, { persist: lastShot }) } if (job.status === 'error' || job.status === 'cancelled' || job.status === 'complete') return if (!lastShot) { const { continueQueuedExtensionsIfNeeded } = await import('~/server/utils/videoChain') await continueQueuedExtensionsIfNeeded(job) } } catch (error) { const message = error instanceof Error ? error.message : String(error) if (job.status === 'error' || job.status === 'cancelled') return job.status = 'error' job.error = message emitChainJob(job, { type: 'error', error: message, message }) } finally { pendingWatches.delete(job.id) } } void resume() } return job }