diff --git a/src/adts/adts-demuxer.ts b/src/adts/adts-demuxer.ts index 43e10a0..30f0ef0 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, Reader } from '../reader2'; +import { readBytes, Reader } from '../reader'; import { FrameHeader, MAX_FRAME_HEADER_SIZE, MIN_FRAME_HEADER_SIZE, readFrameHeader } from './adts-reader'; const SAMPLES_PER_AAC_FRAME = 1024; diff --git a/src/adts/adts-reader.ts b/src/adts/adts-reader.ts index 113dc59..bf21f6a 100644 --- a/src/adts/adts-reader.ts +++ b/src/adts/adts-reader.ts @@ -7,7 +7,7 @@ */ import { Bitstream } from '../misc'; -import { FileSlice, readBytes } from '../reader2'; +import { FileSlice, readBytes } from '../reader'; export const MIN_FRAME_HEADER_SIZE = 7; export const MAX_FRAME_HEADER_SIZE = 9; diff --git a/src/input-format.ts b/src/input-format.ts index 46cb6e5..1d4f17b 100644 --- a/src/input-format.ts +++ b/src/input-format.ts @@ -27,7 +27,7 @@ import { OggDemuxer } from './ogg/ogg-demuxer'; import { WaveDemuxer } from './wave/wave-demuxer'; import { MAX_FRAME_HEADER_SIZE, MIN_FRAME_HEADER_SIZE, readFrameHeader } from './adts/adts-reader'; import { AdtsDemuxer } from './adts/adts-demuxer'; -import { readAscii } from './reader2'; +import { readAscii } from './reader'; /** * Base class representing an input media file format. diff --git a/src/input.ts b/src/input.ts index 792db02..d70eea7 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 { Reader } from './reader2'; +import { Reader } from './reader'; import { Source } from './source'; /** diff --git a/src/isobmff/isobmff-demuxer.ts b/src/isobmff/isobmff-demuxer.ts index 4365e44..0fa784d 100644 --- a/src/isobmff/isobmff-demuxer.ts +++ b/src/isobmff/isobmff-demuxer.ts @@ -82,7 +82,7 @@ import { readU64Be, readU8, readAscii, -} from '../reader2'; +} from '../reader'; type InternalTrack = { id: number; diff --git a/src/isobmff/isobmff-reader.ts b/src/isobmff/isobmff-reader.ts index b3ffa29..f1d0605 100644 --- a/src/isobmff/isobmff-reader.ts +++ b/src/isobmff/isobmff-reader.ts @@ -6,7 +6,7 @@ * file, You can obtain one at https://mozilla.org/MPL/2.0/. */ -import { FileSlice, readAscii, readI32Be, readU32Be, readU64Be, readU8 } from '../reader2'; +import { FileSlice, readAscii, readI32Be, readU32Be, readU64Be, readU8 } from '../reader'; export const MIN_BOX_HEADER_SIZE = 8; export const MAX_BOX_HEADER_SIZE = 16; diff --git a/src/matroska/ebml.ts b/src/matroska/ebml.ts index 0ccbb3f..6b14a42 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, Reader, readF32Be, readF64Be, readU8 } from '../reader2'; +import { FileSlice, readBytes, Reader, readF32Be, readF64Be, readU8 } from '../reader'; import { Writer } from '../writer'; export interface EBMLElement { diff --git a/src/matroska/matroska-demuxer.ts b/src/matroska/matroska-demuxer.ts index 92267c0..0310131 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, Reader, readI16Be, readU8 } from '../reader2'; +import { FileSlice, readBytes, Reader, readI16Be, readU8 } from '../reader'; type Segment = { seekHeadSeen: boolean; diff --git a/src/mp3/mp3-demuxer.ts b/src/mp3/mp3-demuxer.ts index 0d4f1c9..9e9b511 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, Reader, readU32Be } from '../reader2'; +import { readBytes, Reader, readU32Be } from '../reader'; type Sample = { timestamp: number; diff --git a/src/mp3/mp3-reader.ts b/src/mp3/mp3-reader.ts index 884b4db..a28159c 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, Reader, readU32Be } from '../reader2'; +import { FileSlice, readAscii, Reader, readU32Be } from '../reader'; export const readId3 = (slice: FileSlice) => { const tag = readAscii(slice, 3); diff --git a/src/ogg/ogg-demuxer.ts b/src/ogg/ogg-demuxer.ts index c063392..17b70ed 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, Reader } from '../reader2'; +import { readBytes, Reader } from '../reader'; import { buildOggMimeType, computeOggPageCrc, extractSampleMetadata, OggCodecInfo } from './ogg-misc'; import { findNextPageHeader, diff --git a/src/ogg/ogg-reader.ts b/src/ogg/ogg-reader.ts index 34a4eae..011b126 100644 --- a/src/ogg/ogg-reader.ts +++ b/src/ogg/ogg-reader.ts @@ -6,7 +6,7 @@ * file, You can obtain one at https://mozilla.org/MPL/2.0/. */ -import { FileSlice, readI64Le, readU32Le, readU8 } from '../reader2'; +import { FileSlice, readI64Le, readU32Le, readU8 } from '../reader'; import { OGGS } from './ogg-misc'; export const MIN_PAGE_HEADER_SIZE = 27; diff --git a/src/reader2.ts b/src/reader.ts similarity index 95% rename from src/reader2.ts rename to src/reader.ts index 2ff6b33..cfcab58 100644 --- a/src/reader2.ts +++ b/src/reader.ts @@ -1,3 +1,11 @@ +/*! + * 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 { clamp, MaybePromise, toDataView } from './misc'; import { Source } from './source'; diff --git a/src/source.ts b/src/source.ts index 3d841cc..81caa88 100644 --- a/src/source.ts +++ b/src/source.ts @@ -93,165 +93,6 @@ export class BufferSource extends Source { } } -/** - * Options for defining a StreamSource. - * @public - */ -export type StreamSourceOptions = { - /** - * 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'; -}; - -/** - * A general-purpose, callback-driven source that can get its data from anywhere. - * @public - */ -export class StreamSource extends Source { - /** @internal */ - _options: StreamSourceOptions; - /** @internal */ - _orchestrator: ReadOrchestrator; - - constructor(options: StreamSourceOptions) { - if (!options || typeof options !== 'object') { - throw new TypeError('options must be an object.'); - } - if (typeof options.read !== 'function') { - throw new TypeError('options.read must be a function.'); - } - 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 */ - _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 */ - _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; @@ -566,6 +407,165 @@ export class UrlSource extends Source { } } +/** + * Options for defining a StreamSource. + * @public + */ +export type StreamSourceOptions = { + /** + * 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'; +}; + +/** + * A general-purpose, callback-driven source that can get its data from anywhere. + * @public + */ +export class StreamSource extends Source { + /** @internal */ + _options: StreamSourceOptions; + /** @internal */ + _orchestrator: ReadOrchestrator; + + constructor(options: StreamSourceOptions) { + if (!options || typeof options !== 'object') { + throw new TypeError('options must be an object.'); + } + if (typeof options.read !== 'function') { + throw new TypeError('options.read must be a function.'); + } + 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 */ + _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 */ + _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.'); + } + } + } +} + type PrefetchProfile = (start: number, end: number, workers: ReadWorker[]) => { start: number; end: number; diff --git a/src/wave/wave-demuxer.ts b/src/wave/wave-demuxer.ts index d8f4cfc..bfbd669 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, Reader, readU16, readU32, readU64 } from '../reader2'; +import { readAscii, readBytes, Reader, readU16, readU32, readU64 } from '../reader'; export enum WaveFormat { PCM = 0x0001,