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/sharedGpu': { withSharedGpuStart: async fn => fn(), maintainSharedGpu: async () => {}, acquireSharedGpu: async () => false, sharedGpuWaitReason: () => 'GPU is in use. Waiting for availability.' }, '~/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 } } test('standalone YuEGP retains its queue slot while Comfy is empty for longer than the zombie timeout', async () => { const live = liveJob('yuegp-active', 'music', { yueGp: true, promptId: undefined, startedAt: Date.now() - 3600000 }) const f = fixture({ rows: [row('yuegp-active'), row('next', 'waiting', undefined)], lives: [live] }) await f.api.kickStudioQueue() assert.equal(live.status, 'running') assert.equal(f.started.length, 0) assert.equal(f.rows()[1].status, 'waiting') }) 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) }) }