Start I2V when Beast is idle instead of queuing behind leftover jobs.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+33
-20
@@ -3,6 +3,7 @@ import { join } from 'node:path'
|
||||
import { getJob, listJobs, emitJob, type Job } from '~/server/utils/jobs'
|
||||
import { listPendingJobs, patchPendingJob, readPendingJob, deletePendingJob } from '~/server/utils/pending'
|
||||
import { getShotQueue } from '~/server/utils/shotQueue'
|
||||
import { fetchLiveQueue } from '~/server/utils/comfy'
|
||||
import { isLtxWorkflow, LTX_DISABLED_MESSAGE, ltxWorkflowEnabled, parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
|
||||
import { persistLoraFields } from '~/server/utils/loras'
|
||||
import { resolveLoraStack } from '~/utils/loras'
|
||||
@@ -221,11 +222,19 @@ export function summarizeStudioJob(job: StudioJob) {
|
||||
}
|
||||
|
||||
const QUEUED_GRACE_MS = 4 * 60 * 1000
|
||||
const SUBMIT_WINDOW_MS = 45 * 1000
|
||||
const PENDING_ORPHAN_MS = 2 * 60 * 1000
|
||||
|
||||
function liveJobIsActive(job: Job) {
|
||||
if (job.status === 'uploading' || job.status === 'running') return true
|
||||
if (job.status === 'queued') return Date.now() - job.startedAt < QUEUED_GRACE_MS
|
||||
function jobAgeMs(job: { startedAt: number }) {
|
||||
return Date.now() - job.startedAt
|
||||
}
|
||||
|
||||
/** Uploading, or created just now and not yet on Comfy. Paused leftovers do not count. */
|
||||
export function jobIsLocallySubmitting(job: Job) {
|
||||
if (job.library?.stopAfterCurrent) return false
|
||||
if (job.status === 'uploading') return true
|
||||
if (job.status === 'queued' && !job.promptId && jobAgeMs(job) < SUBMIT_WINDOW_MS) return true
|
||||
if (job.status === 'running' && !job.promptId && jobAgeMs(job) < SUBMIT_WINDOW_MS) return true
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -240,13 +249,20 @@ function sweepStaleLiveJobs() {
|
||||
}
|
||||
}
|
||||
|
||||
export function videoJobsBusy() {
|
||||
export async function videoJobsBusy() {
|
||||
sweepStaleLiveJobs()
|
||||
if (listJobs().some(liveJobIsActive)) return true
|
||||
return listPendingJobs().some((pending) => {
|
||||
if (!pending.promptId) return false
|
||||
if (listJobs().some(jobIsLocallySubmitting)) return true
|
||||
const queue = await fetchLiveQueue()
|
||||
if (queue) return queue.running > 0 || queue.pending > 0
|
||||
return listJobs().some((job) => {
|
||||
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
|
||||
const live = getJob(pending.jobId)
|
||||
if (live) return liveJobIsActive(live)
|
||||
if (live) return jobIsLocallySubmitting(live) || live.status === 'running' || live.status === 'uploading'
|
||||
return Date.now() - pending.startedAt < PENDING_ORPHAN_MS
|
||||
})
|
||||
}
|
||||
@@ -571,26 +587,22 @@ async function dispatchStudioQueue() {
|
||||
jobs: current.jobs.map(item => structuredClone(item))
|
||||
}
|
||||
})
|
||||
if (videoJobsBusy()) return
|
||||
if (await videoJobsBusy()) return
|
||||
const cutIn = pickCutIn(store.jobs)
|
||||
if (cutIn) {
|
||||
await startStudioJob(cutIn)
|
||||
return
|
||||
}
|
||||
const waiting = pickWaiting(store.jobs)
|
||||
const held = pickHeld(store.jobs)
|
||||
if (waiting && (!held || waiting.createdAt >= held.updatedAt)) {
|
||||
await startStudioJob(waiting)
|
||||
return
|
||||
}
|
||||
if (held?.shotQueueId) {
|
||||
const resumed = await resumeHeldStudioJob(held)
|
||||
if (resumed) return
|
||||
}
|
||||
if (waiting) {
|
||||
await startStudioJob(waiting)
|
||||
return
|
||||
}
|
||||
const held = pickHeld(store.jobs)
|
||||
if (held?.shotQueueId) {
|
||||
const resumed = await resumeHeldStudioJob(held)
|
||||
if (resumed) return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1050,10 +1062,11 @@ export async function toggleStudioQueuePause(owner: string, options: {
|
||||
const rowHasQueue = rows.some(job => job.shotQueueId === options.queueId)
|
||||
if (!rowHasQueue) await setShotQueuePause(owner, options.queueId, false).catch(() => null)
|
||||
}
|
||||
if (heldPaused[0] && !videoJobsBusy()) {
|
||||
const busy = await videoJobsBusy()
|
||||
if (heldPaused[0] && !busy) {
|
||||
await resumeHeldStudioJob(heldPaused[0]).catch(() => null)
|
||||
} else {
|
||||
if (heldPaused.length && videoJobsBusy()) {
|
||||
if (heldPaused.length && busy) {
|
||||
await mutate(owner, (jobs) => {
|
||||
for (const job of jobs) {
|
||||
if (!heldPaused.some(item => item.id === job.id)) continue
|
||||
|
||||
Reference in New Issue
Block a user