Fix generation progress getting stuck at 1% after Comfy finishes.
Poll Comfy history and the job snapshot so the browser still gets the video when the live stream is buffered or the websocket misses completion. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+161
-127
@@ -3,7 +3,7 @@ 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')) {
|
||||
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')) {
|
||||
@@ -15,154 +15,188 @@ function classifyError(message: string) {
|
||||
return message
|
||||
}
|
||||
|
||||
export async function watchComfyJob(job: Job) {
|
||||
function sleep(ms: number) {
|
||||
return new Promise(resolve => setTimeout(resolve, ms))
|
||||
}
|
||||
|
||||
export function watchComfyJob(job: Job): Promise<void> {
|
||||
const ws = new WebSocket(comfyWsUrl(job.clientId))
|
||||
let settled = false
|
||||
let pollTimer: ReturnType<typeof setInterval> | null = null
|
||||
let timeout: ReturnType<typeof setTimeout> | null = null
|
||||
|
||||
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
|
||||
return new Promise((resolve) => {
|
||||
const cleanup = () => {
|
||||
if (timeout) clearTimeout(timeout)
|
||||
if (pollTimer) clearInterval(pollTimer)
|
||||
timeout = null
|
||||
pollTimer = null
|
||||
try { ws.close() } catch { /* ignore */ }
|
||||
}
|
||||
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
|
||||
if (job.library) {
|
||||
emitJob(job, { type: 'status', message: 'Saving to library...', progress: 98 })
|
||||
|
||||
const finish = async (error?: string) => {
|
||||
if (settled) return
|
||||
settled = true
|
||||
cleanup()
|
||||
try {
|
||||
const buffer = await downloadComfyVideo(video)
|
||||
const clip = await saveClip({
|
||||
folderId: job.library.folderId,
|
||||
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,
|
||||
video: buffer,
|
||||
thumb: job.library.hideThumbnail ? null : job.library.thumb
|
||||
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
|
||||
if (job.library) {
|
||||
emitJob(job, { type: 'status', message: 'Saving to library...', progress: 98 })
|
||||
try {
|
||||
const buffer = await downloadComfyVideo(video)
|
||||
const clip = await saveClip({
|
||||
folderId: job.library.folderId,
|
||||
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,
|
||||
video: buffer,
|
||||
thumb: job.library.hideThumbnail ? null : job.library.thumb
|
||||
})
|
||||
job.clipId = clip.id
|
||||
job.hideThumbnail = clip.hideThumbnail
|
||||
job.library.thumb = undefined
|
||||
} 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 })
|
||||
return
|
||||
}
|
||||
}
|
||||
job.status = 'complete'
|
||||
emitJob(job, {
|
||||
type: 'complete',
|
||||
message: 'Video ready',
|
||||
progress: 100,
|
||||
filename: video.filename,
|
||||
subfolder: video.subfolder,
|
||||
mediaType: video.type,
|
||||
clipId: job.clipId,
|
||||
hideThumbnail: job.hideThumbnail
|
||||
})
|
||||
job.clipId = clip.id
|
||||
job.hideThumbnail = clip.hideThumbnail
|
||||
job.library.thumb = undefined
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
job.status = 'error'
|
||||
job.error = `Video generated but library save failed: ${message}`
|
||||
emitJob(job, { type: 'error', error: job.error, message: job.error })
|
||||
return
|
||||
} finally {
|
||||
resolve()
|
||||
}
|
||||
}
|
||||
job.status = 'complete'
|
||||
emitJob(job, {
|
||||
type: 'complete',
|
||||
message: 'Video ready',
|
||||
progress: 100,
|
||||
filename: video.filename,
|
||||
subfolder: video.subfolder,
|
||||
mediaType: video.type,
|
||||
clipId: job.clipId,
|
||||
hideThumbnail: job.hideThumbnail
|
||||
})
|
||||
}
|
||||
|
||||
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)
|
||||
const pollHistory = async () => {
|
||||
if (settled || !job.promptId) return
|
||||
try {
|
||||
const inspected = inspectHistory(await fetchHistory(job.promptId), job.promptId)
|
||||
if (inspected.error) {
|
||||
await finish(inspected.error)
|
||||
return
|
||||
}
|
||||
if (inspected.video || inspected.completed) {
|
||||
await finish()
|
||||
}
|
||||
}, 1500)
|
||||
} catch {
|
||||
// Transient ComfyUI history misses are expected while the graph is still running.
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
ws.addEventListener('message', async (event) => {
|
||||
let payload: { type?: string; data?: Record<string, unknown> }
|
||||
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
|
||||
timeout = setTimeout(() => {
|
||||
void finish('Timed out waiting for ComfyUI (15 minutes).')
|
||||
}, 15 * 60 * 1000)
|
||||
|
||||
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'
|
||||
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('error', () => {
|
||||
if (settled) return
|
||||
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)}%)`
|
||||
type: 'status',
|
||||
message: job.promptId
|
||||
? 'Live socket unavailable, polling ComfyUI history...'
|
||||
: 'Waiting for ComfyUI socket...',
|
||||
progress: Math.max(job.progress, 4)
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
if (payload.type === 'executing') {
|
||||
const node = data.node === null || data.node === undefined ? null : String(data.node)
|
||||
if (node === null) {
|
||||
clearTimeout(timeout)
|
||||
await finish()
|
||||
ws.addEventListener('message', async (event) => {
|
||||
let payload: { type?: string; data?: Record<string, unknown> }
|
||||
try {
|
||||
payload = JSON.parse(String(event.data))
|
||||
} catch {
|
||||
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}...`
|
||||
})
|
||||
}
|
||||
const data = payload.data || {}
|
||||
if (job.promptId && data.prompt_id && data.prompt_id !== job.promptId) return
|
||||
|
||||
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 === '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 === 'execution_interrupted') {
|
||||
clearTimeout(timeout)
|
||||
job.status = 'cancelled'
|
||||
await finish('Job interrupted.')
|
||||
}
|
||||
if (payload.type === 'executing') {
|
||||
const node = data.node === null || data.node === undefined ? null : String(data.node)
|
||||
if (node === null) {
|
||||
await finish()
|
||||
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 finish(message)
|
||||
}
|
||||
|
||||
if (payload.type === 'execution_interrupted') {
|
||||
job.status = 'cancelled'
|
||||
await finish('Job interrupted.')
|
||||
}
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
return () => {
|
||||
clearTimeout(timeout)
|
||||
try { ws.close() } catch { /* ignore */ }
|
||||
export async function waitForComfySocket(job: Job, ms = 4000) {
|
||||
const started = Date.now()
|
||||
while (Date.now() - started < ms) {
|
||||
if (job.socketReady) return
|
||||
await sleep(100)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user