Clear Generating as soon as a job saves: settle the studio row on persist and repair finished orphans so the next waiting job can start.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+130
-109
@@ -394,19 +394,21 @@ export async function videoJobsBusy() {
|
|||||||
sweepStaleLiveJobs()
|
sweepStaleLiveJobs()
|
||||||
await reapZombieLiveJobs()
|
await reapZombieLiveJobs()
|
||||||
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 (listJobs().some(liveJobOwnsGpu)) return true
|
||||||
// Studio row still "running" owns the queue slot until onLiveVideoSettled clears it.
|
// A studio "running" row only blocks while real work is attached. Orphans are cleared in repair.
|
||||||
// Do not peek under the live job here — that raced settle and started the next item early.
|
if (listOwnersWithStudioQueues().some(owner => listStudioJobs(owner).some((job) => {
|
||||||
if (listOwnersWithStudioQueues().some(owner => listStudioJobs(owner).some(job => job.status === 'running'))) {
|
if (job.status !== 'running') return false
|
||||||
return true
|
if (pendingAlive(job)) return true
|
||||||
}
|
if (!job.liveJobId) return false
|
||||||
|
const live = getJob(job.liveJobId)
|
||||||
|
if (!live) return false
|
||||||
|
if (live.saving || liveJobOwnsGpu(live)) return true
|
||||||
|
// Live already finished — settle/repair must clear the row; do not pretend the GPU is busy.
|
||||||
|
return false
|
||||||
|
}))) return true
|
||||||
const queue = await fetchLiveQueue()
|
const queue = await fetchLiveQueue()
|
||||||
if (queue) {
|
if (queue) {
|
||||||
if (queue.running > 0) return true
|
if (queue.running > 0) return true
|
||||||
// Ghost queue_pending after a clear/restart used to show Ready in the header
|
|
||||||
// (health only checks running) while forever blocking Start now / auto-dispatch.
|
|
||||||
if (queue.pending > 0) {
|
if (queue.pending > 0) {
|
||||||
if (listJobs().some(job => liveJobOwnsGpu(job) || jobIsLocallySubmitting(job))) return true
|
if (listJobs().some(job => liveJobOwnsGpu(job) || jobIsLocallySubmitting(job))) return true
|
||||||
if (listPendingJobs().some((pending) => {
|
if (listPendingJobs().some((pending) => {
|
||||||
@@ -897,7 +899,6 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
|||||||
job.updatedAt = Date.now()
|
job.updatedAt = Date.now()
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
// GPU sleep mid-chain used to leave held+pausedByUser with a transient lastError — Start now refused those.
|
|
||||||
if (job.status === 'held' && isTransientComfyError(job.lastError)) {
|
if (job.status === 'held' && isTransientComfyError(job.lastError)) {
|
||||||
job.status = 'waiting'
|
job.status = 'waiting'
|
||||||
job.liveJobId = undefined
|
job.liveJobId = undefined
|
||||||
@@ -908,7 +909,6 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
|||||||
job.updatedAt = Date.now()
|
job.updatedAt = Date.now()
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
// Image/music never use shot queues — held without one is a deadlock leftover.
|
|
||||||
if (job.status === 'held' && !job.shotQueueId) {
|
if (job.status === 'held' && !job.shotQueueId) {
|
||||||
job.status = 'complete'
|
job.status = 'complete'
|
||||||
job.liveJobId = undefined
|
job.liveJobId = undefined
|
||||||
@@ -920,14 +920,12 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
|||||||
}
|
}
|
||||||
if (job.status !== 'running') continue
|
if (job.status !== 'running') continue
|
||||||
const live = job.liveJobId ? getJob(job.liveJobId) : undefined
|
const live = job.liveJobId ? getJob(job.liveJobId) : undefined
|
||||||
// Live already complete: onLiveVideoSettled owns finalization. Stealing the row here
|
if (live?.saving) continue
|
||||||
// started the next queue item while purge/settle was still running → next jobs failed.
|
if (live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')) continue
|
||||||
if (live && live.status === 'complete') continue
|
if (pendingAlive(job)) continue
|
||||||
const liveBusy = live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')
|
|
||||||
if (liveBusy || pendingAlive(job)) continue
|
|
||||||
const liveFailed = live && (live.status === 'error' || live.status === 'cancelled')
|
|
||||||
const holdPause = job.pauseAfterCurrent === true || job.pausedByUser === true
|
const holdPause = job.pauseAfterCurrent === true || job.pausedByUser === true
|
||||||
if (liveFailed) {
|
// Live finished, failed, or vanished — clear the studio slot so the next waiting job can start.
|
||||||
|
if (live && (live.status === 'error' || live.status === 'cancelled')) {
|
||||||
job.status = live.status === 'cancelled' ? 'cancelled' : 'error'
|
job.status = live.status === 'cancelled' ? 'cancelled' : 'error'
|
||||||
job.liveJobId = undefined
|
job.liveJobId = undefined
|
||||||
job.lastError = live.error || job.lastError || 'Job failed'
|
job.lastError = live.error || job.lastError || 'Job failed'
|
||||||
@@ -936,31 +934,56 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
|||||||
job.updatedAt = Date.now()
|
job.updatedAt = Date.now()
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
// No live job and no fresh pending — orphaned running row.
|
// live complete, or live missing after a restart — job is done.
|
||||||
if (job.shotQueueId) {
|
if (job.shotQueueId) {
|
||||||
job.status = 'held'
|
const queue = getShotQueue(job.ownerKey, job.shotQueueId)
|
||||||
job.liveJobId = undefined
|
const pending = queue ? pendingSegmentCount(queue) : 0
|
||||||
if (holdPause) {
|
if (pending > 0) {
|
||||||
job.pausedByUser = true
|
job.status = 'held'
|
||||||
job.pauseAfterCurrent = false
|
job.liveJobId = undefined
|
||||||
job.holdForCutIn = false
|
if (holdPause) {
|
||||||
|
job.pausedByUser = true
|
||||||
|
job.pauseAfterCurrent = false
|
||||||
|
job.holdForCutIn = false
|
||||||
|
} else {
|
||||||
|
job.holdForCutIn = true
|
||||||
|
job.resumeAutoRun = job.payload.queueAutoRun === true || job.resumeAutoRun === true
|
||||||
|
}
|
||||||
|
job.updatedAt = Date.now()
|
||||||
} else {
|
} else {
|
||||||
job.holdForCutIn = true
|
job.status = 'complete'
|
||||||
job.resumeAutoRun = job.payload.queueAutoRun === true || job.resumeAutoRun === true
|
job.liveJobId = undefined
|
||||||
|
job.pauseAfterCurrent = false
|
||||||
|
job.pausedByUser = false
|
||||||
|
job.holdForCutIn = false
|
||||||
|
job.cutIn = false
|
||||||
|
job.updatedAt = Date.now()
|
||||||
}
|
}
|
||||||
job.updatedAt = Date.now()
|
|
||||||
} else {
|
} else {
|
||||||
// No shot queue (image/music): never park as held — that blocked waiting jobs.
|
|
||||||
job.status = 'complete'
|
job.status = 'complete'
|
||||||
job.liveJobId = undefined
|
job.liveJobId = undefined
|
||||||
job.pauseAfterCurrent = false
|
job.pauseAfterCurrent = false
|
||||||
job.pausedByUser = false
|
job.pausedByUser = false
|
||||||
job.holdForCutIn = false
|
job.holdForCutIn = false
|
||||||
|
job.cutIn = false
|
||||||
job.updatedAt = Date.now()
|
job.updatedAt = Date.now()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async function settleAttachedTerminalLives() {
|
||||||
|
// If a live job already finished but settle did not clear the studio row, finish it
|
||||||
|
// before deciding the GPU is free. Without this, a missed settle blocks forever.
|
||||||
|
for (const live of listJobs()) {
|
||||||
|
if (live.status !== 'complete' && live.status !== 'error' && live.status !== 'cancelled') continue
|
||||||
|
const owner = live.library?.ownerKey
|
||||||
|
if (!owner) continue
|
||||||
|
const attached = readJobs(owner).some(job => job.status === 'running' && job.liveJobId === live.id)
|
||||||
|
if (!attached) continue
|
||||||
|
await onLiveVideoSettled(live).catch(() => null)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
async function dispatchStudioQueue() {
|
async function dispatchStudioQueue() {
|
||||||
sweepStaleLiveJobs()
|
sweepStaleLiveJobs()
|
||||||
// Repair every owner first. videoJobsBusy() looks at all owners — fixing only the
|
// Repair every owner first. videoJobsBusy() looks at all owners — fixing only the
|
||||||
@@ -979,6 +1002,7 @@ async function dispatchStudioQueue() {
|
|||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
await settleAttachedTerminalLives()
|
||||||
if (await videoJobsBusy()) {
|
if (await videoJobsBusy()) {
|
||||||
if (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting'))) {
|
if (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting'))) {
|
||||||
scheduleKickRetry()
|
scheduleKickRetry()
|
||||||
@@ -1557,7 +1581,10 @@ function remainingStudioShots(job: Job) {
|
|||||||
if (!library) return 0
|
if (!library) return 0
|
||||||
if (library.queueId) {
|
if (library.queueId) {
|
||||||
const queue = getShotQueue(library.ownerKey, library.queueId)
|
const queue = getShotQueue(library.ownerKey, library.queueId)
|
||||||
if (queue) return pendingSegmentCount(queue)
|
// Missing shot-queue file = nothing left to run. Do not fall through to chain math
|
||||||
|
// or a finished single-shot stays "held/running" forever.
|
||||||
|
if (!queue) return 0
|
||||||
|
return pendingSegmentCount(queue)
|
||||||
}
|
}
|
||||||
// Image/music are one studio slot (multi-pass runs inside the live job). Parking them
|
// Image/music are one studio slot (multi-pass runs inside the live job). Parking them
|
||||||
// as held without a shot-queue id blocked every waiting item behind them.
|
// as held without a shot-queue id blocked every waiting item behind them.
|
||||||
@@ -1565,6 +1592,33 @@ function remainingStudioShots(job: Job) {
|
|||||||
return Math.max(0, (library.chainTotal || 1) - (library.chainStep || 1))
|
return Math.max(0, (library.chainTotal || 1) - (library.chainStep || 1))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function findStudioRowForLive(jobs: StudioJob[], job: Job) {
|
||||||
|
const byLive = jobs.find(item => item.liveJobId === job.id)
|
||||||
|
if (byLive) return byLive
|
||||||
|
if (job.library?.queueId) {
|
||||||
|
const byQueue = jobs.find(item => (
|
||||||
|
item.shotQueueId === job.library?.queueId
|
||||||
|
&& item.status === 'running'
|
||||||
|
))
|
||||||
|
if (byQueue) return byQueue
|
||||||
|
}
|
||||||
|
// Last resort: sole running row for this owner (liveJobId lost after a restart).
|
||||||
|
const running = jobs.filter(item => item.status === 'running')
|
||||||
|
if (running.length === 1) return running[0]
|
||||||
|
return undefined
|
||||||
|
}
|
||||||
|
|
||||||
|
function clearStudioRowSlot(row: StudioJob, status: StudioJobStatus, error?: string) {
|
||||||
|
row.status = status
|
||||||
|
row.liveJobId = undefined
|
||||||
|
row.cutIn = false
|
||||||
|
row.holdForCutIn = false
|
||||||
|
row.pauseAfterCurrent = false
|
||||||
|
row.pausedByUser = false
|
||||||
|
row.lastError = error
|
||||||
|
row.updatedAt = Date.now()
|
||||||
|
}
|
||||||
|
|
||||||
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) {
|
||||||
@@ -1574,31 +1628,26 @@ export async function onLiveVideoSettled(job: Job) {
|
|||||||
const remaining = remainingStudioShots(job)
|
const remaining = remainingStudioShots(job)
|
||||||
const wakeFail = job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
|
const wakeFail = job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
|
||||||
const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail
|
const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail
|
||||||
const snapshot = readStore(owner)
|
|
||||||
const rowNow = snapshot.jobs.find(item => item.liveJobId === job.id)
|
await mutateStore(owner, (store) => {
|
||||||
|| (job.library?.queueId
|
const row = findStudioRowForLive(store.jobs, job)
|
||||||
? snapshot.jobs.find(item => (
|
if (!row || row.status === 'cancelled') {
|
||||||
item.status === 'running'
|
syncPausedFlag(store)
|
||||||
&& item.shotQueueId === job.library?.queueId
|
return
|
||||||
))
|
}
|
||||||
: undefined)
|
|
||||||
const userPause = !failed && rowNow?.holdForCutIn !== true && (
|
const userPause = !failed && row.holdForCutIn !== true && (
|
||||||
rowNow?.pauseAfterCurrent === true
|
row.pauseAfterCurrent === true
|
||||||
|| rowNow?.pausedByUser === true
|
|| row.pausedByUser === true
|
||||||
|| job.library?.stopAfterCurrent === true
|
|| job.library?.stopAfterCurrent === true
|
||||||
)
|
)
|
||||||
const cutInHold = !failed && remaining > 0 && rowNow?.holdForCutIn === true && !userPause
|
const cutInHold = !failed && remaining > 0 && row.holdForCutIn === true && !userPause
|
||||||
const interrupted = remaining > 0 && (userPause || cutInHold)
|
|
||||||
if (interrupted && job.status === 'running') {
|
if (interruptedRunning(job, remaining, userPause, cutInHold)) {
|
||||||
job.status = 'complete'
|
job.status = 'complete'
|
||||||
}
|
}
|
||||||
if (userPause && remaining > 0) {
|
|
||||||
await mutateStore(owner, (store) => {
|
if (userPause && remaining > 0) {
|
||||||
const row = store.jobs.find(item => item.liveJobId === job.id)
|
|
||||||
if (!row) {
|
|
||||||
syncPausedFlag(store)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
row.status = 'held'
|
row.status = 'held'
|
||||||
row.liveJobId = undefined
|
row.liveJobId = undefined
|
||||||
row.pauseAfterCurrent = false
|
row.pauseAfterCurrent = false
|
||||||
@@ -1607,27 +1656,18 @@ export async function onLiveVideoSettled(job: Job) {
|
|||||||
row.resumeAutoRun = row.payload.queueAutoRun === true || row.resumeAutoRun === true
|
row.resumeAutoRun = row.payload.queueAutoRun === true || row.resumeAutoRun === true
|
||||||
row.updatedAt = Date.now()
|
row.updatedAt = Date.now()
|
||||||
syncPausedFlag(store)
|
syncPausedFlag(store)
|
||||||
})
|
return
|
||||||
} else if (cutInHold) {
|
}
|
||||||
await mutate(owner, (jobs) => {
|
if (cutInHold) {
|
||||||
const row = jobs.find(item => item.liveJobId === job.id)
|
|
||||||
if (!row) return
|
|
||||||
row.status = 'held'
|
row.status = 'held'
|
||||||
row.liveJobId = undefined
|
row.liveJobId = undefined
|
||||||
row.holdForCutIn = true
|
row.holdForCutIn = true
|
||||||
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()
|
||||||
})
|
syncPausedFlag(store)
|
||||||
} else if (wakeFail) {
|
return
|
||||||
// GPU/host blip mid-chain — put back on waiting so wake/Start now can pick it up.
|
}
|
||||||
// Never park as user-paused held (that made Start now refuse the job).
|
if (wakeFail) {
|
||||||
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 || row.status === 'cancelled') {
|
|
||||||
syncPausedFlag(store)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
row.status = 'waiting'
|
row.status = 'waiting'
|
||||||
row.liveJobId = undefined
|
row.liveJobId = undefined
|
||||||
row.cutIn = true
|
row.cutIn = true
|
||||||
@@ -1638,15 +1678,9 @@ export async function onLiveVideoSettled(job: Job) {
|
|||||||
row.lastError = job.error
|
row.lastError = job.error
|
||||||
row.updatedAt = Date.now()
|
row.updatedAt = Date.now()
|
||||||
syncPausedFlag(store)
|
syncPausedFlag(store)
|
||||||
})
|
return
|
||||||
} else if (!failed && remaining > 0) {
|
}
|
||||||
await mutateStore(owner, (store) => {
|
if (!failed && remaining > 0) {
|
||||||
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 || row.status === 'cancelled') {
|
|
||||||
syncPausedFlag(store)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
const auto = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
|
const auto = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
|
||||||
row.status = 'held'
|
row.status = 'held'
|
||||||
row.liveJobId = undefined
|
row.liveJobId = undefined
|
||||||
@@ -1657,39 +1691,26 @@ export async function onLiveVideoSettled(job: Job) {
|
|||||||
row.lastError = undefined
|
row.lastError = undefined
|
||||||
row.updatedAt = Date.now()
|
row.updatedAt = Date.now()
|
||||||
syncPausedFlag(store)
|
syncPausedFlag(store)
|
||||||
})
|
return
|
||||||
} else if (userPause) {
|
}
|
||||||
await mutateStore(owner, (store) => {
|
|
||||||
const row = store.jobs.find(item => item.liveJobId === job.id)
|
// Normal finish: clear Generating and free the queue for the next waiting job.
|
||||||
if (!row) {
|
const status: StudioJobStatus = job.status === 'cancelled'
|
||||||
syncPausedFlag(store)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
row.status = job.status === 'cancelled'
|
|
||||||
? 'cancelled'
|
|
||||||
: job.status === 'complete'
|
|
||||||
? 'complete'
|
|
||||||
: 'error'
|
|
||||||
row.liveJobId = undefined
|
|
||||||
row.cutIn = false
|
|
||||||
row.holdForCutIn = false
|
|
||||||
row.pauseAfterCurrent = false
|
|
||||||
row.pausedByUser = false
|
|
||||||
row.lastError = job.error
|
|
||||||
row.updatedAt = Date.now()
|
|
||||||
syncPausedFlag(store)
|
|
||||||
})
|
|
||||||
} else {
|
|
||||||
const status = job.status === 'cancelled'
|
|
||||||
? 'cancelled'
|
? 'cancelled'
|
||||||
: job.status === 'complete'
|
: failed
|
||||||
? 'complete'
|
? 'error'
|
||||||
: 'error'
|
: 'complete'
|
||||||
await markStudioSettled(owner, job.id, status, job.error)
|
clearStudioRowSlot(row, status, failed ? job.error : undefined)
|
||||||
}
|
syncPausedFlag(store)
|
||||||
|
})
|
||||||
|
|
||||||
await kickStudioQueue()
|
await kickStudioQueue()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function interruptedRunning(job: Job, remaining: number, userPause: boolean, cutInHold: boolean) {
|
||||||
|
return remaining > 0 && (userPause || cutInHold) && job.status === 'running'
|
||||||
|
}
|
||||||
|
|
||||||
function applyLiveJobPause(live: Job | undefined, pause: boolean, restoreAutoRun?: boolean) {
|
function applyLiveJobPause(live: Job | undefined, pause: boolean, restoreAutoRun?: boolean) {
|
||||||
if (!live?.library) return
|
if (!live?.library) return
|
||||||
live.library.stopAfterCurrent = pause
|
live.library.stopAfterCurrent = pause
|
||||||
|
|||||||
@@ -363,6 +363,12 @@ export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Pr
|
|||||||
})
|
})
|
||||||
deletePendingJob(job.id)
|
deletePendingJob(job.id)
|
||||||
job.saving = false
|
job.saving = false
|
||||||
|
// Persist path: clear the studio "Generating" row as soon as the library save
|
||||||
|
// finishes. Waiting for runGeneration's finally used to leave jobs stuck in the queue.
|
||||||
|
if (persist && job.library?.ownerKey) {
|
||||||
|
const { onLiveVideoSettled } = await import('~/server/utils/studioQueue')
|
||||||
|
await onLiveVideoSettled(job).catch(() => null)
|
||||||
|
}
|
||||||
} catch (saveError) {
|
} catch (saveError) {
|
||||||
job.saving = false
|
job.saving = false
|
||||||
removeExtendTemp(job.library?.extendTmpDir)
|
removeExtendTemp(job.library?.extendTmpDir)
|
||||||
|
|||||||
Reference in New Issue
Block a user