From 1466762af49bd06097242645480c921143edf4f0 Mon Sep 17 00:00:00 2001 From: Towsty Date: Thu, 27 Aug 2026 08:09:20 -0500 Subject: [PATCH] Reattach Comfy watchers after refresh so pending jobs finish and save. Co-authored-by: Cursor --- pages/index.vue | 24 ++++++++- server/api/generate/[id].get.ts | 10 +--- server/api/generate/[id]/stream.get.ts | 6 ++- server/api/generate/active.get.ts | 10 ++-- server/api/generate/recover.post.ts | 12 +---- server/plugins/resume-comfy.ts | 6 ++- server/utils/jobs.ts | 30 +++++++++++ server/utils/watch.ts | 69 ++++++++++++++++++++++++++ 8 files changed, 137 insertions(+), 30 deletions(-) diff --git a/pages/index.vue b/pages/index.vue index a2507ae..1863e2c 100644 --- a/pages/index.vue +++ b/pages/index.vue @@ -1792,7 +1792,7 @@ async function resumeTrack(kind: 'video' | 'edit', id: string, hidden: boolean, editStatusMessage.value = snap.message || 'Reconnecting to edit ComfyUI...' } else { videoBusy.value = true - statusMessage.value = snap.message || 'Reconnecting to ComfyUI...' + statusMessage.value = snap.message || 'Waiting for ComfyUI to finish this job...' } startTimer(kind, Number(snap.elapsedMs || 0)) persistActiveJob(kind, id, hidden, locked) @@ -2615,6 +2615,16 @@ function listen(kind: 'video' | 'edit', id: string, hidden: boolean, folderLocke if (gen !== (isEdit ? editListenGen : videoListenGen)) return applyJobEvent(kind, JSON.parse(event.data), hidden, folderLocked, id) } + source.onerror = () => { + if (gen !== (isEdit ? editListenGen : videoListenGen)) return + source.close() + setTimeout(() => { + if (gen !== (isEdit ? editListenGen : videoListenGen)) return + if (isEdit ? editSettledUi : videoSettledUi) return + listen(kind, id, hidden, folderLocked) + }, 2000) + } + let misses = 0 const tick = async () => { const currentGen = isEdit ? editListenGen : videoListenGen const settled = isEdit ? editSettledUi : videoSettledUi @@ -2624,6 +2634,7 @@ function listen(kind: 'video' | 'edit', id: string, hidden: boolean, folderLocke } try { const snap = await $fetch>(`/api/generate/${id}`) + misses = 0 if (gen !== (isEdit ? editListenGen : videoListenGen)) return applyJobEvent(kind, snap, hidden, folderLocked, id) } catch { @@ -2633,10 +2644,19 @@ function listen(kind: 'video' | 'edit', id: string, hidden: boolean, folderLocke method: 'POST', body: { jobId: id } }) + misses = 0 if (gen !== (isEdit ? editListenGen : videoListenGen)) return applyJobEvent(kind, recovered, hidden, folderLocked, id) } catch { - // Job snapshot can 404 briefly after a restart; keep polling while busy. + misses += 1 + if (misses >= 20) { + applyJobEvent(kind, { + type: 'error', + status: 'error', + error: 'Lost the running job after refresh. Check the library, or generate again.', + message: 'Lost the running job after refresh. Check the library, or generate again.' + }, hidden, folderLocked, id) + } } } } diff --git a/server/api/generate/[id].get.ts b/server/api/generate/[id].get.ts index cae3ec8..875c310 100644 --- a/server/api/generate/[id].get.ts +++ b/server/api/generate/[id].get.ts @@ -8,14 +8,8 @@ export default defineEventHandler(async (event) => { if (pending) { const done = await resolvePendingJob(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 - } + const job = ensurePendingWatch(pending) + return jobSnapshot(job) } } diff --git a/server/api/generate/[id]/stream.get.ts b/server/api/generate/[id]/stream.get.ts index 00ea8e2..a44eb53 100644 --- a/server/api/generate/[id]/stream.get.ts +++ b/server/api/generate/[id]/stream.get.ts @@ -1,6 +1,10 @@ export default defineEventHandler(async (event) => { const id = getRouterParam(event, 'id') - const job = id ? getJob(id) : undefined + let job = id ? getJob(id) : undefined + if (!job && id) { + const pending = readPendingJob(id) + if (pending) job = ensurePendingWatch(pending) + } if (!job) { throw createError({ statusCode: 404, statusMessage: 'Job not found' }) } diff --git a/server/api/generate/active.get.ts b/server/api/generate/active.get.ts index 35e06ae..aeb1038 100644 --- a/server/api/generate/active.get.ts +++ b/server/api/generate/active.get.ts @@ -35,17 +35,13 @@ export default defineEventHandler(async (event) => { ...resolved } } else { + const job = ensurePendingWatch(pending) video = { jobId: pending.jobId, kind: 'video', - type: 'snapshot', - status: 'running', - message: 'Reconnecting to ComfyUI...', - progress: Math.max(8, Math.min(90, Math.round((Date.now() - pending.startedAt) / 1000))), - promptId: pending.promptId, - elapsedMs: Date.now() - pending.startedAt, hideThumbnail: pending.hideThumbnail === true, - folderLocked: pending.folderLocked === true + folderLocked: pending.folderLocked === true, + ...jobSnapshot(job) } } } diff --git a/server/api/generate/recover.post.ts b/server/api/generate/recover.post.ts index 4c198a7..c3a6444 100644 --- a/server/api/generate/recover.post.ts +++ b/server/api/generate/recover.post.ts @@ -9,17 +9,7 @@ export default defineEventHandler(async (event) => { if (pending) { const done = await resolvePendingJob(pending) if (done) return done - return { - type: 'snapshot', - status: 'running', - jobId, - message: 'Reconnecting to ComfyUI...', - progress: 90, - promptId: pending.promptId, - elapsedMs: Date.now() - pending.startedAt, - hideThumbnail: pending.hideThumbnail, - folderLocked: pending.folderLocked - } + return jobSnapshot(ensurePendingWatch(pending)) } throw createError({ statusCode: 404, statusMessage: 'Job not found' }) } diff --git a/server/plugins/resume-comfy.ts b/server/plugins/resume-comfy.ts index aa8b52a..d4ded14 100644 --- a/server/plugins/resume-comfy.ts +++ b/server/plugins/resume-comfy.ts @@ -1,5 +1,9 @@ export default defineNitroPlugin(() => { for (const pending of listPendingJobs()) { - void resolvePendingJob(pending).catch(() => null) + void resolvePendingJob(pending).then((done) => { + if (!done) ensurePendingWatch(pending) + }).catch(() => { + ensurePendingWatch(pending) + }) } }) diff --git a/server/utils/jobs.ts b/server/utils/jobs.ts index 596ece4..00b3294 100644 --- a/server/utils/jobs.ts +++ b/server/utils/jobs.ts @@ -118,6 +118,36 @@ export function getJob(id: string) { return jobs.get(id) } +export function restoreJob(params: { + id: string + clientId: string + promptId: string + startedAt: number + hideThumbnail?: boolean + library: Job['library'] +}) { + const existing = jobs.get(params.id) + if (existing) return existing + const job: Job = { + id: params.id, + clientId: params.clientId, + kind: 'video', + promptId: params.promptId, + status: 'running', + message: 'Waiting for ComfyUI to finish this job...', + progress: 8, + step: 0, + maxStep: params.library?.steps || 0, + startedAt: params.startedAt, + hideThumbnail: params.hideThumbnail === true, + library: params.library, + events: [], + listeners: new Set() + } + jobs.set(job.id, job) + return job +} + export function listJobs() { return [...jobs.values()] } diff --git a/server/utils/watch.ts b/server/utils/watch.ts index 806bc85..cb327f1 100644 --- a/server/utils/watch.ts +++ b/server/utils/watch.ts @@ -1,5 +1,6 @@ import { existsSync } from 'node:fs' import type { Job, JobEvent } from '~/server/utils/jobs' +import { getJob, restoreJob } from '~/server/utils/jobs' function classifyError(message: string) { const lower = message.toLowerCase() @@ -375,3 +376,71 @@ export async function waitForComfySocket(job: Job, ms = 4000) { await sleep(100) } } + +const pendingWatches = new Set() + +export function ensurePendingWatch(pending: { + 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 +}) { + const existing = getJob(pending.jobId) + if (existing) return existing + const job = restoreJob({ + id: pending.jobId, + clientId: pending.clientId, + promptId: pending.promptId, + startedAt: pending.startedAt, + hideThumbnail: pending.hideThumbnail, + library: { + ownerKey: pending.ownerKey, + folderId: pending.folderId, + hideThumbnail: pending.hideThumbnail, + 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, + imageName: pending.imageName, + imageSubfolder: pending.imageSubfolder, + extendTmpDir: pending.extendTmpDir, + extendPart1Path: pending.extendPart1Path, + familyId: pending.familyId, + parentClipId: pending.parentClipId, + chainIndex: pending.chainIndex, + stillId: pending.stillId, + sound: pending.sound + } + }) + if (!pendingWatches.has(job.id)) { + pendingWatches.add(job.id) + void watchComfyJob(job, { persist: true }).finally(() => pendingWatches.delete(job.id)) + } + return job +}