137 lines
9.2 KiB
TypeScript
137 lines
9.2 KiB
TypeScript
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. */ }
|
|
}
|
|
}
|