Make the download deadline an idle deadline and cancel failed bodies (closes #24)
check / check (push) Successful in 17s
check / check (push) Successful in 17s
downloadTimeoutMs now aborts a file or thumbnail download only when no bytes have arrived for that long (default 60 s, was a 600 s cap on the whole transfer), so a slow download that keeps making progress completes. The abort reason is still a TimeoutError, so retry classification is unchanged. A download that fails before reading the whole response body cancels it, including when the temp file cannot be opened or the header is malformed, so a failed file no longer holds its connection. Model: opus-5-5
This commit was merged in pull request #88.
This commit is contained in:
+59
-27
@@ -19,14 +19,12 @@ const DEFAULT_FILES_ORIGIN = "https://files.ente.io";
|
||||
const DEFAULT_THUMBS_ORIGIN = "https://thumbnails.ente.io";
|
||||
const CLIENT_PACKAGE = "berlin.sneak.quak";
|
||||
|
||||
// Two deadlines rather than one, because a single number cannot serve both
|
||||
// jobs. Thirty seconds is generous for a JSON call and short enough that a
|
||||
// hung API connection cannot stall a backup for long. A file body is a
|
||||
// different shape of problem: the deadline has to cover the whole transfer,
|
||||
// which for a large video on a slow link is minutes, so a value sane for JSON
|
||||
// would cancel legitimate downloads.
|
||||
// Two deadlines of different kinds. `requestTimeoutMs` bounds a whole JSON
|
||||
// call. `downloadTimeoutMs` is an idle deadline: a file or thumbnail download
|
||||
// is aborted only when no bytes have arrived for that long, so a large video on
|
||||
// a slow link that keeps making progress is never cut off.
|
||||
export const DEFAULT_REQUEST_TIMEOUT_MS = 30_000;
|
||||
export const DEFAULT_DOWNLOAD_TIMEOUT_MS = 600_000;
|
||||
export const DEFAULT_DOWNLOAD_TIMEOUT_MS = 60_000;
|
||||
|
||||
export interface ApiClientOptions {
|
||||
apiOrigin?: string;
|
||||
@@ -48,18 +46,43 @@ export interface StreamOptions {
|
||||
retry?: boolean;
|
||||
}
|
||||
|
||||
// Enforce a deadline over a response body, not merely over its headers.
|
||||
// An abort signal that fires once `ms` pass without a call to `restart`. It
|
||||
// aborts with a `TimeoutError`, the same reason `AbortSignal.timeout()` gives,
|
||||
// so the retry classifier treats an idle download exactly as it treats any
|
||||
// other deadline. `stop` must be called when the download ends, or the timer
|
||||
// keeps the process alive until it fires.
|
||||
const idleDeadline = (ms: number) => {
|
||||
const controller = new AbortController();
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
const stop = (): void => clearTimeout(timer);
|
||||
const restart = (): void => {
|
||||
stop();
|
||||
timer = setTimeout(() => {
|
||||
controller.abort(
|
||||
new DOMException(
|
||||
`download stalled: no bytes received for ${ms} ms`,
|
||||
"TimeoutError",
|
||||
),
|
||||
);
|
||||
}, ms);
|
||||
};
|
||||
restart();
|
||||
return { signal: controller.signal, restart, stop };
|
||||
};
|
||||
|
||||
// Enforce the idle deadline over a response body, not merely over its headers.
|
||||
//
|
||||
// `getFileStream` returns as soon as headers arrive; the bytes are pulled
|
||||
// later, in the download layer. Whether the signal passed to `fetch` also
|
||||
// tears down the body afterwards is up to the fetch implementation, so this
|
||||
// wrapper makes it a property of quak instead: every read races the signal,
|
||||
// and an abort errors the stream with the abort reason — which the retry
|
||||
// classifier recognises.
|
||||
// each chunk that arrives restarts the deadline, and an abort errors the
|
||||
// stream with the abort reason — which the retry classifier recognises.
|
||||
const deadlineStream = (
|
||||
body: ReadableStream<Uint8Array>,
|
||||
signal: AbortSignal,
|
||||
deadline: ReturnType<typeof idleDeadline>,
|
||||
): ReadableStream<Uint8Array> => {
|
||||
const { signal } = deadline;
|
||||
const reader = body.getReader();
|
||||
let rejectOnAbort: (reason: unknown) => void = () => undefined;
|
||||
const aborted = new Promise<never>((_resolve, reject) => {
|
||||
@@ -73,7 +96,10 @@ const deadlineStream = (
|
||||
const onAbort = (): void => rejectOnAbort(signal.reason);
|
||||
if (signal.aborted) onAbort();
|
||||
else signal.addEventListener("abort", onAbort, { once: true });
|
||||
const release = (): void => signal.removeEventListener("abort", onAbort);
|
||||
const release = (): void => {
|
||||
deadline.stop();
|
||||
signal.removeEventListener("abort", onAbort);
|
||||
};
|
||||
|
||||
return new ReadableStream<Uint8Array>({
|
||||
async pull(controller) {
|
||||
@@ -84,6 +110,7 @@ const deadlineStream = (
|
||||
controller.close();
|
||||
return;
|
||||
}
|
||||
deadline.restart();
|
||||
controller.enqueue(next.value);
|
||||
} catch (err) {
|
||||
release();
|
||||
@@ -362,22 +389,27 @@ export class ApiClient {
|
||||
opts?: StreamOptions,
|
||||
): Promise<ReadableStream<Uint8Array>> {
|
||||
const once = async (): Promise<ReadableStream<Uint8Array>> => {
|
||||
// A fresh deadline per attempt, so a retry gets the whole budget
|
||||
// rather than the remainder of the one that just expired.
|
||||
const signal = AbortSignal.timeout(this.downloadTimeoutMs);
|
||||
const resp = await this._fetch(url, {
|
||||
method: "GET",
|
||||
headers: this.headers(),
|
||||
signal,
|
||||
});
|
||||
await this.throwIfError(resp);
|
||||
if (!resp.body) {
|
||||
// Carries the status, and is not retryable: a response that
|
||||
// arrived without a body is malformed, and asking again
|
||||
// produces the same malformed response.
|
||||
throw new ApiError("response body is null", resp.status);
|
||||
// A fresh deadline per attempt. It also covers the wait for the
|
||||
// headers, when no bytes have arrived either.
|
||||
const deadline = idleDeadline(this.downloadTimeoutMs);
|
||||
try {
|
||||
const resp = await this._fetch(url, {
|
||||
method: "GET",
|
||||
headers: this.headers(),
|
||||
signal: deadline.signal,
|
||||
});
|
||||
await this.throwIfError(resp);
|
||||
if (!resp.body) {
|
||||
// Carries the status, and is not retryable: a response
|
||||
// that arrived without a body is malformed, and asking
|
||||
// again produces the same malformed response.
|
||||
throw new ApiError("response body is null", resp.status);
|
||||
}
|
||||
return deadlineStream(resp.body, deadline);
|
||||
} catch (err) {
|
||||
deadline.stop();
|
||||
throw err;
|
||||
}
|
||||
return deadlineStream(resp.body, signal);
|
||||
};
|
||||
return opts?.retry === false ? once() : withRetry(once, this.retry);
|
||||
}
|
||||
|
||||
+63
-47
@@ -98,47 +98,53 @@ const streamDecrypt = async (
|
||||
onProgress?.(totalPlain);
|
||||
};
|
||||
|
||||
for (;;) {
|
||||
const { done, value } = await reader.read();
|
||||
if (value && value.length > 0) {
|
||||
pending.push(value);
|
||||
pendingBytes += value.length;
|
||||
}
|
||||
try {
|
||||
for (;;) {
|
||||
const { done, value } = await reader.read();
|
||||
if (value && value.length > 0) {
|
||||
pending.push(value);
|
||||
pendingBytes += value.length;
|
||||
}
|
||||
|
||||
while (pendingBytes >= ENC_CHUNK_SIZE) {
|
||||
const encChunk = takeContiguous(ENC_CHUNK_SIZE);
|
||||
// A whole chunk that fails to authenticate while the stream carries
|
||||
// on is corruption, not truncation; that error propagates unchanged.
|
||||
const { plaintext, tag } = pullStreamChunk(state, encChunk);
|
||||
await consume(plaintext, tag);
|
||||
}
|
||||
while (pendingBytes >= ENC_CHUNK_SIZE) {
|
||||
const encChunk = takeContiguous(ENC_CHUNK_SIZE);
|
||||
// A whole chunk that fails to authenticate while the stream
|
||||
// carries on is corruption, not truncation; that error
|
||||
// propagates unchanged.
|
||||
const { plaintext, tag } = pullStreamChunk(state, encChunk);
|
||||
await consume(plaintext, tag);
|
||||
}
|
||||
|
||||
if (done) {
|
||||
if (pendingBytes > 0) {
|
||||
const buffer = takeContiguous(pendingBytes);
|
||||
// Whatever is left over once every whole chunk has been
|
||||
// consumed must be the stream's final chunk, and a final
|
||||
// chunk that actually arrived in full authenticates. If it
|
||||
// does not, the body stopped part-way through a chunk — the
|
||||
// ordinary shape of a dropped connection. Poly1305 cannot
|
||||
// tell a partial chunk from a corrupt one, so this is
|
||||
// reported as the truncation it almost always is, with the
|
||||
// authentication failure kept as the error's cause. Only the
|
||||
// pull is guarded: a sink failure on a chunk that did
|
||||
// authenticate is a disk error, not a truncation.
|
||||
let pulled;
|
||||
try {
|
||||
pulled = pullStreamChunk(state, buffer);
|
||||
} catch (err) {
|
||||
throw new TruncatedStreamError(
|
||||
`download: stream truncated: response body ended with ${buffer.length} trailing bytes that did not authenticate as a final chunk (transfer stopped mid-chunk, or the data is corrupt)`,
|
||||
{ cause: err },
|
||||
);
|
||||
if (done) {
|
||||
if (pendingBytes > 0) {
|
||||
const buffer = takeContiguous(pendingBytes);
|
||||
// Whatever is left over once every whole chunk has been
|
||||
// consumed must be the stream's final chunk, and a final
|
||||
// chunk that actually arrived in full authenticates. If
|
||||
// it does not, the body stopped part-way through a chunk
|
||||
// — the ordinary shape of a dropped connection. Poly1305
|
||||
// cannot tell a partial chunk from a corrupt one, so this
|
||||
// is reported as the truncation it almost always is, with
|
||||
// the authentication failure kept as the error's cause.
|
||||
// Only the pull is guarded: a sink failure on a chunk
|
||||
// that did authenticate is a disk error, not a
|
||||
// truncation.
|
||||
let pulled;
|
||||
try {
|
||||
pulled = pullStreamChunk(state, buffer);
|
||||
} catch (err) {
|
||||
throw new TruncatedStreamError(
|
||||
`download: stream truncated: response body ended with ${buffer.length} trailing bytes that did not authenticate as a final chunk (transfer stopped mid-chunk, or the data is corrupt)`,
|
||||
{ cause: err },
|
||||
);
|
||||
}
|
||||
await consume(pulled.plaintext, pulled.tag);
|
||||
}
|
||||
await consume(pulled.plaintext, pulled.tag);
|
||||
break;
|
||||
}
|
||||
break;
|
||||
}
|
||||
} finally {
|
||||
reader.releaseLock();
|
||||
}
|
||||
|
||||
// Only the last chunk of a secretstream carries TAG_FINAL. Everything a
|
||||
@@ -245,17 +251,27 @@ const decryptToTemp = async (
|
||||
onProgress?: ProgressCallback,
|
||||
): Promise<number> => {
|
||||
let bytesWritten = 0;
|
||||
await stageAtomic(destination, async (handle) => {
|
||||
bytesWritten = await streamDecrypt(
|
||||
stream,
|
||||
header,
|
||||
key,
|
||||
async (plaintext) => {
|
||||
await handle.write(plaintext);
|
||||
},
|
||||
onProgress,
|
||||
);
|
||||
});
|
||||
try {
|
||||
await stageAtomic(destination, async (handle) => {
|
||||
bytesWritten = await streamDecrypt(
|
||||
stream,
|
||||
header,
|
||||
key,
|
||||
async (plaintext) => {
|
||||
await handle.write(plaintext);
|
||||
},
|
||||
onProgress,
|
||||
);
|
||||
});
|
||||
} catch (err) {
|
||||
// Cancel the body so its connection is closed now rather than held
|
||||
// until the stream is garbage collected. A backup run carries on past
|
||||
// a failed file, so without this every failure would hold a socket.
|
||||
// This covers every failure, including a temp file that cannot be
|
||||
// opened and a header that is rejected before the body is read.
|
||||
await stream.cancel(err).catch(() => undefined);
|
||||
throw err;
|
||||
}
|
||||
return bytesWritten;
|
||||
};
|
||||
|
||||
|
||||
Reference in New Issue
Block a user