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 { extractLastFrame, probeHasAudio } from '~/server/utils/ffmpeg' import { ensureComfyReady } from '~/server/utils/comfyLifecycle' import { onLiveVideoSettled, onStudioQueueBurstStarted } from '~/server/utils/studioQueue' import { finishQueueBurst, getShotQueue, isShotQueueCancelled, lastCompletedIndex, liveSegmentPrompt, 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 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 = [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 & { 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 : [], ...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 || queue?.status === 'paused') 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 if (liveQueue?.stopAfterCurrent || liveQueue?.autoRun === false) { if (liveQueue.stopAfterCurrent) job.library.stopAfterCurrent = true if (liveQueue.autoRun === false) 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 extractLastFrame(extractSource, 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 } const 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') } 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' 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) 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 } }