From 9a7120f5f4b83b6b16436801f265a67dcf0e0775 Mon Sep 17 00:00:00 2001 From: Vanilagy <1696106+Vanilagy@users.noreply.github.com> Date: Fri, 17 Apr 2026 18:35:33 +0200 Subject: [PATCH] Add per-playlist "mutices" for HLS muxer --- src/hls/hls-muxer.ts | 113 +++++++++++++++++++++++++------------------ src/media-source.ts | 58 +++++++++++++--------- 2 files changed, 101 insertions(+), 70 deletions(-) diff --git a/src/hls/hls-muxer.ts b/src/hls/hls-muxer.ts index 4012837..adb3498 100644 --- a/src/hls/hls-muxer.ts +++ b/src/hls/hls-muxer.ts @@ -8,7 +8,16 @@ import { MediaCodec, validateAudioChunkMetadata, validateVideoChunkMetadata } from '../codec'; import { EncodedAudioPacketSource, EncodedVideoPacketSource } from '../media-source'; -import { arrayArgmax, assert, findLastIndex, joinPaths, textEncoder, toArray, UNDETERMINED_LANGUAGE } from '../misc'; +import { + arrayArgmax, + assert, + AsyncMutex, + findLastIndex, + joinPaths, + textEncoder, + toArray, + UNDETERMINED_LANGUAGE, +} from '../misc'; import { Muxer } from '../muxer'; import { Output, @@ -79,6 +88,11 @@ type Playlist = { nextOffset: number; info: HlsOutputSegmentInfo; } | null; + + // For HLS, having a single mutex is too coarse. Every playlist is basically independent and therefore we can have + // a per-playlist mutex instead of a per-muxer one. This means two packets from different playlists coming in don't + // block each other. + mutex: AsyncMutex; }; type PlaylistDeclaration = { @@ -131,6 +145,8 @@ export class HlsMuxer extends Muxer { } async start(): Promise { + const release = await this.mutex.acquire(); + const someRelative = this.output._tracks.some(t => t.metadata.isRelativeToUnixEpoch); const someNotRelative = this.output._tracks.some(t => !t.metadata.isRelativeToUnixEpoch); if (someRelative && someNotRelative) { @@ -449,6 +465,7 @@ export class HlsMuxer extends Muxer { mediaSequence: 0, done: false, singleFile: null, + mutex: new AsyncMutex(), }; this.playlists.push(playlist); @@ -491,6 +508,8 @@ export class HlsMuxer extends Muxer { : [], }); } + + release(); } async getMimeType(): Promise { @@ -509,17 +528,17 @@ export class HlsMuxer extends Muxer { // eslint-disable-next-line @typescript-eslint/no-misused-promises override async onTrackClose(track: OutputTrack) { - const release = await this.mutex.acquire(); + const trackData = this.trackDatas.find(x => x.track === track); + if (trackData) { + trackData.closed = true; + } + + const playlist = this.playlists.find(x => x.tracks.includes(track)); + assert(playlist); // If there isn't one then the assignment algo failed innit + + const release = await playlist.mutex.acquire(); try { - const trackData = this.trackDatas.find(x => x.track === track); - if (trackData) { - trackData.closed = true; - } - - const playlist = this.playlists.find(x => x.tracks.includes(track)); - assert(playlist); // If there isn't one then the assignment algo failed innit - await this.advancePlaylist(playlist); } finally { release(); @@ -589,18 +608,17 @@ export class HlsMuxer extends Muxer { packet: EncodedPacket, meta?: EncodedVideoChunkMetadata, ) { - const release = await this.mutex.acquire(); + const trackData = this.getVideoTrackData(track, meta); + const playlist = trackData.playlist; + + const release = await playlist.mutex.acquire(); try { - const trackData = this.getVideoTrackData(track, meta); - const timestamp = this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); const adjustedPacket = packet.clone({ timestamp }); trackData.packets.push(adjustedPacket); - const playlist = trackData.playlist; - if (playlist.currentSegmentStartTimestamp === null) { playlist.currentSegmentStartTimestamp = adjustedPacket.timestamp; } else if (!playlist.currentSegmentStartTimestampIsFixed) { @@ -621,18 +639,17 @@ export class HlsMuxer extends Muxer { packet: EncodedPacket, meta?: EncodedAudioChunkMetadata, ) { - const release = await this.mutex.acquire(); + const trackData = this.getAudioTrackData(track, meta); + const playlist = trackData.playlist; + + const release = await playlist.mutex.acquire(); try { - const trackData = this.getAudioTrackData(track, meta); - const timestamp = this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); const adjustedPacket = packet.clone({ timestamp }); trackData.packets.push(adjustedPacket); - const playlist = trackData.playlist; - if (playlist.currentSegmentStartTimestamp === null) { playlist.currentSegmentStartTimestamp = adjustedPacket.timestamp; } else if (!playlist.currentSegmentStartTimestampIsFixed) { @@ -1402,29 +1419,35 @@ export class HlsMuxer extends Muxer { this.format._options.onMaster?.(masterPlaylistText); - let writer: Writer; - if (this.numWrittenMasterPlaylists === 0) { - // For the first master playlist write, we use the normal root writer getter, so that the target returned by - // Output.target emits valid write events. - writer = await this.output._getRootWriter(); - } else { - // For subsequent master playlist writes, we *must* obtain a different target in order to overwrite - // the file. - const target = await this.output._getTarget({ - path: pathedTarget.rootPath, - isRoot: true, - mimeType: HLS_MIME_TYPE, - }); - writer = new Writer(target); - writer.start(); + const release = await this.mutex.acquire(); + + try { + let writer: Writer; + if (this.numWrittenMasterPlaylists === 0) { + // For the first master playlist write, we use the normal root writer getter, so that the target + // returned by Output.target emits valid write events. + writer = await this.output._getRootWriter(); + } else { + // For subsequent master playlist writes, we *must* obtain a different target in order to overwrite + // the file. + const target = await this.output._getTarget({ + path: pathedTarget.rootPath, + isRoot: true, + mimeType: HLS_MIME_TYPE, + }); + writer = new Writer(target); + writer.start(); + } + + writer.write(textEncoder.encode(masterPlaylistText)); + + await writer.flush(); + await writer.finalize(); + + this.numWrittenMasterPlaylists++; + } finally { + release(); } - - writer.write(textEncoder.encode(masterPlaylistText)); - - await writer.flush(); - await writer.finalize(); - - this.numWrittenMasterPlaylists++; } private async tryWriteMasterPlaylist() { @@ -1441,8 +1464,8 @@ export class HlsMuxer extends Muxer { } async finalize() { - assert(this.output._target instanceof PathedTarget); - const release = await this.mutex.acquire(); + const releases = await Promise.all(this.playlists.map(p => p.mutex.acquire())); + releases.forEach(release => release()); for (const trackData of this.trackDatas) { trackData.closed = true; @@ -1455,8 +1478,6 @@ export class HlsMuxer extends Muxer { if (!this.isLive) { await this.writeMasterPlaylist(); } - - release(); } } diff --git a/src/media-source.ts b/src/media-source.ts index 7806a8b..832cb6a 100644 --- a/src/media-source.ts +++ b/src/media-source.ts @@ -267,6 +267,8 @@ class VideoEncoderWrapper { */ private error: Error | null = null; + private lastMuxerPromise: Promise = Promise.resolve(); + constructor(private source: VideoSource, private encodingConfig: VideoEncodingConfig) { const sizeChangeBehavior = encodingConfig.sizeChangeBehavior ?? 'deny'; if (['fill', 'contain', 'cover'].includes(sizeChangeBehavior) && encodingConfig.transform?.fit !== undefined) { @@ -648,7 +650,7 @@ class VideoEncoderWrapper { } } - await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure + await this.lastMuxerPromise; // Allow the writer to apply backpressure } } finally { for (const sample of samplesToEncode) { @@ -711,10 +713,11 @@ class VideoEncoderWrapper { maybeEnsureIsKeyPacket(this.source._connectedTrack!, packet); this.encodingConfig.onEncodedPacket?.(packet, meta); - void this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta) - .catch((error) => { - this.error ??= error; - }); + this.lastMuxerPromise + = this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta) + .catch((error) => { + this.error ??= error; + }); }; await this.customEncoder.init(); @@ -802,10 +805,11 @@ class VideoEncoderWrapper { maybeEnsureIsKeyPacket(this.source._connectedTrack!, packet); this.encodingConfig.onEncodedPacket?.(packet, meta); - void this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta) - .catch((error) => { - this.error ??= error; - }); + this.lastMuxerPromise + = this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta) + .catch((error) => { + this.error ??= error; + }); }; const stack = new Error('Encoding error').stack; @@ -1787,6 +1791,7 @@ class AudioEncoderWrapper { * So, we keep track of the encoder error and throw it as soon as we get the chance. */ private error: Error | null = null; + private lastMuxerPromise: Promise = Promise.resolve(); constructor(private source: AudioSource, private encodingConfig: AudioEncodingConfig) {} @@ -1950,7 +1955,7 @@ class AudioEncoderWrapper { await promise; } - await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure + await this.lastMuxerPromise; // Allow the writer to apply backpressure } else if (this.isPcmEncoder) { await this.doPcmEncoding(audioSample, shouldClose); } else { @@ -1967,7 +1972,7 @@ class AudioEncoderWrapper { await new Promise(resolve => this.encoder!.addEventListener('dequeue', resolve, { once: true })); } - await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure + await this.lastMuxerPromise; // Allow the writer to apply backpressure } } finally { if (shouldClose) { @@ -2080,10 +2085,11 @@ class AudioEncoderWrapper { } this.encodingConfig.onEncodedPacket?.(packet, meta); - void this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta) - .catch((error) => { - this.error ??= error; - }); + this.lastMuxerPromise + = this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta) + .catch((error) => { + this.error ??= error; + }); }; await this.customEncoder.init(); @@ -2145,10 +2151,11 @@ class AudioEncoderWrapper { }); this.encodingConfig.onEncodedPacket?.(packet, meta); - void this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta) - .catch((error) => { - this.error ??= error; - }); + this.lastMuxerPromise + = this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta) + .catch((error) => { + this.error ??= error; + }); }, error: (error) => { error.stack = stack; // Provide a more useful stack trace @@ -2803,6 +2810,8 @@ export class TextSubtitleSource extends SubtitleSource { private _parser: SubtitleParser; /** @internal */ private _error: Error | null = null; + /** @internal */ + private _lastMuxerPromise: Promise = Promise.resolve(); /** Creates a new {@link TextSubtitleSource} where added text chunks are in the specified `codec`. */ constructor(codec: SubtitleCodec) { @@ -2811,10 +2820,11 @@ export class TextSubtitleSource extends SubtitleSource { this._parser = new SubtitleParser({ codec, output: (cue, metadata) => { - void this._connectedTrack?.output._muxer.addSubtitleCue(this._connectedTrack, cue, metadata) - .catch((error) => { - this._error ??= error; - }); + this._lastMuxerPromise + = this._connectedTrack!.output._muxer.addSubtitleCue(this._connectedTrack!, cue, metadata) + .catch((error) => { + this._error ??= error; + }); }, }); } @@ -2836,7 +2846,7 @@ export class TextSubtitleSource extends SubtitleSource { this._ensureValidAdd(); this._parser.parse(text); - return this._connectedTrack!.output._muxer.mutex.currentPromise; + return this._lastMuxerPromise; // Allow the writer to apply backpressure } /** @internal */