Bluesky multipart video upload sample
Something went wrong. Try again.
multipart-video-upload.ts
· Created 1mo ago ·13 kB · 380 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381/** * Post a video to Bluesky via the multipart upload endpoints on * https://video.bsky.app: startUpload → uploadPart × N (parallel) → * finishUpload → poll getJobStatus → post. * * npm i @bsky/sdk @atproto/lex @atproto/lex-password-session * lex install app.bsky.video.{startUpload,uploadPart,finishUpload,getUploadStatus,abortUpload,getJobStatus,defs} * lex build --clear --import-ext .ts * ATP_USERNAME=… ATP_PASSWORD=… node --experimental-strip-types post-video.ts clip.mp4 * * Uploads and stops there, printing the blob. Pass --post to actually post. */import { openAsBlob } from 'node:fs'import { stat } from 'node:fs/promises'import { basename } from 'node:path'import { fileURLToPath } from 'node:url'import { parseArgs } from 'node:util'
import { Client } from '@atproto/lex'import { PasswordSession } from '@atproto/lex-password-session'import { post } from '@bsky/sdk'import { com } from '@bsky/sdk/lexicons'
import * as videoLex from './src/lexicons/app/bsky/video.ts'
export const VIDEO_SERVICE = 'https://video.bsky.app'
export type JobStatus = videoLex.defs.JobStatusexport type PartPlan = { partNumber: number; start: number; end: number }export type UploadResult = { jobId: string; jobStatus: JobStatus }
/** Callable token getter, plus `.refresh()` to force a re-mint. */export type TokenSource = (() => Promise<string>) & { refresh: () => Promise<string> }
/** `failureCode` from getJobStatus. Surface these to the user. */export const FAILURE_REASONS: Record<string, string> = { validation_failure: 'this kind of video file is not supported', encoding_failure: 'the video could not be encoded', pds_upload_failure: 'the video could not be uploaded to your PDS', pds_upload_unsupported_blob_size: 'your PDS does not accept blobs this large', generic_failure: 'something else went wrong',}
export type UploadOptions = { sizeBytes: number mimeType: string name?: string slice: (part: PartPlan) => Blob | Promise<Blob> durationMs?: number width?: number height?: number /** * Parallel part uploads. 3 is a good default. */ concurrency?: number maxRetries?: number signal?: AbortSignal onProgress?: (done: number, total: number) => void /** Fires once startUpload returns — record the jobId to resume with. */ onStart?: (jobId: string, partCount: number) => void /** Called just before finishUpload; see `serviceTokenSource`. */ refreshToken?: () => Promise<unknown>}
/** * finishUpload writes the blob to your PDS with the token you present at that * time, so it must still be valid when a multi-minute upload finishes. * Mint one before startUpload, then `.refresh()` before finishUpload. * Requests in between reuse it and re-mint near expiry. */export function serviceTokenSource( client: Client, pdsHost: string, { lifetimeSec = 30 * 60, refreshBeforeSec = 120 } = {},): TokenSource { let token: string | undefined let expiresAt = 0
const mint = async () => { const now = Math.floor(Date.now() / 1000) const res = await client.call(com.atproto.server.getServiceAuth, { aud: `did:web:${pdsHost}`, lxm: 'com.atproto.repo.uploadBlob', exp: now + lifetimeSec, }) expiresAt = now + lifetimeSec return (token = res.token) }
const get = async () => { const now = Math.floor(Date.now() / 1000) if (token && now < expiresAt - refreshBeforeSec) return token return mint() }
return Object.assign(get, { refresh: mint })}
/** Resolves the token per request, so a refresh mid-upload is picked up. */export function createVideoClient( getToken: () => string | Promise<string>, service = VIDEO_SERVICE,): Client { return new Client(async (path, init) => { const headers = new Headers(init.headers) headers.set('authorization', `Bearer ${await getToken()}`) return fetch(new URL(path, service), { ...init, headers }) })}
/** * Parts are fixed-size except the last, and Content-Length must match exactly. * Derive ranges from the `partSizeBytes` the service returned. */export function planParts(sizeBytes: number, partSizeBytes: number): PartPlan[] { const parts: PartPlan[] = [] for (let start = 0; start < sizeBytes; start += partSizeBytes) { const end = Math.min(start + partSizeBytes, sizeBytes) parts.push({ partNumber: parts.length + 1, start, end }) } return parts}
/** Returns the job id to poll getJobStatus with. */export async function uploadVideoMultipart( video: Client, options: UploadOptions,): Promise<UploadResult> { const { sizeBytes, mimeType, name, durationMs, width, height, signal } = options
const session = await video.call( videoLex.startUpload, { sizeBytes, mimeType, name, durationMs, width, height }, { signal }, ) const { jobId } = session options.onStart?.(jobId, session.partCount)
try { const parts = planParts(sizeBytes, session.partSizeBytes) await sendParts(video, jobId, parts, options) // The token presented here is the one that writes to your PDS. await options.refreshToken?.() // finishUpload echoes the id back as `completedJobId`; take it from the // response rather than reusing ours, since that's the authoritative value. const res = await video.call(videoLex.finishUpload, { jobId }, { signal }) return { jobId: res.completedJobId, jobStatus: res.jobStatus } } catch (err) { // Frees the quota reservation now rather than at expiry. Open sessions count // against TooManyOpenUploads and the daily allowance. await video.call(videoLex.abortUpload, { jobId }).catch(() => {}) throw err }}
/** Parts are idempotent, so only the gaps get re-sent after a crash or drop. */export async function resumeVideoMultipart( video: Client, jobId: string, options: UploadOptions,): Promise<UploadResult> { const { signal } = options const status = await video.call(videoLex.getUploadStatus, { jobId }, { signal })
if (status.state === 'completed') { return { jobId: status.completedJobId!, jobStatus: status.jobStatus! } } if (status.state !== 'created') { throw new Error(`cannot resume upload in state ${status.state}`) }
const received = new Set(status.receivedParts) const missing = planParts(options.sizeBytes, status.partSizeBytes).filter( (part) => !received.has(part.partNumber), )
await sendParts(video, jobId, missing, options) await options.refreshToken?.() const res = await video.call(videoLex.finishUpload, { jobId }, { signal }) return { jobId: res.completedJobId, jobStatus: res.jobStatus }}
async function sendParts( video: Client, jobId: string, parts: PartPlan[], { slice, concurrency = 3, maxRetries = 4, signal, onProgress }: UploadOptions,): Promise<void> { let next = 0 let done = 0
const worker = async () => { for (;;) { const part = parts[next++] if (!part) return
// Takes a Blob/stream body, so the part streams from disk instead of being buffered whole. await video.xrpc(videoLex.uploadPart, { params: { jobId, partNumber: part.partNumber }, body: await slice(part), encoding: 'application/octet-stream', maxRetries, signal, })
onProgress?.(++done, parts.length) } }
await Promise.all( Array.from({ length: Math.min(concurrency, parts.length) }, worker), )}
/** * finishUpload succeeding means the bytes assembled. The probe runs after, so * poll here: `state` covers the pipeline CREATED → ENCODING → SCANNING → * UPLOADING → COMPLETED, and any unlisted value also means in-progress. * `progress` gives 0-100 within the current state. * * A bad file returns JOB_STATE_FAILED. Surface `failureCode` and see FAILURE_REASONS. */export async function waitForBlob( video: Client, jobId: string, { intervalMs = 1000, signal, onState, }: { intervalMs?: number signal?: AbortSignal onState?: (state: string, progress?: number) => void } = {},): Promise<JobStatus> { for (;;) { const { jobStatus } = await video.call(videoLex.getJobStatus, { jobId }, { signal }) onState?.(jobStatus.state, jobStatus.progress)
if (jobStatus.blob) return jobStatus if (jobStatus.state === 'JOB_STATE_FAILED') { const code = jobStatus.failureCode const reason = (code && FAILURE_REASONS[code]) ?? jobStatus.error ?? jobStatus.message throw new Error(`video processing failed (${code ?? 'unknown'}): ${reason}`) } await sleep(intervalMs, signal) }}
/** * Delay that gives up as soon as `signal` aborts. */function sleep(ms: number, signal?: AbortSignal): Promise<void> { return new Promise((resolve, reject) => { if (signal?.aborted) return reject(signal.reason) const onAbort = () => { clearTimeout(timer) reject(signal!.reason) } const timer = setTimeout(() => { signal?.removeEventListener('abort', onAbort) resolve() }, ms) signal?.addEventListener('abort', onAbort, { once: true }) })}
const USAGE = `usage: post-video <file> [options]
Uploads the video and stops, printing the resulting blob. Nothing is postedunless you pass --post.
--post create the post too (off by default) --text <str> post text, with --post --mime <type> default video/mp4 --concurrency <n> parallel part uploads (default 3) --width, --height advisory dimensions; also sets the embed aspect ratio --resume <jobId> finish an interrupted upload, re-sending only the gaps
env: ATP_USERNAME, ATP_PASSWORD`
async function main() { const { values, positionals } = parseArgs({ allowPositionals: true, options: { post: { type: 'boolean', default: false }, text: { type: 'string', default: '' }, mime: { type: 'string', default: 'video/mp4' }, concurrency: { type: 'string', default: '3' }, width: { type: 'string' }, height: { type: 'string' }, resume: { type: 'string' }, help: { type: 'boolean', short: 'h' }, }, })
const [videoPath] = positionals if (values.help || !videoPath) { console.log(USAGE) process.exit(values.help ? 0 : 1) }
const concurrency = Number(values.concurrency) if (!Number.isSafeInteger(concurrency) || concurrency < 1) { throw new Error(`--concurrency must be a positive integer, got "${values.concurrency}"`) }
const { ATP_USERNAME, ATP_PASSWORD } = process.env if (!ATP_USERNAME || !ATP_PASSWORD) { throw new Error('set ATP_USERNAME and ATP_PASSWORD') }
const session = await PasswordSession.login({ service: 'https://bsky.social', identifier: ATP_USERNAME, password: ATP_PASSWORD, }) const client = new Client(session)
const services: any[] = (session.session.didDoc as any)?.service ?? [] const pds = services.find((s) => s.id === '#atproto_pds')?.serviceEndpoint if (!pds) throw new Error('could not determine PDS endpoint') // Mint before startUpload so a bad token fails before any bytes move. const getToken = serviceTokenSource(client, new URL(pds).host) await getToken() const video = createVideoClient(getToken)
// `openAsBlob` is lazy. Slicing yields a part-sized Blob with a known length, // so each part gets an exact Content-Length and memory stays bounded by // concurrency × partSize. A browser `File` from an <input> slices the same. const file = await openAsBlob(videoPath) const { size: sizeBytes } = await stat(videoPath) const width = values.width ? Number(values.width) : undefined const height = values.height ? Number(values.height) : undefined
const options: UploadOptions = { sizeBytes, mimeType: values.mime, name: basename(videoPath), width, height, concurrency, slice: ({ start, end }) => file.slice(start, end), // Printed before any bytes move, so it's on screen to --resume with if this upload dies. onStart: (jobId, partCount) => console.log(`started upload ${jobId} (${partCount} parts)`), onProgress: (done, total) => console.log(`uploaded ${done}/${total} parts`), refreshToken: getToken.refresh, }
const { jobId } = values.resume ? await resumeVideoMultipart(video, values.resume, options) : await uploadVideoMultipart(video, options)
console.log(`finished upload ${jobId} — now polling getJobStatus`)
const jobStatus = await waitForBlob(video, jobId, { onState: (state, progress) => console.log(`${state}${progress === undefined ? '' : ` (${progress}%)`}`), }) console.log('blob:', JSON.stringify(jobStatus.blob))
if (!values.post) { console.log('\nupload complete — re-run with --post to attach it to a post') return }
const { uri } = await client.call(post, { text: values.text, langs: ['en'], embed: { $type: 'app.bsky.embed.video', video: jobStatus.blob!, ...(width && height ? { aspectRatio: { width, height } } : {}), }, }) console.log('posted:', uri)}
if (process.argv[1] === fileURLToPath(import.meta.url)) await main()