Files
aigen/server/utils/studio2/runner.ts
T
TowstyandCursor 50319e03a9 Map Photos roles to Qwen image_1/image_2 and inject encode tags silently.
Mention phrases become <imageN> only in the TextEncode string; the inspector never shows Comfy tokens.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-26 16:01:49 -05:00

617 lines
33 KiB
TypeScript

import { kleinIdentityPlan, applyKleinIdentity } from '~/shared/studio2/klein-identity.mjs';
import { stylePrompt } from '~/shared/studio2/styles.mjs';
import { purge } from './cleanup';
import { watchProgress } from './progress';
import { fitStill, needsFit } from './media';
import { resolveRequestSize } from './size';
import { queueSeeds } from '~/shared/studio2/seed.mjs';
import { readFileSync, writeFileSync, mkdirSync, existsSync, unlinkSync } from 'node:fs';
import { join } from 'node:path';
import template from '../../assets/studio2_minimax_native.json';
import qwen21Template from '../../assets/studio2_qwen21_t2i.json'
import qwen21EditTemplate from '../../assets/studio2_qwen21_edit.json'
import qwen21T2iTurboTemplate from '../../assets/studio2_qwen21_t2i_turbo.json'
import qwen21EditTurboTemplate from '../../assets/studio2_qwen21_edit_turbo.json'
import qwen21PeT2iTemplate from '../../assets/studio2_qwen21_pe_t2i.json'
import qwen21PeEditTemplate from '../../assets/studio2_qwen21_pe_edit.json';
import { nativeVideoGraph, attachHeroReference, applyResolvedImageSize } from '~/shared/studio2/graphs.mjs';
import { compilePrompt, scopedFile } from '~/shared/studio2/contracts.mjs';
import { resolveQwen21Size } from '~/shared/studio2/qwen21-size.mjs';
import { ensureQwen21EditPrompt, stitchQwen21PeEditPrompt } from '~/shared/studio2/qwen21-edit.mjs';
import { applyPhotosToRequest, injectQwenEditMentionTags } from '~/shared/studio2/photos.mjs';
import { qwen21PeRefusal } from '~/shared/studio2/qwen21-pe.mjs';
import { QWEN21_TURBO_LORA, QWEN21_TURBO_SIGMAS, QWEN21_TURBO_STEPS } from '~/shared/studio2/qwen21-turbo.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, freeComfyVram } 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, probeDuration } 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(); if (['complete','failed','cancelled'].includes(state)) r.finishedAt ??= r.updatedAt; saveRecord(r); }
// Labels whose pixels must never be resampled: the hero identity still, and the
// extend/video start frame (which is either the exact anchor PNG or an already-canvas-sized
// extracted last frame). Fitting these every hop is what destroys identity across chained extends.
const NEVER_FIT = new Set(['hero', 'start']);
function previewAnyText(outputs: Record<string, any> | undefined, nodeId: string): string {
const node = outputs?.[nodeId]
if (!node || typeof node !== 'object') return ''
const raw = node.text ?? node.STRING ?? node.source
if (Array.isArray(raw)) return String(raw[0] ?? '')
if (raw == null) return ''
return String(raw)
}
function harvestQwen21Pe(history: Record<string, any> | null | undefined, promptId: string, edit: boolean) {
const outputs = history?.[promptId]?.outputs as Record<string, any> | undefined
const positive_prompt = previewAnyText(outputs, '30').trim()
const wh_ratio = previewAnyText(outputs, '31').trim()
const ratio_follow = edit ? previewAnyText(outputs, '32').trim() : ''
const thinking = previewAnyText(outputs, edit ? '33' : '32').trim()
const parseRaw = previewAnyText(outputs, edit ? '34' : '33').trim()
const parse_ok = /^(true|1|yes)$/i.test(parseRaw)
return { positive_prompt, wh_ratio, ratio_follow, thinking, parse_ok }
}
async function waitPromptHistory(r: any, job: Job, promptId: string) {
let history: any = null, missingSince = 0
for (;;) {
try { history = await fetchHistory(promptId) }
catch { history = null }
const e = history?.[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] === 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))
}
return history
}
function finalizeQwen21EditSample(rawPrompt: string, request: any) {
const photos = Array.isArray(request?.photos) ? request.photos : []
const tagged = injectQwenEditMentionTags(String(rawPrompt || ''), photos)
return ensureQwen21EditPrompt(tagged, { hasImage2: !!request?.imageBId })
}
async function runQwen21PromptEnhance(r: any, job: Job) {
const q = r.request, s = q.settings
if (!q.enhancePrompt || q.engine !== 'qwen21' || !['generate', 'edit'].includes(q.mode)) return
const edit = q.mode === 'edit'
if (Array.isArray(q.photos) && q.photos.length) {
const bound = applyPhotosToRequest(q)
q.imageAId = bound.imageAId
q.imageBId = bound.imageBId
q.identityStillId = bound.identityStillId
q.photos = bound.photos
}
// Edit → PE-I2I only (pe_i2i). Generate → PE-T2I only. Never cross-feed.
if (edit && !q.imageAId)
throw new Error('Add a photo and mark it Photo to change.')
update(r, 'enhancing')
emitJob(job, { type: 'status', message: 'Enhancing prompt…' })
// promptRaw = exactly what they typed (compiled from sections before PE).
const userPrompt = String(q.compiledPrompt || '')
r.request.promptRaw = userPrompt
// Edit PE gets the raw edit instruction + canvas still — not the keep stanza.
const peInput = userPrompt
const graph = structuredClone(edit ? qwen21PeEditTemplate : qwen21PeT2iTemplate)
graph['9'].inputs.prompt = peInput
graph['9'].inputs.seed = s.seed
if (edit) {
// Photo to change → image_1. Outfit/extra → image_2 when present.
const a = await upload(r, 'pe-source', readFileSync(stillPath(r.owner, getStill(r.owner, q.imageAId).id)))
graph['10'].inputs.image = a
graph['9'].inputs.image_1 = ['10', 0]
if (q.imageBId) {
const b = await upload(r, 'pe-compose', readFileSync(stillPath(r.owner, getStill(r.owner, q.imageBId).id)))
graph['11'] = {
inputs: { image: b },
class_type: 'LoadImage',
_meta: { title: 'image_2 optional' }
}
graph['9'].inputs.image_2 = ['11', 0]
}
}
r.peGraphId = edit ? 'studio2_qwen21_pe_edit.json' : 'studio2_qwen21_pe_t2i.json'
// PE is stage 1 of the same studio job — not a complete/requeue.
r.progress = null
saveRecord(r)
const queued = await queuePrompt(graph, job.clientId)
r.pePromptId = queued.prompt_id
// watchProgress listens on promptId; point it at the PE Comfy prompt for this stage.
r.promptId = queued.prompt_id
job.promptId = queued.prompt_id
saveRecord(r)
const history = await waitPromptHistory(r, job, queued.prompt_id)
const result = harvestQwen21Pe(history, queued.prompt_id, edit)
const refusal = qwen21PeRefusal(userPrompt, result)
if (refusal.refused) {
// Fail open: keep the job, send the raw brief (Edit keep-identity stanza still applies).
const compiled = edit ? finalizeQwen21EditSample(userPrompt, q) : userPrompt
r.request.compiledPrompt = compiled
r.request.prompt = compiled
r.request.enhance = {
used: false,
wh_ratio: '',
...(edit ? { ratio_follow: '' } : {}),
parse_ok: !!result.parse_ok,
refused: true,
skipReason: refusal.reason,
thinking: result.thinking || '',
comfyPromptId: queued.prompt_id,
}
} else if (edit) {
// Edit only: stanza + raw first; PE chunk only if it is an edit directive (never PE-alone).
// Generate keeps the long observer paragraph as the product — no describe sanitize.
const stitched = stitchQwen21PeEditPrompt(userPrompt, result.positive_prompt)
const compiled = finalizeQwen21EditSample(stitched.prompt, q)
r.request.compiledPrompt = compiled
r.request.prompt = compiled
r.request.enhance = {
used: !stitched.skippedAsDescribe,
wh_ratio: result.wh_ratio || '',
ratio_follow: result.ratio_follow || '',
parse_ok: !!result.parse_ok,
refused: false,
...(stitched.skippedAsDescribe ? { skippedAsDescribe: true } : {}),
thinking: result.thinking || '',
comfyPromptId: queued.prompt_id,
}
} else {
r.request.compiledPrompt = result.positive_prompt
r.request.prompt = result.positive_prompt
r.request.enhance = {
used: true,
wh_ratio: result.wh_ratio || '',
parse_ok: !!result.parse_ok,
refused: false,
thinking: result.thinking || '',
comfyPromptId: queued.prompt_id,
}
}
if (edit || q.mode === 'generate') r.request.task = q.mode
// Stage flip PE → generate: clear PE prompt id and reset the bar (no leftover 100%).
r.pePromptId = ''
r.promptId = ''
job.promptId = ''
r.progress = null
saveRecord(r)
await freeComfyVram()
// Same job continues into generate — do not complete or requeue.
update(r, 'submitting')
emitJob(job, { type: 'status', message: 'Generating…' })
}
async function upload(r: any, name: string, data: Buffer): Promise<string> {
const prefix = String(useRuntimeConfig().comfyFilenamePrefix).replace(/\/$/, '') + `/studio2/${r.id}/${r.index}`;
let fitted = false;
if (!NEVER_FIT.has(name) && r.request?.engine !== 'qwen21' && needsFit(data, r.request.settings.width, r.request.settings.height)) {
data = await fitStill(data, r.request.settings.width, r.request.settings.height);
fitted = true;
}
if (name === 'start') { r.handoffFitted = fitted; saveRecord(r); }
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}`;
// Photos roles → graph sockets (imageAId / imageBId / identityStillId). Payload is source of truth.
if (Array.isArray(q.photos) && q.photos.length) {
const bound = applyPhotosToRequest(q)
q.imageAId = bound.imageAId
q.imageBId = bound.imageBId
q.identityStillId = bound.identityStillId
q.photos = bound.photos
}
const load = async (id: string, label: string) => id ? upload(r, label, readFileSync(stillPath(r.owner, getStill(r.owner, id).id))) : '';
const video=['video','extend'].includes(q.mode);
const identityPlan=kleinIdentityPlan(q);
// Qwen Edit: Photo to change (imageAId) → image_1. Never Klein hero-ref.
const sourceId = q.engine === 'qwen21'
? (q.imageAId || '')
: (identityPlan ? identityPlan.sourceId : q.imageAId);
await resolveRequestSize(r.owner,identityPlan ? {...q,imageAId:sourceId} : q);
saveRecord(r);
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 && sourceId && sourceId===q.identityStillId);
const hero=attachHero ? await load(q.identityStillId,'hero') : '';
r.sourceStillId=sourceId || null;
r.heroReferenceAttached=!!hero || (!video && q.engine==='flux' && !!sourceId && sourceId===q.identityStillId);
let graph: any;
if (['video', 'extend'].includes(q.mode)) {
let start = a;
if (q.startClipId) {
const dir = join(studio2Root(), r.id);
mkdirSync(dir, { recursive: true });
const dest = join(dir, 'handoff.png');
const handoff = await resolveExtensionHandoffFrame({ ownerKey: r.owner, sourceClipId: q.startClipId, sourceVideoPath: clipVideoPath(r.owner, q.startClipId), destPath: dest });
r.handoffSource = handoff.source;
start = await upload(r, 'start', readFileSync(dest));
unlinkSync(dest);
}
if (q.engine === 'ltx') {
const { buildWorkflow } = await import('../workflow');
const { frameLength } = await import('../videoChain');
const { ltxWorkflowEnabled, LTX_DISABLED_MESSAGE } = await import('~/utils/videoModels');
if (!ltxWorkflowEnabled()) throw new Error(LTX_DISABLED_MESSAGE);
const textToVideo = q.mode === 'video' && !start;
if (!start && !textToVideo) throw new Error('Choose a start still for LTX video.');
const fps = s.fps || 24;
graph = buildWorkflow({
workflow: textToVideo ? 'ltx-t2v' : 'ltx',
imageName: start || '',
prompt: q.compiledPrompt,
width: s.width,
height: s.height,
length: frameLength(s.duration || 5, fps),
fps,
steps: s.steps,
cfg: s.cfg,
seed: s.seed,
samplerName: 'euler',
scheduler: 'simple',
turbo: false,
duration: s.duration || 5,
filenamePrefix: prefix + '/video',
loraStack: s.loraStack,
useIdentityRefs: false
});
await ensureComfyLoraNames('video');
r.graphId = textToVideo ? 'workflow_ltx_video.json#t2v' : 'workflow_ltx_video.json';
} else {
if (q.engine !== 'minimax')
throw new Error('Studio 2 video supports MiniMax H3 and LTX (xAIGen).');
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 if (q.engine === 'qwen21') {
if (!['generate', 'edit'].includes(q.mode))
throw new Error('Qwen 2.1 supports Generate and Edit.')
const turbo = !!q.turbo
if (turbo) {
q.lora = QWEN21_TURBO_LORA
q.sigmas = QWEN21_TURBO_SIGMAS
s.steps = QWEN21_TURBO_STEPS
s.cfg = 1
}
const negative = turbo ? '' : stylePrompt(q.imageStyles, true)
if (q.mode === 'edit') {
// Photo to change → image_1. Outfit/extra → image_2. Never T2I / EmptyLatent.
if (!q.imageAId || !a)
throw new Error('Add a photo and mark it Photo to change.')
// promptRaw stays typed; prompt/compiledPrompt = what TextEncode samples (tags injected silently).
if (q.promptRaw == null) q.promptRaw = String(q.compiledPrompt || '')
const prompt = finalizeQwen21EditSample(String(q.compiledPrompt || ''), q)
q.compiledPrompt = prompt
q.prompt = prompt
q.task = 'edit'
graph = structuredClone(turbo ? qwen21EditTurboTemplate : qwen21EditTemplate)
graph['10'].inputs.image = a
graph['9'].inputs['images.image_1'] = ['10', 0]
if (b) {
graph['11'] = {
inputs: { image: b },
class_type: 'LoadImage',
_meta: { title: 'image_2 optional' }
}
graph['9'].inputs['images.image_2'] = ['11', 0]
}
graph['9'].inputs.prompt = prompt
graph['9'].inputs.negative_prompt = negative
// Official edit: sampler latent comes from TextEncode (sized from image_1).
graph['9'].inputs.resolution = 1024
if (turbo) {
graph['7'].inputs.lora_name = QWEN21_TURBO_LORA
graph['7'].inputs.strength = 1.0
graph['31'].inputs.noise_seed = s.seed
graph['33'].inputs.nodes = QWEN21_TURBO_SIGMAS
graph['36'].inputs.filename_prefix = prefix + '/image'
r.sampleLatent = 'qwen21 edit latent'
r.graphId = 'studio2_qwen21_edit_turbo.json'
} else {
graph['15'].inputs.seed = s.seed
graph['15'].inputs.steps = s.steps || 25
graph['15'].inputs.cfg = s.cfg ?? 1
graph['15'].inputs.sampler_name = 'euler'
graph['15'].inputs.scheduler = 'simple'
graph['21'].inputs.filename_prefix = prefix + '/image'
r.sampleLatent = 'qwen21 edit latent'
r.graphId = 'studio2_qwen21_edit.json'
}
r.sampleDenoise = null
r.heroReferenceAttached = false
} else {
// Bench aspect is law for Generate. PE wh_ratio is advisory only (never sizes the canvas).
if (q.promptRaw == null) q.promptRaw = String(q.compiledPrompt || '')
q.prompt = String(q.compiledPrompt || '')
q.task = 'generate'
const size = resolveQwen21Size(s.aspect)
s.width = size.width
s.height = size.height
graph = structuredClone(turbo ? qwen21T2iTurboTemplate : qwen21Template)
graph['9'].inputs.prompt = q.compiledPrompt
graph['9'].inputs.negative_prompt = negative
graph['9'].inputs.resolution = 1024
graph['16'].inputs.width = size.width
graph['16'].inputs.height = size.height
if (turbo) {
graph['7'].inputs.lora_name = QWEN21_TURBO_LORA
graph['7'].inputs.strength = 1.0
graph['31'].inputs.noise_seed = s.seed
graph['33'].inputs.nodes = QWEN21_TURBO_SIGMAS
graph['36'].inputs.filename_prefix = prefix + '/image'
r.sampleLatent = 'qwen21 empty latent'
r.graphId = 'studio2_qwen21_t2i_turbo.json'
} else {
graph['15'].inputs.seed = s.seed
graph['15'].inputs.steps = s.steps || 25
graph['15'].inputs.cfg = s.cfg ?? 1
graph['15'].inputs.sampler_name = 'euler'
graph['15'].inputs.scheduler = 'simple'
graph['21'].inputs.filename_prefix = prefix + '/image'
r.sampleLatent = 'qwen21 empty latent'
r.graphId = 'studio2_qwen21_t2i.json'
}
r.sampleDenoise = null
r.heroReferenceAttached = false
}
}
else {
const mode = identityPlan?.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, negative: stylePrompt(q.imageStyles, true), imageAName: a, imageBName: mode === 'compose' ? b : undefined, maskName: mode === 'refine' ? mask : undefined, filenamePrefix: prefix + '/image' }).graph;
applyResolvedImageSize(graph,s);
if (identityPlan) Object.assign(r,applyKleinIdentity(graph,identityPlan));
if (q.engine === 'flux')
attachHeroReference(graph, hero);
if (!identityPlan) {
r.sampleLatent=q.mode==='refine'?'masked source latent':'empty latent';
r.sampleDenoise=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');
await runQwen21PromptEnhance(r, job);
if (job.status === 'cancelled')
throw new Error('Cancelled');
const graph = await prepareGraph(r);
if (job.status === 'cancelled')
throw new Error('Cancelled');
update(r, 'submitting');
emitJob(job, { type: 'status', message: 'Generating…' });
let queued: { prompt_id: string; number?: number }
try {
queued = await queuePrompt(graph, job.clientId);
} catch (error: any) {
const msg = String(error?.statusMessage || error?.message || error)
if (r.request?.turbo && /lora|ViggleTurbo|not in list|value_not_in_list/i.test(msg))
throw new Error(`Qwen Turbo LoRA missing on host (${QWEN21_TURBO_LORA}). Install the r128 file under models/loras and restart Comfy.`)
throw error
}
job.promptId = queued.prompt_id;
r.promptId = queued.prompt_id;
r.generateComfyPromptId = queued.prompt_id;
r.progress = null;
r.renderStartedAt = Date.now();
update(r, 'rendering');
}
const history = await waitPromptHistory(r, job, r.promptId);
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);
if (!data.length || data.length < 64)
throw new Error('Comfy wrote an empty video file (<64 bytes); the job did not render. Nothing was saved.');
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;
if (!assembled.length || assembled.length < 64)
throw new Error('The assembled video file is empty (<64 bytes); the job did not render. Nothing was saved.');
const checkPath = join(tmp, 'duration-check.mp4');
writeFileSync(checkPath, assembled);
let duration = 0;
try { duration = await probeDuration(checkPath); }
finally { try { unlinkSync(checkPath); } catch { /* ignore */ } }
if (!Number.isFinite(duration) || duration <= 0)
throw new Error('The rendered video has no readable duration; the job did not render. Nothing was saved.');
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 };
// prompt = TextEncode string (after PE + stitch); promptRaw = typed brief.
const samplePrompt = String(q.prompt || q.compiledPrompt || '')
asset = await saveStill({
ownerKey: r.owner,
folderId: q.folderId,
filename: file.filename,
data,
...size,
role: 'output',
prompt: samplePrompt,
promptRaw: q.promptRaw != null ? String(q.promptRaw) : undefined,
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;
if (q.engine === 'qwen21' && ['generate', 'edit'].includes(q.mode)) {
q.task = q.mode
q.prompt = String(q.prompt || q.compiledPrompt || '')
if (q.promptRaw == null) q.promptRaw = String(q.compiledPrompt || '')
}
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', handoffSource: r.handoffSource || null, fitted: video && q.startClipId ? !!r.handoffFitted : null };
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();
delete r.finishedAt;
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);
}
}