diff --git a/server/plugins/00-resume-yue2.ts b/server/plugins/00-resume-yue2.ts new file mode 100644 index 0000000..b2a9c6c --- /dev/null +++ b/server/plugins/00-resume-yue2.ts @@ -0,0 +1,4 @@ +import { resumeYue2Jobs } from '~/server/utils/yue2' + +// Restore live music IDs before queue repair examines their durable studio rows. +export default defineNitroPlugin(() => { resumeYue2Jobs() }) diff --git a/server/utils/jobs.ts b/server/utils/jobs.ts index d104159..adf3315 100644 --- a/server/utils/jobs.ts +++ b/server/utils/jobs.ts @@ -34,6 +34,7 @@ export interface JobEvent { export interface Job { upscale?: boolean yueGp?: boolean + yue2?: boolean musicActivity?: { checkedAt: number; running: boolean } id: string kind?: 'video' | 'edit' | 'music' diff --git a/server/utils/library.ts b/server/utils/library.ts index a64febc..1eec00c 100644 --- a/server/utils/library.ts +++ b/server/utils/library.ts @@ -1785,6 +1785,10 @@ export function findYueGpTrack(owner: string, jobId: string) { return readCatalog(owner).tracks.find(track => track.comfyFilename === `yuegp-${jobId}.wav`) } +export function findYue2Track(owner: string, jobId: string) { + return readCatalog(owner).tracks.find(track => track.comfyFilename === `yue2-${jobId}.wav`) +} + export async function saveTrack(params: { ownerKey: string folderId: string diff --git a/server/utils/musicChain.ts b/server/utils/musicChain.ts index a8a103f..c291724 100644 --- a/server/utils/musicChain.ts +++ b/server/utils/musicChain.ts @@ -1,4 +1,5 @@ import { startYueGpJob } from './yueGp' +import { startYue2Job } from './yue2' import { createJob, emitJob, type Job } from '~/server/utils/jobs' import { extractAudio, fetchHistory, fetchHistoryAll, findHistoryAudio, purgeComfyArtifacts, queuePrompt } from '~/server/utils/comfy' import { comfyWsUrl, comfyFetch } from '~/server/utils/comfy' @@ -341,6 +342,7 @@ function watchMusicJob(job: Job): Promise { export async function startMusicJob(params: MusicJobParams) { if (params.engine === 'yue') return startYueGpJob(params) + if (params.engine === 'yue2') return startYue2Job(params) const job = createJob('music') job.library = { ownerKey: params.ownerKey, diff --git a/server/utils/musicWorkflow.ts b/server/utils/musicWorkflow.ts index 8a8bf7d..db1c835 100644 --- a/server/utils/musicWorkflow.ts +++ b/server/utils/musicWorkflow.ts @@ -101,6 +101,7 @@ function buildAce15Workflow(params: MusicWorkflowParams): WorkflowGraph { export function buildMusicWorkflow(params: MusicWorkflowParams): WorkflowGraph { const engine = params.engine || 'ace-step' if (engine === 'yue') throw new Error('YuE requires the standalone YuEGP backend.') + if (engine === 'yue2') throw new Error('YuE2 requires the standalone YuE2 backend.') if (engine === 'ace-step-1.5') return buildAce15Workflow(params) return buildAceV1Workflow(params) } @@ -108,6 +109,7 @@ export function buildMusicWorkflow(params: MusicWorkflowParams): WorkflowGraph { export async function assertMusicEngineNodes(engine: MusicEngine | undefined) { const { comfyHasClassType } = await import('~/server/utils/comfy') if (engine === 'yue') throw new Error('YuE cannot run on Comfy.') + if (engine === 'yue2') throw new Error('YuE2 cannot run on Comfy.') if (engine === 'ace-step-1.5') { const present = await comfyHasClassType('TextEncodeAceStepAudio1.5') if (present === false) { diff --git a/server/utils/yue2.ts b/server/utils/yue2.ts new file mode 100644 index 0000000..3f01fe0 --- /dev/null +++ b/server/utils/yue2.ts @@ -0,0 +1,148 @@ +import { existsSync, mkdirSync, readdirSync, readFileSync, writeFileSync, renameSync, unlinkSync } from 'node:fs' +import { join } from 'node:path' +import { createJob, emitJob, getJob, restoreMusicJob, type Job } from './jobs' +import { saveTrack, findYue2Track } from './library' +import { sharedGpuHeaders } from './sharedGpu' +import type { MusicJobParams } from './musicChain' + +function settings() { + const config = useRuntimeConfig() + return { url: String(config.comfyControlUrl || process.env.COMFY_CONTROL_URL || '').replace(/\/$/, ''), + token: String(config.comfyControlToken || process.env.COMFY_CONTROL_TOKEN || '') } +} +export function yue2Configured() { return Boolean(settings().url) } +async function request(path: string, body?: unknown) { + const { url, token } = settings() + if (!url) throw new Error('YuE2 host is not configured. Set COMFY_CONTROL_URL.') + const response = await fetch(`${url}/yue2/${path}`, { + method: body === undefined ? 'GET' : 'POST', + headers: { ...(token ? { Authorization: `Bearer ${token}` } : {}), + ...(body === undefined ? {} : { ...sharedGpuHeaders(), 'Content-Type': 'application/json' }) }, + body: body === undefined ? undefined : JSON.stringify(body), signal: AbortSignal.timeout(60_000) + }) + if (!response.ok) { + const detail = await response.json().catch(() => ({})) as { message?: string; error?: string } + throw Object.assign(new Error(detail.message || detail.error || `YuE2 host returned ${response.status}`), { statusCode: response.status }) + } + return response +} +function pendingRoot() { return join(String(useRuntimeConfig().libraryDir || process.env.LIBRARY_DIR || '/data/library'), 'yue2-pending') } +function persist(job: Job) { + mkdirSync(pendingRoot(), { recursive: true }) + const path = join(pendingRoot(), `${job.id}.json`) + writeFileSync(path + '.tmp', JSON.stringify({ id: job.id, clientId: job.clientId, startedAt: job.startedAt, library: job.library, trackId: job.trackId })) + renameSync(path + '.tmp', path) +} +async function settle(job: Job) { + const { onLiveVideoSettled } = await import('./studioQueue') + await onLiveVideoSettled(job) + const path = join(pendingRoot(), `${job.id}.json`) + if (existsSync(path)) unlinkSync(path) +} +export async function cancelYue2Job(job: Job) { + await request(`jobs/${job.id}/cancel`, {}) + job.status = 'cancelled' + emitJob(job, { type: 'error', error: 'Cancelled', message: 'Cancelled' }) +} + +async function watch(job: Job) { + let failures = 0 + while (job.status !== 'cancelled') { + try { + const state = await (await request(`jobs/${job.id}`)).json() as { + status: string; message: string; error?: string; stage?: string; progress?: number; step?: number; maxStep?: number; duration?: number + } + if (job.status === 'cancelled') { await settle(job); return } + job.musicActivity = { checkedAt: Date.now(), running: ['starting', 'running', 'cancelling'].includes(state.status) } + if (state.status === 'error' || state.status === 'cancelled') { + job.status = state.status === 'cancelled' ? 'cancelled' : 'error' + job.error = state.error || state.message + if (/out of memory/i.test(job.error || '')) job.error += ' YuE2 does not switch engines automatically.' + emitJob(job, { type: 'error', error: job.error, message: job.error }) + await settle(job) + return + } + if (state.status === 'complete') { + job.saving = true + emitJob(job, { type: 'status', message: 'Saving audio to library…', progress: 98 }) + const lib = job.library! + lib.audioExt = 'wav' + job.trackId ||= findYue2Track(lib.ownerKey, job.id)?.id + if (!job.trackId) { + const audio = Buffer.from(await (await request(`jobs/${job.id}/audio`)).arrayBuffer()) + const track = await saveTrack({ ownerKey: lib.ownerKey, folderId: lib.folderId, name: lib.name, + tags: lib.tags || lib.prompt, lyrics: lib.lyrics || '', duration: state.duration || lib.duration || 60, + seed: lib.seed, steps: 0, cfg: 0, instrumental: lib.instrumental === true, + engine: 'yue2', audio, ext: 'wav', comfyFilename: `yue2-${job.id}.wav` }) + job.trackId = track.id + persist(job) + } + job.status = 'complete'; job.saving = false + emitJob(job, { type: 'complete', message: lib.folderLocked ? 'Saved to the locked folder.' : 'Track ready', + progress: 100, trackId: job.trackId, folderLocked: lib.folderLocked, audioExt: 'wav' }) + await settle(job) + return + } + failures = 0 + job.status = 'running' + const percent = Number(state.progress || 0) + emitJob(job, { type: 'progress', message: `${state.message}${state.maxStep ? ` · ${state.step}/${state.maxStep}` : ''}`, + progress: Math.min(95, Math.max(1, percent)), step: state.step || 0, maxStep: state.maxStep || 0 }) + } catch (error) { + failures++ + job.saving = false + emitJob(job, { type: 'status', message: `YuE2 connection/save check failed; retrying: ${error instanceof Error ? error.message : String(error)}` }) + if (failures >= 10) { + job.status = 'error'; job.error = 'YuE2 host unreachable or audio save failed. Job files are retained for recovery.' + emitJob(job, { type: 'error', error: job.error, message: job.error }) + const { onLiveVideoSettled } = await import('./studioQueue') + await onLiveVideoSettled(job) + return + } + } + await new Promise(resolve => setTimeout(resolve, 2000)) + } + await settle(job) +} + +export function startYue2Job(params: MusicJobParams) { + const job = createJob('music') + job.yue2 = true + job.library = { ...params, prompt: params.tags, engine: 'yue2', aspect: 'audio', width: 0, height: 0, + hideThumbnail: false, turbo: false, sound: true } + persist(job) + setTimeout(() => { void (async () => { + try { + if (job.status === 'cancelled') { await settle(job); return } + await request('jobs', { id: job.id, tags: params.tags, lyrics: params.lyrics, + duration: params.duration, seed: params.seed }) + await watch(job) + } catch (error) { + const statusCode = (error as { statusCode?: number })?.statusCode + if (statusCode && statusCode >= 400 && statusCode < 500) { + job.status = 'error'; job.error = error instanceof Error ? error.message : String(error) + emitJob(job, { type: 'error', error: job.error, message: job.error }) + await settle(job) + return + } + emitJob(job, { type: 'status', message: `Checking YuE2 submission: ${error instanceof Error ? error.message : String(error)}` }) + await watch(job) + } + })() }, 0) + return job +} + +export function resumeYue2Jobs() { + if (!existsSync(pendingRoot())) return + for (const file of readdirSync(pendingRoot()).filter(file => /^[a-zA-Z0-9-]+\.json$/.test(file))) { + try { + const record = JSON.parse(readFileSync(join(pendingRoot(), file), 'utf8')) + if (getJob(record.id)) continue + const job = restoreMusicJob(record) + job.yueGp = false + job.yue2 = true + job.message = 'Reconnecting to YuE2' + void watch(job) + } catch { /* Preserve invalid records for diagnosis. */ } + } +}