From fb748e9b6eb7ad53225a4c0e95da3af267548f1f Mon Sep 17 00:00:00 2001 From: Towsty Date: Fri, 28 Aug 2026 09:04:08 -0500 Subject: [PATCH] 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 --- components/ClipFrameThumbs.vue | 87 ++++---- pages/index.vue | 347 ++++++++++++++++++++++++++----- pages/queue.vue | 201 +++++++++++++++++- server/api/edit.post.ts | 281 ++++--------------------- server/api/queues/[id].delete.ts | 4 +- server/api/studio-queue.get.ts | 16 +- server/utils/comfyLifecycle.ts | 15 +- server/utils/imageChain.ts | 213 +++++++++++++++++++ server/utils/shotQueue.ts | 5 +- server/utils/studioQueue.ts | 196 +++++++++++++++-- 10 files changed, 996 insertions(+), 369 deletions(-) create mode 100644 server/utils/imageChain.ts diff --git a/components/ClipFrameThumbs.vue b/components/ClipFrameThumbs.vue index 7e49ace..77d229c 100644 --- a/components/ClipFrameThumbs.vue +++ b/components/ClipFrameThumbs.vue @@ -1,15 +1,17 @@ @@ -360,10 +407,15 @@ type Queue = { type StudioJobRow = { id: string status: string + kind?: string name: string + prompt?: string shotCount: number liveJobId?: string shotQueueId?: string + stillId?: string + hideThumbnail?: boolean + folderLocked?: boolean cutIn?: boolean holdForCutIn?: boolean pauseAfterCurrent?: boolean @@ -374,6 +426,15 @@ type StudioJobRow = { completedCount?: number } +type PendingRemove = { + title: string + name: string + detail: string + consequence: string + job?: StudioJobRow + queue?: Queue +} + const instanceName = ref(String(useRuntimeConfig().public.instanceName || 'AIGen')) const authEnabled = ref(false) const jobs = ref([]) @@ -386,6 +447,7 @@ const busyId = ref('') const liveMessage = ref>({}) const dirty = new Set() const sources = new Map() +const pendingRemove = ref(null) let pollTimer: ReturnType | null = null const canCutIn = computed(() => jobs.value.some(job => job.status === 'running' || job.status === 'held')) @@ -452,6 +514,24 @@ function progressLine(queue: Queue) { } function jobLine(job: StudioJobRow) { + if (job.kind === 'edit') { + const passes = job.shotCount || 1 + const passLine = passes > 1 ? `${passes} passes` : 'image edit' + if (job.status === 'waiting') { + return job.cutIn + ? `Next after the current shot · ${passLine}` + : `Waiting · ${passLine}` + } + if (job.status === 'held') { + if (jobIsHeldPaused(job)) return 'Paused. This edit will not continue until you resume.' + return job.holdForCutIn ? 'Paused so another job can cut in' : 'Paused' + } + if (job.status === 'running') { + if (jobIsLatched(job)) return `Editing · will pause after this · ${passLine}` + return `Editing · ${passLine}` + } + return job.status + } if (job.status === 'waiting') { return job.cutIn ? `Next after the current shot · ${job.shotCount} shot${job.shotCount === 1 ? '' : 's'} · own clip, not part of the current sequence` @@ -559,6 +639,9 @@ async function loadQueues() { for (const queue of [...nextJobs.map(job => job.shots).filter(Boolean), ...orphans.value] as Queue[]) { if (queue.status === 'running' && queue.currentJobId) listen(queue) } + for (const job of nextJobs) { + if (job.status === 'running' && job.liveJobId && !job.shots) listenJob(job) + } } function listen(queue: Queue) { @@ -593,8 +676,101 @@ async function cutIn(job: StudioJobRow) { } } +function listenJob(job: StudioJobRow) { + if (!job.liveJobId || sources.has(job.id)) return + const source = new EventSource(`/api/generate/${job.liveJobId}/stream`) + sources.set(job.id, source) + source.onmessage = (event) => { + try { + const payload = JSON.parse(event.data) as { message?: string; type?: string } + if (payload.message) liveMessage.value = { ...liveMessage.value, [job.id]: payload.message } + if (payload.type === 'checkpoint' || payload.type === 'complete' || payload.type === 'error') { + void loadQueues() + } + if (payload.type === 'complete' || payload.type === 'error') { + source.close() + sources.delete(job.id) + } + } catch { /* ignore */ } + } +} + +function jobKindLabel(job: StudioJobRow) { + if (job.kind === 'edit') return (job.shotCount || 1) > 1 ? 'EDIT' : 'IMAGE' + return 'Job' +} + +function jobThumb(job: StudioJobRow) { + if (!job.stillId || job.hideThumbnail || job.folderLocked) return '' + return `/api/library/stills/${job.stillId}` +} + +function canRemoveJob(job: StudioJobRow) { + return job.status === 'waiting' || job.status === 'held' || job.status === 'running' +} + +function promptSnippet(text?: string) { + const value = String(text || '').replace(/\s+/g, ' ').trim() + if (!value) return '' + return value.length > 120 ? `${value.slice(0, 117)}…` : value +} + +function shotProgressDetail(job: StudioJobRow) { + const queue = job.shots + if (!queue) return job.kind === 'edit' ? promptSnippet(job.prompt) : '' + const total = queue.totalCount || job.shotCount || 1 + const next = nextSegment(queue) + const shot = next ? next.index + 1 : Math.min(total, (queue.completedCount || 0) + 1) + const snippet = promptSnippet(next?.prompt || job.prompt) + return `Shot ${shot} of ${total}${snippet ? ` · ${snippet}` : ''}` +} + +function askRemoveJob(job: StudioJobRow) { + const generating = jobIsGenerating(job) + const pausedBatch = jobIsHeldPaused(job) || (job.status === 'held' && Boolean(job.shotQueueId)) + pendingRemove.value = { + title: generating ? 'Cancel this job?' : 'Remove this job?', + name: job.name || 'Untitled job', + detail: shotProgressDetail(job), + consequence: generating + ? 'This stops the job that’s generating on Beast and removes only this row. Other queue items stay. Saved clips stay in the library.' + : pausedBatch + ? 'This drops the rest of this shot list only. It will not resume, and it will not stop a different job that’s generating. Saved clips stay in the library.' + : 'This removes only this job from the queue. Other items stay. Saved clips stay in the library.', + job + } +} + +function askRemoveQueue(queue: Queue) { + const generating = queue.status === 'running' + const next = nextSegment(queue) + const shot = next ? next.index + 1 : Math.min(queue.totalCount, (queue.completedCount || 0) + 1) + pendingRemove.value = { + title: generating ? 'Cancel this job?' : 'Remove this job?', + name: queue.name || 'Shot batch', + detail: `Shot ${shot} of ${queue.totalCount}${next?.prompt ? ` · ${promptSnippet(next.prompt)}` : ''}`, + consequence: generating + ? 'This stops the job that’s generating on Beast and removes only this batch. Other queue items stay. Saved clips stay in the library.' + : 'This drops the rest of this shot list only. It will not resume, and it will not stop a different job that’s generating. Saved clips stay in the library.', + queue + } +} + +async function confirmRemove() { + const pending = pendingRemove.value + if (!pending) return + pendingRemove.value = null + busyId.value = pending.job?.id || pending.queue?.id || 'remove' + error.value = '' + try { + if (pending.job) await removeJob(pending.job) + else if (pending.queue) await remove(pending.queue) + } finally { + busyId.value = '' + } +} + async function removeJob(job: StudioJobRow) { - if (!confirm('Remove this job from the queue? Saved clips stay in the library.')) return await $fetch(`/api/studio-queue/${job.id}`, { method: 'DELETE' }).catch((err: any) => { error.value = err?.data?.statusMessage || err?.message || 'Could not remove that job' }) @@ -739,7 +915,6 @@ async function togglePause(target?: { job?: StudioJobRow; queue?: Queue }) { } async function remove(queue: Queue) { - if (!confirm('Remove this batch from the queue? Saved clips stay in the library.')) return await $fetch(`/api/queues/${queue.id}`, { method: 'DELETE' }).catch((err: any) => { error.value = err?.data?.statusMessage || err?.message || 'Could not remove that batch' }) @@ -785,6 +960,7 @@ onMounted(async () => { const me = await $fetch<{ user?: { name?: string }; authEnabled?: boolean; instanceName?: string }>('/api/auth/me').catch(() => ({ user: null })) authEnabled.value = Boolean(me.authEnabled) instanceName.value = me.instanceName || instanceName.value + window.addEventListener('keydown', onRemoveKey) await loadQueues().catch((err: any) => { error.value = err?.data?.statusMessage || err?.message || 'Could not load the queue' }) @@ -794,7 +970,12 @@ onMounted(async () => { }, 4000) }) +function onRemoveKey(event: KeyboardEvent) { + if (event.key === 'Escape' && pendingRemove.value) pendingRemove.value = null +} + onBeforeUnmount(() => { + window.removeEventListener('keydown', onRemoveKey) if (pollTimer) clearInterval(pollTimer) if (saveTimer) clearTimeout(saveTimer) for (const source of sources.values()) source.close() diff --git a/server/api/edit.post.ts b/server/api/edit.post.ts index bafa0ce..158448c 100644 --- a/server/api/edit.post.ts +++ b/server/api/edit.post.ts @@ -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, - 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 - }) - } - }) -} diff --git a/server/api/queues/[id].delete.ts b/server/api/queues/[id].delete.ts index 908aba9..3687003 100644 --- a/server/api/queues/[id].delete.ts +++ b/server/api/queues/[id].delete.ts @@ -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 } }) diff --git a/server/api/studio-queue.get.ts b/server/api/studio-queue.get.ts index e5256e8..5e312b0 100644 --- a/server/api/studio-queue.get.ts +++ b/server/api/studio-queue.get.ts @@ -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 ? { diff --git a/server/utils/comfyLifecycle.ts b/server/utils/comfyLifecycle.ts index 2905614..3ad1ebd 100644 --- a/server/utils/comfyLifecycle.ts +++ b/server/utils/comfyLifecycle.ts @@ -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 } diff --git a/server/utils/imageChain.ts b/server/utils/imageChain.ts new file mode 100644 index 0000000..a9d268b --- /dev/null +++ b/server/utils/imageChain.ts @@ -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 +}) { + const job = createJob() + job.kind = 'edit' + job.maxStep = params.steps + job.hideThumbnail = params.hideThumbnail + job.library = params.library + return job +} diff --git a/server/utils/shotQueue.ts b/server/utils/shotQueue.ts index b5a073a..34db11c 100644 --- a/server/utils/shotQueue.ts +++ b/server/utils/shotQueue.ts @@ -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' }) diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts index 21285e7..e779404 100644 --- a/server/utils/studioQueue.ts +++ b/server/utils/studioQueue.ts @@ -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 | { 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