diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts index c1f5c41..11c19bc 100644 --- a/server/utils/studioQueue.ts +++ b/server/utils/studioQueue.ts @@ -3,7 +3,7 @@ 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, pendingSegmentCount } from '~/server/utils/shotQueue' -import { fetchLiveQueue, isComfyPromptDropped } from '~/server/utils/comfy' +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' import { resolveLoraStack } from '~/utils/loras' @@ -291,21 +291,35 @@ function sweepStaleLiveJobs() { async function reapZombieLiveJobs() { const queue = await fetchLiveQueue() if (!queue) return - const comfyBusy = queue.running > 0 || queue.pending > 0 + if (queue.running > 0 || queue.pending > 0) return + const { fetchHistory } = await import('~/server/utils/comfy') for (const job of listJobs()) { if (job.saving) continue if (job.status !== 'running' && job.status !== 'uploading' && job.status !== 'queued') continue if (jobIsLocallySubmitting(job)) continue if (job.library?.chainContinuing) continue - if (jobAgeMs(job) < ZOMBIE_EMPTY_COMFY_MS) continue + const zombieMs = (job.kind === 'edit' || job.kind === 'music') ? 45_000 : ZOMBIE_EMPTY_COMFY_MS + if (jobAgeMs(job) < zombieMs) continue if (!job.promptId) { if (jobAgeMs(job) >= QUEUED_GRACE_MS) failZombieLiveJob(job, 'Job never started on ComfyUI') continue } - if (comfyBusy) continue - if (await isComfyPromptDropped(job.promptId)) { + const history = await fetchHistory(job.promptId).catch(() => null) + const entry = history?.[job.promptId] as { status?: { status_str?: string; completed?: boolean } } | undefined + if (!entry) { failZombieLiveJob(job, 'ComfyUI lost this job (empty queue after a restart or clear). Generate again.') + continue } + const status = entry.status?.status_str + if (status === 'error' || status === 'interrupted') { + failZombieLiveJob(job, job.error || `ComfyUI ${status}`) + continue + } + // History still present + Comfy idle: waiter never settled (common for image jobs). + // Mark complete and advance the studio queue so waiting items can run. + job.status = 'complete' + job.error = undefined + await onLiveVideoSettled(job).catch(() => null) } } @@ -330,7 +344,16 @@ export async function videoJobsBusy() { // 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 + // Studio "running" only blocks when a live job still owns the GPU, or a fresh + // pending file exists. Orphan disk rows used to block the whole queue forever. + if (listOwnersWithStudioQueues().some(owner => listStudioJobs(owner).some((job) => { + if (job.status !== 'running') return false + if (job.liveJobId) { + const live = getJob(job.liveJobId) + if (live && liveJobOwnsGpu(live)) return true + } + return pendingAlive(job) + }))) return true const queue = await fetchLiveQueue() if (queue) return queue.running > 0 || queue.pending > 0 return listPendingJobs().some((pending) => { @@ -746,12 +769,17 @@ function scheduleKickRetry() { } function pendingAlive(job: StudioJob) { - if (job.liveJobId && readPendingJob(job.liveJobId)?.promptId) return true + const freshEnough = (startedAt: number) => Date.now() - startedAt < PENDING_ORPHAN_MS + if (job.liveJobId) { + const pending = readPendingJob(job.liveJobId) + if (pending?.promptId && freshEnough(pending.startedAt)) return true + } if (!job.shotQueueId) return false return listPendingJobs().some(pending => ( pending.ownerKey === job.ownerKey && pending.queueId === job.shotQueueId && Boolean(pending.promptId) + && freshEnough(pending.startedAt) )) } @@ -775,6 +803,16 @@ 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 + job.holdForCutIn = false + job.pauseAfterCurrent = false + job.pausedByUser = false + job.updatedAt = Date.now() + continue + } if (job.status !== 'running') continue const live = job.liveJobId ? getJob(job.liveJobId) : undefined const liveBusy = live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running') @@ -803,12 +841,12 @@ function repairStaleJobs(jobs: StudioJob[]) { } job.updatedAt = Date.now() } else { - job.status = holdPause ? 'held' : 'complete' + // No shot queue (image/music): never park as held — that blocked waiting jobs. + job.status = 'complete' job.liveJobId = undefined - if (holdPause) { - job.pausedByUser = true - job.pauseAfterCurrent = false - } + job.pauseAfterCurrent = false + job.pausedByUser = false + job.holdForCutIn = false job.updatedAt = Date.now() } } @@ -816,7 +854,11 @@ function repairStaleJobs(jobs: StudioJob[]) { async function dispatchStudioQueue() { sweepStaleLiveJobs() - for (const owner of listOwnersWithStudioQueues()) { + // Repair every owner first. videoJobsBusy() looks at all owners — fixing only the + // current owner left a stuck "running" row on a later owner blocking the GPU forever. + const owners = listOwnersWithStudioQueues() + const stores = new Map() + for (const owner of owners) { const store = await mutateStore(owner, (current) => { migrateGlobalPause(current) repairStaleJobs(current.jobs) @@ -828,12 +870,16 @@ async function dispatchStudioQueue() { jobs: current.jobs.map(item => structuredClone(item)) } }) - if (await videoJobsBusy()) { - if (listOwnersWithStudioQueues().some(owner => readJobs(owner).some(job => job.status === 'waiting'))) { - scheduleKickRetry() - } - return + stores.set(owner, store) + } + if (await videoJobsBusy()) { + if (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting'))) { + scheduleKickRetry() } + return + } + for (const owner of owners) { + const store = stores.get(owner) || readStore(owner) const cutIn = pickCutIn(store.jobs) if (cutIn) { await startStudioJob(cutIn) @@ -1368,6 +1414,9 @@ function remainingStudioShots(job: Job) { const queue = getShotQueue(library.ownerKey, library.queueId) if (queue) 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. + if (job.kind !== 'video') return 0 return Math.max(0, (library.chainTotal || 1) - (library.chainStep || 1)) }