diff --git a/pages/index.vue b/pages/index.vue index 753fda1..2fd2dcd 100644 --- a/pages/index.vue +++ b/pages/index.vue @@ -4009,6 +4009,7 @@ async function pokeGpu() { while (Date.now() < deadline) { await pollComfyHealth() if (comfyOk.value || imageComfyOk.value) { + await refreshStudioQueue().catch(() => null) gpuWaking.value = false return } @@ -5435,7 +5436,7 @@ async function cutInStudioJob(id: string) { ? 'Starting now.' : (result.deferred ? 'This job will run after the current shot, as its own clip.' - : 'Queued to start as soon as the GPU is free.')) + : 'Still waiting — try Start now again in a moment.')) await refreshStudioQueue() } catch (error: any) { toast(error?.data?.statusMessage || error?.statusMessage || 'Could not move that job next') @@ -5444,8 +5445,8 @@ async function cutInStudioJob(id: string) { async function retryStudioJob(id: string) { try { - await $fetch(`/api/studio-queue/${id}/retry`, { method: 'POST' }) - toast('Queued to run as soon as the GPU is free.') + const result = await $fetch<{ started?: boolean }>(`/api/studio-queue/${id}/retry`, { method: 'POST' }) + toast(result.started ? 'Starting now.' : 'Queued to run as soon as the GPU is free.') await refreshStudioQueue() } catch (error: any) { toast(error?.data?.statusMessage || error?.statusMessage || 'Could not retry that job') diff --git a/server/api/comfy/wake.post.ts b/server/api/comfy/wake.post.ts index c53039e..64b5643 100644 --- a/server/api/comfy/wake.post.ts +++ b/server/api/comfy/wake.post.ts @@ -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 }) diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts index 6cf4818..bf4b138 100644 --- a/server/utils/studioQueue.ts +++ b/server/utils/studioQueue.ts @@ -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) })