Add image-to-text Describe on the bench via exclusive llama.cpp caption jobs.

This commit is contained in:
Towsty
2026-09-30 22:50:19 -05:00
parent 9481791ac5
commit 15ef891b26
15 changed files with 937 additions and 35 deletions
+232
View File
@@ -0,0 +1,232 @@
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/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) {
return Boolean(jobId) && 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 } : {}),
// Caption loads a 7B VLM; allow the host POST to finish or fall through to watch
signal: AbortSignal.timeout(method === 'PUT' ? 120_000 : method === 'POST' ? 300_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 */ }
}
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
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 })
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.resultText = text
writeSidecar(record.owner, record.stillId, text)
persist({ ...record, resultText: text, status: 'complete' })
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
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,
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 on 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 })
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'
if (record.resultText) job.resultText = record.resultText
job.message = 'Reconnecting to caption host'
void watch(job, record)
} catch { /* keep bad records */ }
}
}
+4 -1
View File
@@ -36,9 +36,11 @@ export interface Job {
upscale?: boolean
yueGp?: boolean
yue2?: boolean
caption?: boolean
resultText?: string
musicActivity?: { checkedAt: number; running: boolean }
id: string
kind?: 'video' | 'edit' | 'music'
kind?: 'video' | 'edit' | 'music' | 'caption'
promptId?: string
clientId: string
status: JobStatus
@@ -222,6 +224,7 @@ export function jobSnapshot(job: Job) {
clipId: job.clipId,
stillId: job.stillId,
trackId: job.trackId,
resultText: job.resultText,
hideThumbnail: job.hideThumbnail,
error: job.error,
folderLocked: job.library?.folderLocked,
+61 -9
View File
@@ -13,7 +13,7 @@ import { imageV2StackSpecials } from '~/utils/imageV2'
import { allowIdentityRefs, type PermanenceRef } from '~/utils/globalLocks'
export type StudioJobStatus = 'waiting' | 'running' | 'held' | 'complete' | 'error' | 'cancelled'
export type StudioJobKind = 'video' | 'edit' | 'music'
export type StudioJobKind = 'video' | 'edit' | 'music' | 'caption'
export interface StudioJobPayload {
studio2Id?: string
@@ -42,6 +42,7 @@ export interface StudioJobPayload {
useIdentityRefs: boolean
stillId?: string
stillFilename?: string
captionStyle?: string
hideThumbnail: boolean
hideInput?: boolean
folderLocked?: boolean
@@ -105,6 +106,7 @@ export interface StudioJob {
resumeAutoRun?: boolean
lastError?: string
waitReason?: string
resultText?: string
}
type StudioQueueStore = {
@@ -226,6 +228,7 @@ export function listStudioJobs(owner: string) {
export function studioJobKind(job: Pick<StudioJob, 'kind'> | { kind?: string }) {
if (job.kind === 'edit') return 'edit'
if (job.kind === 'music') return 'music'
if (job.kind === 'caption') return 'caption'
return 'video'
}
@@ -245,6 +248,8 @@ export function summarizeStudioJob(job: StudioJob) {
folderId: job.payload.folderId,
musicEngine: job.payload.musicEngine,
stillId: job.payload.stillId,
captionStyle: job.payload.captionStyle,
resultText: job.resultText,
workflow: job.payload.workflow,
imagePipeline: job.payload.imagePipeline || 'v1',
duration: job.payload.duration,
@@ -290,7 +295,7 @@ function failZombieLiveJob(job: Job, error: string) {
function sweepStaleLiveJobs() {
const now = Date.now()
for (const job of listJobs()) {
if (job.studio2 || job.yueGp || job.yue2 || job.upscale) continue
if (job.studio2 || job.yueGp || job.yue2 || job.upscale || job.caption) continue
if (job.saving) continue
if (job.status === 'queued' && !job.promptId && now - job.startedAt >= QUEUED_GRACE_MS) {
failZombieLiveJob(job, 'Job never started')
@@ -310,12 +315,12 @@ 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.yue2 || job.upscale) continue
if (job.studio2 || job.yueGp || job.yue2 || job.upscale || job.caption) continue
if (job.saving) continue
if (job.library?.chainContinuing) continue
if (job.status !== 'running' && job.status !== 'uploading' && job.status !== 'queued') continue
if (jobIsLocallySubmitting(job)) continue
const zombieMs = (job.kind === 'edit' || job.kind === 'music') ? 45_000 : ZOMBIE_EMPTY_COMFY_MS
const zombieMs = (job.kind === 'edit' || job.kind === 'music' || job.kind === 'caption') ? 45_000 : ZOMBIE_EMPTY_COMFY_MS
if (jobAgeMs(job) < zombieMs) continue
if (!job.promptId) {
if (jobAgeMs(job) >= QUEUED_GRACE_MS) failZombieLiveJob(job, 'Job never started on ComfyUI')
@@ -355,7 +360,7 @@ async function reapZombieLiveJobs() {
}
function liveJobOwnsGpu(job: Job) {
if ((job.studio2 || job.yueGp || job.yue2 || job.upscale) && ['running', 'queued', 'uploading'].includes(job.status)) return true
if ((job.studio2 || job.yueGp || job.yue2 || job.upscale || job.caption) && ['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
@@ -375,7 +380,7 @@ function liveJobOwnsGpu(job: Job) {
*/
function clearDeadGpuClaimsForForceStart() {
for (const live of listJobs()) {
if (live.studio2 || live.yueGp || live.yue2 || live.upscale) continue
if (live.studio2 || live.yueGp || live.yue2 || live.upscale || live.caption) continue
if (live.saving) continue
if (jobIsLocallySubmitting(live)) continue
if (live.status !== 'running' && live.status !== 'queued' && live.status !== 'uploading') {
@@ -469,8 +474,11 @@ export async function addStudioJob(params: {
kind?: StudioJobKind
}) {
const now = Date.now()
const kind = params.kind === 'edit' ? 'edit' : params.kind === 'music' ? 'music' : 'video'
const shotCount = kind === 'music'
const kind = params.kind === 'edit' ? 'edit'
: params.kind === 'music' ? 'music'
: params.kind === 'caption' ? 'caption'
: 'video'
const shotCount = kind === 'music' || kind === 'caption'
? 1
: kind === 'edit'
? 1 + (params.payload.passes?.length || 0)
@@ -606,6 +614,10 @@ export async function clearStuckStudioWork(owner: string) {
const { cancelYue2Job } = await import('./yue2')
await cancelYue2Job(job)
}
if (job.caption) {
const { cancelCaptionJob } = await import('./caption')
await cancelCaptionJob(job)
}
job.status = 'cancelled'
job.error = 'Cleared by force reset'
if (job.library) {
@@ -678,6 +690,11 @@ async function stopLiveGeneration(liveJobId?: string, shotQueueId?: string) {
await cancelYue2Job(live)
return
}
if (live?.caption) {
const { cancelCaptionJob } = await import('./caption')
await cancelCaptionJob(live)
return
}
if (live) {
live.status = 'cancelled'
if (live.library) {
@@ -889,6 +906,10 @@ function pendingAlive(job: StudioJob) {
const root = join(String(useRuntimeConfig().libraryDir || process.env.LIBRARY_DIR || '/data/library'), 'yue2-pending', `${job.liveJobId}.json`)
if (existsSync(root)) return true
}
if (job.kind === 'caption') {
const root = join(String(useRuntimeConfig().libraryDir || process.env.LIBRARY_DIR || '/data/library'), 'caption-pending', `${job.liveJobId}.json`)
if (existsSync(root)) return true
}
}
if (!job.shotQueueId) return false
return listPendingJobs().some(pending => (
@@ -1397,6 +1418,32 @@ async function startStudioExtendJob(item: StudioJob) {
}
}
async function startStudioCaptionJob(item: StudioJob) {
let live: import('~/server/utils/jobs').Job | undefined
try {
const { startCaptionJob } = await import('./caption')
const style = String(item.payload.captionStyle || 'descriptive')
const stillId = String(item.payload.stillId || '')
if (!stillId) throw new Error('Caption requires a still.')
live = await startCaptionJob({
ownerKey: item.ownerKey,
folderId: item.payload.folderId,
stillId,
captionStyle: style,
name: item.payload.name
})
await markStudioLive(item.ownerKey, item.id, live.id)
} catch (error) {
const message = error instanceof Error ? error.message : String(error)
if (live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')) {
live.status = 'error'
live.error = message
}
await parkStudioOnStartFailure(item.ownerKey, item.id, message)
kickStudioQueue()
}
}
async function startStudioMusicJob(item: StudioJob) {
let live: import('~/server/utils/jobs').Job | undefined
try {
@@ -1439,6 +1486,10 @@ export async function startStudioJob(item: StudioJob) {
}
async function startStudioJobReserved(item: StudioJob) {
if (studioJobKind(item) === 'caption') {
await startStudioCaptionJob(item)
return
}
if (item.payload.studio2Id) {
try {
const { startStudio2Job } = await import('./studio2/runner')
@@ -1703,7 +1754,7 @@ export async function onLiveVideoSettled(job: Job) {
return
}
const remaining = remainingStudioShots(job)
const wakeFail = !job.studio2 && !job.yueGp && !job.yue2 && !job.upscale && job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
const wakeFail = !job.studio2 && !job.yueGp && !job.yue2 && !job.upscale && !job.caption && job.status === 'error' && remaining > 0 && isTransientComfyError(job.error)
const failed = (job.status === 'error' || job.status === 'cancelled') && !wakeFail
await mutateStore(owner, (store) => {
@@ -1777,6 +1828,7 @@ export async function onLiveVideoSettled(job: Job) {
: failed
? 'error'
: 'complete'
if (job.caption && job.resultText) row.resultText = job.resultText
clearStudioRowSlot(row, status, failed ? job.error : undefined)
syncPausedFlag(store)
})