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

451 lines
15 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 } from 'node:fs'
import { basename, join, resolve, relative, isAbsolute } from 'node:path'
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_HOST || '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)
let lastWorkAt = Date.now()
let lastQueueRunning = 0
let stoppedByAgent = false
function markWork() {
lastWorkAt = Date.now()
stoppedByAgent = false
}
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
}
async function probeStats(portNum) {
try {
const res = await fetch(`http://127.0.0.1:${portNum}/system_stats`, { signal: AbortSignal.timeout(2500) })
if (!res.ok) return null
const stats = await res.json()
if (!stats?.system) return null
return stats
} catch {
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 found = []
for (const next of candidatePorts()) {
const stats = await probeStats(next)
if (!stats) continue
found.push({ port: next, advertised: advertisedPort(stats) })
}
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 ensureProxy(targetPort) {
if (!targetPort) return
if (proxyTarget !== targetPort) {
proxyTarget = targetPort
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'proxy-target', listen: proxyPort, target: proxyTarget }))
}
if (proxyServer) return
proxyServer = net.createServer((client) => {
const upstream = net.connect(proxyTarget, '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 }))
})
}
async function syncProxy() {
const healthyPort = await findHealthyPort()
if (healthyPort) ensureProxy(healthyPort)
return healthyPort
}
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',
"Get-CimInstance Win32_Process | Where-Object { $_.CommandLine -match 'ComfyUI\\\\main.py|ComfyUI/main.py' } | Select-Object -First 1 -ExpandProperty ProcessId"
], { windowsHide: true, timeout: 6000 })
return Boolean(String(stdout).trim())
} catch {
return false
}
}
function launchCommand() {
if (process.env.COMFY_LAUNCH_CMD) return process.env.COMFY_LAUNCH_CMD
const local = process.env.LOCALAPPDATA || ''
const desktop = `${local}\\Programs\\Comfy Desktop\\Comfy Desktop.exe`
return `start "" "${desktop}"`
}
function startComfy() {
markWork()
const command = launchCommand()
const cwd = process.env.COMFY_LAUNCH_CWD || undefined
const child = spawn(command, { shell: true, detached: true, stdio: 'ignore', windowsHide: true, cwd })
child.unref()
return { started: true, command, cwd: 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 { $_.CommandLine -match 'ComfyUI\\\\main.py|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()
stoppedByAgent = true
lastQueueRunning = 0
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')
}
}
function inputRoots() {
const local = process.env.LOCALAPPDATA || ''
const roots = [
process.env.COMFY_INPUT_DIR || join(local, 'Comfy-Desktop', 'ComfyUI-Shared', 'input')
]
const cwd = process.env.COMFY_LAUNCH_CWD
if (cwd) roots.push(join(cwd, 'input'))
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', 'input'))
}
} catch { /* ignore */ }
return [...new Set(roots)]
}
function outputRoots() {
const roots = [sharedRoots().output]
const cwd = process.env.COMFY_LAUNCH_CWD
if (cwd) roots.push(join(cwd, 'output'))
return [...new Set(roots)]
}
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 removeFile(path) {
if (!path || !existsSync(path)) return false
rmSync(path, { force: true })
return true
}
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 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) {
for (const root of inputRoots()) {
if (!existsSync(root)) continue
const nested = safeFile(root, imageSub, imageName)
if (removeFile(nested)) deleted.push(nested)
if (imageSub) {
const flat = safeFile(root, '', imageName)
if (removeFile(flat)) deleted.push(flat)
}
}
}
const video = body?.video || null
const videoName = String(video?.filename || '')
if (videoName) {
for (const root of outputRoots()) {
const nested = safeFile(root, video.subfolder || 'video', videoName)
if (removeFile(nested)) deleted.push(nested)
const flat = safeFile(root, '', videoName)
if (removeFile(flat)) deleted.push(flat)
}
}
const output = body?.output || null
const outputName = String(output?.filename || '')
if (outputName) {
const type = String(output?.type || 'output')
const sub = String(output?.subfolder || '')
const rootsForType = type === 'input' ? inputRoots() : outputRoots()
for (const root of rootsForType) {
if (!existsSync(root)) continue
const nested = safeFile(root, sub, outputName)
if (removeFile(nested)) deleted.push(nested)
if (sub) {
const flat = safeFile(root, '', outputName)
if (removeFile(flat)) deleted.push(flat)
}
}
}
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'purge', deleted: deleted.length, files: deleted.map((path) => basename(path)), imageNames: names }))
return { ok: true, deleted }
}
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') {
const [healthyPort, process] = await Promise.all([syncProxy(), processUp()])
if (healthyPort) await noteQueue(healthyPort)
const processUpNow = Boolean(process || healthyPort)
return json(res, 200, {
ok: true,
http: Boolean(healthyPort),
process: processUpNow,
processUp: processUpNow,
asleep: Boolean(stoppedByAgent || (!healthyPort && !process)),
lastActivityAt: new Date(lastWorkAt).toISOString(),
idleMs,
port: healthyPort || null,
proxyPort: healthyPort ? proxyPort : null
})
}
if (req.method === 'POST' && url.pathname === '/start') {
const healthyPort = await syncProxy()
if (healthyPort) {
markWork()
return json(res, 200, { ok: true, started: false, already: true, asleep: false, port: healthyPort, proxyPort })
}
const launched = startComfy()
return json(res, 200, { ok: true, asleep: false, ...launched })
}
if (req.method === 'POST' && url.pathname === '/stop') {
const healthyPort = await syncProxy()
const queue = healthyPort ? await fetchLocalQueue(healthyPort) : { running: 0, pending: 0, ok: false }
if (queue.ok && (queue.running > 0 || queue.pending > 0)) {
return json(res, 409, {
ok: false,
error: 'queue-busy',
running: queue.running,
pending: queue.pending
})
}
const killed = await stopComfyProcesses()
stoppedByAgent = true
lastQueueRunning = 0
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'stopped', reason: 'stop', killed }))
return json(res, 200, { ok: true, stopped: true, asleep: true, killed })
}
if (req.method === 'POST' && url.pathname === '/purge') {
const body = await readJson(req)
return json(res, 200, purgeDesktopFiles(body))
}
json(res, 404, { ok: false, error: 'not found' })
})
server.listen(port, '0.0.0.0', async () => {
const healthyPort = await syncProxy()
if (healthyPort) await noteQueue(healthyPort)
console.log(JSON.stringify({
src: 'comfy-host-agent',
event: 'listen',
port,
proxyPort,
comfyPort: healthyPort || null,
idleMs,
candidates: candidatePorts()
}))
let lastIdleCheck = 0
setInterval(() => {
syncProxy().then(async (nextPort) => {
const now = Date.now()
if (now - lastIdleCheck < 10_000) return
lastIdleCheck = now
await maybeIdleStop(nextPort)
}).catch(() => {})
}, 3000)
})