Files
aigen/server/utils/studioQueue.ts
T

1871 lines
67 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, pendingSegmentCount } from '~/server/utils/shotQueue'
import { fetchLiveQueue } from '~/server/utils/comfy'
import { isLtxWorkflow, isTextToVideo, isXaigenStudio, LTX_DISABLED_MESSAGE, ltxWorkflowEnabled, parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
import { persistLoraFields } from '~/server/utils/loras'
import { resolveLoraStack } from '~/utils/loras'
import { persistPromptWrappers } from '~/utils/promptParts'
import { imageV2StackSpecials } from '~/utils/imageV2'
import { allowIdentityRefs, type PermanenceRef } from '~/utils/globalLocks'
export type StudioJobStatus = 'waiting' | 'running' | 'held' | 'complete' | 'error' | 'cancelled'
export type StudioJobKind = 'video' | 'edit' | 'music'
export interface StudioJobPayload {
prompt: string
promptMid?: string
promptPre?: string
promptPost?: 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: import('~/server/utils/library').QueuedExtension[]
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 }[]
passMode?: 'batch' | 'chain'
referenceStillId?: string
referenceStillFilename?: string
scaleToTotalPixels?: boolean
scaleMegapixels?: number
imagePipeline?: 'v1' | 'v2'
v2Mode?: 'edit' | 'compose' | 'refine' | 'generate' | 'iterate'
v2Task?: 'scene' | 'identity' | 'outfit' | 'face_lock' | 'refine' | 't2i'
engine?: 'flux' | 'krea'
snofsModel?: number
snofsClip?: number
consistencyModel?: number
consistencyClip?: number
maskStillId?: string
maskStillFilename?: string
refineStrength?: number
extendFromClipId?: string
refineExtensionFrame?: boolean
saveLosslessAnchor?: boolean
refinementDenoise?: number
lyrics?: string
instrumental?: boolean
lyricsStrength?: number
musicEngine?: import('~/utils/music').MusicEngine
}
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 }) {
if (job.kind === 'edit') return 'edit'
if (job.kind === 'music') return 'music'
return '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
const ZOMBIE_EMPTY_COMFY_MS = 90 * 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 failZombieLiveJob(job: Job, error: string) {
if (job.status === 'complete' || job.status === 'error' || job.status === 'cancelled') return
job.status = 'error'
job.error = error
emitJob(job, { type: 'error', error, message: error })
void onLiveVideoSettled(job).catch(() => null)
}
function sweepStaleLiveJobs() {
const now = Date.now()
for (const job of listJobs()) {
if (job.saving) continue
if (job.status === 'queued' && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) {
failZombieLiveJob(job, 'Job never started')
continue
}
if ((job.status === 'running' || job.status === 'uploading') && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) {
failZombieLiveJob(job, 'Job never started on ComfyUI')
}
}
}
/** Kill live jobs that still say "running" after Comfy was restarted or cleared. */
async function reapZombieLiveJobs() {
const queue = await fetchLiveQueue()
if (!queue) return
if (queue.running > 0 || queue.pending > 0) return
const { fetchHistory } = await import('~/server/utils/comfy')
for (const job of listJobs()) {
// Never interrupt download/stitch/library save — Comfy is idle then by design.
if (job.saving) continue
if (job.library?.chainContinuing) continue
if (job.status !== 'running' && job.status !== 'uploading' && job.status !== 'queued') continue
if (jobIsLocallySubmitting(job)) continue
const zombieMs = (job.kind === 'edit' || job.kind === 'music') ? 45_000 : ZOMBIE_EMPTY_COMFY_MS
if (jobAgeMs(job) < zombieMs) continue
if (!job.promptId) {
if (jobAgeMs(job) >= QUEUED_GRACE_MS) failZombieLiveJob(job, 'Job never started on ComfyUI')
continue
}
const history = await fetchHistory(job.promptId).catch(() => null)
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.')
continue
}
const status = entry.status?.status_str
if (status === 'error' || status === 'interrupted') {
failZombieLiveJob(job, job.error || `ComfyUI ${status}`)
continue
}
// History success + Comfy idle: the socket waiter usually finishes the save.
// Never fake-complete video here — that skipped library save and left ghosts in the queue.
// Image/music: only settle if media was already persisted and the waiter died before settle.
if (job.kind === 'edit' && job.stillId) {
job.status = 'complete'
job.error = undefined
await onLiveVideoSettled(job).catch(() => null)
} else if (job.kind === 'music' && (job.clipId || job.audio)) {
job.status = 'complete'
job.error = undefined
await onLiveVideoSettled(job).catch(() => null)
}
}
}
function liveJobOwnsGpu(job: Job) {
if (job.library?.stopAfterCurrent) return false
if (job.library?.chainContinuing) return true
if (job.saving) return true
if (job.status === 'uploading') return true
if (job.status === 'running') {
// Without a prompt id past the submit window, this is a zombie — do not block the queue.
if (!job.promptId && jobAgeMs(job) >= SUBMIT_WINDOW_MS) return false
return true
}
if (job.status === 'queued') return Boolean(job.promptId) || jobAgeMs(job) < SUBMIT_WINDOW_MS
return false
}
/**
* Force-start only: drop dead in-memory claims when Comfy is idle.
* Must never use job startedAt alone — real videos run far longer than ZOMBIE_EMPTY_COMFY_MS
* and sit with an empty Comfy queue while downloading/stitching.
*/
function clearDeadGpuClaimsForForceStart() {
for (const live of listJobs()) {
if (live.saving) continue
if (jobIsLocallySubmitting(live)) continue
if (live.status !== 'running' && live.status !== 'queued' && live.status !== 'uploading') {
if (live.library?.chainContinuing) live.library.chainContinuing = false
continue
}
// Active prompt still owned by a living waiter — leave it alone.
if (live.promptId && jobAgeMs(live) < QUEUED_GRACE_MS) continue
if (live.promptId) continue
if (live.library?.chainContinuing) live.library.chainContinuing = false
if (jobAgeMs(live) < SUBMIT_WINDOW_MS) continue
failZombieLiveJob(live, 'Cleared stuck live job so the studio queue can start')
}
}
/** When Comfy is idle but a waiting cut-in still sits, clear ghosts and start it. */
async function forceStartWaitingIfGpuIdle(owner: string, id: string) {
const after = readJobs(owner).find(job => job.id === id)
if (!after || after.status !== 'waiting') {
return after?.status === 'running'
}
await reapZombieLiveJobs()
const queue = await fetchLiveQueue().catch(() => null)
const comfyRendering = Boolean(queue && queue.running > 0)
if (comfyRendering) return false
clearDeadGpuClaimsForForceStart()
if (listJobs().some(liveJobOwnsGpu)) return false
try {
await startStudioJob(after)
return readJobs(owner).find(job => job.id === id)?.status === 'running'
} catch {
return false
}
}
export async function videoJobsBusy() {
sweepStaleLiveJobs()
await reapZombieLiveJobs()
if (listJobs().some(jobIsLocallySubmitting)) return true
// An empty Comfy queue is not idle: the 3s buffer + last-frame extract between
// shots leaves Comfy empty while the current video still owns the GPU.
if (listJobs().some(liveJobOwnsGpu)) return true
// Studio "running" only blocks when a live job still owns the GPU, or a fresh
// pending file exists. Orphan disk rows used to block the whole queue forever.
if (listOwnersWithStudioQueues().some(owner => listStudioJobs(owner).some((job) => {
if (job.status !== 'running') return false
if (job.liveJobId) {
const live = getJob(job.liveJobId)
if (live && liveJobOwnsGpu(live)) return true
}
return pendingAlive(job)
}))) return true
const queue = await fetchLiveQueue()
if (queue) {
if (queue.running > 0) return true
// Ghost queue_pending after a clear/restart used to show Ready in the header
// (health only checks running) while forever blocking Start now / auto-dispatch.
if (queue.pending > 0) {
if (listJobs().some(job => liveJobOwnsGpu(job) || jobIsLocallySubmitting(job))) return true
if (listPendingJobs().some((pending) => {
if (!pending.promptId || pending.stopAfterCurrent) return false
return Date.now() - pending.startedAt < PENDING_ORPHAN_MS
})) return true
return false
}
return false
}
return 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' : params.kind === 'music' ? 'music' : 'video'
const shotCount = kind === 'music'
? 1
: 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
})
void import('~/server/utils/generationLog')
.then(({ recordGenerationQueued }) => recordGenerationQueued({
ownerKey: params.ownerKey,
studioJobId: job.id,
kind,
payload: params.payload
}))
.catch(() => null)
return job
}
export async function patchStudioJob(owner: string, id: string, patch: (job: StudioJob) => void) {
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' })
patch(job)
job.updatedAt = Date.now()
return structuredClone(job)
})
if (updated.status === 'error' || updated.status === 'running' || updated.status === 'cancelled' || updated.status === 'complete') {
void import('~/server/utils/generationLog')
.then(({ updateGenerationLogByStudioJob }) => updateGenerationLogByStudioJob(owner, id, {
status: updated.status === 'running' ? 'running' : updated.status,
lastError: updated.lastError
}))
.catch(() => null)
}
return updated
}
function canCancelStudioStatus(status: StudioJobStatus) {
return status === 'waiting' || status === 'held' || status === 'running' || status === 'error'
}
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 (!canCancelStudioStatus(before.status)) {
throw createError({ statusCode: 409, statusMessage: 'That job is no longer in the queue' })
}
const liveJobId = before.liveJobId
const shotQueueId = before.shotQueueId
if (shotQueueId) {
const { markShotQueueCancelled, setQueueJob } = await import('~/server/utils/shotQueue')
markShotQueueCancelled(shotQueueId)
setQueueJob(shotQueueId, null)
}
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 (!canCancelStudioStatus(job.status)) {
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.shotQueueId = undefined
job.lastError = 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)
})
await stopLiveGeneration(liveJobId, shotQueueId)
if (shotQueueId) {
const { deleteShotQueue } = await import('~/server/utils/shotQueue')
await deleteShotQueue(owner, 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()
void import('~/server/utils/generationLog')
.then(({ updateGenerationLogByStudioJob }) => updateGenerationLogByStudioJob(owner, id, { status: 'cancelled' }))
.catch(() => null)
return result
}
/** Wipe waiting/running/held studio rows + live/pending/shot queues for this owner (stuck recovery). */
export async function clearStuckStudioWork(owner: string) {
const liveIds = new Set<string>()
const cancelledStudio: string[] = []
const clearedPending: string[] = []
for (const job of listJobs()) {
if (job.library?.ownerKey !== owner) continue
if (job.status === 'complete' || job.status === 'error' || job.status === 'cancelled') continue
liveIds.add(job.id)
job.status = 'cancelled'
job.error = 'Cleared by force reset'
if (job.library) {
job.library.stopAfterCurrent = true
job.library.queueAutoRun = false
}
emitJob(job, { type: 'error', error: 'Cleared by force reset', message: 'Cleared by force reset' })
}
for (const pending of listPendingJobs()) {
if (pending.ownerKey !== owner) continue
deletePendingJob(pending.jobId)
clearedPending.push(pending.jobId)
}
const before = readJobs(owner)
for (const job of before) {
if (!canCancelStudioStatus(job.status)) continue
try {
await cancelStudioJob(owner, job.id)
cancelledStudio.push(job.id)
} catch {
await mutateStore(owner, (store) => {
const row = store.jobs.find(item => item.id === job.id)
if (!row) return
row.status = 'cancelled'
row.liveJobId = undefined
row.lastError = 'Cleared by force reset'
row.updatedAt = Date.now()
syncPausedFlag(store)
}).catch(() => null)
cancelledStudio.push(job.id)
}
}
for (const id of liveIds) deletePendingJob(id)
const { forceClearAllShotQueues } = await import('~/server/utils/shotQueue')
const shots = await forceClearAllShotQueues(owner).catch(() => ({ removed: 0 }))
kickStudioQueue()
return {
cancelledStudio: cancelledStudio.length,
clearedPending: clearedPending.length,
cancelledLive: liveIds.size,
clearedShotQueues: Number(shots?.removed || 0)
}
}
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, markShotQueueCancelled } = await import('~/server/utils/shotQueue')
const queue = getShotQueue(owner, queueId)
if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
markShotQueueCancelled(queueId)
const generating = queue.status === 'running' || queueIsProcessing(queueId)
const liveJobId = getQueueJobId(queueId) || queue.currentJobId
if (generating || liveJobId) {
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)
// Only a live "running" row is generating. Held leftovers must not steal Start now.
const running = jobsNow.find(job => job.status === 'running')
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()
// GPU idle + still waiting = dispatcher missed it (ghost live ownership, etc). Force start.
let started = readJobs(owner).find(job => job.id === id)?.status === 'running'
if (!started) {
started = await forceStartWaitingIfGpuIdle(owner, id)
}
const latest = readJobs(owner).find(job => job.id === id) || updated
return {
...summarizeStudioJob(latest),
started,
deferred: Boolean(running) && !started
}
}
export async function retryStudioJob(owner: string, id: string) {
const jobsNow = readJobs(owner)
const target = jobsNow.find(job => job.id === id)
if (!target) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' })
if (target.status !== 'error' && target.status !== 'held' && target.status !== 'waiting') {
throw createError({ statusCode: 409, statusMessage: 'Only a failed or paused job can be retried' })
}
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 !== 'error' && job.status !== 'held' && job.status !== 'waiting') {
throw createError({ statusCode: 409, statusMessage: 'Only a failed or paused job can be retried' })
}
job.status = 'waiting'
job.cutIn = true
job.lastError = undefined
job.liveJobId = undefined
job.pausedByUser = false
job.pauseAfterCurrent = false
job.holdForCutIn = false
job.updatedAt = Date.now()
for (const item of jobs) {
if (item.id !== id && item.status === 'waiting') item.cutIn = false
}
return structuredClone(job)
})
await kickStudioQueue()
const started = await forceStartWaitingIfGpuIdle(owner, id)
return {
...summarizeStudioJob(readJobs(owner).find(job => job.id === id) || updated),
started
}
}
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) {
const settled = await 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)
})
if (settled) {
void import('~/server/utils/generationLog')
.then(({ updateGenerationLogByStudioJob }) => updateGenerationLogByStudioJob(owner, settled.id, {
status: settled.status === 'held' ? 'queued' : status,
lastError: error
}))
.catch(() => null)
}
return settled
}
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.pausedByUser
&& (job.holdForCutIn === true || job.resumeAutoRun === true)
))
.sort((a, b) => a.createdAt - b.createdAt)[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) {
const freshEnough = (startedAt: number) => Date.now() - startedAt < PENDING_ORPHAN_MS
if (job.liveJobId) {
const pending = readPendingJob(job.liveJobId)
if (pending?.promptId && freshEnough(pending.startedAt)) return true
}
if (!job.shotQueueId) return false
return listPendingJobs().some(pending => (
pending.ownerKey === job.ownerKey
&& pending.queueId === job.shotQueueId
&& Boolean(pending.promptId)
&& freshEnough(pending.startedAt)
))
}
function isTransientComfyError(error?: string) {
return /still busy|COMFY_BUSY|offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker|refused to start|API never answered|cold start|ECONNREFUSED|ETIMEDOUT|fetch failed|502|503/i.test(String(error || ''))
}
async function parkStudioOnStartFailure(owner: string, id: string, message: string) {
const transient = isTransientComfyError(message)
await patchStudioJob(owner, id, (job) => {
if (transient) {
job.status = 'waiting'
job.cutIn = true
job.liveJobId = undefined
job.lastError = message
job.pausedByUser = false
job.pauseAfterCurrent = false
job.holdForCutIn = false
} else {
job.status = 'error'
job.lastError = message
}
job.updatedAt = Date.now()
}).catch(() => null)
}
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 === 'error' && isTransientComfyError(job.lastError)) {
job.status = 'waiting'
job.liveJobId = undefined
job.lastError = undefined
job.pausedByUser = false
job.holdForCutIn = false
job.updatedAt = Date.now()
continue
}
// GPU sleep mid-chain used to leave held+pausedByUser with a transient lastError — Start now refused those.
if (job.status === 'held' && isTransientComfyError(job.lastError)) {
job.status = 'waiting'
job.liveJobId = undefined
job.lastError = undefined
job.pausedByUser = false
job.pauseAfterCurrent = false
job.holdForCutIn = false
job.updatedAt = Date.now()
continue
}
// Image/music never use shot queues — held without one is a deadlock leftover.
if (job.status === 'held' && !job.shotQueueId) {
job.status = 'complete'
job.liveJobId = undefined
job.holdForCutIn = false
job.pauseAfterCurrent = false
job.pausedByUser = false
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 liveFailed = live && (live.status === 'error' || live.status === 'cancelled')
const holdPause = job.pauseAfterCurrent === true || job.pausedByUser === true
if (liveFailed) {
job.status = live.status === 'cancelled' ? 'cancelled' : 'error'
job.liveJobId = undefined
job.lastError = live.error || job.lastError || 'Job failed'
job.cutIn = false
job.holdForCutIn = false
job.updatedAt = Date.now()
continue
}
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 {
// No shot queue (image/music): never park as held — that blocked waiting jobs.
job.status = 'complete'
job.liveJobId = undefined
job.pauseAfterCurrent = false
job.pausedByUser = false
job.holdForCutIn = false
job.updatedAt = Date.now()
}
}
}
async function dispatchStudioQueue() {
sweepStaleLiveJobs()
// Repair every owner first. videoJobsBusy() looks at all owners — fixing only the
// current owner left a stuck "running" row on a later owner blocking the GPU forever.
const owners = listOwnersWithStudioQueues()
for (const owner of owners) {
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 (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting'))) {
scheduleKickRetry()
}
return
}
// Re-read after videoJobsBusy — reap/settle may have rewritten disk rows.
for (const owner of owners) {
const store = readStore(owner)
const cutIn = pickCutIn(store.jobs)
if (cutIn) {
await startStudioJob(cutIn)
return
}
const held = pickHeld(store.jobs)
if (held?.shotQueueId) {
const resumed = await resumeHeldStudioJob(held)
if (resumed) return
// Failed resume used to `return` here and starve every waiting job forever.
scheduleKickRetry()
}
const waiting = pickWaiting(store.jobs)
if (waiting) {
await startStudioJob(waiting)
return
}
}
}
async function resumeHeldStudioJob(item: StudioJob) {
const latest = readJobs(item.ownerKey).find(job => job.id === item.id)
if (!latest || latest.status === 'cancelled') return false
if (!latest.shotQueueId) {
await patchStudioJob(item.ownerKey, item.id, (job) => {
if (job.status === 'cancelled') return
job.status = 'complete'
job.holdForCutIn = false
})
return false
}
const { isShotQueueCancelled, getShotQueue } = await import('~/server/utils/shotQueue')
const shotQueue = getShotQueue(item.ownerKey, latest.shotQueueId)
if (!shotQueue || isShotQueueCancelled(latest.shotQueueId) || pendingSegmentCount(shotQueue) <= 0) {
await patchStudioJob(item.ownerKey, item.id, (job) => {
if (job.status === 'cancelled') return
job.status = 'complete'
job.liveJobId = undefined
job.holdForCutIn = false
job.pauseAfterCurrent = false
job.pausedByUser = false
job.updatedAt = Date.now()
})
return false
}
await patchStudioJob(item.ownerKey, item.id, (job) => {
if (job.status === 'cancelled') return
job.status = 'running'
job.lastError = undefined
})
if (readJobs(item.ownerKey).find(job => job.id === item.id)?.status === 'cancelled') return false
const { startQueueBurst } = await import('~/server/utils/videoChain')
const count = latest.resumeAutoRun ? 'all' as const : 1
try {
const started = await startQueueBurst(item.ownerKey, latest.shotQueueId, count)
await patchStudioJob(item.ownerKey, item.id, (job) => {
if (job.status === 'cancelled') return
job.status = 'running'
job.liveJobId = started.jobId
job.holdForCutIn = false
job.cutIn = false
})
return readJobs(item.ownerKey).find(job => job.id === item.id)?.status === 'running'
} catch (error) {
const message = error instanceof Error ? error.message : String(error)
if (/already processing/i.test(message)) return true
if (/removed|not found/i.test(message)) {
await patchStudioJob(item.ownerKey, item.id, (job) => {
if (job.status === 'cancelled') return
job.status = 'complete'
job.liveJobId = undefined
job.holdForCutIn = false
job.pauseAfterCurrent = false
job.pausedByUser = false
job.lastError = message
job.updatedAt = Date.now()
}).catch(() => null)
return false
}
const transient = isTransientComfyError(message)
await patchStudioJob(item.ownerKey, item.id, (job) => {
if (job.status === 'cancelled') return
if (transient) {
job.status = 'waiting'
job.cutIn = true
job.liveJobId = undefined
job.lastError = message
job.holdForCutIn = false
job.pausedByUser = false
job.pauseAfterCurrent = false
} else {
job.status = 'error'
job.lastError = message
job.holdForCutIn = false
}
job.updatedAt = Date.now()
})
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 iterate = payload.imagePipeline === 'v2' && payload.v2Mode === 'iterate'
const generate = payload.imagePipeline === 'v2' && (payload.v2Mode === 'generate' || (iterate && !payload.stillId))
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 === 'iterate'
? 'iterate'
: payload.v2Mode === 'compose'
? 'compose'
: payload.v2Mode === 'refine'
? 'refine'
: payload.v2Mode === 'generate'
? 'generate'
: 'edit'
const graphMode = mode === 'iterate'
? (image ? (reference ? 'compose' : 'edit') : 'generate')
: mode
const maskId = payload.maskStillId
const mask = graphMode === 'refine' && maskId && existsSync(stillPath(item.ownerKey, maskId))
? {
filename: payload.maskStillFilename || 'refine-mask.png',
data: readFileSync(stillPath(item.ownerKey, maskId)),
type: 'image/png'
}
: null
if (graphMode === 'refine' && !mask) {
throw new Error('Refine requires a mask. Refusing to fall back to Edit.')
}
if (graphMode === '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.')
}
const passCount = (mode === 'edit' || mode === 'compose' || mode === 'iterate') ? (payload.passes?.length || 0) : 0
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,
promptMid: payload.promptMid || payload.prompt,
promptPre: payload.promptPre,
promptPost: payload.promptPost,
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 + passCount,
chainLabel: passCount
? (mode === 'iterate' ? 'Iteration 1' : (payload.passMode === 'chain' ? 'Pass 1' : 'Batch 1'))
: undefined,
passes: payload.passes,
...persistLoraFields(payload.loraStack || payload.loraName)
}
})
await markStudioLive(item.ownerKey, item.id, live.id)
const v2Engine = payload.engine === 'krea' ? 'krea' : 'flux'
const v2Specials = imageV2StackSpecials(payload.loraStack, v2Engine)
void runEditV2(live, {
mode,
engine: v2Engine,
task: payload.v2Task || (graphMode === 'generate' ? 't2i' : 'scene'),
image: graphMode === 'generate' ? null : image,
reference: graphMode === 'compose' ? reference : null,
mask,
prompt: payload.prompt,
negative: payload.negative || '',
steps: payload.steps,
seed: payload.seed || Math.floor(Math.random() * 2_147_483_647),
cfg: payload.cfg,
snofsModel: v2Specials.snofs?.strengthModel ?? payload.snofsModel ?? 0,
snofsClip: v2Specials.snofs?.strengthClip ?? payload.snofsClip ?? 0,
consistencyModel: v2Specials.consistency?.strengthModel ?? payload.consistencyModel ?? 0,
consistencyClip: v2Specials.consistency?.strengthClip ?? payload.consistencyClip ?? 0,
megapixels: payload.scaleMegapixels || 0,
strength: payload.refineStrength,
width: payload.width,
height: payload.height,
turbo: payload.turbo === true,
aspect: payload.aspect,
sourceStillId: payload.stillId,
referenceStillId: payload.referenceStillId,
loraStack: payload.loraStack,
passes: payload.passes,
passMode: payload.passMode === 'chain' ? 'chain' : 'batch',
}).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,
promptMid: payload.promptMid || payload.prompt,
promptPre: payload.promptPre,
promptPost: payload.promptPost,
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 ? (payload.passMode === 'chain' ? 'Pass 1' : 'Batch 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,
passMode: payload.passMode === 'chain' ? 'chain' : 'batch',
negative: payload.negative || '',
steps: payload.steps,
seed: live.library?.seed || payload.seed,
cfg: payload.cfg,
aspect: payload.aspect || 'auto',
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 parkStudioOnStartFailure(item.ownerKey, item.id, message)
kickStudioQueue()
}
}
async function startStudioExtendJob(item: StudioJob) {
let live: Job | undefined
try {
const { beginExtendFromClip } = await import('~/server/utils/extendChain')
live = await beginExtendFromClip({
ownerKey: item.ownerKey,
clipId: String(item.payload.extendFromClipId || ''),
prompt: item.payload.prompt,
promptPre: item.payload.promptPre,
promptPost: item.payload.promptPost,
duration: item.payload.duration,
loraStack: item.payload.loraStack || item.payload.loraName,
folderLocked: item.payload.folderLocked,
name: item.payload.name || item.name,
refineExtensionFrame: item.payload.refineExtensionFrame,
saveLosslessAnchor: item.payload.saveLosslessAnchor,
refinementDenoise: item.payload.refinementDenoise
})
await markStudioLive(item.ownerKey, item.id, live.id)
} 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 parkStudioOnStartFailure(item.ownerKey, item.id, message)
kickStudioQueue()
}
}
async function startStudioMusicJob(item: StudioJob) {
let live: import('~/server/utils/jobs').Job | undefined
try {
const { startMusicJob } = await import('~/server/utils/musicChain')
const payload = item.payload
live = await startMusicJob({
ownerKey: item.ownerKey,
folderId: payload.folderId,
name: payload.name,
tags: payload.prompt,
lyrics: payload.instrumental ? '' : (payload.lyrics || ''),
duration: payload.duration,
steps: payload.steps,
seed: payload.seed || Math.floor(Math.random() * 2_147_483_647),
cfg: payload.cfg,
lyricsStrength: payload.instrumental ? 0 : (payload.lyricsStrength ?? 0.9),
instrumental: payload.instrumental === true,
folderLocked: payload.folderLocked,
engine: payload.musicEngine || 'ace-step',
samplerName: payload.samplerName || 'euler',
scheduler: payload.scheduler || 'simple'
})
await markStudioLive(item.ownerKey, item.id, live.id)
} 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 parkStudioOnStartFailure(item.ownerKey, item.id, message)
kickStudioQueue()
}
}
export async function startStudioJob(item: StudioJob) {
if (studioJobKind(item) === 'music') {
await startStudioMusicJob(item)
return
}
if (studioJobKind(item) === 'edit') {
await startStudioEditJob(item)
return
}
if (item.payload.extendFromClipId) {
await startStudioExtendJob(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.promptMid || payload.prompt,
promptMid: payload.promptMid || payload.prompt,
promptPre: payload.promptPre,
promptPost: payload.promptPost,
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,
refineExtensionFrame: payload.refineExtensionFrame !== false,
saveLosslessAnchor: payload.saveLosslessAnchor === true,
refinementDenoise: payload.refinementDenoise,
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),
...persistPromptWrappers(payload),
initial: {
prompt: payload.promptMid || payload.prompt,
duration: payload.duration,
permanenceRefs: payload.shotPermanenceRefs?.[0],
...persistLoraFields(resolveLoraStack(payload.loraStack || payload.loraName, payload.shotLoraStacks?.[0] || payload.shotLoras?.[0])),
...persistPromptWrappers(payload)
},
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 parkStudioOnStartFailure(item.ownerKey, item.id, message)
kickStudioQueue()
}
}
function remainingStudioShots(job: Job) {
const library = job.library
if (!library) return 0
if (library.queueId) {
const queue = getShotQueue(library.ownerKey, library.queueId)
if (queue) return pendingSegmentCount(queue)
}
// Image/music are one studio slot (multi-pass runs inside the live job). Parking them
// as held without a shot-queue id blocked every waiting item behind them.
if (job.kind !== 'video') return 0
return Math.max(0, (library.chainTotal || 1) - (library.chainStep || 1))
}
export async function onLiveVideoSettled(job: Job) {
const owner = job.library?.ownerKey
if (!owner) {
kickStudioQueue()
return
}
const remaining = remainingStudioShots(job)
const wakeFail = job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail
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 (wakeFail) {
// GPU/host blip mid-chain — put back on waiting so wake/Start now can pick it up.
// Never park as user-paused held (that made Start now refuse the job).
await mutateStore(owner, (store) => {
const row = store.jobs.find(item => item.liveJobId === job.id)
|| (job.library?.queueId ? store.jobs.find(item => item.shotQueueId === job.library?.queueId) : undefined)
if (!row || row.status === 'cancelled') {
syncPausedFlag(store)
return
}
row.status = 'waiting'
row.liveJobId = undefined
row.cutIn = true
row.holdForCutIn = false
row.pauseAfterCurrent = false
row.pausedByUser = false
row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
row.lastError = job.error
row.updatedAt = Date.now()
syncPausedFlag(store)
})
} else if (!failed && remaining > 0) {
await mutateStore(owner, (store) => {
const row = store.jobs.find(item => item.liveJobId === job.id)
|| (job.library?.queueId ? store.jobs.find(item => item.shotQueueId === job.library?.queueId) : undefined)
if (!row || row.status === 'cancelled') {
syncPausedFlag(store)
return
}
const auto = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
row.status = 'held'
row.liveJobId = undefined
row.holdForCutIn = auto
row.pauseAfterCurrent = false
row.pausedByUser = !auto
row.resumeAutoRun = auto
row.lastError = undefined
row.updatedAt = Date.now()
syncPausedFlag(store)
})
} 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) {
const { isShotQueueCancelled } = await import('~/server/utils/shotQueue')
if (isShotQueueCancelled(queueId)) return
await mutateStore(owner, (store) => {
const row = store.jobs.find(job => job.shotQueueId === queueId)
if (!row || row.status === 'cancelled') {
syncPausedFlag(store)
return
}
row.status = 'running'
row.liveJobId = liveJobId
row.pauseAfterCurrent = false
row.pausedByUser = false
row.holdForCutIn = false
row.updatedAt = Date.now()
syncPausedFlag(store)
})
}