Files
aigen/scripts/comfy-host-agent.mjs
T

1044 lines
38 KiB
JavaScript

import { createUpscaleHost } from './upscale-host.mjs'
import { createGpuReservation } from './gpu-reservation.mjs'
import { createGpuProxy } from './gpu-proxy.mjs'
import { createYueGpHost } from './yuegp-host.mjs'
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 = createGpuProxy({ target: () => proxyTarget, reservation: gpuReservation, authorized, markWork, externalBusy: () => yueGp.busy() || upscale.busy() })
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 (gpuReservation.availability().busy) return
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 = false // Never sweep another studio's files during a job cleanup.
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 gpuReservation = createGpuReservation({ idle: async () => {
if (yueGp.busy() || upscale.busy()) return false
if ((await trainingLock()).busy) return false
const healthy = await syncProxy()
if (healthy) {
const queue = await fetchLocalQueue(healthy)
return queue.ok && queue.running === 0 && queue.pending === 0
}
// A stopped service may be reserved before waking; a silent live process may not.
return !(await processUp()) && !(await pythonMainUp().catch(() => true))
} })
const yueGp = createYueGpHost({
leaseValid: lease => gpuReservation.isOwner(lease),
prepare: async () => {
if ((await trainingLock()).busy) throw new Error('GPU is busy with training.')
const healthy = await syncProxy()
if (healthy) {
const queue = await fetchLocalQueue(healthy)
if (!queue.ok || queue.running || queue.pending) throw new Error('Comfy is busy; YuEGP cannot start.')
await stopComfyProcesses()
markAsleep()
}
if (await pythonMainUp()) throw new Error('Comfy has not stopped; retry after the GPU is free.')
}
})
const upscale = createUpscaleHost({ leaseValid: token => gpuReservation.isOwner(token) })
async function handleControl(req, res) {
if (!authorized(req)) return json(res, 401, { ok: false, error: 'unauthorized' })
const url = new URL(req.url || '/', 'http://localhost')
if (url.pathname.startsWith('/upscale/')) {
const match = url.pathname.match(/^\/upscale\/jobs\/([a-zA-Z0-9-]{12,80})(\/input|\/video|\/cancel)?$/)
if (req.method === 'POST' && url.pathname === '/upscale/jobs') return json(res, 200, upscale.start(await readJson(req), String(req.headers['x-aigen-gpu-lease'] || '')))
if (match && req.method === 'PUT' && match[2] === '/input') { await upscale.upload(match[1], req); return json(res, 200, { uploaded: true }) }
if (match && req.method === 'POST' && match[2] === '/cancel') return json(res, 200, await upscale.cancel(match[1]))
if (match && req.method === 'GET' && match[2] === '/video') { const path = upscale.output(match[1]); return path ? streamFile(res, path) : json(res, 404, { error: 'Video not ready' }) }
if (match && req.method === 'GET' && !match[2]) { const job = upscale.read(match[1]); return json(res, job ? 200 : 404, job || { error: 'Upscale job not found' }) }
return json(res, 404, { error: 'Unknown upscale endpoint' })
}
if (url.pathname.startsWith('/yuegp/')) {
const match = url.pathname.match(/^\/yuegp\/jobs\/([a-zA-Z0-9-]{12,80})(\/audio|\/cancel)?$/)
if (req.method === 'GET' && url.pathname === '/yuegp/status') return json(res, 200, { configured: yueGp.configured(), busy: yueGp.busy(), backend: 'yuegp' })
if (req.method === 'POST' && url.pathname === '/yuegp/jobs') return json(res, 200, await yueGp.start(await readJson(req), String(req.headers['x-aigen-gpu-lease'] || '')))
if (match && req.method === 'POST' && match[2] === '/cancel') return json(res, 200, await yueGp.cancel(match[1]))
if (match && req.method === 'GET' && match[2] === '/audio') {
const path = yueGp.audio(match[1])
return path ? streamFile(res, path) : json(res, 404, { error: 'Audio not ready' })
}
if (match && req.method === 'GET' && !match[2]) {
const job = yueGp.read(match[1])
return json(res, job ? 200 : 404, job || { error: 'YuEGP job not found' })
}
return json(res, 404, { error: 'Unknown YuEGP endpoint' })
}
if (req.method === 'POST' && url.pathname.startsWith('/gpu/')) {
const body = await readJson(req)
let result
if (url.pathname === '/gpu/acquire') result = await gpuReservation.acquire(body.ticket)
else if (url.pathname === '/gpu/renew') result = await gpuReservation.renew(body.token)
else if (url.pathname === '/gpu/release') result = await gpuReservation.release(body.token)
else if (url.pathname === '/gpu/cancel') result = await gpuReservation.cancel(body.ticket)
else return json(res, 404, { ok: false })
return json(res, 200, result)
}
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,
gpu: gpuReservation.availability(),
training: { busy: lastTraining.busy },
yuegp: { busy: yueGp.busy(), configured: yueGp.configured() },
upscale: { busy: upscale.busy(), engine: 'realesrgan-rife', local: true }
})
}
if (req.method === 'POST' && url.pathname === '/start') {
if (upscale.busy()) return json(res, 409, { message: 'Local upscale is using the GPU.' })
if (yueGp.busy()) return json(res, 409, { message: 'YuEGP is using the GPU.' })
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') return json(res, 409, { ok: false, message: 'Global cleanup is disabled while studios share the GPU.' })
json(res, 404, { ok: false, error: 'not found' })
}
const server = http.createServer(async (req, res) => {
if (!authorized(req)) return json(res, 401, { ok: false, error: 'unauthorized' })
try {
const path = new URL(req.url || '/', 'http://localhost').pathname
if (['POST', 'PUT'].includes(req.method) && !path.startsWith('/gpu/')) {
await gpuReservation.permit(String(req.headers['x-aigen-gpu-lease'] || ''), () => handleControl(req, res))
} else await handleControl(req, res)
} catch (error) {
req.resume()
if ((String(req.url || '').startsWith('/yuegp/') || String(req.url || '').startsWith('/upscale/')) && !res.headersSent) return json(res, error.statusCode || 400, { error: error.message || 'YuEGP request failed' })
if (!res.headersSent) json(res, error.statusCode || 400, { ok: false, message: error.statusCode === 409 ? 'GPU is in use. Waiting for availability.' : 'GPU coordination request failed.' })
}
})
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)
})