Coordinate shared GPU access across studio instances
This commit is contained in:
@@ -1,3 +1,5 @@
|
||||
import { createGpuReservation } from './gpu-reservation.mjs'
|
||||
import { createGpuProxy } from './gpu-proxy.mjs'
|
||||
import http from 'node:http'
|
||||
import net from 'node:net'
|
||||
import { execFile, spawn } from 'node:child_process'
|
||||
@@ -141,26 +143,7 @@ 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 = createGpuProxy({ target: () => proxyTarget, reservation: gpuReservation, authorized, markWork })
|
||||
proxyServer.on('error', (error) => {
|
||||
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'proxy-error', error: String(error.message || error) }))
|
||||
})
|
||||
@@ -508,6 +491,7 @@ async function noteQueue(portNum) {
|
||||
}
|
||||
|
||||
async function maybeIdleStop(healthyPort) {
|
||||
if (gpuReservation.availability().busy) return
|
||||
if (healthyPort) {
|
||||
stoppedByAgent = false
|
||||
const queue = await noteQueue(healthyPort)
|
||||
@@ -767,7 +751,7 @@ function readJson(req) {
|
||||
|
||||
function purgeDesktopFiles(body) {
|
||||
const deleted = []
|
||||
const allowSweep = body?.sweep !== false
|
||||
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 : [])
|
||||
@@ -827,9 +811,30 @@ function purgeDesktopFiles(body) {
|
||||
return { ok: true, deleted: unique }
|
||||
}
|
||||
|
||||
const server = http.createServer(async (req, res) => {
|
||||
const gpuReservation = createGpuReservation({ idle: async () => {
|
||||
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))
|
||||
} })
|
||||
|
||||
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 (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.
|
||||
@@ -849,7 +854,8 @@ const server = http.createServer(async (req, res) => {
|
||||
idleMs,
|
||||
port: healthyPort || null,
|
||||
proxyPort,
|
||||
training: lastTraining
|
||||
gpu: gpuReservation.availability(),
|
||||
training: { busy: lastTraining.busy }
|
||||
})
|
||||
}
|
||||
if (req.method === 'POST' && url.pathname === '/start') {
|
||||
@@ -935,21 +941,21 @@ const server = http.createServer(async (req, res) => {
|
||||
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 })
|
||||
}
|
||||
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 (req.method === 'POST' && !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 (!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() {
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
import http from 'node:http'
|
||||
import net from 'node:net'
|
||||
import { GPU_WAIT_MESSAGE } from './gpu-reservation.mjs'
|
||||
|
||||
function reply(res, status, body) {
|
||||
if (res.headersSent) return res.destroy()
|
||||
res.writeHead(status, { 'Content-Type': 'application/json', 'Cache-Control': 'no-store' })
|
||||
res.end(JSON.stringify(body))
|
||||
}
|
||||
|
||||
/** Streams Comfy HTTP/WebSocket traffic; all mutations require the current reservation. */
|
||||
export function createGpuProxy({ target, reservation, authorized = () => true, markWork = () => {} }) {
|
||||
const server = http.createServer(async (req, res) => {
|
||||
const path = new URL(req.url || '/', 'http://localhost').pathname.replace(/\/+$/, '') || '/'
|
||||
// A global history list can expose another studio's prompts and outputs.
|
||||
if (req.method === 'GET' && (path === '/history' || path === '/history/')) return reply(res, 200, {})
|
||||
const forward = () => new Promise(resolve => {
|
||||
const port = target()
|
||||
if (!port) { reply(res, 503, { error: 'GPU service is unavailable.' }); return resolve() }
|
||||
const headers = { ...req.headers }
|
||||
delete headers['x-aigen-gpu-lease']
|
||||
delete headers.authorization
|
||||
const upstream = http.request({ hostname: '127.0.0.1', port, path: req.url, method: req.method, headers }, response => {
|
||||
if (req.method === 'GET' && path === '/queue') {
|
||||
let size = 0
|
||||
const chunks = []
|
||||
response.on('data', chunk => {
|
||||
size += chunk.length
|
||||
if (size > 16 * 1024 * 1024) { response.destroy(); reply(res, 502, { error: 'GPU status unavailable.' }); resolve(); return }
|
||||
chunks.push(chunk)
|
||||
})
|
||||
response.on('end', () => {
|
||||
try {
|
||||
if (response.statusCode !== 200) throw new Error('status unavailable')
|
||||
const queue = JSON.parse(Buffer.concat(chunks).toString('utf8'))
|
||||
if (!Array.isArray(queue.queue_running) || !Array.isArray(queue.queue_pending)) throw new Error('invalid queue')
|
||||
// Preserve count compatibility without revealing IDs, graphs, or prompts.
|
||||
reply(res, 200, { queue_running: queue.queue_running.map(() => null), queue_pending: queue.queue_pending.map(() => null) })
|
||||
} catch { reply(res, 502, { error: 'GPU status unavailable.' }) }
|
||||
resolve()
|
||||
})
|
||||
} else {
|
||||
res.writeHead(response.statusCode || 502, response.headers)
|
||||
response.pipe(res)
|
||||
response.on('end', resolve)
|
||||
}
|
||||
response.on('error', () => { res.destroy(); resolve() })
|
||||
})
|
||||
upstream.on('error', () => { reply(res, 502, { error: 'GPU service is unavailable.' }); resolve() })
|
||||
upstream.setTimeout(5 * 60 * 1000, () => upstream.destroy())
|
||||
req.on('aborted', () => upstream.destroy())
|
||||
res.on('close', () => { upstream.destroy(); resolve() })
|
||||
req.pipe(upstream)
|
||||
})
|
||||
try {
|
||||
if (req.method === 'GET' || req.method === 'HEAD') await forward()
|
||||
else {
|
||||
if (!authorized(req)) { req.resume(); return reply(res, 401, { error: 'Unauthorized' }) }
|
||||
await reservation.permit(String(req.headers['x-aigen-gpu-lease'] || ''), async () => { markWork(); await forward() })
|
||||
}
|
||||
} catch {
|
||||
req.resume()
|
||||
reply(res, 409, { error: { message: GPU_WAIT_MESSAGE }, code: 'GPU_BUSY' })
|
||||
}
|
||||
})
|
||||
server.on('upgrade', (req, socket, head) => {
|
||||
if (new URL(req.url || '/', 'http://localhost').pathname !== '/ws' || !target()) return socket.destroy()
|
||||
const upstream = net.connect(target(), '127.0.0.1', () => {
|
||||
const headers = Object.entries(req.headers).filter(([key]) => key !== 'authorization' && key !== 'x-aigen-gpu-lease')
|
||||
.map(([key, value]) => `${key}: ${value}`).join('\r\n')
|
||||
upstream.write(`${req.method} ${req.url} HTTP/${req.httpVersion}\r\n${headers}\r\n\r\n`)
|
||||
if (head.length) upstream.write(head)
|
||||
socket.pipe(upstream); upstream.pipe(socket)
|
||||
})
|
||||
socket.on('error', () => upstream.destroy())
|
||||
upstream.on('error', () => socket.destroy())
|
||||
socket.on('close', () => upstream.destroy())
|
||||
upstream.on('close', () => socket.destroy())
|
||||
})
|
||||
return server
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
import { randomUUID } from 'node:crypto'
|
||||
|
||||
export const GPU_WAIT_MESSAGE = 'GPU is in use. Waiting for availability.'
|
||||
|
||||
/** Device-local arbitration. Tickets/tokens are opaque; no studio or job metadata is stored. */
|
||||
export function createGpuReservation({ idle, now = Date.now, ttlMs = 60_000, ticketTtlMs = 30_000 } = {}) {
|
||||
if (typeof idle !== 'function') throw new Error('An authoritative GPU idle probe is required')
|
||||
let owner = null
|
||||
let gate = Promise.resolve()
|
||||
let inFlight = 0
|
||||
const waiting = new Map()
|
||||
const serialize = fn => {
|
||||
const run = gate.then(fn)
|
||||
gate = run.catch(() => {})
|
||||
return run
|
||||
}
|
||||
const valid = token => Boolean(owner && !owner.released && owner.token === token && owner.until > now())
|
||||
const busy = () => ({ acquired: false, message: GPU_WAIT_MESSAGE, retryAfterMs: 2500 })
|
||||
const prune = () => {
|
||||
for (const [ticket, seen] of waiting) if (now() - seen >= ticketTtlMs) waiting.delete(ticket)
|
||||
}
|
||||
return {
|
||||
acquire(ticket) {
|
||||
return serialize(async () => {
|
||||
if (typeof ticket !== 'string' || !/^[a-zA-Z0-9-]{16,80}$/.test(ticket)) throw new Error('Invalid reservation ticket')
|
||||
prune()
|
||||
if (owner?.ticket === ticket && valid(owner.token)) {
|
||||
owner.until = now() + ttlMs
|
||||
return { acquired: true, token: owner.token, ttlMs }
|
||||
}
|
||||
waiting.set(ticket, now())
|
||||
if (owner && (valid(owner.token) || inFlight > 0)) return busy()
|
||||
if (waiting.keys().next().value !== ticket) return busy()
|
||||
// Unknown/offline is not idle; the adapter must explicitly account for a stopped GPU.
|
||||
if (await idle() !== true) return busy()
|
||||
owner = { ticket, token: randomUUID(), until: now() + ttlMs, released: false }
|
||||
waiting.delete(ticket)
|
||||
return { acquired: true, token: owner.token, ttlMs }
|
||||
})
|
||||
},
|
||||
renew(token) {
|
||||
return serialize(() => {
|
||||
if (!valid(token)) return { renewed: false }
|
||||
owner.until = now() + ttlMs
|
||||
return { renewed: true, ttlMs }
|
||||
})
|
||||
},
|
||||
release(token) {
|
||||
return serialize(() => {
|
||||
if (!owner || owner.token !== token) return { released: false }
|
||||
owner.released = true
|
||||
// The next acquire still verifies that Comfy has drained before granting access.
|
||||
return { released: true }
|
||||
})
|
||||
},
|
||||
cancel(ticket) { return serialize(() => { waiting.delete(ticket); return { cancelled: true } }) },
|
||||
async permit(token, operation) {
|
||||
const accepted = await serialize(() => {
|
||||
if (!valid(token)) return false
|
||||
inFlight += 1
|
||||
return true
|
||||
})
|
||||
if (!accepted) throw Object.assign(new Error(GPU_WAIT_MESSAGE), { statusCode: 409 })
|
||||
try { return await operation() } finally { inFlight -= 1 }
|
||||
},
|
||||
availability() { return { busy: Boolean(owner && (valid(owner.token) || inFlight > 0)) } }
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user