APP-2793: harden multipart video uploads
This commit is contained in:
@@ -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<AbortUploadResponse>({
|
return request<AbortUploadResponse>({
|
||||||
route: '/xrpc/app.bsky.video.abortUpload',
|
route: '/xrpc/app.bsky.video.abortUpload',
|
||||||
token,
|
token,
|
||||||
|
signal,
|
||||||
body: {jobId},
|
body: {jobId},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,10 +4,19 @@
|
|||||||
*/
|
*/
|
||||||
|
|
||||||
/** Max parts uploaded concurrently. */
|
/** 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. */
|
/** 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. */
|
/** Attempts to begin/continue server-side finalization before checking state. */
|
||||||
export const MULTIPART_FINISH_ATTEMPTS = 3
|
export const MULTIPART_FINISH_ATTEMPTS = 3
|
||||||
|
|||||||
@@ -14,12 +14,16 @@ import {
|
|||||||
MultipartUploadError,
|
MultipartUploadError,
|
||||||
startUpload,
|
startUpload,
|
||||||
} from './api'
|
} 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 {getMissingParts, planParts} from './planParts'
|
||||||
import {createChunkReader} from './readChunk'
|
import {createChunkReader} from './readChunk'
|
||||||
import {createUploadPart} from './uploadPart'
|
import {createUploadPart} from './uploadPart'
|
||||||
import {uploadParts} from './uploadParts'
|
import {uploadParts} from './uploadParts'
|
||||||
import {delay, isRetryableMultipartError} from './utils'
|
import {delay, isRetryableMultipartError, retryDelayMs} from './utils'
|
||||||
|
|
||||||
export class MultipartFallbackError extends Error {}
|
export class MultipartFallbackError extends Error {}
|
||||||
|
|
||||||
@@ -221,7 +225,7 @@ async function abortThenFallbackOrResolve(
|
|||||||
token: string,
|
token: string,
|
||||||
cause: unknown,
|
cause: unknown,
|
||||||
): Promise<AppBskyVideoDefs.JobStatus> {
|
): Promise<AppBskyVideoDefs.JobStatus> {
|
||||||
const result = await abortUpload(jobId, token)
|
const result = await abortUploadWithRetry(jobId, token)
|
||||||
if (result.state === 'aborted') {
|
if (result.state === 'aborted') {
|
||||||
throw new MultipartFallbackError(
|
throw new MultipartFallbackError(
|
||||||
cause instanceof Error ? cause.message : 'Multipart upload failed',
|
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) {
|
function createTokenProvider(agent: AtpAgent, signal: AbortSignal) {
|
||||||
let token: string | undefined
|
let token: string | undefined
|
||||||
let expiresAt = 0
|
let expiresAt = 0
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import {AbortError} from '#/lib/async/cancelable'
|
import {AbortError} from '#/lib/async/cancelable'
|
||||||
import {createVideoEndpointUrl} from '#/lib/media/video/util'
|
import {createVideoEndpointUrl} from '#/lib/media/video/util'
|
||||||
import {MultipartUploadError} from './api'
|
import {MultipartUploadError} from './api'
|
||||||
|
import {MULTIPART_PART_TIMEOUT_MS} from './constants'
|
||||||
import {type UploadPartFn} from './types'
|
import {type UploadPartFn} from './types'
|
||||||
|
|
||||||
export function createUploadPart(
|
export function createUploadPart(
|
||||||
@@ -34,23 +35,40 @@ function sendPart(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
const xhr = new XMLHttpRequest()
|
const xhr = new XMLHttpRequest()
|
||||||
|
xhr.timeout = MULTIPART_PART_TIMEOUT_MS
|
||||||
const abort = () => xhr.abort()
|
const abort = () => xhr.abort()
|
||||||
signal.addEventListener('abort', abort, {once: true})
|
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<ReturnType<UploadPartFn>>) => {
|
||||||
|
if (settled) return
|
||||||
|
settled = true
|
||||||
|
cleanup()
|
||||||
|
resolve(result)
|
||||||
|
}
|
||||||
|
|
||||||
xhr.upload.addEventListener('progress', event => {
|
xhr.upload.addEventListener('progress', event => {
|
||||||
onProgress(event.loaded)
|
onProgress(event.loaded)
|
||||||
})
|
})
|
||||||
xhr.onerror = () => {
|
xhr.onerror = () => {
|
||||||
cleanup()
|
rejectOnce(new TypeError('Network request failed'))
|
||||||
reject(new TypeError('Network request failed'))
|
}
|
||||||
|
xhr.ontimeout = () => {
|
||||||
|
rejectOnce(new TypeError('Multipart part upload timed out'))
|
||||||
}
|
}
|
||||||
xhr.onabort = () => {
|
xhr.onabort = () => {
|
||||||
cleanup()
|
rejectOnce(new AbortError())
|
||||||
reject(new AbortError())
|
|
||||||
}
|
}
|
||||||
xhr.onload = () => {
|
xhr.onload = () => {
|
||||||
cleanup()
|
|
||||||
let data: {
|
let data: {
|
||||||
partNumber?: number
|
partNumber?: number
|
||||||
sizeBytes?: number
|
sizeBytes?: number
|
||||||
@@ -63,7 +81,7 @@ function sendPart(
|
|||||||
data = {}
|
data = {}
|
||||||
}
|
}
|
||||||
if (xhr.status < 200 || xhr.status >= 300) {
|
if (xhr.status < 200 || xhr.status >= 300) {
|
||||||
reject(
|
rejectOnce(
|
||||||
new MultipartUploadError(
|
new MultipartUploadError(
|
||||||
data.message ||
|
data.message ||
|
||||||
data.error ||
|
data.error ||
|
||||||
@@ -74,12 +92,31 @@ function sendPart(
|
|||||||
)
|
)
|
||||||
} else {
|
} else {
|
||||||
onProgress(part.size)
|
onProgress(part.size)
|
||||||
resolve({
|
resolveOnce({
|
||||||
partNumber: data.partNumber ?? part.partNumber,
|
partNumber: data.partNumber ?? part.partNumber,
|
||||||
sizeBytes: data.sizeBytes ?? part.size,
|
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(
|
xhr.open(
|
||||||
'POST',
|
'POST',
|
||||||
createVideoEndpointUrl('/xrpc/app.bsky.video.uploadPart', {
|
createVideoEndpointUrl('/xrpc/app.bsky.video.uploadPart', {
|
||||||
|
|||||||
@@ -115,6 +115,37 @@ describe('uploadParts', () => {
|
|||||||
expect(attempts).toBe(2)
|
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 () => {
|
it('does not retry a non-retryable response', async () => {
|
||||||
const uploadPart = jest.fn<
|
const uploadPart = jest.fn<
|
||||||
ReturnType<UploadPartFn>,
|
ReturnType<UploadPartFn>,
|
||||||
|
|||||||
@@ -7,7 +7,7 @@ import {
|
|||||||
type PartUploadResult,
|
type PartUploadResult,
|
||||||
type UploadPartFn,
|
type UploadPartFn,
|
||||||
} from './types'
|
} 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
|
* Uploads every part with a concurrency cap and per-part retry, aggregating
|
||||||
@@ -121,7 +121,10 @@ async function uploadPartWithRetry({
|
|||||||
lastError = err
|
lastError = err
|
||||||
if (!isRetryableMultipartError(err)) throw err
|
if (!isRetryableMultipartError(err)) throw err
|
||||||
if (attempt < maxAttempts) {
|
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)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -25,3 +25,9 @@ export function delay(ms: number, signal: AbortSignal) {
|
|||||||
signal.addEventListener('abort', onAbort, {once: true})
|
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)
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user