Route engine yue2 through a standalone job helper.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -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. */ }
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user