From ec5791c5952ab8f4e70352b11d82a1e6ca0abf39 Mon Sep 17 00:00:00 2001 From: Towsty Date: Mon, 14 Sep 2026 18:26:55 -0500 Subject: [PATCH] Mount YuE2 on the host agent and share GPU ownership. Co-authored-by: Cursor --- scripts/comfy-host-agent.mjs | 42 +++++++++++++++++++++++++++++++++--- server/utils/studioQueue.ts | 23 ++++++++++++++------ 2 files changed, 55 insertions(+), 10 deletions(-) diff --git a/scripts/comfy-host-agent.mjs b/scripts/comfy-host-agent.mjs index fe57401..b24117b 100644 --- a/scripts/comfy-host-agent.mjs +++ b/scripts/comfy-host-agent.mjs @@ -4,6 +4,7 @@ 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' @@ -147,7 +148,7 @@ let proxyTarget = 0 function ensureProxyListening() { if (proxyServer) return - proxyServer = createGpuProxy({ target: () => proxyTarget, reservation: gpuReservation, authorized, markWork, externalBusy: () => yueGp.busy() || upscale.busy() }) + 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) })) }) @@ -819,7 +820,7 @@ function purgeDesktopFiles(body) { } const gpuReservation = createGpuReservation({ idle: async () => { - if (yueGp.busy() || upscale.busy()) return false + if (yueGp.busy() || yue2.busy() || upscale.busy()) return false if ((await trainingLock()).busy) return false const healthy = await syncProxy() if (healthy) { @@ -834,6 +835,7 @@ 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) @@ -845,6 +847,22 @@ const yueGp = createYueGpHost({ } }) +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.') + 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) { @@ -874,6 +892,22 @@ async function handleControl(req, res) { } 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 @@ -906,12 +940,14 @@ async function handleControl(req, res) { 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, { @@ -1013,7 +1049,7 @@ const server = http.createServer(async (req, res) => { } else await handleControl(req, res) } catch (error) { req.resume() - if ((String(req.url || '').startsWith('/yuegp/') || String(req.url || '').startsWith('/upscale/')) && !res.headersSent) return json(res, error.statusCode || 400, { error: error.message || 'YuEGP request failed' }) + if ((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.' }) } }) diff --git a/server/utils/studioQueue.ts b/server/utils/studioQueue.ts index 775a233..9887665 100644 --- a/server/utils/studioQueue.ts +++ b/server/utils/studioQueue.ts @@ -286,7 +286,7 @@ function failZombieLiveJob(job: Job, error: string) { function sweepStaleLiveJobs() { const now = Date.now() for (const job of listJobs()) { - if (job.studio2 || job.yueGp || job.upscale) continue + if (job.studio2 || job.yueGp || job.yue2 || job.upscale) continue if (job.saving) continue if (job.status === 'queued' && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) { failZombieLiveJob(job, 'Job never started') @@ -306,7 +306,7 @@ async function reapZombieLiveJobs() { const { fetchHistory } = await import('~/server/utils/comfy') for (const job of listJobs()) { // Never interrupt download/stitch/library save — Comfy is idle then by design. - if (job.studio2 || job.yueGp || job.upscale) continue + if (job.studio2 || job.yueGp || job.yue2 || job.upscale) continue if (job.saving) continue if (job.library?.chainContinuing) continue if (job.status !== 'running' && job.status !== 'uploading' && job.status !== 'queued') continue @@ -351,7 +351,7 @@ async function reapZombieLiveJobs() { } function liveJobOwnsGpu(job: Job) { - if ((job.studio2 || job.yueGp || job.upscale) && ['running', 'queued', 'uploading'].includes(job.status)) return true + if ((job.studio2 || job.yueGp || job.yue2 || job.upscale) && ['running', 'queued', 'uploading'].includes(job.status)) return true if (job.library?.stopAfterCurrent) return false if (job.library?.chainContinuing) return true if (job.saving) return true @@ -371,7 +371,7 @@ function liveJobOwnsGpu(job: Job) { */ function clearDeadGpuClaimsForForceStart() { for (const live of listJobs()) { - if (live.studio2 || live.yueGp || live.upscale) continue + if (live.studio2 || live.yueGp || live.yue2 || live.upscale) continue if (live.saving) continue if (jobIsLocallySubmitting(live)) continue if (live.status !== 'running' && live.status !== 'queued' && live.status !== 'uploading') { @@ -598,6 +598,10 @@ export async function clearStuckStudioWork(owner: string) { const { cancelYueGpJob } = await import('./yueGp') await cancelYueGpJob(job) } + if (job.yue2) { + const { cancelYue2Job } = await import('./yue2') + await cancelYue2Job(job) + } job.status = 'cancelled' job.error = 'Cleared by force reset' if (job.library) { @@ -665,6 +669,11 @@ async function stopLiveGeneration(liveJobId?: string, shotQueueId?: string) { await cancelYueGpJob(live) return } + if (live?.yue2) { + const { cancelYue2Job } = await import('./yue2') + await cancelYue2Job(live) + return + } if (live) { live.status = 'cancelled' if (live.library) { @@ -914,7 +923,7 @@ function repairStaleJobs(jobs: StudioJob[]) { job.updatedAt = Date.now() continue } - if (!job.payload.upscale && job.payload.musicEngine !== 'yue' && job.status === 'error' && isTransientComfyError(job.lastError)) { + if (!job.payload.upscale && job.payload.musicEngine !== 'yue' && job.payload.musicEngine !== 'yue2' && job.status === 'error' && isTransientComfyError(job.lastError)) { job.status = 'waiting' job.liveJobId = undefined job.lastError = undefined @@ -923,7 +932,7 @@ function repairStaleJobs(jobs: StudioJob[]) { job.updatedAt = Date.now() continue } - if (!job.payload.upscale && job.payload.musicEngine !== 'yue' && job.status === 'held' && isTransientComfyError(job.lastError)) { + if (!job.payload.upscale && job.payload.musicEngine !== 'yue' && job.payload.musicEngine !== 'yue2' && job.status === 'held' && isTransientComfyError(job.lastError)) { job.status = 'waiting' job.liveJobId = undefined job.lastError = undefined @@ -1687,7 +1696,7 @@ export async function onLiveVideoSettled(job: Job) { return } const remaining = remainingStudioShots(job) - const wakeFail = !job.studio2 && !job.yueGp && !job.upscale && job.status === 'error' && remaining > 0 && isTransientComfyError(job.error) + const wakeFail = !job.studio2 && !job.yueGp && !job.yue2 && !job.upscale && job.status === 'error' && remaining > 0 && isTransientComfyError(job.error) const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail await mutateStore(owner, (store) => {