Queue durable Studio 2 jobs with sticky identity and saved metadata
This commit is contained in:
@@ -0,0 +1,3 @@
|
||||
import { readRecord } from '../../utils/studio2/store';
|
||||
import { cancelStudioJob } from '../../utils/studioQueue';
|
||||
export default defineEventHandler(async (event) => { const { owner } = assertLibraryOwner(event), body = await readBody(event), r = readRecord(owner, String(body.id)); assertFolderAccess(event, r.request.folderId); await cancelStudioJob(owner, r.queueId); return { ok: true }; });
|
||||
@@ -0,0 +1,9 @@
|
||||
import { records } from '../../utils/studio2/store';
|
||||
import { listStudioJobs } from '../../utils/studioQueue';
|
||||
export default defineEventHandler(event => { const { owner } = assertLibraryOwner(event); const rows = listStudioJobs(owner); return records(owner).filter(r => { try {
|
||||
assertFolderAccess(event, r.request.folderId);
|
||||
return true;
|
||||
}
|
||||
catch {
|
||||
return false;
|
||||
} }).map(r => { const q = rows.find(q => q.id === r.queueId); return { ...r, state: q?.status === 'cancelled' ? 'cancelled' : q?.status === 'error' ? 'failed' : r.state, error: r.error || q?.lastError }; }); });
|
||||
@@ -0,0 +1,31 @@
|
||||
import { validateRequest } from '~/shared/studio2/contracts.mjs';
|
||||
import { parsePostedLoraStack, assertImageV2LoraStack } from '../../utils/loras';
|
||||
import { saveRecord, readRecord } from '../../utils/studio2/store';
|
||||
import { addStudioJob, kickStudioQueue, type StudioJobPayload } from '../../utils/studioQueue';
|
||||
export default defineEventHandler(async (event) => {
|
||||
const { owner } = assertLibraryOwner(event), raw = await readBody(event);
|
||||
if (useRuntimeConfig().public.studio === 'xaigen')
|
||||
throw createError({ statusCode: 400, statusMessage: 'Studio 2 is currently an AIGen preview. Continue using the existing LTX studio on xAIGen.' });
|
||||
let request: any;
|
||||
try {
|
||||
request = validateRequest(raw, useRuntimeConfig().public.studio === 'xaigen');
|
||||
}
|
||||
catch (e: any) {
|
||||
throw createError({ statusCode: e.statusCode || 400, statusMessage: e.message });
|
||||
}
|
||||
const video = ['video', 'extend'].includes(request.mode);
|
||||
request.settings.loraStack = parsePostedLoraStack(request.settings.loraStack, video ? 'video' : 'image');
|
||||
if (!video)
|
||||
assertImageV2LoraStack(request.settings.loraStack, request.engine);
|
||||
assertFolderAccess(event, request.folderId);
|
||||
for (const id of [request.identityStillId, request.imageAId, request.imageBId, request.maskId, request.endStillId, ...request.guides.map((g: any) => g.stillId)].filter(Boolean))
|
||||
assertFolderAccess(event, getStill(owner, id).folderId);
|
||||
if (request.startClipId)
|
||||
assertFolderAccess(event, getClip(owner, request.startClipId).folderId);
|
||||
const id = crypto.randomUUID(), record = { id, owner, request, state: 'waiting', queuedAt: Date.now(), familyId: crypto.randomUUID(), outputs: [], purgeResult: 'Not yet saved' };
|
||||
saveRecord(record);
|
||||
const row = await addStudioJob({ ownerKey: owner, kind: ['video', 'extend'].includes(request.mode) ? 'video' : 'edit', familyId: record.familyId, payload: { ...request.settings, studio2Id: id, prompt: request.compiledPrompt, folderId: request.folderId, extensions: [], referenceStillIds: [], useIdentityRefs: false, queueAutoRun: true } as StudioJobPayload });
|
||||
saveRecord({ ...readRecord(owner, id), queueId: row.id });
|
||||
kickStudioQueue();
|
||||
return { id };
|
||||
});
|
||||
@@ -0,0 +1,2 @@
|
||||
import { resumeStudio2Jobs } from '../utils/studio2/runner';
|
||||
export default defineNitroPlugin(() => { resumeStudio2Jobs(); });
|
||||
@@ -32,6 +32,7 @@ export interface JobEvent {
|
||||
}
|
||||
|
||||
export interface Job {
|
||||
studio2?: boolean
|
||||
upscale?: boolean
|
||||
yueGp?: boolean
|
||||
musicActivity?: { checkedAt: number; running: boolean }
|
||||
|
||||
@@ -35,6 +35,7 @@ export interface LibraryFolder {
|
||||
}
|
||||
|
||||
export interface LibraryClip {
|
||||
studio2?: Record<string, any>
|
||||
id: string
|
||||
folderId: string
|
||||
name: string
|
||||
@@ -95,6 +96,7 @@ export interface LibraryTrack {
|
||||
export type StillRole = 'input' | 'output'
|
||||
|
||||
export interface LibraryStill {
|
||||
studio2?: Record<string, any>
|
||||
id: string
|
||||
folderId: string
|
||||
filename: string
|
||||
@@ -270,6 +272,7 @@ function normalizeStill(still: LibraryStill): LibraryStill {
|
||||
manualGroupIndex: still.manualGroupIndex,
|
||||
parentStillId: still.parentStillId,
|
||||
chainIndex: still.chainIndex,
|
||||
studio2: still.studio2,
|
||||
settings: normalizeStillSettings(still.settings)
|
||||
}
|
||||
}
|
||||
@@ -2211,3 +2214,7 @@ export function groupLibraryStills(owner: string, ids: string[], ungroup = false
|
||||
}
|
||||
|
||||
export function findUpscaledClip(owner: string, jobId: string) { return readCatalog(owner).clips.find(clip => clip.upscaleJobId === jobId) }
|
||||
|
||||
export async function attachStudio2Metadata(owner: string, id: string, metadata: Record<string, any>) {
|
||||
await mutate(owner, catalog => { const item = catalog.clips.find(c => c.id === id) || catalog.stills.find(c => c.id === id); if (!item) throw new Error("Saved Studio 2 asset missing"); item.studio2 = metadata })
|
||||
}
|
||||
|
||||
@@ -0,0 +1,274 @@
|
||||
import { readFileSync, mkdirSync, existsSync, unlinkSync } from 'node:fs';
|
||||
import { join } from 'node:path';
|
||||
import template from '../../assets/studio2_minimax_native.json';
|
||||
import { nativeVideoGraph, attachHeroReference } 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<string> {
|
||||
const prefix = String(useRuntimeConfig().comfyFilenamePrefix).replace(/\/$/, '') + `/studio2/${r.id}/${r.index}`;
|
||||
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 purge(r: any) {
|
||||
const config = useRuntimeConfig(), prefix = String(config.comfyFilenamePrefix);
|
||||
const files = (r.files as DiskFile[]).filter(file => scopedFile(file, prefix));
|
||||
r.purgeResult = 'Left on host';
|
||||
if (files.length !== r.files.length || !files.length)
|
||||
return;
|
||||
try {
|
||||
const response = await fetch(String(config.comfyControlUrl).replace(/\/$/, '') + '/studio2/purge', {
|
||||
method: 'POST', headers: { 'Content-Type': 'application/json', ...sharedGpuHeaders() },
|
||||
body: JSON.stringify({ prefix, files }), signal: AbortSignal.timeout(12000)
|
||||
});
|
||||
const result = await response.json() as any;
|
||||
if (response.ok && result.ok && result.cleared === files.length)
|
||||
r.purgeResult = 'Cleared input + output';
|
||||
}
|
||||
catch { /* Keep the durable library copy and report unconfirmed host cleanup. */ }
|
||||
}
|
||||
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 hero = await load(q.identityStillId, 'hero'), a = await load(q.imageAId, 'source'), b = await load(q.imageBId, 'compose'), mask = await load(q.maskId, 'mask');
|
||||
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;
|
||||
if (q.engine === 'flux')
|
||||
attachHeroReference(graph, hero);
|
||||
else if (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) {
|
||||
try {
|
||||
const prompts = r.prompts || [r.request.promptSections, ...r.request.batch];
|
||||
r.prompts = prompts;
|
||||
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.files = [];
|
||||
r.request.promptSections = prompts[r.index];
|
||||
r.request.compiledPrompt = compilePrompt(prompts[r.index], ['video', 'extend'].includes(r.request.mode));
|
||||
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.width}:${s.height}`, hideThumbnail: false, video: assembled, originalSegment: data, sourceSegments, fps: 24, sound: true, familyId: r.familyId, parentClipId: q.startClipId || undefined, chainIndex: r.index, comfyFilename: file.filename });
|
||||
const anchor = history[r.promptId].outputs?.anchor_save?.images?.[0];
|
||||
if (anchor) {
|
||||
r.files.push(anchor);
|
||||
await persistClipAnchorFrame({ ownerKey: r.owner, clipId: asset.id, fromBuffer: await downloadComfyImage(anchor) });
|
||||
}
|
||||
}
|
||||
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, 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);
|
||||
await purge(r);
|
||||
metadata.purgeResult = r.purgeResult;
|
||||
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);
|
||||
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 {
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
import { mkdirSync, readFileSync, writeFileSync, renameSync, readdirSync } from 'node:fs';
|
||||
import { join } from 'node:path';
|
||||
export function studio2Root() { const dir = join(String(useRuntimeConfig().libraryDir), 'studio2'); mkdirSync(dir, { recursive: true }); return dir; }
|
||||
export function saveRecord(record: any) { const path = join(studio2Root(), `${record.id}.json`); writeFileSync(path + '.tmp', JSON.stringify(record)); renameSync(path + '.tmp', path); }
|
||||
export function readRecord(owner: string, id: string) { if (!/^[\w-]+$/.test(id))
|
||||
throw createError({ statusCode: 400, statusMessage: 'Invalid job ID' }); const r = JSON.parse(readFileSync(join(studio2Root(), `${id}.json`), 'utf8')); if (r.owner !== owner)
|
||||
throw createError({ statusCode: 404, statusMessage: 'Job not found' }); return r; }
|
||||
export function records(owner?: string) { return readdirSync(studio2Root()).filter(n => n.endsWith('.json')).flatMap(n => { try {
|
||||
const r = JSON.parse(readFileSync(join(studio2Root(), n), 'utf8'));
|
||||
return !owner || r.owner === owner ? [r] : [];
|
||||
}
|
||||
catch {
|
||||
return [];
|
||||
} }).sort((a, b) => b.queuedAt - a.queuedAt); }
|
||||
@@ -16,6 +16,7 @@ export type StudioJobStatus = 'waiting' | 'running' | 'held' | 'complete' | 'err
|
||||
export type StudioJobKind = 'video' | 'edit' | 'music'
|
||||
|
||||
export interface StudioJobPayload {
|
||||
studio2Id?: string
|
||||
upscale?: { sourceId: string; scale: 2 | 4; target: string; fps: string | number; enhance: string }
|
||||
prompt: string
|
||||
promptMid?: string
|
||||
@@ -285,7 +286,7 @@ function failZombieLiveJob(job: Job, error: string) {
|
||||
function sweepStaleLiveJobs() {
|
||||
const now = Date.now()
|
||||
for (const job of listJobs()) {
|
||||
if (job.yueGp || job.upscale) continue
|
||||
if (job.studio2 || job.yueGp || job.upscale) continue
|
||||
if (job.saving) continue
|
||||
if (job.status === 'queued' && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) {
|
||||
failZombieLiveJob(job, 'Job never started')
|
||||
@@ -305,7 +306,7 @@ async function reapZombieLiveJobs() {
|
||||
const { fetchHistory } = await import('~/server/utils/comfy')
|
||||
for (const job of listJobs()) {
|
||||
// Never interrupt download/stitch/library save — Comfy is idle then by design.
|
||||
if (job.yueGp || job.upscale) continue
|
||||
if (job.studio2 || job.yueGp || job.upscale) continue
|
||||
if (job.saving) continue
|
||||
if (job.library?.chainContinuing) continue
|
||||
if (job.status !== 'running' && job.status !== 'uploading' && job.status !== 'queued') continue
|
||||
@@ -350,7 +351,7 @@ async function reapZombieLiveJobs() {
|
||||
}
|
||||
|
||||
function liveJobOwnsGpu(job: Job) {
|
||||
if ((job.yueGp || job.upscale) && ['running', 'queued', 'uploading'].includes(job.status)) return true
|
||||
if ((job.studio2 || job.yueGp || job.upscale) && ['running', 'queued', 'uploading'].includes(job.status)) return true
|
||||
if (job.library?.stopAfterCurrent) return false
|
||||
if (job.library?.chainContinuing) return true
|
||||
if (job.saving) return true
|
||||
@@ -370,7 +371,7 @@ function liveJobOwnsGpu(job: Job) {
|
||||
*/
|
||||
function clearDeadGpuClaimsForForceStart() {
|
||||
for (const live of listJobs()) {
|
||||
if (live.yueGp || live.upscale) continue
|
||||
if (live.studio2 || live.yueGp || live.upscale) continue
|
||||
if (live.saving) continue
|
||||
if (jobIsLocallySubmitting(live)) continue
|
||||
if (live.status !== 'running' && live.status !== 'queued' && live.status !== 'uploading') {
|
||||
@@ -1422,6 +1423,19 @@ export async function startStudioJob(item: StudioJob) {
|
||||
}
|
||||
|
||||
async function startStudioJobReserved(item: StudioJob) {
|
||||
if (item.payload.studio2Id) {
|
||||
try {
|
||||
const { startStudio2Job } = await import('./studio2/runner')
|
||||
await startStudio2Job(item)
|
||||
} catch (error) {
|
||||
await patchStudioJob(item.ownerKey, item.id, row => {
|
||||
row.status = 'error'
|
||||
row.lastError = error instanceof Error ? error.message : String(error)
|
||||
})
|
||||
scheduleKickRetry()
|
||||
}
|
||||
return
|
||||
}
|
||||
if (item.payload.upscale) {
|
||||
try { const { startUpscaleJob } = await import('./videoUpscale'); await startUpscaleJob(item) }
|
||||
catch (error) { await patchStudioJob(item.ownerKey, item.id, row => { row.status = 'error'; row.lastError = error instanceof Error ? error.message : String(error) }); scheduleKickRetry() }
|
||||
@@ -1673,7 +1687,7 @@ export async function onLiveVideoSettled(job: Job) {
|
||||
return
|
||||
}
|
||||
const remaining = remainingStudioShots(job)
|
||||
const wakeFail = !job.yueGp && !job.upscale && job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
|
||||
const wakeFail = !job.studio2 && !job.yueGp && !job.upscale && job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
|
||||
const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail
|
||||
|
||||
await mutateStore(owner, (store) => {
|
||||
|
||||
Reference in New Issue
Block a user