Files
aigen/server/utils/watch.ts
T
TowstyandCursor dbe45fd2ff Add generation quality controls and bind them to the MiniMax workflow.
Expose aspect, CFG, FPS, sampler, and scheduler in Generation settings, persist them on clips and held jobs, and inject those values into the Comfy graph. Keep the image picker stills-only and allow renaming library clips.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-25 22:55:06 -05:00

223 lines
7.6 KiB
TypeScript

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<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
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,
cfg: job.library.cfg,
fps: job.library.fps,
samplerName: job.library.samplerName,
scheduler: job.library.scheduler
})
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<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
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)
}
}