import { NODE_LABELS } 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 } export async function watchComfyJob(job: Job) { const ws = new WebSocket(comfyWsUrl(job.clientId)) let settled = false const finish = async (error?: string) => { if (settled) return settled = true try { ws.close() } catch { /* ignore */ } if (error) { job.status = job.status === 'cancelled' ? 'cancelled' : 'error' job.error = classifyError(error) emitJob(job, { type: 'error', error: job.error, message: job.error }) return } emitJob(job, { type: 'status', message: 'Fetching output...', progress: 96 }) const history = await fetchHistory(job.promptId || '') const video = extractVideo(history, job.promptId || '') if (!video) { job.status = 'error' job.error = 'Workflow finished but no MP4 was found in ComfyUI history.' emitJob(job, { type: 'error', error: job.error, message: job.error }) return } job.video = video job.status = 'complete' emitJob(job, { type: 'complete', message: 'Video ready', progress: 100, filename: video.filename, subfolder: video.subfolder, mediaType: video.type }) } const timeout = setTimeout(() => { finish('Timed out waiting for ComfyUI (15 minutes).') }, 15 * 60 * 1000) ws.addEventListener('open', () => { emitJob(job, { type: 'status', message: 'Connected to ComfyUI', progress: 8 }) }) ws.addEventListener('error', () => { if (!settled) finish('ComfyUI WebSocket connection dropped.') }) ws.addEventListener('close', () => { if (!settled) { // Fall back to history polling in case completion arrived as HTTP only. setTimeout(async () => { if (settled) return const history = await fetchHistory(job.promptId || '') if (extractVideo(history, job.promptId || '')) { clearTimeout(timeout) await finish() } }, 1500) } }) 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) { clearTimeout(timeout) await finish() return } const label = NODE_LABELS[node] || `Running node ${node}` const encoding = node === '16' || node === '17' emitJob(job, { type: 'executing', node, progress: encoding ? 90 : Math.max(job.progress, 10), message: encoding ? 'Encoding video...' : `${label}...` }) } if (payload.type === 'execution_error') { clearTimeout(timeout) const message = String(data.exception_message || data.message || 'ComfyUI node execution failed') await finish(message) } if (payload.type === 'execution_interrupted') { clearTimeout(timeout) job.status = 'cancelled' await finish('Job interrupted.') } }) return () => { clearTimeout(timeout) try { ws.close() } catch { /* ignore */ } } }