import { NODE_LABELS, isEncodingNode } from '~/server/utils/workflow' import type { Job } from '~/server/utils/jobs' 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 (8 steps) 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)) } export function watchComfyJob(job: Job): Promise { const ws = new WebSocket(comfyWsUrl(job.clientId)) let settled = false let pollTimer: ReturnType | null = null let timeout: ReturnType | null = null return new Promise((resolve) => { const cleanup = () => { if (timeout) clearTimeout(timeout) if (pollTimer) clearInterval(pollTimer) timeout = null pollTimer = null try { ws.close() } catch { /* ignore */ } } const fail = async (error: string) => { if (settled) return settled = true cleanup() job.status = job.status === 'cancelled' ? 'cancelled' : 'error' job.error = classifyError(error) emitJob(job, { type: 'error', error: job.error, message: job.error }) if (job.id) deletePendingJob(job.id) resolve() } const succeed = async () => { if (settled || !job.promptId) return false const history = await fetchHistory(job.promptId) const video = extractVideo(history, job.promptId) if (!video) return false settled = true cleanup() try { job.video = video emitJob(job, { type: 'status', message: 'Saving to library...', progress: 98 }) if (job.library) { const buffer = await downloadComfyVideo(video) 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.thumb, comfyFilename: video.filename }) job.clipId = clip.id job.hideThumbnail = clip.hideThumbnail job.library.thumb = undefined await purgeComfyArtifacts({ video, imageName: job.library.imageName, imageSubfolder: job.library.imageSubfolder, promptId: job.promptId }) } job.status = 'complete' emitJob(job, { type: 'complete', message: job.library?.folderLocked ? 'Saved to the locked folder. Unlock it to view.' : '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 }) deletePendingJob(job.id) } catch (saveError) { const message = saveError instanceof Error ? saveError.message : String(saveError) job.status = 'error' job.error = `Video generated but library save failed: ${message}` emitJob(job, { type: 'error', error: job.error, message: job.error }) } resolve() return true } let dropMisses = 0 const pollHistory = async () => { if (settled || !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 } if (await isComfyPromptDropped(job.promptId)) { dropMisses += 1 if (dropMisses >= 3) { 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. } } timeout = setTimeout(() => { void fail('Timed out waiting for ComfyUI (15 minutes).') }, 15 * 60 * 1000) pollTimer = setInterval(() => { void pollHistory() }, 2000) ws.addEventListener('open', () => { job.socketReady = true emitJob(job, { type: 'status', message: 'Connected to ComfyUI', progress: Math.max(job.progress, 4) }) }) ws.addEventListener('close', () => { if (settled || !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) job.status = 'running' emitJob(job, { type: 'progress', step: value, maxStep: max, progress: pct, node: data.node ? String(data.node) : null, message: `Sampling step ${value}/${max} (${Math.round((value / Math.max(max, 1)) * 100)}%)` }) } 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) { emitJob(job, { 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) emitJob(job, { type: 'executing', node, progress: encoding ? 90 : Math.max(job.progress, 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) } }