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) { 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. */ } } }