1293 lines
44 KiB
TypeScript
1293 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, isXaigenStudio, 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'
|
|
engine?: 'flux' | 'krea'
|
|
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,
|
|
...persistLoraFields(payload.loraStack || payload.loraName)
|
|
}
|
|
})
|
|
await markStudioLive(item.ownerKey, item.id, live.id)
|
|
void runEditV2(live, {
|
|
mode,
|
|
engine: payload.engine === 'krea' ? 'krea' : 'flux',
|
|
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 ?? (isXaigenStudio() ? 0.65 : 0),
|
|
snofsClip: payload.snofsClip ?? (isXaigenStudio() ? 0.35 : 0),
|
|
consistencyModel: payload.consistencyModel ?? (mode === 'generate' || payload.engine === 'krea' ? 0 : 0.7),
|
|
consistencyClip: payload.consistencyClip ?? (mode === 'generate' || payload.engine === 'krea' ? 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,
|
|
loraStack: payload.loraStack
|
|
}).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)
|
|
})
|
|
}
|