From 6fec10febd7a952185dff828536c4f8753a2f26b Mon Sep 17 00:00:00 2001 From: Towsty Date: Thu, 27 Aug 2026 08:42:12 -0500 Subject: [PATCH] Keep queued extension chains running after refresh instead of saving the first clip as final. Co-authored-by: Cursor --- pages/index.vue | 27 +++ server/api/generate.post.ts | 252 +--------------------------- server/utils/jobs.ts | 3 + server/utils/pending.ts | 131 +++++++++++++++ server/utils/videoChain.ts | 323 ++++++++++++++++++++++++++++++++++++ server/utils/watch.ts | 88 ++++------ 6 files changed, 524 insertions(+), 300 deletions(-) create mode 100644 server/utils/videoChain.ts diff --git a/pages/index.vue b/pages/index.vue index 1863e2c..922eb4f 100644 --- a/pages/index.vue +++ b/pages/index.vue @@ -2531,6 +2531,33 @@ function applyVideoEvent(payload: Record, hidden: boolean, folderLo } if (complete && (payload.clipId || payload.filename)) { + const total = Math.max(1, Number(payload.chainTotal || 0) || 0, Number(chainTotal.value || 0) || 0) + const step = Math.max(1, Number(payload.chainStep || 0) || 0, Number(chainStep.value || 0) || 0) + if (total > 1 && step < total) { + completedChainStep.value = Math.max(completedChainStep.value, step) + loadLibrary().catch(() => null) + const dest = folders.value.find(folder => folder.id === folderId.value) + const locked = dest + ? Boolean(dest.protected && !dest.unlocked) + : (folderLocked || payload.folderLocked === true) + if (!watchingLibrary.value && !locked && payload.clipId) { + currentClipId.value = payload.clipId + videoUrl.value = clipVideoUrl(payload.clipId) + concealOutput.value = hidden || payload.hideThumbnail === true + } + statusMessage.value = `Shot ${step} of ${total} saved, but the remaining shots did not start.` + progress.value = Math.min(Number(payload.progress || progress.value || 99), 99) + toast(`Shot ${step} of ${total} saved. Remaining shots did not run — extend this clip from the library.`) + videoBusy.value = false + statusBusy.value = false + stopTimer('video') + stopListen('video') + clearActiveJob('video') + activeDraftId = '' + void pollComfyHealth() + videoSettledUi = true + return + } videoSettledUi = true const dest = folders.value.find(folder => folder.id === folderId.value) const locked = dest diff --git a/server/api/generate.post.ts b/server/api/generate.post.ts index aa389a1..1d114c9 100644 --- a/server/api/generate.post.ts +++ b/server/api/generate.post.ts @@ -1,5 +1,4 @@ -import { copyFileSync, readFileSync, writeFileSync } from 'node:fs' -import { join } from 'node:path' +import { frameLength, runGeneration } from '~/server/utils/videoChain' function isComfyBusyTimeout(error: unknown) { const err = error as { statusCode?: number; status?: number; data?: { code?: string }; message?: string; statusMessage?: string } @@ -37,10 +36,6 @@ function parseExtensions(raw: string | undefined) { const SAMPLERS = new Set(['res_multistep', 'euler', 'dpmpp_2m']) const SCHEDULERS = new Set(['simple', 'ddim_uniform', 'sgm_uniform']) -function frameLength(seconds: number, fps: number) { - return Math.max(5, Math.floor(seconds * fps)) -} - function parseFps(raw: string | undefined) { const fps = Number(raw) return fps === 12 || fps === 30 || fps === 24 ? fps : 24 @@ -62,19 +57,6 @@ function parseScheduler(raw: string | undefined) { return SCHEDULERS.has(raw || '') ? raw! : 'simple' } -function sleep(ms: number) { - return new Promise(resolve => setTimeout(resolve, ms)) -} - -function assertJobActive(job: ReturnType) { - if (job.status === 'cancelled') { - throw new Error('Job interrupted.') - } - if (job.status === 'error') { - throw new Error(job.error || 'Generation failed') - } -} - export default defineEventHandler(async (event) => { const form = await readMultipartFormData(event) if (!form?.length) { @@ -152,10 +134,11 @@ export default defineEventHandler(async (event) => { height, hideInput }) + const referenceStillIds: Array = [null, null, null, null] if (useIdentityRefs) { for (const [index, ref] of referenceImages.entries()) { if (!ref?.data?.length || isPipelineFrameFilename(ref.filename)) continue - await saveStill({ + const saved = await saveStill({ ownerKey, folderId, filename: ref.filename || `identity-ref-${index + 2}.png`, @@ -164,6 +147,7 @@ export default defineEventHandler(async (event) => { height, hideInput }) + if (saved?.id) referenceStillIds[index] = saved.id } } @@ -193,6 +177,7 @@ export default defineEventHandler(async (event) => { thumb: isPipelineFrameFilename(image.filename) ? undefined : image.data, stillId: still?.id, stillFilename: still?.filename, + referenceStillIds, duration: durationSeconds, sound, draftId: (fields.draftId || '').trim() || undefined, @@ -202,7 +187,8 @@ export default defineEventHandler(async (event) => { chainTotal, chainLabel: extensions.length ? 'Initial' : undefined, familyId: crypto.randomUUID(), - workflow + workflow, + useIdentityRefs } emitChainJob(job, { type: 'status', message: 'Checking ComfyUI...', progress: 1 }) @@ -289,227 +275,3 @@ export default defineEventHandler(async (event) => { chainTotal } }) - -type GenerateParams = { - prompt: string - image: { filename: string; data: Buffer; type?: string } - 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 }[] - workflow: 'v1' | 'v2' - duration: number - useIdentityRefs: boolean - referenceImages: Array<{ filename: string; data: Buffer; type?: string } | null> -} - -async function queueMiniMax( - job: ReturnType, - 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 uploading = 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 = (job.library?.chainIndex || 0) > 0 - ? 'Queueing extension on MiniMax H3...' - : 'Queueing MiniMax H3 job...' - emitChainJob(job, { type: 'status', message: uploading, progress: 4 }) - const uploaded = await uploadImage(params.image, job.id) - 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) - - const graph = buildWorkflow({ - prompt: params.prompt, - 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: comfyFilenamePrefix(), - sound: params.sound, - workflow: params.workflow, - duration: params.duration, - useIdentityRefs: params.useIdentityRefs, - referenceImageNames: params.useIdentityRefs ? referenceNames : [] - }) - - 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({ - jobId: job.id, - promptId: job.promptId, - clientId: job.clientId, - ownerKey: job.library.ownerKey, - folderId: job.library.folderId, - hideThumbnail: job.library.hideThumbnail, - folderLocked: job.library.folderLocked, - name: job.library.name, - prompt: job.library.prompt, - aspect: job.library.aspect, - width: job.library.width, - height: job.library.height, - steps: job.library.steps, - turbo: job.library.turbo, - seed: job.library.seed, - startedAt: job.startedAt, - imageName: job.library.imageName, - imageSubfolder: job.library.imageSubfolder, - extendTmpDir: job.library.extendTmpDir, - extendPart1Path: job.library.extendPart1Path, - familyId: job.library.familyId, - parentClipId: job.library.parentClipId, - chainIndex: job.library.chainIndex, - stillId: job.library.stillId, - sound: job.library.sound, - referenceImageNames: job.library.referenceImageNames - }) - } - 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') - } -} - -async function runGeneration( - job: ReturnType, - params: GenerateParams -) { - 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 - }) - } - - await ensureComfyReady(ready) - - const { extensions, ...base } = params - await queueMiniMax(job, { - ...base, - persist: extensions.length === 0 - }) - - if (!extensions.length || !job.library) return - - const tmpDir = extendTempDir(job.library.ownerKey, job.id) - job.library.extendTmpDir = tmpDir - const currentPath = join(tmpDir, 'current.mp4') - const part1Path = join(tmpDir, 'part1.mp4') - const framePath = join(tmpDir, 'last_frame.png') - writeFileSync(currentPath, job.segmentBuffer!) - job.segmentBuffer = undefined - - for (let i = 0; i < extensions.length; i++) { - assertJobActive(job) - const ext = extensions[i] - const isLast = i === extensions.length - 1 - job.library.chainIndex = i + 1 - job.library.chainStep = i + 2 - job.library.chainLabel = `Extension ${i + 1}` - job.library.extendPart1Path = undefined - job.library.thumb = undefined - - emitChainJob(job, { type: 'status', message: 'Waiting 3s buffer...', progress: 1 }) - await sleep(3000) - assertJobActive(job) - - copyFileSync(currentPath, part1Path) - try { - await extractLastFrame(part1Path, 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) - job.library.extendPart1Path = part1Path - job.library.prompt = ext.prompt - const sound = await probeHasAudio(currentPath) - const seed = Math.floor(Math.random() * 2_147_483_647) - - await ensureComfyReady(ready) - await queueMiniMax(job, { - prompt: ext.prompt, - image: params.useIdentityRefs - ? 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: isLast, - workflow: params.workflow, - duration: ext.duration, - useIdentityRefs: params.useIdentityRefs, - referenceImages: params.useIdentityRefs ? params.referenceImages : [] - }) - - if (!isLast) { - if (!job.segmentBuffer?.length) { - throw new Error('Extension finished without a stitched video') - } - writeFileSync(currentPath, job.segmentBuffer) - job.segmentBuffer = undefined - job.library.extendPart1Path = undefined - } - } -} diff --git a/server/utils/jobs.ts b/server/utils/jobs.ts index 25a8877..dcd1bcf 100644 --- a/server/utils/jobs.ts +++ b/server/utils/jobs.ts @@ -67,6 +67,8 @@ export interface Job { stillId?: string stillFilename?: string referenceImageNames?: string[] + referenceStillIds?: Array + useIdentityRefs?: boolean duration?: number sound?: boolean draftId?: string @@ -81,6 +83,7 @@ export interface Job { familyId?: string parentClipId?: string workflow?: 'v1' | 'v2' + chainContinuing?: boolean } error?: string socketReady?: boolean diff --git a/server/utils/pending.ts b/server/utils/pending.ts index db137b3..3d983ef 100644 --- a/server/utils/pending.ts +++ b/server/utils/pending.ts @@ -1,5 +1,6 @@ import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, rmSync, writeFileSync } from 'node:fs' import { join } from 'node:path' +import type { Job } from '~/server/utils/jobs' export interface PendingJob { jobId: string @@ -25,9 +26,25 @@ export interface PendingJob { familyId?: string parentClipId?: string chainIndex?: number + chainStep?: number + chainTotal?: number + chainLabel?: string stillId?: string + stillFilename?: string sound?: boolean referenceImageNames?: string[] + referenceStillIds?: Array + useIdentityRefs?: boolean + workflow?: 'v1' | 'v2' + duration?: number + cfg?: number + fps?: number + samplerName?: string + scheduler?: string + hideInput?: boolean + extensions?: { prompt: string; duration: number }[] + remainingExtensions?: { prompt: string; duration: number }[] + currentClipId?: string } function pendingRoot() { @@ -43,6 +60,118 @@ function ensurePending() { mkdirSync(pendingRoot(), { recursive: true }) } +export function remainingAfterCurrentShot(job: { + chainIndex?: number + extensions?: { prompt: string; duration: number }[] +}) { + const extensions = job.extensions || [] + const index = job.chainIndex || 0 + return index === 0 ? extensions : extensions.slice(index) +} + +export function isLastChainShot(job: { + chainIndex?: number + extensions?: { prompt: string; duration: number }[] + remainingExtensions?: { prompt: string; duration: number }[] +}) { + if (job.remainingExtensions) return job.remainingExtensions.length === 0 + return remainingAfterCurrentShot(job).length === 0 +} + +export function pendingFromJob(job: Job, overrides: Partial = {}): PendingJob { + const library = job.library + if (!library) { + throw new Error('Cannot persist a job without library metadata') + } + const remaining = remainingAfterCurrentShot(library) + return { + jobId: job.id, + promptId: job.promptId || '', + clientId: job.clientId, + ownerKey: library.ownerKey, + folderId: library.folderId, + hideThumbnail: library.hideThumbnail, + folderLocked: library.folderLocked, + name: library.name, + prompt: library.prompt, + aspect: library.aspect, + width: library.width, + height: library.height, + steps: library.steps, + turbo: library.turbo, + seed: library.seed, + startedAt: job.startedAt, + imageName: library.imageName, + imageSubfolder: library.imageSubfolder, + extendTmpDir: library.extendTmpDir, + extendPart1Path: library.extendPart1Path, + familyId: library.familyId, + parentClipId: library.parentClipId, + chainIndex: library.chainIndex, + chainStep: library.chainStep, + chainTotal: library.chainTotal, + chainLabel: library.chainLabel, + stillId: library.stillId, + stillFilename: library.stillFilename, + sound: library.sound, + referenceImageNames: library.referenceImageNames, + referenceStillIds: library.referenceStillIds, + useIdentityRefs: library.useIdentityRefs, + workflow: library.workflow, + duration: library.duration, + cfg: library.cfg, + fps: library.fps, + samplerName: library.samplerName, + scheduler: library.scheduler, + hideInput: library.hideInput, + extensions: library.extensions, + remainingExtensions: remaining, + currentClipId: job.clipId, + ...overrides + } +} + +export function libraryFromPending(pending: PendingJob): NonNullable { + return { + ownerKey: pending.ownerKey, + folderId: pending.folderId, + hideThumbnail: pending.hideThumbnail, + hideInput: pending.hideInput, + folderLocked: pending.folderLocked, + name: pending.name, + prompt: pending.prompt, + aspect: pending.aspect, + width: pending.width, + height: pending.height, + steps: pending.steps, + turbo: pending.turbo, + seed: pending.seed, + cfg: pending.cfg, + fps: pending.fps, + samplerName: pending.samplerName, + scheduler: pending.scheduler, + imageName: pending.imageName, + imageSubfolder: pending.imageSubfolder, + extendTmpDir: pending.extendTmpDir, + extendPart1Path: pending.extendPart1Path, + familyId: pending.familyId, + parentClipId: pending.parentClipId, + chainIndex: pending.chainIndex, + chainStep: pending.chainStep, + chainTotal: pending.chainTotal, + chainLabel: pending.chainLabel, + stillId: pending.stillId, + stillFilename: pending.stillFilename, + sound: pending.sound, + referenceImageNames: pending.referenceImageNames, + referenceStillIds: pending.referenceStillIds, + useIdentityRefs: pending.useIdentityRefs, + workflow: pending.workflow, + duration: pending.duration, + extensions: pending.extensions || pending.remainingExtensions + } +} + export function writePendingJob(job: PendingJob) { ensurePending() const tmp = pendingPath(job.jobId) + '.tmp' @@ -73,6 +202,8 @@ export function deletePendingJob(jobId: string) { } export async function completePendingIfReady(pending: PendingJob) { + if (!isLastChainShot(pending)) return null + if (!pending.promptId) return null const alreadySaved = () => { const live = getJob(pending.jobId) if (live?.savedPromptId === pending.promptId) { diff --git a/server/utils/videoChain.ts b/server/utils/videoChain.ts new file mode 100644 index 0000000..2d61239 --- /dev/null +++ b/server/utils/videoChain.ts @@ -0,0 +1,323 @@ +import { copyFileSync, existsSync, readFileSync, writeFileSync } from 'node:fs' +import { join } from 'node:path' +import 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 { clipVideoPath, deleteRetryDraft, extendTempDir, stillPath } from '~/server/utils/library' +import { extractLastFrame, probeHasAudio } from '~/server/utils/ffmpeg' +import { ensureComfyReady } from '~/server/utils/comfyLifecycle' + +export type ChainImage = { filename: string; data: Buffer; type?: string } + +type VideoChainParams = { + prompt: string + image: ChainImage + 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 }[] + workflow: 'v1' | 'v2' + duration: number + useIdentityRefs: boolean + referenceImages: Array +} + +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) + if (!image?.data.length) { + 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: library.workflow || 'v1', + duration, + useIdentityRefs: library.useIdentityRefs === true, + referenceImages + } +} + +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 uploading = 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 = (job.library?.chainIndex || 0) > 0 + ? 'Queueing extension on MiniMax H3...' + : 'Queueing MiniMax H3 job...' + emitChainJob(job, { type: 'status', message: uploading, progress: 4 }) + const uploaded = await uploadImage(params.image, job.id) + 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) + + const graph = buildWorkflow({ + prompt: params.prompt, + 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: comfyFilenamePrefix(), + sound: params.sound, + workflow: params.workflow, + duration: params.duration, + useIdentityRefs: params.useIdentityRefs, + referenceImageNames: params.useIdentityRefs ? referenceNames : [] + }) + + 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') } +} + +export async function continueQueuedExtensions(job: Job, params: VideoChainParams) { + const extensions = params.extensions || [] + if (!extensions.length || !job.library) return + if (job.library.chainContinuing) return + job.library.chainContinuing = true + + 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 { + const { currentPath, part1Path, framePath } = seedCurrentVideo(job) + const startFrom = job.library.chainIndex || 0 + + for (let i = startFrom; i < extensions.length; i++) { + assertJobActive(job) + const ext = extensions[i] + const isLast = i === extensions.length - 1 + job.library.chainIndex = i + 1 + job.library.chainStep = i + 2 + job.library.chainLabel = `Extension ${i + 1}` + job.library.extendPart1Path = undefined + job.library.thumb = undefined + + emitChainJob(job, { type: 'status', message: 'Waiting 3s buffer...', progress: 1 }) + await sleep(3000) + assertJobActive(job) + + copyFileSync(currentPath, part1Path) + try { + await extractLastFrame(part1Path, 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) + job.library.extendPart1Path = part1Path + job.library.prompt = ext.prompt + const sound = await probeHasAudio(currentPath) + const seed = Math.floor(Math.random() * 2_147_483_647) + + await ensureComfyReady(ready) + await queueMiniMax(job, { + prompt: ext.prompt, + image: params.useIdentityRefs + ? 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: isLast, + workflow: params.workflow, + duration: ext.duration, + useIdentityRefs: params.useIdentityRefs, + referenceImages: params.useIdentityRefs ? params.referenceImages : [] + }) + + if (!isLast) { + if (!job.segmentBuffer?.length) { + throw new Error('Extension finished without a stitched video') + } + writeFileSync(currentPath, job.segmentBuffer) + job.segmentBuffer = undefined + job.library.extendPart1Path = undefined + } + } + } finally { + if (job.library) job.library.chainContinuing = false + } +} + +export async function continueQueuedExtensionsIfNeeded(job: Job) { + if (!job.library) return + if (job.status === 'error' || job.status === 'cancelled' || job.status === 'complete') return + const remaining = remainingAfterCurrentShot(job.library) + if (!remaining.length) return + await continueQueuedExtensions(job, paramsFromJob(job)) +} + +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 + }) + } + + await ensureComfyReady(ready) + + const { extensions, ...base } = params + await queueMiniMax(job, { + ...base, + persist: extensions.length === 0 + }) + + if (!extensions.length || !job.library) return + await continueQueuedExtensions(job, params) +} diff --git a/server/utils/watch.ts b/server/utils/watch.ts index 2eaf086..ba5bbb4 100644 --- a/server/utils/watch.ts +++ b/server/utils/watch.ts @@ -1,6 +1,8 @@ import { existsSync } from 'node:fs' import type { Job, JobEvent } from '~/server/utils/jobs' import { getJob, restoreJob } from '~/server/utils/jobs' +import type { PendingJob } from '~/server/utils/pending' +import { isLastChainShot, libraryFromPending, pendingFromJob, writePendingJob, deletePendingJob } from '~/server/utils/pending' function classifyError(message: string) { const lower = message.toLowerCase() @@ -219,12 +221,16 @@ export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Pr }) if (!persist) { job.segmentBuffer = buffer - deletePendingJob(job.id) + job.status = 'running' + writePendingJob(pendingFromJob(job, { + promptId: '', + currentClipId: clip.id + })) job.saving = false emitLocal({ type: 'checkpoint', message: stitching ? 'Extension checkpoint saved' : 'Initial segment saved', - progress: 99, + progress: chainMeta(job).chained ? mapChainProgress(job, 99) : 99, clipId: clip.id, hideThumbnail: job.hideThumbnail, folderLocked: job.library.folderLocked @@ -385,70 +391,42 @@ export async function waitForComfySocket(job: Job, ms = 4000) { const pendingWatches = new Set() -export function ensurePendingWatch(pending: { - jobId: string - promptId: string - clientId: string - ownerKey: string - folderId: string - hideThumbnail: boolean - folderLocked?: boolean - name?: string - prompt: string - aspect: string - width: number - height: number - steps: number - turbo: boolean - seed: number - startedAt: number - imageName?: string - imageSubfolder?: string - extendTmpDir?: string - extendPart1Path?: string - familyId?: string - parentClipId?: string - chainIndex?: number - stillId?: string - sound?: boolean - referenceImageNames?: string[] -}) { +export function ensurePendingWatch(pending: PendingJob) { const existing = getJob(pending.jobId) if (existing) return existing + const lastShot = isLastChainShot(pending) const job = restoreJob({ id: pending.jobId, clientId: pending.clientId, promptId: pending.promptId, startedAt: pending.startedAt, hideThumbnail: pending.hideThumbnail, - library: { - ownerKey: pending.ownerKey, - folderId: pending.folderId, - hideThumbnail: pending.hideThumbnail, - folderLocked: pending.folderLocked, - name: pending.name, - prompt: pending.prompt, - aspect: pending.aspect, - width: pending.width, - height: pending.height, - steps: pending.steps, - turbo: pending.turbo, - seed: pending.seed, - imageName: pending.imageName, - imageSubfolder: pending.imageSubfolder, - extendTmpDir: pending.extendTmpDir, - extendPart1Path: pending.extendPart1Path, - familyId: pending.familyId, - parentClipId: pending.parentClipId, - chainIndex: pending.chainIndex, - stillId: pending.stillId, - sound: pending.sound, - referenceImageNames: pending.referenceImageNames - } + library: libraryFromPending(pending) }) + if (pending.currentClipId) job.clipId = pending.currentClipId if (!pendingWatches.has(job.id)) { pendingWatches.add(job.id) - void watchComfyJob(job, { persist: true }).finally(() => pendingWatches.delete(job.id)) + const resume = async () => { + try { + if (pending.promptId) { + await watchComfyJob(job, { persist: lastShot }) + } + if (job.status === 'error' || job.status === 'cancelled' || job.status === 'complete') return + if (!lastShot) { + const { continueQueuedExtensionsIfNeeded } = await import('~/server/utils/videoChain') + await continueQueuedExtensionsIfNeeded(job) + } + } catch (error) { + const message = error instanceof Error ? error.message : String(error) + if (job.status === 'error' || job.status === 'cancelled') return + job.status = 'error' + job.error = message + emitChainJob(job, { type: 'error', error: message, message }) + } finally { + pendingWatches.delete(job.id) + } + } + void resume() } return job }