Ship extension queues, library extend, and smoother seam stitching.
Users can pre-queue extensions, extend from the library with p2 titles, blend seams with xfade, and a unreadable catalog is no longer replaced with an empty one. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+97
-16
@@ -1,5 +1,5 @@
|
||||
import { existsSync } from 'node:fs'
|
||||
import type { Job } from '~/server/utils/jobs'
|
||||
import type { Job, JobEvent } from '~/server/utils/jobs'
|
||||
|
||||
function classifyError(message: string) {
|
||||
const lower = message.toLowerCase()
|
||||
@@ -19,13 +19,71 @@ function sleep(ms: number) {
|
||||
return new Promise(resolve => setTimeout(resolve, ms))
|
||||
}
|
||||
|
||||
export function watchComfyJob(job: Job): Promise<void> {
|
||||
function chainMeta(job: Job) {
|
||||
const total = Math.max(1, job.library?.chainTotal || 1)
|
||||
const index = job.library?.chainIndex ?? 0
|
||||
return {
|
||||
total,
|
||||
index,
|
||||
step: index + 1,
|
||||
label: job.library?.chainLabel || (index === 0 ? 'Initial' : `Extension ${index}`),
|
||||
chained: total > 1
|
||||
}
|
||||
}
|
||||
|
||||
export function mapChainProgress(job: Job, localPct: number) {
|
||||
const { total, index, chained } = chainMeta(job)
|
||||
const local = Math.max(0, Math.min(100, localPct))
|
||||
if (!chained) return Math.round(local)
|
||||
const share = 100 / total
|
||||
return Math.min(99, Math.round(index * share + (Math.min(local, 99) / 100) * share))
|
||||
}
|
||||
|
||||
export function formatChainMessage(job: Job, message: string, samplePct?: number) {
|
||||
const { chained, step, total, label } = chainMeta(job)
|
||||
if (!chained) return message
|
||||
if (typeof samplePct === 'number') {
|
||||
return `Step ${step}/${total}: Generating ${label} (${samplePct}%)`
|
||||
}
|
||||
if (/waiting 3s buffer/i.test(message)) {
|
||||
return `Step ${step}/${total}: Waiting 3s buffer...`
|
||||
}
|
||||
return `Step ${step}/${total}: ${message}`
|
||||
}
|
||||
|
||||
export function emitChainJob(job: Job, event: JobEvent, samplePct?: number) {
|
||||
const { chained, step, total, label } = chainMeta(job)
|
||||
const next: JobEvent = { ...event }
|
||||
if (next.message) {
|
||||
next.message = formatChainMessage(job, next.message, samplePct)
|
||||
}
|
||||
if (chained && typeof next.progress === 'number' && next.type !== 'complete') {
|
||||
next.progress = mapChainProgress(job, next.progress)
|
||||
}
|
||||
if (chained) {
|
||||
next.chainStep = step
|
||||
next.chainTotal = total
|
||||
next.chainLabel = label
|
||||
}
|
||||
emitJob(job, next)
|
||||
}
|
||||
|
||||
export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Promise<void> {
|
||||
const persist = options.persist !== false
|
||||
job.socketReady = false
|
||||
const ws = new WebSocket(comfyWsUrl(job.clientId))
|
||||
let settled = false
|
||||
let pollTimer: ReturnType<typeof setInterval> | null = null
|
||||
let timeout: ReturnType<typeof setTimeout> | null = null
|
||||
|
||||
return new Promise((resolve) => {
|
||||
let localProgress = 0
|
||||
|
||||
const emitLocal = (event: JobEvent, samplePct?: number) => {
|
||||
if (typeof event.progress === 'number') localProgress = event.progress
|
||||
emitChainJob(job, event, samplePct)
|
||||
}
|
||||
|
||||
const cleanup = () => {
|
||||
if (timeout) clearTimeout(timeout)
|
||||
if (pollTimer) clearInterval(pollTimer)
|
||||
@@ -40,7 +98,7 @@ export function watchComfyJob(job: Job): Promise<void> {
|
||||
cleanup()
|
||||
job.status = job.status === 'cancelled' ? 'cancelled' : 'error'
|
||||
job.error = classifyError(error)
|
||||
emitJob(job, { type: 'error', error: job.error, message: job.error })
|
||||
emitChainJob(job, { type: 'error', error: job.error, message: job.error })
|
||||
if (job.id) deletePendingJob(job.id)
|
||||
removeExtendTemp(job.library?.extendTmpDir)
|
||||
resolve()
|
||||
@@ -56,9 +114,9 @@ export function watchComfyJob(job: Job): Promise<void> {
|
||||
try {
|
||||
job.video = video
|
||||
const stitching = Boolean(job.library?.extendPart1Path && job.library.extendTmpDir)
|
||||
emitJob(job, {
|
||||
emitLocal({
|
||||
type: 'status',
|
||||
message: stitching ? 'Extracting last frame & stitching extension...' : 'Saving to library...',
|
||||
message: stitching ? 'Extracting last frame & stitching extension...' : (persist ? 'Saving to library...' : 'Downloading segment...'),
|
||||
progress: stitching ? 97 : 98
|
||||
})
|
||||
if (job.library) {
|
||||
@@ -73,6 +131,23 @@ export function watchComfyJob(job: Job): Promise<void> {
|
||||
tmpDir: job.library.extendTmpDir!
|
||||
})
|
||||
}
|
||||
if (!persist) {
|
||||
job.segmentBuffer = buffer
|
||||
await purgeComfyArtifacts({
|
||||
video,
|
||||
imageName: job.library.imageName,
|
||||
imageSubfolder: job.library.imageSubfolder,
|
||||
promptId: job.promptId
|
||||
})
|
||||
deletePendingJob(job.id)
|
||||
emitLocal({
|
||||
type: 'status',
|
||||
message: stitching ? 'Extension stitched' : 'Initial segment ready',
|
||||
progress: 99
|
||||
})
|
||||
resolve()
|
||||
return true
|
||||
}
|
||||
const clip = await saveClip({
|
||||
ownerKey: job.library.ownerKey,
|
||||
folderId: job.library.folderId,
|
||||
@@ -111,22 +186,27 @@ export function watchComfyJob(job: Job): Promise<void> {
|
||||
type: 'complete',
|
||||
message: job.library?.folderLocked
|
||||
? 'Saved to the locked folder. Unlock it to view.'
|
||||
: (job.library?.extendSourceClipId ? 'Extended video ready' : 'Video ready'),
|
||||
: (job.library?.extendSourceClipId || (job.library?.chainTotal || 1) > 1 ? 'Extended video ready' : 'Video ready'),
|
||||
progress: 100,
|
||||
filename: job.library?.folderLocked ? undefined : video.filename,
|
||||
subfolder: job.library?.folderLocked ? undefined : video.subfolder,
|
||||
mediaType: job.library?.folderLocked ? undefined : video.type,
|
||||
clipId: job.clipId,
|
||||
hideThumbnail: job.hideThumbnail,
|
||||
folderLocked: job.library?.folderLocked
|
||||
folderLocked: job.library?.folderLocked,
|
||||
chainStep: job.library?.chainStep,
|
||||
chainTotal: job.library?.chainTotal,
|
||||
chainLabel: job.library?.chainLabel
|
||||
})
|
||||
deletePendingJob(job.id)
|
||||
} catch (saveError) {
|
||||
removeExtendTemp(job.library?.extendTmpDir)
|
||||
const message = saveError instanceof Error ? saveError.message : String(saveError)
|
||||
job.status = 'error'
|
||||
job.error = `Video generated but library save failed: ${message}`
|
||||
emitJob(job, { type: 'error', error: job.error, message: job.error })
|
||||
job.error = persist
|
||||
? `Video generated but library save failed: ${message}`
|
||||
: `Segment finished but could not be prepared: ${message}`
|
||||
emitChainJob(job, { type: 'error', error: job.error, message: job.error })
|
||||
}
|
||||
resolve()
|
||||
return true
|
||||
@@ -169,7 +249,7 @@ export function watchComfyJob(job: Job): Promise<void> {
|
||||
|
||||
ws.addEventListener('open', () => {
|
||||
job.socketReady = true
|
||||
emitJob(job, { type: 'status', message: 'Connected to ComfyUI', progress: Math.max(job.progress, 4) })
|
||||
emitLocal({ type: 'status', message: 'Connected to ComfyUI', progress: Math.max(localProgress, 4) })
|
||||
})
|
||||
|
||||
ws.addEventListener('close', () => {
|
||||
@@ -191,15 +271,16 @@ export function watchComfyJob(job: Job): Promise<void> {
|
||||
const value = Number(data.value || 0)
|
||||
const max = Number(data.max || 1)
|
||||
const pct = 12 + Math.round((value / Math.max(max, 1)) * 70)
|
||||
const samplePct = Math.round((value / Math.max(max, 1)) * 100)
|
||||
job.status = 'running'
|
||||
emitJob(job, {
|
||||
emitLocal({
|
||||
type: 'progress',
|
||||
step: value,
|
||||
maxStep: max,
|
||||
progress: pct,
|
||||
node: data.node ? String(data.node) : null,
|
||||
message: `Sampling step ${value}/${max} (${Math.round((value / Math.max(max, 1)) * 100)}%)`
|
||||
})
|
||||
message: `Sampling step ${value}/${max} (${samplePct}%)`
|
||||
}, samplePct)
|
||||
}
|
||||
|
||||
if (payload.type === 'executing') {
|
||||
@@ -207,16 +288,16 @@ export function watchComfyJob(job: Job): Promise<void> {
|
||||
if (node === null) {
|
||||
const saved = await succeed()
|
||||
if (!saved) {
|
||||
emitJob(job, { type: 'status', message: 'Waiting for ComfyUI to save the MP4...', progress: 96 })
|
||||
emitLocal({ type: 'status', message: 'Waiting for ComfyUI to save the MP4...', progress: 96 })
|
||||
}
|
||||
return
|
||||
}
|
||||
const label = NODE_LABELS[node] || `Running node ${node}`
|
||||
const encoding = isEncodingNode(node)
|
||||
emitJob(job, {
|
||||
emitLocal({
|
||||
type: 'executing',
|
||||
node,
|
||||
progress: encoding ? 90 : Math.max(job.progress, 10),
|
||||
progress: encoding ? 90 : Math.max(localProgress, 10),
|
||||
message: encoding ? 'Encoding video...' : `${label}...`
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user