Fix studio queue stuck after GPU sleep: clear ghost live ownership so Start now and wake can actually dispatch.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -1,5 +1,9 @@
|
||||
import { requestComfyWake } from '~/server/utils/comfyLifecycle'
|
||||
import { kickStudioQueue } from '~/server/utils/studioQueue'
|
||||
|
||||
export default defineEventHandler(async () => {
|
||||
return await requestComfyWake()
|
||||
const result = await requestComfyWake()
|
||||
// After sleep/crash the queue often sits waiting with cutIn and never fires on its own.
|
||||
void kickStudioQueue()
|
||||
return result
|
||||
})
|
||||
|
||||
+148
-45
@@ -298,9 +298,12 @@ async function reapZombieLiveJobs() {
|
||||
const { fetchHistory } = await import('~/server/utils/comfy')
|
||||
for (const job of listJobs()) {
|
||||
if (job.saving) continue
|
||||
// Stuck extension loops claimed the GPU forever after a crash/sleep.
|
||||
if (job.library?.chainContinuing && jobAgeMs(job) >= ZOMBIE_EMPTY_COMFY_MS) {
|
||||
job.library.chainContinuing = false
|
||||
}
|
||||
if (job.status !== 'running' && job.status !== 'uploading' && job.status !== 'queued') continue
|
||||
if (jobIsLocallySubmitting(job)) continue
|
||||
if (job.library?.chainContinuing) continue
|
||||
const zombieMs = (job.kind === 'edit' || job.kind === 'music') ? 45_000 : ZOMBIE_EMPTY_COMFY_MS
|
||||
if (jobAgeMs(job) < zombieMs) continue
|
||||
if (!job.promptId) {
|
||||
@@ -322,13 +325,25 @@ async function reapZombieLiveJobs() {
|
||||
// Mark complete and advance the studio queue so waiting items can run.
|
||||
job.status = 'complete'
|
||||
job.error = undefined
|
||||
if (job.library) job.library.chainContinuing = false
|
||||
await onLiveVideoSettled(job).catch(() => null)
|
||||
}
|
||||
}
|
||||
|
||||
function liveJobOwnsGpu(job: Job) {
|
||||
if (job.library?.stopAfterCurrent) return false
|
||||
if (job.library?.chainContinuing) return true
|
||||
if (job.library?.chainContinuing) {
|
||||
// Never let a stuck chainContinuing flag block the queue after Comfy went idle.
|
||||
if (job.status === 'complete' || job.status === 'error' || job.status === 'cancelled') {
|
||||
job.library.chainContinuing = false
|
||||
return false
|
||||
}
|
||||
if (jobAgeMs(job) >= ZOMBIE_EMPTY_COMFY_MS * 2) {
|
||||
job.library.chainContinuing = false
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
if (job.saving) return true
|
||||
if (job.status === 'uploading') return true
|
||||
if (job.status === 'running') {
|
||||
@@ -340,6 +355,45 @@ function liveJobOwnsGpu(job: Job) {
|
||||
return false
|
||||
}
|
||||
|
||||
/** Drop soft in-memory GPU claims when Comfy itself is idle so Start now can fire. */
|
||||
function clearSoftGpuClaimsWhenComfyIdle(opts?: { aggressive?: boolean }) {
|
||||
const ageCap = opts?.aggressive ? SUBMIT_WINDOW_MS : ZOMBIE_EMPTY_COMFY_MS
|
||||
for (const live of listJobs()) {
|
||||
if (live.library?.chainContinuing && jobAgeMs(live) >= ageCap) {
|
||||
live.library.chainContinuing = false
|
||||
}
|
||||
if (live.saving && jobAgeMs(live) >= ZOMBIE_EMPTY_COMFY_MS) {
|
||||
live.saving = false
|
||||
}
|
||||
if (live.status !== 'running' && live.status !== 'queued' && live.status !== 'uploading') continue
|
||||
if (jobIsLocallySubmitting(live)) continue
|
||||
if (jobAgeMs(live) < ageCap) continue
|
||||
failZombieLiveJob(live, 'Cleared stuck live job so the studio queue can start')
|
||||
}
|
||||
}
|
||||
|
||||
/** When Comfy is idle but a waiting cut-in still sits, clear ghosts and start it. */
|
||||
async function forceStartWaitingIfGpuIdle(owner: string, id: string) {
|
||||
const after = readJobs(owner).find(job => job.id === id)
|
||||
if (!after || after.status !== 'waiting') {
|
||||
return after?.status === 'running'
|
||||
}
|
||||
await reapZombieLiveJobs()
|
||||
const queue = await fetchLiveQueue().catch(() => null)
|
||||
const comfyRendering = Boolean(queue && queue.running > 0)
|
||||
if (comfyRendering) return false
|
||||
// Soft ownership (chainContinuing / orphan running) used to block Start now forever
|
||||
// while Comfy's queue was already empty.
|
||||
clearSoftGpuClaimsWhenComfyIdle({ aggressive: true })
|
||||
if (listJobs().some(liveJobOwnsGpu)) return false
|
||||
try {
|
||||
await startStudioJob(after)
|
||||
return readJobs(owner).find(job => job.id === id)?.status === 'running'
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
export async function videoJobsBusy() {
|
||||
sweepStaleLiveJobs()
|
||||
await reapZombieLiveJobs()
|
||||
@@ -659,23 +713,10 @@ export async function requestCutIn(owner: string, id: string) {
|
||||
}
|
||||
}
|
||||
await kickStudioQueue()
|
||||
// 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
|
||||
}
|
||||
}
|
||||
// GPU idle + still waiting = dispatcher missed it (ghost live ownership, etc). Force start.
|
||||
let started = readJobs(owner).find(job => job.id === id)?.status === 'running'
|
||||
if (!started) {
|
||||
started = await forceStartWaitingIfGpuIdle(owner, id)
|
||||
}
|
||||
const latest = readJobs(owner).find(job => job.id === id) || updated
|
||||
return {
|
||||
@@ -712,7 +753,11 @@ export async function retryStudioJob(owner: string, id: string) {
|
||||
return structuredClone(job)
|
||||
})
|
||||
await kickStudioQueue()
|
||||
return updated
|
||||
const started = await forceStartWaitingIfGpuIdle(owner, id)
|
||||
return {
|
||||
...summarizeStudioJob(readJobs(owner).find(job => job.id === id) || updated),
|
||||
started
|
||||
}
|
||||
}
|
||||
|
||||
function pruneDone(jobs: StudioJob[]) {
|
||||
@@ -823,7 +868,26 @@ function pendingAlive(job: StudioJob) {
|
||||
}
|
||||
|
||||
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 || ''))
|
||||
return /still busy|COMFY_BUSY|offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker|refused to start|API never answered|cold start|ECONNREFUSED|ETIMEDOUT|fetch failed|502|503/i.test(String(error || ''))
|
||||
}
|
||||
|
||||
async function parkStudioOnStartFailure(owner: string, id: string, message: string) {
|
||||
const transient = isTransientComfyError(message)
|
||||
await patchStudioJob(owner, id, (job) => {
|
||||
if (transient) {
|
||||
job.status = 'waiting'
|
||||
job.cutIn = true
|
||||
job.liveJobId = undefined
|
||||
job.lastError = message
|
||||
job.pausedByUser = false
|
||||
job.pauseAfterCurrent = false
|
||||
job.holdForCutIn = false
|
||||
} else {
|
||||
job.status = 'error'
|
||||
job.lastError = message
|
||||
}
|
||||
job.updatedAt = Date.now()
|
||||
}).catch(() => null)
|
||||
}
|
||||
|
||||
function repairStaleJobs(jobs: StudioJob[]) {
|
||||
@@ -839,6 +903,19 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
||||
job.status = 'waiting'
|
||||
job.liveJobId = undefined
|
||||
job.lastError = undefined
|
||||
job.pausedByUser = false
|
||||
job.holdForCutIn = false
|
||||
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
|
||||
job.lastError = undefined
|
||||
job.pausedByUser = false
|
||||
job.pauseAfterCurrent = false
|
||||
job.holdForCutIn = false
|
||||
job.updatedAt = Date.now()
|
||||
continue
|
||||
}
|
||||
@@ -910,10 +987,16 @@ async function dispatchStudioQueue() {
|
||||
})
|
||||
}
|
||||
if (await videoJobsBusy()) {
|
||||
if (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting'))) {
|
||||
scheduleKickRetry()
|
||||
const queue = await fetchLiveQueue().catch(() => null)
|
||||
if (queue && queue.running === 0 && queue.pending === 0) {
|
||||
clearSoftGpuClaimsWhenComfyIdle()
|
||||
}
|
||||
if (await videoJobsBusy()) {
|
||||
if (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting'))) {
|
||||
scheduleKickRetry()
|
||||
}
|
||||
return
|
||||
}
|
||||
return
|
||||
}
|
||||
// Re-read after videoJobsBusy — reap/settle may have rewritten disk rows.
|
||||
for (const owner of owners) {
|
||||
@@ -997,12 +1080,23 @@ async function resumeHeldStudioJob(item: StudioJob) {
|
||||
}).catch(() => null)
|
||||
return false
|
||||
}
|
||||
const transient = /offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker/i.test(message)
|
||||
const transient = isTransientComfyError(message)
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
if (job.status === 'cancelled') return
|
||||
job.status = transient ? 'held' : 'error'
|
||||
job.lastError = message
|
||||
job.holdForCutIn = transient
|
||||
if (transient) {
|
||||
job.status = 'waiting'
|
||||
job.cutIn = true
|
||||
job.liveJobId = undefined
|
||||
job.lastError = message
|
||||
job.holdForCutIn = false
|
||||
job.pausedByUser = false
|
||||
job.pauseAfterCurrent = false
|
||||
} else {
|
||||
job.status = 'error'
|
||||
job.lastError = message
|
||||
job.holdForCutIn = false
|
||||
}
|
||||
job.updatedAt = Date.now()
|
||||
})
|
||||
return false
|
||||
}
|
||||
@@ -1205,10 +1299,7 @@ async function startStudioEditJob(item: StudioJob) {
|
||||
live.status = 'error'
|
||||
live.error = message
|
||||
}
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
job.status = 'error'
|
||||
job.lastError = message
|
||||
}).catch(() => null)
|
||||
await parkStudioOnStartFailure(item.ownerKey, item.id, message)
|
||||
kickStudioQueue()
|
||||
}
|
||||
}
|
||||
@@ -1238,10 +1329,7 @@ async function startStudioExtendJob(item: StudioJob) {
|
||||
live.status = 'error'
|
||||
live.error = message
|
||||
}
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
job.status = 'error'
|
||||
job.lastError = message
|
||||
}).catch(() => null)
|
||||
await parkStudioOnStartFailure(item.ownerKey, item.id, message)
|
||||
kickStudioQueue()
|
||||
}
|
||||
}
|
||||
@@ -1275,10 +1363,7 @@ async function startStudioMusicJob(item: StudioJob) {
|
||||
live.status = 'error'
|
||||
live.error = message
|
||||
}
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
job.status = 'error'
|
||||
job.lastError = message
|
||||
}).catch(() => null)
|
||||
await parkStudioOnStartFailure(item.ownerKey, item.id, message)
|
||||
kickStudioQueue()
|
||||
}
|
||||
}
|
||||
@@ -1475,10 +1560,7 @@ export async function startStudioJob(item: StudioJob) {
|
||||
live.status = 'error'
|
||||
live.error = message
|
||||
}
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
job.status = 'error'
|
||||
job.lastError = message
|
||||
}).catch(() => null)
|
||||
await parkStudioOnStartFailure(item.ownerKey, item.id, message)
|
||||
kickStudioQueue()
|
||||
}
|
||||
}
|
||||
@@ -1549,6 +1631,27 @@ 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 (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
|
||||
}
|
||||
row.status = 'waiting'
|
||||
row.liveJobId = undefined
|
||||
row.cutIn = true
|
||||
row.holdForCutIn = false
|
||||
row.pauseAfterCurrent = false
|
||||
row.pausedByUser = false
|
||||
row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
|
||||
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)
|
||||
@@ -1564,7 +1667,7 @@ export async function onLiveVideoSettled(job: Job) {
|
||||
row.pauseAfterCurrent = false
|
||||
row.pausedByUser = !auto
|
||||
row.resumeAutoRun = auto
|
||||
row.lastError = wakeFail ? job.error : undefined
|
||||
row.lastError = undefined
|
||||
row.updatedAt = Date.now()
|
||||
syncPausedFlag(store)
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user