import { purge } from './cleanup'; import { watchProgress } from './progress'; import { fitStill } from './media'; import { resolveRequestSize } from './size'; import { queueSeeds } from '~/shared/studio2/seed.mjs'; import { readFileSync, mkdirSync, existsSync, unlinkSync } from 'node:fs'; import { join } from 'node:path'; import template from '../../assets/studio2_minimax_native.json'; import { nativeVideoGraph, attachHeroReference, applyResolvedImageSize, sampleKleinSource } from '~/shared/studio2/graphs.mjs'; import { compilePrompt, scopedFile } from '~/shared/studio2/contracts.mjs'; import { createJob, restoreJob, getJob, emitJob, type Job } from '../jobs'; import { markStudioLive, onLiveVideoSettled, type StudioJob } from '../studioQueue'; import { readRecord, saveRecord, studio2Root, records } from './store'; import { ensureComfyReady } from '../comfyLifecycle'; import { comfyFetch, queuePrompt, fetchHistory, extractVideo } from '../comfy'; import { extractEditedImage } from '../imageComfy'; import { buildImageV2Workflow, resolveKreaGenerateAssets } from '../imageWorkflowV2'; import { resolveGraphLoraNames, applyUserLoraToGraph, ensureComfyLoraNames } from '../loras'; import { getClip, getStill, stillPath, clipVideoPath, saveStill, saveClip, downloadComfyImage, downloadComfyVideo, attachStudio2Metadata } from '../library'; import { imageDimensions } from '../resolution'; import { stitchExtension } from '../ffmpeg'; import { videoSourcePaths } from '../videoSources'; import { resolveExtensionHandoffFrame, persistClipAnchorFrame } from '../extensionFrame'; import { acquireSharedGpu, sharedGpuHeaders } from '../sharedGpu'; type DiskFile = { filename: string; subfolder: string; type: string; }; function update(r: any, state: string) { r.state = state; r.updatedAt = Date.now(); saveRecord(r); } async function upload(r: any, name: string, data: Buffer): Promise { const prefix = String(useRuntimeConfig().comfyFilenamePrefix).replace(/\/$/, '') + `/studio2/${r.id}/${r.index}`; if (name !== 'hero') data = await fitStill(data, r.request.settings.width, r.request.settings.height); const body = new FormData(); body.append('image', new Blob([new Uint8Array(data)]), name + '.png'); body.append('subfolder', prefix); body.append('type', 'input'); body.append('overwrite', 'true'); const response = await comfyFetch('/upload/image', { method: 'POST', body }); if (!response.ok) throw new Error('Reference image upload failed'); const file = await response.json() as any; const disk = { filename: file.name, subfolder: file.subfolder || '', type: 'input' }; if (!scopedFile(disk, String(useRuntimeConfig().comfyFilenamePrefix))) throw new Error('Comfy returned an input outside this instance prefix'); r.files.push(disk); saveRecord(r); return `${disk.subfolder}/${disk.filename}`; } async function prepareGraph(r: any) { const q = r.request, s = q.settings, prefix = String(useRuntimeConfig().comfyFilenamePrefix).replace(/\/$/, '') + `/studio2/${r.id}/${r.index}`; const load = async (id: string, label: string) => id ? upload(r, label, readFileSync(stillPath(r.owner, getStill(r.owner, id).id))) : ''; await resolveRequestSize(r.owner,q); saveRecord(r); const video=['video','extend'].includes(q.mode); const sourceId=q.imageAId || (!video && q.engine==='flux' && ['edit','iterate'].includes(q.mode) ? q.identityStillId : ''); const usesSourceLatent=!video && q.engine==='flux' && ['edit','iterate'].includes(q.mode) && !!sourceId; const a=await load(sourceId,'source'),b=await load(q.imageBId,'compose'),mask=await load(q.maskId,'mask'); const attachHero=(video || q.engine==='flux') && (q.lockFace!==false || q.lockOutfit!==false) && !(!video && q.mode!=='generate' && sourceId && sourceId===q.identityStillId); const hero=attachHero ? await load(q.identityStillId,'hero') : ''; r.sourceStillId=sourceId || null; r.heroReferenceAttached=!!hero || (!video && q.engine==='flux' && q.mode!=='generate' && !!sourceId && sourceId===q.identityStillId); let graph: any; if (['video', 'extend'].includes(q.mode)) { if (q.engine !== 'minimax') throw new Error('Studio 2 native identity video currently requires MiniMax. Use the existing xAIGen studio for LTX.'); let start = a; if (q.startClipId) { const dir = join(studio2Root(), r.id); mkdirSync(dir, { recursive: true }); const dest = join(dir, 'handoff.png'); await resolveExtensionHandoffFrame({ ownerKey: r.owner, sourceClipId: q.startClipId, sourceVideoPath: clipVideoPath(r.owner, q.startClipId), destPath: dest }); start = await upload(r, 'start', readFileSync(dest)); unlinkSync(dest); } const end = await load(q.endStillId, 'end'), guides = []; for (const [i, g] of q.guides.entries()) guides.push({ image: await load(g.stillId, `guide${i}`), frame: g.frame }); graph = nativeVideoGraph(template, q, { hero, start, end, guides }, prefix); await ensureComfyLoraNames('video'); resolveGraphLoraNames(graph, 'video'); applyUserLoraToGraph(graph, s.loraStack); r.graphId = 'studio2_minimax_native.json'; } else { const mode = q.mode === 'iterate' ? (a ? 'edit' : 'generate') : q.mode; const found = q.engine === 'krea' ? await resolveKreaGenerateAssets() : null; const assets = found ? { kreaUnetName: found.unet, kreaClipName: found.clip, kreaVaeName: found.vae, kreaConceptLora: found.conceptLora } : {}; graph = buildImageV2Workflow({ ...s, ...assets, engine: q.engine, mode, task: 'scene', prompt: q.compiledPrompt, imageAName: a, imageBName: mode === 'compose' ? b : undefined, maskName: mode === 'refine' ? mask : undefined, filenamePrefix: prefix + '/image' }).graph; applyResolvedImageSize(graph,s); if (usesSourceLatent) sampleKleinSource(graph,.65); if (q.engine === 'flux') attachHeroReference(graph, hero); r.sampleLatent=usesSourceLatent?'source image latent':q.mode==='refine'?'masked source latent':'empty latent'; r.sampleDenoise=usesSourceLatent ? .65 : null; if (q.engine !== 'flux' && hero && ['edit', 'compose', 'iterate'].includes(q.mode)) throw new Error('Krea native hero-reference binding is not supported by this graph. Select Klein explicitly; no fallback is performed.'); await ensureComfyLoraNames('image'); resolveGraphLoraNames(graph, 'image'); r.graphId = `${q.engine === 'flux' ? 'klein' : 'krea'}_v2_${mode}.json`; } saveRecord(r); return graph; } async function run(r: any, job: Job) { const stopProgress = watchProgress(r,job); try { const prompts = r.prompts || [r.request.promptSections, ...r.request.batch]; r.prompts = prompts; r.request.shotSeeds ||= queueSeeds(r.request.settings,prompts.length); saveRecord(r); for (r.index = r.index || 0; r.index < prompts.length; r.index++) { if (job.status === 'cancelled') throw new Error('Cancelled'); if (!r.promptId) { while (!await acquireSharedGpu()) { if (job.status === 'cancelled') throw new Error('Cancelled'); await new Promise(resolve => setTimeout(resolve, 2000)); } r.request.settings.seed = r.request.shotSeeds[r.index]; r.files = []; r.progress = null; r.request.promptSections = prompts[r.index]; r.request.compiledPrompt = compilePrompt(prompts[r.index], ['video', 'extend'].includes(r.request.mode), r.request); update(r, 'waking'); await ensureComfyReady(message => emitJob(job, { type: 'status', message })); if (job.status === 'cancelled') throw new Error('Cancelled'); const graph = await prepareGraph(r); if (job.status === 'cancelled') throw new Error('Cancelled'); update(r, 'submitting'); const queued = await queuePrompt(graph, job.clientId); job.promptId = queued.prompt_id; r.promptId = queued.prompt_id; r.renderStartedAt = Date.now(); update(r, 'rendering'); } let history: any, missingSince = 0; for (;;) { try { history = await fetchHistory(r.promptId); } catch { history = null; } const e = history?.[r.promptId]; if (job.status === 'cancelled') throw new Error('Cancelled'); if (['error', 'interrupted'].includes(e?.status?.status_str)) throw new Error('Comfy generation failed or was interrupted; inspect the host error log.'); if (!e) { try { const response = await comfyFetch('/queue', { signal: AbortSignal.timeout(5000) }); if (response.ok) { const q = await response.json() as any; const present = [...(q.queue_running || []), ...(q.queue_pending || [])].some((entry: any) => entry[1] === r.promptId); if (present) missingSince = 0; else if (!missingSince) missingSince = Date.now(); else if (Date.now() - missingSince > 60000) throw new Error('PROMPT_MISSING'); } } catch (error: any) { if (error.message === 'PROMPT_MISSING') throw new Error('The saved Comfy prompt is no longer in history or the queue. It was not resubmitted.'); } } if (e?.status?.completed || e?.status?.status_str === 'success') break; if (job.status === 'cancelled') throw new Error('Cancelled'); await new Promise(resolve => setTimeout(resolve, 1500)); } const messages = history[r.promptId]?.status?.messages || []; const executionStart = messages.find((m: any) => m[0] === 'execution_start')?.[1]?.timestamp; const executionEnd = messages.find((m: any) => m[0] === 'execution_success')?.[1]?.timestamp; r.gpuSeconds = executionStart && executionEnd ? (executionEnd - executionStart) / 1000 : null; update(r, 'saving'); job.saving = true; const q = r.request, s = q.settings, video = ['video', 'extend'].includes(q.mode); const file = video ? extractVideo(history, r.promptId) : extractEditedImage(history, r.promptId); if (!file) throw new Error('Comfy finished without a media output'); if (!r.files.some((f: DiskFile) => f.filename === file.filename && f.subfolder === file.subfolder)) r.files.push(file); let asset: any = r.savedAssetId ? (video ? getClip(r.owner, r.savedAssetId) : getStill(r.owner, r.savedAssetId)) : null; if (!asset && video) { const data = await downloadComfyVideo(file); const sourceSegments = q.startClipId ? videoSourcePaths(clipVideoPath(r.owner, q.startClipId)) : []; const tmp = join(studio2Root(), r.id); mkdirSync(tmp, { recursive: true }); const assembled = q.startClipId ? await stitchExtension({ part1Path: clipVideoPath(r.owner, q.startClipId), part2: data, tmpDir: tmp, sourcePaths: sourceSegments }) : data; asset = await saveClip({ ownerKey: r.owner, folderId: q.folderId, prompt: q.compiledPrompt, ...s, aspect: s.aspect || 'auto', hideThumbnail: false, video: assembled, originalSegment: data, sourceSegments, fps: s.fps || 24, sound: true, familyId: r.familyId, parentClipId: q.startClipId || undefined, chainIndex: r.index, comfyFilename: file.filename }); } else if (!asset) { const data = await downloadComfyImage(file), size = imageDimensions(data) || { width: s.width, height: s.height }; asset = await saveStill({ ownerKey: r.owner, folderId: q.folderId, filename: file.filename, data, ...size, role: 'output', prompt: q.compiledPrompt, familyId: r.familyId, chainIndex: r.index, settings: s }); } if (!asset) throw new Error('Library did not save the output'); r.savedAssetId = asset.id; saveRecord(r); r.savedAt = Date.now(); r.wallTime = (r.savedAt - r.startedAt) / 1000; const metadata = { ...structuredClone(q), id: r.id, sourceStillId:r.sourceStillId,heroReferenceAttached:r.heroReferenceAttached,sampleLatent:r.sampleLatent,sampleDenoise:r.sampleDenoise,kind: video ? 'video' : 'image', graphId: r.graphId, promptId: r.promptId, queuedAt: r.queuedAt, startedAt: r.startedAt, savedAt: r.savedAt, gpuSeconds: r.gpuSeconds, wallTime: r.wallTime, outputWidth: asset.width || s.width, outputHeight: asset.height || s.height, purgeResult: 'Left on host' }; await attachStudio2Metadata(r.owner, asset.id, metadata); if (!r.outputs.some((o:any)=>o.id===asset.id)) r.outputs.push({id:asset.id,kind:video?'clip':'still',studio2:metadata}); saveRecord(r); const anchor = video && history[r.promptId].outputs?.anchor_save?.images?.[0]; if (anchor) { if (!r.files.some((f: DiskFile)=>f.filename===anchor.filename && f.subfolder===anchor.subfolder)) r.files.push(anchor); await persistClipAnchorFrame({ownerKey:r.owner,clipId:asset.id,fromBuffer:await downloadComfyImage(anchor)}); } await purge(r); metadata.purgeResult = r.purgeResult; await attachStudio2Metadata(r.owner, asset.id, metadata); saveRecord(r); job.saving = false; if (video) { q.startClipId = asset.id; q.startFrameSource = { kind: 'previous-last-frame', clipId: asset.id }; } else if (q.mode === 'iterate') q.imageAId = asset.id; r.promptId = ''; r.savedAssetId = ''; r.index++; saveRecord(r); r.index--; // Identity remains the original hero throughout the batch. } if (job.status === 'cancelled') throw new Error('Cancelled'); job.status = 'complete'; update(r, 'complete'); emitJob(job, { type: 'complete', message: 'Studio 2 output saved' }); } catch (e: any) { job.status = job.status === 'cancelled' ? 'cancelled' : 'error'; job.error = e.message; r.error = e.message; update(r, job.status === 'cancelled' ? 'cancelled' : 'failed'); emitJob(job, { type: 'error', message: e.message }); } finally { stopProgress(); job.saving = false; await onLiveVideoSettled(job); } } export async function startStudio2Job(item: StudioJob) { const r = readRecord(item.ownerKey, item.payload.studio2Id!), job = createJob(['video', 'extend'].includes(r.request.mode) ? 'video' : 'edit'); job.studio2 = true; job.library = { ownerKey: item.ownerKey, folderId: r.request.folderId, extensions: [], queueAutoRun: false } as any; job.status = 'running'; r.startedAt = Date.now(); r.liveId = job.id; r.clientId = job.clientId; saveRecord(r); await markStudioLive(item.ownerKey, item.id, job.id); void run(r, job); } export function resumeStudio2Jobs() { for (const r of records()) { if (!r.liveId || getJob(r.liveId) || ['waiting', 'complete', 'failed', 'cancelled'].includes(r.state)) continue; if (r.state === 'submitting' && !r.promptId) { r.state = 'failed'; r.error = 'Submission outcome unknown after restart. Inspect the host before retrying.'; saveRecord(r); continue; } const job = restoreJob({ id: r.liveId, clientId: r.clientId, startedAt: r.startedAt, promptId: r.promptId || '', library: { ownerKey: r.owner, folderId: r.request.folderId, extensions: [], queueAutoRun: false } as any }); job.studio2 = true; void run(r, job); } }