Files
aigen/server/utils/imageComfy.ts
T
TowstyandCursor c043890ac0 Stop Comfy input folder filling with duplicate still uploads.
Name uploads by content hash so the same still overwrites itself, enable stale input sweeps on purge, and remove the pile already on the desktop Shared input folder.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-04 08:20:39 -05:00

457 lines
16 KiB
TypeScript

import { createHash } from 'node:crypto'
import { AsyncLocalStorage } from 'node:async_hooks'
import { getComfyHost, viaAgentFrontDoor } from '~/server/utils/comfy'
import { comfyJobPrefix } from '~/utils/outputNames'
const imageHostAls = new AsyncLocalStorage<string>()
let imageComfyHostOverride = ''
export type ImageEditBox = 'beast' | 'sidecar'
export function withImageComfyHost<T>(host: string, fn: () => Promise<T>) {
return imageHostAls.run(String(host || '').replace(/\/$/, ''), fn)
}
export function setImageComfyHostOverride(url: string) {
imageComfyHostOverride = String(url || '').replace(/\/$/, '')
}
function normalizeHost(raw: string, port = '') {
let host = String(raw || '').trim().replace(/\/$/, '')
if (!host) return ''
if (!/^https?:\/\//i.test(host)) host = `http://${host}`
try {
const url = new URL(host)
if (port && !url.port) url.port = port
return url.origin
} catch {
return port ? `${host}:${port}` : host
}
}
export function getBeastImageHost() {
const config = useRuntimeConfig()
const explicit = String(config.imageComfyPrimaryHost || process.env.IMAGE_COMFY_PRIMARY_HOST || '').trim()
if (explicit) return viaAgentFrontDoor(normalizeHost(explicit))
try {
return getComfyHost()
} catch {
return ''
}
}
function isRetiredSidecarHost(raw: string) {
return /192\.168\.77\.101(?!\d)/.test(raw)
}
export function getSidecarImageHost() {
const config = useRuntimeConfig()
const port = String(config.imageComfyPort || process.env.IMAGE_COMFY_PORT || '').trim()
const candidates: Array<[string, string]> = [
[String(config.promptComfyHost || process.env.PROMPT_COMFY_HOST || '').trim(), ''],
[String(config.imageComfyHost || process.env.IMAGE_COMFY_HOST || '').trim(), port],
[String(config.imageComfyFallbackHost || process.env.IMAGE_COMFY_FALLBACK_HOST || '').trim(), '']
]
for (const [raw, extraPort] of candidates) {
if (!raw || isRetiredSidecarHost(raw)) continue
const host = viaAgentFrontDoor(normalizeHost(raw, extraPort))
if (host && !isRetiredSidecarHost(host)) return host
}
return getBeastImageHost()
}
export function getImageComfyHost() {
const fromAls = imageHostAls.getStore()
if (fromAls) return viaAgentFrontDoor(fromAls) || fromAls
if (imageComfyHostOverride) return viaAgentFrontDoor(imageComfyHostOverride) || imageComfyHostOverride
return getBeastImageHost() || getSidecarImageHost()
}
export function imageComfyConfigured() {
return Boolean(getBeastImageHost())
}
export function sameImageHost(a: string, b: string) {
return Boolean(a && b && a.replace(/\/$/, '') === b.replace(/\/$/, ''))
}
export async function imageComfyFetch(path: string, init?: RequestInit) {
const host = getImageComfyHost()
if (!host) {
throw createError({ statusCode: 500, statusMessage: 'No image ComfyUI host is configured' })
}
try {
return await fetch(`${host}${path}`, init)
} catch (error) {
throw createError({
statusCode: 502,
statusMessage: `Image ComfyUI host unreachable (${host})`,
data: { cause: error instanceof Error ? error.message : String(error) }
})
}
}
export function imageComfyWsUrl(clientId: string) {
return `${getImageComfyHost().replace(/^http/, 'ws')}/ws?clientId=${encodeURIComponent(clientId)}`
}
function imageInputFilename(original: string, data?: Buffer, jobId?: string) {
const raw = String(original || 'still.png')
const dot = raw.lastIndexOf('.')
const ext = (dot >= 0 ? raw.slice(dot) : '.png').replace(/[^.a-zA-Z0-9]/g, '') || '.png'
const base = (dot >= 0 ? raw.slice(0, dot) : raw).replace(/[^a-zA-Z0-9._-]+/g, '_').slice(0, 48) || 'still'
const hash = data?.length
? createHash('sha1').update(data).digest('hex').slice(0, 12)
: comfyJobPrefix(jobId)
return `aigen_${hash}_${base}${ext}`
}
export async function uploadImageEdit(file: { filename: string; data: Buffer; type?: string }, jobId?: string) {
const body = new FormData()
const blob = new Blob([new Uint8Array(file.data)], { type: file.type || 'application/octet-stream' })
const filename = imageInputFilename(file.filename, file.data, jobId)
body.append('image', blob, filename)
body.append('overwrite', 'true')
body.append('type', 'input')
const res = await imageComfyFetch('/upload/image', { method: 'POST', body })
if (!res.ok) {
throw createError({ statusCode: 502, statusMessage: `Image upload failed (${res.status})` })
}
const uploaded = (await res.json()) as { name: string; subfolder?: string; type?: string }
return {
name: uploaded.name || filename,
subfolder: uploaded.subfolder || '',
type: uploaded.type || 'input'
}
}
export async function queueImagePrompt(graph: unknown, clientId: string) {
const res = await imageComfyFetch('/prompt', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ prompt: graph, client_id: clientId })
})
const payload = await res.json().catch(() => ({}))
if (!res.ok) {
const message = (payload as { error?: { message?: string } }).error?.message
|| (payload as { node_errors?: unknown }).node_errors
|| `Queue failed (${res.status})`
throw createError({ statusCode: 502, statusMessage: String(message), data: payload })
}
return payload as { prompt_id: string; number?: number }
}
export async function interruptImageComfy(host?: string) {
if (!imageComfyConfigured() && !host) return false
const run = async () => {
const res = await imageComfyFetch('/interrupt', { method: 'POST' })
return res.ok
}
if (host) return withImageComfyHost(host, run)
return run()
}
export async function fetchImageHistory(promptId: string) {
const res = await imageComfyFetch(`/history/${encodeURIComponent(promptId)}`)
if (!res.ok) return null
return (await res.json()) as Record<string, unknown>
}
function isStillFile(item: { filename?: string } | null | undefined) {
if (!item?.filename) return false
return /\.(png|jpe?g|webp)$/i.test(item.filename)
}
function findStill(value: unknown, depth = 0): { filename: string; subfolder: string; type: string } | null {
if (!value || typeof value !== 'object' || depth > 8) return null
if (Array.isArray(value)) {
const files = value.filter((item): item is { filename?: string; subfolder?: string; type?: string } => Boolean(item && typeof item === 'object'))
const match = files.find(isStillFile)
if (match?.filename) {
return {
filename: String(match.filename),
subfolder: String(match.subfolder || ''),
type: String(match.type || 'output')
}
}
for (const item of value) {
const nested = findStill(item, depth + 1)
if (nested) return nested
}
return null
}
const record = value as { filename?: string; subfolder?: string; type?: string; images?: unknown }
if (record.filename && isStillFile(record)) {
return {
filename: String(record.filename),
subfolder: String(record.subfolder || ''),
type: String(record.type || 'output')
}
}
for (const nested of Object.values(value as Record<string, unknown>)) {
const found = findStill(nested, depth + 1)
if (found) return found
}
return null
}
export function extractEditedImage(history: Record<string, unknown> | null, promptId: string) {
if (!history) return null
const wrapped = history[promptId] as { outputs?: Record<string, unknown> } | undefined
if (wrapped) return findStill(wrapped.outputs || {}) || findStill(wrapped)
if ((history as { outputs?: unknown }).outputs) {
return findStill((history as { outputs?: unknown }).outputs) || findStill(history)
}
return findStill(history)
}
export async function downloadEditedImage(image: { filename: string; subfolder: string; type: string }) {
const params = new URLSearchParams({
filename: image.filename,
subfolder: image.subfolder || '',
type: image.type || 'output'
})
const res = await imageComfyFetch(`/view?${params.toString()}`)
if (!res.ok) {
throw createError({ statusCode: 502, statusMessage: `Failed to fetch edited still (${res.status})` })
}
return Buffer.from(await res.arrayBuffer())
}
function historyError(history: Record<string, unknown> | null, promptId: string) {
const entry = history?.[promptId] as {
status?: {
status_str?: string
completed?: boolean
messages?: Array<[string, Record<string, unknown>]>
}
} | undefined
const status = entry?.status?.status_str
const fromMessages = formatExecutionError(entry?.status?.messages?.find(([type]) => type === 'execution_error')?.[1])
if (status === 'interrupted') return 'Job interrupted.'
if (status === 'error') return fromMessages || 'Image ComfyUI reported an execution error'
if (entry?.status?.completed && !extractEditedImage(history, promptId)) {
return 'Image ComfyUI finished without an output still'
}
return null
}
function formatExecutionError(data?: Record<string, unknown> | null) {
if (!data) return ''
const node = String(data.node_type || data.node_id || '').trim()
const message = String(data.exception_message || data.message || '').trim()
if (node && message) return `${node}: ${message}`
return message
}
export async function waitForImageEdit(opts: {
promptId: string
clientId: string
timeoutMs?: number
engineLabel?: string
nodeLabel?: (node: string) => string | undefined
onProgress?: (event: { message: string; progress: number; step?: number; maxStep?: number; node?: string | null }) => void
isCancelled?: () => boolean
}) {
const timeoutMs = opts.timeoutMs || 180_000
const engineLabel = String(opts.engineLabel || 'Flux.2 Klein').trim() || 'Flux.2 Klein'
const started = Date.now()
let settled = false
let lastError: string | null = null
let result: { filename: string; subfolder: string; type: string } | null = null
const finish = (image: { filename: string; subfolder: string; type: string } | null, error?: string) => {
if (settled) return
settled = true
lastError = error || null
result = image
}
let ws: WebSocket | null = null
try {
ws = new WebSocket(imageComfyWsUrl(opts.clientId))
ws.addEventListener('message', (event) => {
const payload = JSON.parse(String(event.data || '{}')) as {
type?: string
data?: {
value?: number
max?: number
node?: string | null
prompt_id?: string
exception_message?: string
exception_type?: string
node_id?: string
node_type?: string
message?: string
}
}
if (payload.type === 'progress' && payload.data) {
const max = Math.max(1, Number(payload.data.max || 4))
const value = Number(payload.data.value || 0)
const node = payload.data.node || null
opts.onProgress?.({
message: `Sampling ${engineLabel}...`,
progress: Math.min(90, 20 + Math.round((value / max) * 60)),
step: value,
maxStep: max,
node
})
}
if (payload.type === 'executing' && payload.data?.node) {
const node = payload.data.node
opts.onProgress?.({
message: opts.nodeLabel?.(node) || `Running ${engineLabel}...`,
progress: node === '9' || node === '94' || node === '21' ? 92 : 30,
node
})
}
if (payload.type === 'execution_error') {
finish(null, formatExecutionError(payload.data) || 'Image ComfyUI reported an execution error')
}
})
} catch {
ws = null
}
while (!settled && Date.now() - started < timeoutMs) {
if (opts.isCancelled?.()) {
finish(null, 'Job interrupted.')
break
}
const history = await fetchImageHistory(opts.promptId)
const image = extractEditedImage(history, opts.promptId)
if (image) {
finish(image)
break
}
const error = historyError(history, opts.promptId)
if (error) {
finish(null, error)
break
}
await new Promise(resolve => setTimeout(resolve, 700))
}
try { ws?.close() } catch { /* ignore */ }
if (result) return result
throw createError({
statusCode: 502,
statusMessage: lastError || `Image edit timed out waiting for ${engineLabel} to finish`
})
}
function purgeImageEnabled() {
return useRuntimeConfig().purgeComfyOutputs !== false
}
async function deleteSidecarFile(file: { filename: string; subfolder?: string; type?: string }) {
const payload = {
filename: file.filename,
subfolder: file.subfolder || '',
type: file.type || 'output'
}
for (const path of ['/delete', '/aigen/purge']) {
try {
const res = await imageComfyFetch(path, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(payload),
signal: AbortSignal.timeout(4000)
})
if (res.ok) return true
} catch {
/* Comfy may not expose this route */
}
}
return false
}
async function purgeImageOnDesktop(opts: {
inputName?: string
inputSubfolder?: string
extraInputNames?: string[]
output?: { filename: string; subfolder: string; type: string }
}) {
const config = useRuntimeConfig()
const controlUrl = String(
config.imageComfyControlUrl || process.env.IMAGE_COMFY_CONTROL_URL || config.comfyControlUrl || process.env.COMFY_CONTROL_URL || ''
).replace(/\/$/, '')
const token = String(
config.imageComfyControlToken || process.env.IMAGE_COMFY_CONTROL_TOKEN || config.comfyControlToken || process.env.COMFY_CONTROL_TOKEN || ''
)
if (!controlUrl) return
const imageNames = [
opts.inputName,
...(Array.isArray(opts.extraInputNames) ? opts.extraInputNames : [])
].map(name => String(name || '').trim()).filter(Boolean)
try {
await fetch(`${controlUrl}/purge`, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
Accept: 'application/json',
...(token ? { Authorization: `Bearer ${token}` } : {})
},
body: JSON.stringify({
imageName: imageNames[0] || '',
imageNames,
imageSubfolder: opts.inputSubfolder || '',
output: opts.output || null,
sweep: true,
sweepMaxAgeMs: 2 * 60 * 1000
}),
signal: AbortSignal.timeout(8000)
})
} catch {
/* Host agent purge is optional; Comfy HTTP delete is the fallback. */
}
}
export async function purgeImageComfyArtifacts(opts: {
output?: { filename: string; subfolder: string; type: string }
inputName?: string
inputSubfolder?: string
extraInputNames?: string[]
promptId?: string
}) {
if (!purgeImageEnabled()) return
await purgeImageOnDesktop({
inputName: opts.inputName,
inputSubfolder: opts.inputSubfolder,
extraInputNames: opts.extraInputNames,
output: opts.output
})
const files: { filename: string; subfolder: string; type: string }[] = []
if (opts.output?.filename) {
files.push({
filename: opts.output.filename,
subfolder: opts.output.subfolder || '',
type: opts.output.type || 'output'
})
}
const inputSub = opts.inputSubfolder || ''
if (opts.inputName) {
files.push({ filename: opts.inputName, subfolder: inputSub, type: 'input' })
}
for (const name of opts.extraInputNames || []) {
if (!name) continue
files.push({ filename: name, subfolder: inputSub, type: 'input' })
}
for (const file of files) {
await deleteSidecarFile(file)
}
if (opts.promptId) {
try {
await imageComfyFetch('/history', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ delete: [opts.promptId] }),
signal: AbortSignal.timeout(4000)
})
} catch {
/* ignore */
}
}
}