diff --git a/src/analytics/metrics/types.ts b/src/analytics/metrics/types.ts index ef875ff5bf..dc80eadeb9 100644 --- a/src/analytics/metrics/types.ts +++ b/src/analytics/metrics/types.ts @@ -5,7 +5,10 @@ import {type Platform} from 'react-native' import {type NotificationReason} from '#/lib/hooks/useNotificationHandler' -import {type VideoCompressSkipReason} from '#/lib/media/video/types' +import { + type VideoCompressSkipReason, + type VideoUploadTransport, +} from '#/lib/media/video/types' import {type NotificationType} from '#/state/queries/notifications/types' import {type FeedDescriptor} from '#/state/queries/post-feed' import {type LiveEventFeedMetricContext} from '#/features/liveEvents/types' @@ -1420,6 +1423,7 @@ export type Events = { bytes: number elapsedMs: number throughputBytesPerSec: number + transport: VideoUploadTransport } 'video:upload:uploadFailed': { uploadId: string @@ -1427,6 +1431,7 @@ export type Events = { bytes: number errorClass: string elapsedMs: number + transport: VideoUploadTransport } 'video:upload:processingStarted': { uploadId: string diff --git a/src/lib/media/video/multipart/upload.ts b/src/lib/media/video/multipart/upload.ts index c15ad2003d..1605bf40c9 100644 --- a/src/lib/media/video/multipart/upload.ts +++ b/src/lib/media/video/multipart/upload.ts @@ -26,14 +26,17 @@ export async function uploadVideoMultipart({ agent, setProgress, signal, + onStarted, }: { video: CompressedVideo agent: AtpAgent setProgress: (progress: number) => void signal: AbortSignal + onStarted?: () => void }): Promise { throwIfAborted(signal) - let token = await mintToken(agent) + const tokenProvider = createTokenProvider(agent) + const token = await tokenProvider.get() const name = `${nanoid(12)}.${mimeToExt(video.mimeType)}` let session try { @@ -46,10 +49,14 @@ export async function uploadVideoMultipart({ err instanceof Error ? err.message : 'Multipart upload unavailable', ) } + onStarted?.() const {jobId} = session const abortOnCancel = () => { - void abortUpload(jobId, token).catch(() => {}) + void tokenProvider + .get() + .then(currentToken => abortUpload(jobId, currentToken)) + .catch(() => {}) } signal.addEventListener('abort', abortOnCancel, {once: true}) let reader: ReturnType | undefined @@ -64,24 +71,28 @@ export async function uploadVideoMultipart({ await uploadParts({ parts, reader, - uploadPart: createUploadPart(jobId, token), + uploadPart: createUploadPart(jobId, tokenProvider.get), totalBytes: video.size, setProgress, signal, }) } catch (err) { if (signal.aborted) throw new AbortError() - return await abortThenFallbackOrResolve(jobId, token, err) + return await abortThenFallbackOrResolve( + jobId, + await tokenProvider.get(), + err, + ) } // Finish stores this credential for the later PDS blob upload, so use a // fresh token rather than the one that may have aged during transfer. - token = await mintToken(agent) + await tokenProvider.get(true) const activeReader = reader if (!activeReader) throw new Error('Video chunk reader is unavailable') return await finishAndRecover({ jobId, - token, + getToken: tokenProvider.get, signal, resendMissingParts: async receivedPartNumbers => { const missing = getMissingParts(parts, receivedPartNumbers) @@ -91,7 +102,7 @@ export async function uploadVideoMultipart({ await uploadParts({ parts: missing, reader: activeReader, - uploadPart: createUploadPart(jobId, token), + uploadPart: createUploadPart(jobId, tokenProvider.get), totalBytes: missingBytes, setProgress: progress => setProgress( @@ -110,18 +121,19 @@ export async function uploadVideoMultipart({ async function finishAndRecover({ jobId, - token, + getToken, signal, resendMissingParts, }: { jobId: string - token: string + getToken: () => Promise signal: AbortSignal resendMissingParts: (receivedPartNumbers: number[]) => Promise }): Promise { let createdFailures = 0 while (true) { throwIfAborted(signal) + const token = await getToken() try { const result = await finishUpload(jobId, token, signal) return result.jobStatus @@ -221,12 +233,33 @@ async function abortThenFallbackOrResolve( ) } -function mintToken(agent: AtpAgent) { - return getServiceAuthToken({ - agent, - lxm: 'com.atproto.repo.uploadBlob', - exp: Date.now() / 1000 + 60 * 30, - }) +function createTokenProvider(agent: AtpAgent) { + let token: string | undefined + let expiresAt = 0 + let refresh: Promise | undefined + + async function get(forceRefresh = false) { + if (!forceRefresh && token && Date.now() < expiresAt - 60_000) return token + if (!refresh) { + const exp = Math.floor(Date.now() / 1000) + 60 * 30 + refresh = getServiceAuthToken({ + agent, + lxm: 'com.atproto.repo.uploadBlob', + exp, + }) + .then(nextToken => { + token = nextToken + expiresAt = exp * 1000 + return nextToken + }) + .finally(() => { + refresh = undefined + }) + } + return refresh + } + + return {get} } function throwIfAborted(signal: AbortSignal) { diff --git a/src/lib/media/video/multipart/uploadPart.ts b/src/lib/media/video/multipart/uploadPart.ts index 2ac5c8df31..e3c57fc2b2 100644 --- a/src/lib/media/video/multipart/uploadPart.ts +++ b/src/lib/media/video/multipart/uploadPart.ts @@ -3,69 +3,92 @@ import {createVideoEndpointUrl} from '#/lib/media/video/util' import {MultipartUploadError} from './api' import {type UploadPartFn} from './types' -export function createUploadPart(jobId: string, token: string): UploadPartFn { - return ({part, chunk, onProgress, signal}) => - new Promise((resolve, reject) => { - if (signal.aborted) { - reject(new AbortError()) - return +export function createUploadPart( + jobId: string, + getToken: (forceRefresh?: boolean) => Promise, +): UploadPartFn { + return async args => { + try { + return await sendPart(jobId, await getToken(), args) + } catch (err) { + if ( + err instanceof MultipartUploadError && + (err.status === 401 || err.error === 'AuthRequired') + ) { + args.onProgress(0) + return await sendPart(jobId, await getToken(true), args) } - const xhr = new XMLHttpRequest() - const abort = () => xhr.abort() - signal.addEventListener('abort', abort, {once: true}) - const cleanup = () => signal.removeEventListener('abort', abort) - - xhr.upload.addEventListener('progress', event => { - onProgress(event.loaded) - }) - xhr.onerror = () => { - cleanup() - reject(new TypeError('Network request failed')) - } - xhr.onabort = () => { - cleanup() - reject(new AbortError()) - } - xhr.onload = () => { - cleanup() - let data: { - partNumber?: number - sizeBytes?: number - error?: string - message?: string - } - try { - data = JSON.parse(xhr.responseText) - } catch { - data = {} - } - if (xhr.status < 200 || xhr.status >= 300) { - reject( - new MultipartUploadError( - data.message || - data.error || - `Video service returned ${xhr.status}`, - data.error, - xhr.status, - ), - ) - } else { - onProgress(part.size) - resolve({ - partNumber: data.partNumber ?? part.partNumber, - sizeBytes: data.sizeBytes ?? part.size, - }) - } - } - xhr.open( - 'POST', - createVideoEndpointUrl('/xrpc/app.bsky.video.uploadPart', { - jobId, - partNumber: String(part.partNumber), - }), - ) - xhr.setRequestHeader('Content-Type', 'application/octet-stream') - xhr.setRequestHeader('Authorization', `Bearer ${token}`) - xhr.send(chunk as XMLHttpRequestBodyInit) - }) + throw err + } + } +} + +function sendPart( + jobId: string, + token: string, + {part, chunk, onProgress, signal}: Parameters[0], +) { + return new Promise>>((resolve, reject) => { + if (signal.aborted) { + reject(new AbortError()) + return + } + const xhr = new XMLHttpRequest() + const abort = () => xhr.abort() + signal.addEventListener('abort', abort, {once: true}) + const cleanup = () => signal.removeEventListener('abort', abort) + + xhr.upload.addEventListener('progress', event => { + onProgress(event.loaded) + }) + xhr.onerror = () => { + cleanup() + reject(new TypeError('Network request failed')) + } + xhr.onabort = () => { + cleanup() + reject(new AbortError()) + } + xhr.onload = () => { + cleanup() + let data: { + partNumber?: number + sizeBytes?: number + error?: string + message?: string + } + try { + data = JSON.parse(xhr.responseText) + } catch { + data = {} + } + if (xhr.status < 200 || xhr.status >= 300) { + reject( + new MultipartUploadError( + data.message || + data.error || + `Video service returned ${xhr.status}`, + data.error, + xhr.status, + ), + ) + } else { + onProgress(part.size) + resolve({ + partNumber: data.partNumber ?? part.partNumber, + sizeBytes: data.sizeBytes ?? part.size, + }) + } + } + xhr.open( + 'POST', + createVideoEndpointUrl('/xrpc/app.bsky.video.uploadPart', { + jobId, + partNumber: String(part.partNumber), + }), + ) + xhr.setRequestHeader('Content-Type', 'application/octet-stream') + xhr.setRequestHeader('Authorization', `Bearer ${token}`) + xhr.send(chunk as XMLHttpRequestBodyInit) + }) } diff --git a/src/lib/media/video/telemetry.ts b/src/lib/media/video/telemetry.ts index 67ff77c58e..9181ec002b 100644 --- a/src/lib/media/video/telemetry.ts +++ b/src/lib/media/video/telemetry.ts @@ -5,6 +5,7 @@ import {nanoid} from 'nanoid/non-secure' import { type ProbedMetadata, type VideoCompressSkipReason, + type VideoUploadTransport, } from '#/lib/media/video/types' import {Sentry} from '#/logger/sentry/lib' import {type Metrics} from '#/analytics/metrics' @@ -45,6 +46,7 @@ export type VideoTelemetry = { compressCompleted: (video: {size: number; mimeType: string}) => void compressFailed: (e: unknown) => void uploadStarted: (bytes: number) => void + uploadTransport: (transport: VideoUploadTransport) => void uploadCompleted: (jobId: string) => void uploadFailed: (e: unknown) => void processingStarted: (jobId: string) => void @@ -70,6 +72,7 @@ export function createVideoTelemetry({ let phaseStartedAt = startedAt let jobId: string | undefined let uploadBytes: number | undefined + let uploadTransport: VideoUploadTransport = 'legacy' let txnEnded = false let abortBound = true @@ -226,6 +229,11 @@ export function createVideoTelemetry({ metric('video:upload:uploadStarted', {uploadId, engine, bytes}) }, + uploadTransport(transport) { + uploadTransport = transport + phaseSpan?.setAttribute('video.upload.transport', transport) + }, + uploadCompleted(id) { jobId = id const elapsedMs = Date.now() - phaseStartedAt @@ -238,6 +246,7 @@ export function createVideoTelemetry({ elapsedMs, throughputBytesPerSec: elapsedMs > 0 ? Math.round((bytes * 1000) / elapsedMs) : 0, + transport: uploadTransport, }) endPhaseSpan() phase = undefined @@ -250,6 +259,7 @@ export function createVideoTelemetry({ bytes: uploadBytes ?? 0, errorClass: errorClass(e), elapsedMs: Date.now() - phaseStartedAt, + transport: uploadTransport, }) endTxn('error') detachAbort() diff --git a/src/lib/media/video/types.ts b/src/lib/media/video/types.ts index 6825f03ba3..1d2062ce00 100644 --- a/src/lib/media/video/types.ts +++ b/src/lib/media/video/types.ts @@ -8,6 +8,8 @@ export type VideoCompressSkipReason = | 'no-webcodecs' | 'compress-error-fallback' +export type VideoUploadTransport = 'multipart' | 'legacy' | 'legacy-fallback' + export type CompressedVideo = { uri: string mimeType: string diff --git a/src/lib/media/video/upload.ts b/src/lib/media/video/upload.ts index b045ab26b7..b91ad7a153 100644 --- a/src/lib/media/video/upload.ts +++ b/src/lib/media/video/upload.ts @@ -6,7 +6,10 @@ import {nanoid} from 'nanoid/non-secure' import {AbortError} from '#/lib/async/cancelable' import {ServerError} from '#/lib/media/video/errors' -import {type CompressedVideo} from '#/lib/media/video/types' +import { + type CompressedVideo, + type VideoUploadTransport, +} from '#/lib/media/video/types' import {Features, features} from '#/analytics/features' import {MultipartFallbackError, uploadVideoMultipart} from './multipart/upload' import {getServiceAuthToken, getVideoUploadLimits} from './upload.shared' @@ -19,6 +22,7 @@ export async function uploadVideo({ setProgress, signal, i18n, + onTransport, }: { video: CompressedVideo agent: AtpAgent @@ -26,6 +30,7 @@ export async function uploadVideo({ setProgress: (progress: number) => void signal: AbortSignal i18n: I18n + onTransport?: (transport: VideoUploadTransport) => void }) { if (signal.aborted) { throw new AbortError() @@ -34,11 +39,20 @@ export async function uploadVideo({ if (features.isOn(Features.VideoMultipartUploadEnable)) { try { - return await uploadVideoMultipart({video, agent, setProgress, signal}) + return await uploadVideoMultipart({ + video, + agent, + setProgress, + signal, + onStarted: () => onTransport?.('multipart'), + }) } catch (err) { if (!(err instanceof MultipartFallbackError)) throw err + onTransport?.('legacy-fallback') setProgress(0) } + } else { + onTransport?.('legacy') } const uri = createVideoEndpointUrl('/xrpc/app.bsky.video.uploadVideo', { diff --git a/src/lib/media/video/upload.web.ts b/src/lib/media/video/upload.web.ts index cd56a44ea7..cfefbd797d 100644 --- a/src/lib/media/video/upload.web.ts +++ b/src/lib/media/video/upload.web.ts @@ -5,7 +5,10 @@ import {nanoid} from 'nanoid/non-secure' import {AbortError} from '#/lib/async/cancelable' import {ServerError} from '#/lib/media/video/errors' -import {type CompressedVideo} from '#/lib/media/video/types' +import { + type CompressedVideo, + type VideoUploadTransport, +} from '#/lib/media/video/types' import {Features, features} from '#/analytics/features' import {MultipartFallbackError, uploadVideoMultipart} from './multipart/upload' import {getServiceAuthToken, getVideoUploadLimits} from './upload.shared' @@ -18,6 +21,7 @@ export async function uploadVideo({ setProgress, signal, i18n, + onTransport, }: { video: CompressedVideo agent: AtpAgent @@ -25,6 +29,7 @@ export async function uploadVideo({ setProgress: (progress: number) => void signal: AbortSignal i18n: I18n + onTransport?: (transport: VideoUploadTransport) => void }) { if (signal.aborted) { throw new AbortError() @@ -33,11 +38,20 @@ export async function uploadVideo({ if (features.isOn(Features.VideoMultipartUploadEnable)) { try { - return await uploadVideoMultipart({video, agent, setProgress, signal}) + return await uploadVideoMultipart({ + video, + agent, + setProgress, + signal, + onStarted: () => onTransport?.('multipart'), + }) } catch (err) { if (!(err instanceof MultipartFallbackError)) throw err + onTransport?.('legacy-fallback') setProgress(0) } + } else { + onTransport?.('legacy') } const uri = createVideoEndpointUrl('/xrpc/app.bsky.video.uploadVideo', { diff --git a/src/view/com/composer/state/video.ts b/src/view/com/composer/state/video.ts index 7b1ba77175..ed4384ccda 100644 --- a/src/view/com/composer/state/video.ts +++ b/src/view/com/composer/state/video.ts @@ -326,6 +326,7 @@ export async function processVideo( did, signal, i18n, + onTransport: telemetry.uploadTransport, setProgress: p => { dispatch({type: 'update_progress', progress: p, signal}) },