Files
aigen/server/utils/watch.ts
T
TowstyandCursor ec887a320a Retry library pulls after Comfy finishes and fall back to disk.
Long jobs were succeeding on Beast then failing the save step when /view briefly dropped; retries plus host-agent file serving keep finished MP4s from being orphaned.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-03 21:06:15 -05:00

553 lines
20 KiB
TypeScript

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'
import { isEncodingNode, NODE_LABELS } from '~/server/utils/workflow'
import { extractFirstFrameFromBuffer } from '~/server/utils/ffmpeg'
import { mergePermanenceRefs, type PermanenceRef } from '~/utils/globalLocks'
import { nextFamilyPartName } from '~/server/utils/library'
function clipRefsFromJob(library: NonNullable<Job['library']>): PermanenceRef[] | undefined {
const refs = mergePermanenceRefs(library.permanenceRefs, library.shotPermanenceRefs?.[library.chainIndex || 0])
return refs.length ? refs : undefined
}
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') || lower.includes('connection dropped')) {
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 errorText(error: unknown) {
if (!error) return 'Unknown error'
if (typeof error === 'object') {
const rec = error as { statusMessage?: string; message?: string; data?: { cause?: string } }
const bits = [rec.statusMessage, rec.message, rec.data?.cause].filter(Boolean)
if (bits.length) return String(bits[0])
}
if (error instanceof Error) return error.message
return String(error)
}
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 (/checkpoint saved|segment saved/i.test(message)) {
return `Shot ${step}/${total} saved`
}
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)
}
function nodeLabel(node: string) {
if (!node) return ''
if (NODE_LABELS[node]) return NODE_LABELS[node]
if (node.startsWith('user:lora')) return 'Applying LoRA'
return ''
}
function isNonSamplerProgress(node: string) {
if (!node) return false
if (isEncodingNode(node)) return true
const label = nodeLabel(node)
return /save|checkpoint|video combine|vhs/i.test(label) && !/sampler/i.test(label)
}
export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Promise<void> {
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<typeof setInterval> | null = null
let timeout: ReturnType<typeof setTimeout> | 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 {
let history: Awaited<ReturnType<typeof fetchHistory>> | null = null
let video: ReturnType<typeof extractVideo> = null
for (let attempt = 1; attempt <= 6; attempt++) {
try {
history = await fetchHistory(job.promptId)
video = extractVideo(history, job.promptId)
if (video) break
} catch (error) {
if (attempt >= 6) throw error
}
emitLocal({
type: 'status',
message: `Waiting for ComfyUI history (${attempt}/6)...`,
progress: 96
})
await sleep(1000 * attempt)
}
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)
let segmentFirstFrame: Buffer | undefined
if (stitching) {
if (!existsSync(job.library.extendPart1Path!)) {
throw new Error('Extension source clip was missing during stitch')
}
if (job.library.extendTmpDir) {
try {
segmentFirstFrame = await extractFirstFrameFromBuffer(buffer, job.library.extendTmpDir)
} catch {
// ensureClipFirstFrame will seek to the join if this frame is missing
}
}
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.promptMid || job.library.prompt,
promptPre: job.library.promptPre,
promptPost: job.library.promptPost,
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,
globalLocks: job.library.globalLocks,
permanenceRefs: clipRefsFromJob(job.library),
loraName: job.library.loraName,
loraStack: job.library.loraStack,
segmentFirstFrame
})
job.clipId = clip.id
job.savedPromptId = job.promptId
job.hideThumbnail = clip.hideThumbnail
job.library.thumb = undefined
job.library.name = nextFamilyPartName(job.library.ownerKey, clip)
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 = classifyError(errorText(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 })
// Keep pending so a later Recover / container restart can still pull the MP4.
if (job.library && job.promptId) {
writePendingJob(pendingFromJob(job, {
promptId: job.promptId,
currentClipId: job.clipId || job.library.parentClipId
}))
}
await syncQueueFromJob(job).catch(() => null)
}
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<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 node = data.node ? String(data.node) : ''
if (isNonSamplerProgress(node)) {
markActivity()
return
}
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: 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) {
try {
const saved = await succeed()
if (!saved) {
emitLocal({ type: 'status', message: 'Waiting for ComfyUI to save the MP4...', progress: 96 })
}
} catch (error) {
finishing = false
job.saving = false
emitLocal({
type: 'status',
message: `Save hiccup (${errorText(error)}). Retrying...`,
progress: 96
})
}
return
}
const label = nodeLabel(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<string>()
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 {
const { applyPersistedPauseToJob, onLiveVideoSettled } = await import('~/server/utils/studioQueue')
applyPersistedPauseToJob(job)
const holdNow = () => job.library?.stopAfterCurrent === true
if (pending.promptId) {
await watchComfyJob(job, { persist: lastShot || holdNow() })
}
applyPersistedPauseToJob(job)
if (job.status === 'error' || job.status === 'cancelled') {
await onLiveVideoSettled(job)
return
}
if (holdNow()) {
await onLiveVideoSettled(job)
return
}
if (job.status === 'complete' && lastShot) {
await onLiveVideoSettled(job)
return
}
if (!lastShot) {
if (job.status === 'complete') job.status = 'running'
const { continueQueuedExtensionsIfNeeded } = await import('~/server/utils/videoChain')
await continueQueuedExtensionsIfNeeded(job)
}
applyPersistedPauseToJob(job)
if (job.status === 'error' || job.status === 'cancelled' || job.status === 'complete' || holdNow()) {
await onLiveVideoSettled(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 })
const { onLiveVideoSettled } = await import('~/server/utils/studioQueue')
await onLiveVideoSettled(job).catch(() => null)
} finally {
pendingWatches.delete(job.id)
}
}
void resume()
}
return job
}