diff --git a/pages/index.vue b/pages/index.vue index 2718329..f7e9733 100644 --- a/pages/index.vue +++ b/pages/index.vue @@ -705,7 +705,15 @@ function listen(id: string, hidden: boolean) { const snap = await $fetch>(`/api/generate/${id}`) applyJobEvent(snap, hidden) } catch { - // Job snapshot can 404 briefly after a restart; keep polling while busy. + try { + const recovered = await $fetch>('/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() diff --git a/server/api/generate.post.ts b/server/api/generate.post.ts index 6499e32..929a91d 100644 --- a/server/api/generate.post.ts +++ b/server/api/generate.post.ts @@ -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 } diff --git a/server/api/generate/[id].get.ts b/server/api/generate/[id].get.ts index 7cacf69..9dfe3c5 100644 --- a/server/api/generate/[id].get.ts +++ b/server/api/generate/[id].get.ts @@ -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' }) }) diff --git a/server/api/generate/recover.post.ts b/server/api/generate/recover.post.ts new file mode 100644 index 0000000..88afefc --- /dev/null +++ b/server/api/generate/recover.post.ts @@ -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 + } +}) diff --git a/server/api/library/index.get.ts b/server/api/library/index.get.ts index 580a0b3..b6a7467 100644 --- a/server/api/library/index.get.ts +++ b/server/api/library/index.get.ts @@ -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) }) diff --git a/server/plugins/resume-comfy.ts b/server/plugins/resume-comfy.ts new file mode 100644 index 0000000..fe4e287 --- /dev/null +++ b/server/plugins/resume-comfy.ts @@ -0,0 +1,5 @@ +export default defineNitroPlugin(() => { + for (const pending of listPendingJobs()) { + void completePendingIfReady(pending).catch(() => null) + } +}) diff --git a/server/utils/comfy.ts b/server/utils/comfy.ts index 9422cc8..de6d9fd 100644 --- a/server/utils/comfy.ts +++ b/server/utils/comfy.ts @@ -68,6 +68,12 @@ export async function fetchHistory(promptId: string) { return (await res.json()) as Record } +export async function fetchHistoryAll() { + const res = await comfyFetch('/history') + if (!res.ok) return {} + return (await res.json()) as Record +} + 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 | null, promptId: string) { if (!history) return null - const entry = (history[promptId] || Object.values(history)[0]) as { outputs?: Record } | undefined - return findVideo(entry?.outputs || {}) + const wrapped = history[promptId] as { outputs?: Record } | 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 }>)) { + if (node?.class_type === 'MiniMaxH3ImageToVideo') { + return String(node.inputs?.prompt || node.inputs?.text || '').trim() + } + } + return '' } export function inspectHistory(history: Record | 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() { diff --git a/server/utils/library.ts b/server/utils/library.ts index 2173433..ee71770 100644 --- a/server/utils/library.ts +++ b/server/utils/library.ts @@ -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 params = new URLSearchParams({ - filename: video.filename, - subfolder: video.subfolder, - type: video.type - }) - 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()) + const subfolders = [...new Set([video.subfolder, 'video', ''])] + let lastStatus = 0 + for (const subfolder of subfolders) { + const params = new URLSearchParams({ + filename: video.filename, + subfolder, + type: video.type || 'output' + }) + const res = await comfyFetch(`/view?${params.toString()}`) + 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 }, 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 } diff --git a/server/utils/pending.ts b/server/utils/pending.ts new file mode 100644 index 0000000..ea1ea01 --- /dev/null +++ b/server/utils/pending.ts @@ -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 + } +} diff --git a/server/utils/watch.ts b/server/utils/watch.ts index b8f8681..90a57c3 100644 --- a/server/utils/watch.ts +++ b/server/utils/watch.ts @@ -34,55 +34,47 @@ export function watchComfyJob(job: Job): Promise { try { ws.close() } catch { /* ignore */ } } - const finish = async (error?: string) => { + const fail = async (error: string) => { if (settled) return settled = true cleanup() + job.status = job.status === 'cancelled' ? 'cancelled' : 'error' + job.error = classifyError(error) + emitJob(job, { type: 'error', error: job.error, message: job.error }) + if (job.id) deletePendingJob(job.id) + resolve() + } + + 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 { - if (error) { - job.status = job.status === 'cancelled' ? 'cancelled' : 'error' - job.error = classifyError(error) - emitJob(job, { type: 'error', error: job.error, message: job.error }) - return - } - 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 + emitJob(job, { type: 'status', message: 'Saving to library...', progress: 98 }) if (job.library) { - emitJob(job, { type: 'status', message: 'Saving to library...', progress: 98 }) - try { - const buffer = await downloadComfyVideo(video) - const clip = await saveClip({ - ownerKey: job.library.ownerKey, - folderId: job.library.folderId, - 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, - hideThumbnail: job.library.hideThumbnail, - video: buffer, - thumb: job.library.hideThumbnail ? null : job.library.thumb - }) - 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 - } + const buffer = await downloadComfyVideo(video) + const clip = await saveClip({ + ownerKey: job.library.ownerKey, + folderId: job.library.folderId, + 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, + hideThumbnail: job.library.hideThumbnail, + video: buffer, + thumb: job.library.hideThumbnail ? null : job.library.thumb, + comfyFilename: video.filename + }) + job.clipId = clip.id + job.hideThumbnail = clip.hideThumbnail + job.library.thumb = undefined } job.status = 'complete' emitJob(job, { @@ -95,9 +87,15 @@ export function watchComfyJob(job: Job): Promise { 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 { 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 { } 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 { 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 { 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.') } }) })