From 8ebbc0fc6e82a957c68b034b2e08e3435fa056ca Mon Sep 17 00:00:00 2001 From: Towsty Date: Sat, 5 Sep 2026 12:38:30 -0500 Subject: [PATCH] Clear Generating as soon as a job saves: settle the studio row on persist and repair finished orphans so the next waiting job can start. Co-authored-by: Cursor --- server/utils/studioQueue.ts | 239 ++++++++++++++++++++---------------- server/utils/watch.ts | 6 + 2 files changed, 136 insertions(+), 109 deletions(-) diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts index be9218e..e2a5e67 100644 --- a/server/utils/studioQueue.ts +++ b/server/utils/studioQueue.ts @@ -394,19 +394,21 @@ export async function videoJobsBusy() { sweepStaleLiveJobs() await reapZombieLiveJobs() 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 - // Studio row still "running" owns the queue slot until onLiveVideoSettled clears it. - // Do not peek under the live job here — that raced settle and started the next item early. - if (listOwnersWithStudioQueues().some(owner => listStudioJobs(owner).some(job => job.status === 'running'))) { - return true - } + // A studio "running" row only blocks while real work is attached. Orphans are cleared in repair. + if (listOwnersWithStudioQueues().some(owner => listStudioJobs(owner).some((job) => { + if (job.status !== 'running') return false + if (pendingAlive(job)) return true + if (!job.liveJobId) return false + const live = getJob(job.liveJobId) + if (!live) return false + if (live.saving || liveJobOwnsGpu(live)) return true + // Live already finished — settle/repair must clear the row; do not pretend the GPU is busy. + return false + }))) return true const queue = await fetchLiveQueue() if (queue) { if (queue.running > 0) return true - // Ghost queue_pending after a clear/restart used to show Ready in the header - // (health only checks running) while forever blocking Start now / auto-dispatch. if (queue.pending > 0) { if (listJobs().some(job => liveJobOwnsGpu(job) || jobIsLocallySubmitting(job))) return true if (listPendingJobs().some((pending) => { @@ -897,7 +899,6 @@ function repairStaleJobs(jobs: StudioJob[]) { job.updatedAt = Date.now() continue } - // GPU sleep mid-chain used to leave held+pausedByUser with a transient lastError — Start now refused those. if (job.status === 'held' && isTransientComfyError(job.lastError)) { job.status = 'waiting' job.liveJobId = undefined @@ -908,7 +909,6 @@ function repairStaleJobs(jobs: StudioJob[]) { job.updatedAt = Date.now() continue } - // Image/music never use shot queues — held without one is a deadlock leftover. if (job.status === 'held' && !job.shotQueueId) { job.status = 'complete' job.liveJobId = undefined @@ -920,14 +920,12 @@ function repairStaleJobs(jobs: StudioJob[]) { } if (job.status !== 'running') continue const live = job.liveJobId ? getJob(job.liveJobId) : undefined - // Live already complete: onLiveVideoSettled owns finalization. Stealing the row here - // started the next queue item while purge/settle was still running → next jobs failed. - if (live && live.status === 'complete') continue - const liveBusy = live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running') - if (liveBusy || pendingAlive(job)) continue - const liveFailed = live && (live.status === 'error' || live.status === 'cancelled') + if (live?.saving) continue + if (live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')) continue + if (pendingAlive(job)) continue const holdPause = job.pauseAfterCurrent === true || job.pausedByUser === true - if (liveFailed) { + // Live finished, failed, or vanished — clear the studio slot so the next waiting job can start. + if (live && (live.status === 'error' || live.status === 'cancelled')) { job.status = live.status === 'cancelled' ? 'cancelled' : 'error' job.liveJobId = undefined job.lastError = live.error || job.lastError || 'Job failed' @@ -936,31 +934,56 @@ function repairStaleJobs(jobs: StudioJob[]) { job.updatedAt = Date.now() continue } - // No live job and no fresh pending — orphaned running row. + // live complete, or live missing after a restart — job is done. if (job.shotQueueId) { - job.status = 'held' - job.liveJobId = undefined - if (holdPause) { - job.pausedByUser = true - job.pauseAfterCurrent = false - job.holdForCutIn = false + const queue = getShotQueue(job.ownerKey, job.shotQueueId) + const pending = queue ? pendingSegmentCount(queue) : 0 + if (pending > 0) { + job.status = 'held' + job.liveJobId = undefined + if (holdPause) { + job.pausedByUser = true + job.pauseAfterCurrent = false + job.holdForCutIn = false + } else { + job.holdForCutIn = true + job.resumeAutoRun = job.payload.queueAutoRun === true || job.resumeAutoRun === true + } + job.updatedAt = Date.now() } else { - job.holdForCutIn = true - job.resumeAutoRun = job.payload.queueAutoRun === true || job.resumeAutoRun === true + job.status = 'complete' + job.liveJobId = undefined + job.pauseAfterCurrent = false + job.pausedByUser = false + job.holdForCutIn = false + job.cutIn = false + job.updatedAt = Date.now() } - job.updatedAt = Date.now() } else { - // No shot queue (image/music): never park as held — that blocked waiting jobs. job.status = 'complete' job.liveJobId = undefined job.pauseAfterCurrent = false job.pausedByUser = false job.holdForCutIn = false + job.cutIn = false job.updatedAt = Date.now() } } } +async function settleAttachedTerminalLives() { + // If a live job already finished but settle did not clear the studio row, finish it + // before deciding the GPU is free. Without this, a missed settle blocks forever. + for (const live of listJobs()) { + if (live.status !== 'complete' && live.status !== 'error' && live.status !== 'cancelled') continue + const owner = live.library?.ownerKey + if (!owner) continue + const attached = readJobs(owner).some(job => job.status === 'running' && job.liveJobId === live.id) + if (!attached) continue + await onLiveVideoSettled(live).catch(() => null) + } +} + async function dispatchStudioQueue() { sweepStaleLiveJobs() // Repair every owner first. videoJobsBusy() looks at all owners — fixing only the @@ -979,6 +1002,7 @@ async function dispatchStudioQueue() { } }) } + await settleAttachedTerminalLives() if (await videoJobsBusy()) { if (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting'))) { scheduleKickRetry() @@ -1557,7 +1581,10 @@ function remainingStudioShots(job: Job) { if (!library) return 0 if (library.queueId) { const queue = getShotQueue(library.ownerKey, library.queueId) - if (queue) return pendingSegmentCount(queue) + // Missing shot-queue file = nothing left to run. Do not fall through to chain math + // or a finished single-shot stays "held/running" forever. + if (!queue) return 0 + return pendingSegmentCount(queue) } // Image/music are one studio slot (multi-pass runs inside the live job). Parking them // as held without a shot-queue id blocked every waiting item behind them. @@ -1565,6 +1592,33 @@ function remainingStudioShots(job: Job) { return Math.max(0, (library.chainTotal || 1) - (library.chainStep || 1)) } +function findStudioRowForLive(jobs: StudioJob[], job: Job) { + const byLive = jobs.find(item => item.liveJobId === job.id) + if (byLive) return byLive + if (job.library?.queueId) { + const byQueue = jobs.find(item => ( + item.shotQueueId === job.library?.queueId + && item.status === 'running' + )) + if (byQueue) return byQueue + } + // Last resort: sole running row for this owner (liveJobId lost after a restart). + const running = jobs.filter(item => item.status === 'running') + if (running.length === 1) return running[0] + return undefined +} + +function clearStudioRowSlot(row: StudioJob, status: StudioJobStatus, error?: string) { + row.status = status + row.liveJobId = undefined + row.cutIn = false + row.holdForCutIn = false + row.pauseAfterCurrent = false + row.pausedByUser = false + row.lastError = error + row.updatedAt = Date.now() +} + export async function onLiveVideoSettled(job: Job) { const owner = job.library?.ownerKey if (!owner) { @@ -1574,31 +1628,26 @@ export async function onLiveVideoSettled(job: Job) { 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 - ? snapshot.jobs.find(item => ( - item.status === 'running' - && item.shotQueueId === job.library?.queueId - )) - : undefined) - const userPause = !failed && rowNow?.holdForCutIn !== true && ( - rowNow?.pauseAfterCurrent === true - || rowNow?.pausedByUser === true - || job.library?.stopAfterCurrent === true - ) - const cutInHold = !failed && remaining > 0 && rowNow?.holdForCutIn === true && !userPause - const interrupted = remaining > 0 && (userPause || cutInHold) - if (interrupted && job.status === 'running') { - job.status = 'complete' - } - if (userPause && remaining > 0) { - await mutateStore(owner, (store) => { - const row = store.jobs.find(item => item.liveJobId === job.id) - if (!row) { - syncPausedFlag(store) - return - } + + await mutateStore(owner, (store) => { + const row = findStudioRowForLive(store.jobs, job) + if (!row || row.status === 'cancelled') { + syncPausedFlag(store) + return + } + + const userPause = !failed && row.holdForCutIn !== true && ( + row.pauseAfterCurrent === true + || row.pausedByUser === true + || job.library?.stopAfterCurrent === true + ) + const cutInHold = !failed && remaining > 0 && row.holdForCutIn === true && !userPause + + if (interruptedRunning(job, remaining, userPause, cutInHold)) { + job.status = 'complete' + } + + if (userPause && remaining > 0) { row.status = 'held' row.liveJobId = undefined row.pauseAfterCurrent = false @@ -1607,27 +1656,18 @@ export async function onLiveVideoSettled(job: Job) { row.resumeAutoRun = row.payload.queueAutoRun === true || row.resumeAutoRun === true row.updatedAt = Date.now() syncPausedFlag(store) - }) - } else if (cutInHold) { - await mutate(owner, (jobs) => { - const row = jobs.find(item => item.liveJobId === job.id) - if (!row) return + return + } + if (cutInHold) { row.status = 'held' row.liveJobId = undefined row.holdForCutIn = true row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true row.updatedAt = Date.now() - }) - } else if (wakeFail) { - // GPU/host blip mid-chain — put back on waiting so wake/Start now can pick it up. - // Never park as user-paused held (that made Start now refuse the job). - 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 || row.status === 'cancelled') { - syncPausedFlag(store) - return - } + syncPausedFlag(store) + return + } + if (wakeFail) { row.status = 'waiting' row.liveJobId = undefined row.cutIn = true @@ -1638,15 +1678,9 @@ export async function onLiveVideoSettled(job: Job) { row.lastError = job.error row.updatedAt = Date.now() syncPausedFlag(store) - }) - } 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 || row.status === 'cancelled') { - syncPausedFlag(store) - return - } + return + } + if (!failed && remaining > 0) { const auto = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true row.status = 'held' row.liveJobId = undefined @@ -1657,39 +1691,26 @@ export async function onLiveVideoSettled(job: Job) { row.lastError = undefined row.updatedAt = Date.now() syncPausedFlag(store) - }) - } else if (userPause) { - await mutateStore(owner, (store) => { - const row = store.jobs.find(item => item.liveJobId === job.id) - if (!row) { - syncPausedFlag(store) - return - } - row.status = job.status === 'cancelled' - ? 'cancelled' - : job.status === 'complete' - ? 'complete' - : 'error' - row.liveJobId = undefined - row.cutIn = false - row.holdForCutIn = false - row.pauseAfterCurrent = false - row.pausedByUser = false - row.lastError = job.error - row.updatedAt = Date.now() - syncPausedFlag(store) - }) - } else { - const status = job.status === 'cancelled' + return + } + + // Normal finish: clear Generating and free the queue for the next waiting job. + const status: StudioJobStatus = job.status === 'cancelled' ? 'cancelled' - : job.status === 'complete' - ? 'complete' - : 'error' - await markStudioSettled(owner, job.id, status, job.error) - } + : failed + ? 'error' + : 'complete' + clearStudioRowSlot(row, status, failed ? job.error : undefined) + syncPausedFlag(store) + }) + await kickStudioQueue() } +function interruptedRunning(job: Job, remaining: number, userPause: boolean, cutInHold: boolean) { + return remaining > 0 && (userPause || cutInHold) && job.status === 'running' +} + function applyLiveJobPause(live: Job | undefined, pause: boolean, restoreAutoRun?: boolean) { if (!live?.library) return live.library.stopAfterCurrent = pause diff --git a/server/utils/watch.ts b/server/utils/watch.ts index 9aeece2..acc2bd8 100644 --- a/server/utils/watch.ts +++ b/server/utils/watch.ts @@ -363,6 +363,12 @@ export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Pr }) deletePendingJob(job.id) job.saving = false + // Persist path: clear the studio "Generating" row as soon as the library save + // finishes. Waiting for runGeneration's finally used to leave jobs stuck in the queue. + if (persist && job.library?.ownerKey) { + const { onLiveVideoSettled } = await import('~/server/utils/studioQueue') + await onLiveVideoSettled(job).catch(() => null) + } } catch (saveError) { job.saving = false removeExtendTemp(job.library?.extendTmpDir)