Files
aigen/scripts/yuegp-host.mjs

138 lines
7.7 KiB
JavaScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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)
}
}
}