Fix queue dispatch deadlocks and stale completion recovery

This commit is contained in:
Towsty
2026-09-06 16:08:50 -05:00
parent 069108ae44
commit 911ddeab55
4 changed files with 189 additions and 10 deletions
+1
View File
@@ -3,6 +3,7 @@
"private": true,
"type": "module",
"scripts": {
"test": "node --test tests/*.test.mjs",
"build": "nuxt build",
"dev": "nuxt dev --host 0.0.0.0 --port 3000",
"preview": "nuxt preview --host 0.0.0.0 --port 3000",
+3 -1
View File
@@ -211,7 +211,9 @@ export async function freeComfyVram() {
}
export async function fetchHistory(promptId: string) {
const res = await comfyFetch(`/history/${encodeURIComponent(promptId)}`)
const res = await comfyFetch(`/history/${encodeURIComponent(promptId)}`, {
signal: AbortSignal.timeout(12000)
})
if (!res.ok) return null
return (await res.json()) as Record<string, unknown>
}
+17 -9
View File
@@ -309,7 +309,14 @@ async function reapZombieLiveJobs() {
if (jobAgeMs(job) >= QUEUED_GRACE_MS) failZombieLiveJob(job, 'Job never started on ComfyUI')
continue
}
const history = await fetchHistory(job.promptId).catch(() => null)
let history: Awaited<ReturnType<typeof fetchHistory>>
try {
history = await fetchHistory(job.promptId)
} catch {
// A failed history request is not evidence that Comfy lost the prompt.
continue
}
if (!history) continue
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.')
@@ -755,10 +762,9 @@ export async function retryStudioJob(owner: string, id: string) {
function pruneDone(jobs: StudioJob[]) {
const cutoff = Date.now() - 1000 * 60 * 60 * 24
return jobs.filter(job => {
if (job.status === 'waiting' || job.status === 'running' || job.status === 'held') return true
return job.updatedAt > cutoff
}).slice(-40)
const active = (job: StudioJob) => job.status === 'waiting' || job.status === 'running' || job.status === 'held'
const recentDone = new Set(jobs.filter(job => !active(job) && job.updatedAt > cutoff).slice(-40))
return jobs.filter(job => active(job) || recentDone.has(job))
}
export async function markStudioLive(owner: string, id: string, liveJobId: string, shotQueueId?: string) {
@@ -1599,12 +1605,12 @@ function findStudioRowForLive(jobs: StudioJob[], job: Job) {
const byQueue = jobs.find(item => (
item.shotQueueId === job.library?.queueId
&& item.status === 'running'
&& !item.liveJobId
))
if (byQueue) return byQueue
}
// Last resort: sole running row for this owner (liveJobId lost after a restart).
const running = jobs.filter(item => item.status === 'running')
if (running.length === 1) return running[0]
// A delayed callback from an older job must never settle a newer running job.
// Missing associations are handled by repairStaleJobs, not guessed here.
return undefined
}
@@ -1704,7 +1710,9 @@ export async function onLiveVideoSettled(job: Job) {
syncPausedFlag(store)
})
await kickStudioQueue()
// Settlement is also called from inside dispatchStudioQueue's recovery checks.
// Awaiting a dispatch queued behind that same dispatch deadlocks the entire queue.
void kickStudioQueue()
}
function interruptedRunning(job: Job, remaining: number, userPause: boolean, cutInHold: boolean) {
+168
View File
@@ -0,0 +1,168 @@
import test from 'node:test'
import assert from 'node:assert/strict'
import { createRequire } from 'node:module'
import { mkdtempSync, mkdirSync, readFileSync, writeFileSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { execFileSync } from 'node:child_process'
import ts from 'typescript'
const require = createRequire(import.meta.url)
const file = 'server/utils/studioQueue.ts'
const source = process.env.TEST_BASELINE === '1'
? execFileSync('git', ['show', `HEAD:${file}`], { encoding: 'utf8' })
: readFileSync(new URL(`../${file}`, import.meta.url), 'utf8')
const compiled = ts.transpileModule(source, {
compilerOptions: { module: ts.ModuleKind.CommonJS, target: ts.ScriptTarget.ES2022 }
}).outputText
// Execute the real queue module against temporary stores and a fake GPU.
// No production library, network, Comfy process, or render is touched.
function fixture({ rows = [], lives = [], queue = { running: 0, pending: 0 }, history = {} } = {}) {
const root = mkdtempSync(join(tmpdir(), 'aigen-queue-test-'))
const owner = 'test-owner'
mkdirSync(join(root, 'users', owner), { recursive: true })
const path = join(root, 'users', owner, 'studio-queue.json')
writeFileSync(path, JSON.stringify({ paused: false, jobs: rows }))
const timers = []
const started = []
const modules = {
'~/server/utils/jobs': {
listJobs: () => lives,
getJob: id => lives.find(job => job.id === id),
emitJob: () => {}
},
'~/server/utils/pending': {
listPendingJobs: () => [], readPendingJob: () => undefined,
patchPendingJob: () => {}, deletePendingJob: () => {}
},
'~/server/utils/shotQueue': {
getShotQueue: () => null, pendingSegmentCount: () => 0
},
'~/server/utils/comfy': { fetchLiveQueue: async () => queue, fetchHistory: async () => {
if (history instanceof Error) throw history
return history
} },
'~/utils/videoModels': {
isLtxWorkflow: () => false, isTextToVideo: () => true, isXaigenStudio: () => false,
ltxWorkflowEnabled: () => true, parseVideoWorkflow: value => value
},
'~/server/utils/loras': { persistLoraFields: () => ({}) },
'~/utils/loras': { resolveLoraStack: () => [] },
'~/utils/promptParts': { persistPromptWrappers: () => ({}) },
'~/utils/imageV2': { imageV2StackSpecials: () => ({}) },
'~/utils/globalLocks': { allowIdentityRefs: () => false },
'~/server/utils/generationLog': { updateGenerationLogByStudioJob: async () => {} },
'~/server/utils/musicChain': {
startMusicJob: async () => {
const live = liveJob(`started-${started.length}`, 'music')
lives.push(live)
started.push(live)
return live
}
}
}
const localRequire = id => {
if (id.startsWith('node:')) return require(id)
if (modules[id]) return modules[id]
throw new Error(`Unexpected test dependency: ${id}`)
}
const exports = {}
new Function('require', 'exports', 'useRuntimeConfig', 'setTimeout', 'createError', compiled)(
localRequire, exports, () => ({ libraryDir: root }),
fn => { timers.push(fn); return timers.length },
info => Object.assign(new Error(info.statusMessage), info)
)
return { api: exports, rows: () => JSON.parse(readFileSync(path, 'utf8')).jobs, lives, started, timers }
}
function row(id, status = 'running', liveJobId = id) {
return { id, ownerKey: 'test-owner', status, liveJobId, kind: 'music',
createdAt: Date.now(), updatedAt: Date.now(), payload: { prompt: 'test', extensions: [] } }
}
function liveJob(id, kind = 'edit', extra = {}) {
return { id, kind, status: 'running', promptId: `prompt-${id}`, startedAt: Date.now() - 300_000,
library: { ownerKey: 'test-owner' }, ...extra }
}
async function finishes(promise) {
let timer
try {
await Promise.race([promise, new Promise((_, reject) => {
timer = setTimeout(() => reject(new Error('Queue dispatcher deadlocked')), 500)
})])
} finally { clearTimeout(timer) }
}
for (const kind of ['edit', 'music']) {
test(`${kind}: recover saved output and dispatch the next job without deadlocking`, async () => {
const live = liveJob('old', kind, kind === 'edit' ? { stillId: 'saved' } : { audio: { filename: 'saved.wav' } })
const f = fixture({ rows: [row('old'), row('next', 'waiting', undefined)], lives: [live],
history: { 'prompt-old': { status: { status_str: 'success', completed: true } } } })
await finishes(f.api.kickStudioQueue())
assert.equal(f.rows().find(job => job.id === 'old').status, 'complete')
assert.equal(f.rows().find(job => job.id === 'next').status, 'running')
assert.equal(f.started.length, 1)
await finishes(f.api.kickStudioQueue())
assert.equal(f.started.length, 1)
})
}
test('late settlement from an older job does not complete the next running job', async () => {
const current = liveJob('current')
const f = fixture({ rows: [row('current')], lives: [current], queue: { running: 1, pending: 0 } })
await finishes(f.api.onLiveVideoSettled(liveJob('old', 'music', { status: 'complete' })))
assert.equal(f.rows()[0].status, 'running')
assert.equal(f.rows()[0].liveJobId, 'current')
})
test('history pruning preserves all waiting, held, and running jobs above 40 rows', async () => {
const rows = Array.from({ length: 50 }, (_, i) => row(`waiting-${i}`, 'waiting'))
rows.unshift(row('held', 'held'), row('active'))
const f = fixture({ rows, lives: [liveJob('active')], queue: { running: 1, pending: 0 } })
await finishes(f.api.kickStudioQueue())
assert.equal(f.rows().length, 52)
assert.equal(f.rows()[0].id, 'held')
})
test('an old shot callback cannot settle a new live job attached to the same batch', async () => {
const current = liveJob('current')
current.library.queueId = 'batch'
const currentRow = { ...row('current'), shotQueueId: 'batch' }
const f = fixture({ rows: [currentRow], lives: [current], queue: { running: 1, pending: 0 } })
const old = liveJob('old', 'video', { status: 'complete' })
old.library.queueId = 'batch'
await finishes(f.api.onLiveVideoSettled(old))
assert.equal(f.rows()[0].status, 'running')
assert.equal(f.rows()[0].liveJobId, 'current')
})
test('an empty Comfy queue during output saving does not release the active job', async () => {
const f = fixture({ rows: [row('saving'), row('next', 'waiting')], lives: [liveJob('saving', 'video', { saving: true })] })
await finishes(f.api.kickStudioQueue())
assert.equal(f.started.length, 0)
assert.equal(f.rows()[0].status, 'running')
assert.equal(f.timers.length, 1)
})
test('a missing prompt is failed and the next job starts', async () => {
const f = fixture({ rows: [row('lost'), row('next', 'waiting')], lives: [liveJob('lost')] })
await finishes(f.api.kickStudioQueue())
await finishes(f.api.kickStudioQueue())
assert.equal(f.rows()[0].status, 'error')
assert.equal(f.started.length, 1)
})
test('a connection outage does not fail an attached running job', async () => {
const f = fixture({ rows: [row('active'), row('next', 'waiting')], lives: [liveJob('active')], queue: null })
await finishes(f.api.kickStudioQueue())
assert.equal(f.rows()[0].status, 'running')
assert.equal(f.started.length, 0)
})
for (const history of [new Error('history timed out'), null]) {
test(`a failed history lookup (${history ? 'timeout' : 'HTTP failure'}) does not discard a render`, async () => {
const f = fixture({ rows: [row('active'), row('next', 'waiting')], lives: [liveJob('active')], history })
await finishes(f.api.kickStudioQueue())
assert.equal(f.rows()[0].status, 'running')
assert.equal(f.started.length, 0)
})
}