379 lines
12 KiB
TypeScript
379 lines
12 KiB
TypeScript
import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, rmSync, writeFileSync } from 'node:fs'
|
|
import { join } from 'node:path'
|
|
import type { Job } from '~/server/utils/jobs'
|
|
import { extractFirstFrameFromBuffer } from '~/server/utils/ffmpeg'
|
|
import { mergePermanenceRefs } from '~/utils/globalLocks'
|
|
|
|
export interface PendingJob {
|
|
jobId: string
|
|
promptId: string
|
|
clientId: string
|
|
ownerKey: string
|
|
folderId: string
|
|
hideThumbnail: boolean
|
|
folderLocked?: boolean
|
|
name?: string
|
|
prompt: string
|
|
aspect: string
|
|
width: number
|
|
height: number
|
|
steps: number
|
|
turbo: boolean
|
|
seed: number
|
|
startedAt: number
|
|
imageName?: string
|
|
imageSubfolder?: string
|
|
extendTmpDir?: string
|
|
extendPart1Path?: string
|
|
familyId?: string
|
|
parentClipId?: string
|
|
chainIndex?: number
|
|
chainStep?: number
|
|
chainTotal?: number
|
|
chainLabel?: string
|
|
stillId?: string
|
|
stillFilename?: string
|
|
sound?: boolean
|
|
referenceImageNames?: string[]
|
|
referenceStillIds?: Array<string | null>
|
|
useIdentityRefs?: boolean
|
|
globalLocks?: string
|
|
permanenceRefs?: import('~/utils/globalLocks').PermanenceRef[]
|
|
shotPermanenceRefs?: import('~/utils/globalLocks').PermanenceRef[][]
|
|
workflow?: import('~/utils/videoModels').VideoWorkflowId
|
|
duration?: number
|
|
cfg?: number
|
|
fps?: number
|
|
samplerName?: string
|
|
scheduler?: string
|
|
hideInput?: boolean
|
|
extensions?: { prompt: string; duration: number; permanenceRefs?: import('~/utils/globalLocks').PermanenceRef[]; loraName?: string; loraStack?: import('~/utils/loras').LoraStackItem[] }[]
|
|
remainingExtensions?: { prompt: string; duration: number; permanenceRefs?: import('~/utils/globalLocks').PermanenceRef[]; loraName?: string; loraStack?: import('~/utils/loras').LoraStackItem[] }[]
|
|
currentClipId?: string
|
|
queueId?: string
|
|
queueAutoRun?: boolean
|
|
queueBudget?: number
|
|
stopAfterCurrent?: boolean
|
|
loraName?: string
|
|
loraStack?: import('~/utils/loras').LoraStackItem[]
|
|
shotLoras?: string[]
|
|
shotLoraStacks?: import('~/utils/loras').LoraStackItem[][]
|
|
}
|
|
|
|
function pendingRoot() {
|
|
const config = useRuntimeConfig()
|
|
return join((config.libraryDir || process.env.LIBRARY_DIR || '/data/library').replace(/\/$/, ''), 'pending')
|
|
}
|
|
|
|
function pendingPath(jobId: string) {
|
|
return join(pendingRoot(), `${jobId}.json`)
|
|
}
|
|
|
|
function ensurePending() {
|
|
mkdirSync(pendingRoot(), { recursive: true })
|
|
}
|
|
|
|
export function remainingAfterCurrentShot(job: {
|
|
chainIndex?: number
|
|
extensions?: { prompt: string; duration: number }[]
|
|
}) {
|
|
const extensions = job.extensions || []
|
|
const index = job.chainIndex || 0
|
|
return index === 0 ? extensions : extensions.slice(index)
|
|
}
|
|
|
|
export function isLastChainShot(job: {
|
|
chainIndex?: number
|
|
extensions?: { prompt: string; duration: number }[]
|
|
remainingExtensions?: { prompt: string; duration: number }[]
|
|
}) {
|
|
if (job.remainingExtensions) return job.remainingExtensions.length === 0
|
|
return remainingAfterCurrentShot(job).length === 0
|
|
}
|
|
|
|
export function pendingFromJob(job: Job, overrides: Partial<PendingJob> = {}): PendingJob {
|
|
const library = job.library
|
|
if (!library) {
|
|
throw new Error('Cannot persist a job without library metadata')
|
|
}
|
|
const remaining = remainingAfterCurrentShot(library)
|
|
return {
|
|
jobId: job.id,
|
|
promptId: job.promptId || '',
|
|
clientId: job.clientId,
|
|
ownerKey: library.ownerKey,
|
|
folderId: library.folderId,
|
|
hideThumbnail: library.hideThumbnail,
|
|
folderLocked: library.folderLocked,
|
|
name: library.name,
|
|
prompt: library.prompt,
|
|
aspect: library.aspect,
|
|
width: library.width,
|
|
height: library.height,
|
|
steps: library.steps,
|
|
turbo: library.turbo,
|
|
seed: library.seed,
|
|
startedAt: job.startedAt,
|
|
imageName: library.imageName,
|
|
imageSubfolder: library.imageSubfolder,
|
|
extendTmpDir: library.extendTmpDir,
|
|
extendPart1Path: library.extendPart1Path,
|
|
familyId: library.familyId,
|
|
parentClipId: library.parentClipId,
|
|
chainIndex: library.chainIndex,
|
|
chainStep: library.chainStep,
|
|
chainTotal: library.chainTotal,
|
|
chainLabel: library.chainLabel,
|
|
stillId: library.stillId,
|
|
stillFilename: library.stillFilename,
|
|
sound: library.sound,
|
|
referenceImageNames: library.referenceImageNames,
|
|
referenceStillIds: library.referenceStillIds,
|
|
useIdentityRefs: library.useIdentityRefs,
|
|
workflow: library.workflow,
|
|
duration: library.duration,
|
|
cfg: library.cfg,
|
|
fps: library.fps,
|
|
samplerName: library.samplerName,
|
|
scheduler: library.scheduler,
|
|
hideInput: library.hideInput,
|
|
extensions: library.extensions,
|
|
remainingExtensions: remaining,
|
|
currentClipId: job.clipId,
|
|
globalLocks: library.globalLocks,
|
|
permanenceRefs: library.permanenceRefs,
|
|
shotPermanenceRefs: library.shotPermanenceRefs,
|
|
loraName: library.loraName,
|
|
loraStack: library.loraStack,
|
|
shotLoras: library.shotLoras,
|
|
shotLoraStacks: library.shotLoraStacks,
|
|
queueId: library.queueId,
|
|
queueAutoRun: library.queueAutoRun,
|
|
queueBudget: library.queueBudget,
|
|
stopAfterCurrent: library.stopAfterCurrent,
|
|
...overrides
|
|
}
|
|
}
|
|
|
|
export function libraryFromPending(pending: PendingJob): NonNullable<Job['library']> {
|
|
return {
|
|
ownerKey: pending.ownerKey,
|
|
folderId: pending.folderId,
|
|
hideThumbnail: pending.hideThumbnail,
|
|
hideInput: pending.hideInput,
|
|
folderLocked: pending.folderLocked,
|
|
name: pending.name,
|
|
prompt: pending.prompt,
|
|
aspect: pending.aspect,
|
|
width: pending.width,
|
|
height: pending.height,
|
|
steps: pending.steps,
|
|
turbo: pending.turbo,
|
|
seed: pending.seed,
|
|
cfg: pending.cfg,
|
|
fps: pending.fps,
|
|
samplerName: pending.samplerName,
|
|
scheduler: pending.scheduler,
|
|
imageName: pending.imageName,
|
|
imageSubfolder: pending.imageSubfolder,
|
|
extendTmpDir: pending.extendTmpDir,
|
|
extendPart1Path: pending.extendPart1Path,
|
|
familyId: pending.familyId,
|
|
parentClipId: pending.parentClipId,
|
|
chainIndex: pending.chainIndex,
|
|
chainStep: pending.chainStep,
|
|
chainTotal: pending.chainTotal,
|
|
chainLabel: pending.chainLabel,
|
|
stillId: pending.stillId,
|
|
stillFilename: pending.stillFilename,
|
|
sound: pending.sound,
|
|
referenceImageNames: pending.referenceImageNames,
|
|
referenceStillIds: pending.referenceStillIds,
|
|
useIdentityRefs: pending.useIdentityRefs,
|
|
workflow: pending.workflow,
|
|
duration: pending.duration,
|
|
extensions: pending.extensions || pending.remainingExtensions,
|
|
queueId: pending.queueId,
|
|
queueAutoRun: pending.queueAutoRun,
|
|
queueBudget: pending.queueBudget,
|
|
stopAfterCurrent: pending.stopAfterCurrent,
|
|
globalLocks: pending.globalLocks,
|
|
permanenceRefs: pending.permanenceRefs,
|
|
shotPermanenceRefs: pending.shotPermanenceRefs,
|
|
loraName: pending.loraName,
|
|
loraStack: pending.loraStack,
|
|
shotLoras: pending.shotLoras,
|
|
shotLoraStacks: pending.shotLoraStacks
|
|
}
|
|
}
|
|
|
|
export function writePendingJob(job: PendingJob) {
|
|
ensurePending()
|
|
const tmp = pendingPath(job.jobId) + '.tmp'
|
|
writeFileSync(tmp, JSON.stringify(job, null, 2))
|
|
renameSync(tmp, pendingPath(job.jobId))
|
|
}
|
|
|
|
export function patchPendingJob(jobId: string, patch: Partial<PendingJob>) {
|
|
const existing = readPendingJob(jobId)
|
|
if (!existing) return null
|
|
const next = { ...existing, ...patch }
|
|
writePendingJob(next)
|
|
return next
|
|
}
|
|
|
|
export function readPendingJob(jobId: string): PendingJob | null {
|
|
const path = pendingPath(jobId)
|
|
if (!existsSync(path)) return null
|
|
try {
|
|
return JSON.parse(readFileSync(path, 'utf8')) as PendingJob
|
|
} catch {
|
|
return null
|
|
}
|
|
}
|
|
|
|
export function listPendingJobs(): PendingJob[] {
|
|
ensurePending()
|
|
return readdirSync(pendingRoot())
|
|
.filter(name => name.endsWith('.json'))
|
|
.map(name => readPendingJob(name.replace(/\.json$/, '')))
|
|
.filter((job): job is PendingJob => Boolean(job))
|
|
}
|
|
|
|
export function deletePendingJob(jobId: string) {
|
|
rmSync(pendingPath(jobId), { force: true })
|
|
}
|
|
|
|
export async function completePendingIfReady(pending: PendingJob) {
|
|
if (!isLastChainShot(pending)) return null
|
|
if (!pending.promptId) return null
|
|
const alreadySaved = () => {
|
|
const live = getJob(pending.jobId)
|
|
if (live?.savedPromptId === pending.promptId) {
|
|
deletePendingJob(pending.jobId)
|
|
return {
|
|
type: 'complete' as const,
|
|
status: 'complete' as const,
|
|
message: pending.folderLocked
|
|
? 'Saved to the locked folder. Unlock it to view.'
|
|
: (pending.extendPart1Path ? 'Extended video ready' : 'Video ready'),
|
|
progress: 100,
|
|
jobId: pending.jobId,
|
|
clipId: live.clipId,
|
|
hideThumbnail: pending.hideThumbnail,
|
|
folderLocked: pending.folderLocked
|
|
}
|
|
}
|
|
return null
|
|
}
|
|
const inFlight = () => {
|
|
const current = getJob(pending.jobId)
|
|
return Boolean(current && current.promptId === pending.promptId && current.saving)
|
|
}
|
|
if (inFlight()) return null
|
|
const saved = alreadySaved()
|
|
if (saved) return saved
|
|
const history = await fetchHistory(pending.promptId)
|
|
if (inFlight()) return null
|
|
const afterHistory = alreadySaved()
|
|
if (afterHistory) return afterHistory
|
|
const video = extractVideo(history, pending.promptId)
|
|
if (!video) return null
|
|
let buffer = await downloadComfyVideo(video)
|
|
if (inFlight()) return null
|
|
const afterDownload = alreadySaved()
|
|
if (afterDownload) return afterDownload
|
|
let segmentFirstFrame: Buffer | undefined
|
|
if (pending.extendPart1Path && pending.extendTmpDir) {
|
|
if (!existsSync(pending.extendPart1Path)) {
|
|
removeExtendTemp(pending.extendTmpDir)
|
|
throw createError({ statusCode: 500, statusMessage: 'Extension source clip was missing during stitch' })
|
|
}
|
|
try {
|
|
segmentFirstFrame = await extractFirstFrameFromBuffer(buffer, pending.extendTmpDir)
|
|
} catch {
|
|
// ensureClipFirstFrame will seek to the join if this frame is missing
|
|
}
|
|
try {
|
|
buffer = await stitchExtension({
|
|
part1Path: pending.extendPart1Path,
|
|
part2: buffer,
|
|
tmpDir: pending.extendTmpDir
|
|
})
|
|
} catch (error) {
|
|
removeExtendTemp(pending.extendTmpDir)
|
|
throw error
|
|
}
|
|
removeExtendTemp(pending.extendTmpDir)
|
|
}
|
|
const afterStitch = alreadySaved()
|
|
if (afterStitch) return afterStitch
|
|
const clip = await saveClip({
|
|
ownerKey: pending.ownerKey,
|
|
folderId: pending.folderId,
|
|
name: pending.name,
|
|
prompt: pending.prompt,
|
|
aspect: pending.aspect,
|
|
width: pending.width,
|
|
height: pending.height,
|
|
steps: pending.steps,
|
|
turbo: pending.turbo,
|
|
seed: pending.seed,
|
|
hideThumbnail: pending.hideThumbnail,
|
|
video: buffer,
|
|
thumb: null,
|
|
comfyFilename: video.filename,
|
|
familyId: pending.familyId,
|
|
parentClipId: pending.parentClipId,
|
|
chainIndex: pending.chainIndex,
|
|
stillId: pending.stillId,
|
|
sound: pending.sound,
|
|
globalLocks: pending.globalLocks,
|
|
permanenceRefs: mergePermanenceRefs(pending.permanenceRefs, pending.shotPermanenceRefs?.[pending.chainIndex || 0]),
|
|
loraName: pending.loraName,
|
|
loraStack: pending.loraStack,
|
|
segmentFirstFrame
|
|
})
|
|
await purgeComfyArtifacts({
|
|
video,
|
|
imageName: pending.imageName,
|
|
imageSubfolder: pending.imageSubfolder,
|
|
extraImageNames: pending.referenceImageNames,
|
|
promptId: pending.promptId
|
|
})
|
|
deletePendingJob(pending.jobId)
|
|
return {
|
|
type: 'complete' as const,
|
|
status: 'complete' as const,
|
|
message: pending.folderLocked
|
|
? 'Saved to the locked folder. Unlock it to view.'
|
|
: (pending.extendPart1Path ? 'Extended video ready' : 'Video ready'),
|
|
progress: 100,
|
|
jobId: pending.jobId,
|
|
clipId: clip.id,
|
|
filename: pending.folderLocked ? undefined : video.filename,
|
|
subfolder: pending.folderLocked ? undefined : video.subfolder,
|
|
mediaType: pending.folderLocked ? undefined : video.type,
|
|
hideThumbnail: pending.hideThumbnail,
|
|
folderLocked: pending.folderLocked
|
|
}
|
|
}
|
|
|
|
export async function resolvePendingJob(pending: PendingJob) {
|
|
const done = await completePendingIfReady(pending)
|
|
if (done) return done
|
|
if (await isComfyPromptDropped(pending.promptId)) {
|
|
deletePendingJob(pending.jobId)
|
|
const message = 'ComfyUI dropped this job. Reloading the Comfy interface clears the queue. Generate again.'
|
|
return {
|
|
type: 'error' as const,
|
|
status: 'error' as const,
|
|
jobId: pending.jobId,
|
|
error: message,
|
|
message,
|
|
progress: 0
|
|
}
|
|
}
|
|
return null
|
|
}
|