Implement ADTS demuxer, Fix MP3 demuxer race conditions

This commit is contained in:
Ally
2025-08-18 01:29:52 +02:00
parent 1ce5e108a7
commit 8da619a624
10 changed files with 560 additions and 85 deletions
+66 -64
View File
@@ -32,7 +32,7 @@ export class Mp3Demuxer extends Demuxer {
tracks: InputAudioTrack[] = [];
loadingMutex = new AsyncMutex();
readingMutex = new AsyncMutex();
lastLoadedPos = 0;
fileSize = 0;
nextTimestampInSamples = 0;
@@ -53,9 +53,8 @@ export class Mp3Demuxer extends Demuxer {
await this.loadNextChunk();
}
if (!this.firstFrameHeader) {
throw new Error('No MP3 frames found.');
}
// There has to be a frame if this demuxer got selected
assert(this.firstFrameHeader);
this.tracks = [new InputAudioTrack(new Mp3AudioTrackBacking(this))];
})();
@@ -63,30 +62,24 @@ export class Mp3Demuxer extends Demuxer {
/** Loads the next 0.5 MiB of frames. */
async loadNextChunk() {
const release = await this.loadingMutex.acquire();
assert(this.lastLoadedPos < this.fileSize);
try {
assert(this.lastLoadedPos < this.fileSize);
const chunkSize = 0.5 * 1024 * 1024; // 0.5 MiB
const endPos = Math.min(this.lastLoadedPos + chunkSize, this.fileSize);
await this.reader.reader.loadRange(this.lastLoadedPos, endPos);
const chunkSize = 0.5 * 1024 * 1024; // 0.5 MiB
const endPos = Math.min(this.lastLoadedPos + chunkSize, this.fileSize);
await this.reader.reader.loadRange(this.lastLoadedPos, endPos);
this.lastLoadedPos = endPos;
assert(this.lastLoadedPos <= this.fileSize);
this.lastLoadedPos = endPos;
assert(this.lastLoadedPos <= this.fileSize);
if (this.reader.pos === 0) {
// First time, let's see if there's an ID3 tag
const id3Tag = this.reader.readId3();
if (id3Tag) {
this.reader.pos += id3Tag.size;
}
if (this.reader.pos === 0) {
// First time, let's see if there's an ID3 tag
const id3Tag = this.reader.readId3();
if (id3Tag) {
this.reader.pos += id3Tag.size;
}
this.parseFramesFromLoadedData();
} finally {
release();
}
this.parseFramesFromLoadedData();
}
private parseFramesFromLoadedData() {
@@ -232,58 +225,67 @@ class Mp3AudioTrackBacking implements InputAudioTrackBacking {
}
async getFirstPacket(options: PacketRetrievalOptions) {
// Ensure we have at least one frame loaded
while (this.demuxer.loadedSamples.length === 0 && this.demuxer.lastLoadedPos < this.demuxer.fileSize) {
await this.demuxer.loadNextChunk();
}
return this.getPacketAtIndex(0, options);
}
async getNextPacket(packet: EncodedPacket, options: PacketRetrievalOptions) {
const sampleIndex = binarySearchExact(
this.demuxer.loadedSamples,
packet.timestamp,
x => x.timestamp,
);
if (sampleIndex === -1) {
throw new Error('Packet was not created from this track.');
}
const release = await this.demuxer.readingMutex.acquire();
const nextIndex = sampleIndex + 1;
// Ensure the next sample exists
while (nextIndex >= this.demuxer.loadedSamples.length && this.demuxer.lastLoadedPos < this.demuxer.fileSize) {
await this.demuxer.loadNextChunk();
}
try {
const sampleIndex = binarySearchExact(
this.demuxer.loadedSamples,
packet.timestamp,
x => x.timestamp,
);
if (sampleIndex === -1) {
throw new Error('Packet was not created from this track.');
}
return this.getPacketAtIndex(nextIndex, options);
const nextIndex = sampleIndex + 1;
// Ensure the next sample exists
while (
nextIndex >= this.demuxer.loadedSamples.length
&& this.demuxer.lastLoadedPos < this.demuxer.fileSize
) {
await this.demuxer.loadNextChunk();
}
return this.getPacketAtIndex(nextIndex, options);
} finally {
release();
}
}
async getPacket(timestamp: number, options: PacketRetrievalOptions) {
while (true) {
const index = binarySearchLessOrEqual(
this.demuxer.loadedSamples,
timestamp,
x => x.timestamp,
);
const release = await this.demuxer.readingMutex.acquire();
try {
while (true) {
const index = binarySearchLessOrEqual(
this.demuxer.loadedSamples,
timestamp,
x => x.timestamp,
);
if (index === -1 && this.demuxer.loadedSamples.length > 0) {
// We're before the first sample
return null;
if (index === -1 && this.demuxer.loadedSamples.length > 0) {
// We're before the first sample
return null;
}
if (this.demuxer.lastLoadedPos === this.demuxer.fileSize) {
// All data is loaded, return what we found
return this.getPacketAtIndex(index, options);
}
if (index >= 0 && index + 1 < this.demuxer.loadedSamples.length) {
// The next packet also exists, we're done
return this.getPacketAtIndex(index, options);
}
// Otherwise, keep loading data
await this.demuxer.loadNextChunk();
}
if (this.demuxer.lastLoadedPos === this.demuxer.fileSize) {
// All data is loaded, return what we found
return this.getPacketAtIndex(index, options);
}
if (index >= 0 && index + 1 < this.demuxer.loadedSamples.length) {
// The next packet also exists, we're done
return this.getPacketAtIndex(index, options);
}
// Otherwise, keep loading data
await this.demuxer.loadNextChunk();
} finally {
release();
}
}