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. */ } } }