Files
aigen/tests/studio-queue.test.mjs
T

214 lines
11 KiB
JavaScript

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 = {}, shotQueues = [] } = {}) {
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 = {
'./videoUpscale': { startUpscaleJob: async item => { started.push({ upscale: true, sourceId: item.payload.upscale.sourceId }) } },
'~/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: (_owner, id) => shotQueues.find(q => q.id === id), pendingSegmentCount: queue => queue?.pendingCount || 0,
isShotQueueCancelled: () => false, listShotQueues: () => shotQueues,
setShotQueuePause: async (_owner, id, pause) => { const q = shotQueues.find(q => q.id === id); if (q) { q.stopAfterCurrent = pause; q.autoRun = !pause } }, pauseShotQueue: async () => {}
},
'~/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/videoChain': { startQueueBurst: async (_owner, id, count) => { started.push({ id, count }); return { jobId: 'resumed-video' } } },
'~/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')
})
test('standalone upscale retains its queue slot while Comfy is empty for longer than the zombie timeout', async () => {
const live = liveJob('upscale-active', 'video', { upscale: true, promptId: undefined, startedAt: Date.now() - 3600000 })
const f = fixture({ rows: [row('upscale-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')
})
test('upscale queue row dispatches to the standalone worker before any video generation path', async () => {
const item = { ...row('upscale-waiting', 'waiting', undefined), kind: 'video', payload: { prompt: 'existing video', workflow: 'ltx-cherry', extensions: [], upscale: { sourceId: 'finished-master', scale: 2, fps: 'keep', target: 'preserve', enhance: 'faithful' } } }
const f = fixture({ rows: [item], queue: null })
await f.api.startStudioJob(item)
assert.deepEqual(f.started, [{ upscale: true, sourceId: 'finished-master' }])
})
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)
})
}
test('Resume a manually held video batch runs all remaining shots instead of one', async () => {
const item = { ...row('paused-video', 'held', undefined), kind: 'video', shotQueueId: 'shots', pausedByUser: true, payload: { prompt: 'video', queueAutoRun: false, extensions: [] } }
const shots = { id: 'shots', autoRun: false, stopAfterCurrent: true, pendingCount: 9, status: 'paused' }
const f = fixture({ rows: [item], shotQueues: [shots] })
await f.api.toggleStudioQueuePause('test-owner', { jobId: item.id, resume: true })
assert.deepEqual(f.started, [{ id: 'shots', count: 'all' }])
assert.equal(f.rows()[0].payload.queueAutoRun, true)
assert.equal(shots.autoRun, true)
})
test('explicit bounded burst clears an earlier run-all policy', async () => {
const item = { ...row('bounded', 'held', undefined), kind: 'video', shotQueueId: 'shots', resumeAutoRun: true, payload: { prompt: 'video', queueAutoRun: true, extensions: [] } }
const f = fixture({ rows: [item], shotQueues: [{ id: 'shots', autoRun: false, pendingCount: 9 }] })
await f.api.onStudioQueueBurstStarted('test-owner', 'shots', 'one-shot-live')
assert.equal(f.rows()[0].resumeAutoRun, false)
assert.equal(f.rows()[0].payload.queueAutoRun, false)
})