A BUNCH of HLS (and related) progress

- Add support for fMP4 segments
- Add initInput (supported by ISOBMFF and MPEG-TS)
- Input.isSupported()
- Refactored file size retrieval logic, the file size can now be retrieved alongside the first read.
- Add support for unsized UrlSources
- Add support for multiple trun boxes in a row for the same track (as the spec literally says lol)
- Add support for more M3U8 features
- Fix faulty line reader
- MPEG-TS improvements: support video parameters not in first packet, support better key frame detection, support extension PES packets without PTS
- Fix reorder buffer being 1 element too small
This commit is contained in:
Vanilagy
2026-02-20 19:56:20 +01:00
parent de84d78092
commit 68b3ae7155
73 changed files with 17552 additions and 967 deletions
+258 -149
View File
@@ -35,6 +35,9 @@ export type ReadResult = {
offset: number;
};
export const DEFAULT_MIN_READ_POSITION = 0;
export const DEFAULT_MAX_READ_POSITION = Infinity;
/**
* The source base class, representing a resource from which bytes can be read.
* @group Input sources
@@ -42,9 +45,14 @@ export type ReadResult = {
*/
export abstract class Source {
/** @internal */
abstract _retrieveSize(): MaybePromise<number | null>;
abstract _getFileSize(): number | null | undefined;
/** @internal */
abstract _read(start: number, end: number): MaybePromise<ReadResult | null>;
abstract _read(
start: number,
end: number,
minReadPosition: number,
maxReadPosition: number,
): MaybePromise<ReadResult | null>;
/** @internal */
abstract _dispose(): void;
/** @internal */
@@ -64,7 +72,18 @@ export abstract class Source {
throw new InputDisposedError();
}
return this._sizePromise ??= Promise.resolve(this._retrieveSize());
return this._sizePromise ??= (async () => {
let size = this._getFileSize();
if (size !== undefined) {
return size;
}
await this._read(0, 1, DEFAULT_MIN_READ_POSITION, DEFAULT_MAX_READ_POSITION);
size = this._getFileSize();
assert(size !== undefined);
return size;
})();
}
/**
@@ -134,7 +153,7 @@ export class BufferSource extends Source {
}
/** @internal */
_retrieveSize(): number {
_getFileSize(): number {
return this._bytes.byteLength;
}
@@ -207,19 +226,23 @@ export class BlobSource extends Source {
runWorker: this._runWorker.bind(this),
prefetchProfile: PREFETCH_PROFILES.fileSystem,
});
this._orchestrator.fileSize = blob.size;
}
/** @internal */
_retrieveSize(): number {
const size = this._blob.size;
this._orchestrator.fileSize = size;
return size;
_getFileSize(): number {
return this._orchestrator.fileSize!; // Faster than blob.size
}
/** @internal */
_read(start: number, end: number): MaybePromise<ReadResult> {
return this._orchestrator.read(start, end);
_read(
start: number,
end: number,
minReadPosition: number,
maxReadPosition: number,
): MaybePromise<ReadResult | null> {
return this._orchestrator.read(start, end, minReadPosition, maxReadPosition);
}
/** @internal */
@@ -227,6 +250,8 @@ export class BlobSource extends Source {
/** @internal */
private async _runWorker(worker: ReadWorker) {
assert(worker.strictTarget);
let reader = this._readers.get(worker);
if (reader === undefined) {
// https://github.com/Vanilagy/mediabunny/issues/184
@@ -252,7 +277,7 @@ export class BlobSource extends Source {
if (reader) {
const { done, value } = await reader.read();
if (done) {
this._orchestrator.forgetWorker(worker);
this._orchestrator.onWorkerFinished(worker);
throw new Error('Blob reader stopped unexpectedly before all requested data was read.');
}
@@ -327,6 +352,8 @@ const DEFAULT_RETRY_DELAY
return Math.min(2 ** (previousAttempts - 2), 16);
}) satisfies UrlSourceOptions['getRetryDelay'];
const warnedOrigins = new Set<string>();
/**
* Options for {@link UrlSource}.
* @group Input sources
@@ -380,11 +407,12 @@ export class UrlSource extends Source {
_options: UrlSourceOptions;
/** @internal */
_orchestrator: ReadOrchestrator;
/** @internal */
_existingResponses = new WeakMap<ReadWorker, {
response: Response;
abortController: AbortController;
}>();
/**
* Note that this value being true does NOT mean the file size can't change anymore; it just signals that we have at
* least checked if we know the file size or not.
* @internal
*/
_fileSizeDetermined = false;
/**
* Creates a new {@link UrlSource} backed by the resource at the specified URL.
@@ -445,109 +473,111 @@ export class UrlSource extends Source {
}
/** @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
// one stone.
const abortController = new AbortController();
const response = await retriedFetch(
this._options.fetchFn ?? fetch,
this._url,
mergeRequestInit(this._options.requestInit ?? {}, {
headers: {
// We could also send a non-range request to request the same bytes (all of them), but doing it like
// this is an easy way to check if the server supports range requests in the first place
Range: 'bytes=0-',
},
signal: abortController.signal,
}),
this._getRetryDelay,
() => this._disposed,
);
if (!response.ok) {
// eslint-disable-next-line @typescript-eslint/no-base-to-string
throw new Error(`Error fetching ${String(this._url)}: ${response.status} ${response.statusText}`);
}
let worker: ReadWorker;
let fileSize: number;
if (response.status === 206) {
fileSize = this._getTotalLengthFromRangeResponse(response);
worker = this._orchestrator.createWorker(0, Math.min(fileSize, URL_SOURCE_MIN_LOAD_AMOUNT));
} else {
// Server probably returned a 200.
const contentLength = response.headers.get('Content-Length');
if (contentLength) {
fileSize = Number(contentLength);
worker = this._orchestrator.createWorker(0, fileSize);
this._orchestrator.options.maxCacheSize = Infinity; // 🤷
console.warn(
'HTTP server did not respond with 206 Partial Content, meaning the entire remote resource now has'
+ ' to be downloaded. For efficient media file streaming across a network, please make sure your'
+ ' server supports range requests.',
);
} else {
throw new Error(`HTTP response (status ${response.status}) must surface Content-Length header.`);
}
}
this._orchestrator.fileSize = fileSize;
this._existingResponses.set(worker, { response, abortController });
this._orchestrator.runWorker(worker);
return fileSize;
_getFileSize(): number | null | undefined {
return this._fileSizeDetermined
? this._orchestrator.fileSize
: undefined;
}
/** @internal */
_read(start: number, end: number): MaybePromise<ReadResult> {
return this._orchestrator.read(start, end);
_read(
start: number,
end: number,
minReadPosition: number,
maxReadPosition: number,
): MaybePromise<ReadResult | null> {
return this._orchestrator.read(start, end, minReadPosition, maxReadPosition);
}
/** @internal */
private async _runWorker(worker: ReadWorker) {
// The outer loop is for resuming a request if it dies mid-response
while (true) {
const existing = this._existingResponses.get(worker);
this._existingResponses.delete(worker);
let abortController = existing?.abortController;
let response = existing?.response;
if (!abortController) {
abortController = new AbortController();
response = await retriedFetch(
this._options.fetchFn ?? fetch,
this._url,
mergeRequestInit(this._options.requestInit ?? {}, {
headers: {
Range: `bytes=${worker.currentPos}-`,
},
signal: abortController.signal,
}),
this._getRetryDelay,
() => this._disposed,
);
}
assert(response);
const abortController = new AbortController();
const response = await retriedFetch(
this._options.fetchFn ?? fetch,
this._url,
mergeRequestInit(this._options.requestInit ?? {}, {
headers: {
// Always sending a range request is a good way to probe if the server supports them
Range: `bytes=${worker.currentPos}-`,
},
signal: abortController.signal,
}),
this._getRetryDelay,
() => this._disposed,
);
if (!response.ok) {
// eslint-disable-next-line @typescript-eslint/no-base-to-string
throw new Error(`Error fetching ${String(this._url)}: ${response.status} ${response.statusText}`);
}
if (worker.currentPos > 0 && response.status !== 206) {
throw new Error(
'HTTP server did not respond with 206 Partial Content to a range request. To enable efficient media'
+ ' file streaming across a network, please make sure your server supports range requests.',
);
outer:
if (this._orchestrator.fileSize === null) {
// See if we can deduce the file size from the response
const contentRange = response.headers.get('Content-Range');
if (contentRange) {
const match = /\/(\d+)/.exec(contentRange);
if (match) {
this._orchestrator.supplyFileSize(Number(match[1]));
break outer;
}
}
const contentLength = response.headers.get('Content-Length');
if (contentLength) {
// Note: For range requests, this is _technically_ not correct, as the range response could contain
// less data than was requested. In practice, it seems most servers don't do this though, and the
// Content-Length header actually contains the length until the end of the file.
this._orchestrator.supplyFileSize(worker.currentPos + Number(contentLength));
}
}
this._fileSizeDetermined = true; // Yes, this is correct even if file size is still null
if (response.status !== 206) {
const origin = new URL(
this._url instanceof Request ? this._url.url : this._url,
typeof window !== 'undefined' ? window.location.href : undefined,
).origin;
if (origin !== 'null') {
if (!warnedOrigins.has(origin)) {
console.warn(
`HTTP server (origin ${origin}) did not respond to a range request with 206 Partial`
+ ' Content, meaning the entire resource will now be downloaded. To enable efficient media'
+ ' file streaming across a network, please make sure your server supports range requests.',
);
warnedOrigins.add(origin);
}
}
worker.currentPos = 0;
this._orchestrator.options.maxCacheSize = Infinity; // 🤷
if (this._orchestrator.fileSize !== null) {
worker.targetPos = this._orchestrator.fileSize;
} else {
// The server is dumb, doesn't even surface the content length, but we'll work with it.
worker.targetPos = Infinity;
worker.strictTarget = false;
}
// IN CASE there are other workers (rare), merge their pending slices into this worker and abort them
for (let i = 0; i < this._orchestrator.workers.length; i++) {
const otherWorker = this._orchestrator.workers[i]!;
if (otherWorker === worker) {
continue;
}
worker.pendingSlices.push(...otherWorker.pendingSlices);
otherWorker.aborted = true;
otherWorker.pendingSlices.length = 0;
this._orchestrator.workers.splice(i, 1);
i--;
}
}
if (!response.body) {
@@ -597,15 +627,20 @@ export class UrlSource extends Source {
if (done) {
if (worker.currentPos >= worker.targetPos) {
// All data was delivered, we're good
this._orchestrator.forgetWorker(worker);
worker.running = false;
this._orchestrator.onWorkerFinished(worker);
return;
}
// The response stopped early, before the target. This can happen if server decides to cap range
// requests arbitrarily, even if the request had an uncapped end. In this case, let's fetch the rest
// of the data using a new request.
break;
if (worker.strictTarget) {
// The response stopped early, before the target. This can happen if server decides to cap range
// requests arbitrarily, even if the request had an uncapped end. In this case, let's fetch the
// rest of the data using a new request.
break;
} else {
// Assume we have simply reached the end of the resource
this._orchestrator.onWorkerFinished(worker);
return;
}
}
this.onread?.(worker.currentPos, worker.currentPos + value.length);
@@ -708,13 +743,18 @@ export class FilePathSource extends Source {
}
/** @internal */
_read(start: number, end: number): MaybePromise<ReadResult> {
return this._streamSource._read(start, end);
_read(
start: number,
end: number,
minReadPosition: number,
maxReadPosition: number,
): MaybePromise<ReadResult | null> {
return this._streamSource._read(start, end, minReadPosition, maxReadPosition);
}
/** @internal */
_retrieveSize(): MaybePromise<number> {
return this._streamSource._retrieveSize();
_getFileSize(): number | null | undefined {
return this._streamSource._getFileSize();
}
/** @internal */
@@ -815,7 +855,21 @@ export class StreamSource extends Source {
}
/** @internal */
_retrieveSize(): MaybePromise<number> {
_getFileSize(): number | null | undefined {
return this._orchestrator.fileSize ?? undefined;
}
/** @internal */
_read(
start: number,
end: number,
minReadPosition: number,
maxReadPosition: number,
): MaybePromise<ReadResult | null> {
if (this._orchestrator.fileSize !== null) {
return this._orchestrator.read(start, end, minReadPosition, maxReadPosition);
}
const result = this._options.getSize();
if (result instanceof Promise) {
@@ -825,7 +879,7 @@ export class StreamSource extends Source {
}
this._orchestrator.fileSize = size;
return size;
return this._orchestrator.read(start, end, minReadPosition, maxReadPosition);
});
} else {
if (!Number.isInteger(result) || result < 0) {
@@ -833,15 +887,10 @@ export class StreamSource extends Source {
}
this._orchestrator.fileSize = result;
return result;
return this._orchestrator.read(start, end, minReadPosition, maxReadPosition);
}
}
/** @internal */
_read(start: number, end: number): MaybePromise<ReadResult> {
return this._orchestrator.read(start, end);
}
/** @internal */
private async _runWorker(worker: ReadWorker) {
while (worker.currentPos < worker.targetPos && !worker.aborted) {
@@ -992,7 +1041,7 @@ export class ReadableStreamSource extends Source {
}
/** @internal */
_retrieveSize() {
_getFileSize(): number | null {
return this._endIndex; // Starts out as null, meaning this source is unsized
}
@@ -1251,7 +1300,7 @@ type PendingSlice = {
start: number;
end: number;
}[];
resolve: (bytes: Uint8Array) => void;
resolve: (bytes: Uint8Array | null) => void;
reject: (error: unknown) => void;
};
@@ -1267,6 +1316,8 @@ type ReadWorker = {
startPos: number;
currentPos: number;
targetPos: number;
/** The target is considered _strict_ when it is an error for the worker to terminate before reaching the target. */
strictTarget: boolean;
running: boolean;
aborted: boolean;
pendingSlices: PendingSlice[];
@@ -1296,19 +1347,18 @@ class ReadOrchestrator {
maxWorkerCount: number;
}) {}
read(innerStart: number, innerEnd: number): MaybePromise<ReadResult> {
assert(this.fileSize !== null);
read(
innerStart: number,
innerEnd: number,
minReadPosition: number,
maxReadPosition: number,
): MaybePromise<ReadResult | null> {
const prefetchRange = this.options.prefetchProfile(innerStart, innerEnd, this.workers);
const outerStart = Math.max(prefetchRange.start, 0);
const outerEnd = Math.min(prefetchRange.end, this.fileSize);
const outerStart = Math.max(prefetchRange.start, minReadPosition);
const outerEnd = Math.min(prefetchRange.end, this.fileSize ?? Infinity, maxReadPosition);
assert(outerStart <= innerStart && innerEnd <= outerEnd);
let result: MaybePromise<{
bytes: Uint8Array;
view: DataView;
offset: number;
}> | null = null;
let result: MaybePromise<ReadResult | null> | null = null;
const innerCacheStartIndex = binarySearchLessOrEqual(this.cache, innerStart, x => x.start);
const innerStartEntry = innerCacheStartIndex !== -1 ? this.cache[innerCacheStartIndex] : null;
@@ -1400,7 +1450,7 @@ class ReadOrchestrator {
}
// We need to read more data, so now we're in async land
const { promise, resolve, reject } = promiseWithResolvers<Uint8Array>();
const { promise, resolve, reject } = promiseWithResolvers<Uint8Array | null>();
const innerHoles: typeof outerHoles = [];
for (const outerHole of outerHoles) {
@@ -1454,7 +1504,8 @@ class ReadOrchestrator {
if (!workerFound) {
// We need to spawn a new worker
const newWorker = this.createWorker(outerHole.start, outerHole.end);
const strictTarget = outerHole.end < outerEnd || this.fileSize !== null;
const newWorker = this.createWorker(outerHole.start, outerHole.end, strictTarget);
if (pendingSlice) {
newWorker.pendingSlices = [pendingSlice];
}
@@ -1465,11 +1516,12 @@ class ReadOrchestrator {
if (!result) {
assert(bytes);
result = promise.then(bytes => ({
result = promise.then(bytes => bytes && ({
bytes,
view: toDataView(bytes),
offset: innerStart,
}));
} satisfies ReadResult));
} else {
// The requested region was satisfied by the cache, but the entire prefetch region was not
}
@@ -1477,11 +1529,12 @@ class ReadOrchestrator {
return result;
}
createWorker(startPos: number, targetPos: number) {
createWorker(startPos: number, targetPos: number, strictTarget: boolean) {
const worker: ReadWorker = {
startPos,
currentPos: startPos,
targetPos,
strictTarget,
running: false,
// Due to async shenanigans, it can happen that workers are started after disposal. In this case, instead of
// simply not creating the worker, we allow it to run but immediately label it as aborted, so it can then
@@ -1608,11 +1661,48 @@ class ReadOrchestrator {
}
}
forgetWorker(worker: ReadWorker) {
supplyFileSize(size: number) {
assert(this.fileSize === null);
this.fileSize = size;
// Trim the workers with this new information
for (const worker of this.workers) {
worker.targetPos = Math.min(worker.targetPos, size);
worker.strictTarget = true;
for (let i = 0; i < worker.pendingSlices.length; i++) {
const pendingSlice = worker.pendingSlices[i]!;
for (const hole of pendingSlice.holes) {
if (hole.end > size) {
// Can't satisfy this slice anymore
pendingSlice.resolve(null);
worker.pendingSlices.splice(i, 1);
i--;
break;
}
}
}
}
}
/** Called when a worker reaches the end of the underlying data and must be cleaned up. */
onWorkerFinished(worker: ReadWorker) {
const index = this.workers.indexOf(worker);
assert(index !== -1);
this.workers.splice(index, 1);
if (this.fileSize === null) {
// We can now deduce the file size!
this.supplyFileSize(worker.currentPos);
}
for (const pendingSlice of worker.pendingSlices) {
pendingSlice.resolve(null);
}
}
insertIntoCache(entry: CacheEntry) {
@@ -1723,12 +1813,11 @@ class ReadOrchestrator {
* from another source.
*/
export class NullSource extends Source {
override _retrieveSize(): MaybePromise<number | null> {
override _getFileSize(): number | null {
return null;
}
// eslint-disable-next-line @typescript-eslint/no-unused-vars
override _read(start: number, end: number): MaybePromise<ReadResult | null> {
override _read(): MaybePromise<ReadResult | null> {
return null;
}
@@ -1753,21 +1842,41 @@ export class RangedSource extends Source {
this._length = length ?? null;
}
override async _retrieveSize(): Promise<number | null> {
const baseSize = await this._baseSource.getSizeOrNull(); // Call getSizeOrNull for memoization
override _getFileSize(): number | null | undefined {
const baseSize = this._baseSource._getFileSize();
if (baseSize === undefined) {
return this._length !== null
? this._length
: undefined;
}
if (baseSize === null) {
return null;
if (this._length !== null) {
return this._length;
} else {
return null;
}
}
return clamp(baseSize - this._offset, 0, this._length ?? Infinity);
}
override _read(start: number, end: number): MaybePromise<ReadResult | null> {
override _read(
start: number,
end: number,
minReadPosition: number,
maxReadPosition: number,
): MaybePromise<ReadResult | null> {
if (this._length !== null && end > this._length) {
return null;
}
const result = this._baseSource._read(this._offset + start, this._offset + end);
const result = this._baseSource._read(
this._offset + start,
this._offset + end,
this._offset + minReadPosition,
this._offset + maxReadPosition,
);
if (result instanceof Promise) {
return result.then((result) => {