Finish the current video shot list before starting the next queued video.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -43,7 +43,7 @@ function settings() {
|
|||||||
healthTimeoutMs: 2500,
|
healthTimeoutMs: 2500,
|
||||||
bootPollAttempts: Number(process.env.COMFY_BOOT_POLL_ATTEMPTS || 20),
|
bootPollAttempts: Number(process.env.COMFY_BOOT_POLL_ATTEMPTS || 20),
|
||||||
bootPollDelayMs: 3000,
|
bootPollDelayMs: 3000,
|
||||||
startTimeoutMs: Number(process.env.COMFY_START_TIMEOUT_MS || 180_000),
|
startTimeoutMs: Number(process.env.COMFY_START_TIMEOUT_MS || 300_000),
|
||||||
busyWaitMs: Number(process.env.COMFY_BUSY_WAIT_MS || 180_000)
|
busyWaitMs: Number(process.env.COMFY_BUSY_WAIT_MS || 180_000)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -320,7 +320,11 @@ export async function prepareQueueBurst(owner: string, id: string, count: number
|
|||||||
if (activeJobs.has(queue.id)) {
|
if (activeJobs.has(queue.id)) {
|
||||||
throw createError({ statusCode: 409, statusMessage: 'This batch is already processing' })
|
throw createError({ statusCode: 409, statusMessage: 'This batch is already processing' })
|
||||||
}
|
}
|
||||||
const remaining = queue.segments.filter(segment => segment.status === 'pending' || segment.status === 'error')
|
const remaining = queue.segments.filter(segment => (
|
||||||
|
segment.status === 'pending'
|
||||||
|
|| segment.status === 'error'
|
||||||
|
|| (segment.status === 'running' && !activeJobs.has(queue.id))
|
||||||
|
))
|
||||||
if (!remaining.length) {
|
if (!remaining.length) {
|
||||||
throw createError({ statusCode: 400, statusMessage: 'No remaining shots to process' })
|
throw createError({ statusCode: 400, statusMessage: 'No remaining shots to process' })
|
||||||
}
|
}
|
||||||
|
|||||||
+66
-20
@@ -2,7 +2,7 @@ import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFile
|
|||||||
import { join } from 'node:path'
|
import { join } from 'node:path'
|
||||||
import { getJob, listJobs, emitJob, type Job } from '~/server/utils/jobs'
|
import { getJob, listJobs, emitJob, type Job } from '~/server/utils/jobs'
|
||||||
import { listPendingJobs, patchPendingJob, readPendingJob, deletePendingJob } from '~/server/utils/pending'
|
import { listPendingJobs, patchPendingJob, readPendingJob, deletePendingJob } from '~/server/utils/pending'
|
||||||
import { getShotQueue } from '~/server/utils/shotQueue'
|
import { getShotQueue, pendingSegmentCount } from '~/server/utils/shotQueue'
|
||||||
import { fetchLiveQueue } from '~/server/utils/comfy'
|
import { fetchLiveQueue } from '~/server/utils/comfy'
|
||||||
import { isLtxWorkflow, isTextToVideo, isXaigenStudio, LTX_DISABLED_MESSAGE, ltxWorkflowEnabled, parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
|
import { isLtxWorkflow, isTextToVideo, isXaigenStudio, LTX_DISABLED_MESSAGE, ltxWorkflowEnabled, parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
|
||||||
import { persistLoraFields } from '~/server/utils/loras'
|
import { persistLoraFields } from '~/server/utils/loras'
|
||||||
@@ -262,17 +262,24 @@ function sweepStaleLiveJobs() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function liveJobOwnsGpu(job: Job) {
|
||||||
|
if (job.library?.stopAfterCurrent) return false
|
||||||
|
if (job.library?.chainContinuing) return true
|
||||||
|
if (job.status === 'uploading' || job.status === 'running') return true
|
||||||
|
if (job.status === 'queued') return Boolean(job.promptId) || jobAgeMs(job) < SUBMIT_WINDOW_MS
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
export async function videoJobsBusy() {
|
export async function videoJobsBusy() {
|
||||||
sweepStaleLiveJobs()
|
sweepStaleLiveJobs()
|
||||||
if (listJobs().some(jobIsLocallySubmitting)) return true
|
if (listJobs().some(jobIsLocallySubmitting)) return true
|
||||||
|
// An empty Comfy queue is not idle: the 3s buffer + last-frame extract between
|
||||||
|
// shots leaves Comfy empty while the current video still owns the GPU.
|
||||||
|
if (listJobs().some(liveJobOwnsGpu)) return true
|
||||||
|
if (listOwnersWithStudioQueues().some(owner => listStudioJobs(owner).some(job => job.status === 'running'))) return true
|
||||||
const queue = await fetchLiveQueue()
|
const queue = await fetchLiveQueue()
|
||||||
if (queue) return queue.running > 0 || queue.pending > 0
|
if (queue) return queue.running > 0 || queue.pending > 0
|
||||||
return listJobs().some((job) => {
|
return listPendingJobs().some((pending) => {
|
||||||
if (job.library?.stopAfterCurrent) return false
|
|
||||||
if (job.status === 'uploading' || job.status === 'running') return true
|
|
||||||
if (job.status === 'queued') return jobAgeMs(job) < SUBMIT_WINDOW_MS
|
|
||||||
return false
|
|
||||||
}) || listPendingJobs().some((pending) => {
|
|
||||||
if (!pending.promptId || pending.stopAfterCurrent) return false
|
if (!pending.promptId || pending.stopAfterCurrent) return false
|
||||||
const live = getJob(pending.jobId)
|
const live = getJob(pending.jobId)
|
||||||
if (live) return jobIsLocallySubmitting(live) || live.status === 'running' || live.status === 'uploading'
|
if (live) return jobIsLocallySubmitting(live) || live.status === 'running' || live.status === 'uploading'
|
||||||
@@ -525,10 +532,10 @@ function pickHeld(jobs: StudioJob[]) {
|
|||||||
.filter(job => (
|
.filter(job => (
|
||||||
job.status === 'held'
|
job.status === 'held'
|
||||||
&& job.shotQueueId
|
&& job.shotQueueId
|
||||||
&& job.holdForCutIn
|
|
||||||
&& !job.pausedByUser
|
&& !job.pausedByUser
|
||||||
|
&& (job.holdForCutIn === true || job.resumeAutoRun === true)
|
||||||
))
|
))
|
||||||
.sort((a, b) => b.updatedAt - a.updatedAt)[0] || null
|
.sort((a, b) => a.createdAt - b.createdAt)[0] || null
|
||||||
}
|
}
|
||||||
|
|
||||||
function pickWaiting(jobs: StudioJob[]) {
|
function pickWaiting(jobs: StudioJob[]) {
|
||||||
@@ -628,16 +635,17 @@ async function dispatchStudioQueue() {
|
|||||||
await startStudioJob(cutIn)
|
await startStudioJob(cutIn)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
const held = pickHeld(store.jobs)
|
||||||
|
if (held?.shotQueueId) {
|
||||||
|
const resumed = await resumeHeldStudioJob(held)
|
||||||
|
if (!resumed) scheduleKickRetry()
|
||||||
|
return
|
||||||
|
}
|
||||||
const waiting = pickWaiting(store.jobs)
|
const waiting = pickWaiting(store.jobs)
|
||||||
if (waiting) {
|
if (waiting) {
|
||||||
await startStudioJob(waiting)
|
await startStudioJob(waiting)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
const held = pickHeld(store.jobs)
|
|
||||||
if (held?.shotQueueId) {
|
|
||||||
const resumed = await resumeHeldStudioJob(held)
|
|
||||||
if (resumed) return
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -649,6 +657,10 @@ async function resumeHeldStudioJob(item: StudioJob) {
|
|||||||
})
|
})
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||||
|
job.status = 'running'
|
||||||
|
job.lastError = undefined
|
||||||
|
})
|
||||||
const { startQueueBurst } = await import('~/server/utils/videoChain')
|
const { startQueueBurst } = await import('~/server/utils/videoChain')
|
||||||
const count = item.resumeAutoRun ? 'all' as const : 1
|
const count = item.resumeAutoRun ? 'all' as const : 1
|
||||||
try {
|
try {
|
||||||
@@ -662,10 +674,12 @@ async function resumeHeldStudioJob(item: StudioJob) {
|
|||||||
return true
|
return true
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
const message = error instanceof Error ? error.message : String(error)
|
const message = error instanceof Error ? error.message : String(error)
|
||||||
|
if (/already processing/i.test(message)) return true
|
||||||
|
const transient = /offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker/i.test(message)
|
||||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||||
job.status = 'error'
|
job.status = transient ? 'held' : 'error'
|
||||||
job.lastError = message
|
job.lastError = message
|
||||||
job.holdForCutIn = false
|
job.holdForCutIn = transient
|
||||||
})
|
})
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
@@ -1031,16 +1045,29 @@ export async function startStudioJob(item: StudioJob) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function remainingStudioShots(job: Job) {
|
||||||
|
const library = job.library
|
||||||
|
if (!library) return 0
|
||||||
|
if (library.queueId) {
|
||||||
|
const queue = getShotQueue(library.ownerKey, library.queueId)
|
||||||
|
if (queue) return pendingSegmentCount(queue)
|
||||||
|
}
|
||||||
|
return Math.max(0, (library.chainTotal || 1) - (library.chainStep || 1))
|
||||||
|
}
|
||||||
|
|
||||||
|
function isTransientComfyError(error?: string) {
|
||||||
|
return /still busy|COMFY_BUSY|offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker/i.test(String(error || ''))
|
||||||
|
}
|
||||||
|
|
||||||
export async function onLiveVideoSettled(job: Job) {
|
export async function onLiveVideoSettled(job: Job) {
|
||||||
const owner = job.library?.ownerKey
|
const owner = job.library?.ownerKey
|
||||||
if (!owner) {
|
if (!owner) {
|
||||||
kickStudioQueue()
|
kickStudioQueue()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
const remaining = job.library
|
const remaining = remainingStudioShots(job)
|
||||||
? Math.max(0, (job.library.chainTotal || 1) - (job.library.chainStep || 1))
|
const wakeFail = job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
|
||||||
: 0
|
const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail
|
||||||
const failed = job.status === 'error' || job.status === 'cancelled'
|
|
||||||
const snapshot = readStore(owner)
|
const snapshot = readStore(owner)
|
||||||
const rowNow = snapshot.jobs.find(item => item.liveJobId === job.id)
|
const rowNow = snapshot.jobs.find(item => item.liveJobId === job.id)
|
||||||
|| (job.library?.queueId
|
|| (job.library?.queueId
|
||||||
@@ -1085,6 +1112,25 @@ export async function onLiveVideoSettled(job: Job) {
|
|||||||
row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
|
row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
|
||||||
row.updatedAt = Date.now()
|
row.updatedAt = Date.now()
|
||||||
})
|
})
|
||||||
|
} else if (!failed && remaining > 0) {
|
||||||
|
await mutateStore(owner, (store) => {
|
||||||
|
const row = store.jobs.find(item => item.liveJobId === job.id)
|
||||||
|
|| (job.library?.queueId ? store.jobs.find(item => item.shotQueueId === job.library?.queueId) : undefined)
|
||||||
|
if (!row) {
|
||||||
|
syncPausedFlag(store)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
const auto = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
|
||||||
|
row.status = 'held'
|
||||||
|
row.liveJobId = undefined
|
||||||
|
row.holdForCutIn = auto
|
||||||
|
row.pauseAfterCurrent = false
|
||||||
|
row.pausedByUser = !auto
|
||||||
|
row.resumeAutoRun = auto
|
||||||
|
row.lastError = wakeFail ? job.error : undefined
|
||||||
|
row.updatedAt = Date.now()
|
||||||
|
syncPausedFlag(store)
|
||||||
|
})
|
||||||
} else if (userPause) {
|
} else if (userPause) {
|
||||||
await mutateStore(owner, (store) => {
|
await mutateStore(owner, (store) => {
|
||||||
const row = store.jobs.find(item => item.liveJobId === job.id)
|
const row = store.jobs.find(item => item.liveJobId === job.id)
|
||||||
|
|||||||
@@ -451,7 +451,12 @@ export async function continueQueuedExtensions(
|
|||||||
|
|
||||||
export async function continueQueuedExtensionsIfNeeded(job: Job) {
|
export async function continueQueuedExtensionsIfNeeded(job: Job) {
|
||||||
if (!job.library) return
|
if (!job.library) return
|
||||||
if (job.status === 'error' || job.status === 'cancelled' || job.status === 'complete') 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)) {
|
if (chainShouldHold(job)) {
|
||||||
job.library.stopAfterCurrent = true
|
job.library.stopAfterCurrent = true
|
||||||
return
|
return
|
||||||
@@ -512,6 +517,15 @@ export async function runGeneration(job: Job, params: VideoChainParams) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
export async function startQueueBurst(owner: string, queueId: string, count: number | 'all') {
|
export async function startQueueBurst(owner: string, queueId: string, count: number | 'all') {
|
||||||
|
await ensureComfyReady((status) => {
|
||||||
|
console.log(JSON.stringify({
|
||||||
|
src: 'video-chain',
|
||||||
|
event: 'wake',
|
||||||
|
queueId,
|
||||||
|
state: status.state,
|
||||||
|
message: status.message
|
||||||
|
}))
|
||||||
|
})
|
||||||
const { queue, count: n } = await prepareQueueBurst(owner, queueId, count)
|
const { queue, count: n } = await prepareQueueBurst(owner, queueId, count)
|
||||||
const lastIndex = lastCompletedIndex(queue)
|
const lastIndex = lastCompletedIndex(queue)
|
||||||
const clipId = queue.currentClipId
|
const clipId = queue.currentClipId
|
||||||
@@ -585,15 +599,13 @@ export async function startQueueBurst(owner: string, queueId: string, count: num
|
|||||||
}).catch(async (error) => {
|
}).catch(async (error) => {
|
||||||
removeExtendTemp(job.library?.extendTmpDir)
|
removeExtendTemp(job.library?.extendTmpDir)
|
||||||
const message = error instanceof Error ? error.message : String(error)
|
const message = error instanceof Error ? error.message : String(error)
|
||||||
if (job.status === 'error' || job.status === 'cancelled') {
|
const transient = /still busy|COMFY_BUSY|offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker/i.test(message)
|
||||||
await finishQueueBurst(owner, queue.id, false, message).catch(() => null)
|
|
||||||
setQueueJob(queue.id, null)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
job.status = 'error'
|
|
||||||
job.error = message
|
job.error = message
|
||||||
emitJob(job, { type: 'error', error: message, message })
|
if (!transient && job.status !== 'cancelled') {
|
||||||
await finishQueueBurst(owner, queue.id, false, message).catch(() => null)
|
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)
|
setQueueJob(queue.id, null)
|
||||||
}).finally(() => {
|
}).finally(() => {
|
||||||
void onLiveVideoSettled(job)
|
void onLiveVideoSettled(job)
|
||||||
|
|||||||
@@ -474,11 +474,12 @@ export function ensurePendingWatch(pending: PendingJob) {
|
|||||||
await onLiveVideoSettled(job)
|
await onLiveVideoSettled(job)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if (job.status === 'complete') {
|
if (job.status === 'complete' && lastShot) {
|
||||||
await onLiveVideoSettled(job)
|
await onLiveVideoSettled(job)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if (!lastShot) {
|
if (!lastShot) {
|
||||||
|
if (job.status === 'complete') job.status = 'running'
|
||||||
const { continueQueuedExtensionsIfNeeded } = await import('~/server/utils/videoChain')
|
const { continueQueuedExtensionsIfNeeded } = await import('~/server/utils/videoChain')
|
||||||
await continueQueuedExtensionsIfNeeded(job)
|
await continueQueuedExtensionsIfNeeded(job)
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user