From 18f350c32a718458607982c3a923e4bdf856348e Mon Sep 17 00:00:00 2001 From: Towsty Date: Sun, 30 Aug 2026 07:24:47 -0500 Subject: [PATCH] Finish the current video shot list before starting the next queued video. Co-authored-by: Cursor --- server/utils/comfyLifecycle.ts | 2 +- server/utils/shotQueue.ts | 6 ++- server/utils/studioQueue.ts | 86 ++++++++++++++++++++++++++-------- server/utils/videoChain.ts | 30 ++++++++---- server/utils/watch.ts | 3 +- 5 files changed, 95 insertions(+), 32 deletions(-) diff --git a/server/utils/comfyLifecycle.ts b/server/utils/comfyLifecycle.ts index a5239de..c74fe55 100644 --- a/server/utils/comfyLifecycle.ts +++ b/server/utils/comfyLifecycle.ts @@ -43,7 +43,7 @@ function settings() { healthTimeoutMs: 2500, bootPollAttempts: Number(process.env.COMFY_BOOT_POLL_ATTEMPTS || 20), bootPollDelayMs: 3000, - startTimeoutMs: Number(process.env.COMFY_START_TIMEOUT_MS || 180_000), + startTimeoutMs: Number(process.env.COMFY_START_TIMEOUT_MS || 300_000), busyWaitMs: Number(process.env.COMFY_BUSY_WAIT_MS || 180_000) } } diff --git a/server/utils/shotQueue.ts b/server/utils/shotQueue.ts index 9a51274..e457288 100644 --- a/server/utils/shotQueue.ts +++ b/server/utils/shotQueue.ts @@ -320,7 +320,11 @@ export async function prepareQueueBurst(owner: string, id: string, count: number if (activeJobs.has(queue.id)) { throw createError({ statusCode: 409, statusMessage: 'This batch is already processing' }) } - const remaining = queue.segments.filter(segment => segment.status === 'pending' || segment.status === 'error') + const remaining = queue.segments.filter(segment => ( + segment.status === 'pending' + || segment.status === 'error' + || (segment.status === 'running' && !activeJobs.has(queue.id)) + )) if (!remaining.length) { throw createError({ statusCode: 400, statusMessage: 'No remaining shots to process' }) } diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts index a488274..54f49a0 100644 --- a/server/utils/studioQueue.ts +++ b/server/utils/studioQueue.ts @@ -2,7 +2,7 @@ import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFile import { join } from 'node:path' import { getJob, listJobs, emitJob, type Job } from '~/server/utils/jobs' import { listPendingJobs, patchPendingJob, readPendingJob, deletePendingJob } from '~/server/utils/pending' -import { getShotQueue } from '~/server/utils/shotQueue' +import { getShotQueue, pendingSegmentCount } from '~/server/utils/shotQueue' import { fetchLiveQueue } from '~/server/utils/comfy' import { isLtxWorkflow, isTextToVideo, isXaigenStudio, LTX_DISABLED_MESSAGE, ltxWorkflowEnabled, parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels' import { persistLoraFields } from '~/server/utils/loras' @@ -262,17 +262,24 @@ function sweepStaleLiveJobs() { } } +function liveJobOwnsGpu(job: Job) { + if (job.library?.stopAfterCurrent) return false + if (job.library?.chainContinuing) return true + if (job.status === 'uploading' || job.status === 'running') return true + if (job.status === 'queued') return Boolean(job.promptId) || jobAgeMs(job) < SUBMIT_WINDOW_MS + return false +} + export async function videoJobsBusy() { sweepStaleLiveJobs() if (listJobs().some(jobIsLocallySubmitting)) return true + // An empty Comfy queue is not idle: the 3s buffer + last-frame extract between + // shots leaves Comfy empty while the current video still owns the GPU. + if (listJobs().some(liveJobOwnsGpu)) return true + if (listOwnersWithStudioQueues().some(owner => listStudioJobs(owner).some(job => job.status === 'running'))) return true const queue = await fetchLiveQueue() if (queue) return queue.running > 0 || queue.pending > 0 - return listJobs().some((job) => { - if (job.library?.stopAfterCurrent) return false - if (job.status === 'uploading' || job.status === 'running') return true - if (job.status === 'queued') return jobAgeMs(job) < SUBMIT_WINDOW_MS - return false - }) || listPendingJobs().some((pending) => { + return listPendingJobs().some((pending) => { if (!pending.promptId || pending.stopAfterCurrent) return false const live = getJob(pending.jobId) if (live) return jobIsLocallySubmitting(live) || live.status === 'running' || live.status === 'uploading' @@ -525,10 +532,10 @@ function pickHeld(jobs: StudioJob[]) { .filter(job => ( job.status === 'held' && job.shotQueueId - && job.holdForCutIn && !job.pausedByUser + && (job.holdForCutIn === true || job.resumeAutoRun === true) )) - .sort((a, b) => b.updatedAt - a.updatedAt)[0] || null + .sort((a, b) => a.createdAt - b.createdAt)[0] || null } function pickWaiting(jobs: StudioJob[]) { @@ -628,16 +635,17 @@ async function dispatchStudioQueue() { await startStudioJob(cutIn) return } + const held = pickHeld(store.jobs) + if (held?.shotQueueId) { + const resumed = await resumeHeldStudioJob(held) + if (!resumed) scheduleKickRetry() + return + } const waiting = pickWaiting(store.jobs) if (waiting) { await startStudioJob(waiting) return } - const held = pickHeld(store.jobs) - if (held?.shotQueueId) { - const resumed = await resumeHeldStudioJob(held) - if (resumed) return - } } } @@ -649,6 +657,10 @@ async function resumeHeldStudioJob(item: StudioJob) { }) return false } + await patchStudioJob(item.ownerKey, item.id, (job) => { + job.status = 'running' + job.lastError = undefined + }) const { startQueueBurst } = await import('~/server/utils/videoChain') const count = item.resumeAutoRun ? 'all' as const : 1 try { @@ -662,10 +674,12 @@ async function resumeHeldStudioJob(item: StudioJob) { return true } catch (error) { const message = error instanceof Error ? error.message : String(error) + if (/already processing/i.test(message)) return true + const transient = /offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker/i.test(message) await patchStudioJob(item.ownerKey, item.id, (job) => { - job.status = 'error' + job.status = transient ? 'held' : 'error' job.lastError = message - job.holdForCutIn = false + job.holdForCutIn = transient }) return false } @@ -1031,16 +1045,29 @@ export async function startStudioJob(item: StudioJob) { } } +function remainingStudioShots(job: Job) { + const library = job.library + if (!library) return 0 + if (library.queueId) { + const queue = getShotQueue(library.ownerKey, library.queueId) + if (queue) return pendingSegmentCount(queue) + } + return Math.max(0, (library.chainTotal || 1) - (library.chainStep || 1)) +} + +function isTransientComfyError(error?: string) { + return /still busy|COMFY_BUSY|offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker/i.test(String(error || '')) +} + export async function onLiveVideoSettled(job: Job) { const owner = job.library?.ownerKey if (!owner) { kickStudioQueue() return } - const remaining = job.library - ? Math.max(0, (job.library.chainTotal || 1) - (job.library.chainStep || 1)) - : 0 - const failed = job.status === 'error' || job.status === 'cancelled' + const remaining = remainingStudioShots(job) + const wakeFail = job.status === 'error' && remaining > 0 && isTransientComfyError(job.error) + const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail const snapshot = readStore(owner) const rowNow = snapshot.jobs.find(item => item.liveJobId === job.id) || (job.library?.queueId @@ -1085,6 +1112,25 @@ export async function onLiveVideoSettled(job: Job) { row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true row.updatedAt = Date.now() }) + } else if (!failed && remaining > 0) { + await mutateStore(owner, (store) => { + const row = store.jobs.find(item => item.liveJobId === job.id) + || (job.library?.queueId ? store.jobs.find(item => item.shotQueueId === job.library?.queueId) : undefined) + if (!row) { + syncPausedFlag(store) + return + } + const auto = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true + row.status = 'held' + row.liveJobId = undefined + row.holdForCutIn = auto + row.pauseAfterCurrent = false + row.pausedByUser = !auto + row.resumeAutoRun = auto + row.lastError = wakeFail ? job.error : undefined + row.updatedAt = Date.now() + syncPausedFlag(store) + }) } else if (userPause) { await mutateStore(owner, (store) => { const row = store.jobs.find(item => item.liveJobId === job.id) diff --git a/server/utils/videoChain.ts b/server/utils/videoChain.ts index 27774eb..9efe3c5 100644 --- a/server/utils/videoChain.ts +++ b/server/utils/videoChain.ts @@ -451,7 +451,12 @@ export async function continueQueuedExtensions( export async function continueQueuedExtensionsIfNeeded(job: Job) { if (!job.library) return - if (job.status === 'error' || job.status === 'cancelled' || job.status === 'complete') return + if (job.status === 'error' || job.status === 'cancelled') return + if (job.status === 'complete') { + const leftover = remainingAfterCurrentShot(job.library) + if (!leftover.length) return + job.status = 'running' + } if (chainShouldHold(job)) { job.library.stopAfterCurrent = true return @@ -512,6 +517,15 @@ export async function runGeneration(job: Job, params: VideoChainParams) { } export async function startQueueBurst(owner: string, queueId: string, count: number | 'all') { + await ensureComfyReady((status) => { + console.log(JSON.stringify({ + src: 'video-chain', + event: 'wake', + queueId, + state: status.state, + message: status.message + })) + }) const { queue, count: n } = await prepareQueueBurst(owner, queueId, count) const lastIndex = lastCompletedIndex(queue) const clipId = queue.currentClipId @@ -585,15 +599,13 @@ export async function startQueueBurst(owner: string, queueId: string, count: num }).catch(async (error) => { removeExtendTemp(job.library?.extendTmpDir) const message = error instanceof Error ? error.message : String(error) - if (job.status === 'error' || job.status === 'cancelled') { - await finishQueueBurst(owner, queue.id, false, message).catch(() => null) - setQueueJob(queue.id, null) - return - } - job.status = 'error' + const transient = /still busy|COMFY_BUSY|offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker/i.test(message) job.error = message - emitJob(job, { type: 'error', error: message, message }) - await finishQueueBurst(owner, queue.id, false, message).catch(() => null) + if (!transient && job.status !== 'cancelled') { + job.status = 'error' + emitJob(job, { type: 'error', error: message, message }) + } + await finishQueueBurst(owner, queue.id, true, transient ? undefined : message).catch(() => null) setQueueJob(queue.id, null) }).finally(() => { void onLiveVideoSettled(job) diff --git a/server/utils/watch.ts b/server/utils/watch.ts index 3605400..873b654 100644 --- a/server/utils/watch.ts +++ b/server/utils/watch.ts @@ -474,11 +474,12 @@ export function ensurePendingWatch(pending: PendingJob) { await onLiveVideoSettled(job) return } - if (job.status === 'complete') { + 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) }