Mount YuE2 on the host agent and share GPU ownership.
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -4,6 +4,7 @@ import { stableMemoryArgs } from './comfy-memory-policy.mjs'
|
|||||||
import { createGpuReservation } from './gpu-reservation.mjs'
|
import { createGpuReservation } from './gpu-reservation.mjs'
|
||||||
import { createGpuProxy } from './gpu-proxy.mjs'
|
import { createGpuProxy } from './gpu-proxy.mjs'
|
||||||
import { createYueGpHost } from './yuegp-host.mjs'
|
import { createYueGpHost } from './yuegp-host.mjs'
|
||||||
|
import { createYue2Host } from './yue2-host.mjs'
|
||||||
import http from 'node:http'
|
import http from 'node:http'
|
||||||
import net from 'node:net'
|
import net from 'node:net'
|
||||||
import { execFile, spawn } from 'node:child_process'
|
import { execFile, spawn } from 'node:child_process'
|
||||||
@@ -147,7 +148,7 @@ let proxyTarget = 0
|
|||||||
|
|
||||||
function ensureProxyListening() {
|
function ensureProxyListening() {
|
||||||
if (proxyServer) return
|
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) => {
|
proxyServer.on('error', (error) => {
|
||||||
console.log(JSON.stringify({ src: 'comfy-host-agent', event: 'proxy-error', error: String(error.message || 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 () => {
|
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
|
if ((await trainingLock()).busy) return false
|
||||||
const healthy = await syncProxy()
|
const healthy = await syncProxy()
|
||||||
if (healthy) {
|
if (healthy) {
|
||||||
@@ -834,6 +835,7 @@ const yueGp = createYueGpHost({
|
|||||||
leaseValid: lease => gpuReservation.isOwner(lease),
|
leaseValid: lease => gpuReservation.isOwner(lease),
|
||||||
prepare: async () => {
|
prepare: async () => {
|
||||||
if ((await trainingLock()).busy) throw new Error('GPU is busy with training.')
|
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()
|
const healthy = await syncProxy()
|
||||||
if (healthy) {
|
if (healthy) {
|
||||||
const queue = await fetchLocalQueue(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) })
|
const upscale = createUpscaleHost({ leaseValid: token => gpuReservation.isOwner(token) })
|
||||||
|
|
||||||
async function handleControl(req, res) {
|
async function handleControl(req, res) {
|
||||||
@@ -874,6 +892,22 @@ async function handleControl(req, res) {
|
|||||||
}
|
}
|
||||||
return json(res, 404, { error: 'Unknown YuEGP endpoint' })
|
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/')) {
|
if (req.method === 'POST' && url.pathname.startsWith('/gpu/')) {
|
||||||
const body = await readJson(req)
|
const body = await readJson(req)
|
||||||
let result
|
let result
|
||||||
@@ -906,12 +940,14 @@ async function handleControl(req, res) {
|
|||||||
gpu: gpuReservation.availability(),
|
gpu: gpuReservation.availability(),
|
||||||
training: { busy: lastTraining.busy },
|
training: { busy: lastTraining.busy },
|
||||||
yuegp: { busy: yueGp.busy(), configured: yueGp.configured() },
|
yuegp: { busy: yueGp.busy(), configured: yueGp.configured() },
|
||||||
|
yue2: { busy: yue2.busy(), configured: yue2.configured() },
|
||||||
upscale: { busy: upscale.busy(), engine: 'realesrgan-rife', local: true }
|
upscale: { busy: upscale.busy(), engine: 'realesrgan-rife', local: true }
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
if (req.method === 'POST' && url.pathname === '/start') {
|
if (req.method === 'POST' && url.pathname === '/start') {
|
||||||
if (upscale.busy()) return json(res, 409, { message: 'Local upscale is using the GPU.' })
|
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 (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()
|
const training = await trainingLock()
|
||||||
if (training.busy) {
|
if (training.busy) {
|
||||||
return json(res, 409, {
|
return json(res, 409, {
|
||||||
@@ -1013,7 +1049,7 @@ const server = http.createServer(async (req, res) => {
|
|||||||
} else await handleControl(req, res)
|
} else await handleControl(req, res)
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
req.resume()
|
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.' })
|
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.' })
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -286,7 +286,7 @@ function failZombieLiveJob(job: Job, error: string) {
|
|||||||
function sweepStaleLiveJobs() {
|
function sweepStaleLiveJobs() {
|
||||||
const now = Date.now()
|
const now = Date.now()
|
||||||
for (const job of listJobs()) {
|
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.saving) continue
|
||||||
if (job.status === 'queued' && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) {
|
if (job.status === 'queued' && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) {
|
||||||
failZombieLiveJob(job, 'Job never started')
|
failZombieLiveJob(job, 'Job never started')
|
||||||
@@ -306,7 +306,7 @@ async function reapZombieLiveJobs() {
|
|||||||
const { fetchHistory } = await import('~/server/utils/comfy')
|
const { fetchHistory } = await import('~/server/utils/comfy')
|
||||||
for (const job of listJobs()) {
|
for (const job of listJobs()) {
|
||||||
// Never interrupt download/stitch/library save — Comfy is idle then by design.
|
// 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.saving) continue
|
||||||
if (job.library?.chainContinuing) continue
|
if (job.library?.chainContinuing) continue
|
||||||
if (job.status !== 'running' && job.status !== 'uploading' && job.status !== 'queued') continue
|
if (job.status !== 'running' && job.status !== 'uploading' && job.status !== 'queued') continue
|
||||||
@@ -351,7 +351,7 @@ async function reapZombieLiveJobs() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
function liveJobOwnsGpu(job: Job) {
|
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?.stopAfterCurrent) return false
|
||||||
if (job.library?.chainContinuing) return true
|
if (job.library?.chainContinuing) return true
|
||||||
if (job.saving) return true
|
if (job.saving) return true
|
||||||
@@ -371,7 +371,7 @@ function liveJobOwnsGpu(job: Job) {
|
|||||||
*/
|
*/
|
||||||
function clearDeadGpuClaimsForForceStart() {
|
function clearDeadGpuClaimsForForceStart() {
|
||||||
for (const live of listJobs()) {
|
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 (live.saving) continue
|
||||||
if (jobIsLocallySubmitting(live)) continue
|
if (jobIsLocallySubmitting(live)) continue
|
||||||
if (live.status !== 'running' && live.status !== 'queued' && live.status !== 'uploading') {
|
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')
|
const { cancelYueGpJob } = await import('./yueGp')
|
||||||
await cancelYueGpJob(job)
|
await cancelYueGpJob(job)
|
||||||
}
|
}
|
||||||
|
if (job.yue2) {
|
||||||
|
const { cancelYue2Job } = await import('./yue2')
|
||||||
|
await cancelYue2Job(job)
|
||||||
|
}
|
||||||
job.status = 'cancelled'
|
job.status = 'cancelled'
|
||||||
job.error = 'Cleared by force reset'
|
job.error = 'Cleared by force reset'
|
||||||
if (job.library) {
|
if (job.library) {
|
||||||
@@ -665,6 +669,11 @@ async function stopLiveGeneration(liveJobId?: string, shotQueueId?: string) {
|
|||||||
await cancelYueGpJob(live)
|
await cancelYueGpJob(live)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
if (live?.yue2) {
|
||||||
|
const { cancelYue2Job } = await import('./yue2')
|
||||||
|
await cancelYue2Job(live)
|
||||||
|
return
|
||||||
|
}
|
||||||
if (live) {
|
if (live) {
|
||||||
live.status = 'cancelled'
|
live.status = 'cancelled'
|
||||||
if (live.library) {
|
if (live.library) {
|
||||||
@@ -914,7 +923,7 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
|||||||
job.updatedAt = Date.now()
|
job.updatedAt = Date.now()
|
||||||
continue
|
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.status = 'waiting'
|
||||||
job.liveJobId = undefined
|
job.liveJobId = undefined
|
||||||
job.lastError = undefined
|
job.lastError = undefined
|
||||||
@@ -923,7 +932,7 @@ function repairStaleJobs(jobs: StudioJob[]) {
|
|||||||
job.updatedAt = Date.now()
|
job.updatedAt = Date.now()
|
||||||
continue
|
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.status = 'waiting'
|
||||||
job.liveJobId = undefined
|
job.liveJobId = undefined
|
||||||
job.lastError = undefined
|
job.lastError = undefined
|
||||||
@@ -1687,7 +1696,7 @@ export async function onLiveVideoSettled(job: Job) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
const remaining = remainingStudioShots(job)
|
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
|
const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail
|
||||||
|
|
||||||
await mutateStore(owner, (store) => {
|
await mutateStore(owner, (store) => {
|
||||||
|
|||||||
Reference in New Issue
Block a user