diff --git a/docs/shared-gpu-rollout.md b/docs/shared-gpu-rollout.md new file mode 100644 index 0000000..5e529d5 --- /dev/null +++ b/docs/shared-gpu-rollout.md @@ -0,0 +1,20 @@ +# Shared GPU coordination rollout + +Both Coolify applications must use the same COMFY_CONTROL_URL (the 5080 host agent on port 8199), COMFY_CONTROL_TOKEN, and COMFY_HOST through its proxy on port 8198. The retired 192.168.77.101 sidecar is not involved. + +The agent now requires an opaque reservation for Comfy mutations and host start/stop/purge. Each application process creates a random ticket. No studio name, user identity, job name, prompt, or media is included in reservation requests. A waiting instance receives availability only. The host serializes acquisition, renews leases for active work, respects waiting tickets, and confirms the physical Comfy queue has drained before handing off an expired/released reservation. + +The HTTP proxy strips graphs and IDs from queue status and disables global history enumeration. Exact prompt history remains available for known job IDs. Automatic cleanup deletes named files only; global sweeps are disabled. This coordinates trusted instances sharing one GPU; it is not a separate-account security boundary around Comfy's local files or direct native port. + +Deployment order: + +1. Let active renders and output saves finish in both sites. Do not use Force reset from one site to clear the other. +2. Update/restart the local host agent with all three adjacent files: comfy-host-agent.mjs, gpu-reservation.mjs, gpu-proxy.mjs. The launcher runs the files directly from this repository's scripts folder. Copying only comfy-host-agent.mjs is no longer sufficient. +3. Deploy the same commit to AIGen and xAIGen. Updating the host before the sites temporarily blocks older clients; updating a site first leaves its jobs waiting until the new coordinator is available. There is deliberately no uncoordinated fallback. +4. Verify both sites report available, then queue a render on one and a render on the other. The second should remain waiting with a generic availability message, then start after the first saves. Repeat with sites reversed and test that Force reset on the waiting site is rejected without clearing its jobs. + +Pushing to Gitea alone does not restart the Windows host-agent process. Coolify's website containers and the host agent are separate processes. Do not open direct Comfy port 8188 to either app as a workaround; both must submit via 8198 for fencing to apply. + +Recovery: clients renew every 10 seconds; host leases last 60 seconds. If a site disappears, its waiting ticket expires and its lease may be reassigned only after Comfy confirms idle. A lost connection never counts as an idle GPU. Known running prompt recovery remains in the original site. Global-history fallback recovery is intentionally unavailable because it could import the other site's results. + +Validation: node --test tests/studio-queue.test.mjs tests/gpu-reservation.test.mjs tests/shared-gpu.test.mjs. The HTTP integration test uses two independent application-client modules, the real reservation coordinator and proxy, and a simulated Comfy server. It does not start Comfy, render media, or alter library data. diff --git a/pages/queue.vue b/pages/queue.vue index 412bb56..91b78f4 100644 --- a/pages/queue.vue +++ b/pages/queue.vue @@ -81,6 +81,7 @@

{{ jobLine(job) }}

{{ liveMessage[job.shots.id] }}

{{ liveMessage[job.id] }}

+

{{ job.waitReason || 'In queue. Waiting for GPU availability.' }}

{{ job.lastError }}

@@ -501,6 +502,7 @@ type Queue = { name: string status: string lastError?: string + waitReason?: string pendingCount: number completedCount: number totalCount: number @@ -528,6 +530,7 @@ type StudioJobRow = { pauseAfterCurrent?: boolean pausedByUser?: boolean lastError?: string + waitReason?: string shots?: Queue | null plannedShots?: { prompt: string; duration: number; loraName?: string; loraStack?: LoraStackItem[] }[] completedCount?: number @@ -556,6 +559,7 @@ type GenerationLogEntry = { prompt: string status: string lastError?: string + waitReason?: string payload: Record imagePipeline?: string } diff --git a/scripts/comfy-host-agent.mjs b/scripts/comfy-host-agent.mjs index 02495ba..6bbfdc0 100644 --- a/scripts/comfy-host-agent.mjs +++ b/scripts/comfy-host-agent.mjs @@ -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() { diff --git a/scripts/gpu-proxy.mjs b/scripts/gpu-proxy.mjs new file mode 100644 index 0000000..b03275c --- /dev/null +++ b/scripts/gpu-proxy.mjs @@ -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 +} diff --git a/scripts/gpu-reservation.mjs b/scripts/gpu-reservation.mjs new file mode 100644 index 0000000..07cb497 --- /dev/null +++ b/scripts/gpu-reservation.mjs @@ -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)) } } + } +} diff --git a/server/api/comfy/force-reset.post.ts b/server/api/comfy/force-reset.post.ts index 5c0c828..e29aac5 100644 --- a/server/api/comfy/force-reset.post.ts +++ b/server/api/comfy/force-reset.post.ts @@ -1,15 +1,18 @@ +import { withSharedGpuReset } from '~/server/utils/sharedGpu' import { libraryOwnerKey } from '~/server/utils/library' import { requestComfyForceReset } from '~/server/utils/comfyLifecycle' import { clearStuckStudioWork } from '~/server/utils/studioQueue' export default defineEventHandler(async (event) => { const owner = libraryOwnerKey(event) - const cleared = await clearStuckStudioWork(owner) - const comfy = await requestComfyForceReset() - return { - ok: true, - message: 'Cleared every stuck Studio/music/shot job and force-stopped Comfy. Use Wake / Poke when you want the GPU again.', - cleared, - comfy - } + return withSharedGpuReset(async () => { + const cleared = await clearStuckStudioWork(owner) + const comfy = await requestComfyForceReset() + return { + ok: true, + message: 'Cleared every stuck Studio/music/shot job and force-stopped Comfy. Use Wake / Poke when you want the GPU again.', + cleared, + comfy + } + }) }) diff --git a/server/plugins/shared-gpu.ts b/server/plugins/shared-gpu.ts new file mode 100644 index 0000000..4f7d6e8 --- /dev/null +++ b/server/plugins/shared-gpu.ts @@ -0,0 +1,17 @@ +import { hasActivePromptJobs } from '~/server/utils/promptComfy' +import { listJobs } from '~/server/utils/jobs' +import { registerSharedGpuWork, maintainSharedGpu } from '~/server/utils/sharedGpu' + +export default defineNitroPlugin((nitro) => { + registerSharedGpuWork(() => hasActivePromptJobs() || listJobs().some(job => + job.saving || job.library?.chainContinuing || ['queued', 'uploading', 'running'].includes(job.status) + )) + let running = false + const timer = setInterval(async () => { + if (running) return + running = true + try { await maintainSharedGpu() } finally { running = false } + }, 10_000) + timer.unref() + nitro.hooks.hook('close', () => { clearInterval(timer) }) +}) diff --git a/server/utils/comfy.ts b/server/utils/comfy.ts index f66f6f0..93c9488 100644 --- a/server/utils/comfy.ts +++ b/server/utils/comfy.ts @@ -1,3 +1,4 @@ +import { assertSharedGpu, sharedGpuHeaders } from '~/server/utils/sharedGpu' import { LTX_NEGATIVE } from '~/utils/videoModels' import { createHash } from 'node:crypto' import { comfyJobPrefix } from '~/utils/outputNames' @@ -73,6 +74,9 @@ function configuredComfyBase() { } export function getComfyHost() { + // Shared studios must never bypass the reservation-enforcing proxy. + const sharedHost = agentFrontDoorOrigin() + if (sharedHost) return sharedHost if (comfyHostOverride) return viaAgentFrontDoor(comfyHostOverride) || comfyHostOverride return viaAgentFrontDoor(configuredComfyBase()) || configuredComfyBase() } @@ -83,6 +87,12 @@ export function comfyWsUrl(clientId: string) { export async function comfyFetch(path: string, init?: RequestInit) { const url = `${getComfyHost()}${path}` + if (init?.method && !['GET', 'HEAD'].includes(init.method.toUpperCase())) { + await assertSharedGpu() + const headers = new Headers(init.headers) + for (const [key, value] of Object.entries(sharedGpuHeaders())) headers.set(key, value) + init = { ...init, headers } + } try { return await fetch(url, init) } catch (error) { @@ -623,7 +633,8 @@ async function purgeOnDesktop(opts: { headers: { 'Content-Type': 'application/json', Accept: 'application/json', - ...(token ? { Authorization: `Bearer ${token}` } : {}) + ...(token ? { Authorization: `Bearer ${token}` } : {}), + ...sharedGpuHeaders() }, body: JSON.stringify({ imageName: imageNames[0] || '', diff --git a/server/utils/comfyLifecycle.ts b/server/utils/comfyLifecycle.ts index 147b366..260b47e 100644 --- a/server/utils/comfyLifecycle.ts +++ b/server/utils/comfyLifecycle.ts @@ -1,3 +1,4 @@ +import { assertSharedGpu, sharedGpuHeaders, withSharedGpuStart } from '~/server/utils/sharedGpu' import { execFile, spawn } from 'node:child_process' import { promisify } from 'node:util' import { getComfyHost, setComfyHostOverride, comfyConfigured } from '~/server/utils/comfy' @@ -85,13 +86,15 @@ export async function fetchQueue() { async function controlRequest(path: string, method = 'GET', timeoutMs = 5000, body?: Record) { const { controlUrl, controlToken } = settings() if (!controlUrl) return null + if (method !== 'GET') await assertSharedGpu() try { const res = await fetch(`${controlUrl}${path}`, { method, headers: { Accept: 'application/json', ...(body ? { 'Content-Type': 'application/json' } : {}), - ...(controlToken ? { Authorization: `Bearer ${controlToken}` } : {}) + ...(controlToken ? { Authorization: `Bearer ${controlToken}` } : {}), + ...(method !== 'GET' ? sharedGpuHeaders() : {}) }, ...(body ? { body: JSON.stringify(body) } : {}), signal: AbortSignal.timeout(timeoutMs) @@ -344,7 +347,7 @@ async function waitWhileBusy(onStatus: StatusFn) { const started = Date.now() while (Date.now() - started < busyWaitMs) { const queue = await fetchQueue() - if (queue.running === 0) return + if (queue.ok && queue.running === 0 && queue.pending === 0) return onStatus({ state: 'busy', message: `ComfyUI is busy with another prompt (${queue.running} running, ${queue.pending} queued). Waiting up to 3 minutes, then this job will be saved.`, @@ -488,7 +491,10 @@ async function ensureUnlocked(onStatus: StatusFn, skipBusyWait = false) { export function ensureComfyReady(onStatus: StatusFn = () => undefined, options?: { skipBusyWait?: boolean }) { const skipBusyWait = options?.skipBusyWait === true - const run = gate.then(() => ensureUnlocked(onStatus, skipBusyWait)) + const run = gate.then(() => withSharedGpuStart( + () => ensureUnlocked(onStatus, skipBusyWait), + async () => { await assertSharedGpu() } + )) gate = run.then(() => undefined, () => undefined) return run } diff --git a/server/utils/imageComfy.ts b/server/utils/imageComfy.ts index a194aef..2218321 100644 --- a/server/utils/imageComfy.ts +++ b/server/utils/imageComfy.ts @@ -1,3 +1,4 @@ +import { assertSharedGpu, sharedGpuHeaders, sharedGpuConfigured } from '~/server/utils/sharedGpu' import { createHash } from 'node:crypto' import { AsyncLocalStorage } from 'node:async_hooks' import { getComfyHost, viaAgentFrontDoor } from '~/server/utils/comfy' @@ -61,6 +62,7 @@ export function getSidecarImageHost() { } export function getImageComfyHost() { + if (sharedGpuConfigured()) return getComfyHost() const fromAls = imageHostAls.getStore() if (fromAls) return viaAgentFrontDoor(fromAls) || fromAls if (imageComfyHostOverride) return viaAgentFrontDoor(imageComfyHostOverride) || imageComfyHostOverride @@ -80,6 +82,12 @@ export async function imageComfyFetch(path: string, init?: RequestInit) { if (!host) { throw createError({ statusCode: 500, statusMessage: 'No image ComfyUI host is configured' }) } + if (init?.method && !['GET', 'HEAD'].includes(init.method.toUpperCase())) { + await assertSharedGpu() + const headers = new Headers(init.headers) + for (const [key, value] of Object.entries(sharedGpuHeaders())) headers.set(key, value) + init = { ...init, headers } + } try { return await fetch(`${host}${path}`, init) } catch (error) { @@ -390,7 +398,8 @@ async function purgeImageOnDesktop(opts: { headers: { 'Content-Type': 'application/json', Accept: 'application/json', - ...(token ? { Authorization: `Bearer ${token}` } : {}) + ...(token ? { Authorization: `Bearer ${token}` } : {}), + ...sharedGpuHeaders() }, body: JSON.stringify({ imageName: imageNames[0] || '', diff --git a/server/utils/imageComfyLifecycle.ts b/server/utils/imageComfyLifecycle.ts index 46dc6f8..25a6e29 100644 --- a/server/utils/imageComfyLifecycle.ts +++ b/server/utils/imageComfyLifecycle.ts @@ -1,3 +1,4 @@ +import { assertSharedGpu, sharedGpuHeaders, withSharedGpuStart } from '~/server/utils/sharedGpu' import { getImageComfyHost, setImageComfyHostOverride, imageComfyConfigured, getBeastImageHost, getSidecarImageHost, withImageComfyHost, type ImageEditBox } from '~/server/utils/imageComfy' export type ImageComfyState = 'online' | 'booting' | 'busy' | 'starting' | 'offline' @@ -92,12 +93,14 @@ export async function fetchImageQueue() { async function controlRequest(path: string, method = 'GET', timeoutMs = 5000) { const { controlUrl, controlToken } = settings() if (!controlUrl) return null + if (method !== 'GET') await assertSharedGpu() try { const res = await fetch(`${controlUrl}${path}`, { method, headers: { Accept: 'application/json', - ...(controlToken ? { Authorization: `Bearer ${controlToken}` } : {}) + ...(controlToken ? { Authorization: `Bearer ${controlToken}` } : {}), + ...(method !== 'GET' ? sharedGpuHeaders() : {}) }, signal: AbortSignal.timeout(timeoutMs) }) @@ -238,7 +241,7 @@ async function ensureUnlocked(onStatus: StatusFn) { } export function ensureImageComfyReady(onStatus: StatusFn = () => undefined) { - const run = gate.then(() => ensureUnlocked(onStatus)) + const run = gate.then(() => withSharedGpuStart(() => ensureUnlocked(onStatus), async () => { await assertSharedGpu() })) gate = run.then(() => undefined, () => undefined) return run } diff --git a/server/utils/promptComfy.ts b/server/utils/promptComfy.ts index ba8a4da..fddc52d 100644 --- a/server/utils/promptComfy.ts +++ b/server/utils/promptComfy.ts @@ -1,4 +1,5 @@ -import { comfyInputFilename } from '~/server/utils/comfy' +import { assertSharedGpu, sharedGpuHeaders, sharedGpuConfigured } from '~/server/utils/sharedGpu' +import { comfyInputFilename, getComfyHost } from '~/server/utils/comfy' import { getSidecarImageHost } from '~/server/utils/imageComfy' import { ensureSidecarReady } from '~/server/utils/imageComfyLifecycle' import { buildVisionPromptWorkflow } from '~/server/utils/promptWorkflow' @@ -19,6 +20,7 @@ const jobs = new Map() const MAX_JOBS = 20 export function getPromptComfyHost() { + if (sharedGpuConfigured()) return getComfyHost() return getSidecarImageHost() } @@ -31,6 +33,12 @@ async function promptComfyFetch(path: string, init?: RequestInit) { if (!host) { throw createError({ statusCode: 503, statusMessage: 'Qwen VL is not configured. Set COMFY_HOST or IMAGE_COMFY_HOST to Beast.' }) } + if (init?.method && !['GET', 'HEAD'].includes(init.method.toUpperCase())) { + await assertSharedGpu() + const headers = new Headers(init.headers) + for (const [key, value] of Object.entries(sharedGpuHeaders())) headers.set(key, value) + init = { ...init, headers } + } try { return await fetch(`${host}${path}`, init) } catch (error) { @@ -260,3 +268,7 @@ export function failPromptJob(job: PromptJob, error: unknown) { job.error = err.statusMessage || err.message || 'Prompt recommend failed' job.message = job.error } + +export function hasActivePromptJobs() { + return [...jobs.values()].some(job => job.status === 'queued' || job.status === 'running') +} diff --git a/server/utils/sharedGpu.ts b/server/utils/sharedGpu.ts new file mode 100644 index 0000000..8028eae --- /dev/null +++ b/server/utils/sharedGpu.ts @@ -0,0 +1,106 @@ +import { randomUUID } from 'node:crypto' + +const ticket = randomUUID() +let lease = '' +let expiresAt = 0 +let lost = false +let starting = 0 +let resetting = false +let gate: Promise = Promise.resolve() +let localWork: () => boolean = () => true +let reason = 'GPU is in use. Waiting for availability.' + +function settings() { + const config = useRuntimeConfig() + return { + url: String(config.comfyControlUrl || process.env.COMFY_CONTROL_URL || '').replace(/\/$/, ''), + token: String(config.comfyControlToken || process.env.COMFY_CONTROL_TOKEN || '') + } +} +export function sharedGpuConfigured() { return Boolean(settings().url) } +export function sharedGpuWaitReason() { return reason } +export function registerSharedGpuWork(check: () => boolean) { localWork = check } + +async function request(path: string, body: Record) { + const { url, token } = settings() + const response = await fetch(`${url}/gpu/${path}`, { + method: 'POST', headers: { 'Content-Type': 'application/json', ...(token ? { Authorization: `Bearer ${token}` } : {}) }, + body: JSON.stringify(body), signal: AbortSignal.timeout(15_000) + }) + if (!response.ok) throw new Error('GPU coordinator unavailable') + return await response.json() as { acquired?: boolean; token?: string; ttlMs?: number; renewed?: boolean; released?: boolean } +} +function serialized(fn: () => Promise): Promise { + const run = gate.then(fn) + gate = run.catch(() => {}) + return run +} + +export function acquireSharedGpu() { + return serialized(async () => { + if (!sharedGpuConfigured()) return true + if (lost) return false + if (lease && Date.now() < expiresAt) return true + if (lease) { lost = true; return false } + try { + const result = await request('acquire', { ticket }) + if (!result.acquired || !result.token || !Number.isFinite(result.ttlMs) || Number(result.ttlMs) < 10_000) { + reason = 'GPU is in use. Waiting for availability.' + return false + } + lease = result.token + expiresAt = Date.now() + Math.max(1000, Number(result.ttlMs) - 5000) + return true + } catch { + reason = 'GPU coordinator unavailable. Waiting to reconnect.' + return false + } + }) +} + +export async function withSharedGpuStart(fn: () => Promise, blocked: () => Promise): Promise { + starting += 1 + try { return await (!resetting && await acquireSharedGpu() ? fn() : blocked()) } + finally { starting -= 1 } +} + +export async function assertSharedGpu() { + if (!await acquireSharedGpu()) throw createError({ statusCode: 409, statusMessage: reason, data: { code: 'GPU_BUSY' } }) +} + +export function sharedGpuHeaders() { + if (!sharedGpuConfigured()) return {} as Record + if (!lease || lost || Date.now() >= expiresAt) { + throw createError({ statusCode: 409, statusMessage: 'GPU reservation is unavailable. Waiting to reconnect.', data: { code: 'GPU_BUSY' } }) + } + const { token } = settings() + return { 'x-aigen-gpu-lease': lease, ...(token ? { Authorization: `Bearer ${token}` } : {}) } +} + +export function maintainSharedGpu() { + return serialized(async () => { + if (!sharedGpuConfigured()) return + if (!starting && !resetting && !localWork()) { + if (lease) { + try { await request('release', { token: lease }) } catch { /* expires if the host is unreachable */ } + } + lease = ''; expiresAt = 0; lost = false + return + } + if (!lease || lost) return + try { + const result = await request('renew', { token: lease }) + if (!result.renewed || !Number.isFinite(result.ttlMs) || Number(result.ttlMs) < 10_000) { lost = true; return } + expiresAt = Date.now() + Math.max(1000, Number(result.ttlMs) - 5000) + } catch { + if (Date.now() >= expiresAt) lost = true + } + }) +} + +export async function withSharedGpuReset(fn: () => Promise) { + if (resetting) throw createError({ statusCode: 409, statusMessage: 'GPU reset is already in progress.' }) + resetting = true + try { await assertSharedGpu(); return await fn() } + finally { resetting = false } +} diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts index 059e13f..f50d82f 100644 --- a/server/utils/studioQueue.ts +++ b/server/utils/studioQueue.ts @@ -1,3 +1,4 @@ +import { withSharedGpuStart, sharedGpuWaitReason, maintainSharedGpu, acquireSharedGpu } from '~/server/utils/sharedGpu' import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFileSync } from 'node:fs' import { join } from 'node:path' import { getJob, listJobs, emitJob, type Job } from '~/server/utils/jobs' @@ -98,6 +99,7 @@ export interface StudioJob { pausedByUser?: boolean resumeAutoRun?: boolean lastError?: string + waitReason?: string } type StudioQueueStore = { @@ -246,7 +248,8 @@ export function summarizeStudioJob(job: StudioJob) { holdForCutIn: job.holdForCutIn === true, pauseAfterCurrent: job.pauseAfterCurrent === true, pausedByUser: job.pausedByUser === true, - lastError: job.lastError + lastError: job.lastError, + waitReason: job.status === 'waiting' ? job.waitReason : undefined } } @@ -415,14 +418,11 @@ export async function videoJobsBusy() { }))) return true const queue = await fetchLiveQueue() if (queue) { - if (queue.running > 0) return true - if (queue.pending > 0) { - if (listJobs().some(job => liveJobOwnsGpu(job) || jobIsLocallySubmitting(job))) return true - if (listPendingJobs().some((pending) => { - if (!pending.promptId || pending.stopAfterCurrent) return false - return Date.now() - pending.startedAt < PENDING_ORPHAN_MS - })) return true - return false + if (queue.running > 0 || queue.pending > 0) { + // Join the device queue while another instance is rendering, so it cannot + // reclaim every gap ahead of this waiting instance. + await acquireSharedGpu() + return true } return false } @@ -865,7 +865,7 @@ function pendingAlive(job: StudioJob) { } function isTransientComfyError(error?: string) { - return /still busy|COMFY_BUSY|offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker|refused to start|API never answered|cold start|ECONNREFUSED|ETIMEDOUT|fetch failed|502|503/i.test(String(error || '')) + return /GPU is in use|GPU reservation|GPU coordinator|GPU_BUSY|still busy|COMFY_BUSY|offline|never answered|host agent|unreachable|asleep|Starting Comfy|ComfyUI \(1\)|instance picker|refused to start|API never answered|cold start|ECONNREFUSED|ETIMEDOUT|fetch failed|502|503/i.test(String(error || '')) } async function parkStudioOnStartFailure(owner: string, id: string, message: string) { @@ -1010,11 +1010,12 @@ async function dispatchStudioQueue() { } await settleAttachedTerminalLives() if (await videoJobsBusy()) { - if (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting'))) { + if (owners.some(owner => readJobs(owner).some(job => job.status === 'waiting' || (job.status === 'held' && !job.pausedByUser)))) { scheduleKickRetry() } return } + await maintainSharedGpu() // Re-read after videoJobsBusy — reap/settle may have rewritten disk rows. for (const owner of owners) { const store = readStore(owner) @@ -1039,6 +1040,13 @@ async function dispatchStudioQueue() { } async function resumeHeldStudioJob(item: StudioJob) { + return withSharedGpuStart(() => resumeHeldStudioJobReserved(item), async () => { + scheduleKickRetry() + return false + }) +} + +async function resumeHeldStudioJobReserved(item: StudioJob) { const latest = readJobs(item.ownerKey).find(job => job.id === item.id) if (!latest || latest.status === 'cancelled') return false if (!latest.shotQueueId) { @@ -1386,6 +1394,13 @@ async function startStudioMusicJob(item: StudioJob) { } export async function startStudioJob(item: StudioJob) { + return withSharedGpuStart(() => startStudioJobReserved(item), async () => { + await patchStudioJob(item.ownerKey, item.id, row => { row.waitReason = sharedGpuWaitReason() }) + scheduleKickRetry() + }) +} + +async function startStudioJobReserved(item: StudioJob) { if (studioJobKind(item) === 'music') { await startStudioMusicJob(item) return diff --git a/tests/gpu-reservation.test.mjs b/tests/gpu-reservation.test.mjs new file mode 100644 index 0000000..9dd19ab --- /dev/null +++ b/tests/gpu-reservation.test.mjs @@ -0,0 +1,75 @@ +import test from 'node:test' +import assert from 'node:assert/strict' +import { createGpuReservation, GPU_WAIT_MESSAGE } from '../scripts/gpu-reservation.mjs' +const a = 'aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa' +const b = 'bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb' + +test('simultaneous sites receive exactly one reservation', async () => { + const gpu = createGpuReservation({ idle: async () => true }) + const replies = await Promise.all([gpu.acquire(a), gpu.acquire(b)]) + assert.equal(replies.filter(reply => reply.acquired).length, 1) + assert.deepEqual(replies[1], { acquired: false, message: GPU_WAIT_MESSAGE, retryAfterMs: 2500 }) + assert.deepEqual(gpu.availability(), { busy: true }) +}) +test('another site cannot renew, release, or execute with an invalid token', async () => { + const gpu = createGpuReservation({ idle: async () => true }) + const first = await gpu.acquire(a) + assert.deepEqual(await gpu.renew(b), { renewed: false }) + assert.deepEqual(await gpu.release(b), { released: false }) + await assert.rejects(gpu.permit(b, () => assert.fail('must not execute')), { statusCode: 409 }) + assert.equal((await gpu.acquire(b)).acquired, false) + assert.equal(await gpu.permit(first.token, () => 'allowed'), 'allowed') +}) +test('release never grants another site while Comfy still has work', async () => { + let idle = true + const gpu = createGpuReservation({ idle: async () => idle }) + const first = await gpu.acquire(a) + idle = false + await gpu.release(first.token) + assert.equal((await gpu.acquire(b)).acquired, false) + idle = true + assert.equal((await gpu.acquire(b)).acquired, true) + await assert.rejects(gpu.permit(first.token, () => {}), { statusCode: 409 }) +}) +test('expired reservations recover only after confirmed GPU idle', async () => { + let time = 0, idle = true + const gpu = createGpuReservation({ idle: async () => idle, now: () => time, ttlMs: 100 }) + const first = await gpu.acquire(a) + time = 101; idle = false + assert.equal((await gpu.acquire(b)).acquired, false) + assert.equal((await gpu.renew(first.token)).renewed, false) + idle = true + assert.equal((await gpu.acquire(b)).acquired, true) +}) +test('an unanswered idle probe fails closed', async () => { + const gpu = createGpuReservation({ idle: async () => null }) + assert.equal((await gpu.acquire(a)).acquired, false) +}) +test('a waiting site gets the next reservation before the previous owner can reacquire', async () => { + const gpu = createGpuReservation({ idle: async () => true }) + const first = await gpu.acquire(a) + await gpu.acquire(b) + await gpu.release(first.token) + assert.equal((await gpu.acquire(a)).acquired, false) + assert.equal((await gpu.acquire(b)).acquired, true) +}) +test('an expired reservation cannot transfer while submission is in flight', async () => { + let time = 0, finish + const gpu = createGpuReservation({ idle: async () => true, now: () => time, ttlMs: 100 }) + const first = await gpu.acquire(a) + const work = gpu.permit(first.token, () => new Promise(resolve => { finish = resolve })) + await new Promise(resolve => setImmediate(resolve)) + time = 101 + assert.equal((await gpu.acquire(b)).acquired, false) + finish(); await work + assert.equal((await gpu.acquire(b)).acquired, true) +}) +test('disconnected waiters expire so they cannot block the device forever', async () => { + let time = 0 + const gpu = createGpuReservation({ idle: async () => true, now: () => time, ttlMs: 100, ticketTtlMs: 50 }) + const first = await gpu.acquire(a) + await gpu.acquire(b) + time = 51 + await gpu.release(first.token) + assert.equal((await gpu.acquire(a)).acquired, true) +}) diff --git a/tests/shared-gpu.test.mjs b/tests/shared-gpu.test.mjs new file mode 100644 index 0000000..f2e7904 --- /dev/null +++ b/tests/shared-gpu.test.mjs @@ -0,0 +1,81 @@ +import test from 'node:test' +import assert from 'node:assert/strict' +import http from 'node:http' +import { createRequire } from 'node:module' +import { readFileSync } from 'node:fs' +import ts from 'typescript' +import { createGpuReservation } from '../scripts/gpu-reservation.mjs' +import { createGpuProxy } from '../scripts/gpu-proxy.mjs' +const require = createRequire(import.meta.url) +const compiled = ts.transpileModule(readFileSync(new URL('../server/utils/sharedGpu.ts', import.meta.url), 'utf8'), { + compilerOptions: { module: ts.ModuleKind.CommonJS, target: ts.ScriptTarget.ES2022 } +}).outputText +const listen = server => new Promise(resolve => server.listen(0, '127.0.0.1', () => resolve(server.address().port))) +const close = server => new Promise(resolve => { server.closeAllConnections(); server.close(resolve) }) +const json = (res, body) => { res.setHeader('Content-Type', 'application/json'); res.end(JSON.stringify(body)) } +function client(url) { + const exports = {} + new Function('require', 'exports', 'useRuntimeConfig', 'createError', compiled)(require, exports, + () => ({ comfyControlUrl: url, comfyControlToken: 'test-secret' }), + info => Object.assign(new Error(info.statusMessage), info)) + let active = false + exports.registerSharedGpuWork(() => active) + return { api: exports, active: value => { active = value } } +} + +test('two independent studio clients share the real HTTP gateway without sharing job details', async t => { + let busy = false + const mutations = [] + const comfy = http.createServer((req,res) => { + if (req.method === 'GET') return json(res, { queue_running: busy ? [[1, 'private-id', { prompt: 'private prompt from other site' }]] : [], queue_pending: [] }) + mutations.push(req.url) + req.resume(); req.on('end', () => json(res, { prompt_id: 'opaque-result' })) + }) + const comfyPort = await listen(comfy) + t.after(() => close(comfy)) + const reservation = createGpuReservation({ idle: async () => !busy }) + const control = http.createServer(async (req,res) => { + let text = ''; for await (const chunk of req) text += chunk + const body = JSON.parse(text || '{}') + assert.equal(req.headers.authorization, 'Bearer test-secret') + const op = req.url.split('/').at(-1) + json(res, await reservation[op](body.ticket || body.token)) + }) + const controlPort = await listen(control) + t.after(() => close(control)) + const proxy = createGpuProxy({ target: () => comfyPort, reservation, authorized: req => req.headers.authorization === 'Bearer test-secret' }) + const proxyPort = await listen(proxy) + t.after(() => close(proxy)) + const a = client(`http://127.0.0.1:${controlPort}`) + const b = client(`http://127.0.0.1:${controlPort}`) + const grants = await Promise.all([a.api.acquireSharedGpu(), b.api.acquireSharedGpu()]) + assert.deepEqual(grants, [true, false]) + a.active(true) + await a.api.maintainSharedGpu() + assert.equal((await fetch(`http://127.0.0.1:${proxyPort}/prompt`, { method:'POST', headers:a.api.sharedGpuHeaders(), body:'{}' })).status, 200) + assert.equal((await fetch(`http://127.0.0.1:${proxyPort}/interrupt`, { method:'POST', headers:{Authorization:'Bearer test-secret'}, body:'{}' })).status, 409) + assert.deepEqual(mutations, ['/prompt']) + busy = true + const queue = await (await fetch(`http://127.0.0.1:${proxyPort}/queue`)).json() + assert.deepEqual(queue, { queue_running: [null], queue_pending: [] }) + assert.equal(JSON.stringify(queue).includes('private'), false) + assert.deepEqual(await (await fetch(`http://127.0.0.1:${proxyPort}/history`)).json(), {}) + assert.equal(await b.api.acquireSharedGpu(), false) + // Saving still owns the reservation even after Comfy reports idle. + busy = false + assert.equal(await b.api.acquireSharedGpu(), false) + a.active(false); await a.api.maintainSharedGpu() + assert.equal(await a.api.acquireSharedGpu(), false, 'waiting B must run before another A job') + assert.equal(await b.api.acquireSharedGpu(), true) + assert.equal((await fetch(`http://127.0.0.1:${proxyPort}/prompt`, { method:'POST', headers:b.api.sharedGpuHeaders(), body:'{}' })).status, 200) + // A reset cannot even enter its local clear routine while B owns the device. + await assert.rejects(a.api.withSharedGpuReset(async () => assert.fail('must not clear jobs')), {statusCode:409}) + b.active(false); await b.api.maintainSharedGpu() +}) + +test('missing coordinator cannot silently fall back to uncoordinated generation', async () => { + const c = client('http://127.0.0.1:1') + assert.equal(await c.api.acquireSharedGpu(), false) + await assert.rejects(c.api.assertSharedGpu(), {statusCode:409}) + assert.throws(() => c.api.sharedGpuHeaders(), {statusCode:409}) +}) diff --git a/tests/studio-queue.test.mjs b/tests/studio-queue.test.mjs index b93e2cb..1188f26 100644 --- a/tests/studio-queue.test.mjs +++ b/tests/studio-queue.test.mjs @@ -27,6 +27,7 @@ function fixture({ rows = [], lives = [], queue = { running: 0, pending: 0 }, hi const timers = [] const started = [] const modules = { + '~/server/utils/sharedGpu': { withSharedGpuStart: async fn => fn(), maintainSharedGpu: async () => {}, acquireSharedGpu: async () => false, sharedGpuWaitReason: () => 'GPU is in use. Waiting for availability.' }, '~/server/utils/jobs': { listJobs: () => lives, getJob: id => lives.find(job => job.id === id),