APP-2670: refresh upload auth and track transport
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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<AppBskyVideoDefs.JobStatus> {
|
||||
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<typeof createChunkReader> | 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<string>
|
||||
signal: AbortSignal
|
||||
resendMissingParts: (receivedPartNumbers: number[]) => Promise<boolean>
|
||||
}): Promise<AppBskyVideoDefs.JobStatus> {
|
||||
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<string> | 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) {
|
||||
|
||||
@@ -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<string>,
|
||||
): 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<UploadPartFn>[0],
|
||||
) {
|
||||
return new Promise<Awaited<ReturnType<UploadPartFn>>>((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)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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', {
|
||||
|
||||
@@ -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', {
|
||||
|
||||
@@ -326,6 +326,7 @@ export async function processVideo(
|
||||
did,
|
||||
signal,
|
||||
i18n,
|
||||
onTransport: telemetry.uploadTransport,
|
||||
setProgress: p => {
|
||||
dispatch({type: 'update_progress', progress: p, signal})
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user