From 17a8fe87c209afa4c05d3446b98e1cfd07af84ae Mon Sep 17 00:00:00 2001 From: Spence Pope Date: Wed, 5 Aug 2026 10:15:16 -0400 Subject: [PATCH] APP-2793: harden multipart video uploads (#11366) --- src/lib/media/video/multipart/api.ts | 7 ++- src/lib/media/video/multipart/constants.ts | 13 ++++- src/lib/media/video/multipart/upload.ts | 32 +++++++++-- src/lib/media/video/multipart/uploadPart.ts | 53 ++++++++++++++++--- .../media/video/multipart/uploadParts.test.ts | 31 +++++++++++ src/lib/media/video/multipart/uploadParts.ts | 7 ++- src/lib/media/video/multipart/utils.ts | 6 +++ 7 files changed, 133 insertions(+), 16 deletions(-) diff --git a/src/lib/media/video/multipart/api.ts b/src/lib/media/video/multipart/api.ts index 811696b3e3..c6a268ffe1 100644 --- a/src/lib/media/video/multipart/api.ts +++ b/src/lib/media/video/multipart/api.ts @@ -111,10 +111,15 @@ export function getUploadStatus( }) } -export function abortUpload(jobId: string, token: string) { +export function abortUpload( + jobId: string, + token: string, + signal?: AbortSignal, +) { return request({ route: '/xrpc/app.bsky.video.abortUpload', token, + signal, body: {jobId}, }) } diff --git a/src/lib/media/video/multipart/constants.ts b/src/lib/media/video/multipart/constants.ts index 2aaa37b8c2..2592a9d7db 100644 --- a/src/lib/media/video/multipart/constants.ts +++ b/src/lib/media/video/multipart/constants.ts @@ -4,10 +4,19 @@ */ /** Max parts uploaded concurrently. */ -export const MULTIPART_CONCURRENCY = 3 +export const MULTIPART_CONCURRENCY = 4 /** Per-part upload attempts before the part (and the upload) fails. */ -export const MULTIPART_MAX_ATTEMPTS = 3 +export const MULTIPART_MAX_ATTEMPTS = 5 + +/** Maximum time to wait for an individual part request to settle. */ +export const MULTIPART_PART_TIMEOUT_MS = 120_000 + +/** Maximum time to wait for each best-effort abort request. */ +export const MULTIPART_ABORT_TIMEOUT_MS = 10_000 + +/** Attempts to release a failed multipart upload reservation. */ +export const MULTIPART_ABORT_ATTEMPTS = 3 /** Attempts to begin/continue server-side finalization before checking state. */ export const MULTIPART_FINISH_ATTEMPTS = 3 diff --git a/src/lib/media/video/multipart/upload.ts b/src/lib/media/video/multipart/upload.ts index af2d71d596..2b45e6d2c6 100644 --- a/src/lib/media/video/multipart/upload.ts +++ b/src/lib/media/video/multipart/upload.ts @@ -14,12 +14,16 @@ import { MultipartUploadError, startUpload, } from './api' -import {MULTIPART_FINISH_ATTEMPTS} from './constants' +import { + MULTIPART_ABORT_ATTEMPTS, + MULTIPART_ABORT_TIMEOUT_MS, + MULTIPART_FINISH_ATTEMPTS, +} from './constants' import {getMissingParts, planParts} from './planParts' import {createChunkReader} from './readChunk' import {createUploadPart} from './uploadPart' import {uploadParts} from './uploadParts' -import {delay, isRetryableMultipartError} from './utils' +import {delay, isRetryableMultipartError, retryDelayMs} from './utils' export class MultipartFallbackError extends Error {} @@ -221,7 +225,7 @@ async function abortThenFallbackOrResolve( token: string, cause: unknown, ): Promise { - const result = await abortUpload(jobId, token) + const result = await abortUploadWithRetry(jobId, token) if (result.state === 'aborted') { throw new MultipartFallbackError( cause instanceof Error ? cause.message : 'Multipart upload failed', @@ -238,6 +242,28 @@ async function abortThenFallbackOrResolve( ) } +async function abortUploadWithRetry(jobId: string, token: string) { + let lastError: unknown + for (let attempt = 1; attempt <= MULTIPART_ABORT_ATTEMPTS; attempt++) { + const controller = new AbortController() + const timer = setTimeout( + () => controller.abort(), + MULTIPART_ABORT_TIMEOUT_MS, + ) + try { + return await abortUpload(jobId, token, controller.signal) + } catch (err) { + lastError = err + if (attempt < MULTIPART_ABORT_ATTEMPTS) { + await delay(retryDelayMs(attempt), new AbortController().signal) + } + } finally { + clearTimeout(timer) + } + } + throw lastError +} + function createTokenProvider(agent: AtpAgent, signal: AbortSignal) { let token: string | undefined let expiresAt = 0 diff --git a/src/lib/media/video/multipart/uploadPart.ts b/src/lib/media/video/multipart/uploadPart.ts index e3c57fc2b2..d0a910682c 100644 --- a/src/lib/media/video/multipart/uploadPart.ts +++ b/src/lib/media/video/multipart/uploadPart.ts @@ -1,6 +1,7 @@ import {AbortError} from '#/lib/async/cancelable' import {createVideoEndpointUrl} from '#/lib/media/video/util' import {MultipartUploadError} from './api' +import {MULTIPART_PART_TIMEOUT_MS} from './constants' import {type UploadPartFn} from './types' export function createUploadPart( @@ -34,23 +35,40 @@ function sendPart( return } const xhr = new XMLHttpRequest() + xhr.timeout = MULTIPART_PART_TIMEOUT_MS const abort = () => xhr.abort() signal.addEventListener('abort', abort, {once: true}) - const cleanup = () => signal.removeEventListener('abort', abort) + let settled = false + const cleanup = () => { + signal.removeEventListener('abort', abort) + xhr.onreadystatechange = null + } + const rejectOnce = (err: Error) => { + if (settled) return + settled = true + cleanup() + reject(err) + } + const resolveOnce = (result: Awaited>) => { + if (settled) return + settled = true + cleanup() + resolve(result) + } xhr.upload.addEventListener('progress', event => { onProgress(event.loaded) }) xhr.onerror = () => { - cleanup() - reject(new TypeError('Network request failed')) + rejectOnce(new TypeError('Network request failed')) + } + xhr.ontimeout = () => { + rejectOnce(new TypeError('Multipart part upload timed out')) } xhr.onabort = () => { - cleanup() - reject(new AbortError()) + rejectOnce(new AbortError()) } xhr.onload = () => { - cleanup() let data: { partNumber?: number sizeBytes?: number @@ -63,7 +81,7 @@ function sendPart( data = {} } if (xhr.status < 200 || xhr.status >= 300) { - reject( + rejectOnce( new MultipartUploadError( data.message || data.error || @@ -74,12 +92,31 @@ function sendPart( ) } else { onProgress(part.size) - resolve({ + resolveOnce({ partNumber: data.partNumber ?? part.partNumber, sizeBytes: data.sizeBytes ?? part.size, }) } } + xhr.onreadystatechange = () => { + if ( + xhr.readyState === XMLHttpRequest.HEADERS_RECEIVED && + xhr.status >= 400 + ) { + // React Native does not dispatch `load` until the response body has + // completed. Reject from the headers so a stalled 5xx response body + // cannot prevent the retry loop (or eventual abortUpload) from running. + const status = xhr.status + rejectOnce( + new MultipartUploadError( + `Video service returned ${status}`, + undefined, + status, + ), + ) + xhr.abort() + } + } xhr.open( 'POST', createVideoEndpointUrl('/xrpc/app.bsky.video.uploadPart', { diff --git a/src/lib/media/video/multipart/uploadParts.test.ts b/src/lib/media/video/multipart/uploadParts.test.ts index 7965a28df9..b580f14de2 100644 --- a/src/lib/media/video/multipart/uploadParts.test.ts +++ b/src/lib/media/video/multipart/uploadParts.test.ts @@ -115,6 +115,37 @@ describe('uploadParts', () => { expect(attempts).toBe(2) }) + it('retries service-unavailable parts', async () => { + let attempts = 0 + const uploadPart: UploadPartFn = ({part}) => { + attempts++ + if (attempts === 1) { + return Promise.reject( + new MultipartUploadError( + 'failed to upload multipart part', + 'ServiceUnavailable', + 503, + ), + ) + } + return Promise.resolve({ + partNumber: part.partNumber, + sizeBytes: part.size, + }) + } + + await uploadParts({ + parts: parts.slice(0, 1), + reader: fakeReader(), + uploadPart, + totalBytes: 10, + setProgress: () => {}, + signal: new AbortController().signal, + }) + + expect(attempts).toBe(2) + }) + it('does not retry a non-retryable response', async () => { const uploadPart = jest.fn< ReturnType, diff --git a/src/lib/media/video/multipart/uploadParts.ts b/src/lib/media/video/multipart/uploadParts.ts index cda502fb54..041ecff1b9 100644 --- a/src/lib/media/video/multipart/uploadParts.ts +++ b/src/lib/media/video/multipart/uploadParts.ts @@ -7,7 +7,7 @@ import { type PartUploadResult, type UploadPartFn, } from './types' -import {delay, isRetryableMultipartError} from './utils' +import {delay, isRetryableMultipartError, retryDelayMs} from './utils' /** * Uploads every part with a concurrency cap and per-part retry, aggregating @@ -121,7 +121,10 @@ async function uploadPartWithRetry({ lastError = err if (!isRetryableMultipartError(err)) throw err if (attempt < maxAttempts) { - await delay(500 * 2 ** (attempt - 1), signal) + // XHR progress starts over on a retry, so remove bytes reported by the + // failed attempt from the aggregate while backing off. + onProgress(0) + await delay(retryDelayMs(attempt), signal) } } } diff --git a/src/lib/media/video/multipart/utils.ts b/src/lib/media/video/multipart/utils.ts index 4184ab5666..05c42858f9 100644 --- a/src/lib/media/video/multipart/utils.ts +++ b/src/lib/media/video/multipart/utils.ts @@ -25,3 +25,9 @@ export function delay(ms: number, signal: AbortSignal) { signal.addEventListener('abort', onAbort, {once: true}) }) } + +/** Exponential backoff with 50-100% jitter to avoid synchronized retries. */ +export function retryDelayMs(attempt: number) { + const ceiling = Math.min(500 * 2 ** (attempt - 1), 8_000) + return ceiling * (0.5 + Math.random() * 0.5) +}