Recover finished Comfy videos after relay restarts.

Persist queued jobs to disk, import leftover MP4s from Comfy history, and do not mark a run complete until the saved video exists.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Towsty
2026-08-25 18:43:15 -05:00
co-authored by Cursor
parent f0f80b5873
commit b7196ecde3
10 changed files with 342 additions and 77 deletions
+8
View File
@@ -704,10 +704,18 @@ function listen(id: string, hidden: boolean) {
try {
const snap = await $fetch<Record<string, any>>(`/api/generate/${id}`)
applyJobEvent(snap, hidden)
} catch {
try {
const recovered = await $fetch<Record<string, any>>('/api/generate/recover', {
method: 'POST',
body: { jobId: id }
})
applyJobEvent(recovered, hidden)
} catch {
// Job snapshot can 404 briefly after a restart; keep polling while busy.
}
}
}
void tick()
pollTimer = setInterval(() => {
void tick()
+18
View File
@@ -125,6 +125,24 @@ async function runGeneration(
const queued = await queuePrompt(graph, job.clientId)
job.promptId = queued.prompt_id
job.status = 'running'
if (job.library && job.promptId) {
writePendingJob({
jobId: job.id,
promptId: job.promptId,
clientId: job.clientId,
ownerKey: job.library.ownerKey,
folderId: job.library.folderId,
hideThumbnail: job.library.hideThumbnail,
prompt: job.library.prompt,
aspect: job.library.aspect,
width: job.library.width,
height: job.library.height,
steps: job.library.steps,
turbo: job.library.turbo,
seed: job.library.seed,
startedAt: job.startedAt
})
}
emitJob(job, { type: 'status', message: 'Job queued on ComfyUI', progress: 8 })
await done
}
+20 -5
View File
@@ -1,8 +1,23 @@
export default defineEventHandler((event) => {
export default defineEventHandler(async (event) => {
const id = getRouterParam(event, 'id')
const job = id ? getJob(id) : undefined
if (!job) {
throw createError({ statusCode: 404, statusMessage: 'Job not found' })
const live = id ? getJob(id) : undefined
if (live) return jobSnapshot(live)
if (id) {
const pending = readPendingJob(id)
if (pending) {
const done = await completePendingIfReady(pending).catch(() => null)
if (done) return done
return {
type: 'snapshot',
status: 'running',
message: 'Reconnecting to ComfyUI...',
progress: 90,
promptId: pending.promptId,
elapsedMs: Date.now() - pending.startedAt
}
return jobSnapshot(job)
}
}
throw createError({ statusCode: 404, statusMessage: 'Job not found' })
})
+40
View File
@@ -0,0 +1,40 @@
export default defineEventHandler(async (event) => {
const body = await readBody<{ jobId?: string }>(event).catch(() => ({} as { jobId?: string }))
const jobId = String(body?.jobId || '')
const live = jobId ? getJob(jobId) : undefined
if (live) return jobSnapshot(live)
if (jobId) {
const pending = readPendingJob(jobId)
if (pending) {
const done = await completePendingIfReady(pending)
if (done) return done
return {
type: 'snapshot',
status: 'running',
message: 'Reconnecting to ComfyUI...',
progress: 90,
promptId: pending.promptId,
elapsedMs: Date.now() - pending.startedAt
}
}
}
const { owner } = assertLibraryOwner(event)
const library = publicLibrary(event)
const folderId = library.folders.find(folder => folder.unlocked)?.id
const imported = folderId ? await importMissingComfyVideos(owner, folderId) : []
const newest = imported[imported.length - 1]
if (!newest) {
throw createError({ statusCode: 404, statusMessage: 'No finished ComfyUI video to recover' })
}
return {
type: 'complete',
status: 'complete',
message: 'Recovered video from ComfyUI',
progress: 100,
clipId: newest.id,
filename: newest.comfyFilename,
hideThumbnail: newest.hideThumbnail
}
})
+9 -1
View File
@@ -1,3 +1,11 @@
export default defineEventHandler((event) => {
export default defineEventHandler(async (event) => {
const library = publicLibrary(event)
try {
const { owner } = assertLibraryOwner(event)
const folderId = library.folders.find(folder => folder.unlocked)?.id
if (folderId) await importMissingComfyVideos(owner, folderId)
} catch {
return library
}
return publicLibrary(event)
})
+5
View File
@@ -0,0 +1,5 @@
export default defineNitroPlugin(() => {
for (const pending of listPendingJobs()) {
void completePendingIfReady(pending).catch(() => null)
}
})
+31 -9
View File
@@ -68,6 +68,12 @@ export async function fetchHistory(promptId: string) {
return (await res.json()) as Record<string, unknown>
}
export async function fetchHistoryAll() {
const res = await comfyFetch('/history')
if (!res.ok) return {}
return (await res.json()) as Record<string, unknown>
}
function isVideoFile(item: { filename?: string; format?: string } | null | undefined) {
if (!item) return false
const name = String(item.filename || '').toLowerCase()
@@ -107,22 +113,38 @@ function findVideo(value: unknown, depth = 0): { filename: string; subfolder: st
export function extractVideo(history: Record<string, unknown> | null, promptId: string) {
if (!history) return null
const entry = (history[promptId] || Object.values(history)[0]) as { outputs?: Record<string, unknown> } | undefined
return findVideo(entry?.outputs || {})
const wrapped = history[promptId] as { outputs?: Record<string, unknown> } | undefined
if (wrapped) {
return findVideo(wrapped.outputs || {}) || findVideo(wrapped)
}
if ((history as { outputs?: unknown }).outputs) {
return findVideo((history as { outputs?: unknown }).outputs) || findVideo(history)
}
return findVideo(history)
}
export function extractPromptFromHistory(entry: unknown) {
const prompt = (entry as { prompt?: unknown[] })?.prompt
const graph = Array.isArray(prompt) ? prompt[2] : null
if (!graph || typeof graph !== 'object') return ''
for (const node of Object.values(graph as Record<string, { class_type?: string; inputs?: Record<string, unknown> }>)) {
if (node?.class_type === 'MiniMaxH3ImageToVideo') {
return String(node.inputs?.prompt || node.inputs?.text || '').trim()
}
}
return ''
}
export function inspectHistory(history: Record<string, unknown> | null, promptId: string) {
const video = extractVideo(history, promptId)
if (!history) return { video: null, completed: false, error: null as string | null }
const entry = (history[promptId] || Object.values(history)[0]) as {
status?: { status_str?: string; completed?: boolean; messages?: unknown[] }
} | undefined
const status = entry?.status
const error = status?.status_str === 'error'
const entry = (history[promptId] || history) as {
status?: { status_str?: string; completed?: boolean }
}
const error = entry?.status?.status_str === 'error'
? 'ComfyUI reported an execution error'
: null
const completed = Boolean(video) || status?.completed === true || status?.status_str === 'success'
return { video, completed, error }
return { video, completed: Boolean(video), error }
}
export async function probeComfy() {
+58 -5
View File
@@ -27,6 +27,7 @@ export interface LibraryClip {
seed: number
hideThumbnail: boolean
createdAt: number
comfyFilename?: string
}
export interface LibraryStill {
@@ -460,10 +461,14 @@ export async function saveClip(params: {
hideThumbnail: boolean
video: Buffer
thumb?: Buffer | null
comfyFilename?: string
}) {
const catalog = readCatalog(params.ownerKey)
const folder = catalog.folders.find(item => item.id === params.folderId) || catalog.folders[0]
if (!folder) throw createError({ statusCode: 400, statusMessage: 'No library folder available' })
if (params.comfyFilename && catalog.clips.some(clip => clip.comfyFilename === params.comfyFilename)) {
return catalog.clips.find(clip => clip.comfyFilename === params.comfyFilename) as LibraryClip
}
const clip: LibraryClip = {
id: crypto.randomUUID(),
folderId: folder.id,
@@ -475,7 +480,8 @@ export async function saveClip(params: {
turbo: params.turbo,
seed: params.seed,
hideThumbnail: params.hideThumbnail,
createdAt: Date.now()
createdAt: Date.now(),
comfyFilename: params.comfyFilename
}
mkdirSync(clipDir(params.ownerKey, clip.id), { recursive: true })
writeFileSync(clipVideoPath(params.ownerKey, clip.id), params.video)
@@ -515,12 +521,59 @@ export function deleteClip(owner: string, id: string) {
}
export async function downloadComfyVideo(video: { filename: string; subfolder: string; type: string }) {
const subfolders = [...new Set([video.subfolder, 'video', ''])]
let lastStatus = 0
for (const subfolder of subfolders) {
const params = new URLSearchParams({
filename: video.filename,
subfolder: video.subfolder,
type: video.type
subfolder,
type: video.type || 'output'
})
const res = await comfyFetch(`/view?${params.toString()}`)
if (!res.ok) throw new Error('Failed to fetch completed video from ComfyUI')
return Buffer.from(await res.arrayBuffer())
lastStatus = res.status
if (res.ok) return Buffer.from(await res.arrayBuffer())
}
throw new Error(`Failed to fetch completed video from ComfyUI (${lastStatus})`)
}
export async function importMissingComfyVideos(owner: string, folderId?: string) {
const catalog = readCatalog(owner)
const folder = catalog.folders.find(item => item.id === folderId) || catalog.folders[0]
if (!folder) return []
const known = new Set(catalog.clips.map(clip => clip.comfyFilename).filter(Boolean) as string[])
const unlabeled = catalog.clips.filter(clip => !clip.comfyFilename).length
const history = await fetchHistoryAll()
const found: { promptId: string; video: { filename: string; subfolder: string; type: string }; prompt: string }[] = []
for (const [promptId, entry] of Object.entries(history)) {
const video = extractVideo({ [promptId]: entry as Record<string, unknown> }, promptId)
if (!video || known.has(video.filename)) continue
found.push({
promptId,
video,
prompt: extractPromptFromHistory(entry) || 'Recovered from ComfyUI'
})
}
found.sort((a, b) => a.video.filename.localeCompare(b.video.filename, undefined, { numeric: true }))
const toImport = found.slice(unlabeled)
const imported = []
for (const item of toImport) {
const buffer = await downloadComfyVideo(item.video)
const clip = await saveClip({
ownerKey: owner,
folderId: folder.id,
prompt: item.prompt,
aspect: 'auto',
width: 0,
height: 0,
steps: 0,
turbo: true,
seed: 0,
hideThumbnail: false,
video: buffer,
thumb: null,
comfyFilename: item.video.filename
})
imported.push(clip)
}
return imported
}
+95
View File
@@ -0,0 +1,95 @@
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
prompt: string
aspect: string
width: number
height: number
steps: number
turbo: boolean
seed: number
startedAt: number
}
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 history = await fetchHistory(pending.promptId)
const video = extractVideo(history, pending.promptId)
if (!video) return null
const buffer = await downloadComfyVideo(video)
const clip = await saveClip({
ownerKey: pending.ownerKey,
folderId: pending.folderId,
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
})
deletePendingJob(pending.jobId)
return {
type: 'complete' as const,
status: 'complete' as const,
message: 'Video ready',
progress: 100,
clipId: clip.id,
filename: video.filename,
subfolder: video.subfolder,
mediaType: video.type,
hideThumbnail: pending.hideThumbnail
}
}
+34 -33
View File
@@ -34,30 +34,28 @@ export function watchComfyJob(job: Job): Promise<void> {
try { ws.close() } catch { /* ignore */ }
}
const finish = async (error?: string) => {
const fail = async (error: string) => {
if (settled) return
settled = true
cleanup()
try {
if (error) {
job.status = job.status === 'cancelled' ? 'cancelled' : 'error'
job.error = classifyError(error)
emitJob(job, { type: 'error', error: job.error, message: job.error })
return
if (job.id) deletePendingJob(job.id)
resolve()
}
emitJob(job, { type: 'status', message: 'Fetching output...', progress: 96 })
const history = await fetchHistory(job.promptId || '')
const video = extractVideo(history, job.promptId || '')
if (!video) {
job.status = 'error'
job.error = 'Workflow finished but no MP4 was found in ComfyUI history.'
emitJob(job, { type: 'error', error: job.error, message: job.error })
return
}
job.video = video
if (job.library) {
emitJob(job, { type: 'status', message: 'Saving to library...', progress: 98 })
const succeed = async () => {
if (settled || !job.promptId) return false
const history = await fetchHistory(job.promptId)
const video = extractVideo(history, job.promptId)
if (!video) return false
settled = true
cleanup()
try {
job.video = video
emitJob(job, { type: 'status', message: 'Saving to library...', progress: 98 })
if (job.library) {
const buffer = await downloadComfyVideo(video)
const clip = await saveClip({
ownerKey: job.library.ownerKey,
@@ -71,18 +69,12 @@ export function watchComfyJob(job: Job): Promise<void> {
seed: job.library.seed,
hideThumbnail: job.library.hideThumbnail,
video: buffer,
thumb: job.library.hideThumbnail ? null : job.library.thumb
thumb: job.library.hideThumbnail ? null : job.library.thumb,
comfyFilename: video.filename
})
job.clipId = clip.id
job.hideThumbnail = clip.hideThumbnail
job.library.thumb = undefined
} catch (saveError) {
const message = saveError instanceof Error ? saveError.message : String(saveError)
job.status = 'error'
job.error = `Video generated but library save failed: ${message}`
emitJob(job, { type: 'error', error: job.error, message: job.error })
return
}
}
job.status = 'complete'
emitJob(job, {
@@ -95,9 +87,15 @@ export function watchComfyJob(job: Job): Promise<void> {
clipId: job.clipId,
hideThumbnail: job.hideThumbnail
})
} finally {
resolve()
deletePendingJob(job.id)
} catch (saveError) {
const message = saveError instanceof Error ? saveError.message : String(saveError)
job.status = 'error'
job.error = `Video generated but library save failed: ${message}`
emitJob(job, { type: 'error', error: job.error, message: job.error })
}
resolve()
return true
}
const pollHistory = async () => {
@@ -105,11 +103,11 @@ export function watchComfyJob(job: Job): Promise<void> {
try {
const inspected = inspectHistory(await fetchHistory(job.promptId), job.promptId)
if (inspected.error) {
await finish(inspected.error)
await fail(inspected.error)
return
}
if (inspected.video || inspected.completed) {
await finish()
if (inspected.video) {
await succeed()
}
} catch {
// Transient ComfyUI history misses are expected while the graph is still running.
@@ -117,7 +115,7 @@ export function watchComfyJob(job: Job): Promise<void> {
}
timeout = setTimeout(() => {
void finish('Timed out waiting for ComfyUI (15 minutes).')
void fail('Timed out waiting for ComfyUI (15 minutes).')
}, 15 * 60 * 1000)
pollTimer = setInterval(() => {
@@ -168,7 +166,10 @@ export function watchComfyJob(job: Job): Promise<void> {
if (payload.type === 'executing') {
const node = data.node === null || data.node === undefined ? null : String(data.node)
if (node === null) {
await finish()
const saved = await succeed()
if (!saved) {
emitJob(job, { type: 'status', message: 'Waiting for ComfyUI to save the MP4...', progress: 96 })
}
return
}
const label = NODE_LABELS[node] || `Running node ${node}`
@@ -183,12 +184,12 @@ export function watchComfyJob(job: Job): Promise<void> {
if (payload.type === 'execution_error') {
const message = String(data.exception_message || data.message || 'ComfyUI node execution failed')
await finish(message)
await fail(message)
}
if (payload.type === 'execution_interrupted') {
job.status = 'cancelled'
await finish('Job interrupted.')
await fail('Job interrupted.')
}
})
})