Post stillId with video generate and do not queue I2V behind a phantom busy job.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+58
-16
@@ -220,11 +220,35 @@ export function summarizeStudioJob(job: StudioJob) {
|
||||
}
|
||||
}
|
||||
|
||||
export function videoJobsBusy() {
|
||||
if (listJobs().some(job => job.status === 'queued' || job.status === 'uploading' || job.status === 'running')) {
|
||||
return true
|
||||
const QUEUED_GRACE_MS = 4 * 60 * 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
|
||||
return false
|
||||
}
|
||||
|
||||
function sweepStaleLiveJobs() {
|
||||
const now = Date.now()
|
||||
for (const job of listJobs()) {
|
||||
if (job.status !== 'queued') continue
|
||||
if (job.promptId) continue
|
||||
if (now - job.startedAt < QUEUED_GRACE_MS) continue
|
||||
job.status = 'error'
|
||||
job.error = job.error || 'Job never started'
|
||||
}
|
||||
return listPendingJobs().some(job => Boolean(job.promptId))
|
||||
}
|
||||
|
||||
export function videoJobsBusy() {
|
||||
sweepStaleLiveJobs()
|
||||
if (listJobs().some(liveJobIsActive)) return true
|
||||
return listPendingJobs().some((pending) => {
|
||||
if (!pending.promptId) return false
|
||||
const live = getJob(pending.jobId)
|
||||
if (live) return liveJobIsActive(live)
|
||||
return Date.now() - pending.startedAt < PENDING_ORPHAN_MS
|
||||
})
|
||||
}
|
||||
|
||||
export function applyPersistedPauseToJob(job: Job) {
|
||||
@@ -534,7 +558,7 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
||||
}
|
||||
|
||||
async function dispatchStudioQueue() {
|
||||
if (videoJobsBusy()) return
|
||||
sweepStaleLiveJobs()
|
||||
for (const owner of listOwnersWithStudioQueues()) {
|
||||
const store = await mutateStore(owner, (current) => {
|
||||
migrateGlobalPause(current)
|
||||
@@ -547,17 +571,22 @@ async function dispatchStudioQueue() {
|
||||
jobs: current.jobs.map(item => structuredClone(item))
|
||||
}
|
||||
})
|
||||
if (videoJobsBusy()) return
|
||||
const cutIn = pickCutIn(store.jobs)
|
||||
if (cutIn) {
|
||||
await startStudioJob(cutIn)
|
||||
return
|
||||
}
|
||||
const waiting = pickWaiting(store.jobs)
|
||||
const held = pickHeld(store.jobs)
|
||||
if (held?.shotQueueId) {
|
||||
await resumeHeldStudioJob(held)
|
||||
if (waiting && (!held || waiting.createdAt >= held.updatedAt)) {
|
||||
await startStudioJob(waiting)
|
||||
return
|
||||
}
|
||||
const waiting = pickWaiting(store.jobs)
|
||||
if (held?.shotQueueId) {
|
||||
const resumed = await resumeHeldStudioJob(held)
|
||||
if (resumed) return
|
||||
}
|
||||
if (waiting) {
|
||||
await startStudioJob(waiting)
|
||||
return
|
||||
@@ -571,7 +600,7 @@ async function resumeHeldStudioJob(item: StudioJob) {
|
||||
job.status = 'complete'
|
||||
job.holdForCutIn = false
|
||||
})
|
||||
return
|
||||
return false
|
||||
}
|
||||
const { startQueueBurst } = await import('~/server/utils/videoChain')
|
||||
const count = item.resumeAutoRun ? 'all' as const : 1
|
||||
@@ -583,6 +612,7 @@ async function resumeHeldStudioJob(item: StudioJob) {
|
||||
job.holdForCutIn = false
|
||||
job.cutIn = false
|
||||
})
|
||||
return true
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
@@ -590,10 +620,12 @@ async function resumeHeldStudioJob(item: StudioJob) {
|
||||
job.lastError = message
|
||||
job.holdForCutIn = false
|
||||
})
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
async function startStudioEditJob(item: StudioJob) {
|
||||
let live: Job | undefined
|
||||
try {
|
||||
const { stillPath } = await import('~/server/utils/library')
|
||||
const { existsSync, readFileSync } = await import('node:fs')
|
||||
@@ -618,7 +650,7 @@ async function startStudioEditJob(item: StudioJob) {
|
||||
: null
|
||||
|
||||
const passes = payload.passes || []
|
||||
const job = createEditLiveJob({
|
||||
live = createEditLiveJob({
|
||||
steps: payload.steps,
|
||||
hideThumbnail: payload.hideThumbnail,
|
||||
library: {
|
||||
@@ -648,28 +680,32 @@ async function startStudioEditJob(item: StudioJob) {
|
||||
}
|
||||
})
|
||||
|
||||
await markStudioLive(item.ownerKey, item.id, job.id)
|
||||
void runEdit(job, {
|
||||
await markStudioLive(item.ownerKey, item.id, live.id)
|
||||
void runEdit(live, {
|
||||
image,
|
||||
reference,
|
||||
prompt: payload.prompt,
|
||||
passes,
|
||||
negative: payload.negative || '',
|
||||
steps: payload.steps,
|
||||
seed: job.library?.seed || payload.seed,
|
||||
seed: live.library?.seed || payload.seed,
|
||||
cfg: payload.cfg,
|
||||
scaleToTotalPixels: payload.scaleToTotalPixels === true,
|
||||
scaleMegapixels: payload.scaleMegapixels,
|
||||
...persistLoraFields(payload.loraStack || payload.loraName)
|
||||
}).catch((error) => {
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
if (job.status !== 'error' && job.status !== 'cancelled' && job.status !== 'deferred') {
|
||||
job.status = 'error'
|
||||
job.error = message
|
||||
if (live && live.status !== 'error' && live.status !== 'cancelled' && live.status !== 'deferred') {
|
||||
live.status = 'error'
|
||||
live.error = message
|
||||
}
|
||||
})
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
if (live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')) {
|
||||
live.status = 'error'
|
||||
live.error = message
|
||||
}
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
job.status = 'error'
|
||||
job.lastError = message
|
||||
@@ -683,6 +719,7 @@ export async function startStudioJob(item: StudioJob) {
|
||||
await startStudioEditJob(item)
|
||||
return
|
||||
}
|
||||
let live: Job | undefined
|
||||
try {
|
||||
const { createJob } = await import('~/server/utils/jobs')
|
||||
const { runGeneration, frameLength } = await import('~/server/utils/videoChain')
|
||||
@@ -702,6 +739,7 @@ export async function startStudioJob(item: StudioJob) {
|
||||
payload.shotPermanenceRefs
|
||||
)
|
||||
const job = createJob()
|
||||
live = job
|
||||
job.kind = 'video'
|
||||
job.maxStep = payload.steps
|
||||
job.hideThumbnail = payload.hideThumbnail
|
||||
@@ -840,6 +878,10 @@ export async function startStudioJob(item: StudioJob) {
|
||||
})
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
if (live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')) {
|
||||
live.status = 'error'
|
||||
live.error = message
|
||||
}
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
job.status = 'error'
|
||||
job.lastError = message
|
||||
|
||||
Reference in New Issue
Block a user