Pause only the current shot batch so a new Generate can run when Beast is idle, without auto-resuming the held list.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Towsty
2026-08-28 08:42:55 -05:00
co-authored by Cursor
parent 7d783bd97e
commit 8598dab516
5 changed files with 265 additions and 123 deletions
+13 -7
View File
@@ -1,5 +1,5 @@
import { listShotQueues, summarizeQueue } from '~/server/utils/shotQueue'
import { isStudioQueuePaused, listStudioJobs, summarizeStudioJob } from '~/server/utils/studioQueue'
import { listStudioJobs, summarizeStudioJob } from '~/server/utils/studioQueue'
export default defineEventHandler((event) => {
const { owner } = assertLibraryOwner(event)
@@ -28,14 +28,20 @@ export default defineEventHandler((event) => {
pendingCount: shotSummary?.pendingCount ?? job.shotCount
}
})
const orphans = full
? queues
.filter(queue => !claimed.has(queue.id) && queue.status !== 'complete')
.map(summarizeQueue)
: []
return {
jobs,
orphans: full
? queues
.filter(queue => !claimed.has(queue.id) && queue.status !== 'complete')
.map(summarizeQueue)
: [],
orphans,
waitingCount: jobs.filter(job => job.status === 'waiting').length,
paused: isStudioQueuePaused(owner)
paused: jobs.some(job => job.pausedByUser === true)
|| orphans.some(queue => queue.stopAfterCurrent === true && queue.status !== 'running'),
willPause: jobs.some(job => (
job.pauseAfterCurrent === true
|| (job.status === 'running' && Boolean(job.shotQueueId && byId.get(job.shotQueueId)?.stopAfterCurrent))
))
}
})
+1 -2
View File
@@ -1,6 +1,6 @@
import { listPendingJobs } from '~/server/utils/pending'
import { listOwnersWithQueues, listShotQueues, pendingSegmentCount } from '~/server/utils/shotQueue'
import { isStudioQueuePaused, kickStudioQueue, listOwnersWithStudioQueues, listStudioJobs } from '~/server/utils/studioQueue'
import { kickStudioQueue, listOwnersWithStudioQueues, listStudioJobs } from '~/server/utils/studioQueue'
import { startQueueBurst } from '~/server/utils/videoChain'
export default defineNitroPlugin(() => {
@@ -15,7 +15,6 @@ export default defineNitroPlugin(() => {
}
}
for (const owner of listOwnersWithQueues()) {
if (isStudioQueuePaused(owner)) continue
for (const queue of listShotQueues(owner)) {
if (claimedQueueIds.has(queue.id)) continue
if (!queue.autoRun) continue
+105 -35
View File
@@ -139,8 +139,28 @@ function mutateStore<T>(owner: string, fn: (store: StudioQueueStore) => T): Prom
return run
}
function jobIsUserPaused(job: StudioJob) {
return job.pauseAfterCurrent === true || job.pausedByUser === true
}
function syncPausedFlag(store: StudioQueueStore) {
store.paused = store.jobs.some(jobIsUserPaused)
}
function migrateGlobalPause(store: StudioQueueStore) {
const anyJobFlag = store.jobs.some(jobIsUserPaused)
if (store.paused && !anyJobFlag) {
for (const job of store.jobs) {
if (job.status === 'running') job.pauseAfterCurrent = true
else if (job.status === 'held' && !job.holdForCutIn) job.pausedByUser = true
}
}
syncPausedFlag(store)
}
export function isStudioQueuePaused(owner: string) {
return readStore(owner).paused === true
const store = readStore(owner)
return store.jobs.some(jobIsUserPaused)
}
export function listOwnersWithStudioQueues() {
@@ -196,11 +216,9 @@ export function applyPersistedPauseToJob(job: Job) {
const row = store.jobs.find(item => item.liveJobId === job.id)
|| (library.queueId ? store.jobs.find(item => item.shotQueueId === library.queueId) : undefined)
const queue = library.queueId ? getShotQueue(library.ownerKey, library.queueId) : null
const paused = store.paused === true
|| row?.pauseAfterCurrent === true
const paused = row?.pauseAfterCurrent === true
|| row?.pausedByUser === true
|| queue?.stopAfterCurrent === true
|| queue?.status === 'paused'
if (paused) {
library.stopAfterCurrent = true
library.queueAutoRun = false
@@ -385,7 +403,12 @@ function pickHeld(jobs: StudioJob[]) {
}
function pickWaiting(jobs: StudioJob[]) {
return jobs.find(job => job.status === 'waiting') || null
const waiting = jobs.filter(job => job.status === 'waiting')
if (!waiting.length) return null
if (jobs.some(job => job.pausedByUser === true)) {
return waiting.reduce((latest, job) => job.createdAt >= latest.createdAt ? job : latest)
}
return waiting[0] || null
}
export function kickStudioQueue() {
@@ -404,13 +427,13 @@ function pendingAlive(job: StudioJob) {
))
}
function repairStaleJobs(jobs: StudioJob[], paused: boolean) {
function repairStaleJobs(jobs: StudioJob[]) {
for (const job of jobs) {
if (job.status !== 'running') continue
const live = job.liveJobId ? getJob(job.liveJobId) : undefined
const liveBusy = live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')
if (liveBusy || pendingAlive(job)) continue
const holdPause = paused || job.pauseAfterCurrent === true || job.pausedByUser === true
const holdPause = job.pauseAfterCurrent === true || job.pausedByUser === true
if (job.shotQueueId) {
job.status = 'held'
job.liveJobId = undefined
@@ -439,15 +462,16 @@ async function dispatchStudioQueue() {
if (videoJobsBusy()) return
for (const owner of listOwnersWithStudioQueues()) {
const store = await mutateStore(owner, (current) => {
repairStaleJobs(current.jobs, current.paused === true)
migrateGlobalPause(current)
repairStaleJobs(current.jobs)
const next = pruneDone(current.jobs)
current.jobs.splice(0, current.jobs.length, ...next)
syncPausedFlag(current)
return {
paused: current.paused === true,
jobs: current.jobs.map(item => structuredClone(item))
}
})
if (store.paused) return
const cutIn = pickCutIn(store.jobs)
if (cutIn) {
await startStudioJob(cutIn)
@@ -654,10 +678,10 @@ export async function onLiveVideoSettled(job: Job) {
const rowNow = snapshot.jobs.find(item => item.liveJobId === job.id)
|| snapshot.jobs.find(item => item.status === 'running' && item.shotQueueId && item.shotQueueId === job.library?.queueId)
|| snapshot.jobs.find(item => item.status === 'running' && !job.library?.queueId)
const userPause = !failed && (
const userPause = !failed && rowNow?.holdForCutIn !== true && (
rowNow?.pauseAfterCurrent === true
|| (snapshot.paused && !rowNow?.holdForCutIn)
|| (job.library?.stopAfterCurrent === true && !rowNow?.holdForCutIn)
|| rowNow?.pausedByUser === true
|| job.library?.stopAfterCurrent === true
)
const cutInHold = !failed && remaining > 0 && rowNow?.holdForCutIn === true && !userPause
const interrupted = remaining > 0 && (userPause || cutInHold)
@@ -666,9 +690,11 @@ export async function onLiveVideoSettled(job: Job) {
}
if (userPause && remaining > 0) {
await mutateStore(owner, (store) => {
store.paused = true
const row = store.jobs.find(item => item.liveJobId === job.id) || store.jobs.find(item => item.status === 'running')
if (!row) return
if (!row) {
syncPausedFlag(store)
return
}
row.status = 'held'
row.liveJobId = undefined
row.pauseAfterCurrent = false
@@ -676,6 +702,7 @@ export async function onLiveVideoSettled(job: Job) {
row.holdForCutIn = false
row.resumeAutoRun = row.payload.queueAutoRun === true || row.resumeAutoRun === true
row.updatedAt = Date.now()
syncPausedFlag(store)
})
} else if (cutInHold) {
await mutate(owner, (jobs) => {
@@ -689,9 +716,11 @@ export async function onLiveVideoSettled(job: Job) {
})
} else if (userPause) {
await mutateStore(owner, (store) => {
store.paused = true
const row = store.jobs.find(item => item.liveJobId === job.id) || store.jobs.find(item => item.status === 'running')
if (!row) return
if (!row) {
syncPausedFlag(store)
return
}
row.status = job.status === 'cancelled'
? 'cancelled'
: job.status === 'complete'
@@ -704,6 +733,7 @@ export async function onLiveVideoSettled(job: Job) {
row.pausedByUser = false
row.lastError = job.error
row.updatedAt = Date.now()
syncPausedFlag(store)
})
} else {
const status = job.status === 'cancelled'
@@ -737,11 +767,17 @@ export async function toggleStudioQueuePause(owner: string, options: {
const running = storeNow.jobs.find(job => job.status === 'running')
const target = options.jobId
? storeNow.jobs.find(job => job.id === options.jobId)
: running || storeNow.jobs.find(job => job.pauseAfterCurrent || job.pausedByUser)
: options.queueId
? storeNow.jobs.find(job => job.shotQueueId === options.queueId)
: running || storeNow.jobs.find(job => job.pauseAfterCurrent || job.pausedByUser)
const queueArmed = options.queueId
? listShotQueues(owner).some(queue => queue.id === options.queueId && queue.stopAfterCurrent === true)
: false
const armed = target?.pauseAfterCurrent === true || queueArmed || (storeNow.paused && !running)
: target?.shotQueueId
? listShotQueues(owner).some(queue => queue.id === target.shotQueueId && queue.stopAfterCurrent === true)
: false
const armed = target?.pauseAfterCurrent === true
|| target?.pausedByUser === true
|| queueArmed
const shouldPause = options.resume === true ? false : options.resume === false ? true : !armed
if (!shouldPause) {
@@ -750,55 +786,75 @@ export async function toggleStudioQueuePause(owner: string, options: {
&& Boolean(job.shotQueueId)
&& job.pausedByUser === true
&& (!options.jobId || job.id === options.jobId)
&& (!options.queueId || job.shotQueueId === options.queueId)
))
const restored = await mutateStore(owner, (store) => {
store.paused = false
for (const job of store.jobs) {
if (options.jobId && job.id !== options.jobId) continue
if (options.queueId && job.shotQueueId !== options.queueId) continue
job.pauseAfterCurrent = false
job.pausedByUser = false
}
if (!options.jobId) {
for (const job of store.jobs) {
job.pauseAfterCurrent = false
job.pausedByUser = false
}
}
syncPausedFlag(store)
return structuredClone(store)
})
const rows = options.jobId
? restored.jobs.filter(job => job.id === options.jobId)
: restored.jobs
: options.queueId
? restored.jobs.filter(job => job.shotQueueId === options.queueId)
: restored.jobs
for (const row of rows) {
if (row.liveJobId) applyLiveJobPause(getJob(row.liveJobId), false, row.payload.queueAutoRun === true || row.resumeAutoRun === true)
if (row.shotQueueId) await setShotQueuePause(owner, row.shotQueueId, false).catch(() => null)
}
if (!options.jobId) {
if (!options.jobId && !options.queueId) {
for (const queue of listShotQueues(owner)) {
if (queue.stopAfterCurrent) await setShotQueuePause(owner, queue.id, false).catch(() => null)
}
} else if (options.queueId) {
const rowHasQueue = rows.some(job => job.shotQueueId === options.queueId)
if (!rowHasQueue) await setShotQueuePause(owner, options.queueId, false).catch(() => null)
}
if (heldPaused[0] && !videoJobsBusy()) {
await resumeHeldStudioJob(heldPaused[0]).catch(() => null)
} else {
if (heldPaused.length && videoJobsBusy()) {
await mutate(owner, (jobs) => {
for (const job of jobs) {
if (!heldPaused.some(item => item.id === job.id)) continue
job.holdForCutIn = true
job.pausedByUser = false
job.pauseAfterCurrent = false
job.updatedAt = Date.now()
}
}).catch(() => null)
}
kickStudioQueue()
}
return { paused: false, willPause: false }
}
const updated = await mutateStore(owner, (store) => {
store.paused = true
const row = options.jobId
? store.jobs.find(job => job.id === options.jobId)
: store.jobs.find(job => job.status === 'running')
: options.queueId
? store.jobs.find(job => job.shotQueueId === options.queueId)
: store.jobs.find(job => job.status === 'running')
if (row) {
if (row.status !== 'running' && row.status !== 'held') {
throw createError({ statusCode: 409, statusMessage: 'Nothing is generating on that job' })
}
row.pauseAfterCurrent = row.status === 'running'
if (row.status === 'running') {
row.pauseAfterCurrent = true
} else {
row.pausedByUser = true
row.pauseAfterCurrent = false
row.holdForCutIn = false
}
row.resumeAutoRun = row.payload.queueAutoRun === true || row.resumeAutoRun === true
row.updatedAt = Date.now()
}
syncPausedFlag(store)
return {
paused: store.paused,
job: row ? structuredClone(row) : null
@@ -812,7 +868,13 @@ export async function toggleStudioQueuePause(owner: string, options: {
}
for (const pending of listPendingJobs()) {
if (pending.ownerKey !== owner) continue
if (queueId && pending.queueId && pending.queueId !== queueId) continue
if (queueId) {
if (pending.queueId !== queueId) continue
} else if (updated.job?.liveJobId) {
if (pending.jobId !== updated.job.liveJobId) continue
} else {
continue
}
patchPendingJob(pending.jobId, { stopAfterCurrent: true, queueAutoRun: false })
}
const queued = queueId ? getShotQueue(owner, queueId) : null
@@ -832,19 +894,27 @@ export async function toggleStudioQueuePause(owner: string, options: {
if (!updated.job && !queueId && !anyGenerating) {
throw createError({ statusCode: 409, statusMessage: 'Nothing is generating, so there is nothing to pause after' })
}
return { paused: true, willPause: true, jobId: updated.job?.id, queueId }
return {
paused: true,
willPause: updated.job?.status === 'running' || Boolean(queued?.status === 'running') || anyGenerating,
jobId: updated.job?.id,
queueId
}
}
export async function onStudioQueueBurstStarted(owner: string, queueId: string, liveJobId: string) {
await mutateStore(owner, (store) => {
const row = store.jobs.find(job => job.shotQueueId === queueId)
if (!row) return
store.paused = false
if (!row) {
syncPausedFlag(store)
return
}
row.status = 'running'
row.liveJobId = liveJobId
row.pauseAfterCurrent = false
row.pausedByUser = false
row.holdForCutIn = false
row.updatedAt = Date.now()
syncPausedFlag(store)
})
}
+1 -2
View File
@@ -9,7 +9,7 @@ import { isLtxWorkflow, isTextToVideo, parseVideoWorkflow, workflowForExtension,
import { clipVideoPath, clipTitle, deleteRetryDraft, extendTempDir, getClip, nextClipPartName, removeExtendTemp, stillPath } from '~/server/utils/library'
import { extractLastFrame, probeHasAudio } from '~/server/utils/ffmpeg'
import { ensureComfyReady } from '~/server/utils/comfyLifecycle'
import { onLiveVideoSettled, onStudioQueueBurstStarted, isStudioQueuePaused } from '~/server/utils/studioQueue'
import { onLiveVideoSettled, onStudioQueueBurstStarted } from '~/server/utils/studioQueue'
import {
finishQueueBurst,
getShotQueue,
@@ -241,7 +241,6 @@ function seedCurrentVideo(job: Job) {
function chainShouldHold(job: Job) {
if (!job.library) return false
if (job.library.stopAfterCurrent) return true
if (isStudioQueuePaused(job.library.ownerKey)) return true
const queue = job.library.queueId
? getShotQueue(job.library.ownerKey, job.library.queueId)
: null