import { spawn } from 'node:child_process' import { existsSync, mkdirSync, readFileSync, writeFileSync, renameSync, createWriteStream, appendFileSync, readdirSync } from 'node:fs' import { pipeline } from 'node:stream/promises' import { Transform } from 'node:stream' import { join, resolve } from 'node:path' import { fileURLToPath } from 'node:url' import { upscaleOptions } from '../shared/video-upscale.mjs' export function createUpscaleHost({ leaseValid, spawnProcess = spawn, root = process.env.UPSCALE_ROOT || fileURLToPath(new URL('../../AIGen-Upscale', import.meta.url)) }) { root = resolve(root) const data = join(root, 'jobs') mkdirSync(data, { recursive: true }) let active = null const directory = id => { if (!/^[a-zA-Z0-9-]{12,80}$/.test(id || '')) throw new Error('Invalid upscale ID.') return join(data, id) } const persist = state => { const path = join(directory(state.id), 'status.json'); writeFileSync(path + '.tmp', JSON.stringify(state)); renameSync(path + '.tmp', path) } // A previous worker's parent watchdog terminates its process tree before another GPU job starts. let holdUntil = 0 for (const id of readdirSync(data)) { try { const state = JSON.parse(readFileSync(join(directory(id), 'status.json'), 'utf8')) if (state.status === 'running') { holdUntil = Date.now() + 15000; state.status = 'error'; state.error = 'Upscale host restarted. Original video is safe.'; persist(state) } } catch { /* Not a valid job record. */ } } const read = id => { const path = join(directory(id), 'status.json') return existsSync(path) ? JSON.parse(readFileSync(path, 'utf8')) : null } const stop = () => new Promise((resolve, reject) => { if (!active?.child?.pid) return resolve() const child = active.child child.once('close', resolve) const killer = spawn('taskkill', ['/PID', String(child.pid), '/T', '/F'], { windowsHide: true, stdio: 'ignore' }) killer.once('error', reject) killer.once('close', code => { if (code && child.exitCode === null) reject(new Error('Could not stop upscale worker. GPU remains reserved.')) }) }) return { busy: () => Boolean(active) || Date.now() < holdUntil, read, output: id => read(id)?.status === 'complete' ? join(directory(id), 'output.mp4') : null, async upload(id, stream) { if (read(id)) throw new Error('This upscale job has already been submitted.') const dir = directory(id); mkdirSync(dir, { recursive: true }) let bytes = 0 await pipeline(stream, new Transform({ transform(chunk, _, callback) { bytes += chunk.length callback(bytes > 8 * 1024 ** 3 ? new Error('Video exceeds the 8 GB upload limit.') : null, chunk) } }), createWriteStream(join(dir, 'input.upload'), { flags: 'w' })) if (bytes < 32) throw new Error('Video file is empty.') renameSync(join(dir, 'input.upload'), join(dir, 'input.mp4')) }, start(body, lease) { const options = upscaleOptions(body), dir = directory(body.id) const previous = read(body.id) if (previous) return previous if (active || Date.now() < holdUntil) throw Object.assign(new Error('Local upscaler is busy.'), { statusCode: 409 }) if (!leaseValid(lease)) throw new Error('GPU reservation expired.') if (!existsSync(join(dir, 'input.mp4'))) throw new Error('Upload a completed video first.') const state = { id: body.id, status: 'running', startedAt: Date.now(), message: 'Preparing local upscaler', progress: 0 } writeFileSync(join(dir, 'request.json'), JSON.stringify({ ...options, root, parentPid: process.pid })) persist(state) const child = spawnProcess(process.execPath, [fileURLToPath(new URL('./upscale-worker.mjs', import.meta.url)), join(dir, 'request.json')], { windowsHide: true, shell: false, stdio: ['ignore', 'pipe', 'pipe'] }) active = { child, state } let lines = '', tail = '', cancelled = false const watchdog = setInterval(() => { if (!leaseValid(lease)) { state.error = 'GPU reservation expired.'; void stop().catch(() => {}) } }, 5000) child.stderr.on('data', chunk => { appendFileSync(join(dir, 'worker.log'), chunk); tail = (tail + chunk).slice(-2000) }) child.stdout.on('data', chunk => { appendFileSync(join(dir, 'worker.log'), chunk) lines += chunk const parts = lines.split('\n'); lines = parts.pop().slice(-64000) for (const line of parts) if (line.startsWith('AIGEN_EVENT ')) { try { const event = JSON.parse(line.slice(12)); for (const key of ['message', 'progress', 'error', 'width', 'height', 'sourceWidth', 'sourceHeight', 'fps', 'duration', 'engine', 'elapsedMs']) if (event[key] !== undefined) state[key] = event[key]; persist(state) } catch {} } }) active.cancel = async () => { cancelled = true; await stop() } let finished = false const finish = error => { if (finished) return finished = true; clearInterval(watchdog) state.status = cancelled ? 'cancelled' : !error && !state.error && existsSync(join(dir, 'output.mp4')) ? 'complete' : 'error' state.elapsedMs = Date.now() - state.startedAt if (state.status === 'error') state.error ||= error?.message || tail || 'Upscale failed. Original video is safe.' persist(state); active = null } child.once('error', finish) child.once('close', code => finish(code ? new Error(`Upscaler exited (${code}). ${tail}`) : null)) return state }, async cancel(id) { if (active?.state.id === id) await active.cancel() if (!read(id)) { mkdirSync(directory(id), { recursive: true }) persist({ id, status: 'cancelled', message: 'Cancelled before processing', elapsedMs: 0 }) } return read(id) } } }