import { createUpscaleHost } from './upscale-host.mjs' import { purgeStudio2Files } from './studio2-purge.mjs' import { stableMemoryArgs } from './comfy-memory-policy.mjs' import { createGpuReservation } from './gpu-reservation.mjs' import { createGpuProxy } from './gpu-proxy.mjs' import { createYueGpHost } from './yuegp-host.mjs' import { createYue2Host } from './yue2-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() || yue2.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)) const cliPath = join(inst.installPath, 'ComfyUI', 'comfy', 'cli_args.py') const memoryArgs = stableMemoryArgs(args, { cliSource: existsSync(cliPath) ? readFileSync(cliPath, 'utf8') : '' }) args.splice(0, args.length, ...memoryArgs) 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__… 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() || yue2.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.') if (yue2.busy()) throw new Error('YuE2 is using the GPU.') 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 yue2 = createYue2Host({ leaseValid: lease => gpuReservation.isOwner(lease), prepare: async () => { if ((await trainingLock()).busy) throw new Error('GPU is busy with training.') if (yueGp.busy()) throw new Error('YuEGP is using the GPU.') const healthy = await syncProxy() if (healthy) { const queue = await fetchLocalQueue(healthy) if (!queue.ok || queue.running || queue.pending) throw new Error('Comfy is busy; YuE2 cannot start.') } // Same gate as YuEGP: stop Comfy and refuse to launch while its Python still owns VRAM. if (healthy || await processUp() || await pythonMainUp().catch(() => false)) { 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 (url.pathname.startsWith('/yue2/')) { if (yueGp.busy()) return json(res, 409, { message: 'YuEGP is using the GPU.' }) const match = url.pathname.match(/^\/yue2\/jobs\/([a-zA-Z0-9-]{12,80})(\/audio|\/cancel)?$/) if (req.method === 'GET' && url.pathname === '/yue2/status') return json(res, 200, { configured: yue2.configured(), busy: yue2.busy(), backend: 'yue2' }) if (req.method === 'POST' && url.pathname === '/yue2/jobs') return json(res, 200, await yue2.start(await readJson(req), String(req.headers['x-aigen-gpu-lease'] || ''))) if (match && req.method === 'POST' && match[2] === '/cancel') return json(res, 200, await yue2.cancel(match[1])) if (match && req.method === 'GET' && match[2] === '/audio') { const path = yue2.audio(match[1]) return path ? streamFile(res, path) : json(res, 404, { error: 'Audio not ready' }) } if (match && req.method === 'GET' && !match[2]) { const job = yue2.read(match[1]) return json(res, job ? 200 : 404, job || { error: 'YuE2 job not found' }) } return json(res, 404, { error: 'Unknown YuE2 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() }, yue2: { busy: yue2.busy(), configured: yue2.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.' }) if (yue2.busy()) return json(res, 409, { message: 'YuE2 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 === '/studio2/purge') { if (!gpuReservation.isOwner(String(req.headers['x-aigen-gpu-lease'] || ''))) return json(res, 409, { ok: false, error: 'GPU lease required' }) const body = await readJson(req) try { return json(res, 200, purgeStudio2Files(body, rootsForType)) } catch (error) { return json(res, 400, { ok: false, error: String(error.message || error) }) } } 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('/yue2/') || String(req.url || '').startsWith('/upscale/')) && !res.headersSent) return json(res, error.statusCode || 400, { error: error.message || 'Music host 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) })