Implement improved StreamSource, restructure source code (hah, get it?)

This commit is contained in:
Vanilagy
2025-08-29 22:50:24 +02:00
parent c93755def1
commit c222973571
13 changed files with 362 additions and 170 deletions
+322 -128
View File
@@ -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<number>;
/** @internal */
abstract _retrieveSize(): MaybePromise<number>;
/** @internal */
abstract _read(start: number, end: number): MaybePromise<ReadResult>;
/** @internal */
_sizePromise: Promise<number> | 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<Uint8Array>;
/** Called when the size of the entire file is requested. Should return or resolve to the size in bytes. */
getSize: () => number | Promise<number>;
/**
* 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<Uint8Array | ReadableStream<Uint8Array>>;
/** Called when the size of the entire file is requested. Must return or resolve to the size in bytes. */
getSize: () => MaybePromise<number>;
/** 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<number> {
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<ReadResult> {
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<ReadResult> {
return this._orchestrator.read(start, end);
}
readers = new WeakMap<ReadWorker, ReadableStreamDefaultReader<Uint8Array>>();
/** @internal */
_readers = new WeakMap<ReadWorker, ReadableStreamDefaultReader<Uint8Array>>();
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<ReadWorker, {
response: Response;
abortController: AbortController;
@@ -241,6 +391,12 @@ export class UrlSource2 extends Source {
if (options.getRetryDelay !== undefined && typeof options.getRetryDelay !== 'function') {
throw new TypeError('options.getRetryDelay, when provided, 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.');
}
super();
@@ -248,61 +404,17 @@ export class UrlSource2 extends Source {
this._options = options;
this._orchestrator = new ReadOrchestrator({
maxCacheSize: 64 * 2 ** 20, // 64 MiB
maxCacheSize: options.maxCacheSize ?? (64 * 2 ** 20 /* 64 MiB */),
// Most files in the real-world have a single sequential access pattern, but having two in parallel can
// also happen
maxWorkerCount: 2,
runWorker: this._runWorker.bind(this),
getPrefetchRange(start, end, workers) {
// Add a slight bit of start padding because
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 prefetch strategy employed for UrlSource 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,
};
},
prefetchProfile: PREFETCH_PROFILES.network,
});
}
async _retrieveSize2() {
/** @internal */
async _retrieveSize(): Promise<number> {
// 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<ReadResult> {
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<string, PrefetchProfile>;
type PendingSlice = {
start: number;
bytes: Uint8Array;
@@ -494,18 +676,15 @@ class ReadOrchestrator {
constructor(public options: {
maxCacheSize: number;
runWorker: (worker: ReadWorker) => Promise<void>;
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<ReadResult> {
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<Uint8Array>();
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;
}
}
}
}