Keep image jobs on the studio queue and confirm each job or library delete. Label stitch-frame slots when thumbs are hidden, and open output stills in a lightbox.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
+46
-235
@@ -1,7 +1,5 @@
|
||||
import { ensureComfyReady } from '~/server/utils/comfyLifecycle'
|
||||
import { getComfyHost, comfyConfigured, uploadImage, queuePrompt, purgeComfyArtifacts } from '~/server/utils/comfy'
|
||||
import { withImageComfyHost, waitForImageEdit, downloadEditedImage } from '~/server/utils/imageComfy'
|
||||
import { buildEditWorkflow } from '~/server/utils/imageWorkflow'
|
||||
import { addStudioJob, kickStudioQueue, listStudioJobs } from '~/server/utils/studioQueue'
|
||||
import { comfyConfigured } from '~/server/utils/comfy'
|
||||
import { imageDimensions } from '~/server/utils/resolution'
|
||||
|
||||
function parseSteps(raw: unknown) {
|
||||
@@ -82,6 +80,7 @@ export default defineEventHandler(async (event) => {
|
||||
? Number(fields.seed)
|
||||
: Math.floor(Math.random() * 2_147_483_647)
|
||||
const size = imageDimensions(image.data)
|
||||
const clipName = (fields.name || '').trim().slice(0, 80)
|
||||
|
||||
const still = await rememberInputStill({
|
||||
ownerKey,
|
||||
@@ -92,74 +91,63 @@ export default defineEventHandler(async (event) => {
|
||||
height: size?.height || 0,
|
||||
hideInput
|
||||
})
|
||||
if (reference) {
|
||||
await rememberInputStill({
|
||||
const savedRef = reference
|
||||
? await rememberInputStill({
|
||||
ownerKey,
|
||||
folderId,
|
||||
filename: reference.filename,
|
||||
data: reference.data,
|
||||
hideInput
|
||||
})
|
||||
}
|
||||
: null
|
||||
|
||||
const passes = extraPasses
|
||||
const chainTotal = 1 + passes.length
|
||||
const familyId = chainTotal > 1 ? crypto.randomUUID() : undefined
|
||||
const job = createJob()
|
||||
job.kind = 'edit'
|
||||
job.maxStep = steps
|
||||
job.hideThumbnail = hideThumbnail
|
||||
job.library = {
|
||||
const familyId = crypto.randomUUID()
|
||||
const studio = await addStudioJob({
|
||||
ownerKey,
|
||||
folderId,
|
||||
hideThumbnail,
|
||||
hideInput,
|
||||
folderLocked,
|
||||
name: (fields.name || '').trim().slice(0, 80),
|
||||
prompt,
|
||||
aspect: 'auto',
|
||||
width: size?.width || 0,
|
||||
height: size?.height || 0,
|
||||
steps,
|
||||
turbo: true,
|
||||
seed,
|
||||
cfg,
|
||||
stillId: still?.id,
|
||||
familyId,
|
||||
chainIndex: 0,
|
||||
chainStep: 1,
|
||||
chainTotal,
|
||||
chainLabel: chainTotal > 1 ? 'Pass 1' : undefined,
|
||||
passes
|
||||
}
|
||||
emitJob(job, {
|
||||
type: 'status',
|
||||
message: 'Checking Beast ComfyUI...',
|
||||
progress: 2,
|
||||
chainStep: 1,
|
||||
chainTotal,
|
||||
chainLabel: job.library.chainLabel
|
||||
})
|
||||
|
||||
void runEdit(job, {
|
||||
image,
|
||||
reference,
|
||||
prompt,
|
||||
passes,
|
||||
negative: (fields.negative || '').trim(),
|
||||
steps,
|
||||
seed,
|
||||
cfg
|
||||
}).catch(async (error) => {
|
||||
if (job.status === 'error' || job.status === 'cancelled') return
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
job.status = 'error'
|
||||
job.error = message
|
||||
emitJob(job, { type: 'error', error: message, message })
|
||||
kind: 'edit',
|
||||
payload: {
|
||||
prompt,
|
||||
name: clipName,
|
||||
folderId,
|
||||
aspect: 'auto',
|
||||
width: size?.width || 0,
|
||||
height: size?.height || 0,
|
||||
steps,
|
||||
turbo: true,
|
||||
seed,
|
||||
cfg,
|
||||
fps: 24,
|
||||
samplerName: 'res_multistep',
|
||||
scheduler: 'simple',
|
||||
duration: 0,
|
||||
sound: false,
|
||||
workflow: 'v1',
|
||||
useIdentityRefs: false,
|
||||
stillId: still?.id,
|
||||
stillFilename: still?.filename,
|
||||
hideThumbnail,
|
||||
hideInput,
|
||||
folderLocked,
|
||||
referenceStillIds: [null, null, null, null],
|
||||
extensions: [],
|
||||
queueAutoRun: false,
|
||||
negative: (fields.negative || '').trim(),
|
||||
passes,
|
||||
referenceStillId: savedRef?.id,
|
||||
referenceStillFilename: savedRef?.filename
|
||||
}
|
||||
})
|
||||
await kickStudioQueue()
|
||||
const latest = listStudioJobs(ownerKey).find(item => item.id === studio.id)
|
||||
const liveJobId = latest?.liveJobId || ''
|
||||
|
||||
return {
|
||||
jobId: job.id,
|
||||
jobId: liveJobId || studio.id,
|
||||
studioJobId: studio.id,
|
||||
queued: !liveJobId,
|
||||
seed,
|
||||
steps,
|
||||
hideThumbnail,
|
||||
@@ -168,180 +156,3 @@ export default defineEventHandler(async (event) => {
|
||||
chainTotal
|
||||
}
|
||||
})
|
||||
|
||||
async function runEdit(
|
||||
job: ReturnType<typeof createJob>,
|
||||
params: {
|
||||
image: { filename: string; data: Buffer; type?: string }
|
||||
reference: { filename: string; data: Buffer; type?: string } | null
|
||||
prompt: string
|
||||
passes: { prompt: string }[]
|
||||
negative: string
|
||||
steps: number
|
||||
seed: number
|
||||
cfg: number
|
||||
}
|
||||
) {
|
||||
const library = job.library
|
||||
if (!library) throw new Error('Edit job is missing library metadata')
|
||||
const prompts = [params.prompt, ...params.passes.map(item => item.prompt)]
|
||||
const chainTotal = prompts.length
|
||||
library.chainTotal = chainTotal
|
||||
library.familyId = chainTotal > 1 ? (library.familyId || crypto.randomUUID()) : library.familyId
|
||||
|
||||
await ensureComfyReady((status) => {
|
||||
emitChainJob(job, {
|
||||
type: status.state === 'busy' ? 'busy' : 'status',
|
||||
message: status.message,
|
||||
progress: status.state === 'online' ? Math.max(job.progress, 6) : Math.max(job.progress, 3),
|
||||
busy: status.state === 'busy'
|
||||
})
|
||||
})
|
||||
if (job.status === 'cancelled') throw new Error('Job interrupted.')
|
||||
job.imageComfyHost = getComfyHost()
|
||||
|
||||
await withImageComfyHost(job.imageComfyHost, async () => {
|
||||
let current = params.image
|
||||
let parentStillId: string | undefined
|
||||
for (let index = 0; index < prompts.length; index++) {
|
||||
if (job.status === 'cancelled') throw new Error('Job interrupted.')
|
||||
const prompt = prompts[index]
|
||||
const last = index === prompts.length - 1
|
||||
if (index > 0) {
|
||||
await ensureComfyReady((status) => {
|
||||
emitChainJob(job, {
|
||||
type: status.state === 'busy' ? 'busy' : 'status',
|
||||
message: status.message,
|
||||
progress: status.state === 'online' ? 4 : 2,
|
||||
busy: status.state === 'busy'
|
||||
})
|
||||
})
|
||||
}
|
||||
const seed = index === 0 ? params.seed : Math.floor(Math.random() * 2_147_483_647)
|
||||
library.prompt = prompt
|
||||
library.seed = seed
|
||||
library.chainIndex = index
|
||||
library.chainStep = index + 1
|
||||
library.chainLabel = chainTotal > 1 ? `Pass ${index + 1}` : undefined
|
||||
const reference = index === 0 ? params.reference : null
|
||||
const dual = Boolean(reference)
|
||||
const passName = stillChainName(library.name || '', index)
|
||||
|
||||
emitChainJob(job, {
|
||||
type: 'status',
|
||||
message: dual ? 'Uploading both stills to Beast...' : 'Uploading still to Beast...',
|
||||
progress: 8
|
||||
})
|
||||
const uploaded = await uploadImage(current, job.id)
|
||||
const uploadedRef = reference
|
||||
? await uploadImage({
|
||||
...reference,
|
||||
filename: `ref_${reference.filename || 'image2.png'}`
|
||||
}, job.id)
|
||||
: null
|
||||
if (job.status === 'cancelled') throw new Error('Job interrupted.')
|
||||
|
||||
emitChainJob(job, {
|
||||
type: 'status',
|
||||
message: dual
|
||||
? 'Queueing Flux.2 Klein two-image edit on Beast...'
|
||||
: 'Queueing Flux.2 Klein edit on Beast...',
|
||||
progress: 12
|
||||
})
|
||||
const graph = buildEditWorkflow({
|
||||
prompt,
|
||||
negative: params.negative,
|
||||
imageName: uploaded.name,
|
||||
referenceImageName: uploadedRef?.name,
|
||||
steps: params.steps,
|
||||
seed,
|
||||
cfg: params.cfg,
|
||||
filenamePrefix: `aigen_edit_${job.id.slice(0, 8)}_p${index + 1}`
|
||||
})
|
||||
const queued = await queuePrompt(graph, job.clientId)
|
||||
job.promptId = queued.prompt_id
|
||||
job.status = 'running'
|
||||
emitChainJob(job, {
|
||||
type: 'status',
|
||||
message: dual ? 'Applying image 2 onto image 1...' : 'Editing still on Beast...',
|
||||
progress: 18,
|
||||
maxStep: params.steps
|
||||
})
|
||||
|
||||
const output = await waitForImageEdit({
|
||||
promptId: queued.prompt_id,
|
||||
clientId: job.clientId,
|
||||
timeoutMs: 10 * 60 * 1000,
|
||||
onProgress: (event) => {
|
||||
emitChainJob(job, {
|
||||
type: 'status',
|
||||
message: event.message,
|
||||
progress: event.progress,
|
||||
step: event.step,
|
||||
maxStep: event.maxStep || params.steps,
|
||||
node: event.node
|
||||
})
|
||||
},
|
||||
isCancelled: () => job.status === 'cancelled'
|
||||
})
|
||||
|
||||
emitChainJob(job, { type: 'status', message: 'Saving edited still...', progress: 94 })
|
||||
const buffer = await downloadEditedImage(output)
|
||||
const size = imageDimensions(buffer)
|
||||
const still = await saveStill({
|
||||
ownerKey: library.ownerKey,
|
||||
folderId: library.folderId,
|
||||
filename: passName ? `${passName}.png` : output.filename,
|
||||
data: buffer,
|
||||
width: size?.width || 0,
|
||||
height: size?.height || 0,
|
||||
hideInput: job.hideThumbnail === true,
|
||||
role: 'output',
|
||||
name: passName || undefined,
|
||||
prompt,
|
||||
familyId: library.familyId,
|
||||
parentStillId,
|
||||
chainIndex: index
|
||||
})
|
||||
job.stillId = still?.id
|
||||
parentStillId = still?.id
|
||||
await purgeComfyArtifacts({
|
||||
video: { filename: output.filename, subfolder: output.subfolder, type: output.type },
|
||||
imageName: uploaded.name,
|
||||
imageSubfolder: uploaded.subfolder,
|
||||
extraImageNames: uploadedRef?.name ? [uploadedRef.name] : [],
|
||||
promptId: job.promptId
|
||||
})
|
||||
|
||||
if (!last) {
|
||||
emitChainJob(job, {
|
||||
type: 'checkpoint',
|
||||
message: `Pass ${index + 1} saved`,
|
||||
progress: 100,
|
||||
stillId: still?.id,
|
||||
hideThumbnail: job.hideThumbnail,
|
||||
folderLocked: library.folderLocked
|
||||
})
|
||||
current = {
|
||||
filename: still?.filename || `pass_${index + 1}.png`,
|
||||
data: buffer,
|
||||
type: 'image/png'
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
job.status = 'complete'
|
||||
emitChainJob(job, {
|
||||
type: 'complete',
|
||||
message: chainTotal > 1 ? 'Edit chain finished on Beast' : 'Edit finished on Beast',
|
||||
progress: 100,
|
||||
stillId: still?.id,
|
||||
filename: output.filename,
|
||||
subfolder: output.subfolder,
|
||||
mediaType: 'image',
|
||||
hideThumbnail: job.hideThumbnail,
|
||||
folderLocked: library.folderLocked
|
||||
})
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
import { deleteShotQueue } from '~/server/utils/shotQueue'
|
||||
import { cancelOrphanQueue } from '~/server/utils/studioQueue'
|
||||
|
||||
export default defineEventHandler(async (event) => {
|
||||
const { owner } = assertLibraryOwner(event)
|
||||
const id = String(getRouterParam(event, 'id') || '')
|
||||
await deleteShotQueue(owner, id)
|
||||
await cancelOrphanQueue(owner, id)
|
||||
return { ok: true }
|
||||
})
|
||||
|
||||
@@ -11,13 +11,21 @@ export default defineEventHandler((event) => {
|
||||
const byId = new Map(queues.map(queue => [queue.id, queue]))
|
||||
const jobs = active.map((job) => {
|
||||
const shots = job.shotQueueId ? byId.get(job.shotQueueId) || null : null
|
||||
const kind = job.kind === 'edit' ? 'edit' : 'video'
|
||||
const shotSummary = shots ? summarizeQueue(shots) : null
|
||||
const plannedShots = shotSummary
|
||||
? undefined
|
||||
: [
|
||||
{ prompt: job.payload.prompt, duration: job.payload.duration },
|
||||
...(job.payload.extensions || [])
|
||||
]
|
||||
: kind === 'edit'
|
||||
? (job.payload.passes?.length
|
||||
? [
|
||||
{ prompt: job.payload.prompt, duration: 0 },
|
||||
...job.payload.passes.map(item => ({ prompt: item.prompt, duration: 0 }))
|
||||
]
|
||||
: undefined)
|
||||
: [
|
||||
{ prompt: job.payload.prompt, duration: job.payload.duration },
|
||||
...(job.payload.extensions || [])
|
||||
]
|
||||
return {
|
||||
...summarizeStudioJob(job),
|
||||
...(full ? {
|
||||
|
||||
@@ -220,7 +220,7 @@ function hostWithPort(host: string, port: number) {
|
||||
}
|
||||
}
|
||||
|
||||
async function ensureUnlocked(onStatus: StatusFn) {
|
||||
async function ensureUnlocked(onStatus: StatusFn, skipBusyWait = false) {
|
||||
const cfg = settings()
|
||||
const host = cfg.host
|
||||
onStatus({ state: 'offline', message: 'Checking ComfyUI...', host })
|
||||
@@ -228,7 +228,7 @@ async function ensureUnlocked(onStatus: StatusFn) {
|
||||
const health = await checkComfyHttp(cfg.healthTimeoutMs)
|
||||
if (health.ok) {
|
||||
log('http-ok', { host })
|
||||
await waitWhileBusy(onStatus)
|
||||
if (!skipBusyWait) await waitWhileBusy(onStatus)
|
||||
onStatus({ state: 'online', message: 'ComfyUI online', host })
|
||||
return
|
||||
}
|
||||
@@ -247,7 +247,7 @@ async function ensureUnlocked(onStatus: StatusFn) {
|
||||
if (viaProxy.ok) {
|
||||
setComfyHostOverride(proxyUrl)
|
||||
log('http-ok-proxy', { host: proxyUrl, localPort: remote.port || null })
|
||||
await waitWhileBusy(onStatus)
|
||||
if (!skipBusyWait) await waitWhileBusy(onStatus)
|
||||
onStatus({ state: 'online', message: 'ComfyUI online', host: proxyUrl, processRunning: true })
|
||||
return
|
||||
}
|
||||
@@ -269,7 +269,7 @@ async function ensureUnlocked(onStatus: StatusFn) {
|
||||
})
|
||||
const ready = await pollUntilHealthy(cfg.bootPollAttempts, cfg.bootPollDelayMs, onStatus, 'booting', 'Waiting for ComfyUI to finish booting')
|
||||
if (ready) {
|
||||
await waitWhileBusy(onStatus)
|
||||
if (!skipBusyWait) await waitWhileBusy(onStatus)
|
||||
onStatus({ state: 'online', message: 'ComfyUI online', host, processRunning: true })
|
||||
return
|
||||
}
|
||||
@@ -309,7 +309,7 @@ async function ensureUnlocked(onStatus: StatusFn) {
|
||||
const next = await checkComfyHttp()
|
||||
if (next.ok) {
|
||||
log('started', { attempt, reason: started.reason })
|
||||
await waitWhileBusy(onStatus)
|
||||
if (!skipBusyWait) await waitWhileBusy(onStatus)
|
||||
onStatus({ state: 'online', message: 'ComfyUI online', host, processRunning: true })
|
||||
return
|
||||
}
|
||||
@@ -322,8 +322,9 @@ async function ensureUnlocked(onStatus: StatusFn) {
|
||||
})
|
||||
}
|
||||
|
||||
export function ensureComfyReady(onStatus: StatusFn = () => undefined) {
|
||||
const run = gate.then(() => ensureUnlocked(onStatus))
|
||||
export function ensureComfyReady(onStatus: StatusFn = () => undefined, options?: { skipBusyWait?: boolean }) {
|
||||
const skipBusyWait = options?.skipBusyWait === true
|
||||
const run = gate.then(() => ensureUnlocked(onStatus, skipBusyWait))
|
||||
gate = run.then(() => undefined, () => undefined)
|
||||
return run
|
||||
}
|
||||
|
||||
@@ -0,0 +1,213 @@
|
||||
import { createJob, emitJob, type Job } from '~/server/utils/jobs'
|
||||
import { ensureComfyReady } from '~/server/utils/comfyLifecycle'
|
||||
import { getComfyHost, uploadImage, queuePrompt, purgeComfyArtifacts } from '~/server/utils/comfy'
|
||||
import { withImageComfyHost, waitForImageEdit, downloadEditedImage } from '~/server/utils/imageComfy'
|
||||
import { buildEditWorkflow } from '~/server/utils/imageWorkflow'
|
||||
import { imageDimensions } from '~/server/utils/resolution'
|
||||
import { emitChainJob } from '~/server/utils/watch'
|
||||
import { saveStill, stillChainName } from '~/server/utils/library'
|
||||
|
||||
export type EditImageFile = { filename: string; data: Buffer; type?: string }
|
||||
|
||||
export type EditRunParams = {
|
||||
image: EditImageFile
|
||||
reference: EditImageFile | null
|
||||
prompt: string
|
||||
passes: { prompt: string }[]
|
||||
negative: string
|
||||
steps: number
|
||||
seed: number
|
||||
cfg: number
|
||||
}
|
||||
|
||||
export async function runEdit(job: Job, params: EditRunParams) {
|
||||
const library = job.library
|
||||
if (!library) throw new Error('Edit job is missing library metadata')
|
||||
const prompts = [params.prompt, ...params.passes.map(item => item.prompt)]
|
||||
const chainTotal = prompts.length
|
||||
library.chainTotal = chainTotal
|
||||
library.familyId = chainTotal > 1 ? (library.familyId || crypto.randomUUID()) : library.familyId
|
||||
|
||||
try {
|
||||
await ensureComfyReady((status) => {
|
||||
emitChainJob(job, {
|
||||
type: status.state === 'busy' ? 'busy' : 'status',
|
||||
message: status.message,
|
||||
progress: status.state === 'online' ? Math.max(job.progress, 6) : Math.max(job.progress, 3),
|
||||
busy: status.state === 'busy'
|
||||
})
|
||||
}, { skipBusyWait: true })
|
||||
if (job.status === 'cancelled') throw new Error('Job interrupted.')
|
||||
job.imageComfyHost = getComfyHost()
|
||||
|
||||
await withImageComfyHost(job.imageComfyHost, async () => {
|
||||
let current = params.image
|
||||
let parentStillId: string | undefined
|
||||
for (let index = 0; index < prompts.length; index++) {
|
||||
if (job.status === 'cancelled') throw new Error('Job interrupted.')
|
||||
if (index > 0 && library.stopAfterCurrent === true) break
|
||||
const prompt = prompts[index]
|
||||
const last = index === prompts.length - 1
|
||||
if (index > 0) {
|
||||
await ensureComfyReady((status) => {
|
||||
emitChainJob(job, {
|
||||
type: status.state === 'busy' ? 'busy' : 'status',
|
||||
message: status.message,
|
||||
progress: status.state === 'online' ? 4 : 2,
|
||||
busy: status.state === 'busy'
|
||||
})
|
||||
}, { skipBusyWait: true })
|
||||
}
|
||||
const seed = index === 0 ? params.seed : Math.floor(Math.random() * 2_147_483_647)
|
||||
library.prompt = prompt
|
||||
library.seed = seed
|
||||
library.chainIndex = index
|
||||
library.chainStep = index + 1
|
||||
library.chainLabel = chainTotal > 1 ? `Pass ${index + 1}` : undefined
|
||||
const reference = index === 0 ? params.reference : null
|
||||
const dual = Boolean(reference)
|
||||
const passName = stillChainName(library.name || '', index)
|
||||
|
||||
emitChainJob(job, {
|
||||
type: 'status',
|
||||
message: dual ? 'Uploading both stills to Beast...' : 'Uploading still to Beast...',
|
||||
progress: 8
|
||||
})
|
||||
const uploaded = await uploadImage(current, job.id)
|
||||
const uploadedRef = reference
|
||||
? await uploadImage({
|
||||
...reference,
|
||||
filename: `ref_${reference.filename || 'image2.png'}`
|
||||
}, job.id)
|
||||
: null
|
||||
if (job.status === 'cancelled') throw new Error('Job interrupted.')
|
||||
|
||||
emitChainJob(job, {
|
||||
type: 'status',
|
||||
message: dual
|
||||
? 'Queueing Flux.2 Klein two-image edit on Beast...'
|
||||
: 'Queueing Flux.2 Klein edit on Beast...',
|
||||
progress: 12
|
||||
})
|
||||
const graph = buildEditWorkflow({
|
||||
prompt,
|
||||
negative: params.negative,
|
||||
imageName: uploaded.name,
|
||||
referenceImageName: uploadedRef?.name,
|
||||
steps: params.steps,
|
||||
seed,
|
||||
cfg: params.cfg,
|
||||
filenamePrefix: `aigen_edit_${job.id.slice(0, 8)}_p${index + 1}`
|
||||
})
|
||||
const queued = await queuePrompt(graph, job.clientId)
|
||||
job.promptId = queued.prompt_id
|
||||
job.status = 'running'
|
||||
emitChainJob(job, {
|
||||
type: 'status',
|
||||
message: dual ? 'Applying image 2 onto image 1...' : 'Editing still on Beast...',
|
||||
progress: 18,
|
||||
maxStep: params.steps
|
||||
})
|
||||
|
||||
const output = await waitForImageEdit({
|
||||
promptId: queued.prompt_id,
|
||||
clientId: job.clientId,
|
||||
timeoutMs: 10 * 60 * 1000,
|
||||
onProgress: (event) => {
|
||||
emitChainJob(job, {
|
||||
type: 'status',
|
||||
message: event.message,
|
||||
progress: event.progress,
|
||||
step: event.step,
|
||||
maxStep: event.maxStep || params.steps,
|
||||
node: event.node
|
||||
})
|
||||
},
|
||||
isCancelled: () => job.status === 'cancelled'
|
||||
})
|
||||
|
||||
emitChainJob(job, { type: 'status', message: 'Saving edited still...', progress: 94 })
|
||||
const buffer = await downloadEditedImage(output)
|
||||
const size = imageDimensions(buffer)
|
||||
const still = await saveStill({
|
||||
ownerKey: library.ownerKey,
|
||||
folderId: library.folderId,
|
||||
filename: passName ? `${passName}.png` : output.filename,
|
||||
data: buffer,
|
||||
width: size?.width || 0,
|
||||
height: size?.height || 0,
|
||||
hideInput: job.hideThumbnail === true,
|
||||
role: 'output',
|
||||
name: passName || undefined,
|
||||
prompt,
|
||||
familyId: library.familyId,
|
||||
parentStillId,
|
||||
chainIndex: index
|
||||
})
|
||||
job.stillId = still?.id
|
||||
parentStillId = still?.id
|
||||
await purgeComfyArtifacts({
|
||||
video: { filename: output.filename, subfolder: output.subfolder, type: output.type },
|
||||
imageName: uploaded.name,
|
||||
imageSubfolder: uploaded.subfolder,
|
||||
extraImageNames: uploadedRef?.name ? [uploadedRef.name] : [],
|
||||
promptId: job.promptId
|
||||
})
|
||||
|
||||
if (!last) {
|
||||
emitChainJob(job, {
|
||||
type: 'checkpoint',
|
||||
message: `Pass ${index + 1} saved`,
|
||||
progress: 100,
|
||||
stillId: still?.id,
|
||||
hideThumbnail: job.hideThumbnail,
|
||||
folderLocked: library.folderLocked
|
||||
})
|
||||
current = {
|
||||
filename: still?.filename || `pass_${index + 1}.png`,
|
||||
data: buffer,
|
||||
type: 'image/png'
|
||||
}
|
||||
if (library.stopAfterCurrent === true) break
|
||||
continue
|
||||
}
|
||||
|
||||
job.status = 'complete'
|
||||
emitChainJob(job, {
|
||||
type: 'complete',
|
||||
message: chainTotal > 1 ? 'Edit chain finished on Beast' : 'Edit finished on Beast',
|
||||
progress: 100,
|
||||
stillId: still?.id,
|
||||
filename: output.filename,
|
||||
subfolder: output.subfolder,
|
||||
mediaType: 'image',
|
||||
hideThumbnail: job.hideThumbnail,
|
||||
folderLocked: library.folderLocked
|
||||
})
|
||||
}
|
||||
})
|
||||
} catch (error) {
|
||||
if (job.status !== 'error' && job.status !== 'cancelled') {
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
job.status = 'error'
|
||||
job.error = message
|
||||
emitJob(job, { type: 'error', error: message, message })
|
||||
}
|
||||
} finally {
|
||||
const { onLiveVideoSettled } = await import('~/server/utils/studioQueue')
|
||||
await onLiveVideoSettled(job)
|
||||
}
|
||||
}
|
||||
|
||||
export function createEditLiveJob(params: {
|
||||
steps: number
|
||||
hideThumbnail: boolean
|
||||
library: NonNullable<Job['library']>
|
||||
}) {
|
||||
const job = createJob()
|
||||
job.kind = 'edit'
|
||||
job.maxStep = params.steps
|
||||
job.hideThumbnail = params.hideThumbnail
|
||||
job.library = params.library
|
||||
return job
|
||||
}
|
||||
@@ -388,10 +388,11 @@ export async function pauseShotQueue(owner: string, id: string) {
|
||||
return setShotQueuePause(owner, id, true)
|
||||
}
|
||||
|
||||
export async function deleteShotQueue(owner: string, id: string) {
|
||||
if (activeJobs.has(id)) {
|
||||
export async function deleteShotQueue(owner: string, id: string, options: { force?: boolean } = {}) {
|
||||
if (!options.force && activeJobs.has(id)) {
|
||||
throw createError({ statusCode: 409, statusMessage: 'Stop this batch before deleting it' })
|
||||
}
|
||||
activeJobs.delete(id)
|
||||
return mutate(owner, (queues) => {
|
||||
const index = queues.findIndex(item => item.id === id)
|
||||
if (index < 0) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
|
||||
|
||||
+179
-17
@@ -1,12 +1,13 @@
|
||||
import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFileSync } from 'node:fs'
|
||||
import { join } from 'node:path'
|
||||
import { getJob, listJobs, type Job } from '~/server/utils/jobs'
|
||||
import { listPendingJobs, patchPendingJob, readPendingJob } from '~/server/utils/pending'
|
||||
import { getJob, listJobs, emitJob, type Job } from '~/server/utils/jobs'
|
||||
import { listPendingJobs, patchPendingJob, readPendingJob, deletePendingJob } from '~/server/utils/pending'
|
||||
import { getShotQueue } from '~/server/utils/shotQueue'
|
||||
import { parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
|
||||
import type { PermanenceRef } from '~/utils/globalLocks'
|
||||
|
||||
export type StudioJobStatus = 'waiting' | 'running' | 'held' | 'complete' | 'error' | 'cancelled'
|
||||
export type StudioJobKind = 'video' | 'edit'
|
||||
|
||||
export interface StudioJobPayload {
|
||||
prompt: string
|
||||
@@ -37,6 +38,10 @@ export interface StudioJobPayload {
|
||||
globalLocks?: string
|
||||
permanenceRefs?: PermanenceRef[]
|
||||
shotPermanenceRefs?: PermanenceRef[][]
|
||||
negative?: string
|
||||
passes?: { prompt: string }[]
|
||||
referenceStillId?: string
|
||||
referenceStillFilename?: string
|
||||
}
|
||||
|
||||
export interface StudioJob {
|
||||
@@ -45,6 +50,7 @@ export interface StudioJob {
|
||||
createdAt: number
|
||||
updatedAt: number
|
||||
status: StudioJobStatus
|
||||
kind: StudioJobKind
|
||||
name: string
|
||||
prompt: string
|
||||
shotCount: number
|
||||
@@ -176,12 +182,17 @@ export function listStudioJobs(owner: string) {
|
||||
return readJobs(owner)
|
||||
}
|
||||
|
||||
export function studioJobKind(job: Pick<StudioJob, 'kind'> | { kind?: string }) {
|
||||
return job.kind === 'edit' ? 'edit' : 'video'
|
||||
}
|
||||
|
||||
export function summarizeStudioJob(job: StudioJob) {
|
||||
return {
|
||||
id: job.id,
|
||||
createdAt: job.createdAt,
|
||||
updatedAt: job.updatedAt,
|
||||
status: job.status,
|
||||
kind: studioJobKind(job),
|
||||
name: job.name,
|
||||
prompt: job.prompt,
|
||||
shotCount: job.shotCount,
|
||||
@@ -203,7 +214,7 @@ export function summarizeStudioJob(job: StudioJob) {
|
||||
}
|
||||
|
||||
export function videoJobsBusy() {
|
||||
if (listJobs().some(job => job.kind !== 'edit' && (job.status === 'queued' || job.status === 'uploading' || job.status === 'running'))) {
|
||||
if (listJobs().some(job => job.status === 'queued' || job.status === 'uploading' || job.status === 'running')) {
|
||||
return true
|
||||
}
|
||||
return listPendingJobs().some(job => Boolean(job.promptId))
|
||||
@@ -230,15 +241,20 @@ export async function addStudioJob(params: {
|
||||
ownerKey: string
|
||||
payload: StudioJobPayload
|
||||
familyId: string
|
||||
kind?: StudioJobKind
|
||||
}) {
|
||||
const now = Date.now()
|
||||
const shotCount = 1 + (params.payload.extensions?.length || 0)
|
||||
const kind = params.kind === 'edit' ? 'edit' : 'video'
|
||||
const shotCount = kind === 'edit'
|
||||
? 1 + (params.payload.passes?.length || 0)
|
||||
: 1 + (params.payload.extensions?.length || 0)
|
||||
const job: StudioJob = {
|
||||
id: crypto.randomUUID(),
|
||||
ownerKey: params.ownerKey,
|
||||
createdAt: now,
|
||||
updatedAt: now,
|
||||
status: 'waiting',
|
||||
kind,
|
||||
name: (params.payload.name || '').trim() || params.payload.prompt.slice(0, 80),
|
||||
prompt: params.payload.prompt,
|
||||
shotCount,
|
||||
@@ -264,32 +280,46 @@ export async function patchStudioJob(owner: string, id: string, patch: (job: Stu
|
||||
}
|
||||
|
||||
export async function cancelStudioJob(owner: string, id: string) {
|
||||
const result = await mutate(owner, (jobs) => {
|
||||
const job = jobs.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.status !== 'waiting' && before.status !== 'held' && before.status !== 'running') {
|
||||
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
|
||||
|
||||
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 === 'running') {
|
||||
throw createError({ statusCode: 409, statusMessage: 'Stop the live job from Output if you want to cancel the one that is generating' })
|
||||
}
|
||||
if (job.status !== 'waiting' && job.status !== 'held') {
|
||||
throw createError({ statusCode: 409, statusMessage: 'Only a waiting or paused job can be removed from the queue' })
|
||||
if (job.status !== 'waiting' && job.status !== 'held' && job.status !== 'running') {
|
||||
throw createError({ statusCode: 409, statusMessage: 'That job is no longer in the queue' })
|
||||
}
|
||||
job.status = 'cancelled'
|
||||
job.cutIn = false
|
||||
job.holdForCutIn = false
|
||||
job.pauseAfterCurrent = false
|
||||
job.pausedByUser = false
|
||||
job.liveJobId = undefined
|
||||
job.updatedAt = Date.now()
|
||||
const stillCutIn = jobs.some(item => item.status === 'waiting' && item.cutIn)
|
||||
const stillCutIn = store.jobs.some(item => item.status === 'waiting' && item.cutIn)
|
||||
if (!stillCutIn) {
|
||||
for (const item of jobs) {
|
||||
for (const item of store.jobs) {
|
||||
if (item.status === 'running' || item.status === 'held') {
|
||||
item.holdForCutIn = false
|
||||
}
|
||||
}
|
||||
}
|
||||
syncPausedFlag(store)
|
||||
return structuredClone(job)
|
||||
})
|
||||
if (stopLive) {
|
||||
await stopLiveGeneration(liveJobId, shotQueueId)
|
||||
}
|
||||
if (result.shotQueueId) {
|
||||
const { deleteShotQueue } = await import('~/server/utils/shotQueue')
|
||||
await deleteShotQueue(owner, result.shotQueueId).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)
|
||||
if (!jobs.some(item => item.status === 'waiting' && item.cutIn)) {
|
||||
@@ -303,6 +333,44 @@ export async function cancelStudioJob(owner: string, id: string) {
|
||||
return result
|
||||
}
|
||||
|
||||
async function stopLiveGeneration(liveJobId?: string, shotQueueId?: string) {
|
||||
if (shotQueueId) {
|
||||
const { setQueueJob, pauseShotQueue } = await import('~/server/utils/shotQueue')
|
||||
setQueueJob(shotQueueId, null)
|
||||
const live = liveJobId ? getJob(liveJobId) : undefined
|
||||
if (live?.library?.ownerKey) {
|
||||
await pauseShotQueue(live.library.ownerKey, shotQueueId).catch(() => null)
|
||||
}
|
||||
}
|
||||
if (!liveJobId) return
|
||||
const live = getJob(liveJobId)
|
||||
if (live) {
|
||||
live.status = 'cancelled'
|
||||
if (live.library) {
|
||||
live.library.stopAfterCurrent = true
|
||||
live.library.queueAutoRun = false
|
||||
}
|
||||
emitJob(live, { type: 'error', error: 'Job interrupted.', message: 'Job interrupted.' })
|
||||
}
|
||||
deletePendingJob(liveJobId)
|
||||
const { interruptComfy } = await import('~/server/utils/comfy')
|
||||
await interruptComfy().catch(() => false)
|
||||
}
|
||||
|
||||
export async function cancelOrphanQueue(owner: string, queueId: string) {
|
||||
const { getShotQueue, getQueueJobId, queueIsProcessing, deleteShotQueue } = await import('~/server/utils/shotQueue')
|
||||
const queue = getShotQueue(owner, queueId)
|
||||
if (!queue) throw createError({ statusCode: 404, statusMessage: 'Queue not found' })
|
||||
const generating = queue.status === 'running' || queueIsProcessing(queueId)
|
||||
const liveJobId = getQueueJobId(queueId) || queue.currentJobId
|
||||
if (generating) {
|
||||
await stopLiveGeneration(liveJobId, queueId)
|
||||
}
|
||||
await deleteShotQueue(owner, queueId, { force: true }).catch(() => null)
|
||||
kickStudioQueue()
|
||||
return { ok: true }
|
||||
}
|
||||
|
||||
export async function requestCutIn(owner: string, id: string) {
|
||||
const jobsNow = readJobs(owner)
|
||||
const running = jobsNow.find(job => job.status === 'running')
|
||||
@@ -518,7 +586,92 @@ async function resumeHeldStudioJob(item: StudioJob) {
|
||||
}
|
||||
}
|
||||
|
||||
async function startStudioEditJob(item: StudioJob) {
|
||||
try {
|
||||
const { stillPath } = await import('~/server/utils/library')
|
||||
const { existsSync, readFileSync } = await import('node:fs')
|
||||
const { createEditLiveJob, runEdit } = await import('~/server/utils/imageChain')
|
||||
|
||||
const payload = item.payload
|
||||
if (!payload.stillId || !existsSync(stillPath(item.ownerKey, payload.stillId))) {
|
||||
throw new Error('The input still is missing from the library')
|
||||
}
|
||||
const image = {
|
||||
filename: payload.stillFilename || 'still.png',
|
||||
data: readFileSync(stillPath(item.ownerKey, payload.stillId)),
|
||||
type: 'image/png'
|
||||
}
|
||||
const refId = payload.referenceStillId
|
||||
const reference = refId && existsSync(stillPath(item.ownerKey, refId))
|
||||
? {
|
||||
filename: payload.referenceStillFilename || 'image2.png',
|
||||
data: readFileSync(stillPath(item.ownerKey, refId)),
|
||||
type: 'image/png'
|
||||
}
|
||||
: null
|
||||
|
||||
const passes = payload.passes || []
|
||||
const job = createEditLiveJob({
|
||||
steps: payload.steps,
|
||||
hideThumbnail: payload.hideThumbnail,
|
||||
library: {
|
||||
ownerKey: item.ownerKey,
|
||||
folderId: payload.folderId,
|
||||
hideThumbnail: payload.hideThumbnail,
|
||||
hideInput: payload.hideInput,
|
||||
folderLocked: payload.folderLocked,
|
||||
name: payload.name,
|
||||
prompt: payload.prompt,
|
||||
aspect: payload.aspect || 'auto',
|
||||
width: payload.width,
|
||||
height: payload.height,
|
||||
steps: payload.steps,
|
||||
turbo: true,
|
||||
seed: payload.seed || Math.floor(Math.random() * 2_147_483_647),
|
||||
cfg: payload.cfg,
|
||||
stillId: payload.stillId,
|
||||
stillFilename: payload.stillFilename,
|
||||
familyId: item.familyId,
|
||||
chainIndex: 0,
|
||||
chainStep: 1,
|
||||
chainTotal: 1 + passes.length,
|
||||
chainLabel: passes.length ? 'Pass 1' : undefined,
|
||||
passes
|
||||
}
|
||||
})
|
||||
|
||||
await markStudioLive(item.ownerKey, item.id, job.id)
|
||||
void runEdit(job, {
|
||||
image,
|
||||
reference,
|
||||
prompt: payload.prompt,
|
||||
passes,
|
||||
negative: payload.negative || '',
|
||||
steps: payload.steps,
|
||||
seed: job.library?.seed || payload.seed,
|
||||
cfg: payload.cfg
|
||||
}).catch((error) => {
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
if (job.status !== 'error' && job.status !== 'cancelled' && job.status !== 'deferred') {
|
||||
job.status = 'error'
|
||||
job.error = message
|
||||
}
|
||||
})
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error)
|
||||
await patchStudioJob(item.ownerKey, item.id, (job) => {
|
||||
job.status = 'error'
|
||||
job.lastError = message
|
||||
}).catch(() => null)
|
||||
kickStudioQueue()
|
||||
}
|
||||
}
|
||||
|
||||
export async function startStudioJob(item: StudioJob) {
|
||||
if (studioJobKind(item) === 'edit') {
|
||||
await startStudioEditJob(item)
|
||||
return
|
||||
}
|
||||
try {
|
||||
const { createJob } = await import('~/server/utils/jobs')
|
||||
const { runGeneration, frameLength } = await import('~/server/utils/videoChain')
|
||||
@@ -676,8 +829,17 @@ export async function onLiveVideoSettled(job: Job) {
|
||||
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' && item.shotQueueId && item.shotQueueId === job.library?.queueId)
|
||||
|| snapshot.jobs.find(item => item.status === 'running' && !job.library?.queueId)
|
||||
|| snapshot.jobs.find(item => (
|
||||
item.status === 'running'
|
||||
&& Boolean(item.shotQueueId)
|
||||
&& item.shotQueueId === job.library?.queueId
|
||||
))
|
||||
|| snapshot.jobs.find(item => (
|
||||
item.status === 'running'
|
||||
&& !item.shotQueueId
|
||||
&& !job.library?.queueId
|
||||
&& (!item.liveJobId || item.liveJobId === job.id)
|
||||
))
|
||||
const userPause = !failed && rowNow?.holdForCutIn !== true && (
|
||||
rowNow?.pauseAfterCurrent === true
|
||||
|| rowNow?.pausedByUser === true
|
||||
|
||||
Reference in New Issue
Block a user