Remove older mutexes (mutices??)

This commit is contained in:
Vanilagy
2025-12-29 22:23:44 +01:00
parent 72de8da31b
commit aafb928460
15 changed files with 352 additions and 446 deletions
+2 -2
View File
@@ -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;
+44 -47
View File
@@ -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() {
+2 -2
View File
@@ -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<Sample, TransformedSample>;
private _autoClose: boolean;
private _mutex = new AsyncMutex4();
private _mutex = new AsyncMutex();
private _packetReader: PacketReader;
private _packetCursor: PacketCursor;
+2 -2
View File
@@ -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) {
+54 -58
View File
@@ -205,7 +205,8 @@ export class FlacMuxer extends Muxer {
packet: EncodedPacket,
meta?: EncodedAudioChunkMetadata,
): Promise<void> {
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<void> {
@@ -276,7 +273,8 @@ export class FlacMuxer extends Muxer {
}
async finalize(): Promise<void> {
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();
}
}
+64 -76
View File
@@ -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();
}
}
+58 -72
View File
@@ -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();
}
}
+4 -4
View File
@@ -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 */
+8 -53
View File
@@ -246,24 +246,6 @@ export const isAllowSharedBufferSource = (x: unknown) => {
);
};
export class AsyncMutex {
currentPromise = Promise.resolve();
async acquire() {
let resolver: () => void;
const nextPromise = new Promise<void>((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<T> {
}
}
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<void> | null,
) {}
+2 -2
View File
@@ -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;
+55 -59
View File
@@ -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();
}
}
+2 -2
View File
@@ -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
+23 -28
View File
@@ -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();
}
}
+6 -9
View File
@@ -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();
})();
}
}
+26 -30
View File
@@ -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();
}
}