Files
aigen/server/utils/videoChain.ts
T

604 lines
22 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, nextClipPartName, 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,
lastCompletedIndex,
liveSegmentPrompt,
prepareQueueBurst,
remainingFromQueue,
setQueueJob,
updateShotQueue
} from '~/server/utils/shotQueue'
import { composeShotPrompt, allowIdentityRefs, type PermanenceRef } from '~/utils/globalLocks'
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: { prompt: string; duration: number; loraName?: string; loraStack?: LoraStackItem[]; permanenceRefs?: PermanenceRef[] }[]
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) Object.assign(job.library, persistLoraFields(params.loraStack || params.loraName || job.library.loraStack || job.library.loraName))
const engineName = isLtxWorkflow(graphId) ? 'LTX-2.3' : 'MiniMax H3'
const hasImage = Boolean(params.image?.data?.length)
const uploading = !hasImage
? `Queueing ${engineName} text-to-video…`
: params.useIdentityRefs
? (chainIndex > 0 ? 'Uploading identity stills for next shot...' : 'Uploading image to ComfyUI...')
: (chainIndex > 0 ? '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: 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.library.stopAfterCurrent) return true
const queue = job.library.queueId
? getShotQueue(job.library.ownerKey, job.library.queueId)
: null
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 shotStack = resolveLoraStack(
readLoraStack(liveQueue).length ? readLoraStack(liveQueue) : (params.loraStack || params.loraName),
live.loraStack || live.loraName || extensions[i]?.loraStack || extensions[i]?.loraName
)
if (!ext.prompt.trim()) {
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}`)
}
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 = ext.prompt
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
)
await ensureComfyReady(ready)
await queueMiniMax(job, {
prompt: ext.prompt,
image: identity
? params.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' || job.status === 'complete') return
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') {
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 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: nextClipPartName(clipTitle(source)),
prompt: queue.segments[lastIndex]?.prompt || source.prompt,
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)
if (job.status === 'error' || job.status === 'cancelled') {
await finishQueueBurst(owner, queue.id, false, message).catch(() => null)
setQueueJob(queue.id, null)
return
}
job.status = 'error'
job.error = message
emitJob(job, { type: 'error', error: message, message })
await finishQueueBurst(owner, queue.id, false, 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 }
}