Files
aigen/scripts/comfy-host-agent.mjs
T
TowstyandCursor a9c2c01b35 Delete Comfy input stills after library save, not only outputs.
Purge was skipping Beast input copies and never sending image names to the host agent, so studio temps piled up on the desktop. Also sweep stale hex-/aigen-prefixed inputs.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-03 23:27:14 -05:00

989 lines
33 KiB
JavaScript

import http from 'node:http'
import net from 'node:net'
import { execFile, spawn } from 'node:child_process'
import { promisify } from 'node:util'
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 { tmpdir } from 'node:os'
const execFileAsync = promisify(execFile)
const port = Number(process.env.COMFY_CONTROL_PORT || 8199)
const proxyPort = Number(process.env.COMFY_PROXY_PORT || 8198)
const token = process.env.COMFY_CONTROL_TOKEN || ''
const defaultHttp = (process.env.COMFY_LISTEN || 'http://127.0.0.1:8188').replace(/\/$/, '')
const idleMs = Math.max(60_000, Number(process.env.COMFY_IDLE_MS || 30 * 60 * 1000) || 30 * 60 * 1000)
const trainingUrl = String(process.env.TRAINING_CONTROL_URL || 'http://127.0.0.1:8200').replace(/\/$/, '')
let lastWorkAt = Date.now()
let lastQueueRunning = 0
let lastLaunchAt = 0
let stoppedByAgent = false
let lastHealthyPort = 0
let lastProcessUp = false
let lastTraining = { busy: false, jobId: null, status: null, name: null, message: '' }
function markWork() {
lastWorkAt = Date.now()
stoppedByAgent = false
}
async function trainingLock() {
try {
const res = await fetch(`${trainingUrl}/status`, { signal: AbortSignal.timeout(2500) })
const body = await res.json().catch(() => null)
const job = body?.job
const status = String(job?.status || '')
const busy = Boolean(job && ['preparing', 'queued', 'running', 'stopping'].includes(status))
lastTraining = {
busy,
jobId: job?.id || null,
status: busy ? status : null,
name: busy ? (job?.outputName || null) : null,
message: busy
? `AITraining is ${status}${job?.outputName ? ` (${job.outputName})` : ''}. Stop that job before poking Comfy.`
: ''
}
return lastTraining
} catch {
lastTraining = { busy: false, jobId: null, status: null, name: null, message: '' }
return lastTraining
}
}
function json(res, status, body) {
const payload = JSON.stringify(body)
res.writeHead(status, { 'Content-Type': 'application/json', 'Content-Length': Buffer.byteLength(payload) })
res.end(payload)
}
function authorized(req) {
if (!token) return true
const header = String(req.headers.authorization || '')
return header === `Bearer ${token}`
}
function lockPorts() {
const ports = []
try {
const dir = join(process.env.APPDATA || '', 'Comfy Desktop', 'port-locks')
for (const name of readdirSync(dir)) {
const match = name.match(/^port-(\d+)\.json$/)
if (match) ports.push(Number(match[1]))
}
} catch { /* ignore */ }
return ports
}
function candidatePorts() {
const skip = new Set([port, proxyPort])
const ordered = []
const seen = new Set()
const add = (value) => {
const next = Number(value)
if (!Number.isInteger(next) || next < 1 || next > 65535) return
if (skip.has(next) || seen.has(next)) return
seen.add(next)
ordered.push(next)
}
for (const next of lockPorts()) add(next)
try {
const url = new URL(defaultHttp.includes('://') ? defaultHttp : `http://${defaultHttp}`)
if (url.port) add(url.port)
} catch { /* ignore */ }
for (let next = 8188; next <= 8210; next++) add(next)
return ordered
}
const probeSkipUntil = new Map()
async function probeStats(portNum, timeoutMs = 400) {
if (Date.now() < (probeSkipUntil.get(portNum) || 0)) return null
try {
const res = await fetch(`http://127.0.0.1:${portNum}/system_stats`, { signal: AbortSignal.timeout(timeoutMs) })
if (!res.ok) return null
const stats = await res.json()
if (!stats?.system) return null
probeSkipUntil.delete(portNum)
return stats
} catch (error) {
const timedOut = error?.name === 'TimeoutError' || error?.name === 'AbortError' || /timeout/i.test(String(error?.message || ''))
if (timedOut) probeSkipUntil.set(portNum, Date.now() + 30_000)
return null
}
}
function advertisedPort(stats) {
const argv = stats?.system?.argv
if (!Array.isArray(argv)) return 0
const index = argv.findIndex(item => String(item) === '--port')
if (index >= 0) return Number(argv[index + 1]) || 0
const flag = argv.find(item => /^--port=\d+$/.test(String(item)))
if (flag) return Number(String(flag).split('=')[1]) || 0
return 0
}
async function findHealthyPort() {
const probed = await Promise.all(candidatePorts().map(async (next) => {
const stats = await probeStats(next)
return stats ? { port: next, advertised: advertisedPort(stats) } : null
}))
const found = probed.filter(Boolean)
const native = found.find(item => item.advertised && item.advertised === item.port)
if (native) return native.port
for (const item of found) {
if (item.advertised && found.some(other => other.port === item.advertised)) return item.advertised
}
return found[0]?.port || 0
}
let proxyServer = null
let proxyTarget = 0
function ensureProxyListening() {
if (proxyServer) return
proxyServer = net.createServer((client) => {
const target = proxyTarget
if (!target) {
client.destroy()
return
}
const upstream = net.connect(target, '127.0.0.1')
const fail = () => {
try { client.destroy() } catch { /* ignore */ }
try { upstream.destroy() } catch { /* ignore */ }
}
client.on('error', fail)
upstream.on('error', fail)
client.once('data', (chunk) => {
if (isWorkHttp(chunk)) markWork()
upstream.write(chunk)
client.pipe(upstream)
})
upstream.pipe(client)
})
proxyServer.on('error', (error) => {
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'proxy-error', error: String(error.message || error) }))
})
proxyServer.listen(proxyPort, '0.0.0.0', () => {
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'proxy', listen: proxyPort, target: proxyTarget || null }))
})
}
function setProxyTarget(targetPort) {
ensureProxyListening()
const next = Number(targetPort) || 0
if (proxyTarget === next) return
proxyTarget = next
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'proxy-target', listen: proxyPort, target: proxyTarget || null }))
}
function markAsleep() {
stoppedByAgent = true
lastQueueRunning = 0
lastHealthyPort = 0
lastProcessUp = false
lastLaunchAt = 0
probeSkipUntil.clear()
setProxyTarget(0)
}
async function syncProxy() {
const healthyPort = await findHealthyPort()
lastHealthyPort = healthyPort || 0
if (healthyPort) {
setProxyTarget(healthyPort)
lastProcessUp = true
} else {
setProxyTarget(0)
lastProcessUp = false
}
return healthyPort
}
// PowerShell -Command text used to contain the literal main.py path, so the
// probe process matched itself and /start forever returned already/booting.
const PS_COMFY_MAIN =
"Get-CimInstance Win32_Process | Where-Object { $_.Name -match '^python(w)?\\.exe$' -and $_.CommandLine -match ('ComfyUI' + '[\\\\/]main\\.py') } | Select-Object -First 1 -ExpandProperty ProcessId"
async function processUp() {
try {
const { stdout } = await execFileAsync('tasklist', ['/FO', 'CSV', '/NH'], { windowsHide: true, timeout: 4000 })
if (/Comfy Desktop|ComfyUI/i.test(stdout)) return true
} catch { /* ignore */ }
try {
const { stdout } = await execFileAsync('powershell.exe', [
'-NoProfile', '-Command',
PS_COMFY_MAIN
], { windowsHide: true, timeout: 6000 })
return Boolean(String(stdout).trim())
} catch {
return false
}
}
function readJsonFile(path) {
try {
return JSON.parse(readFileSync(path, 'utf8'))
} catch {
return null
}
}
function desktopListenPort() {
try {
const url = new URL(defaultHttp.includes('://') ? defaultHttp : `http://${defaultHttp}`)
const next = Number(url.port || 8188) || 8188
if (next === port || next === proxyPort) return 8188
return next
} catch {
return 8188
}
}
function pickDesktopInstall() {
const wanted = String(process.env.COMFY_INSTANCE_NAME || 'ComfyUI (1)').trim()
const parsed = readJsonFile(join(process.env.APPDATA || '', 'Comfy Desktop', 'installations.json'))
const installs = (Array.isArray(parsed) ? parsed : [])
.filter(item => item && item.status === 'installed' && item.installPath && item.sourceId !== 'cloud')
if (!installs.length) return null
return installs.find(item => item.name === wanted)
|| installs.slice().sort((a, b) => Number(b.lastLaunchedAt || 0) - Number(a.lastLaunchedAt || 0))[0]
}
function resolveHeadlessLaunch() {
if (process.env.COMFY_LAUNCH_CMD) {
return {
command: process.env.COMFY_LAUNCH_CMD,
cwd: process.env.COMFY_LAUNCH_CWD || undefined,
mode: 'env'
}
}
const inst = pickDesktopInstall()
if (!inst) return null
const venvPython = join(inst.installPath, 'ComfyUI', '.venv', 'Scripts', 'python.exe')
const standalonePython = join(inst.installPath, 'standalone-env', 'python.exe')
const python = existsSync(venvPython) ? venvPython : standalonePython
const cwd = inst.installPath
const main = join(inst.installPath, 'ComfyUI', 'main.py')
if (!existsSync(python) || !existsSync(main)) return null
const extra = join(process.env.APPDATA || '', 'Comfy Desktop', 'instance-model-paths', `${inst.id}.yaml`)
const shared = join(process.env.APPDATA || '', 'Comfy Desktop', 'shared_model_paths.yaml')
const extraConfig = existsSync(extra) ? extra : (existsSync(shared) ? shared : '')
const roots = sharedRoots()
const args = [
'-s',
'-u',
join('ComfyUI', 'main.py'),
'--listen', '0.0.0.0',
'--port', String(desktopListenPort()),
'--disable-auto-launch'
]
const extraArgs = String(inst.launchArgs || '').trim()
if (extraArgs) args.push(...extraArgs.split(/\s+/).filter(Boolean))
if (extraConfig) args.push('--extra-model-paths-config', extraConfig)
if (existsSync(roots.input)) args.push('--input-directory', roots.input)
if (existsSync(roots.output)) args.push('--output-directory', roots.output)
return {
file: python,
args,
cwd,
mode: 'headless',
instance: inst.name,
log: join(inst.installPath, 'logs', 'aigen-headless.log')
}
}
function desktopShellCommand() {
const local = process.env.LOCALAPPDATA || ''
const desktop = `${local}\\Programs\\Comfy Desktop\\Comfy Desktop.exe`
return `start "" "${desktop}"`
}
async function pythonMainUp() {
try {
const { stdout } = await execFileAsync('tasklist', ['/FO', 'CSV', '/NH'], { windowsHide: true, timeout: 2000 })
if (!/python(?:w)?\.exe/i.test(stdout)) return false
const { stdout: verbose } = await execFileAsync('powershell.exe', [
'-NoProfile', '-Command',
PS_COMFY_MAIN
], { windowsHide: true, timeout: 2500 })
return Boolean(String(verbose).trim())
} catch {
return false
}
}
function portInUse(portNum) {
return new Promise((resolve) => {
const probe = net.createServer()
probe.once('error', () => resolve(true))
probe.once('listening', () => {
probe.close(() => resolve(false))
})
probe.listen(portNum, '0.0.0.0')
})
}
async function pickListenPort() {
const preferred = desktopListenPort()
const ordered = [preferred, ...candidatePorts()]
const seen = new Set()
for (const next of ordered) {
if (!next || seen.has(next)) continue
seen.add(next)
if (await probeStats(next, 400)) return next
if (!(await portInUse(next))) return next
}
return preferred
}
function windowlessPython(pythonPath) {
if (process.platform !== 'win32') return pythonPath
if (!/python\.exe$/i.test(pythonPath)) return pythonPath
const pythonw = pythonPath.replace(/python\.exe$/i, 'pythonw.exe')
return existsSync(pythonw) ? pythonw : pythonPath
}
function vbsString(value) {
return `"${String(value).replace(/"/g, '""')}"`
}
function spawnViaWscript(file, args, opts = {}) {
const command = [file, ...args].map((part) => {
const text = String(part)
return /\s|"/.test(text) ? `"${text.replace(/"/g, '\\"')}"` : text
}).join(' ')
const lines = ['Set sh = CreateObject("WScript.Shell")']
if (opts.cwd) lines.push(`sh.CurrentDirectory = ${vbsString(opts.cwd)}`)
lines.push(`sh.Run ${vbsString(command)}, 0, False`)
const tmp = join(tmpdir(), `aigen-launch-${process.pid}-${Date.now()}.vbs`)
writeFileSync(tmp, `${lines.join('\r\n')}\r\n`)
const child = spawn('wscript.exe', ['//B', '//Nologo', tmp], {
detached: true,
stdio: 'ignore',
windowsHide: true
})
child.unref()
const cleanup = () => {
try { unlinkSync(tmp) } catch { /* ignore */ }
}
child.once('exit', cleanup)
setTimeout(cleanup, 15000)
return child
}
function spawnHidden(file, args, opts = {}) {
const exe = windowlessPython(file)
// python.exe + detached:true allocates an empty console on the active display.
// pythonw.exe is windowless. If it is missing, launch via hidden WScript.
if (process.platform === 'win32' && /python\.exe$/i.test(exe)) {
return spawnViaWscript(exe, args, opts)
}
const child = spawn(exe, args, {
detached: true,
stdio: opts.stdio || 'ignore',
windowsHide: true,
cwd: opts.cwd
})
child.unref()
return child
}
function spawnShellHidden(command, cwd) {
if (process.platform === 'win32') {
const match = String(command).match(/start\s+""\s+"([^"]+)"/i)
|| String(command).match(/^"([^"]+\.exe)"\s*$/i)
if (match) return spawnHidden(match[1], [], { cwd })
return spawnViaWscript(process.env.ComSpec || 'cmd.exe', ['/d', '/s', '/c', command], { cwd })
}
const child = spawn(command, { shell: true, detached: true, stdio: 'ignore', cwd })
child.unref()
return child
}
async function startComfy() {
markWork()
const launch = resolveHeadlessLaunch()
if (launch?.file) {
const port = await pickListenPort()
const portIdx = launch.args.indexOf('--port')
if (portIdx >= 0) launch.args[portIdx + 1] = String(port)
let stdio = 'ignore'
try {
mkdirSync(dirname(launch.log), { recursive: true })
stdio = ['ignore', openSync(launch.log, 'a'), openSync(launch.log, 'a')]
} catch { /* ignore */ }
const child = spawnHidden(launch.file, launch.args, {
stdio,
cwd: launch.cwd
})
console.log(JSON.stringify({
src: 'comfy-host-agent',
event: 'start',
mode: launch.mode,
instance: launch.instance,
cwd: launch.cwd,
port,
pid: child.pid || null
}))
return { started: true, mode: launch.mode, instance: launch.instance, cwd: launch.cwd, port, pid: child.pid || null }
}
const fallback = {
command: launch?.command || desktopShellCommand(),
cwd: launch?.cwd || process.env.COMFY_LAUNCH_CWD || undefined,
mode: launch?.mode || 'desktop-shell'
}
spawnShellHidden(fallback.command, fallback.cwd)
console.log(JSON.stringify({
src: 'comfy-host-agent',
event: 'start',
mode: fallback.mode,
cwd: fallback.cwd || null
}))
return { started: true, mode: fallback.mode, instance: null, cwd: fallback.cwd || null }
}
async function fetchLocalQueue(portNum) {
if (!portNum) return { running: 0, pending: 0, ok: false }
try {
const res = await fetch(`http://127.0.0.1:${portNum}/queue`, { signal: AbortSignal.timeout(2500) })
if (!res.ok) return { running: 0, pending: 0, ok: false }
const payload = await res.json()
return {
running: Array.isArray(payload?.queue_running) ? payload.queue_running.length : 0,
pending: Array.isArray(payload?.queue_pending) ? payload.queue_pending.length : 0,
ok: true
}
} catch {
return { running: 0, pending: 0, ok: false }
}
}
function isWorkHttp(chunk) {
const head = chunk.toString('latin1', 0, Math.min(chunk.length, 256))
const line = head.split('\r\n')[0] || ''
const match = line.match(/^(GET|HEAD|POST|PUT|PATCH|DELETE)\s+(\S+)/i)
if (!match) return false
const method = match[1].toUpperCase()
let path = String(match[2] || '').split('?')[0] || ''
try { path = decodeURIComponent(path) } catch { /* keep raw */ }
if (method === 'GET' || method === 'HEAD') return false
if (/^\/(system_stats|queue|object_info|internal)\b/i.test(path)) return false
return true
}
async function taskkillImage(image) {
try {
await execFileAsync('taskkill', ['/IM', image, '/F', '/T'], { windowsHide: true, timeout: 8000 })
return true
} catch {
return false
}
}
async function stopComfyProcesses() {
const killed = []
if (await taskkillImage('Comfy Desktop.exe')) killed.push('Comfy Desktop.exe')
if (await taskkillImage('ComfyUI.exe')) killed.push('ComfyUI.exe')
try {
const { stdout } = await execFileAsync('powershell.exe', [
'-NoProfile', '-Command',
"Get-CimInstance Win32_Process | Where-Object { $_.Name -match '^python(w)?\\.exe$' -and $_.CommandLine -match ('ComfyUI' + '[\\\\/]main\\.py') } | ForEach-Object { Stop-Process -Id $_.ProcessId -Force -ErrorAction SilentlyContinue; $_.ProcessId }"
], { windowsHide: true, timeout: 8000 })
if (String(stdout).trim()) killed.push('main.py')
} catch { /* ignore */ }
return killed
}
async function noteQueue(portNum) {
const queue = await fetchLocalQueue(portNum)
if (!queue.ok) return queue
if (queue.running > 0 || queue.pending > 0) {
markWork()
} else if (lastQueueRunning > 0 && queue.running === 0) {
lastWorkAt = Date.now()
}
lastQueueRunning = queue.running
return queue
}
async function maybeIdleStop(healthyPort) {
if (healthyPort) {
stoppedByAgent = false
const queue = await noteQueue(healthyPort)
if (!queue.ok) return
if (queue.running > 0 || queue.pending > 0) return
} else if (await processUp()) {
stoppedByAgent = false
} else {
stoppedByAgent = true
return
}
if (Date.now() - lastWorkAt < idleMs) return
console.log(JSON.stringify({
src: 'comfy-host-agent',
event: 'idle-stop',
idleMs,
lastActivityAt: new Date(lastWorkAt).toISOString(),
http: Boolean(healthyPort)
}))
const killed = await stopComfyProcesses()
markAsleep()
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'stopped', reason: 'idle', killed }))
}
function sharedRoots() {
const local = process.env.LOCALAPPDATA || ''
return {
input: process.env.COMFY_INPUT_DIR || join(local, 'Comfy-Desktop', 'ComfyUI-Shared', 'input'),
output: process.env.COMFY_OUTPUT_DIR || join(local, 'Comfy-Desktop', 'ComfyUI-Shared', 'output'),
temp: process.env.COMFY_TEMP_DIR || join(local, 'Comfy-Desktop', 'ComfyUI-Shared', 'temp')
}
}
function installLeafRoots(leaf) {
const local = process.env.LOCALAPPDATA || ''
const roots = []
const installs = join(local, 'Comfy-Desktop', 'ComfyUI-Installs')
try {
for (const dir of readdirSync(installs, { withFileTypes: true })) {
if (!dir.isDirectory()) continue
roots.push(join(installs, dir.name, 'ComfyUI', leaf))
}
} catch { /* ignore */ }
return roots
}
function inputRoots() {
const roots = [sharedRoots().input]
const cwd = process.env.COMFY_LAUNCH_CWD
if (cwd) roots.push(join(cwd, 'input'))
roots.push(...installLeafRoots('input'))
return [...new Set(roots)]
}
function outputRoots() {
const roots = [sharedRoots().output]
const cwd = process.env.COMFY_LAUNCH_CWD
if (cwd) {
roots.push(join(cwd, 'output'))
roots.push(join(cwd, 'ComfyUI', 'output'))
}
roots.push(...installLeafRoots('output'))
return [...new Set(roots)]
}
function tempRoots() {
const roots = [sharedRoots().temp]
const cwd = process.env.COMFY_LAUNCH_CWD
if (cwd) {
roots.push(join(cwd, 'temp'))
roots.push(join(cwd, 'ComfyUI', 'temp'))
}
roots.push(...installLeafRoots('temp'))
return [...new Set(roots)]
}
function rootsForType(type) {
if (type === 'input') return inputRoots()
if (type === 'temp') return [...tempRoots(), ...outputRoots()]
return [...outputRoots(), ...tempRoots()]
}
function safeFile(root, subfolder, filename) {
const name = basename(String(filename || ''))
if (!name || name === '.' || name === '..') return null
const rootAbs = resolve(root)
const target = resolve(rootAbs, String(subfolder || '').replace(/\.\./g, ''), name)
const rel = relative(rootAbs, target)
if (!rel || rel.startsWith('..') || isAbsolute(rel)) return null
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) {
if (!path || !existsSync(path)) return false
try {
rmSync(path, { force: true })
return true
} catch {
return false
}
}
function fileStem(filename) {
const name = basename(String(filename || ''))
const dot = name.lastIndexOf('.')
return dot > 0 ? name.slice(0, dot) : name
}
function removeNamedFile(roots, filename, subfolders, { stemSiblings = true } = {}) {
const deleted = []
const name = basename(String(filename || ''))
if (!name) return deleted
const stem = fileStem(name)
const subs = [...new Set((subfolders || []).map(item => String(item || '')))]
for (const root of roots) {
if (!existsSync(root)) continue
for (const sub of subs) {
const dir = sub ? join(root, String(sub).replace(/\.\./g, '')) : root
const exact = safeFile(root, sub, name)
if (removeFile(exact)) deleted.push(exact)
if (!stemSiblings || !stem || !existsSync(dir)) continue
let files = []
try { files = readdirSync(dir) } catch { continue }
for (const file of files) {
if (file === name || file.startsWith(`${stem}.`)) {
const path = join(dir, file)
if (removeFile(path)) deleted.push(path)
}
}
}
}
return deleted
}
function defaultSweepPrefixes() {
return ['MiniMax_H3', 'xAIGen', 'XAIgen', 'LTX23', 'AIGen', 'v2-krea', 'v2-generate']
}
function looksLikeStudioStill(name) {
// Studio uploads: 12-hex job prefix, or legacy aigen_<hex>_…
const n = String(name || '')
return /^[0-9a-f]{12}_.+\.(png|jpe?g|webp|gif)$/i.test(n)
|| /^aigen_[0-9a-f]+_.+\.(png|jpe?g|webp|gif)$/i.test(n)
}
function isStudioOutputName(file, sub, prefixes) {
if (looksLikeStudioStill(file) && (!sub || sub === 'image' || sub === 'still')) return true
if (sub === 'audio') {
return prefixes.some(prefix => file.startsWith(prefix)) || file.startsWith('ComfyUI')
}
if (!sub) {
return prefixes.some(prefix => prefix !== 'LTX23' && file.startsWith(prefix))
}
return prefixes.some(prefix => file.startsWith(prefix))
}
function sweepStaleStudioInputs(maxAgeMs = 15 * 60 * 1000) {
const deleted = []
const cutoff = Date.now() - Math.max(60_000, Number(maxAgeMs) || 900_000)
for (const root of inputRoots()) {
if (!existsSync(root)) continue
let files = []
try { files = readdirSync(root) } catch { continue }
for (const file of files) {
if (!looksLikeStudioStill(file)) continue
const path = join(root, file)
try {
const st = statSync(path)
if (st.isDirectory() || st.mtimeMs > cutoff) continue
} catch {
continue
}
if (removeFile(path)) deleted.push(path)
}
}
return deleted
}
function sweepStaleStudioOutputs(prefixes, subfolders, maxAgeMs = 15 * 60 * 1000) {
const deleted = []
const names = [...new Set([...defaultSweepPrefixes(), ...(prefixes || [])])].map(item => String(item || '').trim()).filter(Boolean)
const subs = [...new Set(['video', 'audio', 'image', 'still', ...(subfolders || [])])].map(item => String(item || ''))
const cutoff = Date.now() - Math.max(60_000, Number(maxAgeMs) || 900_000)
for (const root of outputRoots()) {
if (!existsSync(root)) continue
for (const sub of ['', ...subs]) {
const dir = sub ? join(root, sub) : root
if (!existsSync(dir)) continue
let files = []
try { files = readdirSync(dir) } catch { continue }
for (const file of files) {
if (!isStudioOutputName(file, sub, names)) continue
const path = join(dir, file)
try {
const st = statSync(path)
if (st.isDirectory() || st.mtimeMs > cutoff) continue
} catch {
continue
}
if (removeFile(path)) deleted.push(path)
}
}
}
return deleted
}
function readJson(req) {
return new Promise((resolve) => {
const chunks = []
req.on('data', (chunk) => chunks.push(chunk))
req.on('end', () => {
try {
resolve(JSON.parse(Buffer.concat(chunks).toString('utf8') || '{}'))
} catch {
resolve({})
}
})
req.on('error', () => resolve({}))
})
}
function purgeDesktopFiles(body) {
const deleted = []
const allowSweep = body?.sweep !== false
const names = [
String(body?.imageName || body?.filename || ''),
...(Array.isArray(body?.imageNames) ? body.imageNames : [])
].map(name => String(name || '').trim()).filter(Boolean)
const imageSub = String(body?.imageSubfolder || '')
for (const imageName of names) {
deleted.push(...removeNamedFile(inputRoots(), imageName, [imageSub, '']))
}
const video = body?.video || null
const videoName = String(video?.filename || '')
if (videoName) {
const type = String(video?.type || 'output')
const imageOut = /\.(png|jpe?g|webp|gif)$/i.test(videoName)
deleted.push(...removeNamedFile(
rootsForType(type),
videoName,
[video.subfolder, imageOut ? '' : 'video', ''],
{ stemSiblings: false }
))
}
const audio = body?.audio || null
const audioName = String(audio?.filename || '')
if (audioName) {
deleted.push(...removeNamedFile(
rootsForType(String(audio?.type || 'output')),
audioName,
[audio.subfolder, 'audio', ''],
{ stemSiblings: false }
))
}
const output = body?.output || null
const outputName = String(output?.filename || '')
if (outputName) {
deleted.push(...removeNamedFile(
rootsForType(String(output?.type || 'output')),
outputName,
[output.subfolder, 'image', 'still', ''],
{ stemSiblings: false }
))
}
if (allowSweep) {
deleted.push(...sweepStaleStudioInputs(body?.sweepMaxAgeMs))
deleted.push(...sweepStaleStudioOutputs(body?.sweepPrefixes, body?.sweepSubfolders, body?.sweepMaxAgeMs))
}
const unique = [...new Set(deleted)]
console.log(JSON.stringify({
src: 'comfy-host-agent',
event: 'purge',
deleted: unique.length,
files: unique.map((path) => basename(path)),
imageNames: names,
video: videoName || null,
audio: audioName || null,
output: outputName || null,
sweep: allowSweep
}))
return { ok: true, deleted: unique }
}
const server = http.createServer(async (req, res) => {
if (!authorized(req)) return json(res, 401, { ok: false, error: 'unauthorized' })
const url = new URL(req.url || '/', 'http://localhost')
if (req.method === 'GET' && url.pathname === '/status') {
// Answer immediately. MyMonitor allows ~500ms; probing ports + WMI here
// made the A light stay red even while this process was running.
// Do not treat a stale proxyTarget as proof Comfy is up — that made
// /status report asleep+http at once, so training skipped /stop.
const healthyPort = lastHealthyPort || 0
const processUpNow = Boolean(lastProcessUp && healthyPort)
const asleep = !healthyPort && !processUpNow
if (!asleep) stoppedByAgent = false
return json(res, 200, {
ok: true,
http: Boolean(healthyPort),
process: processUpNow,
processUp: processUpNow,
asleep,
lastActivityAt: new Date(lastWorkAt).toISOString(),
idleMs,
port: healthyPort || null,
proxyPort,
training: lastTraining
})
}
if (req.method === 'POST' && url.pathname === '/start') {
const training = await trainingLock()
if (training.busy) {
return json(res, 409, {
ok: false,
error: 'train-busy',
message: training.message,
training
})
}
const healthyPort = await syncProxy()
if (healthyPort) {
markWork()
return json(res, 200, { ok: true, started: false, already: true, asleep: false, port: healthyPort, proxyPort })
}
const listenPort = desktopListenPort()
// Do not trust lastProcessUp alone — force-stop can leave it sticky while Comfy is dead.
const processAlive = await processUp().catch(() => false) || await pythonMainUp().catch(() => false)
const portBusy = await portInUse(listenPort)
if (processAlive || portBusy) {
markWork()
lastProcessUp = true
if (portBusy) setProxyTarget(listenPort)
return json(res, 200, {
ok: true,
started: false,
already: true,
booting: true,
asleep: false,
port: portBusy ? listenPort : null,
proxyPort
})
}
// Only honor launch cooldown if we actually spawned recently; never block a real cold start.
if (lastLaunchAt && Date.now() - lastLaunchAt < 45_000 && await processUp().catch(() => false)) {
return json(res, 200, { ok: true, started: false, already: true, booting: true, asleep: false, proxyPort })
}
stoppedByAgent = false
probeSkipUntil.clear()
lastLaunchAt = Date.now()
const launched = await startComfy()
if (launched.port) {
lastProcessUp = true
setProxyTarget(launched.port)
}
return json(res, 200, { ok: true, asleep: false, proxyPort, ...launched })
}
if (req.method === 'POST' && url.pathname === '/stop') {
const body = await readJson(req).catch(() => ({}))
const force = body?.force === true || url.searchParams.get('force') === '1'
const healthyPort = await syncProxy()
const queue = healthyPort ? await fetchLocalQueue(healthyPort) : { running: 0, pending: 0, ok: false }
if (!force && queue.ok && (queue.running > 0 || queue.pending > 0)) {
return json(res, 409, {
ok: false,
error: 'queue-busy',
running: queue.running,
pending: queue.pending
})
}
if (force && healthyPort) {
try {
await fetch(`http://127.0.0.1:${healthyPort}/interrupt`, { method: 'POST', signal: AbortSignal.timeout(4000) })
} catch { /* ignore */ }
}
const killed = await stopComfyProcesses()
markAsleep()
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'stopped', reason: force ? 'force-stop' : 'stop', killed, queue }))
return json(res, 200, { ok: true, stopped: true, asleep: true, killed, forced: force })
}
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') {
const body = await readJson(req)
return json(res, 200, purgeDesktopFiles(body))
}
if (req.method === 'POST' && url.pathname === '/sweep') {
const body = await readJson(req)
const deleted = [
...sweepStaleStudioInputs(body?.sweepMaxAgeMs),
...sweepStaleStudioOutputs(body?.sweepPrefixes, body?.sweepSubfolders, body?.sweepMaxAgeMs)
]
console.log(JSON.stringify({
src: 'comfy-host-agent',
event: 'sweep',
deleted: deleted.length,
files: deleted.map((path) => basename(path))
}))
return json(res, 200, { ok: true, deleted })
}
json(res, 404, { ok: false, error: 'not found' })
})
function logStartupSweep() {
// Disabled: sweeping studio outputs on agent start was deleting Comfy
// files that are still part of the generation / recover workflow.
}
server.listen(port, '0.0.0.0', async () => {
ensureProxyListening()
lastHealthyPort = await syncProxy() || 0
lastProcessUp = Boolean(lastHealthyPort) || await processUp().catch(() => false)
if (lastHealthyPort) await noteQueue(lastHealthyPort)
logStartupSweep()
console.log(JSON.stringify({
src: 'comfy-host-agent',
event: 'listen',
port,
proxyPort,
comfyPort: lastHealthyPort || null,
idleMs,
candidates: candidatePorts()
}))
void trainingLock()
setInterval(() => { void trainingLock() }, 3000)
let lastIdleCheck = 0
setInterval(() => {
syncProxy().then(async (nextPort) => {
lastHealthyPort = nextPort || 0
const now = Date.now()
if (now - lastIdleCheck < 10_000) return
lastIdleCheck = now
await maybeIdleStop(nextPort)
lastProcessUp = Boolean(nextPort) || await processUp().catch(() => false)
}).catch(() => {})
}, 3000)
})