Pause the queue after the current clip for real, and carry global locks plus labeled permanence stills through generate and shot prompts.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+246
-32
@@ -2,6 +2,7 @@ import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFile
|
||||
import { join } from 'node:path'
|
||||
import { getJob, listJobs, type Job } from '~/server/utils/jobs'
|
||||
import { parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
|
||||
import type { PermanenceRef } from '~/utils/globalLocks'
|
||||
|
||||
export type StudioJobStatus = 'waiting' | 'running' | 'held' | 'complete' | 'error' | 'cancelled'
|
||||
|
||||
@@ -29,8 +30,11 @@ export interface StudioJobPayload {
|
||||
hideInput?: boolean
|
||||
folderLocked?: boolean
|
||||
referenceStillIds: Array<string | null>
|
||||
extensions: { prompt: string; duration: number }[]
|
||||
extensions: { prompt: string; duration: number; permanenceRefs?: PermanenceRef[] }[]
|
||||
queueAutoRun: boolean
|
||||
globalLocks?: string
|
||||
permanenceRefs?: PermanenceRef[]
|
||||
shotPermanenceRefs?: PermanenceRef[][]
|
||||
}
|
||||
|
||||
export interface StudioJob {
|
||||
@@ -48,10 +52,17 @@ export interface StudioJob {
|
||||
payload: StudioJobPayload
|
||||
cutIn?: boolean
|
||||
holdForCutIn?: boolean
|
||||
pauseAfterCurrent?: boolean
|
||||
pausedByUser?: boolean
|
||||
resumeAutoRun?: boolean
|
||||
lastError?: string
|
||||
}
|
||||
|
||||
type StudioQueueStore = {
|
||||
paused: boolean
|
||||
jobs: StudioJob[]
|
||||
}
|
||||
|
||||
const writeChains = new Map<string, Promise<unknown>>()
|
||||
let dispatchChain: Promise<unknown> = Promise.resolve()
|
||||
|
||||
@@ -68,38 +79,68 @@ function ensureOwner(owner: string) {
|
||||
mkdirSync(join(libraryRoot(), 'users', owner), { recursive: true })
|
||||
}
|
||||
|
||||
function readJobs(owner: string): StudioJob[] {
|
||||
function readStore(owner: string): StudioQueueStore {
|
||||
ensureOwner(owner)
|
||||
const path = queuePath(owner)
|
||||
if (!existsSync(path)) return []
|
||||
if (!existsSync(path)) return { paused: false, jobs: [] }
|
||||
try {
|
||||
const parsed = JSON.parse(readFileSync(path, 'utf8'))
|
||||
return Array.isArray(parsed) ? parsed : []
|
||||
if (Array.isArray(parsed)) return { paused: false, jobs: parsed }
|
||||
if (parsed && typeof parsed === 'object') {
|
||||
return {
|
||||
paused: parsed.paused === true,
|
||||
jobs: Array.isArray(parsed.jobs) ? parsed.jobs : []
|
||||
}
|
||||
}
|
||||
return { paused: false, jobs: [] }
|
||||
} catch {
|
||||
return []
|
||||
return { paused: false, jobs: [] }
|
||||
}
|
||||
}
|
||||
|
||||
function writeJobs(owner: string, jobs: StudioJob[]) {
|
||||
function writeStore(owner: string, store: StudioQueueStore) {
|
||||
ensureOwner(owner)
|
||||
const path = queuePath(owner)
|
||||
const tmp = `${path}.tmp`
|
||||
writeFileSync(tmp, JSON.stringify(jobs, null, 2))
|
||||
writeFileSync(tmp, JSON.stringify({
|
||||
paused: store.paused === true,
|
||||
jobs: store.jobs
|
||||
}, null, 2))
|
||||
renameSync(tmp, path)
|
||||
}
|
||||
|
||||
function readJobs(owner: string): StudioJob[] {
|
||||
return readStore(owner).jobs
|
||||
}
|
||||
|
||||
function mutate<T>(owner: string, fn: (jobs: StudioJob[]) => T): Promise<T> {
|
||||
const prev = writeChains.get(owner) || Promise.resolve()
|
||||
const run = prev.then(() => {
|
||||
const jobs = readJobs(owner)
|
||||
const result = fn(jobs)
|
||||
writeJobs(owner, jobs)
|
||||
const store = readStore(owner)
|
||||
const result = fn(store.jobs)
|
||||
writeStore(owner, store)
|
||||
return result
|
||||
})
|
||||
writeChains.set(owner, run.then(() => undefined, () => undefined))
|
||||
return run
|
||||
}
|
||||
|
||||
function mutateStore<T>(owner: string, fn: (store: StudioQueueStore) => T): Promise<T> {
|
||||
const prev = writeChains.get(owner) || Promise.resolve()
|
||||
const run = prev.then(() => {
|
||||
const store = readStore(owner)
|
||||
const result = fn(store)
|
||||
writeStore(owner, store)
|
||||
return result
|
||||
})
|
||||
writeChains.set(owner, run.then(() => undefined, () => undefined))
|
||||
return run
|
||||
}
|
||||
|
||||
export function isStudioQueuePaused(owner: string) {
|
||||
return readStore(owner).paused === true
|
||||
}
|
||||
|
||||
export function listOwnersWithStudioQueues() {
|
||||
const root = join(libraryRoot(), 'users')
|
||||
if (!existsSync(root)) return [] as string[]
|
||||
@@ -133,6 +174,8 @@ export function summarizeStudioJob(job: StudioJob) {
|
||||
queueAutoRun: job.payload.queueAutoRun === true,
|
||||
cutIn: job.cutIn === true,
|
||||
holdForCutIn: job.holdForCutIn === true,
|
||||
pauseAfterCurrent: job.pauseAfterCurrent === true,
|
||||
pausedByUser: job.pausedByUser === true,
|
||||
lastError: job.lastError
|
||||
}
|
||||
}
|
||||
@@ -307,11 +350,14 @@ function pickCutIn(jobs: StudioJob[]) {
|
||||
}
|
||||
|
||||
function pickHeld(jobs: StudioJob[]) {
|
||||
const held = jobs
|
||||
.filter(job => job.status === 'held' && job.shotQueueId)
|
||||
.sort((a, b) => b.updatedAt - a.updatedAt)
|
||||
const cutInHold = held.find(job => job.holdForCutIn)
|
||||
return cutInHold || held[0] || null
|
||||
return jobs
|
||||
.filter(job => (
|
||||
job.status === 'held'
|
||||
&& job.shotQueueId
|
||||
&& job.holdForCutIn
|
||||
&& !job.pausedByUser
|
||||
))
|
||||
.sort((a, b) => b.updatedAt - a.updatedAt)[0] || null
|
||||
}
|
||||
|
||||
function pickWaiting(jobs: StudioJob[]) {
|
||||
@@ -347,23 +393,27 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
||||
async function dispatchStudioQueue() {
|
||||
if (videoJobsBusy()) return
|
||||
for (const owner of listOwnersWithStudioQueues()) {
|
||||
const jobs = await mutate(owner, (list) => {
|
||||
repairStaleJobs(list)
|
||||
const next = pruneDone(list)
|
||||
list.splice(0, list.length, ...next)
|
||||
return list.map(item => structuredClone(item))
|
||||
const store = await mutateStore(owner, (current) => {
|
||||
repairStaleJobs(current.jobs)
|
||||
const next = pruneDone(current.jobs)
|
||||
current.jobs.splice(0, current.jobs.length, ...next)
|
||||
return {
|
||||
paused: current.paused === true,
|
||||
jobs: current.jobs.map(item => structuredClone(item))
|
||||
}
|
||||
})
|
||||
const cutIn = pickCutIn(jobs)
|
||||
const cutIn = pickCutIn(store.jobs)
|
||||
if (cutIn) {
|
||||
await startStudioJob(cutIn)
|
||||
return
|
||||
}
|
||||
const held = pickHeld(jobs)
|
||||
if (store.paused) return
|
||||
const held = pickHeld(store.jobs)
|
||||
if (held?.shotQueueId) {
|
||||
await resumeHeldStudioJob(held)
|
||||
return
|
||||
}
|
||||
const waiting = pickWaiting(jobs)
|
||||
const waiting = pickWaiting(store.jobs)
|
||||
if (waiting) {
|
||||
await startStudioJob(waiting)
|
||||
return
|
||||
@@ -445,7 +495,10 @@ export async function startStudioJob(item: StudioJob) {
|
||||
workflow: parseVideoWorkflow(payload.workflow),
|
||||
useIdentityRefs: payload.useIdentityRefs,
|
||||
queueAutoRun: extensions.length ? payload.queueAutoRun : false,
|
||||
queueBudget: extensions.length && payload.queueAutoRun ? extensions.length : 0
|
||||
queueBudget: extensions.length && payload.queueAutoRun ? extensions.length : 0,
|
||||
globalLocks: payload.globalLocks,
|
||||
permanenceRefs: payload.permanenceRefs,
|
||||
shotPermanenceRefs: payload.shotPermanenceRefs
|
||||
}
|
||||
|
||||
if (extensions.length) {
|
||||
@@ -472,7 +525,9 @@ export async function startStudioJob(item: StudioJob) {
|
||||
sound: payload.sound,
|
||||
useIdentityRefs: payload.useIdentityRefs,
|
||||
referenceStillIds: payload.referenceStillIds,
|
||||
initial: { prompt: payload.prompt, duration: payload.duration },
|
||||
globalLocks: payload.globalLocks,
|
||||
permanenceRefs: payload.permanenceRefs,
|
||||
initial: { prompt: payload.prompt, duration: payload.duration, permanenceRefs: payload.shotPermanenceRefs?.[0] },
|
||||
extensions,
|
||||
jobId: job.id
|
||||
})
|
||||
@@ -549,22 +604,58 @@ export async function onLiveVideoSettled(job: Job) {
|
||||
const remaining = job.library
|
||||
? Math.max(0, (job.library.chainTotal || 1) - (job.library.chainStep || 1))
|
||||
: 0
|
||||
const interrupted = job.library?.stopAfterCurrent === true
|
||||
&& remaining > 0
|
||||
&& job.status !== 'error'
|
||||
&& job.status !== 'cancelled'
|
||||
const failed = job.status === 'error' || job.status === 'cancelled'
|
||||
const snapshot = readStore(owner)
|
||||
const rowNow = snapshot.jobs.find(item => item.liveJobId === job.id) || snapshot.jobs.find(item => item.status === 'running')
|
||||
const userPause = !failed && (
|
||||
rowNow?.pauseAfterCurrent === true
|
||||
|| (snapshot.paused && !rowNow?.holdForCutIn)
|
||||
|| (job.library?.stopAfterCurrent === true && !rowNow?.holdForCutIn)
|
||||
)
|
||||
const cutInHold = !failed && remaining > 0 && rowNow?.holdForCutIn === true && !userPause
|
||||
const interrupted = remaining > 0 && (userPause || cutInHold)
|
||||
if (interrupted && job.status === 'running') {
|
||||
job.status = 'complete'
|
||||
}
|
||||
const hold = interrupted
|
||||
if (hold) {
|
||||
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
|
||||
row.status = 'held'
|
||||
row.liveJobId = undefined
|
||||
row.pauseAfterCurrent = false
|
||||
row.pausedByUser = true
|
||||
row.holdForCutIn = false
|
||||
row.resumeAutoRun = row.payload.queueAutoRun === true || row.resumeAutoRun === true
|
||||
row.updatedAt = Date.now()
|
||||
})
|
||||
} else if (cutInHold) {
|
||||
await mutate(owner, (jobs) => {
|
||||
const row = jobs.find(item => item.liveJobId === job.id) || jobs.find(item => item.status === 'running')
|
||||
if (!row) return
|
||||
row.status = 'held'
|
||||
row.liveJobId = undefined
|
||||
row.holdForCutIn = true
|
||||
row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true
|
||||
row.resumeAutoRun = row.payload.queueAutoRun === true || job.library?.queueAutoRun === true || row.resumeAutoRun === true
|
||||
row.updatedAt = Date.now()
|
||||
})
|
||||
} 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
|
||||
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()
|
||||
})
|
||||
} else {
|
||||
@@ -577,3 +668,126 @@ export async function onLiveVideoSettled(job: Job) {
|
||||
}
|
||||
kickStudioQueue()
|
||||
}
|
||||
|
||||
function applyLiveJobPause(live: Job | undefined, pause: boolean, restoreAutoRun?: boolean) {
|
||||
if (!live?.library) return
|
||||
live.library.stopAfterCurrent = pause
|
||||
if (pause) live.library.queueAutoRun = false
|
||||
else if (restoreAutoRun) live.library.queueAutoRun = true
|
||||
}
|
||||
|
||||
export async function toggleStudioQueuePause(owner: string, options: {
|
||||
resume?: boolean
|
||||
jobId?: string
|
||||
queueId?: string
|
||||
} = {}) {
|
||||
const { pauseShotQueue, setShotQueuePause, listShotQueues, getShotQueue } = await import('~/server/utils/shotQueue')
|
||||
const storeNow = readStore(owner)
|
||||
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)
|
||||
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)
|
||||
const shouldPause = options.resume === true ? false : options.resume === false ? true : !armed
|
||||
|
||||
if (!shouldPause) {
|
||||
const heldPaused = storeNow.jobs.filter(job => (
|
||||
job.status === 'held'
|
||||
&& Boolean(job.shotQueueId)
|
||||
&& job.pausedByUser === true
|
||||
&& (!options.jobId || job.id === options.jobId)
|
||||
))
|
||||
const restored = await mutateStore(owner, (store) => {
|
||||
store.paused = false
|
||||
for (const job of store.jobs) {
|
||||
if (options.jobId && job.id !== options.jobId) continue
|
||||
job.pauseAfterCurrent = false
|
||||
job.pausedByUser = false
|
||||
}
|
||||
if (!options.jobId) {
|
||||
for (const job of store.jobs) {
|
||||
job.pauseAfterCurrent = false
|
||||
job.pausedByUser = false
|
||||
}
|
||||
}
|
||||
return structuredClone(store)
|
||||
})
|
||||
const rows = options.jobId
|
||||
? restored.jobs.filter(job => job.id === options.jobId)
|
||||
: 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) {
|
||||
for (const queue of listShotQueues(owner)) {
|
||||
if (queue.stopAfterCurrent) await setShotQueuePause(owner, queue.id, false).catch(() => null)
|
||||
}
|
||||
}
|
||||
if (heldPaused[0] && !videoJobsBusy()) {
|
||||
await resumeHeldStudioJob(heldPaused[0]).catch(() => null)
|
||||
} else {
|
||||
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')
|
||||
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'
|
||||
row.resumeAutoRun = row.payload.queueAutoRun === true || row.resumeAutoRun === true
|
||||
row.updatedAt = Date.now()
|
||||
}
|
||||
return {
|
||||
paused: store.paused,
|
||||
job: row ? structuredClone(row) : null
|
||||
}
|
||||
})
|
||||
|
||||
if (updated.job?.liveJobId) {
|
||||
applyLiveJobPause(getJob(updated.job.liveJobId), true)
|
||||
}
|
||||
const queueId = updated.job?.shotQueueId || options.queueId
|
||||
const queued = queueId ? getShotQueue(owner, queueId) : null
|
||||
if (updated.job && queued?.autoRun) {
|
||||
await patchStudioJob(owner, updated.job.id, (job) => {
|
||||
job.resumeAutoRun = true
|
||||
}).catch(() => null)
|
||||
}
|
||||
if (queueId) await pauseShotQueue(owner, queueId).catch(() => null)
|
||||
if (!updated.job && !queueId) {
|
||||
for (const queue of listShotQueues(owner)) {
|
||||
if (queue.status === 'running') await pauseShotQueue(owner, queue.id).catch(() => null)
|
||||
}
|
||||
}
|
||||
const anyGenerating = storeNow.jobs.some(job => job.status === 'running')
|
||||
|| listShotQueues(owner).some(queue => queue.status === 'running')
|
||||
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 }
|
||||
}
|
||||
|
||||
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
|
||||
row.status = 'running'
|
||||
row.liveJobId = liveJobId
|
||||
row.pauseAfterCurrent = false
|
||||
row.pausedByUser = false
|
||||
row.holdForCutIn = false
|
||||
row.updatedAt = Date.now()
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user