mirror of
https://github.com/Chocobozzz/PeerTube.git
synced 2026-09-03 20:53:09 -05:00
Fix infinite non-settle when client closes the stream
This commit is contained in:
@@ -30,6 +30,7 @@ export class VideoDownload {
|
|||||||
private readonly video: MVideoThumbnails
|
private readonly video: MVideoThumbnails
|
||||||
private readonly videoFiles: MVideoFile[]
|
private readonly videoFiles: MVideoFile[]
|
||||||
|
|
||||||
|
private cleanupPromise: Promise<void>
|
||||||
private allowDirectSending = true
|
private allowDirectSending = true
|
||||||
|
|
||||||
constructor (options: {
|
constructor (options: {
|
||||||
@@ -45,99 +46,114 @@ export class VideoDownload {
|
|||||||
bytesPerIpPerSecond: number
|
bytesPerIpPerSecond: number
|
||||||
ip: string
|
ip: string
|
||||||
}) {
|
}) {
|
||||||
return new Promise<void>(async (res, rej) => {
|
const totalBytesPerSecond = options?.totalBytesPerSecond
|
||||||
const totalBytesPerSecond = options?.totalBytesPerSecond
|
const bytesPerIpPerSecond = options?.bytesPerIpPerSecond
|
||||||
const bytesPerIpPerSecond = options?.bytesPerIpPerSecond
|
const ip = options?.ip
|
||||||
const ip = options?.ip
|
|
||||||
|
let rejectOnStreamError: (err: Error) => void
|
||||||
|
const streamErrorPromise = new Promise<never>((_, rej) => {
|
||||||
|
rejectOnStreamError = rej
|
||||||
|
})
|
||||||
|
|
||||||
|
const run = async () => {
|
||||||
|
VideoDownload.totalDownloads++
|
||||||
|
|
||||||
|
const maxResolution = await this.buildMuxInputs(rejectOnStreamError)
|
||||||
|
|
||||||
|
// Include cover to audio file?
|
||||||
|
const { coverPath, isTmpDestination } = maxResolution === 0
|
||||||
|
? await this.buildCoverInput()
|
||||||
|
: { coverPath: undefined, isTmpDestination: false }
|
||||||
|
|
||||||
|
if (coverPath && isTmpDestination) {
|
||||||
|
this.tmpDestinations.push(coverPath)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Prefer sending the file directly if possible
|
||||||
|
if (this.allowDirectSending && !coverPath && this.inputs.length === 1) {
|
||||||
|
logger.info(`Piping single file for video ${this.video.url}`, { input: this.inputsToLog()[0], ...lTags(this.video.uuid) })
|
||||||
|
|
||||||
|
const input = typeof this.inputs[0] === 'string'
|
||||||
|
? createReadStream(this.inputs[0])
|
||||||
|
: this.inputs[0]
|
||||||
|
|
||||||
|
const throttleStream = totalBytesPerSecond || bytesPerIpPerSecond
|
||||||
|
? new ThrottleStream({ totalBytesPerSecond, bytesPerIpPerSecond, ip })
|
||||||
|
: new PassThrough()
|
||||||
|
|
||||||
|
try {
|
||||||
|
await pipeline(input, throttleStream, output)
|
||||||
|
} catch (err) {
|
||||||
|
if ((err?.message || '').includes('Output stream closed')) {
|
||||||
|
logger.info(`Client aborted direct download for video ${this.video.url}`, lTags(this.video.uuid))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
throw err
|
||||||
|
}
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
logger.info(`Muxing files for video ${this.video.url}`, { inputs: this.inputsToLog(), ...lTags(this.video.uuid) })
|
||||||
|
|
||||||
|
this.ffmpegContainer = new FFmpegContainer(getFFmpegCommandWrapperOptions('vod'))
|
||||||
|
|
||||||
|
const throttleStream = totalBytesPerSecond || bytesPerIpPerSecond
|
||||||
|
? new ThrottleStream({ totalBytesPerSecond, bytesPerIpPerSecond, ip })
|
||||||
|
: undefined
|
||||||
|
|
||||||
|
const finalOutput = throttleStream ?? output
|
||||||
|
|
||||||
|
const throttlePipeline = throttleStream
|
||||||
|
? pipeline(throttleStream, output)
|
||||||
|
: Promise.resolve()
|
||||||
|
|
||||||
try {
|
try {
|
||||||
VideoDownload.totalDownloads++
|
// Run in parallel to prevent throttlePipeline unhandled rejection if an input stream errors
|
||||||
|
await Promise.all([
|
||||||
|
this.ffmpegContainer.mergeInputs({
|
||||||
|
inputs: this.inputs,
|
||||||
|
output: finalOutput,
|
||||||
|
logError: false,
|
||||||
|
|
||||||
const maxResolution = await this.buildMuxInputs(rej)
|
// Include a cover if this is an audio file
|
||||||
|
coverPath
|
||||||
|
}),
|
||||||
|
|
||||||
// Include cover to audio file?
|
throttlePipeline
|
||||||
const { coverPath, isTmpDestination } = maxResolution === 0
|
])
|
||||||
? await this.buildCoverInput()
|
|
||||||
: { coverPath: undefined, isTmpDestination: false }
|
|
||||||
|
|
||||||
if (coverPath && isTmpDestination) {
|
logger.info(`Mux ended for video ${this.video.url}`, { inputs: this.inputsToLog(), ...lTags(this.video.uuid) })
|
||||||
this.tmpDestinations.push(coverPath)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Prefer sending the file directly if possible
|
|
||||||
if (this.allowDirectSending && !coverPath && this.inputs.length === 1) {
|
|
||||||
logger.info(`Piping single file for video ${this.video.url}`, { input: this.inputsToLog()[0], ...lTags(this.video.uuid) })
|
|
||||||
|
|
||||||
const input = typeof this.inputs[0] === 'string'
|
|
||||||
? createReadStream(this.inputs[0])
|
|
||||||
: this.inputs[0]
|
|
||||||
|
|
||||||
const throttleStream = totalBytesPerSecond || bytesPerIpPerSecond
|
|
||||||
? new ThrottleStream({ totalBytesPerSecond, bytesPerIpPerSecond, ip })
|
|
||||||
: new PassThrough()
|
|
||||||
|
|
||||||
await pipeline(input, throttleStream, output)
|
|
||||||
|
|
||||||
res()
|
|
||||||
} else {
|
|
||||||
logger.info(`Muxing files for video ${this.video.url}`, { inputs: this.inputsToLog(), ...lTags(this.video.uuid) })
|
|
||||||
|
|
||||||
this.ffmpegContainer = new FFmpegContainer(getFFmpegCommandWrapperOptions('vod'))
|
|
||||||
|
|
||||||
const throttleStream = totalBytesPerSecond || bytesPerIpPerSecond
|
|
||||||
? new ThrottleStream({ totalBytesPerSecond, bytesPerIpPerSecond, ip })
|
|
||||||
: undefined
|
|
||||||
|
|
||||||
const finalOutput = throttleStream ?? output
|
|
||||||
|
|
||||||
const throttlePipeline = throttleStream
|
|
||||||
? pipeline(throttleStream, output)
|
|
||||||
: Promise.resolve()
|
|
||||||
|
|
||||||
try {
|
|
||||||
// Run in parallel to prevent throttlePipeline unhandled rejection if an input stream errors
|
|
||||||
await Promise.all([
|
|
||||||
this.ffmpegContainer.mergeInputs({
|
|
||||||
inputs: this.inputs,
|
|
||||||
output: finalOutput,
|
|
||||||
logError: false,
|
|
||||||
|
|
||||||
// Include a cover if this is an audio file
|
|
||||||
coverPath
|
|
||||||
}),
|
|
||||||
|
|
||||||
throttlePipeline
|
|
||||||
])
|
|
||||||
|
|
||||||
logger.info(`Mux ended for video ${this.video.url}`, { inputs: this.inputsToLog(), ...lTags(this.video.uuid) })
|
|
||||||
|
|
||||||
res()
|
|
||||||
} catch (err) {
|
|
||||||
const message = err?.message || ''
|
|
||||||
|
|
||||||
if (message.includes('Output stream closed')) {
|
|
||||||
logger.info(`Client aborted mux for video ${this.video.url}`, lTags(this.video.uuid))
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
if (err.inputStreamError) {
|
|
||||||
err.inputStreamError = buildRequestError(err.inputStreamError)
|
|
||||||
}
|
|
||||||
|
|
||||||
logger.warn(`Cannot mux files of video ${this.video.url}`, { err, inputs: this.inputsToLog(), ...lTags(this.video.uuid) })
|
|
||||||
|
|
||||||
throw err
|
|
||||||
} finally {
|
|
||||||
this.ffmpegContainer.forceKill()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
rej(err)
|
const message = err?.message || ''
|
||||||
|
|
||||||
|
if (message.includes('Output stream closed')) {
|
||||||
|
logger.info(`Client aborted mux for video ${this.video.url}`, lTags(this.video.uuid))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
if (err.inputStreamError) {
|
||||||
|
err.inputStreamError = buildRequestError(err.inputStreamError)
|
||||||
|
}
|
||||||
|
|
||||||
|
logger.warn(`Cannot mux files of video ${this.video.url}`, { err, inputs: this.inputsToLog(), ...lTags(this.video.uuid) })
|
||||||
|
|
||||||
|
throw err
|
||||||
} finally {
|
} finally {
|
||||||
this.cleanup()
|
// cleanup() may already have force killed and unset ffmpegContainer if a stream errored and won the race below
|
||||||
.catch(cleanupErr => logger.error('Cannot cleanup after mux error', { err: cleanupErr, ...lTags(this.video.uuid) }))
|
this.ffmpegContainer?.forceKill()
|
||||||
}
|
}
|
||||||
})
|
}
|
||||||
|
|
||||||
|
try {
|
||||||
|
// If a stream errors, don't wait for run() to notice: reject immediately so the caller isn't stuck
|
||||||
|
// run() keeps executing in the background but will abort once cleanup() destroys its streams/ffmpeg process
|
||||||
|
await Promise.race([ run(), streamErrorPromise ])
|
||||||
|
} finally {
|
||||||
|
this.cleanup()
|
||||||
|
.catch(cleanupErr => logger.error('Cannot cleanup after mux error', { err: cleanupErr, ...lTags(this.video.uuid) }))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
@@ -167,9 +183,7 @@ export class VideoDownload {
|
|||||||
...lTags(this.video.uuid)
|
...lTags(this.video.uuid)
|
||||||
})
|
})
|
||||||
|
|
||||||
this.cleanup()
|
// cleanup() is already run by muxToMergeVideoFiles's finally block once this rejection wins the race
|
||||||
.catch(cleanupErr => logger.error('Cannot cleanup after mux error', { err: cleanupErr, ...lTags(this.video.uuid) }))
|
|
||||||
|
|
||||||
rej(err)
|
rej(err)
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
@@ -290,6 +304,14 @@ export class VideoDownload {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private async cleanup () {
|
private async cleanup () {
|
||||||
|
if (this.cleanupPromise === undefined) {
|
||||||
|
this.cleanupPromise = this._doCleanup()
|
||||||
|
}
|
||||||
|
|
||||||
|
return this.cleanupPromise
|
||||||
|
}
|
||||||
|
|
||||||
|
private async _doCleanup () {
|
||||||
VideoDownload.totalDownloads--
|
VideoDownload.totalDownloads--
|
||||||
|
|
||||||
for (const destination of this.tmpDestinations) {
|
for (const destination of this.tmpDestinations) {
|
||||||
|
|||||||
Reference in New Issue
Block a user