import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFileSync } from 'node:fs' import { join } from 'node:path' import { getJob, listJobs, type Job } from '~/server/utils/jobs' import { parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels' export type StudioJobStatus = 'waiting' | 'running' | 'held' | 'complete' | 'error' | 'cancelled' export interface StudioJobPayload { prompt: string name: string folderId: string aspect: string width: number height: number steps: number turbo: boolean seed: number cfg: number fps: number samplerName: string scheduler: string duration: number sound: boolean workflow: VideoWorkflowId useIdentityRefs: boolean stillId?: string stillFilename?: string hideThumbnail: boolean hideInput?: boolean folderLocked?: boolean referenceStillIds: Array extensions: { prompt: string; duration: number }[] queueAutoRun: boolean } export interface StudioJob { id: string ownerKey: string createdAt: number updatedAt: number status: StudioJobStatus name: string prompt: string shotCount: number familyId: string shotQueueId?: string liveJobId?: string payload: StudioJobPayload cutIn?: boolean holdForCutIn?: boolean resumeAutoRun?: boolean lastError?: string } const writeChains = new Map>() let dispatchChain: Promise = Promise.resolve() function libraryRoot() { const config = useRuntimeConfig() return (config.libraryDir || process.env.LIBRARY_DIR || '/data/library').replace(/\/$/, '') } function queuePath(owner: string) { return join(libraryRoot(), 'users', owner, 'studio-queue.json') } function ensureOwner(owner: string) { mkdirSync(join(libraryRoot(), 'users', owner), { recursive: true }) } function readJobs(owner: string): StudioJob[] { ensureOwner(owner) const path = queuePath(owner) if (!existsSync(path)) return [] try { const parsed = JSON.parse(readFileSync(path, 'utf8')) return Array.isArray(parsed) ? parsed : [] } catch { return [] } } function writeJobs(owner: string, jobs: StudioJob[]) { ensureOwner(owner) const path = queuePath(owner) const tmp = `${path}.tmp` writeFileSync(tmp, JSON.stringify(jobs, null, 2)) renameSync(tmp, path) } function mutate(owner: string, fn: (jobs: StudioJob[]) => T): Promise { const prev = writeChains.get(owner) || Promise.resolve() const run = prev.then(() => { const jobs = readJobs(owner) const result = fn(jobs) writeJobs(owner, jobs) return result }) writeChains.set(owner, run.then(() => undefined, () => undefined)) return run } export function listOwnersWithStudioQueues() { const root = join(libraryRoot(), 'users') if (!existsSync(root)) return [] as string[] return readdirSync(root, { withFileTypes: true }) .filter(entry => entry.isDirectory()) .map(entry => entry.name) .filter(owner => existsSync(queuePath(owner))) } export function listStudioJobs(owner: string) { return readJobs(owner) } export function summarizeStudioJob(job: StudioJob) { return { id: job.id, createdAt: job.createdAt, updatedAt: job.updatedAt, status: job.status, name: job.name, prompt: job.prompt, shotCount: job.shotCount, familyId: job.familyId, shotQueueId: job.shotQueueId, liveJobId: job.liveJobId, stillId: job.payload.stillId, workflow: job.payload.workflow, duration: job.payload.duration, hideThumbnail: job.payload.hideThumbnail, folderLocked: job.payload.folderLocked === true, queueAutoRun: job.payload.queueAutoRun === true, cutIn: job.cutIn === true, holdForCutIn: job.holdForCutIn === true, lastError: job.lastError } } export function videoJobsBusy() { return listJobs().some(job => job.kind !== 'edit' && (job.status === 'queued' || job.status === 'uploading' || job.status === 'running')) } export async function addStudioJob(params: { ownerKey: string payload: StudioJobPayload familyId: string }) { const now = Date.now() const shotCount = 1 + (params.payload.extensions?.length || 0) const job: StudioJob = { id: crypto.randomUUID(), ownerKey: params.ownerKey, createdAt: now, updatedAt: now, status: 'waiting', name: (params.payload.name || '').trim() || params.payload.prompt.slice(0, 80), prompt: params.payload.prompt, shotCount, familyId: params.familyId, payload: params.payload, cutIn: false } await mutate(params.ownerKey, (jobs) => { jobs.push(job) return job }) return job } export async function patchStudioJob(owner: string, id: string, patch: (job: StudioJob) => void) { return mutate(owner, (jobs) => { const job = jobs.find(item => item.id === id) if (!job) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' }) patch(job) job.updatedAt = Date.now() return structuredClone(job) }) } export async function cancelStudioJob(owner: string, id: string) { const result = await mutate(owner, (jobs) => { const job = jobs.find(item => item.id === id) if (!job) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' }) if (job.status === 'running') { throw createError({ statusCode: 409, statusMessage: 'Stop the live job from Output if you want to cancel the one that is generating' }) } if (job.status !== 'waiting' && job.status !== 'held') { throw createError({ statusCode: 409, statusMessage: 'Only a waiting or paused job can be removed from the queue' }) } job.status = 'cancelled' job.cutIn = false job.holdForCutIn = false job.updatedAt = Date.now() const stillCutIn = jobs.some(item => item.status === 'waiting' && item.cutIn) if (!stillCutIn) { for (const item of jobs) { if (item.status === 'running' || item.status === 'held') { item.holdForCutIn = false } } } return structuredClone(job) }) const jobs = readJobs(owner) if (!jobs.some(item => item.status === 'waiting' && item.cutIn)) { const running = jobs.find(item => item.status === 'running') if (running?.liveJobId) { const live = getJob(running.liveJobId) if (live?.library) live.library.stopAfterCurrent = false } } kickStudioQueue() return result } export async function requestCutIn(owner: string, id: string) { const jobsNow = readJobs(owner) const running = jobsNow.find(job => job.status === 'running') || jobsNow.find(job => job.status === 'held') if (!running) { throw createError({ statusCode: 409, statusMessage: 'Nothing is generating, so this job can just wait its turn' }) } if (running.id === id) { throw createError({ statusCode: 409, statusMessage: 'That job is already running' }) } const updated = await mutate(owner, (jobs) => { const job = jobs.find(item => item.id === id) if (!job) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' }) if (job.status !== 'waiting') { throw createError({ statusCode: 409, statusMessage: 'Only a waiting job can cut in' }) } job.cutIn = true job.updatedAt = Date.now() for (const item of jobs) { if (item.id !== id && item.status === 'waiting') item.cutIn = false } const active = jobs.find(item => item.id === running.id) if (active && active.status === 'running') { active.holdForCutIn = true active.resumeAutoRun = active.payload.queueAutoRun === true active.updatedAt = Date.now() } return structuredClone(job) }) const liveId = running.liveJobId if (liveId) { const live = getJob(liveId) if (live?.library) { live.library.stopAfterCurrent = true } } kickStudioQueue() return updated } 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) } export async function markStudioLive(owner: string, id: string, liveJobId: string, shotQueueId?: string) { return patchStudioJob(owner, id, (job) => { job.status = 'running' job.liveJobId = liveJobId if (shotQueueId) job.shotQueueId = shotQueueId job.lastError = undefined }) } export async function markStudioHeld(owner: string, id: string) { return patchStudioJob(owner, id, (job) => { job.status = 'held' job.liveJobId = undefined }) } export async function markStudioSettled(owner: string, liveJobId: string, status: 'complete' | 'error' | 'cancelled', error?: string) { return mutate(owner, (jobs) => { const job = jobs.find(item => item.liveJobId === liveJobId) || jobs.find(item => item.status === 'running') if (!job) return null if (job.holdForCutIn && status === 'complete') { job.status = 'held' job.liveJobId = undefined job.updatedAt = Date.now() return structuredClone(job) } job.status = status job.liveJobId = undefined job.cutIn = false job.holdForCutIn = false job.lastError = error job.updatedAt = Date.now() return structuredClone(job) }) } function pickCutIn(jobs: StudioJob[]) { return jobs.find(job => job.status === 'waiting' && job.cutIn) || null } function pickHeld(jobs: StudioJob[]) { const held = jobs .filter(job => job.status === 'held' && job.shotQueueId) .sort((a, b) => b.updatedAt - a.updatedAt) const cutInHold = held.find(job => job.holdForCutIn) return cutInHold || held[0] || null } function pickWaiting(jobs: StudioJob[]) { return jobs.find(job => job.status === 'waiting') || null } export function kickStudioQueue() { const run = dispatchChain.then(() => dispatchStudioQueue()).catch(() => undefined) dispatchChain = run return run } function repairStaleJobs(jobs: StudioJob[]) { for (const job of jobs) { if (job.status !== 'running') continue const live = job.liveJobId ? getJob(job.liveJobId) : undefined const liveBusy = live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running') if (liveBusy) continue if (job.shotQueueId) { job.status = 'held' job.liveJobId = undefined job.holdForCutIn = true job.resumeAutoRun = job.payload.queueAutoRun === true || job.resumeAutoRun === true job.updatedAt = Date.now() } else { job.status = 'complete' job.liveJobId = undefined job.updatedAt = Date.now() } } } async function dispatchStudioQueue() { if (videoJobsBusy()) return for (const owner of listOwnersWithStudioQueues()) { const jobs = await mutate(owner, (list) => { repairStaleJobs(list) const next = pruneDone(list) list.splice(0, list.length, ...next) return list.map(item => structuredClone(item)) }) const cutIn = pickCutIn(jobs) if (cutIn) { await startStudioJob(cutIn) return } const held = pickHeld(jobs) if (held?.shotQueueId) { await resumeHeldStudioJob(held) return } const waiting = pickWaiting(jobs) if (waiting) { await startStudioJob(waiting) return } } } async function resumeHeldStudioJob(item: StudioJob) { if (!item.shotQueueId) { await patchStudioJob(item.ownerKey, item.id, (job) => { job.status = 'complete' job.holdForCutIn = false }) return } const { startQueueBurst } = await import('~/server/utils/videoChain') const count = item.resumeAutoRun ? 'all' as const : 1 try { const started = await startQueueBurst(item.ownerKey, item.shotQueueId, count) await patchStudioJob(item.ownerKey, item.id, (job) => { job.status = 'running' job.liveJobId = started.jobId job.holdForCutIn = false job.cutIn = false }) } catch (error) { const message = error instanceof Error ? error.message : String(error) await patchStudioJob(item.ownerKey, item.id, (job) => { job.status = 'error' job.lastError = message job.holdForCutIn = false }) } } export async function startStudioJob(item: StudioJob) { try { const { createJob } = await import('~/server/utils/jobs') const { runGeneration, frameLength } = await import('~/server/utils/videoChain') const { createShotQueue, setQueueJob } = await import('~/server/utils/shotQueue') const { stillPath } = await import('~/server/utils/library') const { existsSync, readFileSync } = await import('node:fs') const payload = item.payload const job = createJob() job.kind = 'video' job.maxStep = payload.steps job.hideThumbnail = payload.hideThumbnail const extensions = payload.extensions || [] job.library = { ownerKey: item.ownerKey, folderId: payload.folderId, hideThumbnail: payload.hideThumbnail, hideInput: payload.hideInput, folderLocked: payload.folderLocked, name: payload.name, prompt: payload.prompt, aspect: payload.aspect, width: payload.width, height: payload.height, steps: payload.steps, turbo: payload.turbo, seed: payload.seed || Math.floor(Math.random() * 2_147_483_647), cfg: payload.cfg, fps: payload.fps, samplerName: payload.samplerName, scheduler: payload.scheduler, stillId: payload.stillId, stillFilename: payload.stillFilename, referenceStillIds: payload.referenceStillIds, duration: payload.duration, sound: payload.sound, extensions, chainIndex: 0, chainStep: 1, chainTotal: 1 + extensions.length, chainLabel: extensions.length ? 'Initial' : undefined, familyId: item.familyId, workflow: parseVideoWorkflow(payload.workflow), useIdentityRefs: payload.useIdentityRefs, queueAutoRun: extensions.length ? payload.queueAutoRun : false, queueBudget: extensions.length && payload.queueAutoRun ? extensions.length : 0 } if (extensions.length) { const queue = await createShotQueue({ ownerKey: item.ownerKey, name: payload.name, familyId: item.familyId, folderId: payload.folderId, autoRun: payload.queueAutoRun, stillId: payload.stillId, stillFilename: payload.stillFilename, hideThumbnail: payload.hideThumbnail, hideInput: payload.hideInput, aspect: payload.aspect, width: payload.width, height: payload.height, steps: payload.steps, turbo: payload.turbo, cfg: payload.cfg, fps: payload.fps, samplerName: payload.samplerName, scheduler: payload.scheduler, workflow: parseVideoWorkflow(payload.workflow), sound: payload.sound, useIdentityRefs: payload.useIdentityRefs, referenceStillIds: payload.referenceStillIds, initial: { prompt: payload.prompt, duration: payload.duration }, extensions, jobId: job.id }) job.library.queueId = queue.id setQueueJob(queue.id, job.id) await patchStudioJob(item.ownerKey, item.id, (row) => { row.shotQueueId = queue.id }) } await markStudioLive(item.ownerKey, item.id, job.id, job.library.queueId) const stillFile = payload.stillId && existsSync(stillPath(item.ownerKey, payload.stillId)) ? { filename: payload.stillFilename || 'still.png', data: readFileSync(stillPath(item.ownerKey, payload.stillId)), type: 'image/png' } : null const referenceImages: Array<{ filename: string; data: Buffer; type?: string } | null> = [null, null, null, null] for (const [index, stillId] of (payload.referenceStillIds || []).entries()) { if (!stillId || index >= 4) continue const path = stillPath(item.ownerKey, stillId) if (!existsSync(path)) continue referenceImages[index] = { filename: `identity-ref-${index + 2}.png`, data: readFileSync(path), type: 'image/png' } } void runGeneration(job, { prompt: payload.prompt, image: stillFile, width: payload.width, height: payload.height, steps: payload.steps, seed: job.library.seed, turbo: payload.turbo, length: frameLength(payload.duration, payload.fps), sound: payload.sound, cfg: payload.cfg, fps: payload.fps, samplerName: payload.samplerName, scheduler: payload.scheduler, extensions, workflow: parseVideoWorkflow(payload.workflow), duration: payload.duration, useIdentityRefs: payload.useIdentityRefs, referenceImages }).catch((error) => { const message = error instanceof Error ? error.message : String(error) if (job.status !== 'error' && job.status !== 'cancelled' && job.status !== 'deferred') { job.status = 'error' job.error = message } }) } catch (error) { const message = error instanceof Error ? error.message : String(error) await patchStudioJob(item.ownerKey, item.id, (job) => { job.status = 'error' job.lastError = message }).catch(() => null) kickStudioQueue() } } export async function onLiveVideoSettled(job: Job) { const owner = job.library?.ownerKey if (!owner) { kickStudioQueue() return } const remaining = job.library ? Math.max(0, (job.library.chainTotal || 1) - (job.library.chainStep || 1)) : 0 const interrupted = job.library?.stopAfterCurrent === true && remaining > 0 && job.status !== 'error' && job.status !== 'cancelled' if (interrupted && job.status === 'running') { job.status = 'complete' } const hold = interrupted if (hold) { await mutate(owner, (jobs) => { const row = jobs.find(item => item.liveJobId === job.id) || jobs.find(item => item.status === 'running') if (!row) return row.status = 'held' row.liveJobId = undefined row.holdForCutIn = true row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true row.updatedAt = Date.now() }) } else { const status = job.status === 'cancelled' ? 'cancelled' : job.status === 'complete' ? 'complete' : 'error' await markStudioSettled(owner, job.id, status, job.error) } kickStudioQueue() }