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 <cursoragent@cursor.com>
This commit is contained in:
+2
-2
@@ -1108,12 +1108,12 @@
|
|||||||
{{ job.cutIn ? 'Next' : (canCutInStudioJob ? 'Run next' : 'Start now') }}
|
{{ job.cutIn ? 'Next' : (canCutInStudioJob ? 'Run next' : 'Start now') }}
|
||||||
</button>
|
</button>
|
||||||
<button
|
<button
|
||||||
v-if="job.status === 'waiting' || job.status === 'held'"
|
v-if="job.status === 'waiting' || job.status === 'held' || job.status === 'running' || job.status === 'error'"
|
||||||
type="button"
|
type="button"
|
||||||
class="shrink-0 text-xs text-zinc-400 hover:text-red-200"
|
class="shrink-0 text-xs text-zinc-400 hover:text-red-200"
|
||||||
@click="dismissStudioJob(job.id)"
|
@click="dismissStudioJob(job.id)"
|
||||||
>
|
>
|
||||||
Remove
|
{{ job.status === 'running' ? 'Cancel' : 'Remove' }}
|
||||||
</button>
|
</button>
|
||||||
</li>
|
</li>
|
||||||
</ul>
|
</ul>
|
||||||
|
|||||||
+9
-6
@@ -537,7 +537,7 @@ const pauseDisabledReason = computed(() => {
|
|||||||
return ''
|
return ''
|
||||||
})
|
})
|
||||||
const clearableCount = computed(() => (
|
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
|
+ orphans.value.filter(queue => queue.status !== 'running').length
|
||||||
))
|
))
|
||||||
|
|
||||||
@@ -776,7 +776,7 @@ function jobThumb(job: StudioJobRow) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
function canRemoveJob(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) {
|
function promptSnippet(text?: string) {
|
||||||
@@ -841,9 +841,12 @@ async function confirmRemove() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async function removeJob(job: StudioJobRow) {
|
async function removeJob(job: StudioJobRow) {
|
||||||
await $fetch(`/api/studio-queue/${job.id}`, { method: 'DELETE' }).catch((err: any) => {
|
try {
|
||||||
error.value = err?.data?.statusMessage || err?.message || 'Could not remove that job'
|
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()
|
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
|
if (!confirm(`Clear waiting and paused jobs from the queue? Saved clips stay in the library.${extra}`)) return
|
||||||
error.value = ''
|
error.value = ''
|
||||||
try {
|
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/studio-queue/${job.id}`, { method: 'DELETE' }).catch(() => null)
|
||||||
}
|
}
|
||||||
await $fetch('/api/queues', { method: 'DELETE' })
|
await $fetch('/api/queues', { method: 'DELETE' })
|
||||||
|
|||||||
@@ -58,6 +58,7 @@ export interface ShotQueue {
|
|||||||
|
|
||||||
const writeChains = new Map<string, Promise<unknown>>()
|
const writeChains = new Map<string, Promise<unknown>>()
|
||||||
const activeJobs = new Map<string, string>()
|
const activeJobs = new Map<string, string>()
|
||||||
|
const cancelledQueueIds = new Set<string>()
|
||||||
|
|
||||||
function libraryRoot() {
|
function libraryRoot() {
|
||||||
const config = useRuntimeConfig()
|
const config = useRuntimeConfig()
|
||||||
@@ -135,9 +136,23 @@ export function queueIsProcessing(id: string) {
|
|||||||
return activeJobs.has(id)
|
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) {
|
export function setQueueJob(id: string, jobId: string | null) {
|
||||||
if (jobId) activeJobs.set(id, jobId)
|
if (jobId) {
|
||||||
else activeJobs.delete(id)
|
if (cancelledQueueIds.has(id)) return
|
||||||
|
activeJobs.set(id, jobId)
|
||||||
|
} else {
|
||||||
|
activeJobs.delete(id)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export function getQueueJobId(id: string) {
|
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) {
|
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) => {
|
return mutate(owner, (queues) => {
|
||||||
const queue = queues.find(item => item.id === id)
|
const queue = queues.find(item => item.id === id)
|
||||||
if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
|
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') {
|
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) => {
|
return mutate(owner, (queues) => {
|
||||||
const queue = queues.find(item => item.id === id)
|
const queue = queues.find(item => item.id === id)
|
||||||
if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
|
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)) {
|
if (!options.force && activeJobs.has(id)) {
|
||||||
throw createError({ statusCode: 409, statusMessage: 'Stop this batch before deleting it' })
|
throw createError({ statusCode: 409, statusMessage: 'Stop this batch before deleting it' })
|
||||||
}
|
}
|
||||||
activeJobs.delete(id)
|
markShotQueueCancelled(id)
|
||||||
return mutate(owner, (queues) => {
|
return mutate(owner, (queues) => {
|
||||||
const index = queues.findIndex(item => item.id === id)
|
const index = queues.findIndex(item => item.id === id)
|
||||||
if (index < 0) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
|
if (index < 0) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
|
||||||
|
|||||||
+38
-18
@@ -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) {
|
export async function cancelStudioJob(owner: string, id: string) {
|
||||||
const before = readJobs(owner).find(item => item.id === id)
|
const before = readJobs(owner).find(item => item.id === id)
|
||||||
if (!before) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' })
|
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' })
|
throw createError({ statusCode: 409, statusMessage: 'That job is no longer in the queue' })
|
||||||
}
|
}
|
||||||
const stopLive = before.status === 'running'
|
|
||||||
const liveJobId = before.liveJobId
|
const liveJobId = before.liveJobId
|
||||||
const shotQueueId = before.shotQueueId
|
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 result = await mutateStore(owner, (store) => {
|
||||||
const job = store.jobs.find(item => item.id === id)
|
const job = store.jobs.find(item => item.id === id)
|
||||||
if (!job) throw createError({ statusCode: 404, statusMessage: 'Queued job not found' })
|
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' })
|
throw createError({ statusCode: 409, statusMessage: 'That job is no longer in the queue' })
|
||||||
}
|
}
|
||||||
job.status = 'cancelled'
|
job.status = 'cancelled'
|
||||||
@@ -368,6 +376,8 @@ export async function cancelStudioJob(owner: string, id: string) {
|
|||||||
job.pauseAfterCurrent = false
|
job.pauseAfterCurrent = false
|
||||||
job.pausedByUser = false
|
job.pausedByUser = false
|
||||||
job.liveJobId = undefined
|
job.liveJobId = undefined
|
||||||
|
job.shotQueueId = undefined
|
||||||
|
job.lastError = undefined
|
||||||
job.updatedAt = Date.now()
|
job.updatedAt = Date.now()
|
||||||
const stillCutIn = store.jobs.some(item => item.status === 'waiting' && item.cutIn)
|
const stillCutIn = store.jobs.some(item => item.status === 'waiting' && item.cutIn)
|
||||||
if (!stillCutIn) {
|
if (!stillCutIn) {
|
||||||
@@ -380,13 +390,10 @@ export async function cancelStudioJob(owner: string, id: string) {
|
|||||||
syncPausedFlag(store)
|
syncPausedFlag(store)
|
||||||
return structuredClone(job)
|
return structuredClone(job)
|
||||||
})
|
})
|
||||||
if (stopLive) {
|
await stopLiveGeneration(liveJobId, shotQueueId)
|
||||||
await stopLiveGeneration(liveJobId, shotQueueId)
|
if (shotQueueId) {
|
||||||
}
|
const { deleteShotQueue } = await import('~/server/utils/shotQueue')
|
||||||
if (result.shotQueueId) {
|
await deleteShotQueue(owner, shotQueueId, { force: true }).catch(() => null)
|
||||||
const { deleteShotQueue, setQueueJob } = await import('~/server/utils/shotQueue')
|
|
||||||
setQueueJob(result.shotQueueId, null)
|
|
||||||
await deleteShotQueue(owner, result.shotQueueId, { force: true }).catch(() => null)
|
|
||||||
}
|
}
|
||||||
const jobs = readJobs(owner)
|
const jobs = readJobs(owner)
|
||||||
if (!jobs.some(item => item.status === 'waiting' && item.cutIn)) {
|
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) {
|
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)
|
const queue = getShotQueue(owner, queueId)
|
||||||
if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
|
if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
|
||||||
|
markShotQueueCancelled(queueId)
|
||||||
const generating = queue.status === 'running' || queueIsProcessing(queueId)
|
const generating = queue.status === 'running' || queueIsProcessing(queueId)
|
||||||
const liveJobId = getQueueJobId(queueId) || queue.currentJobId
|
const liveJobId = getQueueJobId(queueId) || queue.currentJobId
|
||||||
if (generating) {
|
if (generating || liveJobId) {
|
||||||
await stopLiveGeneration(liveJobId, queueId)
|
await stopLiveGeneration(liveJobId, queueId)
|
||||||
}
|
}
|
||||||
await deleteShotQueue(owner, queueId, { force: true }).catch(() => null)
|
await deleteShotQueue(owner, queueId, { force: true }).catch(() => null)
|
||||||
@@ -650,33 +658,43 @@ async function dispatchStudioQueue() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async function resumeHeldStudioJob(item: StudioJob) {
|
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) => {
|
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||||
|
if (job.status === 'cancelled') return
|
||||||
job.status = 'complete'
|
job.status = 'complete'
|
||||||
job.holdForCutIn = false
|
job.holdForCutIn = false
|
||||||
})
|
})
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
const { isShotQueueCancelled } = await import('~/server/utils/shotQueue')
|
||||||
|
if (isShotQueueCancelled(latest.shotQueueId)) return false
|
||||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||||
|
if (job.status === 'cancelled') return
|
||||||
job.status = 'running'
|
job.status = 'running'
|
||||||
job.lastError = undefined
|
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 { startQueueBurst } = await import('~/server/utils/videoChain')
|
||||||
const count = item.resumeAutoRun ? 'all' as const : 1
|
const count = latest.resumeAutoRun ? 'all' as const : 1
|
||||||
try {
|
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) => {
|
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||||
|
if (job.status === 'cancelled') return
|
||||||
job.status = 'running'
|
job.status = 'running'
|
||||||
job.liveJobId = started.jobId
|
job.liveJobId = started.jobId
|
||||||
job.holdForCutIn = false
|
job.holdForCutIn = false
|
||||||
job.cutIn = false
|
job.cutIn = false
|
||||||
})
|
})
|
||||||
return true
|
return readJobs(item.ownerKey).find(job => job.id === item.id)?.status === 'running'
|
||||||
} 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
|
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)
|
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) => {
|
||||||
|
if (job.status === 'cancelled') return
|
||||||
job.status = transient ? 'held' : 'error'
|
job.status = transient ? 'held' : 'error'
|
||||||
job.lastError = message
|
job.lastError = message
|
||||||
job.holdForCutIn = transient
|
job.holdForCutIn = transient
|
||||||
@@ -1116,7 +1134,7 @@ export async function onLiveVideoSettled(job: Job) {
|
|||||||
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)
|
||||||
|| (job.library?.queueId ? store.jobs.find(item => item.shotQueueId === job.library?.queueId) : undefined)
|
|| (job.library?.queueId ? store.jobs.find(item => item.shotQueueId === job.library?.queueId) : undefined)
|
||||||
if (!row) {
|
if (!row || row.status === 'cancelled') {
|
||||||
syncPausedFlag(store)
|
syncPausedFlag(store)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -1321,9 +1339,11 @@ export async function toggleStudioQueuePause(owner: string, options: {
|
|||||||
}
|
}
|
||||||
|
|
||||||
export async function onStudioQueueBurstStarted(owner: string, queueId: string, liveJobId: string) {
|
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) => {
|
await mutateStore(owner, (store) => {
|
||||||
const row = store.jobs.find(job => job.shotQueueId === queueId)
|
const row = store.jobs.find(job => job.shotQueueId === queueId)
|
||||||
if (!row) {
|
if (!row || row.status === 'cancelled') {
|
||||||
syncPausedFlag(store)
|
syncPausedFlag(store)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ import { onLiveVideoSettled, onStudioQueueBurstStarted } from '~/server/utils/st
|
|||||||
import {
|
import {
|
||||||
finishQueueBurst,
|
finishQueueBurst,
|
||||||
getShotQueue,
|
getShotQueue,
|
||||||
|
isShotQueueCancelled,
|
||||||
lastCompletedIndex,
|
lastCompletedIndex,
|
||||||
liveSegmentPrompt,
|
liveSegmentPrompt,
|
||||||
prepareQueueBurst,
|
prepareQueueBurst,
|
||||||
@@ -515,6 +516,9 @@ 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') {
|
||||||
|
if (isShotQueueCancelled(queueId)) {
|
||||||
|
throw createError({ statusCode: 410, statusMessage: 'This batch was removed' })
|
||||||
|
}
|
||||||
await ensureComfyReady((status) => {
|
await ensureComfyReady((status) => {
|
||||||
console.log(JSON.stringify({
|
console.log(JSON.stringify({
|
||||||
src: 'video-chain',
|
src: 'video-chain',
|
||||||
@@ -524,6 +528,9 @@ export async function startQueueBurst(owner: string, queueId: string, count: num
|
|||||||
message: status.message
|
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 { queue, count: n } = await prepareQueueBurst(owner, queueId, count)
|
||||||
const lastIndex = lastCompletedIndex(queue)
|
const lastIndex = lastCompletedIndex(queue)
|
||||||
const clipId = queue.currentClipId
|
const clipId = queue.currentClipId
|
||||||
|
|||||||
Reference in New Issue
Block a user