Fail zombie studio jobs when Comfy's queue is empty after a restart.
A live job stuck at 8% with no Comfy prompt was treated as still owning the GPU, so waiting work never started. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -480,9 +480,11 @@ export async function isComfyPromptDropped(promptId: string) {
|
|||||||
if (queue.running > 0 || queue.pending > 0) return false
|
if (queue.running > 0 || queue.pending > 0) return false
|
||||||
const history = await fetchHistory(promptId)
|
const history = await fetchHistory(promptId)
|
||||||
const entry = history?.[promptId] as { status?: { status_str?: string; completed?: boolean } } | undefined
|
const entry = history?.[promptId] as { status?: { status_str?: string; completed?: boolean } } | undefined
|
||||||
if (!entry) return false
|
// Empty Comfy + no history for this id: wiped by restart/clear, or never accepted.
|
||||||
|
if (!entry) return true
|
||||||
const inspected = inspectHistory(history, promptId)
|
const inspected = inspectHistory(history, promptId)
|
||||||
if (inspected.video) return false
|
if (inspected.video) return false
|
||||||
|
if (extractAudio(history, promptId)) return false
|
||||||
const status = entry.status?.status_str
|
const status = entry.status?.status_str
|
||||||
if (status === 'error' || status === 'interrupted') return true
|
if (status === 'error' || status === 'interrupted') return true
|
||||||
if (entry.status?.completed) return true
|
if (entry.status?.completed) return true
|
||||||
|
|||||||
@@ -3,7 +3,7 @@ import { join } from 'node:path'
|
|||||||
import { getJob, listJobs, emitJob, type Job } from '~/server/utils/jobs'
|
import { getJob, listJobs, emitJob, type Job } from '~/server/utils/jobs'
|
||||||
import { listPendingJobs, patchPendingJob, readPendingJob, deletePendingJob } from '~/server/utils/pending'
|
import { listPendingJobs, patchPendingJob, readPendingJob, deletePendingJob } from '~/server/utils/pending'
|
||||||
import { getShotQueue, pendingSegmentCount } from '~/server/utils/shotQueue'
|
import { getShotQueue, pendingSegmentCount } from '~/server/utils/shotQueue'
|
||||||
import { fetchLiveQueue } from '~/server/utils/comfy'
|
import { fetchLiveQueue, isComfyPromptDropped } from '~/server/utils/comfy'
|
||||||
import { isLtxWorkflow, isTextToVideo, isXaigenStudio, LTX_DISABLED_MESSAGE, ltxWorkflowEnabled, parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
|
import { isLtxWorkflow, isTextToVideo, isXaigenStudio, LTX_DISABLED_MESSAGE, ltxWorkflowEnabled, parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
|
||||||
import { persistLoraFields } from '~/server/utils/loras'
|
import { persistLoraFields } from '~/server/utils/loras'
|
||||||
import { resolveLoraStack } from '~/utils/loras'
|
import { resolveLoraStack } from '~/utils/loras'
|
||||||
@@ -250,6 +250,7 @@ export function summarizeStudioJob(job: StudioJob) {
|
|||||||
const QUEUED_GRACE_MS = 4 * 60 * 1000
|
const QUEUED_GRACE_MS = 4 * 60 * 1000
|
||||||
const SUBMIT_WINDOW_MS = 45 * 1000
|
const SUBMIT_WINDOW_MS = 45 * 1000
|
||||||
const PENDING_ORPHAN_MS = 2 * 60 * 1000
|
const PENDING_ORPHAN_MS = 2 * 60 * 1000
|
||||||
|
const ZOMBIE_EMPTY_COMFY_MS = 90 * 1000
|
||||||
|
|
||||||
function jobAgeMs(job: { startedAt: number }) {
|
function jobAgeMs(job: { startedAt: number }) {
|
||||||
return Date.now() - job.startedAt
|
return Date.now() - job.startedAt
|
||||||
@@ -264,27 +265,67 @@ export function jobIsLocallySubmitting(job: Job) {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function failZombieLiveJob(job: Job, error: string) {
|
||||||
|
if (job.status === 'complete' || job.status === 'error' || job.status === 'cancelled') return
|
||||||
|
job.status = 'error'
|
||||||
|
job.error = error
|
||||||
|
emitJob(job, { type: 'error', error, message: error })
|
||||||
|
void onLiveVideoSettled(job).catch(() => null)
|
||||||
|
}
|
||||||
|
|
||||||
function sweepStaleLiveJobs() {
|
function sweepStaleLiveJobs() {
|
||||||
const now = Date.now()
|
const now = Date.now()
|
||||||
for (const job of listJobs()) {
|
for (const job of listJobs()) {
|
||||||
if (job.status !== 'queued') continue
|
if (job.saving) continue
|
||||||
if (job.promptId) continue
|
if (job.status === 'queued' && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) {
|
||||||
if (now - job.startedAt < QUEUED_GRACE_MS) continue
|
failZombieLiveJob(job, 'Job never started')
|
||||||
job.status = 'error'
|
continue
|
||||||
job.error = job.error || 'Job never started'
|
}
|
||||||
|
if ((job.status === 'running' || job.status === 'uploading') && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) {
|
||||||
|
failZombieLiveJob(job, 'Job never started on ComfyUI')
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Kill live jobs that still say "running" after Comfy was restarted or cleared. */
|
||||||
|
async function reapZombieLiveJobs() {
|
||||||
|
const queue = await fetchLiveQueue()
|
||||||
|
if (!queue) return
|
||||||
|
const comfyBusy = queue.running > 0 || queue.pending > 0
|
||||||
|
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
|
||||||
|
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)) {
|
||||||
|
failZombieLiveJob(job, 'ComfyUI lost this job (empty queue after a restart or clear). Generate again.')
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
function liveJobOwnsGpu(job: Job) {
|
function liveJobOwnsGpu(job: Job) {
|
||||||
if (job.library?.stopAfterCurrent) return false
|
if (job.library?.stopAfterCurrent) return false
|
||||||
if (job.library?.chainContinuing) return true
|
if (job.library?.chainContinuing) return true
|
||||||
if (job.status === 'uploading' || job.status === 'running') return true
|
if (job.saving) return true
|
||||||
|
if (job.status === 'uploading') return true
|
||||||
|
if (job.status === 'running') {
|
||||||
|
// Without a prompt id past the submit window, this is a zombie — do not block the queue.
|
||||||
|
if (!job.promptId && jobAgeMs(job) >= SUBMIT_WINDOW_MS) return false
|
||||||
|
return true
|
||||||
|
}
|
||||||
if (job.status === 'queued') return Boolean(job.promptId) || jobAgeMs(job) < SUBMIT_WINDOW_MS
|
if (job.status === 'queued') return Boolean(job.promptId) || jobAgeMs(job) < SUBMIT_WINDOW_MS
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function videoJobsBusy() {
|
export async function videoJobsBusy() {
|
||||||
sweepStaleLiveJobs()
|
sweepStaleLiveJobs()
|
||||||
|
await reapZombieLiveJobs()
|
||||||
if (listJobs().some(jobIsLocallySubmitting)) return true
|
if (listJobs().some(jobIsLocallySubmitting)) return true
|
||||||
// An empty Comfy queue is not idle: the 3s buffer + last-frame extract between
|
// 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.
|
// shots leaves Comfy empty while the current video still owns the GPU.
|
||||||
@@ -649,7 +690,17 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
|||||||
const live = job.liveJobId ? getJob(job.liveJobId) : undefined
|
const live = job.liveJobId ? getJob(job.liveJobId) : undefined
|
||||||
const liveBusy = live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')
|
const liveBusy = live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')
|
||||||
if (liveBusy || pendingAlive(job)) continue
|
if (liveBusy || pendingAlive(job)) continue
|
||||||
|
const liveFailed = live && (live.status === 'error' || live.status === 'cancelled')
|
||||||
const holdPause = job.pauseAfterCurrent === true || job.pausedByUser === true
|
const holdPause = job.pauseAfterCurrent === true || job.pausedByUser === true
|
||||||
|
if (liveFailed) {
|
||||||
|
job.status = live.status === 'cancelled' ? 'cancelled' : 'error'
|
||||||
|
job.liveJobId = undefined
|
||||||
|
job.lastError = live.error || job.lastError || 'Job failed'
|
||||||
|
job.cutIn = false
|
||||||
|
job.holdForCutIn = false
|
||||||
|
job.updatedAt = Date.now()
|
||||||
|
continue
|
||||||
|
}
|
||||||
if (job.shotQueueId) {
|
if (job.shotQueueId) {
|
||||||
job.status = 'held'
|
job.status = 'held'
|
||||||
job.liveJobId = undefined
|
job.liveJobId = undefined
|
||||||
|
|||||||
Reference in New Issue
Block a user