Implement ReadOrchestrator class, add highly optimized UrlSource

This commit is contained in:
Vanilagy
2025-08-29 19:39:50 +02:00
parent 641cdd02e3
commit c93755def1
15 changed files with 704 additions and 824 deletions
+1 -6
View File
@@ -44,7 +44,6 @@ export class AdtsDemuxer extends Demuxer {
readingMutex = new AsyncMutex();
lastSampleLoaded = false;
lastLoadedPos = 0;
fileSize = 0;
nextTimestampInSamples = 0;
constructor(input: Input) {
@@ -55,10 +54,6 @@ export class AdtsDemuxer extends Demuxer {
async readMetadata() {
return this.metadataPromise ??= (async () => {
let fileSize = this.reader.requestSize();
if (fileSize instanceof Promise) fileSize = await fileSize;
this.fileSize = fileSize;
// Keep loading until we find the first frame header
while (!this.firstFrameHeader && !this.lastSampleLoaded) {
await this.advanceReader();
@@ -86,7 +81,7 @@ export class AdtsDemuxer extends Demuxer {
return;
}
if (header.startPos + header.frameLength > this.fileSize) {
if (header.startPos + header.frameLength > this.reader.fileSize) {
// Frame doesn't fit in the rest of the file
this.lastSampleLoaded = true;
return;
+1
View File
@@ -98,6 +98,7 @@ export {
StreamSourceOptions,
BlobSource,
UrlSource,
UrlSource2,
UrlSourceOptions,
} from './source';
export {
+2 -5
View File
@@ -9,7 +9,6 @@
import { Demuxer } from './demuxer';
import { InputFormat } from './input-format';
import { assert } from './misc';
import { Reader } from './reader';
import { Reader2 } from './reader2';
import { Source } from './source';
@@ -34,13 +33,12 @@ export class Input<S extends Source = Source> {
/** @internal */
_formats: InputFormat[];
/** @internal */
_mainReader: Reader;
/** @internal */
_demuxerPromise: Promise<Demuxer> | null = null;
/** @internal */
_format: InputFormat | null = null;
_reader2: Reader2;
_size!: number;
constructor(options: InputOptions<S>) {
if (!options || typeof options !== 'object') {
@@ -55,14 +53,13 @@ export class Input<S extends Source = Source> {
this._formats = options.formats;
this._source = options.source;
this._mainReader = new Reader(options.source);
this._reader2 = new Reader2(options.source);
}
/** @internal */
_getDemuxer() {
return this._demuxerPromise ??= (async () => {
await this._mainReader.loadRange(0, 4096); // Load the first 4 kiB so we can determine the format
this._reader2.fileSize = await this._source.getSize();
for (const format of this._formats) {
const canRead = await format._canReadInput(this);
+5 -11
View File
@@ -271,11 +271,8 @@ export class IsobmffDemuxer extends Demuxer {
readMetadata() {
return this.metadataPromise ??= (async () => {
let sourceSize = this.reader.requestSize();
if (sourceSize instanceof Promise) sourceSize = await sourceSize;
let currentPos = 0;
while (currentPos < sourceSize) {
while (currentPos < this.reader.fileSize) {
let slice = this.reader.requestSliceRange(currentPos, MIN_BOX_HEADER_SIZE, MAX_BOX_HEADER_SIZE);
if (slice instanceof Promise) slice = await slice;
if (!slice) break;
@@ -315,14 +312,14 @@ export class IsobmffDemuxer extends Demuxer {
if (this.isFragmented) {
// The last 4 bytes may contain the size of the mfra box at the end of the file
let lastWordSlice = this.reader.requestSlice(sourceSize - 4, 4);
let lastWordSlice = this.reader.requestSlice(this.reader.fileSize - 4, 4);
if (lastWordSlice instanceof Promise) lastWordSlice = await lastWordSlice;
assert(lastWordSlice);
const lastWord = readU32Be(lastWordSlice);
const potentialMfraPos = sourceSize - lastWord;
const potentialMfraPos = this.reader.fileSize - lastWord;
if (potentialMfraPos >= 0 && potentialMfraPos <= sourceSize - MAX_BOX_HEADER_SIZE) {
if (potentialMfraPos >= 0 && potentialMfraPos <= this.reader.fileSize - MAX_BOX_HEADER_SIZE) {
let mfraHeaderSlice = this.reader.requestSliceRange(
potentialMfraPos,
MIN_BOX_HEADER_SIZE,
@@ -2417,9 +2414,6 @@ abstract class IsobmffTrackBacking implements InputTrackBacking {
return this.fetchPacketInFragment(fragment, sampleIndex, options);
}
let sourceSize = demuxer.reader.requestSize();
if (sourceSize instanceof Promise) sourceSize = await sourceSize;
let prevFragment: Fragment | null = null;
let bestFragmentIndex = fragmentIndex;
let bestSampleIndex = sampleIndex;
@@ -2455,7 +2449,7 @@ abstract class IsobmffTrackBacking implements InputTrackBacking {
}
}
while (currentPos < sourceSize) {
while (currentPos < demuxer.reader.fileSize) {
if (prevFragment) {
const trackData = prevFragment.trackData.get(this.internalTrack.id);
if (trackData && trackData.startTimestamp > latestTimestamp) {
+13 -11
View File
@@ -234,12 +234,9 @@ export class MatroskaDemuxer extends Demuxer {
readMetadata() {
return this.readMetadataPromise ??= (async () => {
let fileSize = this.reader.requestSize();
if (fileSize instanceof Promise) fileSize = await fileSize;
// Loop over all top-level elements in the file
let currentPos = 0;
while (currentPos < fileSize) {
while (currentPos < this.reader.fileSize) {
let slice = this.reader.requestSliceRange(currentPos, MIN_HEADER_SIZE, MAX_HEADER_SIZE);
if (slice instanceof Promise) slice = await slice;
if (!slice) break;
@@ -281,9 +278,9 @@ export class MatroskaDemuxer extends Demuxer {
this.reader,
dataStartPos,
LEVEL_0_AND_1_EBML_IDS,
fileSize,
this.reader.fileSize,
);
size = (nextElementPos ?? fileSize) - dataStartPos;
size = (nextElementPos ?? this.reader.fileSize) - dataStartPos;
}
const lastSegment = last(this.segments);
@@ -422,12 +419,17 @@ export class MatroskaDemuxer extends Demuxer {
}
}
// Use the seek head to read missing metadata elements
for (const target of METADATA_ELEMENTS) {
if (this.currentSegment[target.flag]) continue;
// Sort the seek entries by file position so reading them exhibits a sequential pattern
this.currentSegment.seekEntries.sort((a, b) => a.segmentPosition - b.segmentPosition);
const seekEntry = this.currentSegment.seekEntries.find(entry => entry.id === target.id);
if (!seekEntry) continue;
// Use the seek head to read missing metadata elements
for (const seekEntry of this.currentSegment.seekEntries) {
const target = METADATA_ELEMENTS.find(x => x.id === seekEntry.id);
if (!target) {
continue;
}
if (this.currentSegment[target.flag]) continue;
let slice = this.reader.requestSliceRange(
segmentDataStart + seekEntry.segmentPosition,
+4
View File
@@ -633,3 +633,7 @@ export const isSafari = () => {
* @public
*/
export type MaybePromise<T> = T | Promise<T>;
export const closedIntervalsOverlap = (startA: number, endA: number, startB: number, endB: number) => {
return startA <= endB && startB <= endA;
};
+2 -7
View File
@@ -36,7 +36,6 @@ export class Mp3Demuxer extends Demuxer {
readingMutex = new AsyncMutex();
lastSampleLoaded = false;
lastLoadedPos = 0;
fileSize = 0;
nextTimestampInSamples = 0;
constructor(input: Input) {
@@ -47,12 +46,8 @@ export class Mp3Demuxer extends Demuxer {
async readMetadata() {
return this.metadataPromise ??= (async () => {
let fileSize = this.reader.requestSize();
if (fileSize instanceof Promise) fileSize = await fileSize;
this.fileSize = fileSize;
// Keep loading until we find the first frame header
while (!this.firstFrameHeader && this.lastLoadedPos < this.fileSize) {
while (!this.firstFrameHeader && this.lastLoadedPos < this.reader.fileSize) {
await this.advanceReader();
}
@@ -82,7 +77,7 @@ export class Mp3Demuxer extends Demuxer {
const startPos = this.lastLoadedPos;
const result = await readNextFrameHeader(this.reader, startPos, this.fileSize);
const result = await readNextFrameHeader(this.reader, startPos, this.reader.fileSize);
if (!result) {
this.lastSampleLoaded = true;
return;
+1 -1
View File
@@ -59,7 +59,7 @@ export class Mp3Muxer extends Muxer {
}
const word = view.getUint32(0, false);
const header = readFrameHeader(word, { pos: 0, fileSize: null });
const header = readFrameHeader(word, null).header;
if (!header) {
throw new Error('Invalid MP3 header in sample.');
}
+1 -4
View File
@@ -26,9 +26,6 @@ export const readNextFrameHeader = async (reader: Reader2, startPos: number, unt
header: FrameHeader;
startPos: number;
} | null> => {
let fileSize = reader.requestSize();
if (fileSize instanceof Promise) fileSize = await fileSize;
let currentPos = startPos;
while (currentPos < until) {
@@ -38,7 +35,7 @@ export const readNextFrameHeader = async (reader: Reader2, startPos: number, unt
const word = readU32Be(slice);
const result = readFrameHeader(word, fileSize - currentPos);
const result = readFrameHeader(word, reader.fileSize - currentPos);
if (result.header) {
return { header: result.header, startPos: currentPos };
}
+4 -14
View File
@@ -46,7 +46,6 @@ export class OggDemuxer extends Demuxer {
reader: Reader2;
metadataPromise: Promise<void> | null = null;
fileSize: number | null = null;
bitstreams: LogicalBitstream[] = [];
tracks: InputAudioTrack[] = [];
@@ -58,12 +57,8 @@ export class OggDemuxer extends Demuxer {
async readMetadata() {
return this.metadataPromise ??= (async () => {
let fileSize = this.reader.requestSize();
if (fileSize instanceof Promise) fileSize = await fileSize;
this.fileSize = fileSize;
let currentPos = 0;
while (currentPos <= this.fileSize - MIN_PAGE_HEADER_SIZE) {
while (currentPos <= this.reader.fileSize - MIN_PAGE_HEADER_SIZE) {
let slice = this.reader.requestSliceRange(currentPos, MIN_PAGE_HEADER_SIZE, MAX_PAGE_HEADER_SIZE);
if (slice instanceof Promise) slice = await slice;
if (!slice) break;
@@ -249,7 +244,6 @@ export class OggDemuxer extends Demuxer {
async readPacket(startPage: Page, startSegmentIndex: number): Promise<Packet | null> {
assert(startSegmentIndex < startPage.lacingValues.length);
assert(this.fileSize);
let startDataOffset = 0;
for (let i = 0; i < startSegmentIndex; i++) {
@@ -290,7 +284,7 @@ export class OggDemuxer extends Demuxer {
// The packet extends to the next page; let's find it
let currentPos = currentPage.headerStartPos + currentPage.totalSize;
while (true) {
if (currentPos > this.fileSize - MIN_PAGE_HEADER_SIZE) {
if (currentPos > this.reader.fileSize - MIN_PAGE_HEADER_SIZE) {
return null;
}
@@ -335,8 +329,6 @@ export class OggDemuxer extends Demuxer {
}
async findNextPacketStart(lastPacket: Packet) {
assert(this.fileSize !== null);
// If there's another segment in the same page, return it
if (lastPacket.endSegmentIndex < lastPacket.endPage.lacingValues.length - 1) {
return { startPage: lastPacket.endPage, startSegmentIndex: lastPacket.endSegmentIndex + 1 };
@@ -351,7 +343,7 @@ export class OggDemuxer extends Demuxer {
// Otherwise, search for the next page belonging to the same bitstream
let currentPos = lastPacket.endPage.headerStartPos + lastPacket.endPage.totalSize;
while (true) {
if (currentPos >= this.fileSize - MIN_PAGE_HEADER_SIZE) {
if (currentPos >= this.reader.fileSize - MIN_PAGE_HEADER_SIZE) {
return null;
}
@@ -565,8 +557,6 @@ class OggAudioTrackBacking implements InputAudioTrackBacking {
}
async getPacket(timestamp: number, options: PacketRetrievalOptions) {
assert(this.demuxer.fileSize !== null);
const timestampInSamples = roundToPrecision(timestamp * this.internalSampleRate, 14);
if (timestampInSamples === 0) {
// Fast path for timestamp 0 - avoids binary search when playing back from the start
@@ -584,7 +574,7 @@ class OggAudioTrackBacking implements InputAudioTrackBacking {
}
let lowPage = startPosition.startPage;
let high = this.demuxer.fileSize;
let high = this.demuxer.reader.fileSize;
const lowPages: Page[] = [lowPage];
-205
View File
@@ -1,205 +0,0 @@
/*!
* Copyright (c) 2025-present, Vanilagy and contributors
*
* This Source Code Form is subject to the terms of the Mozilla Public
* License, v. 2.0. If a copy of the MPL was not distributed with this
* file, You can obtain one at https://mozilla.org/MPL/2.0/.
*/
import { assert, binarySearchLessOrEqual, removeItem } from './misc';
import { Source } from './source';
type ReadSegment = {
start: number;
end: number;
bytes: Uint8Array;
view: DataView;
age: number;
};
type LoadingSegment = {
start: number;
end: number;
promise: Promise<Uint8Array>;
};
export class Reader {
loadedSegments: ReadSegment[] = [];
loadingSegments: LoadingSegment[] = [];
sourceSizePromise: Promise<number> | null = null;
nextAge = 0;
totalStoredBytes = 0;
constructor(public source: Source, public maxStorableBytes = Infinity) {}
async loadRange(start: number, end: number) {
end = Math.min(end, await this.source.getSize());
if (start >= end) {
return;
}
const matchingLoadingSegment = this.loadingSegments.find(x => x.start <= start && x.end >= end);
if (matchingLoadingSegment) {
// Simply wait for the existing promise to finish to avoid loading the same range twice
await matchingLoadingSegment.promise;
return;
}
const index = binarySearchLessOrEqual(
this.loadedSegments,
start,
x => x.start,
);
if (index !== -1) {
for (let i = index; i < this.loadedSegments.length; i++) {
const segment = this.loadedSegments[i]!;
if (segment.start > start) {
break;
}
const segmentEncasesRequestedRange = segment.end >= end;
if (segmentEncasesRequestedRange) {
// Nothing to load
return;
}
}
}
this.source.onread?.(start, end);
const bytesPromise = this.source._read(start, end);
const loadingSegment: LoadingSegment = { start, end, promise: bytesPromise };
this.loadingSegments.push(loadingSegment);
const bytes = await bytesPromise;
removeItem(this.loadingSegments, loadingSegment);
this.insertIntoLoadedSegments(start, bytes);
}
rangeIsLoaded(start: number, end: number) {
if (end <= start) {
return true;
}
const index = binarySearchLessOrEqual(this.loadedSegments, start, x => x.start);
if (index === -1) {
return false;
}
for (let i = index; i < this.loadedSegments.length; i++) {
const segment = this.loadedSegments[i]!;
if (segment.start > start) {
break;
}
const segmentEncasesRequestedRange = segment.end >= end;
if (segmentEncasesRequestedRange) {
return true;
}
}
return false;
}
private insertIntoLoadedSegments(start: number, bytes: Uint8Array) {
const segment: ReadSegment = {
start,
end: start + bytes.byteLength,
bytes,
view: new DataView(bytes.buffer),
age: this.nextAge++,
};
let index = binarySearchLessOrEqual(this.loadedSegments, start, x => x.start);
if (index === -1 || this.loadedSegments[index]!.start < segment.start) {
index++;
}
// Insert the segment at the right place so that the array remains sorted by start offset
this.loadedSegments.splice(index, 0, segment);
this.totalStoredBytes += bytes.byteLength;
// Remove all other segments from the array that are completely covered by the newly-inserted segment
for (let i = index + 1; i < this.loadedSegments.length; i++) {
const otherSegment = this.loadedSegments[i]!;
if (otherSegment.start >= segment.end) {
break;
}
if (segment.start <= otherSegment.start && otherSegment.end <= segment.end) {
this.loadedSegments.splice(i, 1);
i--;
}
}
// If we overshoot the max amount of permitted bytes, let's start evicting the oldest segments
while (this.totalStoredBytes > this.maxStorableBytes && this.loadedSegments.length > 1) {
let oldestSegment: ReadSegment | null = null;
let oldestSegmentIndex = -1;
for (let i = 0; i < this.loadedSegments.length; i++) {
const candidate = this.loadedSegments[i]!;
if (!oldestSegment || candidate.age < oldestSegment.age) {
oldestSegment = candidate;
oldestSegmentIndex = i;
}
}
assert(oldestSegment);
this.totalStoredBytes -= oldestSegment.bytes.byteLength;
this.loadedSegments.splice(oldestSegmentIndex, 1);
}
}
getViewAndOffset(start: number, end: number) {
const startIndex = binarySearchLessOrEqual(this.loadedSegments, start, x => x.start);
let segment: ReadSegment | null = null;
if (startIndex !== -1) {
for (let i = startIndex; i < this.loadedSegments.length; i++) {
const candidate = this.loadedSegments[i]!;
if (candidate.start > start) {
break;
}
if (end <= candidate.end) {
segment = candidate;
break;
}
}
}
if (!segment) {
throw new Error(`No segment loaded for range [${start}, ${end}).`);
}
segment.age = this.nextAge++;
return {
view: segment.view,
offset: segment.bytes.byteOffset + start - segment.start,
};
}
forgetRange(start: number, end: number) {
if (end <= start) {
return;
}
const startIndex = binarySearchLessOrEqual(this.loadedSegments, start, x => x.start);
if (startIndex === -1) {
return;
}
const segment = this.loadedSegments[startIndex]!;
if (segment.start !== start || segment.end !== end) {
return;
}
this.loadedSegments.splice(startIndex, 1);
this.totalStoredBytes -= segment.bytes.byteLength;
}
}
+10 -32
View File
@@ -56,28 +56,15 @@ export class FileSlice {
}
export class Reader2 {
private size: number | null = null;
fileSize!: number;
constructor(public source: Source) {
}
requestSize(): MaybePromise<number> {
if (this.size !== null) {
return this.size;
}
const size = this.source._retrieveSize2();
if (size instanceof Promise) {
void size.then(x => this.size = x);
return size;
} else {
this.size = size;
return size;
}
}
constructor(public source: Source) {}
requestSlice(start: number, length: number): MaybePromise<FileSlice | null> {
if (start + length > this.fileSize) {
return null;
}
const end = start + length;
const result = this.source._read2(start, end);
@@ -99,19 +86,10 @@ export class Reader2 {
}
requestSliceRange(start: number, minLength: number, maxLength: number): MaybePromise<FileSlice | null> {
const fileSize = this.requestSize();
if (fileSize instanceof Promise) {
return fileSize.then(size => this.requestSlice(
start,
clamp(size - start, minLength, maxLength),
));
} else {
return this.requestSlice(
start,
clamp(fileSize - start, minLength, maxLength),
);
}
return this.requestSlice(
start,
clamp(this.fileSize - start, minLength, maxLength),
);
}
}
+658 -460
View File
File diff suppressed because it is too large Load Diff
-63
View File
@@ -1,63 +0,0 @@
/*!
* Copyright (c) 2025-present, Vanilagy and contributors
*
* This Source Code Form is subject to the terms of the Mozilla Public
* License, v. 2.0. If a copy of the MPL was not distributed with this
* file, You can obtain one at https://mozilla.org/MPL/2.0/.
*/
import { Reader } from '../reader';
export class RiffReader {
pos = 0;
littleEndian = true;
constructor(public reader: Reader) {}
readBytes(length: number) {
const { view, offset } = this.reader.getViewAndOffset(this.pos, this.pos + length);
this.pos += length;
return new Uint8Array(view.buffer, offset, length);
}
readU16() {
const { view, offset } = this.reader.getViewAndOffset(this.pos, this.pos + 2);
this.pos += 2;
return view.getUint16(offset, this.littleEndian);
}
readU32() {
const { view, offset } = this.reader.getViewAndOffset(this.pos, this.pos + 4);
this.pos += 4;
return view.getUint32(offset, this.littleEndian);
}
readU64() {
let low: number;
let high: number;
if (this.littleEndian) {
low = this.readU32();
high = this.readU32();
} else {
high = this.readU32();
low = this.readU32();
}
return high * 0x100000000 + low;
}
readAscii(length: number) {
const { view, offset } = this.reader.getViewAndOffset(this.pos, this.pos + length);
this.pos += length;
let str = '';
for (let i = 0; i < length; i++) {
str += String.fromCharCode(view.getUint8(offset + i));
}
return str;
}
}
+2 -5
View File
@@ -47,9 +47,6 @@ export class WaveDemuxer extends Demuxer {
async readMetadata() {
return this.metadataPromise ??= (async () => {
let actualFileSize = this.reader.requestSize();
if (actualFileSize instanceof Promise) actualFileSize = await actualFileSize;
let slice = this.reader.requestSlice(0, 12);
if (slice instanceof Promise) slice = await slice;
assert(slice);
@@ -61,7 +58,7 @@ export class WaveDemuxer extends Demuxer {
const outerChunkSize = readU32(slice, littleEndian);
let totalFileSize = isRf64 ? actualFileSize : Math.min(outerChunkSize + 8, actualFileSize);
let totalFileSize = isRf64 ? this.reader.fileSize : Math.min(outerChunkSize + 8, this.reader.fileSize);
const format = readAscii(slice, 4);
if (format !== 'WAVE') {
@@ -98,7 +95,7 @@ export class WaveDemuxer extends Demuxer {
const riffChunkSize = readU64(slice, littleEndian);
dataChunkSize = readU64(slice, littleEndian);
totalFileSize = Math.min(riffChunkSize + 8, actualFileSize);
totalFileSize = Math.min(riffChunkSize + 8, this.reader.fileSize);
}
currentPos = startPos + chunkSize + (chunkSize & 1); // Handle padding