Files
aigen/server/utils/pending.ts
T

325 lines
9.8 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'
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
workflow?: 'v1' | 'v2'
duration?: number
cfg?: number
fps?: number
samplerName?: string
scheduler?: string
hideInput?: boolean
extensions?: { prompt: string; duration: number }[]
remainingExtensions?: { prompt: string; duration: number }[]
currentClipId?: string
}
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,
...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
}
}
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 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
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 {
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
})
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
}