Files
aigen/server/utils/studioQueue.ts
T
TowstyandCursor be5211f488 Show generator settings in the library and restore them as input.
Collapse tiles so prompts stay behind the info button. Use as input opens the matching Video or Image v2 mode with the saved sliders.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 23:38:23 -05:00

1289 lines
44 KiB
TypeScript

import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFileSync } from 'node:fs'
import { join } from 'node:path'
import { getJob, listJobs, emitJob, type Job } from '~/server/utils/jobs'
import { listPendingJobs, patchPendingJob, readPendingJob, deletePendingJob } from '~/server/utils/pending'
import { getShotQueue } from '~/server/utils/shotQueue'
import { fetchLiveQueue } from '~/server/utils/comfy'
import { isLtxWorkflow, isTextToVideo, LTX_DISABLED_MESSAGE, ltxWorkflowEnabled, parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
import { persistLoraFields } from '~/server/utils/loras'
import { resolveLoraStack } from '~/utils/loras'
import { allowIdentityRefs } from '~/utils/globalLocks'
export type StudioJobStatus = 'waiting' | 'running' | 'held' | 'complete' | 'error' | 'cancelled'
export type StudioJobKind = 'video' | 'edit'
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<string | null>
extensions: { prompt: string; duration: number; permanenceRefs?: PermanenceRef[]; loraName?: string; loraStack?: import('~/utils/loras').LoraStackItem[] }[]
queueAutoRun: boolean
globalLocks?: string
permanenceRefs?: PermanenceRef[]
shotPermanenceRefs?: PermanenceRef[][]
loraName?: string
loraStack?: import('~/utils/loras').LoraStackItem[]
shotLoras?: string[]
shotLoraStacks?: import('~/utils/loras').LoraStackItem[][]
negative?: string
passes?: { prompt: string }[]
referenceStillId?: string
referenceStillFilename?: string
scaleToTotalPixels?: boolean
scaleMegapixels?: number
imagePipeline?: 'v1' | 'v2'
v2Mode?: 'edit' | 'compose' | 'refine' | 'generate'
v2Task?: 'scene' | 'identity' | 'outfit' | 'face_lock' | 'refine' | 't2i'
snofsModel?: number
snofsClip?: number
consistencyModel?: number
consistencyClip?: number
maskStillId?: string
maskStillFilename?: string
refineStrength?: number
}
export interface StudioJob {
id: string
ownerKey: string
createdAt: number
updatedAt: number
status: StudioJobStatus
kind: StudioJobKind
name: string
prompt: string
shotCount: number
familyId: string
shotQueueId?: string
liveJobId?: string
payload: StudioJobPayload
cutIn?: boolean
holdForCutIn?: boolean
pauseAfterCurrent?: boolean
pausedByUser?: boolean
resumeAutoRun?: boolean
lastError?: string
}
type StudioQueueStore = {
paused: boolean
jobs: StudioJob[]
}
const writeChains = new Map<string, Promise<unknown>>()
let dispatchChain: Promise<unknown> = 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 readStore(owner: string): StudioQueueStore {
ensureOwner(owner)
const path = queuePath(owner)
if (!existsSync(path)) return { paused: false, jobs: [] }
try {
const parsed = JSON.parse(readFileSync(path, 'utf8'))
if (Array.isArray(parsed)) return { paused: false, jobs: parsed }
if (parsed && typeof parsed === 'object') {
return {
paused: parsed.paused === true,
jobs: Array.isArray(parsed.jobs) ? parsed.jobs : []
}
}
return { paused: false, jobs: [] }
} catch {
return { paused: false, jobs: [] }
}
}
function writeStore(owner: string, store: StudioQueueStore) {
ensureOwner(owner)
const path = queuePath(owner)
const tmp = `${path}.tmp`
writeFileSync(tmp, JSON.stringify({
paused: store.paused === true,
jobs: store.jobs
}, null, 2))
renameSync(tmp, path)
}
function readJobs(owner: string): StudioJob[] {
return readStore(owner).jobs
}
function mutate<T>(owner: string, fn: (jobs: StudioJob[]) => T): Promise<T> {
const prev = writeChains.get(owner) || Promise.resolve()
const run = prev.then(() => {
const store = readStore(owner)
const result = fn(store.jobs)
writeStore(owner, store)
return result
})
writeChains.set(owner, run.then(() => undefined, () => undefined))
return run
}
function mutateStore<T>(owner: string, fn: (store: StudioQueueStore) => T): Promise<T> {
const prev = writeChains.get(owner) || Promise.resolve()
const run = prev.then(() => {
const store = readStore(owner)
const result = fn(store)
writeStore(owner, store)
return result
})
writeChains.set(owner, run.then(() => undefined, () => undefined))
return run
}
function jobIsUserPaused(job: StudioJob) {
return job.pauseAfterCurrent === true || job.pausedByUser === true
}
function syncPausedFlag(store: StudioQueueStore) {
store.paused = store.jobs.some(jobIsUserPaused)
}
function migrateGlobalPause(store: StudioQueueStore) {
const anyJobFlag = store.jobs.some(jobIsUserPaused)
if (store.paused && !anyJobFlag) {
for (const job of store.jobs) {
if (job.status === 'running') job.pauseAfterCurrent = true
else if (job.status === 'held' && !job.holdForCutIn) job.pausedByUser = true
}
}
syncPausedFlag(store)
}
export function isStudioQueuePaused(owner: string) {
const store = readStore(owner)
return store.jobs.some(jobIsUserPaused)
}
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 studioJobKind(job: Pick<StudioJob, 'kind'> | { kind?: string }) {
return job.kind === 'edit' ? 'edit' : 'video'
}
export function summarizeStudioJob(job: StudioJob) {
return {
id: job.id,
createdAt: job.createdAt,
updatedAt: job.updatedAt,
status: job.status,
kind: studioJobKind(job),
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,
imagePipeline: job.payload.imagePipeline || 'v1',
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,
pauseAfterCurrent: job.pauseAfterCurrent === true,
pausedByUser: job.pausedByUser === true,
lastError: job.lastError
}
}
const QUEUED_GRACE_MS = 4 * 60 * 1000
const SUBMIT_WINDOW_MS = 45 * 1000
const PENDING_ORPHAN_MS = 2 * 60 * 1000
function jobAgeMs(job: { startedAt: number }) {
return Date.now() - job.startedAt
}
/** Uploading, or created just now and not yet on Comfy. Paused leftovers do not count. */
export function jobIsLocallySubmitting(job: Job) {
if (job.library?.stopAfterCurrent) return false
if (job.status === 'uploading') return true
if (job.status === 'queued' && !job.promptId && jobAgeMs(job) < SUBMIT_WINDOW_MS) return true
if (job.status === 'running' && !job.promptId && jobAgeMs(job) < SUBMIT_WINDOW_MS) return true
return false
}
function sweepStaleLiveJobs() {
const now = Date.now()
for (const job of listJobs()) {
if (job.status !== 'queued') continue
if (job.promptId) continue
if (now - job.startedAt < QUEUED_GRACE_MS) continue
job.status = 'error'
job.error = job.error || 'Job never started'
}
}
export async function videoJobsBusy() {
sweepStaleLiveJobs()
if (listJobs().some(jobIsLocallySubmitting)) return true
const queue = await fetchLiveQueue()
if (queue) return queue.running > 0 || queue.pending > 0
return listJobs().some((job) => {
if (job.library?.stopAfterCurrent) return false
if (job.status === 'uploading' || job.status === 'running') return true
if (job.status === 'queued') return jobAgeMs(job) < SUBMIT_WINDOW_MS
return false
}) || listPendingJobs().some((pending) => {
if (!pending.promptId || pending.stopAfterCurrent) return false
const live = getJob(pending.jobId)
if (live) return jobIsLocallySubmitting(live) || live.status === 'running' || live.status === 'uploading'
return Date.now() - pending.startedAt < PENDING_ORPHAN_MS
})
}
export function applyPersistedPauseToJob(job: Job) {
const library = job.library
if (!library) return false
const store = readStore(library.ownerKey)
const row = store.jobs.find(item => item.liveJobId === job.id)
|| (library.queueId ? store.jobs.find(item => item.shotQueueId === library.queueId) : undefined)
const queue = library.queueId ? getShotQueue(library.ownerKey, library.queueId) : null
const paused = row?.pauseAfterCurrent === true
|| row?.pausedByUser === true
|| queue?.stopAfterCurrent === true
if (paused) {
library.stopAfterCurrent = true
library.queueAutoRun = false
}
return paused
}
export async function addStudioJob(params: {
ownerKey: string
payload: StudioJobPayload
familyId: string
kind?: StudioJobKind
}) {
const now = Date.now()
const kind = params.kind === 'edit' ? 'edit' : 'video'
const shotCount = kind === 'edit'
? 1 + (params.payload.passes?.length || 0)
: 1 + (params.payload.extensions?.length || 0)
const job: StudioJob = {
id: crypto.randomUUID(),
ownerKey: params.ownerKey,
createdAt: now,
updatedAt: now,
status: 'waiting',
kind,
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 before = readJobs(owner).find(item => item.id === id)
if (!before) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' })
if (before.status !== 'waiting' && before.status !== 'held' && before.status !== 'running') {
throw createError({ statusCode: 409, statusMessage: 'That job is no longer in the queue' })
}
const stopLive = before.status === 'running'
const liveJobId = before.liveJobId
const shotQueueId = before.shotQueueId
const result = await mutateStore(owner, (store) => {
const job = store.jobs.find(item => item.id === id)
if (!job) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' })
if (job.status !== 'waiting' && job.status !== 'held' && job.status !== 'running') {
throw createError({ statusCode: 409, statusMessage: 'That job is no longer in the queue' })
}
job.status = 'cancelled'
job.cutIn = false
job.holdForCutIn = false
job.pauseAfterCurrent = false
job.pausedByUser = false
job.liveJobId = undefined
job.updatedAt = Date.now()
const stillCutIn = store.jobs.some(item => item.status === 'waiting' && item.cutIn)
if (!stillCutIn) {
for (const item of store.jobs) {
if (item.status === 'running' || item.status === 'held') {
item.holdForCutIn = false
}
}
}
syncPausedFlag(store)
return structuredClone(job)
})
if (stopLive) {
await stopLiveGeneration(liveJobId, shotQueueId)
}
if (result.shotQueueId) {
const { deleteShotQueue, setQueueJob } = await import('~/server/utils/shotQueue')
setQueueJob(result.shotQueueId, null)
await deleteShotQueue(owner, result.shotQueueId, { force: true }).catch(() => null)
}
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
}
async function stopLiveGeneration(liveJobId?: string, shotQueueId?: string) {
if (shotQueueId) {
const { setQueueJob, pauseShotQueue } = await import('~/server/utils/shotQueue')
setQueueJob(shotQueueId, null)
const live = liveJobId ? getJob(liveJobId) : undefined
if (live?.library?.ownerKey) {
await pauseShotQueue(live.library.ownerKey, shotQueueId).catch(() => null)
}
}
if (!liveJobId) return
const live = getJob(liveJobId)
if (live) {
live.status = 'cancelled'
if (live.library) {
live.library.stopAfterCurrent = true
live.library.queueAutoRun = false
}
emitJob(live, { type: 'error', error: 'Job interrupted.', message: 'Job interrupted.' })
}
deletePendingJob(liveJobId)
const { interruptComfy } = await import('~/server/utils/comfy')
await interruptComfy().catch(() => false)
}
export async function cancelOrphanQueue(owner: string, queueId: string) {
const { getShotQueue, getQueueJobId, queueIsProcessing, deleteShotQueue } = await import('~/server/utils/shotQueue')
const queue = getShotQueue(owner, queueId)
if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
const generating = queue.status === 'running' || queueIsProcessing(queueId)
const liveJobId = getQueueJobId(queueId) || queue.currentJobId
if (generating) {
await stopLiveGeneration(liveJobId, queueId)
}
await deleteShotQueue(owner, queueId, { force: true }).catch(() => null)
kickStudioQueue()
return { ok: true }
}
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')
const target = jobsNow.find(job => job.id === id)
if (!target) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' })
if (target.status !== 'waiting') {
throw createError({ statusCode: 409, statusMessage: 'Only a waiting job can cut in' })
}
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 = running ? jobs.find(item => item.id === running.id) : undefined
if (active && active.status === 'running') {
active.holdForCutIn = true
active.resumeAutoRun = active.payload.queueAutoRun === true
active.updatedAt = Date.now()
}
return structuredClone(job)
})
if (running?.liveJobId) {
const live = getJob(running.liveJobId)
if (live?.library) {
live.library.stopAfterCurrent = true
}
}
await 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)
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[]) {
return jobs
.filter(job => (
job.status === 'held'
&& job.shotQueueId
&& job.holdForCutIn
&& !job.pausedByUser
))
.sort((a, b) => b.updatedAt - a.updatedAt)[0] || null
}
function pickWaiting(jobs: StudioJob[]) {
const waiting = jobs.filter(job => job.status === 'waiting')
if (!waiting.length) return null
if (jobs.some(job => job.pausedByUser === true)) {
return waiting.reduce((latest, job) => job.createdAt >= latest.createdAt ? job : latest)
}
return waiting[0] || null
}
export function kickStudioQueue() {
const run = dispatchChain.then(() => dispatchStudioQueue()).catch(() => undefined)
dispatchChain = run
return run
}
let kickRetryTimer: ReturnType<typeof setTimeout> | null = null
function scheduleKickRetry() {
if (kickRetryTimer) return
kickRetryTimer = setTimeout(() => {
kickRetryTimer = null
kickStudioQueue()
}, 2500)
}
function pendingAlive(job: StudioJob) {
if (job.liveJobId && readPendingJob(job.liveJobId)?.promptId) return true
if (!job.shotQueueId) return false
return listPendingJobs().some(pending => (
pending.ownerKey === job.ownerKey
&& pending.queueId === job.shotQueueId
&& Boolean(pending.promptId)
))
}
function repairStaleJobs(jobs: StudioJob[]) {
for (const job of jobs) {
if (job.status === 'error' && /allowIdentityRefs is not defined/i.test(String(job.lastError || ''))) {
job.status = 'waiting'
job.liveJobId = undefined
job.lastError = undefined
job.updatedAt = Date.now()
continue
}
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 || pendingAlive(job)) continue
const holdPause = job.pauseAfterCurrent === true || job.pausedByUser === true
if (job.shotQueueId) {
job.status = 'held'
job.liveJobId = undefined
if (holdPause) {
job.pausedByUser = true
job.pauseAfterCurrent = false
job.holdForCutIn = false
} else {
job.holdForCutIn = true
job.resumeAutoRun = job.payload.queueAutoRun === true || job.resumeAutoRun === true
}
job.updatedAt = Date.now()
} else {
job.status = holdPause ? 'held' : 'complete'
job.liveJobId = undefined
if (holdPause) {
job.pausedByUser = true
job.pauseAfterCurrent = false
}
job.updatedAt = Date.now()
}
}
}
async function dispatchStudioQueue() {
sweepStaleLiveJobs()
for (const owner of listOwnersWithStudioQueues()) {
const store = await mutateStore(owner, (current) => {
migrateGlobalPause(current)
repairStaleJobs(current.jobs)
const next = pruneDone(current.jobs)
current.jobs.splice(0, current.jobs.length, ...next)
syncPausedFlag(current)
return {
paused: current.paused === true,
jobs: current.jobs.map(item => structuredClone(item))
}
})
if (await videoJobsBusy()) {
if (listOwnersWithStudioQueues().some(owner => readJobs(owner).some(job => job.status === 'waiting'))) {
scheduleKickRetry()
}
return
}
const cutIn = pickCutIn(store.jobs)
if (cutIn) {
await startStudioJob(cutIn)
return
}
const waiting = pickWaiting(store.jobs)
if (waiting) {
await startStudioJob(waiting)
return
}
const held = pickHeld(store.jobs)
if (held?.shotQueueId) {
const resumed = await resumeHeldStudioJob(held)
if (resumed) return
}
}
}
async function resumeHeldStudioJob(item: StudioJob) {
if (!item.shotQueueId) {
await patchStudioJob(item.ownerKey, item.id, (job) => {
job.status = 'complete'
job.holdForCutIn = false
})
return false
}
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
})
return true
} 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
})
return false
}
}
async function startStudioEditJob(item: StudioJob) {
let live: Job | undefined
try {
const { stillPath } = await import('~/server/utils/library')
const { existsSync, readFileSync } = await import('node:fs')
const { createEditLiveJob, runEdit } = await import('~/server/utils/imageChain')
const { runEditV2 } = await import('~/server/utils/imageChainV2')
const payload = item.payload
const generate = payload.imagePipeline === 'v2' && payload.v2Mode === 'generate'
if (!generate && (!payload.stillId || !existsSync(stillPath(item.ownerKey, payload.stillId)))) {
throw new Error('The input still is missing from the library')
}
const image = 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 refId = payload.referenceStillId
const reference = refId && existsSync(stillPath(item.ownerKey, refId))
? {
filename: payload.referenceStillFilename || 'image2.png',
data: readFileSync(stillPath(item.ownerKey, refId)),
type: 'image/png'
}
: null
if (payload.imagePipeline === 'v2') {
const mode = payload.v2Mode === 'compose'
? 'compose'
: payload.v2Mode === 'refine'
? 'refine'
: payload.v2Mode === 'generate'
? 'generate'
: 'edit'
const maskId = payload.maskStillId
const mask = mode === 'refine' && maskId && existsSync(stillPath(item.ownerKey, maskId))
? {
filename: payload.maskStillFilename || 'refine-mask.png',
data: readFileSync(stillPath(item.ownerKey, maskId)),
type: 'image/png'
}
: null
if (mode === 'refine' && !mask) {
throw new Error('Refine requires a mask. Refusing to fall back to Edit.')
}
if (mode === 'compose' && !reference) {
throw new Error('Compose requires image B. Refusing to fall back to one-image generation.')
}
if (mode === 'edit' && reference) {
throw new Error('Edit mode takes one image. Use Compose for two stills.')
}
live = createEditLiveJob({
steps: payload.steps,
hideThumbnail: payload.hideThumbnail,
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 || 'auto',
width: payload.width,
height: payload.height,
steps: payload.steps,
turbo: payload.turbo === true,
seed: payload.seed || Math.floor(Math.random() * 2_147_483_647),
cfg: payload.cfg,
stillId: payload.stillId,
stillFilename: payload.stillFilename,
familyId: item.familyId,
chainIndex: 0,
chainStep: 1,
chainTotal: 1
}
})
await markStudioLive(item.ownerKey, item.id, live.id)
void runEditV2(live, {
mode,
task: payload.v2Task || 'scene',
image: mode === 'generate' ? null : image,
reference: mode === 'refine' || mode === 'generate' ? null : reference,
mask,
prompt: payload.prompt,
negative: payload.negative || '',
steps: payload.steps,
seed: live.library?.seed || payload.seed,
cfg: payload.cfg,
snofsModel: payload.snofsModel ?? 0.65,
snofsClip: payload.snofsClip ?? 0.35,
consistencyModel: payload.consistencyModel ?? (mode === 'generate' ? 0 : 0.7),
consistencyClip: payload.consistencyClip ?? (mode === 'generate' ? 0 : 0.7),
megapixels: payload.scaleMegapixels ?? 1,
strength: payload.refineStrength,
width: payload.width,
height: payload.height,
turbo: payload.turbo === true,
sourceStillId: payload.stillId,
referenceStillId: payload.referenceStillId
}).catch((error) => {
const message = error instanceof Error ? error.message : String(error)
if (live && live.status !== 'error' && live.status !== 'cancelled' && live.status !== 'deferred') {
live.status = 'error'
live.error = message
}
})
return
}
const passes = payload.passes || []
live = createEditLiveJob({
steps: payload.steps,
hideThumbnail: payload.hideThumbnail,
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 || 'auto',
width: payload.width,
height: payload.height,
steps: payload.steps,
turbo: true,
seed: payload.seed || Math.floor(Math.random() * 2_147_483_647),
cfg: payload.cfg,
stillId: payload.stillId,
stillFilename: payload.stillFilename,
familyId: item.familyId,
chainIndex: 0,
chainStep: 1,
chainTotal: 1 + passes.length,
chainLabel: passes.length ? 'Pass 1' : undefined,
passes,
...persistLoraFields(payload.loraStack || payload.loraName)
}
})
await markStudioLive(item.ownerKey, item.id, live.id)
void runEdit(live, {
image,
reference,
prompt: payload.prompt,
passes,
negative: payload.negative || '',
steps: payload.steps,
seed: live.library?.seed || payload.seed,
cfg: payload.cfg,
scaleToTotalPixels: payload.scaleToTotalPixels === true,
scaleMegapixels: payload.scaleMegapixels,
...persistLoraFields(payload.loraStack || payload.loraName)
}).catch((error) => {
const message = error instanceof Error ? error.message : String(error)
if (live && live.status !== 'error' && live.status !== 'cancelled' && live.status !== 'deferred') {
live.status = 'error'
live.error = message
}
})
} catch (error) {
const message = error instanceof Error ? error.message : String(error)
if (live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')) {
live.status = 'error'
live.error = message
}
await patchStudioJob(item.ownerKey, item.id, (job) => {
job.status = 'error'
job.lastError = message
}).catch(() => null)
kickStudioQueue()
}
}
export async function startStudioJob(item: StudioJob) {
if (studioJobKind(item) === 'edit') {
await startStudioEditJob(item)
return
}
let live: Job | undefined
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 workflow = parseVideoWorkflow(payload.workflow)
if (isLtxWorkflow(workflow) && !ltxWorkflowEnabled()) {
throw new Error(LTX_DISABLED_MESSAGE)
}
const identity = allowIdentityRefs(
payload.useIdentityRefs,
payload.permanenceRefs,
payload.globalLocks,
payload.shotPermanenceRefs
)
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
if (!isTextToVideo(workflow) && !stillFile) {
throw new Error('The input still is missing from the library')
}
const job = createJob()
live = job
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: workflow,
useIdentityRefs: identity,
queueAutoRun: extensions.length ? payload.queueAutoRun : false,
queueBudget: extensions.length && payload.queueAutoRun ? extensions.length : 0,
globalLocks: payload.globalLocks,
permanenceRefs: payload.permanenceRefs,
shotPermanenceRefs: payload.shotPermanenceRefs,
...persistLoraFields(resolveLoraStack(payload.loraStack || payload.loraName, payload.shotLoraStacks?.[0] || payload.shotLoras?.[0])),
shotLoras: payload.shotLoras,
shotLoraStacks: payload.shotLoraStacks
}
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: workflow,
sound: payload.sound,
useIdentityRefs: identity,
referenceStillIds: payload.referenceStillIds,
globalLocks: payload.globalLocks,
permanenceRefs: payload.permanenceRefs,
...persistLoraFields(payload.loraStack || payload.loraName),
initial: { prompt: payload.prompt, duration: payload.duration, permanenceRefs: payload.shotPermanenceRefs?.[0], ...persistLoraFields(resolveLoraStack(payload.loraStack || payload.loraName, payload.shotLoraStacks?.[0] || payload.shotLoras?.[0])) },
extensions: extensions.map((item, index) => ({
...item,
...persistLoraFields(item.loraStack || item.loraName || payload.shotLoraStacks?.[index + 1] || payload.shotLoras?.[index + 1])
})),
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 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: workflow,
duration: payload.duration,
useIdentityRefs: identity,
referenceImages,
...persistLoraFields(resolveLoraStack(payload.loraStack || payload.loraName, payload.shotLoraStacks?.[0] || payload.shotLoras?.[0])),
shotLoras: payload.shotLoras,
shotLoraStacks: payload.shotLoraStacks
}).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)
if (live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')) {
live.status = 'error'
live.error = message
}
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 failed = job.status === 'error' || job.status === 'cancelled'
const snapshot = readStore(owner)
const rowNow = snapshot.jobs.find(item => item.liveJobId === job.id)
|| (job.library?.queueId
? snapshot.jobs.find(item => (
item.status === 'running'
&& item.shotQueueId === job.library?.queueId
))
: undefined)
const userPause = !failed && rowNow?.holdForCutIn !== true && (
rowNow?.pauseAfterCurrent === true
|| rowNow?.pausedByUser === true
|| job.library?.stopAfterCurrent === true
)
const cutInHold = !failed && remaining > 0 && rowNow?.holdForCutIn === true && !userPause
const interrupted = remaining > 0 && (userPause || cutInHold)
if (interrupted && job.status === 'running') {
job.status = 'complete'
}
if (userPause && remaining > 0) {
await mutateStore(owner, (store) => {
const row = store.jobs.find(item => item.liveJobId === job.id)
if (!row) {
syncPausedFlag(store)
return
}
row.status = 'held'
row.liveJobId = undefined
row.pauseAfterCurrent = false
row.pausedByUser = true
row.holdForCutIn = false
row.resumeAutoRun = row.payload.queueAutoRun === true || row.resumeAutoRun === true
row.updatedAt = Date.now()
syncPausedFlag(store)
})
} else if (cutInHold) {
await mutate(owner, (jobs) => {
const row = jobs.find(item => item.liveJobId === job.id)
if (!row) return
row.status = 'held'
row.liveJobId = undefined
row.holdForCutIn = true
row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
row.updatedAt = Date.now()
})
} else if (userPause) {
await mutateStore(owner, (store) => {
const row = store.jobs.find(item => item.liveJobId === job.id)
if (!row) {
syncPausedFlag(store)
return
}
row.status = job.status === 'cancelled'
? 'cancelled'
: job.status === 'complete'
? 'complete'
: 'error'
row.liveJobId = undefined
row.cutIn = false
row.holdForCutIn = false
row.pauseAfterCurrent = false
row.pausedByUser = false
row.lastError = job.error
row.updatedAt = Date.now()
syncPausedFlag(store)
})
} else {
const status = job.status === 'cancelled'
? 'cancelled'
: job.status === 'complete'
? 'complete'
: 'error'
await markStudioSettled(owner, job.id, status, job.error)
}
await kickStudioQueue()
}
function applyLiveJobPause(live: Job | undefined, pause: boolean, restoreAutoRun?: boolean) {
if (!live?.library) return
live.library.stopAfterCurrent = pause
if (pause) live.library.queueAutoRun = false
else if (restoreAutoRun) live.library.queueAutoRun = true
patchPendingJob(live.id, {
stopAfterCurrent: pause,
queueAutoRun: live.library.queueAutoRun
})
}
export async function toggleStudioQueuePause(owner: string, options: {
resume?: boolean
jobId?: string
queueId?: string
} = {}) {
const { pauseShotQueue, setShotQueuePause, listShotQueues, getShotQueue } = await import('~/server/utils/shotQueue')
const storeNow = readStore(owner)
const running = storeNow.jobs.find(job => job.status === 'running')
const target = options.jobId
? storeNow.jobs.find(job => job.id === options.jobId)
: options.queueId
? storeNow.jobs.find(job => job.shotQueueId === options.queueId)
: running || storeNow.jobs.find(job => job.pauseAfterCurrent || job.pausedByUser)
const queueArmed = options.queueId
? listShotQueues(owner).some(queue => queue.id === options.queueId && queue.stopAfterCurrent === true)
: target?.shotQueueId
? listShotQueues(owner).some(queue => queue.id === target.shotQueueId && queue.stopAfterCurrent === true)
: false
const armed = target?.pauseAfterCurrent === true
|| target?.pausedByUser === true
|| queueArmed
const shouldPause = options.resume === true ? false : options.resume === false ? true : !armed
if (!shouldPause) {
const heldPaused = storeNow.jobs.filter(job => (
job.status === 'held'
&& Boolean(job.shotQueueId)
&& job.pausedByUser === true
&& (!options.jobId || job.id === options.jobId)
&& (!options.queueId || job.shotQueueId === options.queueId)
))
const restored = await mutateStore(owner, (store) => {
for (const job of store.jobs) {
if (options.jobId && job.id !== options.jobId) continue
if (options.queueId && job.shotQueueId !== options.queueId) continue
job.pauseAfterCurrent = false
job.pausedByUser = false
}
syncPausedFlag(store)
return structuredClone(store)
})
const rows = options.jobId
? restored.jobs.filter(job => job.id === options.jobId)
: options.queueId
? restored.jobs.filter(job => job.shotQueueId === options.queueId)
: restored.jobs
for (const row of rows) {
if (row.liveJobId) applyLiveJobPause(getJob(row.liveJobId), false, row.payload.queueAutoRun === true || row.resumeAutoRun === true)
if (row.shotQueueId) await setShotQueuePause(owner, row.shotQueueId, false).catch(() => null)
}
if (!options.jobId && !options.queueId) {
for (const queue of listShotQueues(owner)) {
if (queue.stopAfterCurrent) await setShotQueuePause(owner, queue.id, false).catch(() => null)
}
} else if (options.queueId) {
const rowHasQueue = rows.some(job => job.shotQueueId === options.queueId)
if (!rowHasQueue) await setShotQueuePause(owner, options.queueId, false).catch(() => null)
}
const busy = await videoJobsBusy()
if (heldPaused[0] && !busy) {
await resumeHeldStudioJob(heldPaused[0]).catch(() => null)
} else {
if (heldPaused.length && busy) {
await mutate(owner, (jobs) => {
for (const job of jobs) {
if (!heldPaused.some(item => item.id === job.id)) continue
job.holdForCutIn = true
job.pausedByUser = false
job.pauseAfterCurrent = false
job.updatedAt = Date.now()
}
}).catch(() => null)
}
kickStudioQueue()
}
return { paused: false, willPause: false }
}
const updated = await mutateStore(owner, (store) => {
const row = options.jobId
? store.jobs.find(job => job.id === options.jobId)
: options.queueId
? store.jobs.find(job => job.shotQueueId === options.queueId)
: store.jobs.find(job => job.status === 'running')
if (row) {
if (row.status !== 'running' && row.status !== 'held') {
throw createError({ statusCode: 409, statusMessage: 'Nothing is generating on that job' })
}
if (row.status === 'running') {
row.pauseAfterCurrent = true
} else {
row.pausedByUser = true
row.pauseAfterCurrent = false
row.holdForCutIn = false
}
row.resumeAutoRun = row.payload.queueAutoRun === true || row.resumeAutoRun === true
row.updatedAt = Date.now()
}
syncPausedFlag(store)
return {
paused: store.paused,
job: row ? structuredClone(row) : null
}
})
const queueId = updated.job?.shotQueueId || options.queueId
if (updated.job?.liveJobId) {
applyLiveJobPause(getJob(updated.job.liveJobId), true)
patchPendingJob(updated.job.liveJobId, { stopAfterCurrent: true, queueAutoRun: false })
}
for (const pending of listPendingJobs()) {
if (pending.ownerKey !== owner) continue
if (queueId) {
if (pending.queueId !== queueId) continue
} else if (updated.job?.liveJobId) {
if (pending.jobId !== updated.job.liveJobId) continue
} else {
continue
}
patchPendingJob(pending.jobId, { stopAfterCurrent: true, queueAutoRun: false })
}
const queued = queueId ? getShotQueue(owner, queueId) : null
if (updated.job && queued?.autoRun) {
await patchStudioJob(owner, updated.job.id, (job) => {
job.resumeAutoRun = true
}).catch(() => null)
}
if (queueId) await pauseShotQueue(owner, queueId).catch(() => null)
if (!updated.job && !queueId) {
for (const queue of listShotQueues(owner)) {
if (queue.status === 'running') await pauseShotQueue(owner, queue.id).catch(() => null)
}
}
const anyGenerating = storeNow.jobs.some(job => job.status === 'running')
|| listShotQueues(owner).some(queue => queue.status === 'running')
if (!updated.job && !queueId && !anyGenerating) {
throw createError({ statusCode: 409, statusMessage: 'Nothing is generating, so there is nothing to pause after' })
}
return {
paused: true,
willPause: updated.job?.status === 'running' || Boolean(queued?.status === 'running') || anyGenerating,
jobId: updated.job?.id,
queueId
}
}
export async function onStudioQueueBurstStarted(owner: string, queueId: string, liveJobId: string) {
await mutateStore(owner, (store) => {
const row = store.jobs.find(job => job.shotQueueId === queueId)
if (!row) {
syncPausedFlag(store)
return
}
row.status = 'running'
row.liveJobId = liveJobId
row.pauseAfterCurrent = false
row.pausedByUser = false
row.holdForCutIn = false
row.updatedAt = Date.now()
syncPausedFlag(store)
})
}