Add local queued video upscale with Real-ESRGAN and RIFE
This commit is contained in:
@@ -32,6 +32,7 @@ export interface JobEvent {
|
||||
}
|
||||
|
||||
export interface Job {
|
||||
upscale?: boolean
|
||||
yueGp?: boolean
|
||||
musicActivity?: { checkedAt: number; running: boolean }
|
||||
id: string
|
||||
|
||||
+14
-3
@@ -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) }
|
||||
|
||||
@@ -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) => {
|
||||
|
||||
@@ -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. */ }
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user