Files
aigen/server/utils/studio2/runner.ts
T
TowstyandCursor 99b61e67d7 Keep Enhance + Generate as one queue row with shared progress.
The PE stage now stays on the same jobId so the header and queue badge agree, and jobs behind it do not promote when enhance finishes.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-26 14:42:22 -05:00

532 lines
29 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 { 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
}
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'
// Edit → PE-I2I only (pe_i2i). Generate → PE-T2I only. Never cross-feed.
if (edit && !q.imageAId)
throw new Error('Qwen Edit requires Start still (image_1). Hero is ignored.')
update(r, 'enhancing')
emitJob(job, { type: 'status', message: 'Enhancing prompt…' })
let prompt = String(q.compiledPrompt || '')
if (edit) prompt = ensureQwen21EditPrompt(prompt)
const graph = structuredClone(edit ? qwen21PeEditTemplate : qwen21PeT2iTemplate)
graph['9'].inputs.prompt = prompt
graph['9'].inputs.seed = s.seed
if (edit) {
// Same Start still pixels as Edit image_1 — never Hero.
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)
if (!result.parse_ok || !result.positive_prompt)
throw new Error('Prompt enhance failed to parse a rewrite. Try again or turn Enhance prompt off.')
r.request.promptRaw = q.compiledPrompt
// PE-I2I: if rewrite dropped every <imageN>, stitch keep-identity onto the rewrite.
r.request.compiledPrompt = edit
? stitchQwen21PeEditPrompt(result.positive_prompt)
: result.positive_prompt
r.request.enhance = {
wh_ratio: result.wh_ratio || '',
...(edit ? { ratio_follow: result.ratio_follow || '' } : {}),
parse_ok: true,
thinking: result.thinking || '',
comfyPromptId: queued.prompt_id,
}
// 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}`;
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: Start still (imageAId) only — never Hero / identityStillId.
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') {
// Start still → image_1 only. Never Hero. Never T2I / EmptyLatent.
if (!q.imageAId || !a)
throw new Error('Qwen Edit requires Start still (image_1). Hero is ignored.')
const prompt = ensureQwen21EditPrompt(String(q.compiledPrompt || ''))
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).
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 };
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', 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);
}
}