Add per-playlist "mutices" for HLS muxer

This commit is contained in:
Vanilagy
2026-04-17 18:35:33 +02:00
parent 2ab6bbe3d1
commit 9a7120f5f4
2 changed files with 101 additions and 70 deletions
+67 -46
View File
@@ -8,7 +8,16 @@
import { MediaCodec, validateAudioChunkMetadata, validateVideoChunkMetadata } from '../codec'; import { MediaCodec, validateAudioChunkMetadata, validateVideoChunkMetadata } from '../codec';
import { EncodedAudioPacketSource, EncodedVideoPacketSource } from '../media-source'; 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 { Muxer } from '../muxer';
import { import {
Output, Output,
@@ -79,6 +88,11 @@ type Playlist = {
nextOffset: number; nextOffset: number;
info: HlsOutputSegmentInfo; info: HlsOutputSegmentInfo;
} | null; } | 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 = { type PlaylistDeclaration = {
@@ -131,6 +145,8 @@ export class HlsMuxer extends Muxer {
} }
async start(): Promise<void> { async start(): Promise<void> {
const release = await this.mutex.acquire();
const someRelative = this.output._tracks.some(t => t.metadata.isRelativeToUnixEpoch); const someRelative = this.output._tracks.some(t => t.metadata.isRelativeToUnixEpoch);
const someNotRelative = this.output._tracks.some(t => !t.metadata.isRelativeToUnixEpoch); const someNotRelative = this.output._tracks.some(t => !t.metadata.isRelativeToUnixEpoch);
if (someRelative && someNotRelative) { if (someRelative && someNotRelative) {
@@ -449,6 +465,7 @@ export class HlsMuxer extends Muxer {
mediaSequence: 0, mediaSequence: 0,
done: false, done: false,
singleFile: null, singleFile: null,
mutex: new AsyncMutex(),
}; };
this.playlists.push(playlist); this.playlists.push(playlist);
@@ -491,6 +508,8 @@ export class HlsMuxer extends Muxer {
: [], : [],
}); });
} }
release();
} }
async getMimeType(): Promise<string> { async getMimeType(): Promise<string> {
@@ -509,17 +528,17 @@ export class HlsMuxer extends Muxer {
// eslint-disable-next-line @typescript-eslint/no-misused-promises // eslint-disable-next-line @typescript-eslint/no-misused-promises
override async onTrackClose(track: OutputTrack) { 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 { 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); await this.advancePlaylist(playlist);
} finally { } finally {
release(); release();
@@ -589,18 +608,17 @@ export class HlsMuxer extends Muxer {
packet: EncodedPacket, packet: EncodedPacket,
meta?: EncodedVideoChunkMetadata, 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 { try {
const trackData = this.getVideoTrackData(track, meta);
const timestamp = this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); const timestamp = this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key');
const adjustedPacket = packet.clone({ timestamp }); const adjustedPacket = packet.clone({ timestamp });
trackData.packets.push(adjustedPacket); trackData.packets.push(adjustedPacket);
const playlist = trackData.playlist;
if (playlist.currentSegmentStartTimestamp === null) { if (playlist.currentSegmentStartTimestamp === null) {
playlist.currentSegmentStartTimestamp = adjustedPacket.timestamp; playlist.currentSegmentStartTimestamp = adjustedPacket.timestamp;
} else if (!playlist.currentSegmentStartTimestampIsFixed) { } else if (!playlist.currentSegmentStartTimestampIsFixed) {
@@ -621,18 +639,17 @@ export class HlsMuxer extends Muxer {
packet: EncodedPacket, packet: EncodedPacket,
meta?: EncodedAudioChunkMetadata, 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 { try {
const trackData = this.getAudioTrackData(track, meta);
const timestamp = this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); const timestamp = this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key');
const adjustedPacket = packet.clone({ timestamp }); const adjustedPacket = packet.clone({ timestamp });
trackData.packets.push(adjustedPacket); trackData.packets.push(adjustedPacket);
const playlist = trackData.playlist;
if (playlist.currentSegmentStartTimestamp === null) { if (playlist.currentSegmentStartTimestamp === null) {
playlist.currentSegmentStartTimestamp = adjustedPacket.timestamp; playlist.currentSegmentStartTimestamp = adjustedPacket.timestamp;
} else if (!playlist.currentSegmentStartTimestampIsFixed) { } else if (!playlist.currentSegmentStartTimestampIsFixed) {
@@ -1402,29 +1419,35 @@ export class HlsMuxer extends Muxer {
this.format._options.onMaster?.(masterPlaylistText); this.format._options.onMaster?.(masterPlaylistText);
let writer: Writer; const release = await this.mutex.acquire();
if (this.numWrittenMasterPlaylists === 0) {
// For the first master playlist write, we use the normal root writer getter, so that the target returned by try {
// Output.target emits valid write events. let writer: Writer;
writer = await this.output._getRootWriter(); if (this.numWrittenMasterPlaylists === 0) {
} else { // For the first master playlist write, we use the normal root writer getter, so that the target
// For subsequent master playlist writes, we *must* obtain a different target in order to overwrite // returned by Output.target emits valid write events.
// the file. writer = await this.output._getRootWriter();
const target = await this.output._getTarget({ } else {
path: pathedTarget.rootPath, // For subsequent master playlist writes, we *must* obtain a different target in order to overwrite
isRoot: true, // the file.
mimeType: HLS_MIME_TYPE, const target = await this.output._getTarget({
}); path: pathedTarget.rootPath,
writer = new Writer(target); isRoot: true,
writer.start(); 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() { private async tryWriteMasterPlaylist() {
@@ -1441,8 +1464,8 @@ export class HlsMuxer extends Muxer {
} }
async finalize() { async finalize() {
assert(this.output._target instanceof PathedTarget); const releases = await Promise.all(this.playlists.map(p => p.mutex.acquire()));
const release = await this.mutex.acquire(); releases.forEach(release => release());
for (const trackData of this.trackDatas) { for (const trackData of this.trackDatas) {
trackData.closed = true; trackData.closed = true;
@@ -1455,8 +1478,6 @@ export class HlsMuxer extends Muxer {
if (!this.isLive) { if (!this.isLive) {
await this.writeMasterPlaylist(); await this.writeMasterPlaylist();
} }
release();
} }
} }
+34 -24
View File
@@ -267,6 +267,8 @@ class VideoEncoderWrapper {
*/ */
private error: Error | null = null; private error: Error | null = null;
private lastMuxerPromise: Promise<void> = Promise.resolve();
constructor(private source: VideoSource, private encodingConfig: VideoEncodingConfig) { constructor(private source: VideoSource, private encodingConfig: VideoEncodingConfig) {
const sizeChangeBehavior = encodingConfig.sizeChangeBehavior ?? 'deny'; const sizeChangeBehavior = encodingConfig.sizeChangeBehavior ?? 'deny';
if (['fill', 'contain', 'cover'].includes(sizeChangeBehavior) && encodingConfig.transform?.fit !== undefined) { 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 { } finally {
for (const sample of samplesToEncode) { for (const sample of samplesToEncode) {
@@ -711,10 +713,11 @@ class VideoEncoderWrapper {
maybeEnsureIsKeyPacket(this.source._connectedTrack!, packet); maybeEnsureIsKeyPacket(this.source._connectedTrack!, packet);
this.encodingConfig.onEncodedPacket?.(packet, meta); this.encodingConfig.onEncodedPacket?.(packet, meta);
void this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta) this.lastMuxerPromise
.catch((error) => { = this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta)
this.error ??= error; .catch((error) => {
}); this.error ??= error;
});
}; };
await this.customEncoder.init(); await this.customEncoder.init();
@@ -802,10 +805,11 @@ class VideoEncoderWrapper {
maybeEnsureIsKeyPacket(this.source._connectedTrack!, packet); maybeEnsureIsKeyPacket(this.source._connectedTrack!, packet);
this.encodingConfig.onEncodedPacket?.(packet, meta); this.encodingConfig.onEncodedPacket?.(packet, meta);
void this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta) this.lastMuxerPromise
.catch((error) => { = this.muxer!.addEncodedVideoPacket(this.source._connectedTrack!, packet, meta)
this.error ??= error; .catch((error) => {
}); this.error ??= error;
});
}; };
const stack = new Error('Encoding error').stack; 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. * So, we keep track of the encoder error and throw it as soon as we get the chance.
*/ */
private error: Error | null = null; private error: Error | null = null;
private lastMuxerPromise: Promise<void> = Promise.resolve();
constructor(private source: AudioSource, private encodingConfig: AudioEncodingConfig) {} constructor(private source: AudioSource, private encodingConfig: AudioEncodingConfig) {}
@@ -1950,7 +1955,7 @@ class AudioEncoderWrapper {
await promise; 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) { } else if (this.isPcmEncoder) {
await this.doPcmEncoding(audioSample, shouldClose); await this.doPcmEncoding(audioSample, shouldClose);
} else { } else {
@@ -1967,7 +1972,7 @@ class AudioEncoderWrapper {
await new Promise(resolve => this.encoder!.addEventListener('dequeue', resolve, { once: true })); 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 { } finally {
if (shouldClose) { if (shouldClose) {
@@ -2080,10 +2085,11 @@ class AudioEncoderWrapper {
} }
this.encodingConfig.onEncodedPacket?.(packet, meta); this.encodingConfig.onEncodedPacket?.(packet, meta);
void this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta) this.lastMuxerPromise
.catch((error) => { = this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta)
this.error ??= error; .catch((error) => {
}); this.error ??= error;
});
}; };
await this.customEncoder.init(); await this.customEncoder.init();
@@ -2145,10 +2151,11 @@ class AudioEncoderWrapper {
}); });
this.encodingConfig.onEncodedPacket?.(packet, meta); this.encodingConfig.onEncodedPacket?.(packet, meta);
void this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta) this.lastMuxerPromise
.catch((error) => { = this.muxer!.addEncodedAudioPacket(this.source._connectedTrack!, packet, meta)
this.error ??= error; .catch((error) => {
}); this.error ??= error;
});
}, },
error: (error) => { error: (error) => {
error.stack = stack; // Provide a more useful stack trace error.stack = stack; // Provide a more useful stack trace
@@ -2803,6 +2810,8 @@ export class TextSubtitleSource extends SubtitleSource {
private _parser: SubtitleParser; private _parser: SubtitleParser;
/** @internal */ /** @internal */
private _error: Error | null = null; private _error: Error | null = null;
/** @internal */
private _lastMuxerPromise: Promise<void> = Promise.resolve();
/** Creates a new {@link TextSubtitleSource} where added text chunks are in the specified `codec`. */ /** Creates a new {@link TextSubtitleSource} where added text chunks are in the specified `codec`. */
constructor(codec: SubtitleCodec) { constructor(codec: SubtitleCodec) {
@@ -2811,10 +2820,11 @@ export class TextSubtitleSource extends SubtitleSource {
this._parser = new SubtitleParser({ this._parser = new SubtitleParser({
codec, codec,
output: (cue, metadata) => { output: (cue, metadata) => {
void this._connectedTrack?.output._muxer.addSubtitleCue(this._connectedTrack, cue, metadata) this._lastMuxerPromise
.catch((error) => { = this._connectedTrack!.output._muxer.addSubtitleCue(this._connectedTrack!, cue, metadata)
this._error ??= error; .catch((error) => {
}); this._error ??= error;
});
}, },
}); });
} }
@@ -2836,7 +2846,7 @@ export class TextSubtitleSource extends SubtitleSource {
this._ensureValidAdd(); this._ensureValidAdd();
this._parser.parse(text); this._parser.parse(text);
return this._connectedTrack!.output._muxer.mutex.currentPromise; return this._lastMuxerPromise; // Allow the writer to apply backpressure
} }
/** @internal */ /** @internal */