Files
aigen/server/utils/jobs.ts
T

255 lines
6.7 KiB
TypeScript

export type JobStatus = 'queued' | 'uploading' | 'running' | 'complete' | 'error' | 'cancelled' | 'deferred'
export interface JobEvent {
type: string
message?: string
progress?: number
step?: number
maxStep?: number
node?: string | null
filename?: string
subfolder?: string
mediaType?: string
clipId?: string
stillId?: string
trackId?: string
hideThumbnail?: boolean
error?: string
elapsedMs?: number
busy?: boolean
queueRunning?: number
queuePending?: number
folderLocked?: boolean
jobId?: string
draftId?: string
chainStep?: number
chainTotal?: number
chainLabel?: string
overallProgress?: number
queuedRemaining?: number
queueId?: string
}
export interface Job {
musicActivity?: { checkedAt: number; running: boolean }
id: string
kind?: 'video' | 'edit' | 'music'
promptId?: string
clientId: string
status: JobStatus
message: string
progress: number
step: number
maxStep: number
startedAt: number
video?: { filename: string; subfolder: string; type: string }
audio?: { filename: string; subfolder: string; type: string }
clipId?: string
stillId?: string
trackId?: string
hideThumbnail?: boolean
imageComfyHost?: string
library?: {
ownerKey: string
folderId: string
hideThumbnail: boolean
hideInput?: boolean
folderLocked?: boolean
name?: string
prompt: string
promptMid?: string
promptPre?: string
promptPost?: string
aspect: string
width: number
height: number
steps: number
turbo: boolean
seed: number
cfg?: number
fps?: number
samplerName?: string
scheduler?: string
thumb?: Buffer
imageName?: string
imageSubfolder?: string
stillId?: string
stillFilename?: string
referenceImageNames?: string[]
referenceStillIds?: Array<string | null>
useIdentityRefs?: boolean
globalLocks?: string
permanenceRefs?: import('~/utils/globalLocks').PermanenceRef[]
shotPermanenceRefs?: import('~/utils/globalLocks').PermanenceRef[][]
duration?: number
sound?: boolean
draftId?: string
extensions?: import('~/server/utils/library').QueuedExtension[]
chainIndex?: number
chainStep?: number
chainTotal?: number
chainLabel?: string
extendTmpDir?: string
extendPart1Path?: string
extendSourceClipId?: string
familyId?: string
parentClipId?: string
parentStillId?: string
workflow?: import('~/utils/videoModels').VideoWorkflowId
chainContinuing?: boolean
passes?: import('~/utils/imageIterations').ImageIteration[]
queueId?: string
queueAutoRun?: boolean
queueBudget?: number
stopAfterCurrent?: boolean
refineExtensionFrame?: boolean
saveLosslessAnchor?: boolean
refinementDenoise?: number
loraName?: string
loraStack?: import('~/utils/loras').LoraStackItem[]
shotLoras?: string[]
shotLoraStacks?: import('~/utils/loras').LoraStackItem[][]
tags?: string
lyrics?: string
instrumental?: boolean
audioExt?: string
engine?: string
lyricsStrength?: number
}
error?: string
socketReady?: boolean
saving?: boolean
savedPromptId?: string
segmentBuffer?: Buffer
events: JobEvent[]
listeners: Set<(event: JobEvent) => void>
}
const jobs = new Map<string, Job>()
const MAX_JOBS = 40
export function createJob(kind: Job['kind'] = 'video'): Job {
const job: Job = {
id: crypto.randomUUID(),
kind,
clientId: crypto.randomUUID(),
status: 'queued',
message: 'Queued',
progress: 0,
step: 0,
maxStep: 0,
startedAt: Date.now(),
events: [],
listeners: new Set()
}
jobs.set(job.id, job)
while (jobs.size > MAX_JOBS) {
const oldest = jobs.keys().next().value
if (oldest) jobs.delete(oldest)
}
return job
}
export function getJob(id: string) {
return jobs.get(id)
}
export function restoreJob(params: {
id: string
clientId: string
promptId: string
startedAt: number
hideThumbnail?: boolean
library: Job['library']
}) {
const existing = jobs.get(params.id)
if (existing) return existing
const job: Job = {
id: params.id,
clientId: params.clientId,
kind: 'video',
promptId: params.promptId,
status: 'running',
message: 'Waiting for ComfyUI to finish this job...',
progress: 8,
step: 0,
maxStep: params.library?.steps || 0,
startedAt: params.startedAt,
hideThumbnail: params.hideThumbnail === true,
library: params.library,
events: [],
listeners: new Set()
}
jobs.set(job.id, job)
return job
}
export function listJobs() {
return [...jobs.values()]
}
export function jobSnapshot(job: Job) {
const media = job.kind === 'music' ? job.audio : job.video
return {
jobId: job.id,
kind: job.kind || 'video',
engine: job.library?.engine,
musicActivity: job.musicActivity,
type: 'snapshot' as const,
status: job.status,
message: job.message,
progress: job.progress,
step: job.step,
maxStep: job.maxStep,
promptId: job.promptId,
elapsedMs: Date.now() - job.startedAt,
filename: job.library?.folderLocked ? undefined : media?.filename,
subfolder: job.library?.folderLocked ? undefined : media?.subfolder,
mediaType: job.library?.folderLocked ? undefined : media?.type,
clipId: job.clipId,
stillId: job.stillId,
trackId: job.trackId,
hideThumbnail: job.hideThumbnail,
error: job.error,
folderLocked: job.library?.folderLocked,
draftId: job.library?.draftId,
chainStep: job.library?.chainStep,
chainTotal: job.library?.chainTotal,
chainLabel: job.library?.chainLabel,
queueId: job.library?.queueId,
queuedRemaining: job.library?.queueId
? Math.max(0, (job.library.chainTotal || 1) - (job.library.chainStep || 1))
: undefined
}
}
export function emitJob(job: Job, event: JobEvent) {
event.jobId = job.id
if (event.message) job.message = event.message
if (typeof event.progress === 'number') job.progress = event.progress
if (typeof event.step === 'number') job.step = event.step
if (typeof event.maxStep === 'number') job.maxStep = event.maxStep
if (typeof event.chainStep === 'number') {
if (job.library) job.library.chainStep = event.chainStep
}
if (typeof event.chainTotal === 'number') {
if (job.library) job.library.chainTotal = event.chainTotal
}
if (event.chainLabel && job.library) job.library.chainLabel = event.chainLabel
event.elapsedMs = Date.now() - job.startedAt
job.events.push(event)
if (job.events.length > 200) job.events.shift()
for (const listener of job.listeners) {
try {
listener(event)
} catch {
// ignore broken SSE listeners
}
}
}
export function subscribeJob(job: Job, listener: (event: JobEvent) => void) {
job.listeners.add(listener)
return () => job.listeners.delete(listener)
}