Retry library pulls after Comfy finishes and fall back to disk.
Long jobs were succeeding on Beast then failing the save step when /view briefly dropped; retries plus host-agent file serving keep finished MP4s from being orphaned. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -2,7 +2,7 @@ import http from 'node:http'
|
|||||||
import net from 'node:net'
|
import net from 'node:net'
|
||||||
import { execFile, spawn } from 'node:child_process'
|
import { execFile, spawn } from 'node:child_process'
|
||||||
import { promisify } from 'node:util'
|
import { promisify } from 'node:util'
|
||||||
import { readdirSync, existsSync, rmSync, readFileSync, openSync, mkdirSync, writeFileSync, unlinkSync, statSync } from 'node:fs'
|
import { readdirSync, existsSync, rmSync, readFileSync, openSync, mkdirSync, writeFileSync, unlinkSync, statSync, createReadStream } from 'node:fs'
|
||||||
import { basename, dirname, join, resolve, relative, isAbsolute } from 'node:path'
|
import { basename, dirname, join, resolve, relative, isAbsolute } from 'node:path'
|
||||||
import { tmpdir } from 'node:os'
|
import { tmpdir } from 'node:os'
|
||||||
|
|
||||||
@@ -594,6 +594,42 @@ function safeFile(root, subfolder, filename) {
|
|||||||
return target
|
return target
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function findOutputFile(filename, subfolder, type) {
|
||||||
|
const roots = rootsForType(type || 'output')
|
||||||
|
const subs = [...new Set([String(subfolder || ''), type === 'input' ? '' : 'video', type === 'input' ? '' : 'audio', ''])]
|
||||||
|
for (const root of roots) {
|
||||||
|
if (!existsSync(root)) continue
|
||||||
|
for (const sub of subs) {
|
||||||
|
const path = safeFile(root, sub, filename)
|
||||||
|
if (path && existsSync(path)) {
|
||||||
|
try {
|
||||||
|
if (statSync(path).isFile()) return path
|
||||||
|
} catch { /* ignore */ }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return null
|
||||||
|
}
|
||||||
|
|
||||||
|
function streamFile(res, path) {
|
||||||
|
const size = statSync(path).size
|
||||||
|
const lower = path.toLowerCase()
|
||||||
|
const type = lower.endsWith('.mp4') ? 'video/mp4'
|
||||||
|
: lower.endsWith('.webm') ? 'video/webm'
|
||||||
|
: lower.endsWith('.flac') ? 'audio/flac'
|
||||||
|
: lower.endsWith('.wav') ? 'audio/wav'
|
||||||
|
: lower.endsWith('.mp3') ? 'audio/mpeg'
|
||||||
|
: lower.endsWith('.png') ? 'image/png'
|
||||||
|
: lower.endsWith('.jpg') || lower.endsWith('.jpeg') ? 'image/jpeg'
|
||||||
|
: 'application/octet-stream'
|
||||||
|
res.writeHead(200, {
|
||||||
|
'Content-Type': type,
|
||||||
|
'Content-Length': size,
|
||||||
|
'Cache-Control': 'no-store'
|
||||||
|
})
|
||||||
|
createReadStream(path).pipe(res)
|
||||||
|
}
|
||||||
|
|
||||||
function removeFile(path) {
|
function removeFile(path) {
|
||||||
if (!path || !existsSync(path)) return false
|
if (!path || !existsSync(path)) return false
|
||||||
try {
|
try {
|
||||||
@@ -845,6 +881,15 @@ const server = http.createServer(async (req, res) => {
|
|||||||
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'stopped', reason: 'stop', killed }))
|
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'stopped', reason: 'stop', killed }))
|
||||||
return json(res, 200, { ok: true, stopped: true, asleep: true, killed })
|
return json(res, 200, { ok: true, stopped: true, asleep: true, killed })
|
||||||
}
|
}
|
||||||
|
if (req.method === 'GET' && url.pathname === '/view') {
|
||||||
|
const filename = String(url.searchParams.get('filename') || '')
|
||||||
|
const subfolder = String(url.searchParams.get('subfolder') || '')
|
||||||
|
const type = String(url.searchParams.get('type') || 'output')
|
||||||
|
const path = findOutputFile(filename, subfolder, type)
|
||||||
|
if (!path) return json(res, 404, { ok: false, error: 'not-found', filename, subfolder, type })
|
||||||
|
markWork()
|
||||||
|
return streamFile(res, path)
|
||||||
|
}
|
||||||
if (req.method === 'POST' && url.pathname === '/purge') {
|
if (req.method === 'POST' && url.pathname === '/purge') {
|
||||||
const body = await readJson(req)
|
const body = await readJson(req)
|
||||||
return json(res, 200, purgeDesktopFiles(body))
|
return json(res, 200, purgeDesktopFiles(body))
|
||||||
|
|||||||
+75
-22
@@ -1964,36 +1964,89 @@ export function renameClipFamily(owner: string, id: string, name: string) {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function downloadComfyVideo(video: { filename: string; subfolder: string; type: string }) {
|
function sleep(ms: number) {
|
||||||
const subfolders = [...new Set([video.subfolder, 'video', ''])]
|
return new Promise(resolve => setTimeout(resolve, ms))
|
||||||
let lastStatus = 0
|
}
|
||||||
|
|
||||||
|
function downloadErrorText(error: unknown) {
|
||||||
|
if (!error) return 'Unknown download error'
|
||||||
|
if (typeof error === 'object') {
|
||||||
|
const rec = error as { statusMessage?: string; message?: string; data?: { cause?: string } }
|
||||||
|
const bits = [rec.statusMessage, rec.message, rec.data?.cause].filter(Boolean)
|
||||||
|
if (bits.length) return bits[0] as string
|
||||||
|
}
|
||||||
|
if (error instanceof Error) return error.message
|
||||||
|
return String(error)
|
||||||
|
}
|
||||||
|
|
||||||
|
async function downloadViaHostAgent(file: { filename: string; subfolder: string; type: string }) {
|
||||||
|
const config = useRuntimeConfig()
|
||||||
|
const controlUrl = String(config.comfyControlUrl || process.env.COMFY_CONTROL_URL || '').replace(/\/$/, '')
|
||||||
|
const token = String(config.comfyControlToken || process.env.COMFY_CONTROL_TOKEN || '')
|
||||||
|
if (!controlUrl) return null
|
||||||
|
const subfolders = [...new Set([file.subfolder, file.type === 'input' ? '' : 'video', file.type === 'input' ? '' : 'audio', ''])]
|
||||||
for (const subfolder of subfolders) {
|
for (const subfolder of subfolders) {
|
||||||
const params = new URLSearchParams({
|
const params = new URLSearchParams({
|
||||||
filename: video.filename,
|
filename: file.filename,
|
||||||
subfolder,
|
subfolder,
|
||||||
type: video.type || 'output'
|
type: file.type || 'output'
|
||||||
})
|
})
|
||||||
const res = await comfyFetch(`/view?${params.toString()}`)
|
try {
|
||||||
lastStatus = res.status
|
const res = await fetch(`${controlUrl}/view?${params.toString()}`, {
|
||||||
if (res.ok) return Buffer.from(await res.arrayBuffer())
|
headers: {
|
||||||
|
Accept: '*/*',
|
||||||
|
...(token ? { Authorization: `Bearer ${token}` } : {})
|
||||||
|
},
|
||||||
|
signal: AbortSignal.timeout(5 * 60 * 1000)
|
||||||
|
})
|
||||||
|
if (res.ok) return Buffer.from(await res.arrayBuffer())
|
||||||
|
} catch {
|
||||||
|
/* try next subfolder / fall through */
|
||||||
|
}
|
||||||
}
|
}
|
||||||
throw new Error(`Failed to fetch completed video from ComfyUI (${lastStatus})`)
|
return null
|
||||||
|
}
|
||||||
|
|
||||||
|
async function downloadComfyMedia(
|
||||||
|
file: { filename: string; subfolder: string; type: string },
|
||||||
|
kind: 'video' | 'audio'
|
||||||
|
) {
|
||||||
|
const extras = kind === 'video' ? ['video', ''] : ['audio', '']
|
||||||
|
const subfolders = [...new Set([file.subfolder, ...extras])]
|
||||||
|
let lastStatus = 0
|
||||||
|
let lastError = ''
|
||||||
|
const attempts = 8
|
||||||
|
for (let attempt = 1; attempt <= attempts; attempt++) {
|
||||||
|
for (const subfolder of subfolders) {
|
||||||
|
const params = new URLSearchParams({
|
||||||
|
filename: file.filename,
|
||||||
|
subfolder,
|
||||||
|
type: file.type || 'output'
|
||||||
|
})
|
||||||
|
try {
|
||||||
|
const res = await comfyFetch(`/view?${params.toString()}`, {
|
||||||
|
signal: AbortSignal.timeout(5 * 60 * 1000)
|
||||||
|
})
|
||||||
|
lastStatus = res.status
|
||||||
|
if (res.ok) return Buffer.from(await res.arrayBuffer())
|
||||||
|
} catch (error) {
|
||||||
|
lastError = downloadErrorText(error)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
const disk = await downloadViaHostAgent(file)
|
||||||
|
if (disk?.length) return disk
|
||||||
|
if (attempt < attempts) await sleep(1500 * attempt)
|
||||||
|
}
|
||||||
|
const detail = lastError || `HTTP ${lastStatus || 'n/a'}`
|
||||||
|
throw new Error(`Failed to fetch completed ${kind} from ComfyUI (${detail}). File may still be on the desktop — use Recover.`)
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function downloadComfyVideo(video: { filename: string; subfolder: string; type: string }) {
|
||||||
|
return downloadComfyMedia(video, 'video')
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function downloadComfyAudio(audio: { filename: string; subfolder: string; type: string }) {
|
export async function downloadComfyAudio(audio: { filename: string; subfolder: string; type: string }) {
|
||||||
const subfolders = [...new Set([audio.subfolder, 'audio', ''])]
|
return downloadComfyMedia(audio, 'audio')
|
||||||
let lastStatus = 0
|
|
||||||
for (const subfolder of subfolders) {
|
|
||||||
const params = new URLSearchParams({
|
|
||||||
filename: audio.filename,
|
|
||||||
subfolder,
|
|
||||||
type: audio.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 audio from ComfyUI (${lastStatus})`)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function importMissingComfyVideos(owner: string, folderId?: string) {
|
export async function importMissingComfyVideos(owner: string, folderId?: string) {
|
||||||
|
|||||||
+51
-7
@@ -19,7 +19,7 @@ function classifyError(message: string) {
|
|||||||
if (lower.includes('out of memory') || (lower.includes('cuda') && lower.includes('alloc')) || lower.includes('vram')) {
|
if (lower.includes('out of memory') || (lower.includes('cuda') && lower.includes('alloc')) || lower.includes('vram')) {
|
||||||
return 'ComfyUI VRAM allocation failed. Try Turbo or a smaller aspect ratio.'
|
return 'ComfyUI VRAM allocation failed. Try Turbo or a smaller aspect ratio.'
|
||||||
}
|
}
|
||||||
if (lower.includes('econnrefused') || lower.includes('unreachable') || lower.includes('fetch failed')) {
|
if (lower.includes('econnrefused') || lower.includes('unreachable') || lower.includes('fetch failed') || lower.includes('connection dropped')) {
|
||||||
return 'ComfyUI host connection dropped. Confirm the desktop instance is running and reachable on the LAN.'
|
return 'ComfyUI host connection dropped. Confirm the desktop instance is running and reachable on the LAN.'
|
||||||
}
|
}
|
||||||
if (lower.includes('timeout')) {
|
if (lower.includes('timeout')) {
|
||||||
@@ -28,6 +28,17 @@ function classifyError(message: string) {
|
|||||||
return message
|
return message
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function errorText(error: unknown) {
|
||||||
|
if (!error) return 'Unknown error'
|
||||||
|
if (typeof error === 'object') {
|
||||||
|
const rec = error as { statusMessage?: string; message?: string; data?: { cause?: string } }
|
||||||
|
const bits = [rec.statusMessage, rec.message, rec.data?.cause].filter(Boolean)
|
||||||
|
if (bits.length) return String(bits[0])
|
||||||
|
}
|
||||||
|
if (error instanceof Error) return error.message
|
||||||
|
return String(error)
|
||||||
|
}
|
||||||
|
|
||||||
function sleep(ms: number) {
|
function sleep(ms: number) {
|
||||||
return new Promise(resolve => setTimeout(resolve, ms))
|
return new Promise(resolve => setTimeout(resolve, ms))
|
||||||
}
|
}
|
||||||
@@ -175,8 +186,23 @@ export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Pr
|
|||||||
finishing = true
|
finishing = true
|
||||||
job.saving = true
|
job.saving = true
|
||||||
try {
|
try {
|
||||||
const history = await fetchHistory(job.promptId)
|
let history: Awaited<ReturnType<typeof fetchHistory>> | null = null
|
||||||
const video = extractVideo(history, job.promptId)
|
let video: ReturnType<typeof extractVideo> = null
|
||||||
|
for (let attempt = 1; attempt <= 6; attempt++) {
|
||||||
|
try {
|
||||||
|
history = await fetchHistory(job.promptId)
|
||||||
|
video = extractVideo(history, job.promptId)
|
||||||
|
if (video) break
|
||||||
|
} catch (error) {
|
||||||
|
if (attempt >= 6) throw error
|
||||||
|
}
|
||||||
|
emitLocal({
|
||||||
|
type: 'status',
|
||||||
|
message: `Waiting for ComfyUI history (${attempt}/6)...`,
|
||||||
|
progress: 96
|
||||||
|
})
|
||||||
|
await sleep(1000 * attempt)
|
||||||
|
}
|
||||||
if (!video) {
|
if (!video) {
|
||||||
finishing = false
|
finishing = false
|
||||||
job.saving = false
|
job.saving = false
|
||||||
@@ -312,12 +338,20 @@ export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Pr
|
|||||||
} catch (saveError) {
|
} catch (saveError) {
|
||||||
job.saving = false
|
job.saving = false
|
||||||
removeExtendTemp(job.library?.extendTmpDir)
|
removeExtendTemp(job.library?.extendTmpDir)
|
||||||
const message = saveError instanceof Error ? saveError.message : String(saveError)
|
const message = classifyError(errorText(saveError))
|
||||||
job.status = 'error'
|
job.status = 'error'
|
||||||
job.error = persist
|
job.error = persist
|
||||||
? `Video generated but library save failed: ${message}`
|
? `Video generated but library save failed: ${message}`
|
||||||
: `Segment finished but could not be prepared: ${message}`
|
: `Segment finished but could not be prepared: ${message}`
|
||||||
emitChainJob(job, { type: 'error', error: job.error, message: job.error })
|
emitChainJob(job, { type: 'error', error: job.error, message: job.error })
|
||||||
|
// Keep pending so a later Recover / container restart can still pull the MP4.
|
||||||
|
if (job.library && job.promptId) {
|
||||||
|
writePendingJob(pendingFromJob(job, {
|
||||||
|
promptId: job.promptId,
|
||||||
|
currentClipId: job.clipId || job.library.parentClipId
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
await syncQueueFromJob(job).catch(() => null)
|
||||||
}
|
}
|
||||||
resolve()
|
resolve()
|
||||||
return true
|
return true
|
||||||
@@ -406,9 +440,19 @@ export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Pr
|
|||||||
if (payload.type === 'executing') {
|
if (payload.type === 'executing') {
|
||||||
const node = data.node === null || data.node === undefined ? null : String(data.node)
|
const node = data.node === null || data.node === undefined ? null : String(data.node)
|
||||||
if (node === null) {
|
if (node === null) {
|
||||||
const saved = await succeed()
|
try {
|
||||||
if (!saved) {
|
const saved = await succeed()
|
||||||
emitLocal({ type: 'status', message: 'Waiting for ComfyUI to save the MP4...', progress: 96 })
|
if (!saved) {
|
||||||
|
emitLocal({ type: 'status', message: 'Waiting for ComfyUI to save the MP4...', progress: 96 })
|
||||||
|
}
|
||||||
|
} catch (error) {
|
||||||
|
finishing = false
|
||||||
|
job.saving = false
|
||||||
|
emitLocal({
|
||||||
|
type: 'status',
|
||||||
|
message: `Save hiccup (${errorText(error)}). Retrying...`,
|
||||||
|
progress: 96
|
||||||
|
})
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user