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 validateYueGpRequest(body) { const profile = body.profile ?? 1 const duration = body.duration ?? 60 if (![1, 3].includes(profile)) throw new Error('YuEGP profile must be 1 or manual fallback 3.') if (!Number.isInteger(duration) || duration < 30 || duration > 150) throw new Error('YuEGP 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('YuEGP requires one non-empty lyric section.') return { id: body.id, profile, duration, seed: body.seed, tags, lyrics } } /** One isolated Python process per song. A process exit is the GPU release boundary. */ export function createYueGpHost({ prepare, leaseValid, spawnProcess = spawn, root, python, dataDir } = {}) { const repo = resolve(root || process.env.YUEGP_ROOT || join(fileURLToPath(new URL('../..', import.meta.url)), 'YuEGP')) const executable = python || process.env.YUEGP_PYTHON || join(repo, '.venv', 'Scripts', 'python.exe') const data = resolve(dataDir || process.env.YUEGP_JOBS_DIR || join(repo, 'aigen-jobs')) const worker = fileURLToPath(new URL('./yuegp-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 = 'YuEGP 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 = validateYueGpRequest(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('YuEGP is already running or releasing a previous worker.'), { statusCode: 409 }) if (!existsSync(executable) || !existsSync(join(repo, 'aigen-ready.json'))) throw new Error('YuEGP is not installed. Run scripts/setup-yuegp.ps1 on the GPU host.') const job = { id: request.id, status: 'starting', message: 'Preparing GPU for YuEGP', profile: request.profile, 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('YuEGP start cancelled or GPU reservation expired.') const installation = JSON.parse(readFileSync(join(repo, 'aigen-ready.json'), 'utf8').replace(/^\uFEFF/, '')) const localModels = process.env.YUEGP_MODELS_DIR || installation.modelsDir const model = (key, name) => process.env[key] || (localModels && existsSync(join(localModels, name)) ? join(localModels, name) : `m-a-p/${name}`) writeFileSync(join(dir(job.id), 'request.json'), JSON.stringify({ ...request, stage1Model: model('YUEGP_STAGE1_MODEL', 'YuE-s1-7B-anneal-en-cot'), stage2Model: model('YUEGP_STAGE2_MODEL', 'YuE-s2-1B-general') })) const child = spawnProcess(executable, ['-u', worker, '--root', repo, '--request', join(dir(job.id), 'request.json'), '--profile', String(request.profile)], { cwd: join(repo, 'inference'), 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 YuEGP'; 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. */ } } }) // The worker also exits if this host dies, using its parent PID watchdog. const watchdog = setInterval(() => { if (!leaseValid(lease)) { job.error = 'GPU reservation expired; YuEGP 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 || `YuEGP 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) } } }