Files
aigen/server/utils/pending.ts
T
TowstyandCursor e334cb7a7a Move library collections as a group and map v2 identity refs.
Collapsed collection actions now move, rename, delete, and rerun every part so chains stay together. Identity mode follows the v2u3 Ref2VA picture wiring.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-26 22:56:37 -05:00

187 lines
5.6 KiB
TypeScript

import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, rmSync, writeFileSync } from 'node:fs'
import { join } from 'node:path'
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
stillId?: string
sound?: boolean
}
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 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) {
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, 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
}