node.type === 'clip').length }} parts hidden with preview
+
{{ focused.kind === 'track' ? 'Open in Music' : 'Open in Studio' }}
Use as input
proxyTarget, reservation: gpuReservation, authorized, markWork, externalBusy: () => yueGp.busy() })
+ proxyServer = createGpuProxy({ target: () => proxyTarget, reservation: gpuReservation, authorized, markWork, externalBusy: () => yueGp.busy() || upscale.busy() })
proxyServer.on('error', (error) => {
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'proxy-error', error: String(error.message || error) }))
})
@@ -813,7 +814,7 @@ function purgeDesktopFiles(body) {
}
const gpuReservation = createGpuReservation({ idle: async () => {
- if (yueGp.busy()) return false
+ if (yueGp.busy() || upscale.busy()) return false
if ((await trainingLock()).busy) return false
const healthy = await syncProxy()
if (healthy) {
@@ -839,9 +840,20 @@ const yueGp = createYueGpHost({
}
})
+const upscale = createUpscaleHost({ leaseValid: token => gpuReservation.isOwner(token) })
+
async function handleControl(req, res) {
if (!authorized(req)) return json(res, 401, { ok: false, error: 'unauthorized' })
const url = new URL(req.url || '/', 'http://localhost')
+ if (url.pathname.startsWith('/upscale/')) {
+ const match = url.pathname.match(/^\/upscale\/jobs\/([a-zA-Z0-9-]{12,80})(\/input|\/video|\/cancel)?$/)
+ if (req.method === 'POST' && url.pathname === '/upscale/jobs') return json(res, 200, upscale.start(await readJson(req), String(req.headers['x-aigen-gpu-lease'] || '')))
+ if (match && req.method === 'PUT' && match[2] === '/input') { await upscale.upload(match[1], req); return json(res, 200, { uploaded: true }) }
+ if (match && req.method === 'POST' && match[2] === '/cancel') return json(res, 200, await upscale.cancel(match[1]))
+ if (match && req.method === 'GET' && match[2] === '/video') { const path = upscale.output(match[1]); return path ? streamFile(res, path) : json(res, 404, { error: 'Video not ready' }) }
+ if (match && req.method === 'GET' && !match[2]) { const job = upscale.read(match[1]); return json(res, job ? 200 : 404, job || { error: 'Upscale job not found' }) }
+ return json(res, 404, { error: 'Unknown upscale endpoint' })
+ }
if (url.pathname.startsWith('/yuegp/')) {
const match = url.pathname.match(/^\/yuegp\/jobs\/([a-zA-Z0-9-]{12,80})(\/audio|\/cancel)?$/)
if (req.method === 'GET' && url.pathname === '/yuegp/status') return json(res, 200, { configured: yueGp.configured(), busy: yueGp.busy(), backend: 'yuegp' })
@@ -888,10 +900,12 @@ async function handleControl(req, res) {
proxyPort,
gpu: gpuReservation.availability(),
training: { busy: lastTraining.busy },
- yuegp: { busy: yueGp.busy(), configured: yueGp.configured() }
+ yuegp: { busy: yueGp.busy(), configured: yueGp.configured() },
+ upscale: { busy: upscale.busy(), engine: 'realesrgan-rife', local: true }
})
}
if (req.method === 'POST' && url.pathname === '/start') {
+ if (upscale.busy()) return json(res, 409, { message: 'Local upscale is using the GPU.' })
if (yueGp.busy()) return json(res, 409, { message: 'YuEGP is using the GPU.' })
const training = await trainingLock()
if (training.busy) {
@@ -983,12 +997,12 @@ const server = http.createServer(async (req, res) => {
if (!authorized(req)) return json(res, 401, { ok: false, error: 'unauthorized' })
try {
const path = new URL(req.url || '/', 'http://localhost').pathname
- if (req.method === 'POST' && !path.startsWith('/gpu/')) {
+ if (['POST', 'PUT'].includes(req.method) && !path.startsWith('/gpu/')) {
await gpuReservation.permit(String(req.headers['x-aigen-gpu-lease'] || ''), () => handleControl(req, res))
} else await handleControl(req, res)
} catch (error) {
req.resume()
- if (String(req.url || '').startsWith('/yuegp/') && !res.headersSent) return json(res, error.statusCode || 400, { error: error.message || 'YuEGP request failed' })
+ if ((String(req.url || '').startsWith('/yuegp/') || String(req.url || '').startsWith('/upscale/')) && !res.headersSent) return json(res, error.statusCode || 400, { error: error.message || 'YuEGP request failed' })
if (!res.headersSent) json(res, error.statusCode || 400, { ok: false, message: error.statusCode === 409 ? 'GPU is in use. Waiting for availability.' : 'GPU coordination request failed.' })
}
})
diff --git a/scripts/setup-upscale.ps1 b/scripts/setup-upscale.ps1
new file mode 100644
index 0000000..3441970
--- /dev/null
+++ b/scripts/setup-upscale.ps1
@@ -0,0 +1,30 @@
+param([string]$Root = $(if ($env:UPSCALE_ROOT) { $env:UPSCALE_ROOT } else { Join-Path (Split-Path $PSScriptRoot) '../AIGen-Upscale' }))
+$ErrorActionPreference = 'Stop'
+$ProgressPreference = 'SilentlyContinue'
+$Root = [IO.Path]::GetFullPath($Root)
+New-Item -ItemType Directory -Force -Path $Root | Out-Null
+function Install-Package($Name, $Url, $Executable) {
+ $Destination = Join-Path $Root $Name
+ $Ready = Test-Path (Join-Path $Destination $Executable)
+ if ($Name -eq 'esrgan') { $Ready = $Ready -and (Test-Path (Join-Path $Destination 'models/realesrgan-x4plus.bin')) -and (Test-Path (Join-Path $Destination 'models/realesrgan-x4plus.param')) }
+ if ($Name -eq 'rife') { $Ready = $Ready -and (Test-Path (Join-Path $Destination 'rife-v4*/flownet.bin')) -and (Test-Path (Join-Path $Destination 'rife-v4*/flownet.param')) }
+ if ($Name -eq 'ffmpeg') { $Ready = $Ready -and (Test-Path (Join-Path $Destination 'ffprobe.exe')) }
+ if ($Ready) { return }
+ $Zip = Join-Path $Root "$Name.zip"
+ Write-Output "Downloading $Name tools and weights..."
+ Invoke-WebRequest -UseBasicParsing -Uri $Url -OutFile $Zip -TimeoutSec 1800
+ $Unpack = Join-Path $Root "$Name-unpack"
+ New-Item -ItemType Directory -Force -Path $Unpack | Out-Null
+ Expand-Archive -LiteralPath $Zip -DestinationPath $Unpack -Force
+ $Exe = Get-ChildItem -LiteralPath $Unpack -Recurse -Filter $Executable | Select-Object -First 1
+ if (!$Exe) { throw "Package is missing $Executable" }
+ New-Item -ItemType Directory -Force -Path $Destination | Out-Null
+ Copy-Item -Path (Join-Path $Exe.DirectoryName '*') -Destination $Destination -Recurse -Force
+}
+Install-Package 'esrgan' 'https://github.com/xinntao/Real-ESRGAN/releases/download/v0.2.5.0/realesrgan-ncnn-vulkan-20220424-windows.zip' 'realesrgan-ncnn-vulkan.exe'
+Install-Package 'rife' 'https://github.com/nihui/rife-ncnn-vulkan/releases/download/20221029/rife-ncnn-vulkan-20221029-windows.zip' 'rife-ncnn-vulkan.exe'
+Install-Package 'ffmpeg' 'https://www.gyan.dev/ffmpeg/builds/ffmpeg-release-essentials.zip' 'ffmpeg.exe'
+if (!(Test-Path (Join-Path $Root 'ffmpeg/ffprobe.exe'))) { throw 'FFprobe is missing.' }
+if (!(Test-Path (Join-Path $Root 'esrgan/models/realesrgan-x4plus.bin'))) { throw 'Real-ESRGAN weights are missing.' }
+if (!(Get-ChildItem -LiteralPath (Join-Path $Root 'rife') -Directory -Filter 'rife-v4*')) { throw 'RIFE v4 weights are missing.' }
+Write-Output "Local upscale tools and models ready: $Root"
diff --git a/scripts/upscale-host.mjs b/scripts/upscale-host.mjs
new file mode 100644
index 0000000..ef18790
--- /dev/null
+++ b/scripts/upscale-host.mjs
@@ -0,0 +1,100 @@
+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)
+ }
+ }
+}
diff --git a/scripts/upscale-worker.mjs b/scripts/upscale-worker.mjs
new file mode 100644
index 0000000..9c522bc
--- /dev/null
+++ b/scripts/upscale-worker.mjs
@@ -0,0 +1,84 @@
+import { spawn } from 'node:child_process'
+import { readFileSync, mkdirSync, readdirSync, rmSync, existsSync, writeFileSync } from 'node:fs'
+import { join, dirname, resolve } from 'node:path'
+import { fileURLToPath } from 'node:url'
+import { upscaleOptions, upscaleDimensions } from '../shared/video-upscale.mjs'
+
+export async function upscaleVideo(requestPath, { runCommand, report } = {}) {
+const request = JSON.parse(readFileSync(requestPath, 'utf8'))
+const root = request.root, dir = dirname(requestPath), source = join(dir, 'input.mp4')
+const options = upscaleOptions(request)
+const started = Date.now()
+let child
+function event(message, extra = {}) { if (report) report({ message, ...extra }); else process.stdout.write(`AIGEN_EVENT ${JSON.stringify({ message, ...extra })}\n`) }
+function execute(exe, args, cwd = dir) {
+ return new Promise((resolve, reject) => {
+ child = spawn(exe, args, { cwd, windowsHide: true, shell: false, stdio: ['ignore', 'pipe', 'pipe'] })
+ let out = '', tail = ''
+ child.stdout.on('data', b => { out = (out + b).slice(-2_000_000) })
+ child.stderr.on('data', b => { tail = (tail + b).slice(-6000) })
+ child.once('error', reject)
+ child.once('close', code => { child = null; code === 0 ? resolve(out) : reject(new Error(/out of memory|allocate|vkAllocateMemory/i.test(tail) ? 'GPU out of memory. Try 2x with Keep original FPS.' : `${exe.split(/[\\/]/).pop()} failed: ${tail.slice(-1500)}`)) })
+ })
+}
+const run = runCommand || execute
+// Host cancellation kills this entire tree. Also stop descendants if the host disappears.
+const watchdog = setInterval(() => {
+ try { process.kill(Number(request.parentPid), 0) } catch {
+ spawn('taskkill', ['/PID', String(process.pid), '/T', '/F'], { windowsHide: true })
+ }
+}, 2000)
+try {
+ event('Checking local tools and weights', { progress: 1 })
+ await run('powershell.exe', ['-NoProfile', '-NonInteractive', '-ExecutionPolicy', 'Bypass', '-File', fileURLToPath(new URL('./setup-upscale.ps1', import.meta.url)), '-Root', root])
+ const ffmpeg = join(root, 'ffmpeg', 'ffmpeg.exe'), ffprobe = join(root, 'ffmpeg', 'ffprobe.exe')
+ const probe = JSON.parse(await run(ffprobe, ['-v', 'error', '-show_streams', '-show_format', '-of', 'json', source]))
+ const video = probe.streams.find(s => s.codec_type === 'video')
+ if (!video) throw new Error('Input has no video stream.')
+ const rate = String(video.avg_frame_rate).split('/').map(Number)
+ const fps = rate[0] / (rate[1] || 1), duration = Number(video.duration || probe.format.duration)
+ if (!Number.isFinite(fps) || fps <= 0 || fps > 240 || !Number.isFinite(duration) || duration <= 0) throw new Error('Cannot read source frame rate or duration.')
+ const target = upscaleDimensions(video.width, video.height, options)
+ const outputFps = options.fps === 'keep' ? fps : options.fps
+ const interpolate = options.fps !== 'keep'
+ const multiplier = interpolate ? Math.max(2, Math.ceil(outputFps / fps)) : 1
+ const totalFrames = Math.max(1, Math.round(duration * fps))
+ const chunkFrames = Math.max(2, Math.round(fps * 2))
+ const rifeDir = join(root, 'rife')
+ const rifeModel = readdirSync(rifeDir).filter(n => /^rife-v4/.test(n)).sort().at(-1)
+ const info = { sourceWidth: video.width, sourceHeight: video.height, ...target, duration, fps: outputFps, engine: 'Real-ESRGAN' + (interpolate ? ' + RIFE' : '') }
+ event('Upscaling video', { ...info, progress: 3 })
+ const segments = []
+ for (let start = 0, index = 0; start < totalFrames; start += chunkFrames, index++) {
+ const count = Math.min(chunkFrames, totalFrames - start)
+ const work = join(dir, `chunk-${index}`)
+ for (const folder of ['in', 'esrgan', 'sized', 'rife']) mkdirSync(join(work, folder), { recursive: true })
+ const pattern = folder => join(work, folder, '%08d.png')
+ await run(ffmpeg, ['-v', 'error', '-ss', String(start / fps), '-i', source, '-an', '-vf', `fps=${fps},tpad=stop_mode=clone:stop_duration=1`, '-frames:v', String(count + 1), '-start_number', '1', pattern('in')])
+ event(`Upscaling ${Math.round(start / fps)}–${Math.min(Math.round((start + count) / fps), Math.ceil(duration))}s`, { progress: 3 + 90 * start / totalFrames })
+ await run(join(root, 'esrgan', 'realesrgan-ncnn-vulkan.exe'), ['-i', join(work, 'in'), '-o', join(work, 'esrgan'), '-m', join(root, 'esrgan', 'models'), '-n', 'realesrgan-x4plus', '-s', '4', '-t', String(process.env.UPSCALE_TILE || 128), '-g', process.env.UPSCALE_GPU_ID || '0', '-j', '1:1:1', '-f', 'png'])
+ await run(ffmpeg, ['-v', 'error', '-i', pattern('esrgan'), '-vf', `scale=${target.width}:${target.height}:flags=lanczos${options.enhance === 'sharpen' ? ',unsharp=5:5:0.25:5:5:0' : ''}`, '-start_number', '1', pattern('sized')])
+ if (interpolate) {
+ event(`Smoothing motion at ${Math.round(start / fps)}s`, { progress: 3 + 90 * (start + count / 2) / totalFrames })
+ await run(join(rifeDir, 'rife-ncnn-vulkan.exe'), ['-i', join(work, 'sized'), '-o', join(work, 'rife'), '-m', join(rifeDir, rifeModel), '-n', String((count + 1) * multiplier), '-g', process.env.UPSCALE_GPU_ID || '0', '-j', '1:1:1', '-u', '-f', '%08d.png'], rifeDir)
+ }
+ const outputCount = Math.round((start + count) / fps * outputFps) - Math.round(start / fps * outputFps)
+ const segment = `part-${index}.mp4`
+ await run(ffmpeg, ['-v', 'error', '-framerate', String(fps * multiplier), '-i', pattern(interpolate ? 'rife' : 'sized'), '-vf', `fps=${outputFps},setsar=1`, '-frames:v', String(outputCount), '-an', '-c:v', 'libx264', '-preset', 'medium', '-crf', '17', '-pix_fmt', 'yuv420p', '-video_track_timescale', '90000', join(dir, segment)])
+ segments.push(`file '${segment}'`)
+ rmSync(work, { recursive: true, force: true }) // Fixed job-local directory, never an input path.
+ }
+ writeFileSync(join(dir, 'concat.txt'), segments.join('\n'))
+ event('Joining video and preserving original audio', { progress: 96 })
+ await run(ffmpeg, ['-v', 'error', '-f', 'concat', '-safe', '1', '-i', join(dir, 'concat.txt'), '-i', source, '-map', '0:v:0', '-map', '1:a?', '-c:v', 'copy', '-c:a', 'copy', '-map_metadata', '1', '-t', String(duration), '-movflags', '+faststart', join(dir, 'output.mp4')])
+ if (!existsSync(join(dir, 'output.mp4'))) throw new Error('Upscaler did not produce a video.')
+ event('Upscale complete', { ...info, progress: 100, elapsedMs: Date.now() - started })
+} catch (error) {
+ event(error.message, { error: error.message, elapsedMs: Date.now() - started })
+ throw error
+} finally { clearInterval(watchdog) }
+
+}
+if (process.argv[1] && resolve(process.argv[1]) === fileURLToPath(import.meta.url)) {
+ upscaleVideo(process.argv[2]).catch(() => { process.exitCode = 1 })
+}
diff --git a/server/api/generate/active.get.ts b/server/api/generate/active.get.ts
index 9c71d4d..440372e 100644
--- a/server/api/generate/active.get.ts
+++ b/server/api/generate/active.get.ts
@@ -18,6 +18,7 @@ function livePayload(job: Job) {
}
function jobLooksLive(job: Job, comfyBusy: boolean) {
+ if (job.upscale) return false // Clip-level post-processing must not occupy the generation form.
if (job.library?.stopAfterCurrent) return false
if (jobIsLocallySubmitting(job)) return true
if (job.status === 'running' && job.promptId) return comfyBusy
diff --git a/server/api/generate/recover.post.ts b/server/api/generate/recover.post.ts
index 3df8690..9e8e1ef 100644
--- a/server/api/generate/recover.post.ts
+++ b/server/api/generate/recover.post.ts
@@ -82,7 +82,7 @@ export default defineEventHandler(async (event) => {
if (live.status === 'complete' || live.status === 'error' || live.trackId || live.clipId) {
return jobSnapshot(live)
}
- if (live.yueGp) return jobSnapshot(live)
+ if (live.yueGp || live.upscale) return jobSnapshot(live)
if (live.kind === 'music' && Date.now() - live.startedAt > 15_000) {
const recovered = await recoverFinishedMedia(event, { tags: live.library?.tags })
if (recovered?.trackId) return recovered
diff --git a/server/api/interrupt.post.ts b/server/api/interrupt.post.ts
index 785782f..dfaa9a3 100644
--- a/server/api/interrupt.post.ts
+++ b/server/api/interrupt.post.ts
@@ -28,9 +28,19 @@ export default defineEventHandler(async (event) => {
const jobId = String(body?.jobId || '')
const kind = String(body?.kind || '')
const named = jobId ? getJob(jobId) : undefined
+ if (named?.upscale) {
+ const { owner } = assertLibraryOwner(event)
+ if (named.library?.ownerKey !== owner) throw createError({ statusCode: 404, statusMessage: 'Job not found' })
+ const { cancelUpscaleJob } = await import('~/server/utils/videoUpscale')
+ const { onLiveVideoSettled } = await import('~/server/utils/studioQueue')
+ await cancelUpscaleJob(named)
+ await onLiveVideoSettled(named)
+ return { ok: true, cancelled: 1 }
+ }
const targets = named
? [named]
: listJobs().filter((job) => {
+ if (job.upscale) return false
if (job.status === 'complete' || job.status === 'error' || job.status === 'cancelled') return false
if (kind === 'edit') return job.kind === 'edit'
if (kind === 'music') return job.kind === 'music'
@@ -46,6 +56,7 @@ export default defineEventHandler(async (event) => {
if (jobId) deletePendingJob(jobId)
const { owner } = assertLibraryOwner(event)
for (const row of listStudioJobs(owner)) {
+ if (row.payload.upscale) continue
if (row.status !== 'running' && row.status !== 'held') continue
if (kind === 'edit' && row.kind !== 'edit') continue
if (kind === 'music' && row.kind !== 'music') continue
diff --git a/server/api/library/clips/[id]/upscale.get.ts b/server/api/library/clips/[id]/upscale.get.ts
new file mode 100644
index 0000000..b28ce96
--- /dev/null
+++ b/server/api/library/clips/[id]/upscale.get.ts
@@ -0,0 +1,24 @@
+import { basename, parse } from 'node:path'
+import { upscaleRecords } from '../../../../utils/videoUpscale'
+import { listStudioJobs } from '../../../../utils/studioQueue'
+import { listJobs } from '../../../../utils/jobs'
+
+export default defineEventHandler(event => {
+ const { owner } = assertLibraryOwner(event)
+ const clip = getClip(owner, String(getRouterParam(event, 'id') || ''))
+ assertFolderAccess(event, clip.folderId)
+ const rows = listStudioJobs(owner)
+ const active = rows.find(row => ['waiting', 'held', 'running'].includes(row.status) && row.payload.upscale?.sourceId === clip.id)
+ const records = upscaleRecords(owner, clip.id)
+ const latest = active || rows.filter(row => row.payload.upscale?.sourceId === clip.id).sort((a, b) => b.createdAt - a.createdAt)[0]
+ const record = records.find(r => !latest || r.id === latest.id)
+ const status = active?.status || record?.status || latest?.status
+ const labels: Record = { waiting: 'queued', held: 'queued', complete: 'done', error: 'failed', cancelled: 'failed', running: 'running' }
+ const generationBusy = rows.some(row => row.status === 'running' && row.payload.extendFromClipId === clip.id)
+ || listJobs().some(job => job.library?.ownerKey === owner && job.clipId === clip.id && ['queued', 'uploading', 'running'].includes(job.status))
+ return { source: { id: clip.id, name: clip.name, upscaledFromClipId: clip.upscaledFromClipId, stem: parse(basename(clipVideoPath(owner, clip.id))).name, width: clip.width, height: clip.height },
+ busy: Boolean(active || record?.status === 'running' || generationBusy),
+ job: latest || record ? { id: latest?.id || record.id, status: labels[status || ''] || status,
+ message: record?.message || latest?.waitReason, error: record?.error || latest?.lastError || (status === 'cancelled' ? 'Cancelled. Original video is safe.' : undefined), progress: record?.progress || 0,
+ width: record?.width, height: record?.height, elapsedMs: record?.elapsedMs || (record?.createdAt ? Date.now() - record.createdAt : 0), outputClipId: record?.outputClipId } : null }
+})
diff --git a/server/api/library/clips/[id]/upscale.post.ts b/server/api/library/clips/[id]/upscale.post.ts
new file mode 100644
index 0000000..aca6e77
--- /dev/null
+++ b/server/api/library/clips/[id]/upscale.post.ts
@@ -0,0 +1,29 @@
+import { existsSync } from 'node:fs'
+import { upscaleOptions } from '~/shared/video-upscale.mjs'
+import { addStudioJob, listStudioJobs, kickStudioQueue, type StudioJobPayload } from '../../../../utils/studioQueue'
+import { listJobs } from '../../../../utils/jobs'
+
+const submitting = new Set()
+export default defineEventHandler(async event => {
+ const { owner } = assertLibraryOwner(event)
+ const clip = getClip(owner, String(getRouterParam(event, 'id') || ''))
+ assertFolderAccess(event, clip.folderId)
+ const body = await readBody(event)
+ let options
+ try { options = upscaleOptions(body) } catch { throw createError({ statusCode: 400, statusMessage: 'Invalid upscale settings.' }) }
+ const key = `${owner}:${clip.id}`
+ if (submitting.has(key)) throw createError({ statusCode: 409, statusMessage: 'This clip is already being queued.' })
+ submitting.add(key)
+ try {
+ if (!existsSync(clipVideoPath(owner, clip.id))) throw createError({ statusCode: 404, statusMessage: 'Completed video file is missing.' })
+ if (listStudioJobs(owner).some(row => ['waiting', 'held', 'running'].includes(row.status) && (row.payload.upscale?.sourceId === clip.id || row.status === 'running' && row.payload.extendFromClipId === clip.id))
+ || listJobs().some(job => job.library?.ownerKey === owner && job.clipId === clip.id && ['running', 'queued', 'uploading'].includes(job.status))) throw createError({ statusCode: 409, statusMessage: 'This clip already has a generation or upscale in progress.' })
+ const payload = { ...clip, name: `Upscale · ${clip.name || 'Video'}`, prompt: clip.prompt || 'Upscale video', folderId: clip.folderId,
+ steps: 0, cfg: 0, seed: 0, turbo: false, fps: clip.fps || 24, duration: clip.duration || 0, sound: clip.sound !== false,
+ samplerName: '', scheduler: '', useIdentityRefs: false, referenceStillIds: [], extensions: [], queueAutoRun: true,
+ upscale: { ...options, sourceId: clip.id } } as StudioJobPayload
+ const job = await addStudioJob({ ownerKey: owner, payload, familyId: crypto.randomUUID(), kind: 'video' })
+ kickStudioQueue()
+ return { id: job.id, status: 'queued' }
+ } finally { submitting.delete(key) }
+})
diff --git a/server/api/studio-queue/[id].patch.ts b/server/api/studio-queue/[id].patch.ts
index 6b26c50..2b38f1c 100644
--- a/server/api/studio-queue/[id].patch.ts
+++ b/server/api/studio-queue/[id].patch.ts
@@ -28,6 +28,7 @@ export default defineEventHandler(async (event) => {
if (!current) {
throw createError({ statusCode: 404, statusMessage: 'Queued job not found' })
}
+ if (current.payload.upscale) throw createError({ statusCode: 409, statusMessage: 'Cancel this upscale and use the clip dialog to change settings.' })
if (!EDITABLE.has(current.status)) {
throw createError({ statusCode: 409, statusMessage: 'That job is already generating. Copy it or save a template.' })
}
diff --git a/server/plugins/00-resume-upscale.ts b/server/plugins/00-resume-upscale.ts
new file mode 100644
index 0000000..0d39e90
--- /dev/null
+++ b/server/plugins/00-resume-upscale.ts
@@ -0,0 +1,2 @@
+import { resumeUpscaleJobs } from '../utils/videoUpscale'
+export default defineNitroPlugin(() => { resumeUpscaleJobs() })
diff --git a/server/utils/jobs.ts b/server/utils/jobs.ts
index 5b49193..d104159 100644
--- a/server/utils/jobs.ts
+++ b/server/utils/jobs.ts
@@ -32,6 +32,7 @@ export interface JobEvent {
}
export interface Job {
+ upscale?: boolean
yueGp?: boolean
musicActivity?: { checkedAt: number; running: boolean }
id: string
diff --git a/server/utils/library.ts b/server/utils/library.ts
index d3b6973..a64febc 100644
--- a/server/utils/library.ts
+++ b/server/utils/library.ts
@@ -1,7 +1,7 @@
import { preserveVideoSources } from './videoSources'
import { createHash } from 'node:crypto'
import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, rmSync, writeFileSync, statSync, createReadStream, openSync, readSync, closeSync } from 'node:fs'
-import { writeFile } from 'node:fs/promises'
+import { writeFile, copyFile } from 'node:fs/promises'
import { join } from 'node:path'
import { spawn } from 'node:child_process'
import { randomBytes, scrypt, timingSafeEqual } from 'node:crypto'
@@ -54,6 +54,8 @@ export interface LibraryClip {
scheduler?: string
duration?: number
familyId?: string
+ upscaledFromClipId?: string
+ upscaleJobId?: string
parentClipId?: string
chainIndex?: number
workflow?: import('~/utils/videoModels').VideoWorkflowId
@@ -1669,7 +1671,8 @@ export async function saveClip(params: {
seed: number
hideThumbnail: boolean
hideInput?: boolean
- video: Buffer
+ video?: Buffer
+ videoFile?: string
thumb?: Buffer | null
comfyFilename?: string
cfg?: number
@@ -1678,6 +1681,8 @@ export async function saveClip(params: {
scheduler?: string
duration?: number
familyId?: string
+ upscaledFromClipId?: string
+ upscaleJobId?: string
parentClipId?: string
chainIndex?: number
workflow?: import('~/utils/videoModels').VideoWorkflowId
@@ -1715,6 +1720,8 @@ export async function saveClip(params: {
samplerName: params.samplerName,
scheduler: params.scheduler,
familyId: params.familyId,
+ upscaledFromClipId: params.upscaledFromClipId,
+ upscaleJobId: params.upscaleJobId,
parentClipId: params.parentClipId,
chainIndex: params.chainIndex,
workflow: params.workflow,
@@ -1730,7 +1737,9 @@ export async function saveClip(params: {
}
mkdirSync(clipDir(params.ownerKey, clip.id), { recursive: true })
const videoPath = clipVideoPath(params.ownerKey, clip.id)
- await writeFile(videoPath, params.video)
+ if (params.videoFile) await copyFile(params.videoFile, videoPath)
+ else if (params.video) await writeFile(videoPath, params.video)
+ else throw new Error('Video data is missing')
await remuxFaststart(videoPath)
if (params.originalSegment) preserveVideoSources(videoPath, params.sourceSegments || [], params.originalSegment)
if (params.segmentFirstFrame?.length && params.segmentFirstFrame.length >= 64) {
@@ -2200,3 +2209,5 @@ export function groupLibraryStills(owner: string, ids: string[], ungroup = false
})
})
}
+
+export function findUpscaledClip(owner: string, jobId: string) { return readCatalog(owner).clips.find(clip => clip.upscaleJobId === jobId) }
diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts
index ff48d6d..58f8065 100644
--- a/server/utils/studioQueue.ts
+++ b/server/utils/studioQueue.ts
@@ -16,6 +16,7 @@ export type StudioJobStatus = 'waiting' | 'running' | 'held' | 'complete' | 'err
export type StudioJobKind = 'video' | 'edit' | 'music'
export interface StudioJobPayload {
+ upscale?: { sourceId: string; scale: 2 | 4; target: string; fps: string | number; enhance: string }
prompt: string
promptMid?: string
promptPre?: string
@@ -284,7 +285,7 @@ function failZombieLiveJob(job: Job, error: string) {
function sweepStaleLiveJobs() {
const now = Date.now()
for (const job of listJobs()) {
- if (job.yueGp) continue
+ if (job.yueGp || job.upscale) continue
if (job.saving) continue
if (job.status === 'queued' && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) {
failZombieLiveJob(job, 'Job never started')
@@ -304,7 +305,7 @@ async function reapZombieLiveJobs() {
const { fetchHistory } = await import('~/server/utils/comfy')
for (const job of listJobs()) {
// Never interrupt download/stitch/library save — Comfy is idle then by design.
- if (job.yueGp) continue
+ if (job.yueGp || job.upscale) continue
if (job.saving) continue
if (job.library?.chainContinuing) continue
if (job.status !== 'running' && job.status !== 'uploading' && job.status !== 'queued') continue
@@ -349,7 +350,7 @@ async function reapZombieLiveJobs() {
}
function liveJobOwnsGpu(job: Job) {
- if (job.yueGp && ['running', 'queued', 'uploading'].includes(job.status)) return true
+ if ((job.yueGp || job.upscale) && ['running', 'queued', 'uploading'].includes(job.status)) return true
if (job.library?.stopAfterCurrent) return false
if (job.library?.chainContinuing) return true
if (job.saving) return true
@@ -369,7 +370,7 @@ function liveJobOwnsGpu(job: Job) {
*/
function clearDeadGpuClaimsForForceStart() {
for (const live of listJobs()) {
- if (live.yueGp) continue
+ if (live.yueGp || live.upscale) continue
if (live.saving) continue
if (jobIsLocallySubmitting(live)) continue
if (live.status !== 'running' && live.status !== 'queued' && live.status !== 'uploading') {
@@ -591,6 +592,7 @@ export async function clearStuckStudioWork(owner: string) {
if (job.library?.ownerKey !== owner) continue
if (job.status === 'complete' || job.status === 'error' || job.status === 'cancelled') continue
liveIds.add(job.id)
+ if (job.upscale) { const { cancelUpscaleJob } = await import('./videoUpscale'); await cancelUpscaleJob(job) }
if (job.yueGp) {
const { cancelYueGpJob } = await import('./yueGp')
await cancelYueGpJob(job)
@@ -656,6 +658,7 @@ async function stopLiveGeneration(liveJobId?: string, shotQueueId?: string) {
}
if (!liveJobId) return
const live = getJob(liveJobId)
+ if (live?.upscale) { const { cancelUpscaleJob } = await import('./videoUpscale'); await cancelUpscaleJob(live); return }
if (live?.yueGp) {
const { cancelYueGpJob } = await import('./yueGp')
await cancelYueGpJob(live)
@@ -910,7 +913,7 @@ function repairStaleJobs(jobs: StudioJob[]) {
job.updatedAt = Date.now()
continue
}
- if (job.payload.musicEngine !== 'yue' && job.status === 'error' && isTransientComfyError(job.lastError)) {
+ if (!job.payload.upscale && job.payload.musicEngine !== 'yue' && job.status === 'error' && isTransientComfyError(job.lastError)) {
job.status = 'waiting'
job.liveJobId = undefined
job.lastError = undefined
@@ -919,7 +922,7 @@ function repairStaleJobs(jobs: StudioJob[]) {
job.updatedAt = Date.now()
continue
}
- if (job.payload.musicEngine !== 'yue' && job.status === 'held' && isTransientComfyError(job.lastError)) {
+ if (!job.payload.upscale && job.payload.musicEngine !== 'yue' && job.status === 'held' && isTransientComfyError(job.lastError)) {
job.status = 'waiting'
job.liveJobId = undefined
job.lastError = undefined
@@ -1419,6 +1422,11 @@ export async function startStudioJob(item: StudioJob) {
}
async function startStudioJobReserved(item: StudioJob) {
+ if (item.payload.upscale) {
+ try { const { startUpscaleJob } = await import('./videoUpscale'); await startUpscaleJob(item) }
+ catch (error) { await patchStudioJob(item.ownerKey, item.id, row => { row.status = 'error'; row.lastError = error instanceof Error ? error.message : String(error) }); scheduleKickRetry() }
+ return
+ }
if (studioJobKind(item) === 'music') {
await startStudioMusicJob(item)
return
@@ -1665,7 +1673,7 @@ export async function onLiveVideoSettled(job: Job) {
return
}
const remaining = remainingStudioShots(job)
- const wakeFail = !job.yueGp && job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
+ const wakeFail = !job.yueGp && !job.upscale && job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail
await mutateStore(owner, (store) => {
diff --git a/server/utils/videoUpscale.ts b/server/utils/videoUpscale.ts
new file mode 100644
index 0000000..b22a75b
--- /dev/null
+++ b/server/utils/videoUpscale.ts
@@ -0,0 +1,136 @@
+import { createReadStream, createWriteStream, existsSync, mkdirSync, readFileSync, writeFileSync, renameSync, readdirSync, constants, copyFileSync, unlinkSync } from 'node:fs'
+import { basename, dirname, join, parse } from 'node:path'
+import { Readable } from 'node:stream'
+import { pipeline } from 'node:stream/promises'
+import { upscaleOptions, upscaleName } from '~/shared/video-upscale.mjs'
+import { createJob, restoreJob, getJob, emitJob, type Job } from './jobs'
+import { clipVideoPath, getClip, saveClip, findUpscaledClip } from './library'
+import { sharedGpuHeaders } from './sharedGpu'
+import { patchStudioJob, onLiveVideoSettled, type StudioJob } from './studioQueue'
+
+function root() { return join(String(useRuntimeConfig().libraryDir || process.env.LIBRARY_DIR || '/data/library'), 'upscale-jobs') }
+function path(id: string) { if (!/^[\w-]+$/.test(id)) throw new Error('Invalid job ID'); return join(root(), `${id}.json`) }
+function persist(record: any) { mkdirSync(root(), { recursive: true }); writeFileSync(path(record.id) + '.tmp', JSON.stringify(record)); renameSync(path(record.id) + '.tmp', path(record.id)) }
+export function upscaleRecords(owner: string, sourceId?: string) {
+ if (!existsSync(root())) return []
+ return readdirSync(root()).filter(n => n.endsWith('.json')).flatMap(n => { try { const r = JSON.parse(readFileSync(join(root(), n), 'utf8')); return r.owner === owner && (!sourceId || r.sourceId === sourceId) ? [r] : [] } catch { return [] } }).sort((a, b) => b.createdAt - a.createdAt)
+}
+async function request(endpoint: string, method = 'GET', body?: any) {
+ const config = useRuntimeConfig()
+ const url = String(config.comfyControlUrl || process.env.COMFY_CONTROL_URL || '').replace(/\/$/, '')
+ if (!url) throw new Error('Local GPU host is not configured.')
+ const token = String(config.comfyControlToken || process.env.COMFY_CONTROL_TOKEN || '')
+ const response = await fetch(`${url}/upscale/${endpoint}`, { method,
+ headers: { ...(token ? { Authorization: `Bearer ${token}` } : {}), ...(method === 'GET' ? {} : sharedGpuHeaders()), ...(method === 'POST' ? { 'Content-Type': 'application/json' } : {}) },
+ body: method === 'POST' ? JSON.stringify(body) : body, ...(method === 'PUT' ? { duplex: 'half' } : {}), signal: AbortSignal.timeout(method === 'PUT' ? 600_000 : 60_000)
+ } as RequestInit)
+ if (!response.ok) { const detail = await response.json().catch(() => ({})); throw Object.assign(new Error(detail.error || detail.message || `Local upscale host returned ${response.status}`), { statusCode: response.status }) }
+ return response
+}
+export async function cancelUpscaleJob(job: Job) {
+ await request(`jobs/${job.id}/cancel`, 'POST', {})
+ job.status = 'cancelled'; emitJob(job, { type: 'error', message: 'Cancelled', error: 'Cancelled' })
+}
+async function watch(job: Job, record: any) {
+ let saveFailures = 0
+ for (;;) {
+ let hostFinished = false
+ try {
+ const state = await (await request(`jobs/${job.id}`)).json()
+ Object.assign(record, { ...state, id: record.id, liveId: job.id })
+ if (['error', 'cancelled'].includes(state.status)) {
+ job.status = state.status; job.error = String(state.error || 'Upscale cancelled. Original video is safe.').slice(0, 500)
+ record.error = job.error; persist(record); emitJob(job, { type: 'error', error: job.error, message: job.error })
+ await onLiveVideoSettled(job); return
+ }
+ if (state.status === 'complete') {
+ hostFinished = true
+ // Saving is part of the same durable job; retries reuse the sibling and catalog entry.
+ record.status = 'running'; persist(record); job.saving = true
+ const source = getClip(record.owner, record.sourceId)
+ const sourcePath = clipVideoPath(record.owner, source.id)
+ if (!record.sibling) {
+ const download = join(root(), `${record.id}.mp4.partial`)
+ const response = await request(`jobs/${job.id}/video`)
+ if (!response.body) throw new Error('Upscaled video download is empty.')
+ await pipeline(Readable.fromWeb(response.body as any), createWriteStream(download))
+ let index = 1
+ for (;;) {
+ const sibling = join(dirname(sourcePath), upscaleName(parse(sourcePath).name, record.options.scale, index))
+ try { copyFileSync(download, sibling, constants.COPYFILE_EXCL); record.sibling = sibling; break } catch (error) { if ((error as NodeJS.ErrnoException).code !== 'EEXIST') throw error; index++ }
+ }
+ persist(record)
+ unlinkSync(download)
+ }
+ if (!record.outputClipId) {
+ const existing = findUpscaledClip(record.owner, record.id)
+ const clip = existing || await saveClip({ ...source, ownerKey: record.owner, folderId: source.folderId,
+ name: `${source.name || 'Video'} · Upscale ${record.options.scale}x`, width: state.width, height: state.height,
+ fps: state.fps, duration: state.duration, videoFile: record.sibling,
+ familyId: undefined, parentClipId: undefined, chainIndex: undefined, stillId: undefined,
+ upscaledFromClipId: source.id, upscaleJobId: record.id, comfyFilename: basename(record.sibling) })
+ record.outputClipId = clip.id
+ }
+ record.status = 'complete'; persist(record)
+ job.clipId = record.outputClipId; job.status = 'complete'; job.saving = false
+ emitJob(job, { type: 'complete', progress: 100, clipId: job.clipId, message: 'Upscale ready', elapsedMs: state.elapsedMs })
+ console.info('[upscale]', JSON.stringify({ source: `${state.sourceWidth}x${state.sourceHeight}`, target: `${state.width}x${state.height}`, duration: state.duration, engine: state.engine, elapsedMs: state.elapsedMs }))
+ await onLiveVideoSettled(job); return
+ }
+ record.status = 'running'; persist(record); job.status = 'running'
+ emitJob(job, { type: 'progress', message: state.message, progress: state.progress })
+ } catch (error) {
+ // A lost connection is not proof the GPU stopped. Keep its queue slot until the host confirms exit.
+ record.message = `Checking local worker: ${error instanceof Error ? error.message : String(error)}`
+ if (hostFinished && ++saveFailures >= 3) {
+ record.status = 'error'; record.error = `Could not save upscale: ${error instanceof Error ? error.message : String(error)}`.slice(0, 500)
+ job.status = 'error'; job.saving = false; job.error = record.error; persist(record)
+ emitJob(job, { type: 'error', error: job.error, message: job.error }); await onLiveVideoSettled(job); return
+ }
+ if ((error as any).statusCode === 404) {
+ record.status = 'error'; record.error = 'Local upscale job was not found. Original video is safe.'
+ job.status = 'error'; job.error = record.error; persist(record); await onLiveVideoSettled(job); return
+ }
+ persist(record); emitJob(job, { type: 'status', message: record.message })
+ }
+ if (job.status === 'cancelled') {
+ Object.assign(record, { status: 'cancelled', error: 'Cancelled. Original video is safe.' }); persist(record); await onLiveVideoSettled(job); return
+ }
+ await new Promise(resolve => setTimeout(resolve, 2000))
+ }
+}
+export async function startUpscaleJob(item: StudioJob) {
+ const options = upscaleOptions(item.payload.upscale), sourceId = item.payload.upscale!.sourceId
+ const source = getClip(item.ownerKey, sourceId)
+ const job = createJob('video'); job.upscale = true
+ job.library = { ...source, ownerKey: item.ownerKey, extensions: [], queueAutoRun: false }
+ const record = { id: item.id, liveId: job.id, clientId: job.clientId, owner: item.ownerKey, sourceId, options, createdAt: Date.now(), status: 'running', library: job.library }
+ persist(record)
+ const attached = await patchStudioJob(item.ownerKey, item.id, row => {
+ if (row.status === 'cancelled') return
+ row.status = 'running'; row.liveJobId = job.id; row.lastError = undefined
+ })
+ if (attached.status === 'cancelled') { job.status = 'cancelled'; record.status = 'cancelled'; persist(record); return }
+ void (async () => {
+ try {
+ await request(`jobs/${job.id}/input`, 'PUT', createReadStream(clipVideoPath(item.ownerKey, sourceId)))
+ } catch (error) {
+ job.status = 'error'; job.error = `Video upload failed: ${error instanceof Error ? error.message : String(error)}`
+ Object.assign(record, { status: 'error', error: job.error }); persist(record); await onLiveVideoSettled(job); return
+ }
+ if (job.status === 'cancelled') { Object.assign(record, { status: 'cancelled' }); persist(record); await onLiveVideoSettled(job); return }
+ try { await request('jobs', 'POST', { id: job.id, ...options }) } catch { /* Confirm by ID; never generate a duplicate on an HTTP timeout. */ }
+ await watch(job, record)
+ })()
+}
+export function resumeUpscaleJobs() {
+ if (!existsSync(root())) return
+ for (const name of readdirSync(root()).filter(n => n.endsWith('.json'))) {
+ try {
+ const r = JSON.parse(readFileSync(join(root(), name), 'utf8'))
+ if (r.status !== 'running' || getJob(r.liveId)) continue
+ const job = restoreJob({ id: r.liveId, clientId: r.clientId, startedAt: r.createdAt, promptId: '', library: r.library })
+ job.upscale = true; void watch(job, r)
+ } catch { /* Preserve malformed records for diagnosis. */ }
+ }
+}
diff --git a/shared/video-upscale.mjs b/shared/video-upscale.mjs
new file mode 100644
index 0000000..1d9f85f
--- /dev/null
+++ b/shared/video-upscale.mjs
@@ -0,0 +1,12 @@
+export function upscaleOptions(value = {}) {
+ const options = { scale: value.scale ?? 2, target: value.target ?? 'preserve', fps: value.fps ?? 'keep', enhance: value.enhance ?? 'faithful' }
+ if (![2, 4].includes(options.scale) || !['preserve', '1080p', '4k'].includes(options.target)
+ || !['keep', 48, 50, 60].includes(options.fps) || !['faithful', 'sharpen'].includes(options.enhance)) throw new Error('Invalid upscale settings.')
+ return options
+}
+export function upscaleDimensions(width, height, options) {
+ if (![width, height].every(n => Number.isFinite(n) && n > 0)) throw new Error('Invalid source resolution.')
+ const factor = options.target === 'preserve' ? options.scale : (options.target === '1080p' ? 1920 : 3840) / Math.max(width, height)
+ return { width: Math.max(2, Math.round(width * factor / 2) * 2), height: Math.max(2, Math.round(height * factor / 2) * 2) }
+}
+export function upscaleName(stem, scale, index = 1) { return `${stem}_up${scale}x${index > 1 ? `-${index}` : ''}.mp4` }
diff --git a/tests/studio-queue.test.mjs b/tests/studio-queue.test.mjs
index 0089971..e661b6b 100644
--- a/tests/studio-queue.test.mjs
+++ b/tests/studio-queue.test.mjs
@@ -27,6 +27,7 @@ function fixture({ rows = [], lives = [], queue = { running: 0, pending: 0 }, hi
const timers = []
const started = []
const modules = {
+ './videoUpscale': { startUpscaleJob: async item => { started.push({ upscale: true, sourceId: item.payload.upscale.sourceId }) } },
'~/server/utils/sharedGpu': { withSharedGpuStart: async fn => fn(), maintainSharedGpu: async () => {}, acquireSharedGpu: async () => false, sharedGpuWaitReason: () => 'GPU is in use. Waiting for availability.' },
'~/server/utils/jobs': {
listJobs: () => lives,
@@ -93,6 +94,20 @@ test('standalone YuEGP retains its queue slot while Comfy is empty for longer th
assert.equal(f.started.length, 0)
assert.equal(f.rows()[1].status, 'waiting')
})
+test('standalone upscale retains its queue slot while Comfy is empty for longer than the zombie timeout', async () => {
+ const live = liveJob('upscale-active', 'video', { upscale: true, promptId: undefined, startedAt: Date.now() - 3600000 })
+ const f = fixture({ rows: [row('upscale-active'), row('next', 'waiting', undefined)], lives: [live] })
+ await f.api.kickStudioQueue()
+ assert.equal(live.status, 'running')
+ assert.equal(f.started.length, 0)
+ assert.equal(f.rows()[1].status, 'waiting')
+})
+test('upscale queue row dispatches to the standalone worker before any video generation path', async () => {
+ const item = { ...row('upscale-waiting', 'waiting', undefined), kind: 'video', payload: { prompt: 'existing video', workflow: 'ltx-cherry', extensions: [], upscale: { sourceId: 'finished-master', scale: 2, fps: 'keep', target: 'preserve', enhance: 'faithful' } } }
+ const f = fixture({ rows: [item], queue: null })
+ await f.api.startStudioJob(item)
+ assert.deepEqual(f.started, [{ upscale: true, sourceId: 'finished-master' }])
+})
async function finishes(promise) {
let timer
try {
diff --git a/tests/video-upscale.test.mjs b/tests/video-upscale.test.mjs
new file mode 100644
index 0000000..79633af
--- /dev/null
+++ b/tests/video-upscale.test.mjs
@@ -0,0 +1,108 @@
+import test from 'node:test'
+import assert from 'node:assert/strict'
+import { EventEmitter } from 'node:events'
+import { PassThrough, Readable } from 'node:stream'
+import { mkdtempSync, mkdirSync, writeFileSync, readFileSync, existsSync } from 'node:fs'
+import { tmpdir } from 'node:os'
+import { join } from 'node:path'
+import { upscaleOptions, upscaleDimensions, upscaleName } from '../shared/video-upscale.mjs'
+import { createUpscaleHost } from '../scripts/upscale-host.mjs'
+import { upscaleVideo } from '../scripts/upscale-worker.mjs'
+import { execFile } from 'node:child_process'
+import { promisify } from 'node:util'
+const exec = promisify(execFile)
+
+test('upscale defaults, delivery sizes, portrait and square aspect, and collision names', () => {
+ assert.deepEqual(upscaleOptions(), { scale: 2, target: 'preserve', fps: 'keep', enhance: 'faithful' })
+ assert.deepEqual(upscaleDimensions(1344, 768, upscaleOptions()), { width: 2688, height: 1536 })
+ assert.deepEqual(upscaleDimensions(768, 1344, upscaleOptions({ target: '1080p' })), { width: 1098, height: 1920 })
+ assert.deepEqual(upscaleDimensions(960, 960, upscaleOptions({ target: '4k' })), { width: 3840, height: 3840 })
+ assert.equal(upscaleName('master', 4, 3), 'master_up4x-3.mp4')
+ for (const value of [{ scale: 3 }, { fps: 30 }, { target: 'crop' }, { enhance: 'creative' }]) assert.throws(() => upscaleOptions(value))
+})
+function fixture() {
+ const root = mkdtempSync(join(tmpdir(), 'aigen-upscale-'))
+ const calls = [], child = new EventEmitter()
+ child.stdout = new PassThrough(); child.stderr = new PassThrough()
+ const host = createUpscaleHost({ root, leaseValid: token => token === 'lease', spawnProcess: (...args) => { calls.push(args); return child } })
+ return { root, calls, child, host }
+}
+test('standalone upscale requires a lease, uploads locally, deduplicates, and reports OOM after exit', async () => {
+ const f = fixture(), id = 'test-upscale-12345'
+ await f.host.upload(id, Readable.from(Buffer.alloc(64)))
+ assert.throws(() => f.host.start({ id }, 'invalid'), /reservation/)
+ f.host.start({ id }, 'lease'); f.host.start({ id }, 'lease')
+ assert.equal(f.calls.length, 1)
+ assert.equal(f.calls[0][2].windowsHide, true)
+ assert.match(f.calls[0][1][0], /upscale-worker/)
+ assert.equal(f.host.busy(), true)
+ f.child.stdout.write('AIGEN_EVENT {"error":"GPU out of ')
+ f.child.stdout.write('memory"}\n')
+ assert.equal(f.host.output(id), null)
+ f.child.emit('close', 1)
+ assert.equal(f.host.busy(), false)
+ assert.equal(f.host.read(id).status, 'error')
+ assert.match(f.host.read(id).error, /out of memory/)
+ assert.equal(readFileSync(join(f.root, 'jobs', id, 'input.mp4')).length, 64)
+})
+test('cancel before submission is a durable tombstone and cannot later start a GPU worker', async () => {
+ const f = fixture(), id = 'cancel-upscale-12345'
+ await f.host.cancel(id)
+ assert.equal(f.host.start({ id }, 'lease').status, 'cancelled')
+ assert.equal(f.calls.length, 0)
+ await assert.rejects(f.host.upload('../escape', Readable.from('x')), /Invalid/)
+})
+test('two-minute fractional-FPS master stays one job; ESRGAN precedes RIFE and audio is stream-copied', async () => {
+ const root = mkdtempSync(join(tmpdir(), 'aigen-upscale-plan-'))
+ mkdirSync(join(root, 'rife', 'rife-v4.6'), { recursive: true })
+ const request = join(root, 'request.json'), events = [], calls = []
+ writeFileSync(request, JSON.stringify({ root, parentPid: process.pid, fps: 50 }))
+ let outputFrames = 0
+ await upscaleVideo(request, { report: e => events.push(e), runCommand: async (exe, args) => {
+ calls.push({ exe, args })
+ if (exe.endsWith('ffprobe.exe')) return JSON.stringify({ streams: [{ codec_type: 'video', width: 1344, height: 768, avg_frame_rate: '30000/1001', duration: '120.02' }] })
+ if (args.includes('libx264')) outputFrames += Number(args[args.indexOf('-frames:v') + 1])
+ if (args.at(-1) === join(root, 'output.mp4')) writeFileSync(args.at(-1), 'mock result')
+ return ''
+ } })
+ assert.equal(outputFrames, Math.round(Math.round(120.02 * 30000 / 1001) / (30000 / 1001) * 50))
+ const engines = calls.filter(c => /ncnn-vulkan/.test(c.exe))
+ assert.ok(engines.length > 100)
+ for (let i = 0; i < engines.length; i += 2) { assert.match(engines[i].exe, /realesrgan/); assert.match(engines[i + 1].exe, /rife/) }
+ const mux = calls.at(-1).args
+ assert.equal(mux[mux.indexOf('-c:a') + 1], 'copy')
+ assert.ok(mux.includes('1:a?'))
+ assert.ok(!mux.includes('-shortest'))
+ assert.equal(events.at(-1).progress, 100)
+ assert.ok(existsSync(join(root, 'output.mp4')))
+})
+
+const cpuFfmpeg = process.env.UPSCALE_TEST_FFMPEG || join(process.cwd(), '.data', 'video-test-tools', 'ffmpeg.exe')
+for (const hasAudio of [true, false]) test(`CPU-only fixture: chunk boundary, frame count and ${hasAudio ? 'bit-identical audio' : 'video-only output'}`, { skip: !existsSync(cpuFfmpeg) }, async () => {
+ const root = mkdtempSync(join(tmpdir(), 'aigen-upscale-cpu-'))
+ mkdirSync(join(root, 'rife', 'rife-v4.6'), { recursive: true })
+ const source = join(root, 'input.mp4'), result = join(root, 'output.mp4')
+ const ffmpeg = async args => (await exec(cpuFfmpeg, args, { windowsHide: true, maxBuffer: 4_000_000 })).stdout
+ await ffmpeg(['-v', 'error', '-f', 'lavfi', '-i', 'testsrc2=size=48x32:rate=24:duration=2.5', ...(hasAudio ? ['-f', 'lavfi', '-i', 'sine=frequency=440:duration=2.5'] : []), '-c:v', 'libx264', ...(hasAudio ? ['-c:a', 'aac'] : []), source])
+ const original = readFileSync(source)
+ writeFileSync(join(root, 'request.json'), JSON.stringify({ root, parentPid: process.pid, fps: hasAudio ? 50 : 'keep' }))
+ await upscaleVideo(join(root, 'request.json'), { report: () => {}, runCommand: async (exe, args) => {
+ if (exe === 'powershell.exe') return ''
+ if (exe.endsWith('ffprobe.exe')) return JSON.stringify({ streams: [{ codec_type: 'video', width: 48, height: 32, avg_frame_rate: '24/1', duration: '2.5' }] })
+ if (exe.endsWith('realesrgan-ncnn-vulkan.exe')) {
+ // Substitute a CPU scaler only in this test. Never execute the neural engine.
+ return ffmpeg(['-v', 'error', '-i', join(args[args.indexOf('-i') + 1], '%08d.png'), '-vf', 'scale=192:128', '-threads', '1', join(args[args.indexOf('-o') + 1], '%08d.png')])
+ }
+ if (exe.endsWith('rife-ncnn-vulkan.exe')) {
+ return ffmpeg(['-v', 'error', '-framerate', '24', '-i', join(args[args.indexOf('-i') + 1], '%08d.png'), '-vf', 'fps=72', '-frames:v', args[args.indexOf('-n') + 1], '-threads', '1', join(args[args.indexOf('-o') + 1], '%08d.png')])
+ }
+ return ffmpeg(args)
+ } })
+ assert.deepEqual(readFileSync(source), original)
+ const hashAudio = file => ffmpeg(['-v', 'error', '-i', file, '-map', '0:a:0', '-c:a', 'copy', '-f', 'hash', '-'])
+ if (hasAudio) assert.equal(await hashAudio(result), await hashAudio(source))
+ else await assert.rejects(hashAudio(result), /matches no streams/)
+ const frames = await ffmpeg(['-v', 'error', '-i', result, '-map', '0:v:0', '-f', 'framemd5', '-'])
+ assert.equal(frames.split('\n').filter(line => /^0,/.test(line)).length, hasAudio ? 125 : 60)
+ assert.match(frames, /dimensions 0: 96x64/)
+})
diff --git a/utils/queuedJob.ts b/utils/queuedJob.ts
index 447ba8b..ba6d787 100644
--- a/utils/queuedJob.ts
+++ b/utils/queuedJob.ts
@@ -22,6 +22,7 @@ export type QueuedShotDraft = import('~/utils/imageIterations').ImageIteration &
}
export type QueuedInspectPayload = {
+ upscale?: { sourceId: string; scale: number; target: string; fps: string | number; enhance: string }
yueProfile?: 1 | 3
prompt?: string
promptMid?: string