Fix studio queue stalling after image jobs finish.
Repair all owners before the busy check, stop treating orphan running/held rows as GPU locks, and settle finished Comfy ghosts so waiting items actually start. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+67
-18
@@ -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<string, StudioQueueStore>()
|
||||
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))
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user