A finished job no longer closes a draft or jumps the page, and a second extend from the same part saves as p8b instead of another p8. Co-authored-by: Cursor <cursoragent@cursor.com>
509 lines
19 KiB
TypeScript
509 lines
19 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')) {
|
|
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))
|
|
}
|
|
|
|
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 {
|
|
const history = await fetchHistory(job.promptId)
|
|
const video = extractVideo(history, job.promptId)
|
|
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 = saveError instanceof Error ? saveError.message : String(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 })
|
|
}
|
|
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) {
|
|
const saved = await succeed()
|
|
if (!saved) {
|
|
emitLocal({ type: 'status', message: 'Waiting for ComfyUI to save the MP4...', 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
|
|
}
|