698 lines
26 KiB
TypeScript
698 lines
26 KiB
TypeScript
import { copyFileSync, existsSync, readFileSync, writeFileSync } from 'node:fs'
|
|
import { join } from 'node:path'
|
|
import { createJob, emitJob, type Job } from '~/server/utils/jobs'
|
|
import { pendingFromJob, remainingAfterCurrentShot, writePendingJob } from '~/server/utils/pending'
|
|
import { emitChainJob, waitForComfySocket, watchComfyJob } from '~/server/utils/watch'
|
|
import { comfyFilenamePrefix, queuePrompt, uploadImage } from '~/server/utils/comfy'
|
|
import { buildWorkflow } from '~/server/utils/workflow'
|
|
import { isLtxWorkflow, isTextToVideo, parseVideoWorkflow, workflowForExtension, type VideoWorkflowId } from '~/utils/videoModels'
|
|
import { clipVideoPath, clipTitle, deleteRetryDraft, extendTempDir, getClip, nextFamilyPartName, removeExtendTemp, stillPath } from '~/server/utils/library'
|
|
import { probeHasAudio } from '~/server/utils/ffmpeg'
|
|
import { ensureComfyReady } from '~/server/utils/comfyLifecycle'
|
|
import { onLiveVideoSettled, onStudioQueueBurstStarted } from '~/server/utils/studioQueue'
|
|
import {
|
|
refineExtensionHandoffFrame,
|
|
resolveExtensionHandoffFrame
|
|
} from '~/server/utils/extensionFrame'
|
|
import {
|
|
finishQueueBurst,
|
|
getShotQueue,
|
|
isShotQueueCancelled,
|
|
lastCompletedIndex,
|
|
liveSegmentPrompt,
|
|
pendingSegmentCount,
|
|
prepareQueueBurst,
|
|
remainingFromQueue,
|
|
setQueueJob,
|
|
updateShotQueue
|
|
} from '~/server/utils/shotQueue'
|
|
import { composeShotPrompt, allowIdentityRefs, type PermanenceRef } from '~/utils/globalLocks'
|
|
import { composePromptParts, joinPromptParts, resolvePromptWrappers, wrappedPromptForComfy } from '~/utils/promptParts'
|
|
import type { QueuedExtension } from '~/server/utils/library'
|
|
import { persistLoraFields, ensureComfyLoraNames } from '~/server/utils/loras'
|
|
import { readLoraStack, resolveLoraStack } from '~/utils/loras'
|
|
import type { LoraStackItem } from '~/utils/loras'
|
|
|
|
export type ChainImage = { filename: string; data: Buffer; type?: string }
|
|
|
|
type VideoChainParams = {
|
|
prompt: string
|
|
image?: ChainImage | null
|
|
width: number
|
|
height: number
|
|
steps: number
|
|
seed: number
|
|
turbo: boolean
|
|
length: number
|
|
sound: boolean
|
|
cfg: number
|
|
fps: number
|
|
samplerName: string
|
|
scheduler: string
|
|
extensions: QueuedExtension[]
|
|
workflow: VideoWorkflowId
|
|
duration: number
|
|
useIdentityRefs: boolean
|
|
referenceImages: Array<ChainImage | null>
|
|
globalLocks?: string
|
|
permanenceRefs?: PermanenceRef[]
|
|
shotPermanenceRefs?: PermanenceRef[][]
|
|
loraName?: string
|
|
loraStack?: LoraStackItem[]
|
|
shotLoras?: string[]
|
|
shotLoraStacks?: LoraStackItem[][]
|
|
}
|
|
|
|
function sleep(ms: number) {
|
|
return new Promise(resolve => setTimeout(resolve, ms))
|
|
}
|
|
|
|
export function frameLength(seconds: number, fps: number) {
|
|
return Math.max(5, Math.floor(seconds * fps))
|
|
}
|
|
|
|
function assertJobActive(job: Job) {
|
|
if (job.status === 'cancelled') {
|
|
throw new Error('Job interrupted.')
|
|
}
|
|
if (job.status === 'error') {
|
|
throw new Error(job.error || 'Generation failed')
|
|
}
|
|
}
|
|
|
|
function loadStillImage(owner: string, stillId?: string, filename?: string): ChainImage | null {
|
|
if (!stillId) return null
|
|
const path = stillPath(owner, stillId)
|
|
if (!existsSync(path)) return null
|
|
return {
|
|
filename: filename || 'still.png',
|
|
data: readFileSync(path),
|
|
type: 'image/png'
|
|
}
|
|
}
|
|
|
|
function paramsFromJob(job: Job): VideoChainParams {
|
|
const library = job.library
|
|
if (!library) {
|
|
throw new Error('Cannot continue the shot chain: job metadata is missing')
|
|
}
|
|
const image = loadStillImage(library.ownerKey, library.stillId, library.stillFilename)
|
|
const workflow = parseVideoWorkflow(library.workflow)
|
|
if (!image?.data.length && !isTextToVideo(workflow)) {
|
|
throw new Error('Cannot continue the shot chain: the start still is missing from the library')
|
|
}
|
|
const referenceImages: Array<ChainImage | null> = [null, null, null, null]
|
|
for (const [index, stillId] of (library.referenceStillIds || []).entries()) {
|
|
if (!stillId || index >= 4) continue
|
|
referenceImages[index] = loadStillImage(library.ownerKey, stillId, `identity-ref-${index + 2}.png`)
|
|
}
|
|
const fps = library.fps || 24
|
|
const duration = library.duration || 5
|
|
return {
|
|
prompt: library.prompt,
|
|
image,
|
|
width: library.width,
|
|
height: library.height,
|
|
steps: library.steps,
|
|
seed: library.seed,
|
|
turbo: library.turbo,
|
|
length: frameLength(duration, fps),
|
|
sound: library.sound !== false,
|
|
cfg: library.cfg || (library.turbo ? 1.5 : 4),
|
|
fps,
|
|
samplerName: library.samplerName || 'res_multistep',
|
|
scheduler: library.scheduler || 'simple',
|
|
extensions: library.extensions || [],
|
|
workflow,
|
|
duration,
|
|
useIdentityRefs: allowIdentityRefs(
|
|
library.useIdentityRefs === true,
|
|
library.permanenceRefs,
|
|
library.globalLocks,
|
|
library.shotPermanenceRefs
|
|
),
|
|
referenceImages,
|
|
globalLocks: library.globalLocks,
|
|
permanenceRefs: library.permanenceRefs,
|
|
shotPermanenceRefs: library.shotPermanenceRefs,
|
|
...persistLoraFields(readLoraStack(library)),
|
|
shotLoras: library.shotLoras,
|
|
shotLoraStacks: library.shotLoraStacks
|
|
}
|
|
}
|
|
|
|
export async function queueMiniMax(
|
|
job: Job,
|
|
params: Omit<VideoChainParams, 'extensions'> & { persist: boolean }
|
|
) {
|
|
assertJobActive(job)
|
|
job.socketReady = false
|
|
job.promptId = undefined
|
|
job.video = undefined
|
|
job.segmentBuffer = undefined
|
|
|
|
const done = watchComfyJob(job, { persist: params.persist })
|
|
job.status = 'uploading'
|
|
const chainIndex = job.library?.chainIndex || 0
|
|
const graphId = chainIndex > 0 ? workflowForExtension(params.workflow) : params.workflow
|
|
if (job.library) {
|
|
const stack = resolveLoraStack(job.library.loraStack || job.library.loraName, params.loraStack || params.loraName)
|
|
if (stack.length) Object.assign(job.library, persistLoraFields(stack))
|
|
}
|
|
const engineName = isLtxWorkflow(graphId) ? 'LTX-2.3' : 'MiniMax H3'
|
|
const hasImage = Boolean(params.image?.data?.length)
|
|
const uploading = !hasImage
|
|
? `Queueing ${engineName} text-to-video…`
|
|
: chainIndex > 0
|
|
? (params.useIdentityRefs ? 'Uploading last frame and character stills...' : 'Uploading last frame to ComfyUI...')
|
|
: 'Uploading image to ComfyUI...'
|
|
const queueing = chainIndex > 0
|
|
? `Queueing extension on ${engineName}...`
|
|
: `Queueing ${engineName} job...`
|
|
emitChainJob(job, { type: 'status', message: uploading, progress: 4 })
|
|
const uploaded = hasImage && params.image
|
|
? await uploadImage(params.image, job.id)
|
|
: { name: '', subfolder: '' }
|
|
const referenceNames: string[] = ['', '', '', '']
|
|
if (params.useIdentityRefs) {
|
|
for (const [index, ref] of (params.referenceImages || []).entries()) {
|
|
if (!ref?.data?.length) continue
|
|
const next = await uploadImage({
|
|
...ref,
|
|
filename: `ref${index + 1}_${ref.filename || 'identity.png'}`
|
|
}, job.id)
|
|
referenceNames[index] = next.name
|
|
}
|
|
}
|
|
if (job.library) {
|
|
job.library.imageName = uploaded.name
|
|
job.library.imageSubfolder = uploaded.subfolder
|
|
job.library.referenceImageNames = referenceNames.filter(Boolean)
|
|
}
|
|
emitChainJob(job, { type: 'status', message: queueing, progress: 6 })
|
|
await waitForComfySocket(job, 4000)
|
|
await ensureComfyLoraNames('video')
|
|
|
|
const shotIndex = job.library?.chainIndex || 0
|
|
const composedPrompt = composeShotPrompt({
|
|
globalLocks: job.library?.globalLocks || params.globalLocks,
|
|
prompt: wrappedPromptForComfy(job.library, job.library?.promptMid, params.prompt),
|
|
shotIndex,
|
|
familyRefs: job.library?.permanenceRefs || params.permanenceRefs,
|
|
shotRefs: job.library?.shotPermanenceRefs?.[shotIndex] || params.shotPermanenceRefs?.[shotIndex]
|
|
})
|
|
|
|
const graph = buildWorkflow({
|
|
prompt: composedPrompt,
|
|
imageName: uploaded.name,
|
|
width: params.width,
|
|
height: params.height,
|
|
steps: params.steps,
|
|
seed: params.seed,
|
|
turbo: params.turbo,
|
|
length: params.length,
|
|
cfg: params.cfg,
|
|
fps: params.fps,
|
|
samplerName: params.samplerName,
|
|
scheduler: params.scheduler,
|
|
filenamePrefix: isLtxWorkflow(graphId) ? 'video/LTX23' : comfyFilenamePrefix(),
|
|
sound: params.sound && !isLtxWorkflow(graphId),
|
|
workflow: graphId,
|
|
duration: params.duration,
|
|
useIdentityRefs: params.useIdentityRefs,
|
|
referenceImageNames: params.useIdentityRefs ? referenceNames : [],
|
|
saveLosslessAnchor: job.library?.saveLosslessAnchor !== false,
|
|
...persistLoraFields(params.loraStack || params.loraName)
|
|
})
|
|
|
|
const queued = await queuePrompt(graph, job.clientId)
|
|
job.promptId = queued.prompt_id
|
|
job.status = 'running'
|
|
if (job.library?.draftId) {
|
|
await deleteRetryDraft(job.library.ownerKey, job.library.draftId).catch(() => null)
|
|
job.library.draftId = undefined
|
|
}
|
|
if (job.library && job.promptId) {
|
|
writePendingJob(pendingFromJob(job))
|
|
}
|
|
emitChainJob(job, { type: 'status', message: 'Job queued on ComfyUI', progress: 8 })
|
|
await done
|
|
assertJobActive(job)
|
|
if (!params.persist && !job.segmentBuffer?.length) {
|
|
throw new Error('Segment finished without a video')
|
|
}
|
|
}
|
|
|
|
function seedCurrentVideo(job: Job) {
|
|
const library = job.library
|
|
if (!library) {
|
|
throw new Error('Cannot continue the shot chain: job metadata is missing')
|
|
}
|
|
const tmpDir = library.extendTmpDir || extendTempDir(library.ownerKey, job.id)
|
|
library.extendTmpDir = tmpDir
|
|
const currentPath = join(tmpDir, 'current.mp4')
|
|
if (job.segmentBuffer?.length) {
|
|
writeFileSync(currentPath, job.segmentBuffer)
|
|
job.segmentBuffer = undefined
|
|
} else if (!existsSync(currentPath)) {
|
|
const clipId = job.clipId
|
|
if (!clipId) {
|
|
throw new Error('Cannot continue the shot chain: the previous clip is missing')
|
|
}
|
|
const source = clipVideoPath(library.ownerKey, clipId)
|
|
if (!existsSync(source)) {
|
|
throw new Error('Cannot continue the shot chain: the previous clip file is missing')
|
|
}
|
|
copyFileSync(source, currentPath)
|
|
}
|
|
return { tmpDir, currentPath, part1Path: join(tmpDir, 'part1.mp4'), framePath: join(tmpDir, 'last_frame.png') }
|
|
}
|
|
|
|
function chainShouldHold(job: Job) {
|
|
if (!job.library) return false
|
|
if (job.status === 'cancelled' || job.status === 'error') return true
|
|
if (job.library.stopAfterCurrent) return true
|
|
if (job.library.queueId && isShotQueueCancelled(job.library.queueId)) return true
|
|
const queue = job.library.queueId
|
|
? getShotQueue(job.library.ownerKey, job.library.queueId)
|
|
: null
|
|
if (job.library.queueId && !queue) return true
|
|
if (queue?.stopAfterCurrent) return true
|
|
// A stale "paused" flag used to abort auto-run mid-batch after finishQueueBurst.
|
|
// Only honor paused when this live job is not explicitly auto-running.
|
|
if (queue?.status === 'paused' && job.library.queueAutoRun !== true) return true
|
|
return false
|
|
}
|
|
|
|
export async function continueQueuedExtensions(
|
|
job: Job,
|
|
params: VideoChainParams,
|
|
options: { maxShots?: number; autoRun?: boolean } = {}
|
|
) {
|
|
const extensions = params.extensions || []
|
|
if (!extensions.length || !job.library) return
|
|
if (job.library.chainContinuing) return
|
|
job.library.chainContinuing = true
|
|
const autoRun = options.autoRun === true || job.library.queueAutoRun === true
|
|
let shotsLeft = options.maxShots ?? (autoRun ? extensions.length : (job.library.queueBudget || 0))
|
|
|
|
const ready = (status: { state: string; message: string; queueRunning?: number; queuePending?: number }) => {
|
|
emitChainJob(job, {
|
|
type: status.state === 'busy' ? 'busy' : 'status',
|
|
message: status.message,
|
|
progress: status.state === 'online' ? 3 : 1,
|
|
busy: status.state === 'busy',
|
|
queueRunning: status.queueRunning,
|
|
queuePending: status.queuePending
|
|
})
|
|
}
|
|
|
|
try {
|
|
if (shotsLeft <= 0) return
|
|
const { currentPath, part1Path, framePath } = seedCurrentVideo(job)
|
|
const startFrom = job.library.chainIndex || 0
|
|
|
|
for (let i = startFrom; i < extensions.length; i++) {
|
|
if (shotsLeft <= 0) break
|
|
assertJobActive(job)
|
|
if (chainShouldHold(job)) {
|
|
job.library.stopAfterCurrent = true
|
|
break
|
|
}
|
|
|
|
const liveQueue = job.library.queueId
|
|
? getShotQueue(job.library.ownerKey, job.library.queueId)
|
|
: null
|
|
// User pause / cut-in only. Do not let a stale autoRun=false on disk cancel an
|
|
// in-flight "run remaining" / checkbox auto-run batch (options.autoRun).
|
|
if (liveQueue?.stopAfterCurrent) {
|
|
job.library.stopAfterCurrent = true
|
|
} else if (liveQueue?.autoRun === false && options.autoRun !== true) {
|
|
job.library.queueAutoRun = false
|
|
}
|
|
if (chainShouldHold(job)) {
|
|
job.library.stopAfterCurrent = true
|
|
break
|
|
}
|
|
const live = liveQueue ? liveSegmentPrompt(liveQueue, i) : extensions[i]
|
|
const ext = {
|
|
prompt: live.prompt || extensions[i].prompt,
|
|
duration: live.duration || extensions[i].duration
|
|
}
|
|
const wrappers = resolvePromptWrappers(live, extensions[i], job.library)
|
|
const shotStack = resolveLoraStack(
|
|
readLoraStack(liveQueue).length ? readLoraStack(liveQueue) : (params.loraStack || params.loraName),
|
|
live.loraStack || live.loraName || extensions[i]?.loraStack || extensions[i]?.loraName
|
|
)
|
|
if (!composePromptParts(wrappers.promptPre, ext.prompt, wrappers.promptPost)) {
|
|
throw new Error(`Shot ${i + 2} needs a prompt`)
|
|
}
|
|
|
|
emitChainJob(job, { type: 'status', message: 'Waiting 3s buffer...', progress: 1 })
|
|
await sleep(3000)
|
|
assertJobActive(job)
|
|
if (chainShouldHold(job)) {
|
|
job.library.stopAfterCurrent = true
|
|
break
|
|
}
|
|
|
|
const remainingAfter = extensions.length - (i + 1)
|
|
const lastOfChain = remainingAfter === 0
|
|
const lastOfBurst = !autoRun && shotsLeft === 1
|
|
const persist = lastOfChain || lastOfBurst || job.library.stopAfterCurrent === true
|
|
|
|
job.library.chainIndex = i + 1
|
|
job.library.chainStep = i + 2
|
|
job.library.chainLabel = `Extension ${i + 1}`
|
|
job.library.extendPart1Path = undefined
|
|
job.library.thumb = undefined
|
|
job.library.duration = ext.duration
|
|
Object.assign(job.library, persistLoraFields(shotStack))
|
|
if (liveQueue) {
|
|
await updateShotQueue(job.library.ownerKey, liveQueue.id, (queue) => {
|
|
const segment = queue.segments.find(item => item.index === i + 1)
|
|
if (segment) {
|
|
segment.status = 'running'
|
|
delete segment.error
|
|
}
|
|
queue.status = 'running'
|
|
queue.currentJobId = job.id
|
|
}).catch(() => null)
|
|
}
|
|
|
|
copyFileSync(currentPath, part1Path)
|
|
const parentId = job.library.parentClipId
|
|
const parentVideo = parentId ? clipVideoPath(job.library.ownerKey, parentId) : ''
|
|
const extractSource = parentVideo && existsSync(parentVideo) ? parentVideo : part1Path
|
|
try {
|
|
await resolveExtensionHandoffFrame({
|
|
ownerKey: job.library.ownerKey,
|
|
sourceClipId: parentId,
|
|
sourceVideoPath: extractSource,
|
|
destPath: framePath
|
|
})
|
|
} catch (error) {
|
|
const detail = error instanceof Error ? error.message : String(error)
|
|
throw new Error(detail.includes('last frame')
|
|
? detail
|
|
: `Could not extract the last frame for the next extension: ${detail}`)
|
|
}
|
|
assertJobActive(job)
|
|
if (chainShouldHold(job)) {
|
|
job.library.stopAfterCurrent = true
|
|
break
|
|
}
|
|
let frame = readFileSync(framePath)
|
|
if (!frame.length || frame.length < 64) {
|
|
throw new Error('Could not extract the last frame for the next extension: the frame file was empty')
|
|
}
|
|
if (job.library.refineExtensionFrame !== false) {
|
|
emitChainJob(job, { type: 'status', message: 'Refining extension handoff frame...', progress: 2 })
|
|
try {
|
|
frame = await refineExtensionHandoffFrame({
|
|
frame,
|
|
denoise: job.library.refinementDenoise,
|
|
clientId: job.clientId,
|
|
jobId: job.id,
|
|
onQueued: (promptId) => { job.promptId = promptId },
|
|
onStatus: (message) => emitChainJob(job, { type: 'status', message, progress: 3 })
|
|
})
|
|
writeFileSync(framePath, frame)
|
|
} catch (error) {
|
|
throw new Error(`Extension frame refinement failed: ${error instanceof Error ? error.message : String(error)}. Retry after fixing refinement, or explicitly turn it off.`)
|
|
}
|
|
}
|
|
assertJobActive(job)
|
|
job.library.extendPart1Path = part1Path
|
|
job.library.prompt = composePromptParts(wrappers.promptPre, ext.prompt, wrappers.promptPost)
|
|
job.library.promptMid = ext.prompt
|
|
job.library.promptPre = wrappers.promptPre || undefined
|
|
job.library.promptPost = wrappers.promptPost || undefined
|
|
if (live.permanenceRefs?.length) {
|
|
const refs = [...(job.library.shotPermanenceRefs || [])]
|
|
refs[i + 1] = live.permanenceRefs
|
|
job.library.shotPermanenceRefs = refs
|
|
}
|
|
const sound = await probeHasAudio(currentPath)
|
|
const seed = Math.floor(Math.random() * 2_147_483_647)
|
|
const identity = allowIdentityRefs(
|
|
params.useIdentityRefs,
|
|
job.library.permanenceRefs || params.permanenceRefs,
|
|
job.library.globalLocks || params.globalLocks,
|
|
job.library.shotPermanenceRefs || params.shotPermanenceRefs
|
|
)
|
|
|
|
assertJobActive(job)
|
|
if (chainShouldHold(job)) {
|
|
job.library.stopAfterCurrent = true
|
|
break
|
|
}
|
|
await ensureComfyReady(ready)
|
|
assertJobActive(job)
|
|
if (chainShouldHold(job)) {
|
|
job.library.stopAfterCurrent = true
|
|
break
|
|
}
|
|
await queueMiniMax(job, {
|
|
prompt: joinPromptParts(wrappers.promptPre, ext.prompt, wrappers.promptPost),
|
|
image: { filename: 'last_frame.png', data: frame, type: 'image/png' },
|
|
width: params.width,
|
|
height: params.height,
|
|
steps: params.steps,
|
|
seed,
|
|
turbo: params.turbo,
|
|
length: frameLength(ext.duration, params.fps),
|
|
sound,
|
|
cfg: params.cfg,
|
|
fps: params.fps,
|
|
samplerName: params.samplerName,
|
|
scheduler: params.scheduler,
|
|
persist,
|
|
workflow: workflowForExtension(params.workflow),
|
|
duration: ext.duration,
|
|
useIdentityRefs: identity,
|
|
referenceImages: identity ? params.referenceImages : [],
|
|
...persistLoraFields(shotStack)
|
|
})
|
|
|
|
shotsLeft -= 1
|
|
if (job.library) job.library.queueBudget = shotsLeft
|
|
|
|
if (!lastOfChain && !persist) {
|
|
if (!job.segmentBuffer?.length) {
|
|
throw new Error('Extension finished without a stitched video')
|
|
}
|
|
writeFileSync(currentPath, job.segmentBuffer)
|
|
job.segmentBuffer = undefined
|
|
job.library.extendPart1Path = undefined
|
|
}
|
|
|
|
if (shotsLeft <= 0 || chainShouldHold(job)) {
|
|
if (chainShouldHold(job)) job.library.stopAfterCurrent = true
|
|
break
|
|
}
|
|
}
|
|
} finally {
|
|
if (job.library) job.library.chainContinuing = false
|
|
if (job.library?.queueId) {
|
|
const paused = job.status !== 'error' && job.status !== 'cancelled'
|
|
const hold = job.library.stopAfterCurrent === true && paused
|
|
if (hold && job.status === 'running') job.status = 'complete'
|
|
// Keep autoRun on the shot queue when this burst was an explicit run-all and more
|
|
// shots remain — finishQueueBurst used to clear it and force one-click-per-shot.
|
|
const keepAuto = autoRun && !hold && job.status !== 'error' && job.status !== 'cancelled'
|
|
if (keepAuto && job.library) job.library.queueAutoRun = true
|
|
await finishQueueBurst(
|
|
job.library.ownerKey,
|
|
job.library.queueId,
|
|
hold || (paused && !job.library.queueAutoRun),
|
|
job.status === 'error' || job.status === 'cancelled' ? (job.error || 'Stopped') : undefined
|
|
).catch(() => null)
|
|
if (keepAuto) {
|
|
await updateShotQueue(job.library.ownerKey, job.library.queueId, (queue) => {
|
|
if (pendingSegmentCount(queue) > 0 && queue.status !== 'complete' && queue.status !== 'error') {
|
|
queue.autoRun = true
|
|
if (queue.status === 'paused') queue.status = 'idle'
|
|
}
|
|
}).catch(() => null)
|
|
}
|
|
setQueueJob(job.library.queueId, null)
|
|
}
|
|
}
|
|
}
|
|
|
|
export async function continueQueuedExtensionsIfNeeded(job: Job) {
|
|
if (!job.library) return
|
|
if (job.status === 'error' || job.status === 'cancelled') return
|
|
if (job.status === 'complete') {
|
|
const leftover = remainingAfterCurrentShot(job.library)
|
|
if (!leftover.length) return
|
|
job.status = 'running'
|
|
}
|
|
if (chainShouldHold(job)) {
|
|
job.library.stopAfterCurrent = true
|
|
return
|
|
}
|
|
const remaining = remainingAfterCurrentShot(job.library)
|
|
if (!remaining.length) return
|
|
const queue = job.library.queueId
|
|
? getShotQueue(job.library.ownerKey, job.library.queueId)
|
|
: null
|
|
if (queue?.stopAfterCurrent) return
|
|
const autoRun = queue ? queue.autoRun === true : job.library.queueAutoRun === true
|
|
const budget = job.library.queueBudget || 0
|
|
if (queue && !autoRun && budget <= 0) return
|
|
if (!queue && !autoRun && budget <= 0) {
|
|
// Legacy chains without a queue keep the old auto-continue behavior.
|
|
await continueQueuedExtensions(job, paramsFromJob(job), { autoRun: true })
|
|
return
|
|
}
|
|
const maxShots = autoRun ? remaining.length : budget
|
|
if (maxShots <= 0) return
|
|
await continueQueuedExtensions(job, paramsFromJob(job), { maxShots, autoRun })
|
|
}
|
|
|
|
export async function runGeneration(job: Job, params: VideoChainParams) {
|
|
const ready = (status: { state: string; message: string; queueRunning?: number; queuePending?: number }) => {
|
|
emitChainJob(job, {
|
|
type: status.state === 'busy' ? 'busy' : 'status',
|
|
message: status.message,
|
|
progress: status.state === 'online' ? 3 : 1,
|
|
busy: status.state === 'busy',
|
|
queueRunning: status.queueRunning,
|
|
queuePending: status.queuePending
|
|
})
|
|
}
|
|
|
|
try {
|
|
await ensureComfyReady(ready)
|
|
|
|
const { extensions, ...base } = params
|
|
const autoRun = job.library?.queueAutoRun === true
|
|
await queueMiniMax(job, {
|
|
...base,
|
|
persist: extensions.length === 0 || !autoRun
|
|
})
|
|
|
|
const stopNow = chainShouldHold(job)
|
|
const keepGoing = Boolean(autoRun && !stopNow && extensions.length && job.library)
|
|
if (job.library?.queueId && !keepGoing) {
|
|
await finishQueueBurst(job.library.ownerKey, job.library.queueId, true).catch(() => null)
|
|
setQueueJob(job.library.queueId, null)
|
|
}
|
|
|
|
if (!keepGoing) return
|
|
await continueQueuedExtensions(job, params, { maxShots: extensions.length, autoRun: true })
|
|
} finally {
|
|
await onLiveVideoSettled(job)
|
|
}
|
|
}
|
|
|
|
export async function startQueueBurst(owner: string, queueId: string, count: number | 'all') {
|
|
if (isShotQueueCancelled(queueId)) {
|
|
throw createError({ statusCode: 410, statusMessage: 'This batch was removed' })
|
|
}
|
|
await ensureComfyReady((status) => {
|
|
console.log(JSON.stringify({
|
|
src: 'video-chain',
|
|
event: 'wake',
|
|
queueId,
|
|
state: status.state,
|
|
message: status.message
|
|
}))
|
|
})
|
|
if (isShotQueueCancelled(queueId)) {
|
|
throw createError({ statusCode: 410, statusMessage: 'This batch was removed' })
|
|
}
|
|
const { queue, count: n } = await prepareQueueBurst(owner, queueId, count)
|
|
const lastIndex = lastCompletedIndex(queue)
|
|
const clipId = queue.currentClipId
|
|
if (!clipId) {
|
|
throw createError({ statusCode: 409, statusMessage: 'The previous clip is missing, so the next shot cannot start' })
|
|
}
|
|
const source = getClip(owner, clipId)
|
|
const destFolderLocked = false
|
|
const nextSegment = queue.segments.find(segment => segment.index === lastIndex + 1)
|
|
const job = createJob()
|
|
job.kind = 'video'
|
|
job.maxStep = queue.steps
|
|
job.hideThumbnail = queue.hideThumbnail
|
|
job.clipId = clipId
|
|
job.library = {
|
|
ownerKey: owner,
|
|
folderId: queue.folderId,
|
|
hideThumbnail: queue.hideThumbnail,
|
|
hideInput: queue.hideInput,
|
|
folderLocked: destFolderLocked,
|
|
name: nextFamilyPartName(owner, source),
|
|
prompt: nextSegment?.prompt || queue.segments[lastIndex]?.prompt || source.prompt,
|
|
promptMid: nextSegment?.prompt || queue.segments[lastIndex]?.prompt || source.prompt,
|
|
promptPre: queue.promptPre || source.promptPre,
|
|
promptPost: queue.promptPost || source.promptPost,
|
|
aspect: queue.aspect,
|
|
width: queue.width,
|
|
height: queue.height,
|
|
steps: queue.steps,
|
|
turbo: queue.turbo,
|
|
seed: Math.floor(Math.random() * 2_147_483_647),
|
|
cfg: queue.cfg,
|
|
fps: queue.fps,
|
|
samplerName: queue.samplerName,
|
|
scheduler: queue.scheduler,
|
|
duration: queue.segments[lastIndex]?.duration || source.duration,
|
|
stillId: queue.stillId,
|
|
stillFilename: queue.stillFilename,
|
|
referenceStillIds: queue.referenceStillIds,
|
|
sound: queue.sound,
|
|
extensions: remainingFromQueue(queue),
|
|
chainIndex: Math.max(0, lastIndex),
|
|
chainStep: lastIndex + 1,
|
|
chainTotal: queue.segments.length,
|
|
chainLabel: lastIndex <= 0 ? 'Initial' : `Extension ${lastIndex}`,
|
|
familyId: queue.familyId,
|
|
parentClipId: clipId,
|
|
workflow: queue.workflow,
|
|
useIdentityRefs: allowIdentityRefs(
|
|
queue.useIdentityRefs === true,
|
|
queue.permanenceRefs,
|
|
queue.globalLocks,
|
|
queue.segments.map(segment => segment.permanenceRefs || [])
|
|
),
|
|
queueId: queue.id,
|
|
queueAutoRun: count === 'all',
|
|
queueBudget: n,
|
|
globalLocks: queue.globalLocks,
|
|
permanenceRefs: queue.permanenceRefs,
|
|
shotPermanenceRefs: queue.segments.map(segment => segment.permanenceRefs || []),
|
|
...persistLoraFields(readLoraStack(queue))
|
|
}
|
|
setQueueJob(queue.id, job.id)
|
|
await updateShotQueue(owner, queue.id, (next) => {
|
|
next.currentJobId = job.id
|
|
next.status = 'running'
|
|
next.stopAfterCurrent = false
|
|
}).catch(() => null)
|
|
await onStudioQueueBurstStarted(owner, queue.id, job.id).catch(() => null)
|
|
|
|
emitChainJob(job, { type: 'status', message: 'Checking ComfyUI...', progress: 1 })
|
|
void continueQueuedExtensions(job, paramsFromJob(job), {
|
|
maxShots: n,
|
|
autoRun: count === 'all'
|
|
}).catch(async (error) => {
|
|
removeExtendTemp(job.library?.extendTmpDir)
|
|
const message = error instanceof Error ? error.message : String(error)
|
|
const transient = /still busy|COMFY_BUSY|offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker/i.test(message)
|
|
job.error = message
|
|
if (!transient && job.status !== 'cancelled') {
|
|
job.status = 'error'
|
|
emitJob(job, { type: 'error', error: message, message })
|
|
}
|
|
await finishQueueBurst(owner, queue.id, true, transient ? undefined : message).catch(() => null)
|
|
setQueueJob(queue.id, null)
|
|
}).finally(() => {
|
|
void onLiveVideoSettled(job)
|
|
})
|
|
|
|
return { jobId: job.id, queueId: queue.id, count: n, chainTotal: queue.segments.length, chainStep: lastIndex + 1 }
|
|
}
|