diff --git a/src/adts/adts-demuxer.ts b/src/adts/adts-demuxer.ts index 438b1cd..43e10a0 100644 --- a/src/adts/adts-demuxer.ts +++ b/src/adts/adts-demuxer.ts @@ -20,7 +20,7 @@ import { UNDETERMINED_LANGUAGE, } from '../misc'; import { EncodedPacket, PLACEHOLDER_DATA } from '../packet'; -import { readBytes, Reader2 } from '../reader2'; +import { readBytes, Reader } from '../reader2'; import { FrameHeader, MAX_FRAME_HEADER_SIZE, MIN_FRAME_HEADER_SIZE, readFrameHeader } from './adts-reader'; const SAMPLES_PER_AAC_FRAME = 1024; @@ -33,7 +33,7 @@ type Sample = { }; export class AdtsDemuxer extends Demuxer { - reader: Reader2; + reader: Reader; metadataPromise: Promise | null = null; firstFrameHeader: FrameHeader | null = null; @@ -49,7 +49,7 @@ export class AdtsDemuxer extends Demuxer { constructor(input: Input) { super(input); - this.reader = input._reader2; + this.reader = input._reader; } async readMetadata() { diff --git a/src/index.ts b/src/index.ts index a218ddb..9354d29 100644 --- a/src/index.ts +++ b/src/index.ts @@ -98,7 +98,6 @@ export { StreamSourceOptions, BlobSource, UrlSource, - UrlSource2, UrlSourceOptions, } from './source'; export { diff --git a/src/input-format.ts b/src/input-format.ts index bc14a1c..46cb6e5 100644 --- a/src/input-format.ts +++ b/src/input-format.ts @@ -53,7 +53,7 @@ export abstract class InputFormat { export abstract class IsobmffInputFormat extends InputFormat { /** @internal */ protected async _getMajorBrand(input: Input) { - let slice = input._reader2.requestSlice(0, 12); + let slice = input._reader.requestSlice(0, 12); if (slice instanceof Promise) slice = await slice; if (!slice) return null; @@ -120,7 +120,7 @@ export class QuickTimeInputFormat extends IsobmffInputFormat { export class MatroskaInputFormat extends InputFormat { /** @internal */ protected async isSupportedEBMLOfDocType(input: Input, desiredDocType: string) { - let headerSlice = input._reader2.requestSlice(0, MAX_HEADER_SIZE); + let headerSlice = input._reader.requestSlice(0, MAX_HEADER_SIZE); if (headerSlice instanceof Promise) headerSlice = await headerSlice; if (!headerSlice) return false; @@ -143,7 +143,7 @@ export class MatroskaInputFormat extends InputFormat { return false; // Miss me with that shit } - let dataSlice = input._reader2.requestSlice(headerSlice.filePos, dataSize); + let dataSlice = input._reader.requestSlice(headerSlice.filePos, dataSize); if (dataSlice instanceof Promise) dataSlice = await dataSlice; if (!dataSlice) return false; @@ -235,7 +235,7 @@ export class WebMInputFormat extends MatroskaInputFormat { export class Mp3InputFormat extends InputFormat { /** @internal */ async _canReadInput(input: Input) { - let slice = input._reader2.requestSlice(0, 10); + let slice = input._reader.requestSlice(0, 10); if (slice instanceof Promise) slice = await slice; if (!slice) return false; @@ -247,7 +247,7 @@ export class Mp3InputFormat extends InputFormat { currentPos = slice.filePos + id3Tag.size; } - const firstResult = await readNextFrameHeader(input._reader2, currentPos, currentPos + 4096); + const firstResult = await readNextFrameHeader(input._reader, currentPos, currentPos + 4096); if (!firstResult) { return false; } @@ -261,7 +261,7 @@ export class Mp3InputFormat extends InputFormat { // Fine, we found one frame header, but we're still not entirely sure this is MP3. Let's check if we can find // another header right after it: - const secondResult = await readNextFrameHeader(input._reader2, currentPos, currentPos + FRAME_HEADER_SIZE); + const secondResult = await readNextFrameHeader(input._reader, currentPos, currentPos + FRAME_HEADER_SIZE); if (!secondResult) { return false; } @@ -299,7 +299,7 @@ export class Mp3InputFormat extends InputFormat { export class WaveInputFormat extends InputFormat { /** @internal */ async _canReadInput(input: Input) { - let slice = input._reader2.requestSlice(0, 12); + let slice = input._reader.requestSlice(0, 12); if (slice instanceof Promise) slice = await slice; if (!slice) return false; @@ -335,7 +335,7 @@ export class WaveInputFormat extends InputFormat { export class OggInputFormat extends InputFormat { /** @internal */ async _canReadInput(input: Input) { - let slice = input._reader2.requestSlice(0, 4); + let slice = input._reader.requestSlice(0, 4); if (slice instanceof Promise) slice = await slice; if (!slice) return false; @@ -363,7 +363,7 @@ export class OggInputFormat extends InputFormat { export class AdtsInputFormat extends InputFormat { /** @internal */ async _canReadInput(input: Input) { - let slice = input._reader2.requestSliceRange(0, MIN_FRAME_HEADER_SIZE, MAX_FRAME_HEADER_SIZE); + let slice = input._reader.requestSliceRange(0, MIN_FRAME_HEADER_SIZE, MAX_FRAME_HEADER_SIZE); if (slice instanceof Promise) slice = await slice; if (!slice) return false; @@ -372,7 +372,7 @@ export class AdtsInputFormat extends InputFormat { return false; } - slice = input._reader2.requestSliceRange(firstHeader.frameLength, MIN_FRAME_HEADER_SIZE, MAX_FRAME_HEADER_SIZE); + slice = input._reader.requestSliceRange(firstHeader.frameLength, MIN_FRAME_HEADER_SIZE, MAX_FRAME_HEADER_SIZE); if (slice instanceof Promise) slice = await slice; if (!slice) return false; diff --git a/src/input.ts b/src/input.ts index 6e2857e..792db02 100644 --- a/src/input.ts +++ b/src/input.ts @@ -9,7 +9,7 @@ import { Demuxer } from './demuxer'; import { InputFormat } from './input-format'; import { assert } from './misc'; -import { Reader2 } from './reader2'; +import { Reader } from './reader2'; import { Source } from './source'; /** @@ -36,9 +36,8 @@ export class Input { _demuxerPromise: Promise | null = null; /** @internal */ _format: InputFormat | null = null; - - _reader2: Reader2; - _size!: number; + /** @internal */ + _reader: Reader; constructor(options: InputOptions) { if (!options || typeof options !== 'object') { @@ -53,13 +52,13 @@ export class Input { this._formats = options.formats; this._source = options.source; - this._reader2 = new Reader2(options.source); + this._reader = new Reader(options.source); } /** @internal */ _getDemuxer() { return this._demuxerPromise ??= (async () => { - this._reader2.fileSize = await this._source.getSize(); + this._reader.fileSize = await this._source.getSize(); for (const format of this._formats) { const canRead = await format._canReadInput(this); diff --git a/src/isobmff/isobmff-demuxer.ts b/src/isobmff/isobmff-demuxer.ts index 4acda9d..4365e44 100644 --- a/src/isobmff/isobmff-demuxer.ts +++ b/src/isobmff/isobmff-demuxer.ts @@ -75,7 +75,7 @@ import { readI16Be, readI32Be, readI64Be, - Reader2, + Reader, readU16Be, readU24Be, readU32Be, @@ -223,7 +223,7 @@ type Fragment = { }; export class IsobmffDemuxer extends Demuxer { - reader: Reader2; + reader: Reader; moovSlice: FileSlice | null = null; currentTrack: InternalTrack | null = null; @@ -242,7 +242,7 @@ export class IsobmffDemuxer extends Demuxer { constructor(input: Input) { super(input); - this.reader = input._reader2; + this.reader = input._reader; } override async computeDuration() { diff --git a/src/matroska/ebml.ts b/src/matroska/ebml.ts index 276795a..0ccbb3f 100644 --- a/src/matroska/ebml.ts +++ b/src/matroska/ebml.ts @@ -8,7 +8,7 @@ import { MediaCodec } from '../codec'; import { assertNever, textDecoder, textEncoder } from '../misc'; -import { FileSlice, readBytes, Reader2, readF32Be, readF64Be, readU8 } from '../reader2'; +import { FileSlice, readBytes, Reader, readF32Be, readF64Be, readU8 } from '../reader2'; import { Writer } from '../writer'; export interface EBMLElement { @@ -554,7 +554,7 @@ export const readFloat = (slice: FileSlice, width: number) => { }; /** Returns the byte offset in the file of the next element with a matching ID. */ -export const searchForNextElementId = async (reader: Reader2, startPos: number, ids: EBMLId[], until: number) => { +export const searchForNextElementId = async (reader: Reader, startPos: number, ids: EBMLId[], until: number) => { const idsSet = new Set(ids); let currentPos = startPos; @@ -581,7 +581,7 @@ export const searchForNextElementId = async (reader: Reader2, startPos: number, }; /** Searches for the next occurrence of an element ID using a naive byte-wise search. */ -export const resync = async (reader: Reader2, startPos: number, ids: EBMLId[], until: number) => { +export const resync = async (reader: Reader, startPos: number, ids: EBMLId[], until: number) => { const CHUNK_SIZE = 2 ** 16; // So we don't need to grab thousands of slices const idsSet = new Set(ids); let currentPos = startPos; diff --git a/src/matroska/matroska-demuxer.ts b/src/matroska/matroska-demuxer.ts index 43b2503..92267c0 100644 --- a/src/matroska/matroska-demuxer.ts +++ b/src/matroska/matroska-demuxer.ts @@ -69,7 +69,7 @@ import { searchForNextElementId, } from './ebml'; import { buildMatroskaMimeType } from './matroska-misc'; -import { FileSlice, readBytes, Reader2, readI16Be, readU8 } from '../reader2'; +import { FileSlice, readBytes, Reader, readI16Be, readU8 } from '../reader2'; type Segment = { seekHeadSeen: boolean; @@ -188,7 +188,7 @@ const METADATA_ELEMENTS = [ const MAX_RESYNC_LENGTH = 10 * 2 ** 20; // 10 MiB export class MatroskaDemuxer extends Demuxer { - reader: Reader2; + reader: Reader; readMetadataPromise: Promise | null = null; @@ -204,7 +204,7 @@ export class MatroskaDemuxer extends Demuxer { constructor(input: Input) { super(input); - this.reader = input._reader2; + this.reader = input._reader; } override async computeDuration() { diff --git a/src/mp3/mp3-demuxer.ts b/src/mp3/mp3-demuxer.ts index 9d188dd..0d4f1c9 100644 --- a/src/mp3/mp3-demuxer.ts +++ b/src/mp3/mp3-demuxer.ts @@ -15,7 +15,7 @@ import { assert, AsyncMutex, binarySearchExact, binarySearchLessOrEqual, UNDETER import { EncodedPacket, PLACEHOLDER_DATA } from '../packet'; import { FrameHeader, getXingOffset, INFO, XING } from '../../shared/mp3-misc'; import { readId3, readNextFrameHeader } from './mp3-reader'; -import { readBytes, Reader2, readU32Be } from '../reader2'; +import { readBytes, Reader, readU32Be } from '../reader2'; type Sample = { timestamp: number; @@ -25,7 +25,7 @@ type Sample = { }; export class Mp3Demuxer extends Demuxer { - reader: Reader2; + reader: Reader; metadataPromise: Promise | null = null; firstFrameHeader: FrameHeader | null = null; @@ -41,7 +41,7 @@ export class Mp3Demuxer extends Demuxer { constructor(input: Input) { super(input); - this.reader = input._reader2; + this.reader = input._reader; } async readMetadata() { diff --git a/src/mp3/mp3-reader.ts b/src/mp3/mp3-reader.ts index 90cf143..884b4db 100644 --- a/src/mp3/mp3-reader.ts +++ b/src/mp3/mp3-reader.ts @@ -7,7 +7,7 @@ */ import { FRAME_HEADER_SIZE, FrameHeader, readFrameHeader } from '../../shared/mp3-misc'; -import { FileSlice, readAscii, Reader2, readU32Be } from '../reader2'; +import { FileSlice, readAscii, Reader, readU32Be } from '../reader2'; export const readId3 = (slice: FileSlice) => { const tag = readAscii(slice, 3); @@ -22,7 +22,7 @@ export const readId3 = (slice: FileSlice) => { return { size }; }; -export const readNextFrameHeader = async (reader: Reader2, startPos: number, until: number): Promise<{ +export const readNextFrameHeader = async (reader: Reader, startPos: number, until: number): Promise<{ header: FrameHeader; startPos: number; } | null> => { diff --git a/src/ogg/ogg-demuxer.ts b/src/ogg/ogg-demuxer.ts index 559deec..c063392 100644 --- a/src/ogg/ogg-demuxer.ts +++ b/src/ogg/ogg-demuxer.ts @@ -14,7 +14,7 @@ import { InputAudioTrack, InputAudioTrackBacking } from '../input-track'; import { PacketRetrievalOptions } from '../media-sink'; import { assert, findLast, roundToPrecision, toDataView, UNDETERMINED_LANGUAGE } from '../misc'; import { EncodedPacket, PLACEHOLDER_DATA } from '../packet'; -import { readBytes, Reader2 } from '../reader2'; +import { readBytes, Reader } from '../reader2'; import { buildOggMimeType, computeOggPageCrc, extractSampleMetadata, OggCodecInfo } from './ogg-misc'; import { findNextPageHeader, @@ -43,7 +43,7 @@ type Packet = { }; export class OggDemuxer extends Demuxer { - reader: Reader2; + reader: Reader; metadataPromise: Promise | null = null; bitstreams: LogicalBitstream[] = []; @@ -52,7 +52,7 @@ export class OggDemuxer extends Demuxer { constructor(input: Input) { super(input); - this.reader = input._reader2; + this.reader = input._reader; } async readMetadata() { diff --git a/src/reader2.ts b/src/reader2.ts index ff000fd..2ff6b33 100644 --- a/src/reader2.ts +++ b/src/reader2.ts @@ -55,7 +55,7 @@ export class FileSlice { } } -export class Reader2 { +export class Reader { fileSize!: number; constructor(public source: Source) {} @@ -66,7 +66,7 @@ export class Reader2 { } const end = start + length; - const result = this.source._read2(start, end); + const result = this.source._read(start, end); if (result instanceof Promise) { return result.then((x) => { diff --git a/src/source.ts b/src/source.ts index cdb4e03..3d841cc 100644 --- a/src/source.ts +++ b/src/source.ts @@ -17,17 +17,22 @@ import { toDataView, } from './misc'; +export type ReadResult = { + bytes: Uint8Array; + view: DataView; + /** The offset of the bytes in the file. */ + offset: number; +}; + /** * The source base class, representing a resource from which bytes can be read. * @public */ export abstract class Source { - abstract _read2(start: number, end: number): MaybePromise<{ - bytes: Uint8Array; - view: DataView; - offset: number; - }>; - abstract _retrieveSize2(): MaybePromise; + /** @internal */ + abstract _retrieveSize(): MaybePromise; + /** @internal */ + abstract _read(start: number, end: number): MaybePromise; /** @internal */ _sizePromise: Promise | null = null; @@ -37,10 +42,10 @@ export abstract class Source { * will retrieve the size. */ async getSize() { - return this._sizePromise ??= Promise.resolve(this._retrieveSize2()); + return this._sizePromise ??= Promise.resolve(this._retrieveSize()); } - /** Called each time data is requested from the source. */ + /** Called each time data is retrieved from the source. Will be called with the retrieved range. */ onread: ((start: number, end: number) => unknown) | null = null; } @@ -53,6 +58,8 @@ export class BufferSource extends Source { _bytes: Uint8Array; /** @internal */ _view: DataView; + /** @internal */ + _onreadCalled = false; constructor(buffer: ArrayBuffer | Uint8Array) { if (!(buffer instanceof ArrayBuffer) && !(buffer instanceof Uint8Array)) { @@ -65,11 +72,19 @@ export class BufferSource extends Source { this._view = toDataView(this._bytes); } - _retrieveSize2() { + /** @internal */ + _retrieveSize(): number { return this._bytes.byteLength; } - _read2() { + /** @internal */ + _read(): ReadResult { + if (!this._onreadCalled) { + // We just say the first read retrives all bytes from the source (which, I mean, it does) + this.onread?.(0, this._bytes.byteLength); + this._onreadCalled = true; + } + return { bytes: this._bytes, view: this._view, @@ -83,10 +98,29 @@ export class BufferSource extends Source { * @public */ export type StreamSourceOptions = { - /** Called when data is requested. Should return or resolve to the bytes from the specified byte range. */ - read: (start: number, end: number) => Uint8Array | Promise; - /** Called when the size of the entire file is requested. Should return or resolve to the size in bytes. */ - getSize: () => number | Promise; + /** + * Called when data is requested. Must return or resolve to the bytes from the specified byte range, or a stream + * that yields these bytes. + */ + read: (start: number, end: number) => MaybePromise>; + + /** Called when the size of the entire file is requested. Must return or resolve to the size in bytes. */ + getSize: () => MaybePromise; + + /** The maximum number of bytes the cache is allowed to hold in memory. Defaults to 8 MiB. */ + maxCacheSize?: number; + + /** + * Specifies the prefetch profile that the reader should use with this source. A prefetch propfile specifies the + * pattern with which bytes outside of the requested range are preloaded to reduce latency for future reads. + * + * - `'none'` (default): No prefetching; only the data needed in the moment is requested. + * - `'fileSystem'`: File system-optimized prefetching: a small amount of data is prefetched bidirectionally. + * - `'network'`: Network-optimized prefetching, or more generally, prefetching optimized for any high-latency + * environment: tries to minimize the amount of read calls and aggressively prefetches data when sequential access + * patterns are detected. + */ + prefetchProfile?: 'none' | 'fileSystem' | 'network'; }; /** @@ -96,6 +130,8 @@ export type StreamSourceOptions = { export class StreamSource extends Source { /** @internal */ _options: StreamSourceOptions; + /** @internal */ + _orchestrator: ReadOrchestrator; constructor(options: StreamSourceOptions) { if (!options || typeof options !== 'object') { @@ -107,23 +143,120 @@ export class StreamSource extends Source { if (typeof options.getSize !== 'function') { throw new TypeError('options.getSize must be a function.'); } + if ( + options.maxCacheSize !== undefined + && (!Number.isInteger(options.maxCacheSize) || options.maxCacheSize < 0) + ) { + throw new TypeError('options.maxCacheSize, when provided, must be a non-negative integer.'); + } + if (options.prefetchProfile && !['none', 'fileSystem', 'network'].includes(options.prefetchProfile)) { + throw new TypeError( + 'options.prefetchProfile, when provided, must be one of \'none\', \'fileSystem\' or \'network\'.', + ); + } super(); this._options = options; + + this._orchestrator = new ReadOrchestrator({ + maxCacheSize: options.maxCacheSize ?? (8 * 2 ** 20 /* 8 MiB */), + maxWorkerCount: 2, // Fixed for now, *should* be fine + prefetchProfile: PREFETCH_PROFILES[options.prefetchProfile ?? 'none'], + runWorker: this._runWorker.bind(this), + }); } /** @internal */ - async _read(start: number, end: number) { - return this._options.read(start, end); + _retrieveSize(): MaybePromise { + const result = this._options.getSize(); + + if (result instanceof Promise) { + return result.then((size) => { + if (!Number.isInteger(size) || size < 0) { + throw new TypeError('options.getSize must return or resolve to a non-negative integer.'); + } + + this._orchestrator.fileSize = size; + return size; + }); + } else { + if (!Number.isInteger(result) || result < 0) { + throw new TypeError('options.getSize must return or resolve to a non-negative integer.'); + } + + this._orchestrator.fileSize = result; + return result; + } } /** @internal */ - async _retrieveSize() { - return this._options.getSize(); + _read(start: number, end: number): MaybePromise { + return this._orchestrator.read(start, end); + } + + private async _runWorker(worker: ReadWorker) { + while (worker.currentPos < worker.targetPos && !worker.aborted) { + const originalCurrentPos = worker.currentPos; + const originalTargetPos = worker.targetPos; + + let data = this._options.read(worker.currentPos, originalTargetPos); + if (data instanceof Promise) data = await data; + + if (data instanceof Uint8Array) { + if (data.length !== originalTargetPos - worker.currentPos) { + // Yes, we're that strict + throw new Error( + `options.read returned a Uint8Array with unexpected length: Requested ${ + originalTargetPos - worker.currentPos + } bytes, but got ${data.length}.`, + ); + } + + this.onread?.(worker.currentPos, worker.currentPos + data.length); + this._orchestrator.supplyWorkerData(worker, data); + } else if (data instanceof ReadableStream) { + const reader = data.getReader(); + + while (true) { + const { done, value } = await reader.read(); + if (done) { + if (worker.currentPos < originalTargetPos) { + // Yes, we're *that* strict + throw new Error( + `ReadableStream returned by options.read ended before supplying enough data.` + + ` Requested ${originalTargetPos - originalCurrentPos} bytes, but got ${ + worker.currentPos - originalCurrentPos + }`, + ); + } + + break; + } + + if (!(value instanceof Uint8Array)) { + throw new TypeError('ReadableStream returned by options.read must yield Uint8Array chunks.'); + } + + this.onread?.(worker.currentPos, worker.currentPos + value.length); + this._orchestrator.supplyWorkerData(worker, value); + + if (worker.currentPos >= originalTargetPos || worker.aborted) { + break; + } + } + } else { + throw new TypeError('options.read must return or resolve to a Uint8Array or a ReadableStream.'); + } + } } } +export type BlobSourceOptions = { + /** The maximum number of bytes the cache is allowed to hold in memory. Defaults to 8 MiB. */ + maxCacheSize?: number; +}; + /** * A source backed by a Blob. Since Files are also Blobs, this is the source to use when reading files off the disk. * @public @@ -131,60 +264,71 @@ export class StreamSource extends Source { export class BlobSource extends Source { /** @internal */ _blob: Blob; + /** @internal */ _orchestrator: ReadOrchestrator; - constructor(blob: Blob) { + constructor(blob: Blob, options: BlobSourceOptions = {}) { if (!(blob instanceof Blob)) { throw new TypeError('blob must be a Blob.'); } + if (!options || typeof options !== 'object') { + throw new TypeError('options must be an object.'); + } + if ( + options.maxCacheSize !== undefined + && (!Number.isInteger(options.maxCacheSize) || options.maxCacheSize < 0) + ) { + throw new TypeError('options.maxCacheSize, when provided, must be a non-negative integer.'); + } super(); this._blob = blob; this._orchestrator = new ReadOrchestrator({ - maxCacheSize: 8 * 2 ** 20, // 8 MiB + maxCacheSize: options.maxCacheSize ?? (8 * 2 ** 20 /* 8 MiB */), maxWorkerCount: 4, runWorker: this._runWorker.bind(this), - getPrefetchRange(start, end) { - const paddingStart = 2 ** 16; - const paddingEnd = 2 ** 17; - - start = Math.max(0, Math.floor((start - paddingStart) / paddingStart) * paddingStart); - end += paddingEnd; // Preload a tad into the future - - return { start, end }; - }, + prefetchProfile: PREFETCH_PROFILES.fileSystem, }); } - _retrieveSize2() { + /** @internal */ + _retrieveSize(): number { const size = this._blob.size; this._orchestrator.fileSize = size; return size; } - _read2(start: number, end: number) { + /** @internal */ + _read(start: number, end: number): MaybePromise { return this._orchestrator.read(start, end); } - readers = new WeakMap>(); + /** @internal */ + _readers = new WeakMap>(); - async _runWorker(worker: ReadWorker) { - let reader = this.readers.get(worker); + private async _runWorker(worker: ReadWorker) { + let reader = this._readers.get(worker); if (!reader) { // Get a reader of the blob starting at the required offset, and then keep it around reader = this._blob.slice(worker.currentPos).stream().getReader(); - this.readers.set(worker, reader); + this._readers.set(worker, reader); } while (worker.currentPos < worker.targetPos && !worker.aborted) { const { done, value } = await reader.read(); if (done) { this._orchestrator.forgetWorker(worker); + + if (worker.currentPos < worker.targetPos) { // I think this `if` should always hit? + throw new Error('Blob reader stopped unexpectedly before all requested data was read.'); + } + break; } + this.onread?.(worker.currentPos, worker.currentPos + value.length); this._orchestrator.supplyWorkerData(worker, value); } } @@ -208,6 +352,9 @@ export type UrlSourceOptions = { * with the number of previous, unsuccessful attempts. If the function returns `null`, no more retries will be made. */ getRetryDelay?: (previousAttempts: number) => number | null; + + /** The maximum number of bytes the cache is allowed to hold in memory. Defaults to 64 MiB. */ + maxCacheSize?: number; }; /** @@ -215,11 +362,14 @@ export type UrlSourceOptions = { * as it typically comes with increased latency. * @beta */ -export class UrlSource2 extends Source { +export class UrlSource extends Source { + /** @internal */ _url: URL; + /** @internal */ _options: UrlSourceOptions; + /** @internal */ _orchestrator: ReadOrchestrator; - + /** @internal */ _existingResponses = new WeakMap { // Retrieving the resource size for UrlSource is optimized: Almost always (= always), the first bytes we have to // read are the start of the file. This means it's smart to combine size fetching with fetching the start of the // file. We additionally use this step to probe if the server supports range requests, killing three birds with @@ -359,7 +471,8 @@ export class UrlSource2 extends Source { return fileSize; } - async _read2(start: number, end: number) { + /** @internal */ + _read(start: number, end: number): MaybePromise { return this._orchestrator.read(start, end); } @@ -407,9 +520,15 @@ export class UrlSource2 extends Source { const { done, value } = await reader.read(); if (done) { this._orchestrator.forgetWorker(worker); + + if (worker.currentPos < worker.targetPos) { + throw new Error('Response stream reader stopped unexpectedly before all requested data was read.'); + } + break; } + this.onread?.(worker.currentPos, worker.currentPos + value.length); this._orchestrator.supplyWorkerData(worker, value); if (worker.currentPos >= worker.targetPos || worker.aborted) { @@ -447,6 +566,69 @@ export class UrlSource2 extends Source { } } +type PrefetchProfile = (start: number, end: number, workers: ReadWorker[]) => { + start: number; + end: number; +}; + +const PREFETCH_PROFILES = { + none: (start, end) => ({ start, end }), + fileSystem: (start, end) => { + const padding = 2 ** 16; + + start = Math.floor((start - padding) / padding) * padding; + end = Math.ceil((end + padding) / padding) * padding; + + return { start, end }; + }, + network: (start, end, workers) => { + // Add a slight bit of start padding because backwards reading is painful + const paddingStart = 2 ** 16; + start = Math.max(0, Math.floor((start - paddingStart) / paddingStart) * paddingStart); + + // Remote resources have extreme latency (relatively speaking), so the benefit from intelligent + // prefetching is great. The network prefetch strategy is as follows: When we notice + // successive reads to a worker's read region, we prefetch more data at the end of that region, + // growing exponentially (up to a cap). This performs well for real-world use cases: Either we read a + // small part of the file once and then never need it again, in which case the requested about of data + // is small. Or, we're repeatedly doing a sequential access pattern (common in media files), in which + // case we can become more and more confident to prefetch more and more data. + for (const worker of workers) { + const maxExtensionAmount = 8 * 2 ** 20; // 8 MiB + + // When the read region cross the threshold point, we trigger a prefetch. This point is typically + // in the middle of the worker's read region, or a fixed offset from the end if the region has grown + // really large. + const thresholdPoint = Math.max( + (worker.startPos + worker.targetPos) / 2, + worker.targetPos - maxExtensionAmount, + ); + + if (closedIntervalsOverlap( + start, end, + thresholdPoint, worker.targetPos, + )) { + const size = worker.targetPos - worker.startPos; + + // If we extend by maxExtensionAmount + const a = Math.ceil((size + 1) / maxExtensionAmount) * maxExtensionAmount; + // If we extend to the next power of 2 + const b = 2 ** Math.ceil(Math.log2(size + 1)); + + const extent = Math.min(b, a); + end = Math.max(end, worker.startPos + extent); + } + } + + end = Math.max(end, start + URL_SOURCE_MIN_LOAD_AMOUNT); + + return { + start, + end, + }; + }, +} satisfies Record; + type PendingSlice = { start: number; bytes: Uint8Array; @@ -494,18 +676,15 @@ class ReadOrchestrator { constructor(public options: { maxCacheSize: number; runWorker: (worker: ReadWorker) => Promise; - getPrefetchRange: (start: number, end: number, workers: ReadWorker[]) => { - start: number; - end: number; - }; + prefetchProfile: PrefetchProfile; maxWorkerCount: number; }) {} - read(innerStart: number, innerEnd: number) { + read(innerStart: number, innerEnd: number): MaybePromise { assert(this.fileSize !== null); - const prefetchRange = this.options.getPrefetchRange(innerStart, innerEnd, this.workers); - const outerStart = prefetchRange.start; + const prefetchRange = this.options.prefetchProfile(innerStart, innerEnd, this.workers); + const outerStart = Math.max(prefetchRange.start, 0); const outerEnd = Math.min(prefetchRange.end, this.fileSize); assert(outerStart <= innerStart && innerEnd <= outerEnd); @@ -537,7 +716,7 @@ class ReadOrchestrator { let lastEnd = outerStart; // The "holes" in the cache (the parts we need to load) - const holes: { + const outerHoles: { start: number; end: number; }[] = []; @@ -558,7 +737,7 @@ class ReadOrchestrator { assert(cappedOuterStart <= cappedOuterEnd); if (lastEnd < cappedOuterStart) { - holes.push({ start: lastEnd, end: cappedOuterStart }); + outerHoles.push({ start: lastEnd, end: cappedOuterStart }); } lastEnd = cappedOuterEnd; @@ -584,10 +763,10 @@ class ReadOrchestrator { } if (lastEnd < outerEnd) { - holes.push({ start: lastEnd, end: outerEnd }); + outerHoles.push({ start: lastEnd, end: outerEnd }); } } else { - holes.push({ start: outerStart, end: outerEnd }); + outerHoles.push({ start: outerStart, end: outerEnd }); } if (bytes && contiguousBytesWriteEnd >= bytes.length) { @@ -599,7 +778,7 @@ class ReadOrchestrator { }; } - if (holes.length === 0) { + if (outerHoles.length === 0) { assert(result); return result; } @@ -607,12 +786,24 @@ class ReadOrchestrator { // We need to read more data, so now we're in async land const { promise, resolve, reject } = promiseWithResolvers(); + const innerHoles: typeof outerHoles = []; + for (const outerHole of outerHoles) { + const cappedStart = Math.max(innerStart, outerHole.start); + const cappedEnd = Math.min(innerEnd, outerHole.end); + + if (cappedStart === outerHole.start && cappedEnd === outerHole.end) { + innerHoles.push(outerHole); // Can reuse without allocating a new object + } else if (cappedStart < cappedEnd) { + innerHoles.push({ start: cappedStart, end: cappedEnd }); + } + } + // Fire off workers to take care of patching the holes - for (const hole of holes) { + for (const outerHole of outerHoles) { const pendingSlice: PendingSlice | null = bytes && { start: innerStart, bytes, - holes, // Not yet correct! These are the outer holes, not the inner holes. Will be fixed further down! + holes: innerHoles, resolve, reject, }; @@ -622,13 +813,13 @@ class ReadOrchestrator { // A small tolerance in the case that the requested region is *just* after the target position of an // existing worker. In that case, it's probably more efficient to repurpose that worker than to spawn // another one so close to it - const gapCloserTolerance = 2 ** 17; + const gapTolerance = 2 ** 17; if (closedIntervalsOverlap( - hole.start - gapCloserTolerance, hole.start, + outerHole.start - gapTolerance, outerHole.start, worker.currentPos, worker.targetPos, )) { - worker.targetPos = Math.max(worker.targetPos, hole.end); // Update the worker's target position + worker.targetPos = Math.max(worker.targetPos, outerHole.end); // Update the worker's target position workerFound = true; if (pendingSlice && !worker.pendingSlices.includes(pendingSlice)) { @@ -646,7 +837,7 @@ class ReadOrchestrator { if (!workerFound) { // We need to spawn a new worker - const newWorker = this.createWorker(hole.start, hole.end); + const newWorker = this.createWorker(outerHole.start, outerHole.end); if (pendingSlice) { newWorker.pendingSlices = [pendingSlice]; } @@ -655,19 +846,6 @@ class ReadOrchestrator { } } - // Turn the outer holes into inner holes - for (let i = 0; i < holes.length; i++) { - const hole = holes[i]!; - hole.start = Math.max(innerStart, hole.start); - hole.end = Math.min(innerEnd, hole.end); - - if (hole.end <= hole.start) { - // Empty hole - holes.splice(i, 1); - i--; - } - } - if (!result) { assert(bytes); result = promise.then(bytes => ({ @@ -797,6 +975,10 @@ class ReadOrchestrator { } insertIntoCache(entry: CacheEntry) { + if (this.options.maxCacheSize === 0) { + return; // No caching + } + let insertionIndex = binarySearchLessOrEqual(this.cache, entry.start, x => x.start) + 1; if (insertionIndex > 0) { @@ -861,21 +1043,33 @@ class ReadOrchestrator { } // LRU eviction of cache entries - while (this.currentCacheSize > this.options.maxCacheSize && this.cache.length > 1) { - let oldestIndex = 0; - let oldestEntry = this.cache[0]!; + while (this.currentCacheSize > this.options.maxCacheSize) { + if (this.cache.length > 1) { + let oldestIndex = 0; + let oldestEntry = this.cache[0]!; - for (let i = 1; i < this.cache.length; i++) { - const entry = this.cache[i]!; + for (let i = 1; i < this.cache.length; i++) { + const entry = this.cache[i]!; - if (entry.age < oldestEntry.age) { - oldestIndex = i; - oldestEntry = entry; + if (entry.age < oldestEntry.age) { + oldestIndex = i; + oldestEntry = entry; + } } - } - this.cache.splice(oldestIndex, 1); - this.currentCacheSize -= oldestEntry.bytes.length; + this.cache.splice(oldestIndex, 1); + this.currentCacheSize -= oldestEntry.bytes.length; + } else { + // The single entry that's left is too big for the cache, let's trim it + const entry = this.cache[0]!; + assert(entry.bytes.length > this.options.maxCacheSize); + + entry.bytes = entry.bytes.slice(0, this.options.maxCacheSize); + entry.view = toDataView(entry.bytes); + entry.end = entry.start + entry.bytes.length; + + this.currentCacheSize = entry.bytes.length; + } } } } diff --git a/src/wave/wave-demuxer.ts b/src/wave/wave-demuxer.ts index 6b5f65a..d8f4cfc 100644 --- a/src/wave/wave-demuxer.ts +++ b/src/wave/wave-demuxer.ts @@ -13,7 +13,7 @@ import { InputAudioTrack, InputAudioTrackBacking } from '../input-track'; import { PacketRetrievalOptions } from '../media-sink'; import { assert, UNDETERMINED_LANGUAGE } from '../misc'; import { EncodedPacket, PLACEHOLDER_DATA } from '../packet'; -import { readAscii, readBytes, Reader2, readU16, readU32, readU64 } from '../reader2'; +import { readAscii, readBytes, Reader, readU16, readU32, readU64 } from '../reader2'; export enum WaveFormat { PCM = 0x0001, @@ -24,7 +24,7 @@ export enum WaveFormat { } export class WaveDemuxer extends Demuxer { - reader: Reader2; + reader: Reader; metadataPromise: Promise | null = null; dataStart = -1; @@ -42,7 +42,7 @@ export class WaveDemuxer extends Demuxer { constructor(input: Input) { super(input); - this.reader = input._reader2; + this.reader = input._reader; } async readMetadata() {