Add logic to handle HTTP requests dying mid-response, fix a few ReadOrchestrator things

This commit is contained in:
Vanilagy
2025-08-30 14:57:27 +02:00
parent 99051e1848
commit 896d41d60a
+108 -55
View File
@@ -212,6 +212,8 @@ export class UrlSource extends Source {
/** @internal */ /** @internal */
_url: URL; _url: URL;
/** @internal */ /** @internal */
_getRetryDelay: (previousAttempts: number) => number | null;
/** @internal */
_options: UrlSourceOptions; _options: UrlSourceOptions;
/** @internal */ /** @internal */
_orchestrator: ReadOrchestrator; _orchestrator: ReadOrchestrator;
@@ -248,6 +250,7 @@ export class UrlSource extends Source {
this._url = url instanceof URL ? url : new URL(url, location.href); this._url = url instanceof URL ? url : new URL(url, location.href);
this._options = options; this._options = options;
this._getRetryDelay = options.getRetryDelay ?? (previousAttempts => Math.min(2 ** (previousAttempts - 2), 8));
this._orchestrator = new ReadOrchestrator({ this._orchestrator = new ReadOrchestrator({
maxCacheSize: options.maxCacheSize ?? (64 * 2 ** 20 /* 64 MiB */), maxCacheSize: options.maxCacheSize ?? (64 * 2 ** 20 /* 64 MiB */),
@@ -277,7 +280,7 @@ export class UrlSource extends Source {
}, },
signal: abortController.signal, signal: abortController.signal,
}), }),
this._options.getRetryDelay ?? (() => null), this._getRetryDelay,
); );
if (!response.ok) { if (!response.ok) {
@@ -324,64 +327,94 @@ export class UrlSource extends Source {
/** @internal */ /** @internal */
private async _runWorker(worker: ReadWorker) { private async _runWorker(worker: ReadWorker) {
const existing = this._existingResponses.get(worker); // The outer loop is for resuming a request if it dies mid-response
while (!worker.aborted) {
const existing = this._existingResponses.get(worker);
this._existingResponses.delete(worker);
let abortController = existing?.abortController; let abortController = existing?.abortController;
let response = existing?.response; let response = existing?.response;
if (!abortController) { if (!abortController) {
abortController = new AbortController(); abortController = new AbortController();
response = await retriedFetch( response = await retriedFetch(
this._url, this._url,
mergeObjectsDeeply(this._options.requestInit ?? {}, { mergeObjectsDeeply(this._options.requestInit ?? {}, {
headers: { headers: {
Range: `bytes=${worker.currentPos}-`, Range: `bytes=${worker.currentPos}-`,
}, },
signal: abortController.signal, signal: abortController.signal,
}), }),
this._options.getRetryDelay ?? (() => null), this._getRetryDelay,
); );
}
assert(response);
if (!response.ok) {
throw new Error(`Error fetching ${this._url}: ${response.status} ${response.statusText}`);
}
const length = this._getPartialLengthFromRangeResponse(response);
const required = worker.targetPos - worker.currentPos;
if (length < required) {
throw new Error(
`HTTP response unexpectedly too short: Needed at least ${required} bytes, got only ${length}.`,
);
}
if (!response.body) {
throw new Error('Missing HTTP response body.');
}
const reader = response.body.getReader();
while (true) {
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); assert(response);
this._orchestrator.supplyWorkerData(worker, value);
if (worker.currentPos >= worker.targetPos || worker.aborted) { if (!response.ok) {
abortController.abort(); throw new Error(`Error fetching ${this._url}: ${response.status} ${response.statusText}`);
this._existingResponses.delete(worker); }
break;
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.',
);
}
const length = this._getPartialLengthFromRangeResponse(response);
const required = worker.targetPos - worker.currentPos;
if (length < required) {
throw new Error(
`HTTP response unexpectedly too short: Needed at least ${required} bytes, got only ${length}.`,
);
}
if (!response.body) {
throw new Error('Missing HTTP response body.');
}
const reader = response.body.getReader();
while (true) {
let readResult: ReadableStreamReadResult<Uint8Array>;
try {
readResult = await reader.read();
} catch (error) {
const retryDelayInSeconds = this._getRetryDelay(1);
if (retryDelayInSeconds !== null) {
console.error('Error while reading response stream. Attempting to resume.', error);
await new Promise(resolve => setTimeout(resolve, 1000 * retryDelayInSeconds));
break;
} else {
throw error;
}
}
const { done, value } = readResult;
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.',
);
}
return;
}
this.onread?.(worker.currentPos, worker.currentPos + value.length);
this._orchestrator.supplyWorkerData(worker, value);
if (worker.currentPos >= worker.targetPos || worker.aborted) {
abortController.abort();
return;
}
} }
} }
@@ -913,13 +946,16 @@ class ReadOrchestrator {
worker.age = this.nextAge++; worker.age = this.nextAge++;
void this.options.runWorker(worker) void this.options.runWorker(worker)
.then(() => worker.running = false)
.catch((error) => { .catch((error) => {
if (worker.pendingSlices.length > 0) { if (worker.pendingSlices.length > 0) {
worker.pendingSlices.forEach(x => x.reject(error)); // Make sure to propagate any errors worker.pendingSlices.forEach(x => x.reject(error)); // Make sure to propagate any errors
worker.pendingSlices.length = 0;
} else { } else {
throw error; // So it doesn't get swallowed throw error; // So it doesn't get swallowed
} }
})
.finally(() => {
worker.running = false;
}); });
} }
@@ -936,6 +972,7 @@ class ReadOrchestrator {
age: this.nextAge++, age: this.nextAge++,
}); });
worker.currentPos += bytes.length; worker.currentPos += bytes.length;
worker.targetPos = Math.max(worker.targetPos, worker.currentPos); // In case it overshoots
// Now, let's see if we can use the read bytes to fill any pending slice // Now, let's see if we can use the read bytes to fill any pending slice
for (let i = 0; i < worker.pendingSlices.length; i++) { for (let i = 0; i < worker.pendingSlices.length; i++) {
@@ -973,6 +1010,22 @@ class ReadOrchestrator {
i--; i--;
} }
} }
// Remove other idle workers if we "ate" into their territory
for (let i = 0; i < this.workers.length; i++) {
const otherWorker = this.workers[i]!;
if (worker === otherWorker || otherWorker.running) {
continue;
}
if (closedIntervalsOverlap(
start, end,
otherWorker.currentPos, otherWorker.targetPos, // These should typically be equal when the worker's idle
)) {
this.workers.splice(i, 1);
i--;
}
}
} }
forgetWorker(worker: ReadWorker) { forgetWorker(worker: ReadWorker) {