Add YuE2 host, worker, and setup scripts.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Towsty
2026-09-14 18:45:57 -05:00
co-authored by Cursor
parent 064f1db933
commit 85ef087364
3 changed files with 261 additions and 0 deletions
+134
View File
@@ -0,0 +1,134 @@
import { spawn } from 'node:child_process'
import { existsSync, mkdirSync, readFileSync, writeFileSync, renameSync, appendFileSync, readdirSync } from 'node:fs'
import { join, resolve } from 'node:path'
import { fileURLToPath } from 'node:url'
export function validateYue2Request(body) {
const duration = body.duration ?? 60
if (!Number.isInteger(duration) || duration < 30 || duration > 150) throw new Error('YuE2 duration must be 30–150 seconds.')
if (!Number.isInteger(body.seed) || body.seed < 0 || body.seed > 2147483647) throw new Error('Invalid seed.')
if (!/^[a-zA-Z0-9-]{12,80}$/.test(body.id || '')) throw new Error('Invalid job ID.')
const tags = String(body.tags || '').trim()
const lyrics = String(body.lyrics || '').trim()
if (!tags || tags.length > 2000 || !lyrics || lyrics.length > 8000) throw new Error('Genre tags and one lyric section are required.')
const sections = [...lyrics.matchAll(/\[([^\]]+)\]\s*([^\[]*)/gs)].filter(m => m[2].trim())
if (sections.length !== 1) throw new Error('YuE2 requires one non-empty lyric section.')
return { id: body.id, duration, seed: body.seed, tags, lyrics }
}
/** One isolated Python process per song. A process exit is the GPU release boundary. */
export function createYue2Host({ prepare, leaseValid, spawnProcess = spawn, root, python, dataDir } = {}) {
const repo = resolve(root || process.env.YUE2_ROOT || 'C:\\Users\\ianjm\\Development\\YuE2')
const executable = python || process.env.YUE2_PYTHON || join(repo, '.venv', 'Scripts', 'python.exe')
const data = resolve(dataDir || process.env.YUE2_JOBS_DIR || join(repo, 'aigen-jobs'))
const worker = fileURLToPath(new URL('./yue2-worker.py', import.meta.url))
let active = null
// A worker whose parent died exits within two seconds. Block new GPU owners
// across a host restart until that watchdog has had time to run.
const restartHoldUntil = existsSync(data) && readdirSync(data).some(id => {
try { return ['running', 'starting', 'cancelling'].includes(JSON.parse(readFileSync(join(data, id, 'status.json'), 'utf8')).status) }
catch { return false }
}) ? Date.now() + 10000 : 0
const dir = id => {
if (!/^[a-zA-Z0-9-]{12,80}$/.test(id || '')) throw new Error('Invalid job ID.')
return join(data, id)
}
const persist = job => {
const target = join(dir(job.id), 'status.json')
writeFileSync(target + '.tmp', JSON.stringify(job))
renameSync(target + '.tmp', target)
}
const read = id => {
if (active?.job.id === id) return { ...active.job }
const path = join(dir(id), 'status.json')
if (!existsSync(path)) return null
const job = JSON.parse(readFileSync(path, 'utf8'))
if (['running', 'starting', 'cancelling'].includes(job.status)) {
job.status = 'error'; job.error = 'YuE2 host restarted. Intermediate files are preserved.'
}
return job
}
return {
busy: () => Boolean(active) || Date.now() < restartHoldUntil,
configured: () => existsSync(executable) && existsSync(join(repo, 'aigen-ready.json')),
read,
audio(id) {
return read(id)?.status === 'complete' ? join(dir(id), 'audio.wav') : null
},
async start(body, lease) {
const request = validateYue2Request(body)
const previous = read(request.id)
if (previous) return previous // POST retries cannot generate the same song twice.
if (active || Date.now() < restartHoldUntil) throw Object.assign(new Error('YuE2 is already running or releasing a previous worker.'), { statusCode: 409 })
if (!existsSync(executable) || !existsSync(join(repo, 'aigen-ready.json'))) throw new Error('YuE2 is not installed. Run scripts/setup-yue2.ps1 on the GPU host.')
const job = { id: request.id, status: 'starting', message: 'Preparing GPU for YuE2', progress: 0, startedAt: Date.now(), checkedAt: Date.now() }
active = { job, child: null, cancelled: false }
const run = active
mkdirSync(dir(job.id), { recursive: true })
persist(job)
try {
await prepare()
if (run.cancelled || !leaseValid(lease)) throw new Error('YuE2 start cancelled or GPU reservation expired.')
const installation = JSON.parse(readFileSync(join(repo, 'aigen-ready.json'), 'utf8').replace(/^\uFEFF/, ''))
writeFileSync(join(dir(job.id), 'request.json'), JSON.stringify({
...request,
model: process.env.YUE2_MODEL || installation.model || 'm-a-p/YuE2-3B',
vae: process.env.YUE2_VAE || installation.vae || 'm-a-p/YuE2-Vae'
}))
const child = spawnProcess(executable, ['-u', worker, '--root', repo, '--request', join(dir(job.id), 'request.json')],
{ cwd: repo, windowsHide: true, shell: false, stdio: ['ignore', 'pipe', 'pipe'], env: { ...process.env, PYTHONUTF8: '1', TORCH_FORCE_NO_WEIGHTS_ONLY_LOAD: '1' } })
run.child = child
job.status = 'running'; job.message = 'Loading YuE2'; persist(job)
let tail = ''
let lines = ''
const log = chunk => { appendFileSync(join(dir(job.id), 'worker.log'), chunk); tail = (tail + chunk.toString()).slice(-6000) }
child.stderr.on('data', log)
child.stdout.on('data', chunk => {
log(chunk); lines += chunk.toString()
const parts = lines.split(/\r?\n/); lines = parts.pop().slice(-64000)
for (const line of parts) {
if (!line.startsWith('AIGEN_EVENT ')) continue
try {
const event = JSON.parse(line.slice(12))
for (const key of ['stage', 'message', 'progress', 'step', 'maxStep', 'error', 'duration']) if (event[key] !== undefined) job[key] = event[key]
job.checkedAt = Date.now(); persist(job)
} catch { /* Invalid log lines are not state transitions. */ }
}
})
const watchdog = setInterval(() => {
if (!leaseValid(lease)) { job.error = 'GPU reservation expired; YuE2 stopped.'; child.kill() }
}, 5000)
let finished = false
const finish = (code, error) => {
if (finished) return
finished = true
clearInterval(watchdog)
job.status = run.cancelled ? 'cancelled' : !error && code === 0 && existsSync(join(dir(job.id), 'audio.wav')) ? 'complete' : 'error'
if (job.status === 'error') job.error ||= error?.message || tail || `YuE2 exited with code ${code}`
job.message = job.status === 'complete' ? 'Audio ready' : job.status === 'cancelled' ? 'Cancelled' : job.error
job.checkedAt = Date.now(); persist(job)
if (active === run) active = null
}
child.once('error', error => finish(null, error))
child.once('close', code => finish(code))
return { ...job }
} catch (error) {
job.status = run.cancelled ? 'cancelled' : 'error'; job.error = error.message; persist(job)
if (active === run) active = null
throw error
}
},
async cancel(id) {
if (active?.job.id !== id) return read(id)
const run = active
run.cancelled = true
run.job.status = 'cancelling'
persist(run.job)
if (run.child) {
const child = run.child
await new Promise(resolve => { child.once('close', resolve); child.kill() })
}
return read(id)
}
}
}