Keep pause across Coolify restarts, apply the chosen still from the image picker, and remux clips with Range/faststart so finished video can be scrubbed.

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Towsty
2026-08-28 08:15:52 -05:00
co-authored by Cursor
parent 870150f052
commit 7d783bd97e
12 changed files with 338 additions and 49 deletions
+54 -2
View File
@@ -1,4 +1,4 @@
import { existsSync, writeFileSync, readFileSync, unlinkSync } from 'node:fs'
import { existsSync, writeFileSync, readFileSync, unlinkSync, statSync, openSync, readSync, closeSync, copyFileSync } from 'node:fs'
import { dirname, join } from 'node:path'
import { spawn } from 'node:child_process'
@@ -171,6 +171,58 @@ async function probeSize(filePath: string) {
}
}
function mp4MoovIsAtStart(filePath: string) {
const fd = openSync(filePath, 'r')
try {
let offset = 0
const header = Buffer.alloc(16)
for (let i = 0; i < 12; i++) {
if (readSync(fd, header, 0, 8, offset) < 8) return false
let size = header.readUInt32BE(0)
const type = header.toString('ascii', 4, 8)
if (size === 1) {
if (readSync(fd, header, 8, 8, offset + 8) < 8) return false
size = Number(header.readBigUInt64BE(8))
}
if (!size || size < 8) return false
if (type === 'moov') return true
if (type === 'mdat') return false
offset += size
}
return false
} finally {
closeSync(fd)
}
}
export async function remuxFaststart(inputPath: string) {
if (!existsSync(inputPath)) return
try {
if (mp4MoovIsAtStart(inputPath)) return
} catch {
// Fall through and try a remux anyway.
}
const tmp = `${inputPath}.faststart.tmp.mp4`
try {
await runFfmpeg([
'-y',
'-i', inputPath,
'-c', 'copy',
'-map', '0',
'-movflags', '+faststart',
tmp
], 120000)
if (!existsSync(tmp) || statSync(tmp).size < 32) {
try { unlinkSync(tmp) } catch { /* ignore */ }
return
}
copyFileSync(tmp, inputPath)
try { unlinkSync(tmp) } catch { /* ignore */ }
} catch {
try { unlinkSync(tmp) } catch { /* ignore */ }
}
}
function encodeArgs(outputPath: string, withAudio: boolean) {
const args = [
'-c:v', 'libx264',
@@ -271,7 +323,7 @@ export async function concatMp4(part1Path: string, part2Path: string, outputPath
const listPath = join(dirname(outputPath), 'concat_list.txt')
writeFileSync(listPath, `file '${concatPath(part1Path)}'\nfile '${concatPath(part2Path)}'\n`)
try {
await runFfmpeg(['-y', '-f', 'concat', '-safe', '0', '-i', listPath, '-c', 'copy', outputPath], 180000)
await runFfmpeg(['-y', '-f', 'concat', '-safe', '0', '-i', listPath, '-c', 'copy', '-movflags', '+faststart', outputPath], 180000)
await assertStitchLength(part1Path, part2Path, outputPath)
} catch {
await concatReencode(part1Path, part2Path, outputPath)
+69
View File
@@ -0,0 +1,69 @@
import { createReadStream, statSync } from 'node:fs'
import { createError, getHeader, sendStream, setHeader, setResponseStatus, type H3Event } from 'h3'
function parseBytesRange(header: string | undefined, size: number) {
if (!header) return null
const match = /^bytes=(\d*)-(\d*)$/i.exec(header.trim())
if (!match) return null
const startToken = match[1]
const endToken = match[2]
if (!startToken && !endToken) return { unsatisfiable: true as const, start: 0, end: 0 }
let start: number
let end: number
if (!startToken) {
const suffix = Number(endToken)
if (!Number.isFinite(suffix) || suffix <= 0) return { unsatisfiable: true as const, start: 0, end: 0 }
start = Math.max(0, size - suffix)
end = size - 1
} else {
start = Number(startToken)
end = endToken ? Number(endToken) : size - 1
}
if (!Number.isFinite(start) || !Number.isFinite(end) || start < 0 || start >= size || end < start) {
return { unsatisfiable: true as const, start: 0, end: 0 }
}
return { start, end: Math.min(end, size - 1) }
}
function applyRangeHeaders(event: H3Event, contentType: string, extraHeaders?: Record<string, string>) {
setHeader(event, 'Accept-Ranges', 'bytes')
setHeader(event, 'Content-Type', contentType)
for (const [key, value] of Object.entries(extraHeaders || {})) {
setHeader(event, key, value)
}
}
function rejectUnsatisfiable(event: H3Event, size: number) {
setHeader(event, 'Content-Range', `bytes */${size}`)
throw createError({ statusCode: 416, statusMessage: 'Range Not Satisfiable' })
}
export function sendPathWithRange(event: H3Event, path: string, contentType: string, extraHeaders?: Record<string, string>) {
const size = statSync(path).size
applyRangeHeaders(event, contentType, extraHeaders)
const range = parseBytesRange(getHeader(event, 'range'), size)
if (!range) {
setHeader(event, 'Content-Length', String(size))
return sendStream(event, createReadStream(path))
}
if ('unsatisfiable' in range && range.unsatisfiable) rejectUnsatisfiable(event, size)
setResponseStatus(event, 206)
setHeader(event, 'Content-Range', `bytes ${range.start}-${range.end}/${size}`)
setHeader(event, 'Content-Length', String(range.end - range.start + 1))
return sendStream(event, createReadStream(path, { start: range.start, end: range.end }))
}
export function sendBufferWithRange(event: H3Event, buf: Buffer, contentType: string, extraHeaders?: Record<string, string>) {
const size = buf.length
applyRangeHeaders(event, contentType, extraHeaders)
const range = parseBytesRange(getHeader(event, 'range'), size)
if (!range) {
setHeader(event, 'Content-Length', String(size))
return buf
}
if ('unsatisfiable' in range && range.unsatisfiable) rejectUnsatisfiable(event, size)
setResponseStatus(event, 206)
setHeader(event, 'Content-Range', `bytes ${range.start}-${range.end}/${size}`)
setHeader(event, 'Content-Length', String(range.end - range.start + 1))
return buf.subarray(range.start, range.end + 1)
}
+2 -1
View File
@@ -5,7 +5,7 @@ import { spawn } from 'node:child_process'
import { randomBytes, scrypt, timingSafeEqual } from 'node:crypto'
import { promisify } from 'node:util'
import type { H3Event } from 'h3'
import { extractFirstFrame, extractLastFrame } from '~/server/utils/ffmpeg'
import { extractFirstFrame, extractLastFrame, remuxFaststart } from '~/server/utils/ffmpeg'
import { normalizePermanenceRefs, type PermanenceRef } from '~/utils/globalLocks'
const scryptAsync = promisify(scrypt)
@@ -1348,6 +1348,7 @@ export async function saveClip(params: {
mkdirSync(clipDir(params.ownerKey, clip.id), { recursive: true })
const videoPath = clipVideoPath(params.ownerKey, clip.id)
writeFileSync(videoPath, params.video)
await remuxFaststart(videoPath)
if (typeof params.duration === 'number' && params.duration > 0) {
clip.duration = params.duration
} else {
+11
View File
@@ -52,6 +52,7 @@ export interface PendingJob {
queueId?: string
queueAutoRun?: boolean
queueBudget?: number
stopAfterCurrent?: boolean
}
function pendingRoot() {
@@ -140,6 +141,7 @@ export function pendingFromJob(job: Job, overrides: Partial<PendingJob> = {}): P
queueId: library.queueId,
queueAutoRun: library.queueAutoRun,
queueBudget: library.queueBudget,
stopAfterCurrent: library.stopAfterCurrent,
...overrides
}
}
@@ -185,6 +187,7 @@ export function libraryFromPending(pending: PendingJob): NonNullable<Job['librar
queueId: pending.queueId,
queueAutoRun: pending.queueAutoRun,
queueBudget: pending.queueBudget,
stopAfterCurrent: pending.stopAfterCurrent,
globalLocks: pending.globalLocks,
permanenceRefs: pending.permanenceRefs,
shotPermanenceRefs: pending.shotPermanenceRefs
@@ -198,6 +201,14 @@ export function writePendingJob(job: PendingJob) {
renameSync(tmp, pendingPath(job.jobId))
}
export function patchPendingJob(jobId: string, patch: Partial<PendingJob>) {
const existing = readPendingJob(jobId)
if (!existing) return null
const next = { ...existing, ...patch }
writePendingJob(next)
return next
}
export function readPendingJob(jobId: string): PendingJob | null {
const path = pendingPath(jobId)
if (!existsSync(path)) return null
+67 -10
View File
@@ -1,6 +1,8 @@
import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFileSync } from 'node:fs'
import { join } from 'node:path'
import { getJob, listJobs, type Job } from '~/server/utils/jobs'
import { listPendingJobs, patchPendingJob, readPendingJob } from '~/server/utils/pending'
import { getShotQueue } from '~/server/utils/shotQueue'
import { parseVideoWorkflow, type VideoWorkflowId } from '~/utils/videoModels'
import type { PermanenceRef } from '~/utils/globalLocks'
@@ -181,7 +183,29 @@ export function summarizeStudioJob(job: StudioJob) {
}
export function videoJobsBusy() {
return listJobs().some(job => job.kind !== 'edit' && (job.status === 'queued' || job.status === 'uploading' || job.status === 'running'))
if (listJobs().some(job => job.kind !== 'edit' && (job.status === 'queued' || job.status === 'uploading' || job.status === 'running'))) {
return true
}
return listPendingJobs().some(job => Boolean(job.promptId))
}
export function applyPersistedPauseToJob(job: Job) {
const library = job.library
if (!library) return false
const store = readStore(library.ownerKey)
const row = store.jobs.find(item => item.liveJobId === job.id)
|| (library.queueId ? store.jobs.find(item => item.shotQueueId === library.queueId) : undefined)
const queue = library.queueId ? getShotQueue(library.ownerKey, library.queueId) : null
const paused = store.paused === true
|| row?.pauseAfterCurrent === true
|| row?.pausedByUser === true
|| queue?.stopAfterCurrent === true
|| queue?.status === 'paused'
if (paused) {
library.stopAfterCurrent = true
library.queueAutoRun = false
}
return paused
}
export async function addStudioJob(params: {
@@ -370,21 +394,42 @@ export function kickStudioQueue() {
return run
}
function repairStaleJobs(jobs: StudioJob[]) {
function pendingAlive(job: StudioJob) {
if (job.liveJobId && readPendingJob(job.liveJobId)?.promptId) return true
if (!job.shotQueueId) return false
return listPendingJobs().some(pending => (
pending.ownerKey === job.ownerKey
&& pending.queueId === job.shotQueueId
&& Boolean(pending.promptId)
))
}
function repairStaleJobs(jobs: StudioJob[], paused: boolean) {
for (const job of jobs) {
if (job.status !== 'running') continue
const live = job.liveJobId ? getJob(job.liveJobId) : undefined
const liveBusy = live && (live.status === 'queued' || live.status === 'uploading' || live.status === 'running')
if (liveBusy) continue
if (liveBusy || pendingAlive(job)) continue
const holdPause = paused || job.pauseAfterCurrent === true || job.pausedByUser === true
if (job.shotQueueId) {
job.status = 'held'
job.liveJobId = undefined
job.holdForCutIn = true
job.resumeAutoRun = job.payload.queueAutoRun === true || job.resumeAutoRun === true
if (holdPause) {
job.pausedByUser = true
job.pauseAfterCurrent = false
job.holdForCutIn = false
} else {
job.holdForCutIn = true
job.resumeAutoRun = job.payload.queueAutoRun === true || job.resumeAutoRun === true
}
job.updatedAt = Date.now()
} else {
job.status = 'complete'
job.status = holdPause ? 'held' : 'complete'
job.liveJobId = undefined
if (holdPause) {
job.pausedByUser = true
job.pauseAfterCurrent = false
}
job.updatedAt = Date.now()
}
}
@@ -394,7 +439,7 @@ async function dispatchStudioQueue() {
if (videoJobsBusy()) return
for (const owner of listOwnersWithStudioQueues()) {
const store = await mutateStore(owner, (current) => {
repairStaleJobs(current.jobs)
repairStaleJobs(current.jobs, current.paused === true)
const next = pruneDone(current.jobs)
current.jobs.splice(0, current.jobs.length, ...next)
return {
@@ -402,12 +447,12 @@ async function dispatchStudioQueue() {
jobs: current.jobs.map(item => structuredClone(item))
}
})
if (store.paused) return
const cutIn = pickCutIn(store.jobs)
if (cutIn) {
await startStudioJob(cutIn)
return
}
if (store.paused) return
const held = pickHeld(store.jobs)
if (held?.shotQueueId) {
await resumeHeldStudioJob(held)
@@ -606,7 +651,9 @@ export async function onLiveVideoSettled(job: Job) {
: 0
const failed = job.status === 'error' || job.status === 'cancelled'
const snapshot = readStore(owner)
const rowNow = snapshot.jobs.find(item => item.liveJobId === job.id) || snapshot.jobs.find(item => item.status === 'running')
const rowNow = snapshot.jobs.find(item => item.liveJobId === job.id)
|| snapshot.jobs.find(item => item.status === 'running' && item.shotQueueId && item.shotQueueId === job.library?.queueId)
|| snapshot.jobs.find(item => item.status === 'running' && !job.library?.queueId)
const userPause = !failed && (
rowNow?.pauseAfterCurrent === true
|| (snapshot.paused && !rowNow?.holdForCutIn)
@@ -674,6 +721,10 @@ function applyLiveJobPause(live: Job | undefined, pause: boolean, restoreAutoRun
live.library.stopAfterCurrent = pause
if (pause) live.library.queueAutoRun = false
else if (restoreAutoRun) live.library.queueAutoRun = true
patchPendingJob(live.id, {
stopAfterCurrent: pause,
queueAutoRun: live.library.queueAutoRun
})
}
export async function toggleStudioQueuePause(owner: string, options: {
@@ -754,10 +805,16 @@ export async function toggleStudioQueuePause(owner: string, options: {
}
})
const queueId = updated.job?.shotQueueId || options.queueId
if (updated.job?.liveJobId) {
applyLiveJobPause(getJob(updated.job.liveJobId), true)
patchPendingJob(updated.job.liveJobId, { stopAfterCurrent: true, queueAutoRun: false })
}
for (const pending of listPendingJobs()) {
if (pending.ownerKey !== owner) continue
if (queueId && pending.queueId && pending.queueId !== queueId) continue
patchPendingJob(pending.jobId, { stopAfterCurrent: true, queueAutoRun: false })
}
const queueId = updated.job?.shotQueueId || options.queueId
const queued = queueId ? getShotQueue(owner, queueId) : null
if (updated.job && queued?.autoRun) {
await patchStudioJob(owner, updated.job.id, (job) => {
+34 -15
View File
@@ -9,7 +9,7 @@ import { isLtxWorkflow, isTextToVideo, parseVideoWorkflow, workflowForExtension,
import { clipVideoPath, clipTitle, deleteRetryDraft, extendTempDir, getClip, nextClipPartName, removeExtendTemp, stillPath } from '~/server/utils/library'
import { extractLastFrame, probeHasAudio } from '~/server/utils/ffmpeg'
import { ensureComfyReady } from '~/server/utils/comfyLifecycle'
import { onLiveVideoSettled, onStudioQueueBurstStarted } from '~/server/utils/studioQueue'
import { onLiveVideoSettled, onStudioQueueBurstStarted, isStudioQueuePaused } from '~/server/utils/studioQueue'
import {
finishQueueBurst,
getShotQueue,
@@ -238,6 +238,17 @@ function seedCurrentVideo(job: Job) {
return { tmpDir, currentPath, part1Path: join(tmpDir, 'part1.mp4'), framePath: join(tmpDir, 'last_frame.png') }
}
function chainShouldHold(job: Job) {
if (!job.library) return false
if (job.library.stopAfterCurrent) return true
if (isStudioQueuePaused(job.library.ownerKey)) return true
const queue = job.library.queueId
? getShotQueue(job.library.ownerKey, job.library.queueId)
: null
if (queue?.stopAfterCurrent || queue?.status === 'paused') return true
return false
}
export async function continueQueuedExtensions(
job: Job,
params: VideoChainParams,
@@ -269,14 +280,22 @@ export async function continueQueuedExtensions(
for (let i = startFrom; i < extensions.length; i++) {
if (shotsLeft <= 0) break
assertJobActive(job)
if (job.library.stopAfterCurrent) break
if (chainShouldHold(job)) {
job.library.stopAfterCurrent = true
break
}
const liveQueue = job.library.queueId
? getShotQueue(job.library.ownerKey, job.library.queueId)
: null
if (liveQueue?.stopAfterCurrent) job.library.stopAfterCurrent = true
if (liveQueue?.autoRun === false) job.library.queueAutoRun = false
if (job.library.stopAfterCurrent) break
if (liveQueue?.stopAfterCurrent || liveQueue?.autoRun === false) {
if (liveQueue.stopAfterCurrent) job.library.stopAfterCurrent = true
if (liveQueue.autoRun === false) job.library.queueAutoRun = false
}
if (chainShouldHold(job)) {
job.library.stopAfterCurrent = true
break
}
const live = liveQueue ? liveSegmentPrompt(liveQueue, i) : extensions[i]
const ext = {
prompt: live.prompt || extensions[i].prompt,
@@ -289,10 +308,7 @@ export async function continueQueuedExtensions(
emitChainJob(job, { type: 'status', message: 'Waiting 3s buffer...', progress: 1 })
await sleep(3000)
assertJobActive(job)
const pausedQueue = job.library.queueId
? getShotQueue(job.library.ownerKey, job.library.queueId)
: null
if (job.library.stopAfterCurrent || pausedQueue?.stopAfterCurrent) {
if (chainShouldHold(job)) {
job.library.stopAfterCurrent = true
break
}
@@ -376,7 +392,10 @@ export async function continueQueuedExtensions(
job.library.extendPart1Path = undefined
}
if (shotsLeft <= 0 || job.library.stopAfterCurrent) break
if (shotsLeft <= 0 || chainShouldHold(job)) {
if (chainShouldHold(job)) job.library.stopAfterCurrent = true
break
}
}
} finally {
if (job.library) job.library.chainContinuing = false
@@ -398,7 +417,10 @@ export async function continueQueuedExtensions(
export async function continueQueuedExtensionsIfNeeded(job: Job) {
if (!job.library) return
if (job.status === 'error' || job.status === 'cancelled' || job.status === 'complete') return
if (job.library.stopAfterCurrent) return
if (chainShouldHold(job)) {
job.library.stopAfterCurrent = true
return
}
const remaining = remainingAfterCurrentShot(job.library)
if (!remaining.length) return
const queue = job.library.queueId
@@ -440,10 +462,7 @@ export async function runGeneration(job: Job, params: VideoChainParams) {
persist: extensions.length === 0 || !autoRun
})
const stopNow = job.library?.stopAfterCurrent === true
|| (job.library?.queueId
? getShotQueue(job.library.ownerKey, job.library.queueId)?.stopAfterCurrent === true
: false)
const stopNow = chainShouldHold(job)
const keepGoing = Boolean(autoRun && !stopNow && extensions.length && job.library)
if (job.library?.queueId && !keepGoing) {
await finishQueueBurst(job.library.ownerKey, job.library.queueId, true).catch(() => null)
+40 -3
View File
@@ -4,6 +4,7 @@ import { getJob, restoreJob } from '~/server/utils/jobs'
import type { PendingJob } from '~/server/utils/pending'
import { isLastChainShot, libraryFromPending, pendingFromJob, remainingAfterCurrentShot, writePendingJob, deletePendingJob } from '~/server/utils/pending'
import { syncQueueFromJob } from '~/server/utils/shotQueue'
import { isEncodingNode, NODE_LABELS } from '~/server/utils/workflow'
import { mergePermanenceRefs, type PermanenceRef } from '~/utils/globalLocks'
function clipRefsFromJob(library: NonNullable<Job['library']>): PermanenceRef[] | undefined {
@@ -52,6 +53,9 @@ export function mapChainProgress(job: Job, localPct: number) {
export function formatChainMessage(job: Job, message: string, samplePct?: number) {
const { chained, step, total, label } = chainMeta(job)
if (!chained) return message
if (/checkpoint saved|segment saved/i.test(message)) {
return `Shot ${step}/${total} saved`
}
if (typeof samplePct === 'number') {
return `Step ${step}/${total}: Generating ${label} (${samplePct}%)`
}
@@ -78,6 +82,13 @@ export function emitChainJob(job: Job, event: JobEvent, samplePct?: number) {
emitJob(job, next)
}
function isNonSamplerProgress(node: string) {
if (!node) return false
if (isEncodingNode(node)) return true
const label = NODE_LABELS[node] || ''
return /save|checkpoint|video combine|vhs/i.test(label) && !/sampler/i.test(label)
}
export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Promise<void> {
const persist = options.persist !== false
job.socketReady = false
@@ -350,6 +361,11 @@ export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Pr
if (job.promptId && data.prompt_id && data.prompt_id !== job.promptId) return
if (payload.type === 'progress') {
const node = data.node ? String(data.node) : ''
if (isNonSamplerProgress(node)) {
markActivity()
return
}
const value = Number(data.value || 0)
const max = Number(data.max || 1)
const pct = 12 + Math.round((value / Math.max(max, 1)) * 70)
@@ -360,7 +376,7 @@ export function watchComfyJob(job: Job, options: { persist?: boolean } = {}): Pr
step: value,
maxStep: max,
progress: pct,
node: data.node ? String(data.node) : null,
node: node || null,
message: `Sampling step ${value}/${max} (${samplePct}%)`
}, samplePct)
}
@@ -424,20 +440,41 @@ export function ensurePendingWatch(pending: PendingJob) {
pendingWatches.add(job.id)
const resume = async () => {
try {
const { applyPersistedPauseToJob, onLiveVideoSettled } = await import('~/server/utils/studioQueue')
applyPersistedPauseToJob(job)
const holdNow = () => job.library?.stopAfterCurrent === true
if (pending.promptId) {
await watchComfyJob(job, { persist: lastShot })
await watchComfyJob(job, { persist: lastShot || holdNow() })
}
applyPersistedPauseToJob(job)
if (job.status === 'error' || job.status === 'cancelled') {
await onLiveVideoSettled(job)
return
}
if (holdNow()) {
await onLiveVideoSettled(job)
return
}
if (job.status === 'complete') {
await onLiveVideoSettled(job)
return
}
if (job.status === 'error' || job.status === 'cancelled' || job.status === 'complete') return
if (!lastShot) {
const { continueQueuedExtensionsIfNeeded } = await import('~/server/utils/videoChain')
await continueQueuedExtensionsIfNeeded(job)
}
applyPersistedPauseToJob(job)
if (job.status === 'error' || job.status === 'cancelled' || job.status === 'complete' || holdNow()) {
await onLiveVideoSettled(job)
}
} catch (error) {
const message = error instanceof Error ? error.message : String(error)
if (job.status === 'error' || job.status === 'cancelled') return
job.status = 'error'
job.error = message
emitChainJob(job, { type: 'error', error: message, message })
const { onLiveVideoSettled } = await import('~/server/utils/studioQueue')
await onLiveVideoSettled(job).catch(() => null)
} finally {
pendingWatches.delete(job.id)
}