diff --git a/src/adts/adts-demuxer.ts b/src/adts/adts-demuxer.ts index e85dfe2..cb7a669 100644 --- a/src/adts/adts-demuxer.ts +++ b/src/adts/adts-demuxer.ts @@ -12,7 +12,7 @@ import { Input } from '../input'; import { InputAudioTrack, InputAudioTrackBacking } from '../input-track'; import { assert, - AsyncMutex4, + AsyncMutex, binarySearchLessOrEqual, Bitstream, MaybeRelevantPromise, @@ -43,7 +43,7 @@ export class AdtsDemuxer extends Demuxer { tracks: InputAudioTrack[] = []; - readingMutex = new AsyncMutex4(); + readingMutex = new AsyncMutex(); lastSampleLoaded = false; lastLoadedPos = 0; nextTimestampInSamples = 0; diff --git a/src/adts/adts-muxer.ts b/src/adts/adts-muxer.ts index 395a998..3e93b4b 100644 --- a/src/adts/adts-muxer.ts +++ b/src/adts/adts-muxer.ts @@ -47,60 +47,57 @@ export class AdtsMuxer extends Muxer { ) { // https://wiki.multimedia.cx/index.php/ADTS (last visited: 2025/08/17) - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); + this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); - if (!this.audioSpecificConfig) { - validateAudioChunkMetadata(meta); + if (!this.audioSpecificConfig) { + validateAudioChunkMetadata(meta); - const description = meta?.decoderConfig?.description; - assert(description); + const description = meta?.decoderConfig?.description; + assert(description); - this.audioSpecificConfig = parseAacAudioSpecificConfig(toUint8Array(description)); + this.audioSpecificConfig = parseAacAudioSpecificConfig(toUint8Array(description)); - const { objectType, frequencyIndex, channelConfiguration } = this.audioSpecificConfig; - const profile = objectType - 1; + const { objectType, frequencyIndex, channelConfiguration } = this.audioSpecificConfig; + const profile = objectType - 1; - this.headerBitstream.writeBits(12, 0b1111_11111111); // Syncword - this.headerBitstream.writeBits(1, 0); // MPEG Version - this.headerBitstream.writeBits(2, 0); // Layer - this.headerBitstream.writeBits(1, 1); // Protection absence - this.headerBitstream.writeBits(2, profile); // Profile - this.headerBitstream.writeBits(4, frequencyIndex); // MPEG-4 Sampling Frequency Index - this.headerBitstream.writeBits(1, 0); // Private bit - this.headerBitstream.writeBits(3, channelConfiguration); // MPEG-4 Channel Configuration - this.headerBitstream.writeBits(1, 0); // Originality - this.headerBitstream.writeBits(1, 0); // Home - this.headerBitstream.writeBits(1, 0); // Copyright ID bit - this.headerBitstream.writeBits(1, 0); // Copyright ID start - this.headerBitstream.skipBits(13); // Frame length - this.headerBitstream.writeBits(11, 0x7ff); // Buffer fullness - this.headerBitstream.writeBits(2, 0); // Number of AAC frames minus 1 - // Omit CRC check - } - - const frameLength = packet.data.byteLength + this.header.byteLength; - this.headerBitstream.pos = 30; - this.headerBitstream.writeBits(13, frameLength); - - const startPos = this.writer.getPos(); - this.writer.write(this.header); - this.writer.write(packet.data); - - if (this.format._options.onFrame) { - const frameBytes = new Uint8Array(frameLength); - frameBytes.set(this.header, 0); - frameBytes.set(packet.data, this.header.byteLength); - - this.format._options.onFrame(frameBytes, startPos); - } - - await this.writer.flush(); - } finally { - release(); + this.headerBitstream.writeBits(12, 0b1111_11111111); // Syncword + this.headerBitstream.writeBits(1, 0); // MPEG Version + this.headerBitstream.writeBits(2, 0); // Layer + this.headerBitstream.writeBits(1, 1); // Protection absence + this.headerBitstream.writeBits(2, profile); // Profile + this.headerBitstream.writeBits(4, frequencyIndex); // MPEG-4 Sampling Frequency Index + this.headerBitstream.writeBits(1, 0); // Private bit + this.headerBitstream.writeBits(3, channelConfiguration); // MPEG-4 Channel Configuration + this.headerBitstream.writeBits(1, 0); // Originality + this.headerBitstream.writeBits(1, 0); // Home + this.headerBitstream.writeBits(1, 0); // Copyright ID bit + this.headerBitstream.writeBits(1, 0); // Copyright ID start + this.headerBitstream.skipBits(13); // Frame length + this.headerBitstream.writeBits(11, 0x7ff); // Buffer fullness + this.headerBitstream.writeBits(2, 0); // Number of AAC frames minus 1 + // Omit CRC check } + + const frameLength = packet.data.byteLength + this.header.byteLength; + this.headerBitstream.pos = 30; + this.headerBitstream.writeBits(13, frameLength); + + const startPos = this.writer.getPos(); + this.writer.write(this.header); + this.writer.write(packet.data); + + if (this.format._options.onFrame) { + const frameBytes = new Uint8Array(frameLength); + frameBytes.set(this.header, 0); + frameBytes.set(packet.data, this.header.byteLength); + + this.format._options.onFrame(frameBytes, startPos); + } + + await this.writer.flush(); } async addSubtitleCue() { diff --git a/src/cursors.ts b/src/cursors.ts index cf56d59..d0232f3 100644 --- a/src/cursors.ts +++ b/src/cursors.ts @@ -12,7 +12,7 @@ import { InputDisposedError } from './input'; import { InputAudioTrack, InputTrack, InputVideoTrack } from './input-track'; import { assert, - AsyncMutex4, + AsyncMutex, AsyncMutexLock, CallSerializer2, defer, @@ -469,7 +469,7 @@ export abstract class SampleCursor< private _transform: SampleTransformer; private _autoClose: boolean; - private _mutex = new AsyncMutex4(); + private _mutex = new AsyncMutex(); private _packetReader: PacketReader; private _packetCursor: PacketCursor; diff --git a/src/flac/flac-demuxer.ts b/src/flac/flac-demuxer.ts index 3bd9e83..ac2fb15 100644 --- a/src/flac/flac-demuxer.ts +++ b/src/flac/flac-demuxer.ts @@ -12,7 +12,7 @@ import { Input } from '../input'; import { InputAudioTrack, InputAudioTrackBacking } from '../input-track'; import { assert, - AsyncMutex4, + AsyncMutex, binarySearchLessOrEqual, Bitstream, MaybeRelevantPromise, @@ -79,7 +79,7 @@ export class FlacDemuxer extends Demuxer { lastLoadedPos: number | null = null; blockingBit: number | null = null; - readingMutex = new AsyncMutex4(); + readingMutex = new AsyncMutex(); lastSampleLoaded = false; constructor(input: Input) { diff --git a/src/flac/flac-muxer.ts b/src/flac/flac-muxer.ts index b444f0b..0559cf7 100644 --- a/src/flac/flac-muxer.ts +++ b/src/flac/flac-muxer.ts @@ -205,7 +205,8 @@ export class FlacMuxer extends Muxer { packet: EncodedPacket, meta?: EncodedAudioChunkMetadata, ): Promise { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; validateAudioChunkMetadata(meta); @@ -213,62 +214,58 @@ export class FlacMuxer extends Muxer { assert(meta.decoderConfig); assert(meta.decoderConfig.description); - try { - this.validateAndNormalizeTimestamp( - track, - packet.timestamp, - packet.type === 'key', - ); + this.validateAndNormalizeTimestamp( + track, + packet.timestamp, + packet.type === 'key', + ); - if (this.sampleRate === null) { - this.sampleRate = meta.decoderConfig.sampleRate; - } - - if (this.channels === null) { - this.channels = meta.decoderConfig.numberOfChannels; - } - - if (this.bitsPerSample === null) { - const descriptionBitstream = new Bitstream( - toUint8Array(meta.decoderConfig.description), - ); - // skip 'fLaC' + block size + frame size + sample rate + number of channels - // See demuxer for the exact structure - descriptionBitstream.skipBits(103 + 64); - const bitsPerSample = descriptionBitstream.readBits(5) + 1; - this.bitsPerSample = bitsPerSample; - } - - if (!this.metadataWritten) { - this.writeVorbisCommentAndPictureBlock(); - } - - const slice = FileSlice.tempFromBytes(packet.data); - readBytes(slice, 2); - const bytes = readBytes(slice, 2); - const bitstream = new Bitstream(bytes); - const blockSizeOrUncommon = getBlockSizeOrUncommon(bitstream.readBits(4)); - if (blockSizeOrUncommon === null) { - throw new Error('Invalid FLAC frame: Invalid block size.'); - } - - readCodedNumber(slice); // num - const blockSize = readBlockSize(slice, blockSizeOrUncommon); - - this.blockSizes.push(blockSize); - this.frameSizes.push(packet.data.length); - - const startPos = this.writer.getPos(); - this.writer.write(packet.data); - - if (this.format._options.onFrame) { - this.format._options.onFrame(packet.data, startPos); - } - - await this.writer.flush(); - } finally { - release(); + if (this.sampleRate === null) { + this.sampleRate = meta.decoderConfig.sampleRate; } + + if (this.channels === null) { + this.channels = meta.decoderConfig.numberOfChannels; + } + + if (this.bitsPerSample === null) { + const descriptionBitstream = new Bitstream( + toUint8Array(meta.decoderConfig.description), + ); + // skip 'fLaC' + block size + frame size + sample rate + number of channels + // See demuxer for the exact structure + descriptionBitstream.skipBits(103 + 64); + const bitsPerSample = descriptionBitstream.readBits(5) + 1; + this.bitsPerSample = bitsPerSample; + } + + if (!this.metadataWritten) { + this.writeVorbisCommentAndPictureBlock(); + } + + const slice = FileSlice.tempFromBytes(packet.data); + readBytes(slice, 2); + const bytes = readBytes(slice, 2); + const bitstream = new Bitstream(bytes); + const blockSizeOrUncommon = getBlockSizeOrUncommon(bitstream.readBits(4)); + if (blockSizeOrUncommon === null) { + throw new Error('Invalid FLAC frame: Invalid block size.'); + } + + readCodedNumber(slice); // num + const blockSize = readBlockSize(slice, blockSizeOrUncommon); + + this.blockSizes.push(blockSize); + this.frameSizes.push(packet.data.length); + + const startPos = this.writer.getPos(); + this.writer.write(packet.data); + + if (this.format._options.onFrame) { + this.format._options.onFrame(packet.data, startPos); + } + + await this.writer.flush(); } override addSubtitleCue(): Promise { @@ -276,7 +273,8 @@ export class FlacMuxer extends Muxer { } async finalize(): Promise { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; let minimumBlockSize = Infinity; let maximumBlockSize = 0; @@ -314,7 +312,5 @@ export class FlacMuxer extends Muxer { bitsPerSample: this.bitsPerSample, totalSamples, }); - - release(); } } diff --git a/src/isobmff/isobmff-muxer.ts b/src/isobmff/isobmff-muxer.ts index fd73b50..d2d1fb7 100644 --- a/src/isobmff/isobmff-muxer.ts +++ b/src/isobmff/isobmff-muxer.ts @@ -188,7 +188,8 @@ export class IsobmffMuxer extends Muxer { } async start() { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; const holdsAvc = this.output._tracks.some(x => x.type === 'video' && x.source._codec === 'avc'); @@ -238,8 +239,6 @@ export class IsobmffMuxer extends Muxer { } await this.writer.flush(); - - release(); } private allTracksAreKnown() { @@ -458,73 +457,67 @@ export class IsobmffMuxer extends Muxer { } async addEncodedVideoPacket(track: OutputVideoTrack, packet: EncodedPacket, meta?: EncodedVideoChunkMetadata) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - const trackData = this.getVideoTrackData(track, packet, meta); + const trackData = this.getVideoTrackData(track, packet, meta); - let packetData = packet.data; - if (trackData.info.requiresAnnexBTransformation) { - const nalUnits = findNalUnitsInAnnexB(packetData); - if (nalUnits.length === 0) { - // It's not valid Annex B data - throw new Error( - 'Failed to transform packet data. Make sure all packets are provided in Annex B format, as' - + ' specified in ITU-T-REC-H.264 and ITU-T-REC-H.265.', - ); - } - - // We don't strip things like SPS or PPS NALUs here, mainly because they can also appear in the middle - // of a stream and potentially modify the parameters of it. So, let's just leave them in to be sure. - packetData = concatNalUnitsInLengthPrefixed(nalUnits, 4); + let packetData = packet.data; + if (trackData.info.requiresAnnexBTransformation) { + const nalUnits = findNalUnitsInAnnexB(packetData); + if (nalUnits.length === 0) { + // It's not valid Annex B data + throw new Error( + 'Failed to transform packet data. Make sure all packets are provided in Annex B format, as' + + ' specified in ITU-T-REC-H.264 and ITU-T-REC-H.265.', + ); } - const timestamp = this.validateAndNormalizeTimestamp( - trackData.track, - packet.timestamp, - packet.type === 'key', - ); - const internalSample = this.createSampleForTrack( - trackData, - packetData, - timestamp, - packet.duration, - packet.type, - ); - - await this.registerSample(trackData, internalSample); - } finally { - release(); + // We don't strip things like SPS or PPS NALUs here, mainly because they can also appear in the middle + // of a stream and potentially modify the parameters of it. So, let's just leave them in to be sure. + packetData = concatNalUnitsInLengthPrefixed(nalUnits, 4); } + + const timestamp = this.validateAndNormalizeTimestamp( + trackData.track, + packet.timestamp, + packet.type === 'key', + ); + const internalSample = this.createSampleForTrack( + trackData, + packetData, + timestamp, + packet.duration, + packet.type, + ); + + await this.registerSample(trackData, internalSample); } async addEncodedAudioPacket(track: OutputAudioTrack, packet: EncodedPacket, meta?: EncodedAudioChunkMetadata) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - const trackData = this.getAudioTrackData(track, meta); + const trackData = this.getAudioTrackData(track, meta); - const timestamp = this.validateAndNormalizeTimestamp( - trackData.track, - packet.timestamp, - packet.type === 'key', - ); - const internalSample = this.createSampleForTrack( - trackData, - packet.data, - timestamp, - packet.duration, - packet.type, - ); + const timestamp = this.validateAndNormalizeTimestamp( + trackData.track, + packet.timestamp, + packet.type === 'key', + ); + const internalSample = this.createSampleForTrack( + trackData, + packet.data, + timestamp, + packet.duration, + packet.type, + ); - if (trackData.info.requiresPcmTransformation) { - await this.maybePadWithSilence(trackData, timestamp); - } - - await this.registerSample(trackData, internalSample); - } finally { - release(); + if (trackData.info.requiresPcmTransformation) { + await this.maybePadWithSilence(trackData, timestamp); } + + await this.registerSample(trackData, internalSample); } private async maybePadWithSilence(trackData: IsobmffAudioTrackData, untilTimestamp: number) { @@ -559,21 +552,18 @@ export class IsobmffMuxer extends Muxer { } async addSubtitleCue(track: OutputSubtitleTrack, cue: SubtitleCue, meta?: SubtitleMetadata) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - const trackData = this.getSubtitleTrackData(track, meta); + const trackData = this.getSubtitleTrackData(track, meta); - this.validateAndNormalizeTimestamp(trackData.track, cue.timestamp, true); + this.validateAndNormalizeTimestamp(trackData.track, cue.timestamp, true); - if (track.source._codec === 'webvtt') { - trackData.cueQueue.push(cue); - await this.processWebVTTCues(trackData, cue.timestamp); - } else { - // TODO - } - } finally { - release(); + if (track.source._codec === 'webvtt') { + trackData.cueQueue.push(cue); + await this.processWebVTTCues(trackData, cue.timestamp); + } else { + // TODO } } @@ -1176,7 +1166,8 @@ export class IsobmffMuxer extends Muxer { // eslint-disable-next-line @typescript-eslint/no-misused-promises override async onTrackClose(track: OutputTrack) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; if (track.type === 'subtitle' && track.source._codec === 'webvtt') { const trackData = this.trackDatas.find(x => x.track === track) as IsobmffSubtitleTrackData; @@ -1193,13 +1184,12 @@ export class IsobmffMuxer extends Muxer { // Since a track is now closed, we may be able to write out chunks that were previously waiting await this.interleaveSamples(); } - - release(); } /** Finalizes the file, making it ready for use. Must be called after all video and audio chunks have been added. */ async finalize() { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; this.allTracksKnown.resolve(); @@ -1339,7 +1329,5 @@ export class IsobmffMuxer extends Muxer { this.format._options.onMoov(data, start); } } - - release(); } } diff --git a/src/matroska/matroska-muxer.ts b/src/matroska/matroska-muxer.ts index 73af407..57ab89c 100644 --- a/src/matroska/matroska-muxer.ts +++ b/src/matroska/matroska-muxer.ts @@ -157,7 +157,8 @@ export class MatroskaMuxer extends Muxer { } async start() { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; this.writeEBMLHeader(); @@ -165,8 +166,6 @@ export class MatroskaMuxer extends Muxer { this.createCues(); await this.writer.flush(); - - release(); } private writeEBMLHeader() { @@ -837,88 +836,79 @@ export class MatroskaMuxer extends Muxer { } async addEncodedVideoPacket(track: OutputVideoTrack, packet: EncodedPacket, meta?: EncodedVideoChunkMetadata) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - const trackData = this.getVideoTrackData(track, packet, meta); + const trackData = this.getVideoTrackData(track, packet, meta); - const isKeyFrame = packet.type === 'key'; - let timestamp = this.validateAndNormalizeTimestamp(trackData.track, packet.timestamp, isKeyFrame); - let duration = packet.duration; + const isKeyFrame = packet.type === 'key'; + let timestamp = this.validateAndNormalizeTimestamp(trackData.track, packet.timestamp, isKeyFrame); + let duration = packet.duration; - if (track.metadata.frameRate !== undefined) { - // Constrain the time values to the frame rate - timestamp = roundToMultiple(timestamp, 1 / track.metadata.frameRate); - duration = roundToMultiple(duration, 1 / track.metadata.frameRate); - } - - const additions = trackData.info.alphaMode - ? packet.sideData.alpha ?? null - : null; - - const videoChunk = this.createInternalChunk(packet.data, timestamp, duration, packet.type, additions); - if (track.source._codec === 'vp9') this.fixVP9ColorSpace(trackData, videoChunk); - - trackData.chunkQueue.push(videoChunk); - await this.interleaveChunks(); - } finally { - release(); + if (track.metadata.frameRate !== undefined) { + // Constrain the time values to the frame rate + timestamp = roundToMultiple(timestamp, 1 / track.metadata.frameRate); + duration = roundToMultiple(duration, 1 / track.metadata.frameRate); } + + const additions = trackData.info.alphaMode + ? packet.sideData.alpha ?? null + : null; + + const videoChunk = this.createInternalChunk(packet.data, timestamp, duration, packet.type, additions); + if (track.source._codec === 'vp9') this.fixVP9ColorSpace(trackData, videoChunk); + + trackData.chunkQueue.push(videoChunk); + await this.interleaveChunks(); } async addEncodedAudioPacket(track: OutputAudioTrack, packet: EncodedPacket, meta?: EncodedAudioChunkMetadata) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - const trackData = this.getAudioTrackData(track, meta); + const trackData = this.getAudioTrackData(track, meta); - const isKeyFrame = packet.type === 'key'; - const timestamp = this.validateAndNormalizeTimestamp(trackData.track, packet.timestamp, isKeyFrame); - const audioChunk = this.createInternalChunk(packet.data, timestamp, packet.duration, packet.type); + const isKeyFrame = packet.type === 'key'; + const timestamp = this.validateAndNormalizeTimestamp(trackData.track, packet.timestamp, isKeyFrame); + const audioChunk = this.createInternalChunk(packet.data, timestamp, packet.duration, packet.type); - trackData.chunkQueue.push(audioChunk); - await this.interleaveChunks(); - } finally { - release(); - } + trackData.chunkQueue.push(audioChunk); + await this.interleaveChunks(); } async addSubtitleCue(track: OutputSubtitleTrack, cue: SubtitleCue, meta?: SubtitleMetadata) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - const trackData = this.getSubtitleTrackData(track, meta); + const trackData = this.getSubtitleTrackData(track, meta); - const timestamp = this.validateAndNormalizeTimestamp(trackData.track, cue.timestamp, true); + const timestamp = this.validateAndNormalizeTimestamp(trackData.track, cue.timestamp, true); - let bodyText = cue.text; - const timestampMs = Math.round(timestamp * 1000); + let bodyText = cue.text; + const timestampMs = Math.round(timestamp * 1000); - // Replace in-body timestamps so that they're relative to the cue start time - inlineTimestampRegex.lastIndex = 0; - bodyText = bodyText.replace(inlineTimestampRegex, (match) => { - const time = parseSubtitleTimestamp(match.slice(1, -1)); - const offsetTime = time - timestampMs; + // Replace in-body timestamps so that they're relative to the cue start time + inlineTimestampRegex.lastIndex = 0; + bodyText = bodyText.replace(inlineTimestampRegex, (match) => { + const time = parseSubtitleTimestamp(match.slice(1, -1)); + const offsetTime = time - timestampMs; - return `<${formatSubtitleTimestamp(offsetTime)}>`; - }); + return `<${formatSubtitleTimestamp(offsetTime)}>`; + }); - const body = textEncoder.encode(bodyText); - const additions = `${cue.settings ?? ''}\n${cue.identifier ?? ''}\n${cue.notes ?? ''}`; + const body = textEncoder.encode(bodyText); + const additions = `${cue.settings ?? ''}\n${cue.identifier ?? ''}\n${cue.notes ?? ''}`; - const subtitleChunk = this.createInternalChunk( - body, - timestamp, - cue.duration, - 'key', - additions.trim() ? textEncoder.encode(additions) : null, - ); + const subtitleChunk = this.createInternalChunk( + body, + timestamp, + cue.duration, + 'key', + additions.trim() ? textEncoder.encode(additions) : null, + ); - trackData.chunkQueue.push(subtitleChunk); - await this.interleaveChunks(); - } finally { - release(); - } + trackData.chunkQueue.push(subtitleChunk); + await this.interleaveChunks(); } private async interleaveChunks(isFinalCall = false) { @@ -1208,7 +1198,8 @@ export class MatroskaMuxer extends Muxer { // eslint-disable-next-line @typescript-eslint/no-misused-promises override async onTrackClose() { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; if (this.allTracksAreKnown()) { this.allTracksKnown.resolve(); @@ -1216,15 +1207,12 @@ export class MatroskaMuxer extends Muxer { // Since a track is now closed, we may be able to write out chunks that were previously waiting await this.interleaveChunks(); - - release(); } /** Finalizes the file, making it ready for use. Must be called after all media chunks have been added. */ async finalize() { - const release = await this.mutex.acquire(); - - this.allTracksKnown.resolve(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; if (!this.segment) { this.createSegment(); @@ -1261,7 +1249,5 @@ export class MatroskaMuxer extends Muxer { this.writer.seek(endPos); } - - release(); } } diff --git a/src/media-source.ts b/src/media-source.ts index 6a73abf..b215f40 100644 --- a/src/media-source.ts +++ b/src/media-source.ts @@ -405,7 +405,7 @@ class VideoEncoderWrapper { } } - await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure + await this.muxer!.mutex.waitForUnlock(); // Allow the writer to apply backpressure } finally { if (shouldClose) { // Make sure it's always closed, even if there was an error @@ -1449,7 +1449,7 @@ class AudioEncoderWrapper { await promise; } - await this.muxer!.mutex.currentPromise; // Allow the writer to apply backpressure + await this.muxer!.mutex.waitForUnlock(); // Allow the writer to apply backpressure } else if (this.isPcmEncoder) { await this.doPcmEncoding(audioSample, shouldClose); } else { @@ -1466,7 +1466,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.muxer!.mutex.waitForUnlock(); // Allow the writer to apply backpressure } } finally { if (shouldClose) { @@ -2320,7 +2320,7 @@ export class TextSubtitleSource extends SubtitleSource { this._ensureValidAdd(); this._parser.parse(text); - return this._connectedTrack!.output._muxer.mutex.currentPromise; + return this._connectedTrack!.output._muxer.mutex.waitForUnlock(); } /** @internal */ diff --git a/src/misc.ts b/src/misc.ts index bb625df..aef839c 100644 --- a/src/misc.ts +++ b/src/misc.ts @@ -246,24 +246,6 @@ export const isAllowSharedBufferSource = (x: unknown) => { ); }; -export class AsyncMutex { - currentPromise = Promise.resolve(); - - async acquire() { - let resolver: () => void; - const nextPromise = new Promise((resolve) => { - resolver = resolve; - }); - - const currentPromiseAlias = this.currentPromise; - this.currentPromise = nextPromise; - - await currentPromiseAlias; - - return resolver!; - } -} - export const bytesToHexString = (bytes: Uint8Array) => { return [...bytes].map(x => x.toString(16).padStart(2, '0')).join(''); }; @@ -900,40 +882,7 @@ export class ResultValue { } } -export class AsyncMutex2 { - locked = false; - promise = Promise.resolve(); - - lock() { - if (this.locked) { - throw new Error('Mutex already locked.'); - } - - this.locked = true; - - const { promise, resolve } = promiseWithResolvers(); - this.promise = promise; - - let released = false; - - return { - release: () => { - if (released) { - return; - } - released = true; - - resolve(); - this.locked = false; - }, - [Symbol.dispose]() { - this.release(); - }, - }; - } -} - -export class AsyncMutex4 { +export class AsyncMutex { locked = false; private resolverQueue: (() => void)[] = []; @@ -958,13 +907,19 @@ export class AsyncMutex4 { this.locked = false; } } + + async waitForUnlock() { + const lock = this.lock(); + await lock.ready; + lock.release(); + } } export class AsyncMutexLock implements Disposable { private released = false; constructor( - private readonly mutex: AsyncMutex4, + private readonly mutex: AsyncMutex, public readonly pending: boolean, public readonly ready: Promise | null, ) {} diff --git a/src/mp3/mp3-demuxer.ts b/src/mp3/mp3-demuxer.ts index 4c982ff..1809eb3 100644 --- a/src/mp3/mp3-demuxer.ts +++ b/src/mp3/mp3-demuxer.ts @@ -13,7 +13,7 @@ import { InputAudioTrack, InputAudioTrackBacking } from '../input-track'; import { DEFAULT_TRACK_DISPOSITION, MetadataTags } from '../metadata'; import { assert, - AsyncMutex4, + AsyncMutex, binarySearchLessOrEqual, MaybeRelevantPromise, ResultValue, @@ -49,7 +49,7 @@ export class Mp3Demuxer extends Demuxer { tracks: InputAudioTrack[] = []; - readingMutex = new AsyncMutex4(); + readingMutex = new AsyncMutex(); lastSampleLoaded = false; lastLoadedPos = 0; nextTimestampInSamples = 0; diff --git a/src/mp3/mp3-muxer.ts b/src/mp3/mp3-muxer.ts index f2c86e8..0f36096 100644 --- a/src/mp3/mp3-muxer.ts +++ b/src/mp3/mp3-muxer.ts @@ -53,70 +53,67 @@ export class Mp3Muxer extends Muxer { track: OutputAudioTrack, packet: EncodedPacket, ) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - const writeXingHeader = this.format._options.xingHeader !== false; + const writeXingHeader = this.format._options.xingHeader !== false; - if (!this.xingFrameData && writeXingHeader) { - const view = toDataView(packet.data); - if (view.byteLength < 4) { - throw new Error('Invalid MP3 header in sample.'); - } - - const word = view.getUint32(0, false); - const header = readFrameHeader(word, null).header; - if (!header) { - throw new Error('Invalid MP3 header in sample.'); - } - - const xingOffset = getXingOffset(header.mpegVersionId, header.channel); - if (view.byteLength >= xingOffset + 4) { - const word = view.getUint32(xingOffset, false); - const isXing = word === XING || word === INFO; - - if (isXing) { - // This is not a data frame, so let's completely ignore this sample - return; - } - } - - this.xingFrameData = { - mpegVersionId: header.mpegVersionId, - layer: header.layer, - frequencyIndex: header.frequencyIndex, - sampleRate: header.sampleRate, - channel: header.channel, - modeExtension: header.modeExtension, - copyright: header.copyright, - original: header.original, - emphasis: header.emphasis, - - frameCount: null, - fileSize: null, - toc: null, - }; - - // Write a Xing frame because this muxer doesn't make any bitrate constraints, meaning we don't know if - // this will be a constant or variable bitrate file. Therefore, always write the Xing frame. - this.xingFramePos = this.writer.getPos(); - this.mp3Writer.writeXingFrame(this.xingFrameData); - - this.frameCount++; + if (!this.xingFrameData && writeXingHeader) { + const view = toDataView(packet.data); + if (view.byteLength < 4) { + throw new Error('Invalid MP3 header in sample.'); } - this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); + const word = view.getUint32(0, false); + const header = readFrameHeader(word, null).header; + if (!header) { + throw new Error('Invalid MP3 header in sample.'); + } + + const xingOffset = getXingOffset(header.mpegVersionId, header.channel); + if (view.byteLength >= xingOffset + 4) { + const word = view.getUint32(xingOffset, false); + const isXing = word === XING || word === INFO; + + if (isXing) { + // This is not a data frame, so let's completely ignore this sample + return; + } + } + + this.xingFrameData = { + mpegVersionId: header.mpegVersionId, + layer: header.layer, + frequencyIndex: header.frequencyIndex, + sampleRate: header.sampleRate, + channel: header.channel, + modeExtension: header.modeExtension, + copyright: header.copyright, + original: header.original, + emphasis: header.emphasis, + + frameCount: null, + fileSize: null, + toc: null, + }; + + // Write a Xing frame because this muxer doesn't make any bitrate constraints, meaning we don't know if + // this will be a constant or variable bitrate file. Therefore, always write the Xing frame. + this.xingFramePos = this.writer.getPos(); + this.mp3Writer.writeXingFrame(this.xingFrameData); - this.writer.write(packet.data); this.frameCount++; + } - await this.writer.flush(); + this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); - if (writeXingHeader) { - this.framePositions.push(this.writer.getPos()); - } - } finally { - release(); + this.writer.write(packet.data); + this.frameCount++; + + await this.writer.flush(); + + if (writeXingHeader) { + this.framePositions.push(this.writer.getPos()); } } @@ -129,7 +126,8 @@ export class Mp3Muxer extends Muxer { return; } - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; const endPos = this.writer.getPos(); @@ -160,7 +158,5 @@ export class Mp3Muxer extends Muxer { } this.writer.seek(endPos); - - release(); } } diff --git a/src/ogg/ogg-demuxer.ts b/src/ogg/ogg-demuxer.ts index bb3d027..64e5b65 100644 --- a/src/ogg/ogg-demuxer.ts +++ b/src/ogg/ogg-demuxer.ts @@ -15,7 +15,7 @@ import { InputAudioTrack, InputAudioTrackBacking } from '../input-track'; import { DEFAULT_TRACK_DISPOSITION, MetadataTags } from '../metadata'; import { assert, - AsyncMutex4, + AsyncMutex, binarySearchLessOrEqual, findLast, last, @@ -439,7 +439,7 @@ type EncodedPacketMetadata = { class OggAudioTrackBacking implements InputAudioTrackBacking { internalSampleRate: number; sequentialScanCache: EncodedPacketMetadata[] = []; - sequentialScanMutex = new AsyncMutex4(); + sequentialScanMutex = new AsyncMutex(); constructor(public bitstream: LogicalBitstream, public demuxer: OggDemuxer) { // Opus always uses a fixed sample rate for its internal calculations, even if the actual rate is different diff --git a/src/ogg/ogg-muxer.ts b/src/ogg/ogg-muxer.ts index 91ee535..c9f9e5b 100644 --- a/src/ogg/ogg-muxer.ts +++ b/src/ogg/ogg-muxer.ts @@ -256,34 +256,31 @@ export class OggMuxer extends Muxer { } async addEncodedAudioPacket(track: OutputAudioTrack, packet: EncodedPacket, meta?: EncodedAudioChunkMetadata) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - const trackData = this.getTrackData(track, meta); + const trackData = this.getTrackData(track, meta); - this.validateAndNormalizeTimestamp(trackData.track, packet.timestamp, packet.type === 'key'); + this.validateAndNormalizeTimestamp(trackData.track, packet.timestamp, packet.type === 'key'); - const currentTimestampInSamples = trackData.currentTimestampInSamples; + const currentTimestampInSamples = trackData.currentTimestampInSamples; - const { durationInSamples, vorbisBlockSize } = extractSampleMetadata( - packet.data, - trackData.codecInfo, - trackData.vorbisLastBlocksize, - ); - trackData.currentTimestampInSamples += durationInSamples; - trackData.vorbisLastBlocksize = vorbisBlockSize; + const { durationInSamples, vorbisBlockSize } = extractSampleMetadata( + packet.data, + trackData.codecInfo, + trackData.vorbisLastBlocksize, + ); + trackData.currentTimestampInSamples += durationInSamples; + trackData.vorbisLastBlocksize = vorbisBlockSize; - trackData.packetQueue.push({ - data: packet.data, - endGranulePosition: trackData.currentTimestampInSamples, - timestamp: currentTimestampInSamples / trackData.internalSampleRate, - forcePageFlush: false, - }); + trackData.packetQueue.push({ + data: packet.data, + endGranulePosition: trackData.currentTimestampInSamples, + timestamp: currentTimestampInSamples / trackData.internalSampleRate, + forcePageFlush: false, + }); - await this.interleavePages(); - } finally { - release(); - } + await this.interleavePages(); } addSubtitleCue(): never { @@ -467,7 +464,8 @@ export class OggMuxer extends Muxer { // eslint-disable-next-line @typescript-eslint/no-misused-promises override async onTrackClose() { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; if (this.allTracksAreKnown()) { this.allTracksKnown.resolve(); @@ -475,12 +473,11 @@ export class OggMuxer extends Muxer { // Since a track is now closed, we may be able to write out chunks that were previously waiting await this.interleavePages(); - - release(); } async finalize() { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; this.allTracksKnown.resolve(); @@ -491,7 +488,5 @@ export class OggMuxer extends Muxer { this.writePage(trackData, true); } } - - release(); } } diff --git a/src/output.ts b/src/output.ts index 9362805..17cf0d0 100644 --- a/src/output.ts +++ b/src/output.ts @@ -402,14 +402,13 @@ export class Output< this.state = 'started'; this._writer.start(); - const release = await this._mutex.acquire(); + using lock = this._mutex.lock(); + if (lock.pending) await lock.ready; await this._muxer.start(); const promises = this._tracks.map(track => track.source._start()); await Promise.all(promises); - - release(); })(); } @@ -440,14 +439,13 @@ export class Output< return this._cancelPromise = (async () => { this.state = 'canceled'; - const release = await this._mutex.acquire(); + using lock = this._mutex.lock(); + if (lock.pending) await lock.ready; const promises = this._tracks.map(x => x.source._flushOrWaitForOngoingClose(true)); // Force close await Promise.all(promises); await this._writer.close(); - - release(); })(); } @@ -470,7 +468,8 @@ export class Output< return this._finalizePromise = (async () => { this.state = 'finalizing'; - const release = await this._mutex.acquire(); + using lock = this._mutex.lock(); + if (lock.pending) await lock.ready; const promises = this._tracks.map(x => x.source._flushOrWaitForOngoingClose(false)); await Promise.all(promises); @@ -481,8 +480,6 @@ export class Output< await this._writer.finalize(); this.state = 'finalized'; - - release(); })(); } } diff --git a/src/wave/wave-muxer.ts b/src/wave/wave-muxer.ts index 74b6abf..96aa448 100644 --- a/src/wave/wave-muxer.ts +++ b/src/wave/wave-muxer.ts @@ -60,37 +60,34 @@ export class WaveMuxer extends Muxer { packet: EncodedPacket, meta?: EncodedAudioChunkMetadata, ) { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; - try { - if (!this.headerWritten) { - validateAudioChunkMetadata(meta); + if (!this.headerWritten) { + validateAudioChunkMetadata(meta); - assert(meta); - assert(meta.decoderConfig); + assert(meta); + assert(meta.decoderConfig); - this.writeHeader(track, meta.decoderConfig); - this.sampleRate = meta.decoderConfig.sampleRate; - this.headerWritten = true; - } - - this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); - - if (!this.isRf64 && this.writer.getPos() + packet.data.byteLength >= 2 ** 32) { - throw new Error( - 'Adding more audio data would exceed the maximum RIFF size of 4 GiB. To write larger files, use' - + ' RF64 by setting `large: true` in the WavOutputFormatOptions.', - ); - } - - this.writer.write(packet.data); - this.dataSize += packet.data.byteLength; - this.sampleCount += Math.round(packet.duration * this.sampleRate!); - - await this.writer.flush(); - } finally { - release(); + this.writeHeader(track, meta.decoderConfig); + this.sampleRate = meta.decoderConfig.sampleRate; + this.headerWritten = true; } + + this.validateAndNormalizeTimestamp(track, packet.timestamp, packet.type === 'key'); + + if (!this.isRf64 && this.writer.getPos() + packet.data.byteLength >= 2 ** 32) { + throw new Error( + 'Adding more audio data would exceed the maximum RIFF size of 4 GiB. To write larger files, use' + + ' RF64 by setting `large: true` in the WavOutputFormatOptions.', + ); + } + + this.writer.write(packet.data); + this.dataSize += packet.data.byteLength; + this.sampleCount += Math.round(packet.duration * this.sampleRate!); + + await this.writer.flush(); } async addSubtitleCue() { @@ -333,7 +330,8 @@ export class WaveMuxer extends Muxer { } async finalize() { - const release = await this.mutex.acquire(); + using lock = this.mutex.lock(); + if (lock.pending) await lock.ready; const endPos = this.writer.getPos(); @@ -365,7 +363,5 @@ export class WaveMuxer extends Muxer { } this.writer.seek(endPos); - - release(); } }