Unstick studio queue when Ready but waiting jobs will not start.
A failed held shot-batch no longer blocks waiting work, ghost Comfy pending no longer freezes Start now, and Start now force-starts when the GPU is idle. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+72
-12
@@ -358,7 +358,20 @@ export async function videoJobsBusy() {
|
||||
return pendingAlive(job)
|
||||
}))) return true
|
||||
const queue = await fetchLiveQueue()
|
||||
if (queue) return queue.running > 0 || queue.pending > 0
|
||||
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) => {
|
||||
if (!pending.promptId || pending.stopAfterCurrent) return false
|
||||
return Date.now() - pending.startedAt < PENDING_ORPHAN_MS
|
||||
})) return true
|
||||
return false
|
||||
}
|
||||
return false
|
||||
}
|
||||
return listPendingJobs().some((pending) => {
|
||||
if (!pending.promptId || pending.stopAfterCurrent) return false
|
||||
const live = getJob(pending.jobId)
|
||||
@@ -610,8 +623,8 @@ export async function cancelOrphanQueue(owner: string, queueId: string) {
|
||||
|
||||
export async function requestCutIn(owner: string, id: string) {
|
||||
const jobsNow = readJobs(owner)
|
||||
// Only a live "running" row is generating. Held leftovers must not steal Start now.
|
||||
const running = jobsNow.find(job => job.status === 'running')
|
||||
|| jobsNow.find(job => job.status === 'held')
|
||||
const target = jobsNow.find(job => job.id === id)
|
||||
if (!target) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' })
|
||||
if (target.status !== 'waiting') {
|
||||
@@ -646,7 +659,30 @@ export async function requestCutIn(owner: string, id: string) {
|
||||
}
|
||||
}
|
||||
await kickStudioQueue()
|
||||
return updated
|
||||
// GPU idle + still waiting = dispatcher missed it (held deadlock, false busy, etc). Force start.
|
||||
let started = false
|
||||
const after = readJobs(owner).find(job => job.id === id)
|
||||
if (after?.status === 'running') {
|
||||
started = true
|
||||
} else if (after?.status === 'waiting') {
|
||||
const queue = await fetchLiveQueue().catch(() => null)
|
||||
const comfyRendering = Boolean(queue && queue.running > 0)
|
||||
const liveOwns = listJobs().some(liveJobOwnsGpu)
|
||||
if (!comfyRendering && !liveOwns) {
|
||||
try {
|
||||
await startStudioJob(after)
|
||||
started = readJobs(owner).find(job => job.id === id)?.status === 'running'
|
||||
} catch {
|
||||
started = false
|
||||
}
|
||||
}
|
||||
}
|
||||
const latest = readJobs(owner).find(job => job.id === id) || updated
|
||||
return {
|
||||
...summarizeStudioJob(latest),
|
||||
started,
|
||||
deferred: Boolean(running) && !started
|
||||
}
|
||||
}
|
||||
|
||||
export async function retryStudioJob(owner: string, id: string) {
|
||||
@@ -860,9 +896,8 @@ async function dispatchStudioQueue() {
|
||||
// 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) => {
|
||||
await mutateStore(owner, (current) => {
|
||||
migrateGlobalPause(current)
|
||||
repairStaleJobs(current.jobs)
|
||||
const next = pruneDone(current.jobs)
|
||||
@@ -873,7 +908,6 @@ async function dispatchStudioQueue() {
|
||||
jobs: current.jobs.map(item => structuredClone(item))
|
||||
}
|
||||
})
|
||||
stores.set(owner, store)
|
||||
}
|
||||
if (await videoJobsBusy()) {
|
||||
if (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting'))) {
|
||||
@@ -881,8 +915,9 @@ async function dispatchStudioQueue() {
|
||||
}
|
||||
return
|
||||
}
|
||||
// Re-read after videoJobsBusy — reap/settle may have rewritten disk rows.
|
||||
for (const owner of owners) {
|
||||
const store = stores.get(owner) || readStore(owner)
|
||||
const store = readStore(owner)
|
||||
const cutIn = pickCutIn(store.jobs)
|
||||
if (cutIn) {
|
||||
await startStudioJob(cutIn)
|
||||
@@ -891,8 +926,9 @@ async function dispatchStudioQueue() {
|
||||
const held = pickHeld(store.jobs)
|
||||
if (held?.shotQueueId) {
|
||||
const resumed = await resumeHeldStudioJob(held)
|
||||
if (!resumed) scheduleKickRetry()
|
||||
return
|
||||
if (resumed) return
|
||||
// Failed resume used to `return` here and starve every waiting job forever.
|
||||
scheduleKickRetry()
|
||||
}
|
||||
const waiting = pickWaiting(store.jobs)
|
||||
if (waiting) {
|
||||
@@ -913,8 +949,20 @@ async function resumeHeldStudioJob(item: StudioJob) {
|
||||
})
|
||||
return false
|
||||
}
|
||||
const { isShotQueueCancelled } = await import('~/server/utils/shotQueue')
|
||||
if (isShotQueueCancelled(latest.shotQueueId)) return false
|
||||
const { isShotQueueCancelled, getShotQueue } = await import('~/server/utils/shotQueue')
|
||||
const shotQueue = getShotQueue(item.ownerKey, latest.shotQueueId)
|
||||
if (!shotQueue || isShotQueueCancelled(latest.shotQueueId) || pendingSegmentCount(shotQueue) <= 0) {
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
if (job.status === 'cancelled') return
|
||||
job.status = 'complete'
|
||||
job.liveJobId = undefined
|
||||
job.holdForCutIn = false
|
||||
job.pauseAfterCurrent = false
|
||||
job.pausedByUser = false
|
||||
job.updatedAt = Date.now()
|
||||
})
|
||||
return false
|
||||
}
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
if (job.status === 'cancelled') return
|
||||
job.status = 'running'
|
||||
@@ -936,7 +984,19 @@ async function resumeHeldStudioJob(item: StudioJob) {
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
if (/already processing/i.test(message)) return true
|
||||
if (/removed|not found/i.test(message)) return false
|
||||
if (/removed|not found/i.test(message)) {
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
if (job.status === 'cancelled') return
|
||||
job.status = 'complete'
|
||||
job.liveJobId = undefined
|
||||
job.holdForCutIn = false
|
||||
job.pauseAfterCurrent = false
|
||||
job.pausedByUser = false
|
||||
job.lastError = message
|
||||
job.updatedAt = Date.now()
|
||||
}).catch(() => null)
|
||||
return false
|
||||
}
|
||||
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) => {
|
||||
if (job.status === 'cancelled') return
|
||||
|
||||
Reference in New Issue
Block a user