Exclusive llama-server load/unload on the 5080, queued Describe UI with resultText, Copy, and Use as prompt. Co-authored-by: Cursor <cursoragent@cursor.com>
269 lines
9.1 KiB
TypeScript
269 lines
9.1 KiB
TypeScript
import { createReadStream, existsSync, mkdirSync, readdirSync, readFileSync, writeFileSync, renameSync, unlinkSync } from 'node:fs'
|
|
import { join } from 'node:path'
|
|
import { createJob, emitJob, getJob, restoreJob, type Job } from './jobs'
|
|
import { getStill, stillPath } from './library'
|
|
import { sharedGpuHeaders } from './sharedGpu'
|
|
import { CAPTION_STYLES } from '~/shared/studio2/caption.mjs'
|
|
|
|
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 captionConfigured() { return Boolean(settings().url) }
|
|
|
|
function pendingRoot() {
|
|
return join(String(useRuntimeConfig().libraryDir || process.env.LIBRARY_DIR || '/data/library'), 'caption-pending')
|
|
}
|
|
|
|
export function captionPendingAlive(jobId: string) {
|
|
if (!jobId) return false
|
|
return existsSync(join(pendingRoot(), `${jobId}.json`))
|
|
}
|
|
|
|
async function request(path: string, method = 'GET', body?: unknown) {
|
|
const { url, token } = settings()
|
|
if (!url) throw new Error('Caption host is not configured. Set COMFY_CONTROL_URL.')
|
|
const response = await fetch(`${url}/caption/${path}`, {
|
|
method,
|
|
headers: {
|
|
...(token ? { Authorization: `Bearer ${token}` } : {}),
|
|
...(method === 'GET' ? {} : sharedGpuHeaders()),
|
|
...(method === 'POST' ? { 'Content-Type': 'application/json' } : {})
|
|
},
|
|
body: method === 'POST' ? JSON.stringify(body) : body as BodyInit | undefined,
|
|
...(method === 'PUT' ? { duplex: 'half' as const } : {}),
|
|
signal: AbortSignal.timeout(method === 'PUT' ? 120_000 : 60_000)
|
|
} as RequestInit)
|
|
if (!response.ok) {
|
|
const detail = await response.json().catch(() => ({})) as { message?: string; error?: string }
|
|
throw Object.assign(new Error(detail.message || detail.error || `Caption host returned ${response.status}`), { statusCode: response.status })
|
|
}
|
|
return response
|
|
}
|
|
|
|
function persist(record: Record<string, unknown>) {
|
|
mkdirSync(pendingRoot(), { recursive: true })
|
|
const path = join(pendingRoot(), `${record.id}.json`)
|
|
writeFileSync(path + '.tmp', JSON.stringify(record))
|
|
renameSync(path + '.tmp', path)
|
|
}
|
|
|
|
function clearPending(id: string) {
|
|
const path = join(pendingRoot(), `${id}.json`)
|
|
if (existsSync(path)) unlinkSync(path)
|
|
}
|
|
|
|
function writeSidecar(owner: string, stillId: string, text: string) {
|
|
try {
|
|
const image = stillPath(owner, stillId)
|
|
if (!existsSync(image)) return
|
|
writeFileSync(`${image}.txt`, text, 'utf8')
|
|
} catch { /* optional sidecar */ }
|
|
}
|
|
|
|
export async function cancelCaptionJob(job: Job) {
|
|
await request(`jobs/${job.id}/cancel`, 'POST', {})
|
|
job.status = 'cancelled'
|
|
emitJob(job, { type: 'error', error: 'Cancelled', message: 'Cancelled' })
|
|
}
|
|
|
|
async function settle(job: Job) {
|
|
const { onLiveVideoSettled } = await import('./studioQueue')
|
|
await onLiveVideoSettled(job)
|
|
clearPending(job.id)
|
|
}
|
|
|
|
async function watch(job: Job, record: {
|
|
id: string
|
|
owner: string
|
|
stillId: string
|
|
folderId: string
|
|
captionStyle: string
|
|
studio2Id?: string
|
|
liveId: string
|
|
}) {
|
|
let failures = 0
|
|
while (job.status !== 'cancelled') {
|
|
try {
|
|
const state = await (await request(`jobs/${job.id}`)).json() as {
|
|
status: string
|
|
message?: string
|
|
error?: string
|
|
progress?: number
|
|
text?: string
|
|
}
|
|
if (job.status === 'cancelled') { await settle(job); return }
|
|
if (state.status === 'error' || state.status === 'cancelled') {
|
|
job.status = state.status === 'cancelled' ? 'cancelled' : 'error'
|
|
job.error = state.error || state.message || 'Caption failed'
|
|
emitJob(job, { type: 'error', error: job.error, message: job.error })
|
|
if (record.studio2Id) {
|
|
const { readRecord, saveRecord } = await import('./studio2/store')
|
|
try {
|
|
const r = readRecord(record.owner, record.studio2Id)
|
|
r.state = job.status === 'cancelled' ? 'cancelled' : 'failed'
|
|
r.error = job.error
|
|
r.finishedAt = Date.now()
|
|
saveRecord(r)
|
|
} catch { /* record may already be gone */ }
|
|
}
|
|
await settle(job)
|
|
return
|
|
}
|
|
if (state.status === 'complete') {
|
|
const text = String(state.text || '').trim()
|
|
job.message = 'Caption ready'
|
|
job.progress = 100
|
|
job.status = 'complete'
|
|
;(job as Job & { resultText?: string }).resultText = text
|
|
if (record.studio2Id) {
|
|
const { readRecord, saveRecord } = await import('./studio2/store')
|
|
const r = readRecord(record.owner, record.studio2Id)
|
|
r.state = 'complete'
|
|
r.resultText = text
|
|
r.finishedAt = Date.now()
|
|
r.savedAt = Date.now()
|
|
saveRecord(r)
|
|
}
|
|
writeSidecar(record.owner, record.stillId, text)
|
|
emitJob(job, { type: 'complete', message: 'Caption ready', progress: 100 })
|
|
await settle(job)
|
|
return
|
|
}
|
|
failures = 0
|
|
job.status = 'running'
|
|
emitJob(job, {
|
|
type: 'progress',
|
|
message: state.message || 'Captioning',
|
|
progress: Math.min(95, Math.max(1, Number(state.progress || 10)))
|
|
})
|
|
} catch (error) {
|
|
failures++
|
|
emitJob(job, { type: 'status', message: `Caption host check failed; retrying: ${error instanceof Error ? error.message : String(error)}` })
|
|
if (failures >= 10) {
|
|
job.status = 'error'
|
|
job.error = 'Caption host unreachable.'
|
|
emitJob(job, { type: 'error', error: job.error, message: job.error })
|
|
await settle(job)
|
|
return
|
|
}
|
|
}
|
|
await new Promise(resolve => setTimeout(resolve, 1500))
|
|
}
|
|
await settle(job)
|
|
}
|
|
|
|
export async function startCaptionJob(params: {
|
|
ownerKey: string
|
|
folderId: string
|
|
stillId: string
|
|
captionStyle: string
|
|
studio2Id?: string
|
|
name?: string
|
|
}) {
|
|
if (!CAPTION_STYLES.includes(params.captionStyle)) throw new Error('Unknown caption style.')
|
|
getStill(params.ownerKey, params.stillId)
|
|
const job = createJob('caption')
|
|
job.caption = true
|
|
job.library = {
|
|
ownerKey: params.ownerKey,
|
|
folderId: params.folderId,
|
|
hideThumbnail: false,
|
|
name: params.name || `Describe · ${params.captionStyle}`,
|
|
prompt: params.captionStyle,
|
|
aspect: 'image',
|
|
width: 0,
|
|
height: 0,
|
|
steps: 1,
|
|
turbo: true,
|
|
seed: 0,
|
|
stillId: params.stillId,
|
|
engine: 'caption'
|
|
}
|
|
const record = {
|
|
id: job.id,
|
|
liveId: job.id,
|
|
owner: params.ownerKey,
|
|
stillId: params.stillId,
|
|
folderId: params.folderId,
|
|
captionStyle: params.captionStyle,
|
|
studio2Id: params.studio2Id,
|
|
createdAt: Date.now(),
|
|
status: 'running'
|
|
}
|
|
persist(record)
|
|
void (async () => {
|
|
try {
|
|
if (job.status === 'cancelled') { await settle(job); return }
|
|
await request(`jobs/${job.id}/input`, 'PUT', createReadStream(stillPath(params.ownerKey, params.stillId)))
|
|
if (job.status === 'cancelled') { await settle(job); return }
|
|
try {
|
|
await request('jobs', 'POST', { id: job.id, style: params.captionStyle })
|
|
} catch {
|
|
/* Confirm by ID; never caption twice on an HTTP timeout. */
|
|
}
|
|
await watch(job, record)
|
|
} catch (error) {
|
|
const statusCode = (error as { statusCode?: number })?.statusCode
|
|
if (statusCode && statusCode >= 400 && statusCode < 500) {
|
|
job.status = 'error'
|
|
job.error = error instanceof Error ? error.message : String(error)
|
|
emitJob(job, { type: 'error', error: job.error, message: job.error })
|
|
if (params.studio2Id) {
|
|
try {
|
|
const { readRecord, saveRecord } = await import('./studio2/store')
|
|
const r = readRecord(params.ownerKey, params.studio2Id)
|
|
r.state = 'failed'
|
|
r.error = job.error
|
|
r.finishedAt = Date.now()
|
|
saveRecord(r)
|
|
} catch { /* ignore */ }
|
|
}
|
|
await settle(job)
|
|
return
|
|
}
|
|
emitJob(job, { type: 'status', message: `Checking caption submission: ${error instanceof Error ? error.message : String(error)}` })
|
|
await watch(job, record)
|
|
}
|
|
})()
|
|
return job
|
|
}
|
|
|
|
export function resumeCaptionJobs() {
|
|
if (!existsSync(pendingRoot())) return
|
|
for (const file of readdirSync(pendingRoot()).filter(file => /^[a-zA-Z0-9-]+\.json$/.test(file))) {
|
|
try {
|
|
const record = JSON.parse(readFileSync(join(pendingRoot(), file), 'utf8'))
|
|
if (getJob(record.id)) continue
|
|
const job = restoreJob({
|
|
id: record.liveId || record.id,
|
|
clientId: record.clientId || crypto.randomUUID(),
|
|
promptId: '',
|
|
startedAt: record.createdAt || Date.now(),
|
|
library: {
|
|
ownerKey: record.owner,
|
|
folderId: record.folderId,
|
|
hideThumbnail: false,
|
|
prompt: record.captionStyle || 'descriptive',
|
|
aspect: 'image',
|
|
width: 0,
|
|
height: 0,
|
|
steps: 1,
|
|
turbo: true,
|
|
seed: 0,
|
|
stillId: record.stillId,
|
|
engine: 'caption'
|
|
}
|
|
})
|
|
job.caption = true
|
|
job.kind = 'caption'
|
|
job.message = 'Reconnecting to caption host'
|
|
void watch(job, record)
|
|
} catch { /* Preserve invalid records for diagnosis. */ }
|
|
}
|
|
}
|