From c88a8494d50cf9f95b44ef431fc2fbe8db0d575f Mon Sep 17 00:00:00 2001 From: Towsty Date: Sun, 30 Aug 2026 12:47:59 -0500 Subject: [PATCH] Stop cancelled video jobs from coming back after delete. A leftover process-remaining wake was flipping the row back to running or error, so Remove looked like it did nothing. Co-authored-by: Cursor --- pages/index.vue | 4 +-- pages/queue.vue | 15 ++++++---- server/utils/shotQueue.ts | 27 ++++++++++++++++-- server/utils/studioQueue.ts | 56 +++++++++++++++++++++++++------------ server/utils/videoChain.ts | 7 +++++ 5 files changed, 80 insertions(+), 29 deletions(-) diff --git a/pages/index.vue b/pages/index.vue index 2b22d0a..29c643b 100644 --- a/pages/index.vue +++ b/pages/index.vue @@ -1108,12 +1108,12 @@ {{ job.cutIn ? 'Next' : (canCutInStudioJob ? 'Run next' : 'Start now') }} diff --git a/pages/queue.vue b/pages/queue.vue index 226c9aa..87667ad 100644 --- a/pages/queue.vue +++ b/pages/queue.vue @@ -537,7 +537,7 @@ const pauseDisabledReason = computed(() => { return '' }) const clearableCount = computed(() => ( - jobs.value.filter(job => job.status === 'waiting' || job.status === 'held').length + jobs.value.filter(job => job.status === 'waiting' || job.status === 'held' || job.status === 'error').length + orphans.value.filter(queue => queue.status !== 'running').length )) @@ -776,7 +776,7 @@ function jobThumb(job: StudioJobRow) { } function canRemoveJob(job: StudioJobRow) { - return job.status === 'waiting' || job.status === 'held' || job.status === 'running' + return job.status === 'waiting' || job.status === 'held' || job.status === 'running' || job.status === 'error' } function promptSnippet(text?: string) { @@ -841,9 +841,12 @@ async function confirmRemove() { } async function removeJob(job: StudioJobRow) { - await $fetch(`/api/studio-queue/${job.id}`, { method: 'DELETE' }).catch((err: any) => { - error.value = err?.data?.statusMessage || err?.message || 'Could not remove that job' - }) + try { + await $fetch(`/api/studio-queue/${job.id}`, { method: 'DELETE' }) + } catch (err: any) { + error.value = err?.data?.statusMessage || err?.statusMessage || err?.message || 'Could not remove that job' + throw err + } await loadQueues() } @@ -1002,7 +1005,7 @@ async function clearAll() { if (!confirm(`Clear waiting and paused jobs from the queue? Saved clips stay in the library.${extra}`)) return error.value = '' try { - for (const job of jobs.value.filter(item => item.status === 'waiting' || item.status === 'held')) { + for (const job of jobs.value.filter(item => item.status === 'waiting' || item.status === 'held' || item.status === 'error')) { await $fetch(`/api/studio-queue/${job.id}`, { method: 'DELETE' }).catch(() => null) } await $fetch('/api/queues', { method: 'DELETE' }) diff --git a/server/utils/shotQueue.ts b/server/utils/shotQueue.ts index e457288..d54aeaf 100644 --- a/server/utils/shotQueue.ts +++ b/server/utils/shotQueue.ts @@ -58,6 +58,7 @@ export interface ShotQueue { const writeChains = new Map>() const activeJobs = new Map() +const cancelledQueueIds = new Set() function libraryRoot() { const config = useRuntimeConfig() @@ -135,9 +136,23 @@ export function queueIsProcessing(id: string) { return activeJobs.has(id) } +export function markShotQueueCancelled(id: string) { + if (!id) return + cancelledQueueIds.add(id) + activeJobs.delete(id) +} + +export function isShotQueueCancelled(id: string) { + return Boolean(id) && cancelledQueueIds.has(id) +} + export function setQueueJob(id: string, jobId: string | null) { - if (jobId) activeJobs.set(id, jobId) - else activeJobs.delete(id) + if (jobId) { + if (cancelledQueueIds.has(id)) return + activeJobs.set(id, jobId) + } else { + activeJobs.delete(id) + } } export function getQueueJobId(id: string) { @@ -229,6 +244,9 @@ export async function createShotQueue(params: { } export async function updateShotQueue(owner: string, id: string, patch: (queue: ShotQueue) => void) { + if (cancelledQueueIds.has(id)) { + throw createError({ statusCode: 410, statusMessage: 'This batch was removed' }) + } return mutate(owner, (queues) => { const queue = queues.find(item => item.id === id) if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' }) @@ -314,6 +332,9 @@ export async function syncQueueFromJob(job: Job, clipId?: string) { } export async function prepareQueueBurst(owner: string, id: string, count: number | 'all') { + if (cancelledQueueIds.has(id)) { + throw createError({ statusCode: 410, statusMessage: 'This batch was removed' }) + } return mutate(owner, (queues) => { const queue = queues.find(item => item.id === id) if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' }) @@ -418,7 +439,7 @@ export async function deleteShotQueue(owner: string, id: string, options: { forc if (!options.force && activeJobs.has(id)) { throw createError({ statusCode: 409, statusMessage: 'Stop this batch before deleting it' }) } - activeJobs.delete(id) + markShotQueueCancelled(id) return mutate(owner, (queues) => { const index = queues.findIndex(item => item.id === id) if (index < 0) throw createError({ statusCode: 404, statusMessage: 'Queue not found' }) diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts index 54f49a0..3ea18ae 100644 --- a/server/utils/studioQueue.ts +++ b/server/utils/studioQueue.ts @@ -346,20 +346,28 @@ export async function patchStudioJob(owner: string, id: string, patch: (job: Stu }) } +function canCancelStudioStatus(status: StudioJobStatus) { + return status === 'waiting' || status === 'held' || status === 'running' || status === 'error' +} + export async function cancelStudioJob(owner: string, id: string) { const before = readJobs(owner).find(item => item.id === id) if (!before) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' }) - if (before.status !== 'waiting' && before.status !== 'held' && before.status !== 'running') { + if (!canCancelStudioStatus(before.status)) { throw createError({ statusCode: 409, statusMessage: 'That job is no longer in the queue' }) } - const stopLive = before.status === 'running' const liveJobId = before.liveJobId const shotQueueId = before.shotQueueId + if (shotQueueId) { + const { markShotQueueCancelled, setQueueJob } = await import('~/server/utils/shotQueue') + markShotQueueCancelled(shotQueueId) + setQueueJob(shotQueueId, null) + } const result = await mutateStore(owner, (store) => { const job = store.jobs.find(item => item.id === id) if (!job) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' }) - if (job.status !== 'waiting' && job.status !== 'held' && job.status !== 'running') { + if (!canCancelStudioStatus(job.status)) { throw createError({ statusCode: 409, statusMessage: 'That job is no longer in the queue' }) } job.status = 'cancelled' @@ -368,6 +376,8 @@ export async function cancelStudioJob(owner: string, id: string) { job.pauseAfterCurrent = false job.pausedByUser = false job.liveJobId = undefined + job.shotQueueId = undefined + job.lastError = undefined job.updatedAt = Date.now() const stillCutIn = store.jobs.some(item => item.status === 'waiting' && item.cutIn) if (!stillCutIn) { @@ -380,13 +390,10 @@ export async function cancelStudioJob(owner: string, id: string) { syncPausedFlag(store) return structuredClone(job) }) - if (stopLive) { - await stopLiveGeneration(liveJobId, shotQueueId) - } - if (result.shotQueueId) { - const { deleteShotQueue, setQueueJob } = await import('~/server/utils/shotQueue') - setQueueJob(result.shotQueueId, null) - await deleteShotQueue(owner, result.shotQueueId, { force: true }).catch(() => null) + await stopLiveGeneration(liveJobId, shotQueueId) + if (shotQueueId) { + const { deleteShotQueue } = await import('~/server/utils/shotQueue') + await deleteShotQueue(owner, shotQueueId, { force: true }).catch(() => null) } const jobs = readJobs(owner) if (!jobs.some(item => item.status === 'waiting' && item.cutIn)) { @@ -425,12 +432,13 @@ async function stopLiveGeneration(liveJobId?: string, shotQueueId?: string) { } export async function cancelOrphanQueue(owner: string, queueId: string) { - const { getShotQueue, getQueueJobId, queueIsProcessing, deleteShotQueue } = await import('~/server/utils/shotQueue') + const { getShotQueue, getQueueJobId, queueIsProcessing, deleteShotQueue, markShotQueueCancelled } = await import('~/server/utils/shotQueue') const queue = getShotQueue(owner, queueId) if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' }) + markShotQueueCancelled(queueId) const generating = queue.status === 'running' || queueIsProcessing(queueId) const liveJobId = getQueueJobId(queueId) || queue.currentJobId - if (generating) { + if (generating || liveJobId) { await stopLiveGeneration(liveJobId, queueId) } await deleteShotQueue(owner, queueId, { force: true }).catch(() => null) @@ -650,33 +658,43 @@ async function dispatchStudioQueue() { } async function resumeHeldStudioJob(item: StudioJob) { - if (!item.shotQueueId) { + const latest = readJobs(item.ownerKey).find(job => job.id === item.id) + if (!latest || latest.status === 'cancelled') return false + if (!latest.shotQueueId) { await patchStudioJob(item.ownerKey, item.id, (job) => { + if (job.status === 'cancelled') return job.status = 'complete' job.holdForCutIn = false }) return false } + const { isShotQueueCancelled } = await import('~/server/utils/shotQueue') + if (isShotQueueCancelled(latest.shotQueueId)) return false await patchStudioJob(item.ownerKey, item.id, (job) => { + if (job.status === 'cancelled') return job.status = 'running' job.lastError = undefined }) + if (readJobs(item.ownerKey).find(job => job.id === item.id)?.status === 'cancelled') return false const { startQueueBurst } = await import('~/server/utils/videoChain') - const count = item.resumeAutoRun ? 'all' as const : 1 + const count = latest.resumeAutoRun ? 'all' as const : 1 try { - const started = await startQueueBurst(item.ownerKey, item.shotQueueId, count) + const started = await startQueueBurst(item.ownerKey, latest.shotQueueId, count) await patchStudioJob(item.ownerKey, item.id, (job) => { + if (job.status === 'cancelled') return job.status = 'running' job.liveJobId = started.jobId job.holdForCutIn = false job.cutIn = false }) - return true + return readJobs(item.ownerKey).find(job => job.id === item.id)?.status === 'running' } catch (error) { const message = error instanceof Error ? error.message : String(error) if (/already processing/i.test(message)) return true + if (/removed|not found/i.test(message)) return false 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) => { + if (job.status === 'cancelled') return job.status = transient ? 'held' : 'error' job.lastError = message job.holdForCutIn = transient @@ -1116,7 +1134,7 @@ export async function onLiveVideoSettled(job: Job) { 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) { + if (!row || row.status === 'cancelled') { syncPausedFlag(store) return } @@ -1321,9 +1339,11 @@ export async function toggleStudioQueuePause(owner: string, options: { } export async function onStudioQueueBurstStarted(owner: string, queueId: string, liveJobId: string) { + const { isShotQueueCancelled } = await import('~/server/utils/shotQueue') + if (isShotQueueCancelled(queueId)) return await mutateStore(owner, (store) => { const row = store.jobs.find(job => job.shotQueueId === queueId) - if (!row) { + if (!row || row.status === 'cancelled') { syncPausedFlag(store) return } diff --git a/server/utils/videoChain.ts b/server/utils/videoChain.ts index d5a865d..7140b79 100644 --- a/server/utils/videoChain.ts +++ b/server/utils/videoChain.ts @@ -13,6 +13,7 @@ import { onLiveVideoSettled, onStudioQueueBurstStarted } from '~/server/utils/st import { finishQueueBurst, getShotQueue, + isShotQueueCancelled, lastCompletedIndex, liveSegmentPrompt, prepareQueueBurst, @@ -515,6 +516,9 @@ export async function runGeneration(job: Job, params: VideoChainParams) { } export async function startQueueBurst(owner: string, queueId: string, count: number | 'all') { + if (isShotQueueCancelled(queueId)) { + throw createError({ statusCode: 410, statusMessage: 'This batch was removed' }) + } await ensureComfyReady((status) => { console.log(JSON.stringify({ src: 'video-chain', @@ -524,6 +528,9 @@ export async function startQueueBurst(owner: string, queueId: string, count: num message: status.message })) }) + if (isShotQueueCancelled(queueId)) { + throw createError({ statusCode: 410, statusMessage: 'This batch was removed' }) + } const { queue, count: n } = await prepareQueueBurst(owner, queueId, count) const lastIndex = lastCompletedIndex(queue) const clipId = queue.currentClipId