101 lines
5.7 KiB
JavaScript
101 lines
5.7 KiB
JavaScript
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)
|
|
}
|
|
}
|
|
}
|