diff --git a/package.json b/package.json index a623fc3..9c994a3 100644 --- a/package.json +++ b/package.json @@ -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", diff --git a/server/utils/comfy.ts b/server/utils/comfy.ts index d95c2bf..f66f6f0 100644 --- a/server/utils/comfy.ts +++ b/server/utils/comfy.ts @@ -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 } diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts index f73de79..059e13f 100644 --- a/server/utils/studioQueue.ts +++ b/server/utils/studioQueue.ts @@ -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> + 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) { diff --git a/tests/studio-queue.test.mjs b/tests/studio-queue.test.mjs new file mode 100644 index 0000000..b93e2cb --- /dev/null +++ b/tests/studio-queue.test.mjs @@ -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) + }) +}