diff --git a/README.md b/README.md index 41f7e62..42bfbbb 100644 --- a/README.md +++ b/README.md @@ -70,9 +70,11 @@ upstash qstash stats --qstash-id $QSTASH_ID --period 7d # Blob upstash blob create --name my-bucket --visibility private -upstash blob list -upstash blob credentials --bucket-id $BUCKET_ID -upstash blob upload ./assets --bucket-id $BUCKET_ID --prefix assets +upstash blob ls +upstash blob ls my-bucket +upstash blob cp ./assets blob://my-bucket/assets -r +upstash blob sync ./site blob://my-bucket/site -d +upstash blob credentials my-bucket # Team upstash team list @@ -81,60 +83,59 @@ upstash team add-member --team-id $TEAM_ID --member-email you@example.com --role Run `upstash --help` (or `--help` on any subcommand) to discover everything else, and check the [full docs](https://upstash.com/docs/agent-resources/cli) for the complete catalog. `upstash blob credentials` returns temporary S3 credentials for use with AWS CLI, rclone, or an S3 SDK. -## Uploading Blob files and folders +## Working with Blob buckets and objects -Set `UPSTASH_BLOB_TOKEN` in your environment or `.env` file, then run: +The object commands mirror `aws s3`, with `blob:///` in place of +`s3://`. `` is a bucket name or id. ```bash -upstash blob upload ./assets --prefix assets +upstash blob ls # buckets +upstash blob ls my-bucket/images/ # one level; -r for all +upstash blob cp ./photo.png blob://my-bucket/images/ +upstash blob cp ./assets blob://my-bucket/assets -r +upstash blob cp blob://my-bucket/images ./images -r --exclude "*.tmp" +upstash blob cp blob://my-bucket/config.json - | jq . +upstash blob mv blob://my-bucket/a.txt blob://other-bucket/a.txt +upstash blob sync ./site blob://my-bucket/site -d +upstash blob rm my-bucket/tmp -r -n +upstash blob presign my-bucket/report.pdf --expires-in 3600 +upstash blob mb blob://new-bucket +upstash blob rb new-bucket -f ``` -No Upstash login, account email, or management API key is required when using a -bucket token. You can also provide the token explicitly or select another env file: +`cp`, `mv` and `sync` need `blob://` to tell bucket paths from local ones. The +commands that only take bucket paths (`ls`, `rm`, `presign`, `mb`, `rb`) accept +`my-bucket/key` without it, as do `get`, `delete` and `credentials`, which take a +bucket name or id. -```bash -upstash blob upload ./assets --token "$BLOB_TOKEN" --prefix assets -upstash --env-path ./uploads.env blob upload ./assets --prefix assets -upstash blob credentials --token "$BLOB_TOKEN" -``` +Flags follow `aws s3`: `-r/--recursive`, `--exclude`/`--include` (applied in +order, last match wins), `-n/--dryrun`, `-d/--delete`, `--size-only`, +`--exact-timestamps`, `--content-type`, `--cache-control`, `--metadata`, +`--expected-size`, `--concurrency` and `-q/--quiet`. Copies between buckets reset +Cache-Control to the default unless `--cache-control` is given. Local symbolic +links are followed. -`--token` overrides `UPSTASH_BLOB_TOKEN`. Exported environment variables take -precedence over values loaded from `.env` or `--env-path`. Use the Blob bucket -token, not temporary S3 credentials, so the CLI can refresh credentials throughout -the transfer. - -Alternatively, use `--bucket-id $BUCKET_ID` with your saved Upstash login or -Developer API credentials. An explicit bucket ID overrides the ambient token; -`--token` and `--bucket-id` cannot be combined. AWS CLI and manually exported S3 -credentials are not needed. A directory uploads its contents recursively: `./assets/images/logo.png` -becomes `assets/images/logo.png` with the prefix above, or `images/logo.png` without -a prefix. A single file uploads under its filename. Content types are inferred -from filenames, falling back to `application/octet-stream`. - -The Blob SDK streams files, uses multipart for large files, and refreshes temporary -credentials throughout the upload, including between parts of one large file. -Transient failures are retried. Four files upload concurrently by default; use -`--concurrency 1` to reduce memory usage. Progress goes to stderr and the final JSON -summary goes to stdout. `--quiet` suppresses progress. +Progress goes to stderr and a JSON summary to stdout. Transfers retry transient +failures, keep going past a failed file, and exit unsuccessfully at the end. Large +files use multipart uploads, and the Blob SDK refreshes temporary S3 credentials +throughout, even between parts of one file. + +### Using a bucket token instead of a login + +Bucket names need an Upstash login. A Blob bucket token (`--token`, or +`UPSTASH_BLOB_TOKEN` in the environment or `.env`) works without one, but only for +its own bucket, addressed by id. A token is never used for a bucket it wasn't +issued for. ```bash -upstash blob upload ./assets --prefix assets --dry-run -upstash blob upload ./assets --prefix assets --skip-existing +upstash blob cp ./assets blob://$BUCKET_ID/assets -r --token "$BLOB_TOKEN" +upstash --env-path ./uploads.env blob sync ./assets blob://$BUCKET_ID/assets +upstash blob credentials --token "$BLOB_TOKEN" ``` -`--dry-run` lists local files and destination paths without authenticating or -making network requests. By default existing keys are overwritten. `--skip-existing` -skips any existing key **without comparing size or contents**; use it to rerun an -interrupted upload only when the already uploaded objects are the versions you want. -An incomplete individual file starts again on rerun. Files are not deleted from the -bucket. Symlinks and empty directories are skipped. - -On a failed file, the command stops scheduling more files, waits for active uploads, -prints a summary with failed and remaining files, and exits unsuccessfully. Ctrl+C -stops scheduling work and closes local streams; in-flight requests may take time -to settle. Completed objects remain in the bucket. A process kill, or a network -failure that also blocks cleanup, may leave an incomplete multipart upload. It does -not expire on its own; remove it with the Blob SDK's `abortStaleMultipartUploads`. +`--token` and `UPSTASH_BLOB_TOKEN` are each used only for their own bucket, so they +can point at different buckets. Exported environment variables take precedence +over values loaded from `.env` or `--env-path`. ## Telemetry diff --git a/src/commands/blob/buckets.ts b/src/commands/blob/buckets.ts new file mode 100644 index 0000000..08e18ce --- /dev/null +++ b/src/commands/blob/buckets.ts @@ -0,0 +1,206 @@ +import { Bucket } from "@upstash/blob"; +import type { Command } from "commander"; +import { createHash, createHmac } from "node:crypto"; +import { resolveAuth } from "../../auth.js"; +import type { Auth } from "../../auth.js"; +import { request } from "../../client.js"; +import { telemetryStatus } from "../../telemetry.js"; +import type { BlobBucket } from "../../types.js"; +import { fetchBlobCredentials } from "./credentials.js"; +import { isFreshlyCreated, PROVISIONING_MAX_RETRIES, sleep } from "./retry.js"; +import { parseBucket } from "./transfer.js"; + +/** The bucket id a Blob token was issued for, read from the token itself. */ +export function tokenBucketId(token: string): string | undefined { + const raw = Buffer.from(token.trim(), "base64url"); + if (raw.length < 6) return undefined; + const id = raw.subarray(6, 6 + (raw[2] ?? 0)).toString(); + return id || undefined; +} + +export function findAccountBucket(buckets: BlobBucket[], name: string): BlobBucket { + const match = buckets.find((bucket) => bucket.id === name) ?? buckets.find((bucket) => bucket.name === name); + if (!match) throw new Error(`Blob bucket "${name}" not found`); + return match; +} + +/** + * The id behind the bucket argument of get, delete and credentials. The hidden --bucket-id flag + * they used to require is still accepted, and used as is. + */ +export async function bucketIdArgument(command: Command, bucket: string | undefined, flags: { bucketId?: string }): Promise { + if (bucket !== undefined && flags.bucketId !== undefined) throw new Error("Name the bucket once, not also with --bucket-id"); + if (flags.bucketId !== undefined) return flags.bucketId; + if (bucket === undefined) throw new Error(`Name a bucket: upstash blob ${command.name()} `); + const name = parseBucket(bucket); + return findAccountBucket(await request(resolveAuth(command), "GET", "/v2/blob/bucket"), name).id; +} + +/** + * Turns the bucket part of a blob:// URI into a Bucket. A Blob token is used only for the bucket it + * was issued for, matched by id; anything else is looked up by name or id with account credentials, + * so a stray UPSTASH_BLOB_TOKEN in .env never redirects a command to another bucket. + */ +export class BucketResolver { + private readonly buckets = new Map>(); + private accountBuckets?: Promise; + + constructor(private readonly command: Command, private readonly token?: string) {} + + open(name: string): Promise { + let bucket = this.buckets.get(name); + if (!bucket) { + bucket = this.resolve(name); + this.buckets.set(name, bucket); + } + return bucket; + } + + auth(): Auth { + return resolveAuth(this.command); + } + + listAccountBuckets(): Promise { + this.accountBuckets ??= request(this.auth(), "GET", "/v2/blob/bucket"); + return this.accountBuckets; + } + + private async resolve(name: string): Promise { + if (this.token !== undefined && !this.token.trim()) throw new Error("--token must be a non-empty Blob bucket token"); + const tokens = [this.token?.trim(), process.env.UPSTASH_BLOB_TOKEN?.trim()] + .filter((token): token is string => typeof token === "string" && token.length > 0); + const direct = tokens.find((token) => tokenBucketId(token) === name); + if (direct) return this.bucket(direct); + + let auth: Auth; + try { + auth = this.auth(); + } catch (error) { + if (tokens.length === 0) throw error; + const ids = [...new Set(tokens.map((token) => tokenBucketId(token) ?? "unknown"))]; + throw new Error( + `The Blob token is for bucket ${ids.join(", ")}, not "${name}". Address that bucket as blob://${ids[0]}/..., or run \`upstash login\` to use bucket names`, + ); + } + const match = findAccountBucket(await this.listAccountBuckets(), name); + const bucket = await request(auth, "GET", `/v2/blob/bucket/${match.id}`); + if (typeof bucket.token !== "string" || bucket.token.length === 0) { + throw new Error(`Blob bucket ${match.id} did not return a current token`); + } + // A bucket created moments ago answers 401 until provisioning finishes. + if (isFreshlyCreated(bucket.creation_time)) { + await fetchBlobCredentials(bucket.token, sleep, { unauthorizedRetries: PROVISIONING_MAX_RETRIES }); + } + return this.bucket(bucket.token); + } + + private bucket(token: string): Bucket { + return new Bucket({ token, enableTelemetry: telemetryStatus().enabled }); + } +} + +export interface ListedObject { + key: string; + size: number; + last_modified: string; + etag: string; +} + +export interface DirectoryPage { + prefixes: string[]; + objects: ListedObject[]; + cursor?: string; +} + +const sha256 = (data: string): string => createHash("sha256").update(data).digest("hex"); +const hmac = (key: string | Buffer, data: string): Buffer => createHmac("sha256", key).update(data).digest(); +const uriEncode = (value: string): string => + encodeURIComponent(value).replace(/[!'()*]/g, (char) => `%${char.charCodeAt(0).toString(16).toUpperCase()}`); + +function decodeEntities(value: string): string { + return value.replace(/</g, "<").replace(/>/g, ">").replace(/"/g, '"') + .replace(/'|'/g, "'").replace(/&/g, "&"); +} + +function tag(xml: string, name: string): string | undefined { + const match = new RegExp(`<${name}>([\\s\\S]*?)`).exec(xml); + return match?.[1] === undefined ? undefined : decodeEntities(match[1]); +} + +function blocks(xml: string, name: string): string[] { + return [...xml.matchAll(new RegExp(`<${name}>([\\s\\S]*?)`, "g"))].map((match) => match[1] ?? ""); +} + +/** + * One page of a delimiter listing: the keys directly under `prefix` and the "folders" below it. + * The SDK's list() has no delimiter, so this signs its own ListObjectsV2 with the bucket's + * refreshing S3 credentials. + */ +export async function listDirectory(bucket: Bucket, prefix: string, cursor?: string): Promise { + const s3 = bucket.s3(); + for (let attempt = 0; ; attempt++) { + const [{ url: endpoint }, credentials] = await Promise.all([s3.endpoint(), s3.credentials()]); + const query: Record = { delimiter: "/", "list-type": "2" }; + if (prefix) query.prefix = prefix; + if (cursor) query["continuation-token"] = cursor; + const canonicalQuery = Object.keys(query).sort() + .map((key) => `${uriEncode(key)}=${uriEncode(query[key] ?? "")}`).join("&"); + const path = `${endpoint.pathname.replace(/\/+$/, "")}/${uriEncode(s3.bucket)}`; + const amzDate = new Date().toISOString().replace(/[-:]|\.\d{3}/g, ""); + const date = amzDate.slice(0, 8); + const payloadHash = sha256(""); + const headers: Record = { + host: endpoint.host, + "x-amz-content-sha256": payloadHash, + "x-amz-date": amzDate, + "x-amz-security-token": credentials.sessionToken, + }; + const names = Object.keys(headers).sort(); + const canonicalRequest = [ + "GET", path, canonicalQuery, names.map((name) => `${name}:${headers[name]}\n`).join(""), names.join(";"), payloadHash, + ].join("\n"); + const scope = `${date}/${s3.region}/s3/aws4_request`; + let key = hmac(`AWS4${credentials.secretAccessKey}`, date); + for (const part of [s3.region, "s3", "aws4_request"]) key = hmac(key, part); + const signature = createHmac("sha256", key) + .update(["AWS4-HMAC-SHA256", amzDate, scope, sha256(canonicalRequest)].join("\n")).digest("hex"); + const { host: _host, ...sent } = headers; + + let response: Response; + try { + response = await fetch(`${endpoint.origin}${path}?${canonicalQuery}`, { + headers: { + ...sent, + authorization: `AWS4-HMAC-SHA256 Credential=${credentials.accessKeyId}/${scope}, SignedHeaders=${names.join(";")}, Signature=${signature}`, + }, + }); + } catch (error) { + if (attempt >= 2) throw error; + await sleep(500 * 2 ** attempt); + continue; + } + const xml = await response.text(); + if ((response.status === 429 || response.status >= 500) && attempt < 2) { + await sleep(500 * 2 ** attempt); + continue; + } + if (!response.ok) { + throw new Error(`listing failed: ${tag(xml, "Message") ?? tag(xml, "Code") ?? `HTTP ${response.status}`}`); + } + const next = tag(xml, "NextContinuationToken"); + return { + prefixes: blocks(xml, "CommonPrefixes").map((block) => tag(block, "Prefix")).filter((value): value is string => value !== undefined), + objects: blocks(xml, "Contents").flatMap((block) => { + const objectKey = tag(block, "Key"); + if (objectKey === undefined) return []; + return [{ + key: objectKey, + size: Number(tag(block, "Size") ?? 0), + last_modified: new Date(tag(block, "LastModified") ?? 0).toISOString(), + etag: tag(block, "ETag") ?? "", + }]; + }), + cursor: tag(xml, "IsTruncated") === "true" && next ? next : undefined, + }; + } +} diff --git a/src/commands/blob/cp.ts b/src/commands/blob/cp.ts new file mode 100644 index 0000000..d0b2944 --- /dev/null +++ b/src/commands/blob/cp.ts @@ -0,0 +1,173 @@ +import { Command, InvalidArgumentError } from "commander"; +import { Readable } from "node:stream"; +import { pipeline } from "node:stream/promises"; +import type { ReadableStream as NodeReadableStream } from "node:stream/web"; +import { printJSON } from "../../output.js"; +import { BucketResolver } from "./buckets.js"; +import { + addCommonOptions, + addObjectOptions, + checkOverwritesSource, + claimLocal, + foldsCase, + destinationFor, + executePlan, + filtersOf, + formatLocation, + isDryRun, + isIncluded, + isLocalDirectory, + listSource, + nestedPrefix, + parseLocation, + transferAction, +} from "./transfer.js"; +import type { CommonOptions, Failure, Location, ObjectOptions, Operation } from "./transfer.js"; + +interface CopyOptions extends CommonOptions, ObjectOptions { + recursive?: boolean; + expectedSize?: number; +} + +function byteCount(value: string): number { + const number = Number(value); + if (!Number.isSafeInteger(number) || number < 0) throw new InvalidArgumentError("must be a whole number of bytes"); + return number; +} + +/** `cp - blob://...` uploads stdin and `cp blob://... -` writes the object to stdout, as aws s3 cp does. */ +async function copyStream(source: Location, destination: Location, options: CopyOptions, resolver: BucketResolver): Promise { + if (options.recursive) throw new Error("- cannot be combined with --recursive"); + if (source.type === "local" && destination.type === "blob") { + if (!destination.key || destination.key.endsWith("/")) { + throw new Error("uploading stdin needs a full blob:/// destination"); + } + if (isDryRun(options)) { + printJSON({ dry_run: true, operations: [{ action: "upload", source: "-", destination: formatLocation(destination) }] }); + return; + } + const bucket = await resolver.open(destination.bucket); + const body = Readable.toWeb(process.stdin) as ReadableStream; + // Failing the stream makes the SDK abort a multipart upload instead of leaving it behind. + const interrupt = (): void => { + process.stdin.destroy(new Error("upload interrupted")); + }; + process.once("SIGINT", interrupt); + process.once("SIGTERM", interrupt); + let result; + try { + // Without a declared size the SDK has to buffer the stream to learn its length. + result = await bucket.put(destination.key, body, { + contentType: options.contentType ?? "application/octet-stream", + cache: options.cacheControl, + metadata: options.metadata, + // A declared size of 0 would make the SDK drop the stream, so it is buffered like no size. + ...(options.expectedSize ? { size: options.expectedSize } : { maxSize: "5gb" }), + }); + } finally { + process.removeListener("SIGINT", interrupt); + process.removeListener("SIGTERM", interrupt); + } + if (!options.quiet) console.error(`upload: - to ${formatLocation(destination)}`); + printJSON({ completed: 1, bytes: result.size, failed: [], remaining: 0 }); + return; + } + if (source.type === "blob" && destination.type === "local") { + if (!source.key || source.key.endsWith("/")) throw new Error("writing to stdout needs a full blob:/// source"); + if (isDryRun(options)) { + printJSON({ dry_run: true, operations: [{ action: "download", source: formatLocation(source), destination: "-" }] }); + return; + } + const blob = await (await resolver.open(source.bucket)).get(source.key); + await pipeline(Readable.fromWeb(blob.body as NodeReadableStream), process.stdout); + return; + } + throw new Error("- streams between stdin or stdout and a blob:// object"); +} + +function registerTransfer(blob: Command, name: "cp" | "mv"): void { + const command = blob + .command(`${name} `) + .description(name === "cp" + ? "Copy a file or object, or a directory or prefix with --recursive, like aws s3 cp" + : "Move a file or object, or a directory or prefix with --recursive, like aws s3 mv") + .option("-r, --recursive", "Copy everything below a local directory or blob:// prefix"); + addObjectOptions(command); + addCommonOptions(command); + if (name === "cp") { + command.option("--expected-size ", "Size of the stdin stream for `cp - blob://...`; without it stdin is buffered in memory", byteCount); + } + command + .addHelpText("after", ` +Locations are local paths or blob:///, where is a bucket +name or id. At least one side must be blob://. A destination ending in "/" (or an +existing local directory) keeps the source's file name. ${name === "cp" + ? `Use - as the source or +destination to stream stdin or stdout.` + : "Sources, local files included,\nare deleted after each successful copy."} +Copies between buckets keep content type and metadata; Cache-Control is reset +to the default unless --cache-control is given. + +Examples: + upstash blob ${name} ./photo.png blob://my-bucket/images/ + upstash blob ${name} ./site blob://my-bucket/site -r + upstash blob ${name} blob://my-bucket/images ./images -r --exclude "*.tmp" + upstash blob ${name} blob://my-bucket/a.txt blob://other-bucket/b.txt +`) + .action(async (sourceArg: string, destinationArg: string, options: CopyOptions, cmd: Command) => { + const source = parseLocation(sourceArg); + const destination = parseLocation(destinationArg); + if (source.type === "local" && destination.type === "local") { + throw new Error("source or destination must be a blob:/// URI"); + } + const resolver = new BucketResolver(cmd, options.token); + if (sourceArg === "-" || destinationArg === "-") { + if (name === "mv") throw new Error("mv cannot stream stdin or stdout; use cp"); + await copyStream(source, destination, options, resolver); + return; + } + + const recursive = Boolean(options.recursive); + const filters = filtersOf(options); + const [skip, sourceInside] = recursive + ? await Promise.all([nestedPrefix(source, destination, resolver), nestedPrefix(destination, source, resolver)]) + : []; + const entries = (await listSource(source, recursive, resolver)).filter((entry) => + isIncluded(entry.rel, filters) && !(skip !== undefined && entry.location.type === "blob" && entry.location.key.startsWith(skip))); + const into = recursive || (destination.type === "blob" + ? !destination.key || destination.key.endsWith("/") + : await isLocalDirectory(destination.path)); + const operations: Operation[] = []; + const failures: Failure[] = []; + const claimed = new Set(); + for (const entry of entries) { + // Zero-byte "folder" markers have no local file to become. + if (destination.type === "local" && entry.rel.endsWith("/")) continue; + try { + const target = destinationFor(entry, destination, into); + // An mv folds case everywhere: on a case-insensitive Linux mount, two keys writing one + // file would otherwise both be deleted. + await claimLocal(claimed, target, foldsCase || name === "mv"); + checkOverwritesSource(target, sourceInside); + operations.push({ + action: transferAction(source, destination), + source: entry.location, + destination: target, + size: entry.size, + move: name === "mv", + }); + } catch (error) { + failures.push({ source: formatLocation(entry.location), error: (error as Error).message }); + } + } + await executePlan(operations, failures, options, resolver); + }); +} + +export function registerBlobCp(blob: Command): void { + registerTransfer(blob, "cp"); +} + +export function registerBlobMv(blob: Command): void { + registerTransfer(blob, "mv"); +} diff --git a/src/commands/blob/create.ts b/src/commands/blob/create.ts index 8e36afb..dd87b53 100644 --- a/src/commands/blob/create.ts +++ b/src/commands/blob/create.ts @@ -5,7 +5,7 @@ import { printJSON } from "../../output.js"; import { BLOB_VISIBILITIES } from "../../types.js"; import type { BlobBucket, BlobVisibility } from "../../types.js"; -function parseVisibility(value: string): BlobVisibility { +export function parseVisibility(value: string): BlobVisibility { if ((BLOB_VISIBILITIES as readonly string[]).includes(value)) { return value as BlobVisibility; } diff --git a/src/commands/blob/credentials.ts b/src/commands/blob/credentials.ts index 20fea58..1137173 100644 --- a/src/commands/blob/credentials.ts +++ b/src/commands/blob/credentials.ts @@ -1,4 +1,4 @@ -import { Command } from "commander"; +import { Command, Option } from "commander"; import { resolveAuth } from "../../auth.js"; import { HttpError, request } from "../../client.js"; import { printJSON } from "../../output.js"; @@ -10,6 +10,7 @@ import { sleep, } from "./retry.js"; import type { Sleep } from "./retry.js"; +import { bucketIdArgument } from "./buckets.js"; const BLOB_CREDENTIALS_URL = "https://blob.upstash.io/v1/credentials"; const RETRYABLE_STATUSES = new Set([429, 503]); @@ -161,7 +162,7 @@ export function resolveBucketToken( ): Promise { if (flags.token !== undefined) { if (flags.bucketId !== undefined) { - return Promise.reject(new Error("Use either --token or --bucket-id, not both")); + return Promise.reject(new Error("Use either --token or a bucket, not both")); } const token = flags.token.trim(); if (!token) return Promise.reject(new Error("--token must be a non-empty Blob bucket token")); @@ -187,28 +188,31 @@ export function resolveBucketToken( return Promise.reject( new Error( - "Provide --token, UPSTASH_BLOB_TOKEN in the environment or .env, or --bucket-id with Upstash account authentication", + "Name a bucket (with Upstash account authentication), or provide --token or UPSTASH_BLOB_TOKEN in the environment or .env", ), ); } export function registerBlobCredentials(blob: Command): void { blob - .command("credentials") + .command("credentials [bucket]") .description( "Get temporary S3 credentials for a Blob bucket; expiresAt is the credential expiry", ) - .option("--bucket-id ", "Blob bucket ID") + .addOption(new Option("--bucket-id ").hideHelp()) .option("--token ", "Blob bucket token; no management API key needed (overrides UPSTASH_BLOB_TOKEN)") .addHelpText( "after", ` -With --bucket-id, a bucket created in the last few minutes is polled for up +With a bucket name or id, a bucket created in the last few minutes is polled for up to ~30s until provisioning finishes, so it is safe to run right after create. `, ) - .action(async (flags: { bucketId?: string; token?: string }, command: Command) => { - const source = await resolveBucketToken(flags, command); + .action(async (name: string | undefined, flags: { bucketId?: string; token?: string }, command: Command) => { + const named = name !== undefined || flags.bucketId !== undefined; + if (named && flags.token !== undefined) throw new Error("Use either --token or a bucket, not both"); + const bucketId = named ? await bucketIdArgument(command, name, flags) : undefined; + const source = await resolveBucketToken({ token: flags.token, bucketId }, command); const credentials = await fetchBlobCredentials(source.token, sleep, { unauthorizedRetries: source.unauthorizedRetries, }); diff --git a/src/commands/blob/delete.ts b/src/commands/blob/delete.ts index 5d3fcef..bb9ea0e 100644 --- a/src/commands/blob/delete.ts +++ b/src/commands/blob/delete.ts @@ -1,7 +1,8 @@ -import { Command } from "commander"; +import { Command, Option } from "commander"; import { resolveAuth } from "../../auth.js"; import { HttpError, request } from "../../client.js"; import { printJSON } from "../../output.js"; +import { bucketIdArgument } from "./buckets.js"; import { sleep } from "./retry.js"; import type { Sleep } from "./retry.js"; import type { Auth } from "../../auth.js"; @@ -39,17 +40,17 @@ export async function deleteBlobBucket( export function registerBlobDelete(blob: Command): void { blob - .command("delete") - .description("Delete a Blob bucket") - .requiredOption("--bucket-id ", "Blob bucket ID") - .option("--dry-run", "Preview the action without executing it") - .action(async (flags: { bucketId: string; dryRun?: boolean }, command: Command) => { + .command("delete [bucket]") + .description("Delete an empty Blob bucket, by name or id") + .addOption(new Option("--bucket-id ").hideHelp()) + .option("-n, --dry-run", "Preview the action without executing it") + .action(async (name: string | undefined, flags: { bucketId?: string; dryRun?: boolean }, command: Command) => { + const id = await bucketIdArgument(command, name, flags); if (flags.dryRun) { - printJSON({ action: "delete", bucket_id: flags.bucketId, dry_run: true }); + printJSON({ action: "delete", bucket_id: id, dry_run: true }); return; } - const auth = resolveAuth(command); - await deleteBlobBucket(auth, flags.bucketId); - printJSON({ deleted: true, bucket_id: flags.bucketId }); + await deleteBlobBucket(resolveAuth(command), id); + printJSON({ deleted: true, bucket_id: id }); }); } diff --git a/src/commands/blob/get.ts b/src/commands/blob/get.ts index c560201..10e2e6d 100644 --- a/src/commands/blob/get.ts +++ b/src/commands/blob/get.ts @@ -1,18 +1,19 @@ -import { Command } from "commander"; +import { Command, Option } from "commander"; import { resolveAuth } from "../../auth.js"; import { request } from "../../client.js"; import { printJSON } from "../../output.js"; import type { BlobBucket } from "../../types.js"; +import { bucketIdArgument } from "./buckets.js"; export function registerBlobGet(blob: Command): void { blob - .command("get") - .description("Get details of a Blob bucket") - .requiredOption("--bucket-id ", "Blob bucket ID") + .command("get [bucket]") + .description("Get details of a Blob bucket, by name or id") + .addOption(new Option("--bucket-id ").hideHelp()) .option("--hide-credentials", "Omit bucket tokens from output") - .action(async (flags: { bucketId: string; hideCredentials?: boolean }, command: Command) => { - const auth = resolveAuth(command); - const bucket = await request(auth, "GET", `/v2/blob/bucket/${flags.bucketId}`); + .action(async (name: string | undefined, flags: { bucketId?: string; hideCredentials?: boolean }, command: Command) => { + const id = await bucketIdArgument(command, name, flags); + const bucket = await request(resolveAuth(command), "GET", `/v2/blob/bucket/${id}`); if (!flags.hideCredentials) { printJSON(bucket); return; diff --git a/src/commands/blob/index.ts b/src/commands/blob/index.ts index f2678c4..3ca224c 100644 --- a/src/commands/blob/index.ts +++ b/src/commands/blob/index.ts @@ -1,18 +1,28 @@ import { Command } from "commander"; import { registerBlobCreate } from "./create.js"; -import { registerBlobList } from "./list.js"; import { registerBlobGet } from "./get.js"; import { registerBlobDelete } from "./delete.js"; import { registerBlobCredentials } from "./credentials.js"; -import { registerBlobUpload } from "./upload.js"; +import { registerBlobLs } from "./ls.js"; +import { registerBlobCp, registerBlobMv } from "./cp.js"; +import { registerBlobRm } from "./rm.js"; +import { registerBlobSync } from "./sync.js"; +import { registerBlobPresign } from "./presign.js"; +import { registerBlobMb, registerBlobRb } from "./mb.js"; export function registerBlob(program: Command): void { - const blob = program.command("blob").description("Manage Blob buckets"); + const blob = program.command("blob").description("Manage Blob buckets and objects"); registerBlobCreate(blob); - registerBlobList(blob); registerBlobGet(blob); registerBlobDelete(blob); registerBlobCredentials(blob); - registerBlobUpload(blob); + registerBlobLs(blob); + registerBlobCp(blob); + registerBlobMv(blob); + registerBlobRm(blob); + registerBlobSync(blob); + registerBlobPresign(blob); + registerBlobMb(blob); + registerBlobRb(blob); } diff --git a/src/commands/blob/list.ts b/src/commands/blob/list.ts deleted file mode 100644 index 562641f..0000000 --- a/src/commands/blob/list.ts +++ /dev/null @@ -1,16 +0,0 @@ -import { Command } from "commander"; -import { resolveAuth } from "../../auth.js"; -import { request } from "../../client.js"; -import { printJSON } from "../../output.js"; -import type { BlobBucket } from "../../types.js"; - -export function registerBlobList(blob: Command): void { - blob - .command("list") - .description("List Blob buckets") - .action(async (flags: Record, command: Command) => { - const auth = resolveAuth(command); - const buckets = await request(auth, "GET", "/v2/blob/bucket"); - printJSON(buckets); - }); -} diff --git a/src/commands/blob/ls.ts b/src/commands/blob/ls.ts new file mode 100644 index 0000000..b89d1f7 --- /dev/null +++ b/src/commands/blob/ls.ts @@ -0,0 +1,67 @@ +import { Command } from "commander"; +import { printJSON } from "../../output.js"; +import { BucketResolver, listDirectory } from "./buckets.js"; +import type { ListedObject } from "./buckets.js"; +import { parseBlobLocation } from "./transfer.js"; + +interface ListOptions { + recursive?: boolean; + summarize?: boolean; + token?: string; +} + +export function registerBlobLs(blob: Command): void { + blob + .command("ls [path]") + .alias("list") + .description("List buckets, or the objects and prefixes under /, like aws s3 ls") + .option("-r, --recursive", "List every object under the prefix instead of one level") + .option("--summarize", "Add total_objects and total_size") + .option("--token ", "Blob bucket token, used for the bucket it was issued for (default: UPSTASH_BLOB_TOKEN)") + .addHelpText("after", ` + is my-bucket/prefix or blob://my-bucket/prefix, by bucket name or id. +Without --recursive the prefix is matched as typed, so my-bucket/img lists "img/" +as a prefix; add the slash to list inside it. + +Examples: + upstash blob ls + upstash blob ls my-bucket + upstash blob ls my-bucket/images/ -r --summarize +`) + .action(async (uri: string | undefined, options: ListOptions, cmd: Command) => { + const resolver = new BucketResolver(cmd, options.token); + if (uri === undefined) { + printJSON(await resolver.listAccountBuckets()); + return; + } + + const location = parseBlobLocation(uri); + const bucket = await resolver.open(location.bucket); + const prefixes: string[] = []; + const objects: ListedObject[] = []; + let cursor: string | undefined; + do { + if (options.recursive) { + const page = await bucket.list({ prefix: location.key, limit: 1000, cursor }); + for (const object of page.blobs) { + objects.push({ key: object.path, size: object.size, last_modified: object.uploadedAt.toISOString(), etag: object.etag }); + } + cursor = page.cursor; + } else { + const page = await listDirectory(bucket, location.key, cursor); + prefixes.push(...page.prefixes); + objects.push(...page.objects); + cursor = page.cursor; + } + } while (cursor); + + printJSON({ + ...(!options.recursive && { prefixes }), + objects, + ...(options.summarize && { + total_objects: objects.length, + total_size: objects.reduce((sum, object) => sum + object.size, 0), + }), + }); + }); +} diff --git a/src/commands/blob/mb.ts b/src/commands/blob/mb.ts new file mode 100644 index 0000000..66690b1 --- /dev/null +++ b/src/commands/blob/mb.ts @@ -0,0 +1,87 @@ +import { Command } from "commander"; +import { resolveAuth } from "../../auth.js"; +import { HttpError, request } from "../../client.js"; +import { printJSON } from "../../output.js"; +import { BLOB_VISIBILITIES } from "../../types.js"; +import type { BlobBucket, BlobVisibility } from "../../types.js"; +import { BucketResolver, findAccountBucket } from "./buckets.js"; +import { parseVisibility } from "./create.js"; +import { deleteBlobBucket } from "./delete.js"; +import { listBlobs, parseBucket, runOperations } from "./transfer.js"; + +export function registerBlobMb(blob: Command): void { + blob + .command("mb ") + .description("Create a bucket, like aws s3 mb (same as blob create)") + .option( + "--visibility ", + `Bucket visibility. Available: ${BLOB_VISIBILITIES.join(", ")}`, + parseVisibility, + "private", + ) + .option("--cors ", "Allowed CORS origins (space-separated)") + .action(async (uri: string, flags: { visibility: BlobVisibility; cors?: string[] }, cmd: Command) => { + const name = parseBucket(uri); + const bucket = await request(resolveAuth(cmd), "POST", "/v2/blob/bucket", { + name, + visibility: flags.visibility, + cors: flags.cors, + }); + printJSON(bucket); + }); +} + +export function registerBlobRb(blob: Command): void { + blob + .command("rb ") + .description("Delete an empty bucket, or any bucket with --force, like aws s3 rb") + .option("-f, --force", "First delete every object and abort incomplete multipart uploads") + .option("-q, --quiet", "Suppress progress on stderr") + .action(async (uri: string, flags: { force?: boolean; quiet?: boolean }, cmd: Command) => { + const name = parseBucket(uri); + const resolver = new BucketResolver(cmd); + const auth = resolver.auth(); + const match = findAccountBucket(await resolver.listAccountBuckets(), name); + let deletedObjects = 0; + let abortedUploads = 0; + if (flags.force) { + const bucket = await resolver.open(match.id); + // Addressed by id, which the deletes below reuse from the resolver's cache. + const location = { type: "blob" as const, bucket: match.id, key: "" }; + const entries = await listBlobs(bucket, location, "", true); + const summary = await runOperations( + entries.map((entry) => ({ action: "delete", source: entry.location, size: 0 })), + { + resolver, + concurrency: 1, + signal: new AbortController().signal, + progress: (line) => { + if (!flags.quiet) console.error(line); + }, + }, + ); + if (summary.failed.length > 0) { + throw new Error(`${summary.failed.length} objects could not be deleted, so the bucket was kept`); + } + deletedObjects = summary.completed; + for (const upload of await bucket.listMultipartUploads()) { + await bucket.abortMultipartUpload(upload); + abortedUploads++; + } + } + try { + await deleteBlobBucket(auth, match.id); + } catch (error) { + if (error instanceof HttpError && error.status === 400 && /not empty/i.test(error.message) && !flags.force) { + throw new Error("bucket is not empty; rerun with --force to delete its objects first"); + } + throw error; + } + printJSON({ + deleted: true, + bucket: match.name, + bucket_id: match.id, + ...(flags.force && { objects_deleted: deletedObjects, multipart_uploads_aborted: abortedUploads }), + }); + }); +} diff --git a/src/commands/blob/presign.ts b/src/commands/blob/presign.ts new file mode 100644 index 0000000..46b901f --- /dev/null +++ b/src/commands/blob/presign.ts @@ -0,0 +1,37 @@ +import { Command, InvalidArgumentError } from "commander"; +import { printJSON } from "../../output.js"; +import { BucketResolver } from "./buckets.js"; +import { formatLocation, parseBlobLocation } from "./transfer.js"; + +const MAX_EXPIRES_IN = 7 * 24 * 60 * 60; + +function expiresIn(value: string): number { + const seconds = Number(value); + if (!Number.isInteger(seconds) || seconds < 1 || seconds > MAX_EXPIRES_IN) { + throw new InvalidArgumentError(`must be a whole number of seconds from 1 to ${MAX_EXPIRES_IN}`); + } + return seconds; +} + +export function registerBlobPresign(blob: Command): void { + blob + .command("presign ") + .description("Create a temporary download URL for an object, like aws s3 presign") + .option("--expires-in ", "How long the URL should work", expiresIn, 3600) + .option("--token ", "Blob bucket token, used for the bucket it was issued for (default: UPSTASH_BLOB_TOKEN)") + .addHelpText("after", ` +A URL never outlives the temporary credential that signs it, so it can expire +sooner than requested; expires_at is when it actually stops working. +`) + .action(async (uri: string, options: { expiresIn: number; token?: string }, cmd: Command) => { + const location = parseBlobLocation(uri); + if (!location.key || location.key.endsWith("/")) throw new Error(`${formatLocation(location)} names no object`); + const bucket = await new BucketResolver(cmd, options.token).open(location.bucket); + const signed = await bucket.signedReadUrl(location.key, { expiresIn: options.expiresIn }); + const requested = Date.now() + options.expiresIn * 1000; + if (signed.expiresAt.getTime() < requested - 60_000) { + console.error(`note: the URL expires at ${signed.expiresAt.toISOString()}, sooner than requested, because its signing credential expires then`); + } + printJSON({ url: signed.url, expires_at: signed.expiresAt.toISOString() }); + }); +} diff --git a/src/commands/blob/rm.ts b/src/commands/blob/rm.ts new file mode 100644 index 0000000..06035ea --- /dev/null +++ b/src/commands/blob/rm.ts @@ -0,0 +1,52 @@ +import { Command } from "commander"; +import { BucketResolver } from "./buckets.js"; +import { + addCommonOptions, + dirPrefix, + executePlan, + filtersOf, + formatLocation, + isIncluded, + listBlobs, + parseBlobLocation, +} from "./transfer.js"; +import type { CommonOptions, Operation } from "./transfer.js"; + +interface RemoveOptions extends CommonOptions { + recursive?: boolean; +} + +export function registerBlobRm(blob: Command): void { + const command = blob + .command("rm ") + .description("Delete an object, or every object below a prefix with --recursive, like aws s3 rm") + .option("-r, --recursive", "Delete every object below the prefix; a bucket alone empties it"); + addCommonOptions(command); + command + .addHelpText("after", ` +Examples: + upstash blob rm my-bucket/old.txt + upstash blob rm my-bucket/tmp -r -n +`) + .action(async (uri: string, options: RemoveOptions, cmd: Command) => { + const location = parseBlobLocation(uri); + const resolver = new BucketResolver(cmd, options.token); + let operations: Operation[]; + if (options.recursive) { + // A bare `rm -r build` is too easily meant as a local folder to empty a whole bucket. + if (!location.key && !uri.startsWith("blob://")) { + throw new Error(`to delete every object in ${location.bucket}, write blob://${location.bucket}, or remove the bucket with rb -f`); + } + const bucket = await resolver.open(location.bucket); + const filters = filtersOf(options); + const entries = await listBlobs(bucket, location, dirPrefix(location.key), true); + operations = entries + .filter((entry) => isIncluded(entry.rel, filters)) + .map((entry) => ({ action: "delete", source: entry.location, size: 0 })); + } else { + if (!location.key) throw new Error(`${formatLocation(location)} names no object; add a key, or --recursive to delete everything`); + operations = [{ action: "delete", source: location, size: 0 }]; + } + await executePlan(operations, [], options, resolver); + }); +} diff --git a/src/commands/blob/sync.ts b/src/commands/blob/sync.ts new file mode 100644 index 0000000..8d25a00 --- /dev/null +++ b/src/commands/blob/sync.ts @@ -0,0 +1,133 @@ +import { Command } from "commander"; +import { BucketResolver } from "./buckets.js"; +import { + addCommonOptions, + addObjectOptions, + checkOverwritesSource, + claimLocal, + destinationFor, + executePlan, + filtersOf, + formatLocation, + isIncluded, + listDestination, + listSource, + localFileKey, + localRel, + nestedPrefix, + realLocalPath, + parseLocation, + transferAction, +} from "./transfer.js"; +import type { Action, CommonOptions, Entry, Failure, ObjectOptions, Operation } from "./transfer.js"; + +interface SyncOptions extends CommonOptions, ObjectOptions { + delete?: boolean; + sizeOnly?: boolean; + exactTimestamps?: boolean; +} + +/** + * aws s3 sync's rule: copy when sizes differ or the source is newer. Downloads are the exception, + * skipping same-sized files unless the local copy is newer, or with --exact-timestamps unless the + * times differ at all. Times compare at whole seconds, the precision LastModified has. + */ +export function needsSync( + source: Entry, + destination: Entry, + action: Exclude, + options: Pick, +): boolean { + if (source.size !== destination.size) return true; + if (options.sizeOnly) return false; + const sourceTime = Math.floor(source.mtime / 1000); + const destinationTime = Math.floor(destination.mtime / 1000); + if (action === "download") { + return options.exactTimestamps ? sourceTime !== destinationTime : destinationTime > sourceTime; + } + return sourceTime > destinationTime; +} + +export function registerBlobSync(blob: Command): void { + const command = blob + .command("sync ") + .description("Copy new and changed files between a directory and a blob:// prefix, or two prefixes, like aws s3 sync") + .option("-d, --delete", "Delete destination files that are not in the source (filters still apply)") + .option("--size-only", "Compare sizes only, ignoring modification times") + .option("--exact-timestamps", "When downloading, also copy same-sized files whose times differ"); + addObjectOptions(command); + addCommonOptions(command); + command + .addHelpText("after", ` +A file is copied when it is missing at the destination, its size differs, or the +source is newer. Downloads set local modification times to the object's upload +time, so a second sync copies nothing. + +Examples: + upstash blob sync ./site blob://my-bucket/site -d + upstash blob sync blob://my-bucket/backups ./backups --exclude "*" --include "*.gz" +`) + .action(async (sourceArg: string, destinationArg: string, options: SyncOptions, cmd: Command) => { + const source = parseLocation(sourceArg); + const destination = parseLocation(destinationArg); + if (source.type === "local" && destination.type === "local") { + throw new Error("source or destination must be a blob:/// URI"); + } + const resolver = new BucketResolver(cmd, options.token); + const filters = filtersOf(options); + const action = transferAction(source, destination); + const [sources, existing, destinationInside, sourceInside] = await Promise.all([ + listSource(source, true, resolver), + listDestination(destination, resolver), + nestedPrefix(source, destination, resolver), + nestedPrefix(destination, source, resolver), + ]); + const outside = (entry: Entry, prefix: string | undefined): boolean => + prefix === undefined || entry.location.type !== "blob" || !entry.location.key.startsWith(prefix); + const included = (entry: Entry): boolean => + isIncluded(entry.rel, filters) && !(destination.type === "local" && entry.rel.endsWith("/")); + // With nested prefixes in one bucket, neither side's listing includes the other side's keys, + // and a copy onto a source key is refused. + const current = new Map(existing.filter((entry) => included(entry) && outside(entry, sourceInside)) + .map((entry) => [entry.rel, entry])); + + const operations: Operation[] = []; + const failures: Failure[] = []; + let unchanged = 0; + const wanted = new Set(); + const claimed = new Set(); + const written: string[] = []; + for (const entry of sources.filter((entry) => included(entry) && outside(entry, destinationInside))) { + // Keys like "a//b" land on local "a/b", which is what the local listing reports. + const rel = destination.type === "local" ? localRel(entry.rel) : entry.rel; + wanted.add(rel); + const target = current.get(rel); + try { + const to = target?.location ?? destinationFor(entry, destination, true); + await claimLocal(claimed, to); + if (to.type === "local") written.push(to.path); + checkOverwritesSource(to, sourceInside); + if (target && !needsSync(entry, target, action, options)) { + unchanged++; + continue; + } + operations.push({ action, source: entry.location, destination: to, size: entry.size }); + } catch (error) { + failures.push({ source: formatLocation(entry.location), error: (error as Error).message }); + } + } + if (options.delete) { + // A local file may be one a source was just written to, or kept in, under another name: + // "a.txt" for "A.txt" on a case-insensitive disk, or a path through a symbolic link. Case is + // folded on every OS, since Linux can mount such disks; at worst a stale file stays. + const fileKey = async (path: string): Promise => localFileKey(await realLocalPath(path), true); + const kept = new Set(await Promise.all(written.map(fileKey))); + for (const [rel, entry] of current) { + if (wanted.has(rel)) continue; + if (entry.location.type === "local" && kept.has(await fileKey(entry.location.path))) continue; + operations.push({ action: "delete", source: entry.location, size: 0 }); + } + } + await executePlan(operations, failures, options, resolver, { unchanged }); + }); +} diff --git a/src/commands/blob/transfer.ts b/src/commands/blob/transfer.ts new file mode 100644 index 0000000..b1abb28 --- /dev/null +++ b/src/commands/blob/transfer.ts @@ -0,0 +1,761 @@ +import { BlobError } from "@upstash/blob"; +import type { Bucket, PutOptions } from "@upstash/blob"; +import { InvalidArgumentError, Option } from "commander"; +import type { Command } from "commander"; +import { randomUUID } from "node:crypto"; +import { setMaxListeners } from "node:events"; +import { createReadStream, createWriteStream } from "node:fs"; +import { mkdir, opendir, realpath, rename, rm, stat, unlink, utimes } from "node:fs/promises"; +import { basename, dirname, isAbsolute, join, posix, relative, resolve, sep } from "node:path"; +import { Readable } from "node:stream"; +import { pipeline } from "node:stream/promises"; +import type { ReadableStream as NodeReadableStream } from "node:stream/web"; +import mime from "mime"; +import { printJSON } from "../../output.js"; +import type { BucketResolver } from "./buckets.js"; +import { sleep } from "./retry.js"; + +export type Location = + | { type: "local"; path: string } + | { type: "blob"; bucket: string; key: string }; + +export type BlobLocation = Extract; + +const hasScheme = (value: string): boolean => /^[a-z][a-z0-9+.-]*:\/\//i.test(value); + +export function parseLocation(value: string): Location { + if (value.startsWith("blob://")) { + const rest = value.slice("blob://".length); + const slash = rest.indexOf("/"); + const bucket = slash === -1 ? rest : rest.slice(0, slash); + if (!bucket) throw new Error(`"${value}" has no bucket: use blob:///`); + return { type: "blob", bucket, key: slash === -1 ? "" : rest.slice(slash + 1) }; + } + if (hasScheme(value)) { + throw new Error(`unsupported location "${value}": use blob:/// or a local path`); + } + return { type: "local", path: value }; +} + +/** For arguments that can only be in a bucket, where `my-bucket/key` means `blob://my-bucket/key`. */ +export function parseBlobLocation(value: string): BlobLocation { + const location = parseLocation(hasScheme(value) ? value : `blob://${value}`); + if (location.type !== "blob") throw new Error(`"${value}" is not a blob:/// URI`); + return location; +} + +/** A whole bucket: `my-bucket`, or `blob://my-bucket`, by name or id. */ +export function parseBucket(value: string): string { + const location = parseBlobLocation(value); + if (location.key.replace(/\/+$/, "")) throw new Error(`${formatLocation(location)} includes a key; name only the bucket`); + return location.bucket; +} + +export function formatLocation(location: Location): string { + return location.type === "blob" ? `blob://${location.bucket}/${location.key}` : location.path; +} + +/** How aws s3 treats the source of a recursive command: as a directory, whether or not it ends in a slash. */ +export function dirPrefix(key: string): string { + return key === "" || key.endsWith("/") ? key : `${key}/`; +} + +// ── Filters ───────────────────────────────────────────────────────────────── + +interface FilterFlag { + exclude: boolean; + pattern: string; + order: number; +} + +export interface Filter { + exclude: boolean; + match: RegExp; +} + +let filterOrder = 0; + +function collectFilter(exclude: boolean) { + return (pattern: string, previous: FilterFlag[] = []): FilterFlag[] => [ + ...previous, + { exclude, pattern, order: filterOrder++ }, + ]; +} + +/** fnmatch, as aws s3 uses it: `*` also crosses slashes. */ +export function globToRegExp(glob: string): RegExp { + let source = ""; + for (let i = 0; i < glob.length; i++) { + const char = glob[i] ?? ""; + if (char === "*") source += ".*"; + else if (char === "?") source += "."; + else if (char === "[") { + let end = i + 1; + if (glob[end] === "!") end++; + if (glob[end] === "]") end++; + while (end < glob.length && glob[end] !== "]") end++; + if (end >= glob.length) { + source += "\\["; + continue; + } + let body = glob.slice(i + 1, end).replace(/[\\\]]/g, "\\$&"); + if (body.startsWith("!")) body = `^${body.slice(1)}`; + else if (body.startsWith("^")) body = `\\${body}`; + source += `[${body}]`; + i = end; + } else source += char.replace(/[.*+?^${}()|[\]\\/]/g, "\\$&"); + } + return new RegExp(`^${source}$`, "s"); +} + +export function filtersOf(options: { exclude?: FilterFlag[]; include?: FilterFlag[] }): Filter[] { + return [...(options.exclude ?? []), ...(options.include ?? [])] + .sort((a, b) => a.order - b.order) + .map((flag) => ({ exclude: flag.exclude, match: globToRegExp(flag.pattern) })); +} + +/** Everything is included until a filter matches; the last matching filter decides. */ +export function isIncluded(path: string, filters: Filter[]): boolean { + let included = true; + for (const filter of filters) if (filter.match.test(path)) included = !filter.exclude; + return included; +} + +// ── Options ───────────────────────────────────────────────────────────────── + +function concurrency(value: string): number { + const number = Number(value); + if (!Number.isInteger(number) || number < 1 || number > 16) { + throw new InvalidArgumentError("concurrency must be an integer from 1 to 16"); + } + return number; +} + +function parseMetadata(value: string): Record { + const text = value.trim(); + if (text.startsWith("{")) { + let parsed: unknown; + try { + parsed = JSON.parse(text); + } catch { + throw new InvalidArgumentError("--metadata JSON is invalid"); + } + if (!parsed || typeof parsed !== "object" || Array.isArray(parsed) || + !Object.values(parsed).every((entry) => typeof entry === "string")) { + throw new InvalidArgumentError("--metadata JSON must map keys to string values"); + } + return parsed as Record; + } + const metadata: Record = {}; + for (const pair of text.split(",")) { + const eq = pair.indexOf("="); + if (eq <= 0) throw new InvalidArgumentError("--metadata must be key=value[,key=value] or a JSON object"); + metadata[pair.slice(0, eq).trim()] = pair.slice(eq + 1); + } + return metadata; +} + +export interface CommonOptions { + token?: string; + dryrun?: boolean; + dryRun?: boolean; + quiet?: boolean; + concurrency: number; + exclude?: FilterFlag[]; + include?: FilterFlag[]; +} + +export interface ObjectOptions { + contentType?: string; + cacheControl?: string; + metadata?: Record; +} + +export function addCommonOptions(command: Command): Command { + return command + .option("--exclude ", "Skip paths matching this pattern; repeatable, the last matching filter wins", collectFilter(true)) + .option("--include ", "Keep paths matching this pattern even if an earlier --exclude matched", collectFilter(false)) + .option("-n, --dryrun", "Show what would happen without changing anything") + .addOption(new Option("--dry-run").hideHelp()) + .option("-q, --quiet", "Suppress progress on stderr; still print the JSON summary") + .option("--concurrency ", "Number of concurrent transfers (1-16)", concurrency, 4) + .option("--token ", "Blob bucket token, used for the bucket it was issued for (default: UPSTASH_BLOB_TOKEN)"); +} + +export function addObjectOptions(command: Command): Command { + return command + .option("--content-type ", "Content type for every written object (default: guessed from the file name, or kept on copies)") + .option("--cache-control ", "Cache-Control for written objects: a header value, a duration like 1h, immutable, revalidate or no-store") + .option("--metadata ", "Metadata for written objects, as key=value[,key=value] or a JSON object", parseMetadata); +} + +export function isDryRun(options: { dryrun?: boolean; dryRun?: boolean }): boolean { + return Boolean(options.dryrun || options.dryRun); +} + +// ── Listing ───────────────────────────────────────────────────────────────── + +export interface Entry { + /** Path below the listed directory or prefix, with forward slashes. */ + rel: string; + size: number; + /** Local mtime or the object's LastModified, in milliseconds. */ + mtime: number; + location: Location; +} + +function errorCode(error: unknown): string | undefined { + return (error as NodeJS.ErrnoException | undefined)?.code; +} + +/** Symbolic links are followed as aws s3 does; one that loops back to its own ancestor is skipped. */ +async function walk(root: string): Promise { + const entries: Entry[] = []; + const visit = async (directory: string, ancestors: Set): Promise => { + const real = await realpath(directory); + if (ancestors.has(real)) return; + const chain = new Set(ancestors).add(real); + for await (const dirent of await opendir(directory)) { + const path = join(directory, dirent.name); + // A dangling link has nothing to copy. + const info = dirent.isSymbolicLink() ? await stat(path).catch(() => undefined) : dirent; + if (info?.isDirectory()) await visit(path, chain); + else if (info?.isFile()) { + const { size, mtimeMs } = await stat(path); + entries.push({ + rel: relative(root, path).split(sep).join("/"), + size, + mtime: mtimeMs, + location: { type: "local", path }, + }); + } + } + }; + await visit(root, new Set()); + return entries; +} + +/** + * Every object below `prefix`. Zero-byte "folder" markers ending in "/" are skipped, as aws s3 skips + * them everywhere except deletes; `forDelete` keeps them, the marker at the prefix itself included. + */ +export async function listBlobs(bucket: Bucket, location: BlobLocation, prefix: string, forDelete = false): Promise { + const entries: Entry[] = []; + let cursor: string | undefined; + do { + const page = await bucket.list({ prefix, limit: 1000, cursor }); + for (const blob of page.blobs) { + const rel = blob.path.slice(prefix.length); + if (!forDelete && (!rel || (blob.size === 0 && rel.endsWith("/")))) continue; + entries.push({ + rel, + size: blob.size, + mtime: blob.uploadedAt.getTime(), + location: { type: "blob", bucket: location.bucket, key: blob.path }, + }); + } + cursor = page.cursor; + } while (cursor); + return entries; +} + +/** What a cp, mv or sync reads: one file or object, or everything below a directory or prefix. */ +export async function listSource(location: Location, recursive: boolean, resolver: BucketResolver): Promise { + if (location.type === "local") { + let info; + try { + info = await stat(location.path); + } catch (error) { + if (errorCode(error) === "ENOENT") throw new Error(`${location.path} does not exist`); + throw error; + } + if (info.isDirectory()) { + if (!recursive) throw new Error(`${location.path} is a directory; use --recursive`); + return walk(location.path); + } + if (!info.isFile()) throw new Error(`${location.path} is not a regular file`); + if (recursive) throw new Error(`${location.path} is a file; --recursive needs a directory`); + return [{ rel: basename(location.path), size: info.size, mtime: info.mtimeMs, location }]; + } + + const bucket = await resolver.open(location.bucket); + if (recursive) return listBlobs(bucket, location, dirPrefix(location.key)); + if (!location.key || location.key.endsWith("/")) { + throw new Error(`${formatLocation(location)} is a prefix; use --recursive`); + } + const info = await bucket.info(location.key); + return [{ + rel: location.key.slice(location.key.lastIndexOf("/") + 1), + size: info.size, + mtime: info.uploadedAt.getTime(), + location, + }]; +} + +/** What a sync compares against; a missing local directory is simply empty. */ +export async function listDestination(location: Location, resolver: BucketResolver): Promise { + if (location.type === "blob") { + return listBlobs(await resolver.open(location.bucket), location, dirPrefix(location.key)); + } + try { + const info = await stat(location.path); + if (!info.isDirectory()) throw new Error(`${location.path} is not a directory`); + } catch (error) { + if (errorCode(error) === "ENOENT") return []; + throw error; + } + return walk(location.path); +} + +export async function isLocalDirectory(path: string): Promise { + if (path.endsWith("/") || path.endsWith(sep)) return true; + try { + return (await stat(path)).isDirectory(); + } catch { + return false; + } +} + +// ── Planning ──────────────────────────────────────────────────────────────── + +export type Action = "upload" | "download" | "copy" | "delete"; + +export interface Operation { + action: Action; + source: Location; + destination?: Location; + size: number; + /** Delete the source once the transfer succeeds. */ + move?: boolean; +} + +export interface Failure { + source: string; + destination?: string; + error: string; +} + +export function transferAction(source: Location, destination: Location): Exclude { + if (source.type === "local") return "upload"; + return destination.type === "local" ? "download" : "copy"; +} + +/** + * Where one entry lands. A blob key ending in "/" or empty, a local directory, and every recursive + * command take the source's relative path; otherwise the destination names the object or file. + */ +export function destinationFor(entry: Entry, destination: Location, into: boolean): Location { + if (destination.type === "blob") { + const key = into ? dirPrefix(destination.key) + entry.rel : destination.key; + return { type: "blob", bucket: destination.bucket, key }; + } + if (!into) return destination; + const root = resolve(destination.path); + const rel = relative(root, resolve(root, ...entry.rel.split("/"))); + if (!rel || rel === ".." || rel.startsWith(`..${sep}`) || isAbsolute(rel)) { + throw new Error(`${JSON.stringify(entry.rel)} would be written outside ${destination.path}`); + } + return { type: "local", path: join(destination.path, rel) }; +} + +/** Where `rel` lands below a local directory, as the directory walk would report it. */ +export function localRel(rel: string): string { + return posix.normalize(rel).replace(/^\/+/, ""); +} + +/** macOS and Windows disks ignore case, and macOS also Unicode normalization. */ +export const foldsCase = process.platform === "darwin" || process.platform === "win32"; + +/** Identifies a local file, folding case and Unicode form where the disk does or where a wrong guess loses data. */ +export function localFileKey(path: string, fold = foldsCase): string { + const resolved = resolve(path); + return fold ? resolved.normalize("NFC").toLowerCase() : resolved; +} + +/** The path with symbolic links resolved as far as it exists, so aliases of one file compare equal. */ +export async function realLocalPath(path: string): Promise { + try { + return await realpath(path); + } catch { + const parent = dirname(resolve(path)); + if (parent === resolve(path)) return parent; + return join(await realLocalPath(parent), basename(path)); + } +} + +/** + * Two keys can land on one local file: `a//b` and `a/b`, `A.txt` and `a.txt` on a case-insensitive + * disk, or `real/x` and `link/x` through a symbolic link. Only the first may write it, so an mv + * cannot delete the other's source. + */ +export async function claimLocal(claimed: Set, destination: Location, fold = foldsCase): Promise { + if (destination.type !== "local") return; + const key = localFileKey(await realLocalPath(destination.path), fold); + if (claimed.has(key)) throw new Error(`another object is also written to ${destination.path}`); + claimed.add(key); +} + +/** + * When `inner` is a prefix nested inside `outer` in the same bucket, returns it. With the + * destination inside, the source listing skips those keys: `mv blob://b/ blob://b/archive/` must + * not move its own output. With the source inside, copies onto keys below it are refused. + */ +export async function nestedPrefix(outer: Location, inner: Location, resolver: BucketResolver): Promise { + if (outer.type !== "blob" || inner.type !== "blob") return undefined; + const outerPrefix = dirPrefix(outer.key); + const innerPrefix = dirPrefix(inner.key); + if (innerPrefix === outerPrefix || !innerPrefix.startsWith(outerPrefix)) return undefined; + if (outer.bucket !== inner.bucket) { + const [a, b] = await Promise.all([resolver.open(outer.bucket), resolver.open(inner.bucket)]); + if (a.s3().bucket !== b.s3().bucket) return undefined; + } + return innerPrefix; +} + +/** Refuses a copy onto a key below `sourcePrefix`, which is another object of the same transfer. */ +export function checkOverwritesSource(destination: Location, sourcePrefix: string | undefined): void { + if (sourcePrefix !== undefined && destination.type === "blob" && destination.key.startsWith(sourcePrefix)) { + throw new Error(`${formatLocation(destination)} is also a source; the prefixes overlap`); + } +} + +export function describe(operation: Operation): string { + const verb = operation.move ? "move" : operation.action; + const source = formatLocation(operation.source); + return operation.destination ? `${verb}: ${source} to ${formatLocation(operation.destination)}` : `${verb}: ${source}`; +} + +// ── Execution ─────────────────────────────────────────────────────────────── + +export interface TransferSettings extends ObjectOptions { + resolver: BucketResolver; + concurrency: number; + signal: AbortSignal; + progress: (line: string) => void; +} + +export interface TransferSummary { + completed: number; + bytes: number; + failed: Failure[]; + remaining: number; +} + +/** R2 copies server-side only up to 5 GiB; anything larger is streamed through the CLI. */ +const MAX_SERVER_COPY_BYTES = 5 * 1024 ** 3; + +function message(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} + +/** Uploads also refuse control characters and backslashes, which local names should not carry into keys. */ +function checkKey(key: string, upload: boolean): void { + if (!key) throw new Error("destination key is empty"); + if (Buffer.byteLength(key) > 1024 || (upload && /[\x00-\x1f\x7f\\]/.test(key)) || + key.split("/").some((part) => part === "." || part === "..")) { + throw new Error(`unsupported object key: ${JSON.stringify(key)}`); + } +} + +function retryable(error: unknown): boolean { + return BlobError.is(error) && ( + error.code === "rate_limited" || error.code === "not_ready" || + (error.status !== undefined && error.status >= 500) + ) || error instanceof TypeError && (error.message === "fetch failed" || error.message === "terminated") || + error instanceof Error && error.name === "TimeoutError"; +} + +export interface UploadFile { + source: string; + key: string; + size: number; + contentType: string; +} + +/** Streams one local file to the bucket, retrying transient failures from the start of the file. */ +export async function putFile( + bucket: Pick, + file: UploadFile, + signal?: AbortSignal, + options: Pick = {}, +): Promise { + for (let attempt = 0; ; attempt++) { + if (signal?.aborted) throw new Error("upload interrupted"); + const current = await stat(file.source); + if (!current.isFile() || current.size !== file.size) { + throw new Error("source changed since scanning; rerun the command"); + } + const stream = createReadStream(file.source, { signal, highWaterMark: 64 * 1024 }); + try { + const body = Readable.toWeb(stream, { + strategy: { highWaterMark: 64 * 1024, size: (chunk: Buffer) => chunk.byteLength }, + }) as ReadableStream; + await bucket.put(file.key, body, { ...options, size: file.size, contentType: file.contentType }); + return; + } catch (error) { + if (attempt >= 2 || signal?.aborted || !retryable(error)) throw error; + } finally { + stream.destroy(); + } + await sleep(500 * 2 ** attempt); + } +} + +class TruncatedDownload extends Error {} + +async function withRetries(signal: AbortSignal, run: () => Promise): Promise { + for (let attempt = 0; ; attempt++) { + try { + return await run(); + } catch (error) { + if (attempt >= 2 || signal.aborted || !(retryable(error) || error instanceof TruncatedDownload)) throw error; + } + await sleep(500 * 2 ** attempt); + } +} + +/** Writes to a temporary file beside the target and renames it, so a failed download leaves nothing half-written. */ +async function download(bucket: Bucket, key: string, target: string, signal: AbortSignal): Promise { + await mkdir(dirname(target), { recursive: true }); + await withRetries(signal, async () => { + const temp = join(dirname(target), `.upstash-${randomUUID().slice(0, 8)}.tmp`); + try { + const blob = await bucket.get(key); + await pipeline( + Readable.fromWeb(blob.body as NodeReadableStream), + createWriteStream(temp, { flags: "wx" }), + { signal }, + ); + const written = (await stat(temp)).size; + if (written !== blob.size) throw new TruncatedDownload(`download ended after ${written} of ${blob.size} bytes`); + // aws s3 sync compares local mtimes with LastModified, so keep them equal after a download. + await utimes(temp, blob.uploadedAt, blob.uploadedAt); + await rename(temp, target); + } catch (error) { + await rm(temp, { force: true }); + throw error; + } + }); +} + +async function copyObject( + from: Bucket, + fromKey: string, + to: Bucket, + toKey: string, + size: number, + settings: TransferSettings, +): Promise { + if (from.s3().bucket === to.s3().bucket && size <= MAX_SERVER_COPY_BYTES) { + await from.copy(fromKey, toKey, { + ...(settings.contentType !== undefined && { contentType: settings.contentType }), + ...(settings.cacheControl !== undefined && { cache: settings.cacheControl }), + ...(settings.metadata !== undefined && { metadata: settings.metadata }), + }); + return; + } + await withRetries(settings.signal, async () => { + const blob = await from.get(fromKey); + try { + await to.put(toKey, blob.body, { + size: blob.size, + contentType: settings.contentType ?? blob.contentType, + metadata: settings.metadata ?? blob.metadata, + cache: settings.cacheControl, + }); + } catch (error) { + await blob.body.cancel().catch(() => undefined); + throw error; + } + }); +} + +async function perform(operation: Operation, settings: TransferSettings): Promise { + const { source, destination } = operation; + if (operation.action === "upload" && source.type === "local" && destination?.type === "blob") { + checkKey(destination.key, true); + const bucket = await settings.resolver.open(destination.bucket); + await putFile(bucket, { + source: source.path, + key: destination.key, + size: operation.size, + contentType: settings.contentType ?? mime.getType(source.path) ?? "application/octet-stream", + }, settings.signal, { cache: settings.cacheControl, metadata: settings.metadata }); + if (operation.move) await unlink(source.path); + return; + } + if (operation.action === "download" && source.type === "blob" && destination?.type === "local") { + const bucket = await settings.resolver.open(source.bucket); + await download(bucket, source.key, destination.path, settings.signal); + if (operation.move) await bucket.del(source.key); + return; + } + if (operation.action === "copy" && source.type === "blob" && destination?.type === "blob") { + checkKey(destination.key, false); + const from = await settings.resolver.open(source.bucket); + const to = await settings.resolver.open(destination.bucket); + if (from.s3().bucket === to.s3().bucket && source.key === destination.key) { + throw new Error("source and destination are the same object"); + } + await copyObject(from, source.key, to, destination.key, operation.size, settings); + if (operation.move) await from.del(source.key); + return; + } + throw new Error(`cannot ${operation.action} ${formatLocation(source)}`); +} + +function deletable(key: string): boolean { + return key.length > 0 && !key.split("/").some((part) => part === "." || part === ".."); +} + +async function deleteAll(operations: Operation[], settings: TransferSettings, summary: TransferSummary): Promise { + const fail = (location: Location, error: string): void => { + summary.failed.push({ source: formatLocation(location), error }); + settings.progress(`delete failed: ${formatLocation(location)} ${error}`); + }; + const byBucket = new Map(); + for (const operation of operations) { + if (settings.signal.aborted) return; + const location = operation.source; + if (location.type === "local") { + try { + await unlink(location.path); + summary.completed++; + settings.progress(describe(operation)); + } catch (error) { + fail(location, message(error)); + } + summary.remaining--; + } else { + const keys = byBucket.get(location.bucket) ?? []; + keys.push(location.key); + byBucket.set(location.bucket, keys); + } + } + for (const [name, keys] of byBucket) { + let bucket: Bucket | undefined; + let openError: unknown; + try { + bucket = await settings.resolver.open(name); + } catch (error) { + openError = error; + } + for (let i = 0; i < keys.length; i += 1000) { + if (settings.signal.aborted) return; + const chunk = keys.slice(i, i + 1000); + const valid = chunk.filter(deletable); + let failed = new Map(); + for (const key of chunk) if (!deletable(key)) failed.set(key, "unsupported object key"); + if (!bucket) { + for (const key of valid) failed.set(key, message(openError)); + } else if (valid.length > 0) { + try { + await bucket.del(valid); + } catch (error) { + const survivors = BlobError.is(error) && error.code === "partial_delete" ? error.failed ?? valid : valid; + failed = new Map([...failed, ...survivors.map((key): [string, string] => [key, message(error)])]); + } + } + for (const key of chunk) { + const location: Location = { type: "blob", bucket: name, key }; + const error = failed.get(key); + if (error === undefined) { + summary.completed++; + settings.progress(describe({ action: "delete", source: location, size: 0 })); + } else fail(location, error); + summary.remaining--; + } + } + } +} + +/** Transfers run concurrently and keep going past failures, as aws s3 does; deletes run after, in batches. */ +export async function runOperations(operations: Operation[], settings: TransferSettings): Promise { + const summary: TransferSummary = { completed: 0, bytes: 0, failed: [], remaining: operations.length }; + const transfers = operations.filter((operation) => operation.action !== "delete"); + const deletes = operations.filter((operation) => operation.action === "delete"); + let next = 0; + const worker = async (): Promise => { + while (!settings.signal.aborted) { + const operation = transfers[next++]; + if (!operation) return; + try { + await perform(operation, settings); + summary.completed++; + summary.bytes += operation.size; + settings.progress(describe(operation)); + } catch (error) { + const failure: Failure = { source: formatLocation(operation.source), error: message(error) }; + if (operation.destination) failure.destination = formatLocation(operation.destination); + summary.failed.push(failure); + const verb = operation.move ? "move" : operation.action; + settings.progress(`${verb} failed: ${describe(operation).slice(verb.length + 2)} ${failure.error}`); + } finally { + summary.remaining--; + } + } + }; + await Promise.all(Array.from({ length: Math.min(settings.concurrency, transfers.length) }, () => worker())); + if (!settings.signal.aborted) await deleteAll(deletes, settings, summary); + return summary; +} + +/** + * Shared tail of cp, mv, rm and sync: print the plan on --dryrun, otherwise run it with Ctrl+C + * stopping new work, then print the JSON summary and fail if anything did. + */ +export async function executePlan( + operations: Operation[], + failures: Failure[], + options: CommonOptions & ObjectOptions, + resolver: BucketResolver, + extra: Record = {}, +): Promise { + const progress = (line: string): void => { + if (!options.quiet) console.error(line); + }; + for (const failure of failures) progress(`skipped: ${failure.source} ${failure.error}`); + if (isDryRun(options)) { + for (const operation of operations) progress(`(dryrun) ${describe(operation)}`); + printJSON({ + dry_run: true, + operations: operations.map((operation) => ({ + action: operation.move ? "move" : operation.action, + source: formatLocation(operation.source), + ...(operation.destination && { destination: formatLocation(operation.destination) }), + size: operation.size, + })), + ...extra, + ...(failures.length > 0 && { skipped: failures }), + }); + if (failures.length > 0) throw new Error(`${failures.length} entries would be skipped; see the JSON output`); + return; + } + + const controller = new AbortController(); + // Each read stream and download adds an abort listener, removed only once it closes. + setMaxListeners(0, controller.signal); + const interrupt = (): void => controller.abort(); + process.once("SIGINT", interrupt); + process.once("SIGTERM", interrupt); + try { + const summary = await runOperations(operations, { + resolver, + concurrency: options.concurrency, + signal: controller.signal, + progress, + contentType: options.contentType, + cacheControl: options.cacheControl, + metadata: options.metadata, + }); + summary.failed.unshift(...failures); + printJSON({ ...summary, ...extra }); + if (controller.signal.aborted) throw new Error("interrupted; completed transfers remain"); + if (summary.failed.length > 0) { + throw new Error(`${summary.failed.length} of ${operations.length + failures.length} operations failed; see the JSON summary`); + } + } finally { + process.removeListener("SIGINT", interrupt); + process.removeListener("SIGTERM", interrupt); + } +} diff --git a/src/commands/blob/upload.ts b/src/commands/blob/upload.ts deleted file mode 100644 index 5218d36..0000000 --- a/src/commands/blob/upload.ts +++ /dev/null @@ -1,191 +0,0 @@ -import { BlobError, Bucket } from "@upstash/blob"; -import { Command, InvalidArgumentError } from "commander"; -import { setMaxListeners } from "node:events"; -import { createReadStream } from "node:fs"; -import { lstat, opendir } from "node:fs/promises"; -import { basename, join, relative, resolve, sep } from "node:path"; -import { Readable } from "node:stream"; -import mime from "mime"; -import { printJSON } from "../../output.js"; -import { telemetryStatus } from "../../telemetry.js"; -import { fetchBlobCredentials, resolveBucketToken } from "./credentials.js"; -import { sleep } from "./retry.js"; - -export interface UploadFile { - source: string; - path: string; - size: number; - contentType: string; -} - -export interface UploadSummary { - uploaded: number; - skipped: number; - bytes: number; - failed: { path: string; error: string }[]; - remaining: number; -} - -interface UploadOptions { - bucketId?: string; - token?: string; - prefix: string; - concurrency: number; - skipExisting?: boolean; - dryRun?: boolean; - quiet?: boolean; -} - -function concurrency(value: string): number { - const number = Number(value); - if (!Number.isInteger(number) || number < 1 || number > 16) { - throw new InvalidArgumentError("concurrency must be an integer from 1 to 16"); - } - return number; -} - -export async function planUpload(source: string, prefix: string): Promise { - const root = resolve(source); - const base = prefix.replace(/^\/+|\/+$/g, ""); - if (base.split("/").some((part) => part === "." || part === "..") || /[\x00-\x1f\x7f\\]/.test(base)) { - throw new Error("prefix must not contain dot segments, backslashes, or control characters"); - } - const rootStat = await lstat(root); - if (!rootStat.isFile() && !rootStat.isDirectory()) { - throw new Error("source must be a regular file or directory; symbolic links are not followed"); - } - const files: UploadFile[] = []; - const addFile = (path: string, size: number): void => { - const name = rootStat.isFile() ? basename(path) : relative(root, path).split(sep).join("/"); - const key = base ? `${base}/${name}` : name; - if (Buffer.byteLength(key) > 1024 || /[\x00-\x1f\x7f\\]/.test(key)) { - throw new Error(`unsupported object path: ${JSON.stringify(key)}`); - } - files.push({ source: path, path: key, size, contentType: mime.getType(path) ?? "application/octet-stream" }); - }; - const walk = async (directory: string): Promise => { - for await (const entry of await opendir(directory)) { - const path = join(directory, entry.name); - if (entry.isDirectory()) await walk(path); - else if (entry.isFile()) addFile(path, (await lstat(path)).size); - } - }; - if (rootStat.isFile()) addFile(root, rootStat.size); - else await walk(root); - return files; -} - -function retryable(error: unknown): boolean { - return BlobError.is(error) && ( - error.code === "rate_limited" || error.code === "not_ready" || - (error.status !== undefined && error.status >= 500) - ) || error instanceof TypeError && (error.message === "fetch failed" || error.message === "terminated") || - error instanceof Error && error.name === "TimeoutError"; -} - -export async function uploadFiles( - bucket: Pick, - files: UploadFile[], - options: Pick, - progress: (message: string) => void, - signal?: AbortSignal, -): Promise { - const summary: UploadSummary = { uploaded: 0, skipped: 0, bytes: 0, failed: [], remaining: files.length }; - let next = 0; - const worker = async (): Promise => { - while (summary.failed.length === 0 && !signal?.aborted) { - const file = files[next++]; - if (!file) return; - try { - if (options.skipExisting && await bucket.exists(file.path)) { - summary.skipped++; - progress(`Skipped ${JSON.stringify(file.path)}`); - } else { - progress(`Uploading ${JSON.stringify(file.path)} (${file.size} bytes)`); - for (let attempt = 0; ; attempt++) { - if (signal?.aborted) throw new Error("upload interrupted"); - const current = await lstat(file.source); - if (!current.isFile() || current.size !== file.size) { - throw new Error("source changed since scanning; rerun the command"); - } - const stream = createReadStream(file.source, { signal, highWaterMark: 64 * 1024 }); - try { - const body = Readable.toWeb(stream, { - strategy: { highWaterMark: 64 * 1024, size: (chunk: Buffer) => chunk.byteLength }, - }) as ReadableStream; - await bucket.put(file.path, body, { - size: file.size, - contentType: file.contentType, - }); - break; - } catch (error) { - if (attempt >= 2 || signal?.aborted || !retryable(error)) throw error; - } finally { - stream.destroy(); - } - await sleep(500 * 2 ** attempt); - } - summary.uploaded++; - summary.bytes += file.size; - progress(`Uploaded ${JSON.stringify(file.path)}`); - } - } catch (error) { - const message = error instanceof Error ? error.message : "upload failed"; - summary.failed.push({ path: file.path, error: message }); - progress(`Failed ${JSON.stringify(file.path)}: ${message}`); - } finally { - summary.remaining--; - } - } - }; - await Promise.all(Array.from({ length: options.concurrency }, () => worker())); - return summary; -} - -export function registerBlobUpload(blob: Command): void { - blob.command("upload ") - .description("Upload a file or directory with automatic credential refresh") - .option("--bucket-id ", "Blob bucket ID (otherwise uses UPSTASH_BLOB_TOKEN)") - .option("--token ", "Blob bucket token; no management API key needed (overrides UPSTASH_BLOB_TOKEN)") - .option("--prefix ", "Destination prefix; directories upload their contents", "") - .option("--concurrency ", "Number of files uploaded concurrently (1-16)", concurrency, 4) - .option("--skip-existing", "Skip keys already present, without comparing contents") - .option("--dry-run", "List local files and destination paths without contacting Upstash") - .option("--quiet", "Suppress progress on stderr; still print the JSON summary") - .addHelpText("after", ` -Directories are recursive. Symlinks and empty directories are skipped. -Existing keys are overwritten unless --skip-existing is set. Nothing is deleted. -Large files use multipart uploads; credentials refresh during the transfer. -An interrupted file restarts on rerun. --skip-existing skips completed keys, -so use it only when existing objects are already the versions you want. -`) - .action(async (source: string, options: UploadOptions, command: Command) => { - const files = await planUpload(source, options.prefix); - if (options.dryRun) { - printJSON({ dry_run: true, files, bytes: files.reduce((sum, file) => sum + file.size, 0) }); - return; - } - if (files.length === 0) throw new Error("source contains no regular files to upload"); - const { token, unauthorizedRetries } = await resolveBucketToken(options, command); - // A bucket created moments ago answers 401 until provisioning finishes. - if (unauthorizedRetries > 0) await fetchBlobCredentials(token, sleep, { unauthorizedRetries }); - const bucket = new Bucket({ token, enableTelemetry: telemetryStatus().enabled }); - const controller = new AbortController(); - // Each read stream adds an abort listener, removed only once it closes. - setMaxListeners(0, controller.signal); - const interrupt = (): void => controller.abort(); - process.once("SIGINT", interrupt); - process.once("SIGTERM", interrupt); - try { - const summary = await uploadFiles(bucket, files, options, (message) => { - if (!options.quiet) console.error(message); - }, controller.signal); - printJSON(summary); - if (controller.signal.aborted) throw new Error("upload interrupted; completed objects remain in the bucket"); - if (summary.failed.length > 0) throw new Error("upload incomplete; see failed paths in the JSON summary"); - } finally { - process.removeListener("SIGINT", interrupt); - process.removeListener("SIGTERM", interrupt); - } - }); -} diff --git a/tests/unit/blob-s3.test.ts b/tests/unit/blob-s3.test.ts new file mode 100644 index 0000000..644f9a7 --- /dev/null +++ b/tests/unit/blob-s3.test.ts @@ -0,0 +1,265 @@ +import { Command } from "commander"; +import { randomUUID } from "node:crypto"; +import { mkdtemp, mkdir, rm, symlink, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { BucketResolver, tokenBucketId } from "../../src/commands/blob/buckets.js"; +import { needsSync } from "../../src/commands/blob/sync.js"; +import { + checkOverwritesSource, + claimLocal, + destinationFor, + dirPrefix, + globToRegExp, + isIncluded, + listBlobs, + localRel, + nestedPrefix, + parseLocation, +} from "../../src/commands/blob/transfer.js"; +import type { Entry } from "../../src/commands/blob/transfer.js"; +import { createBlobProgram, runCommand } from "../helpers/program.js"; + +let directory: string; +const originalEnv = { ...process.env }; + +beforeEach(async () => { + directory = await mkdtemp(join(tmpdir(), "blob-s3-test-")); + delete process.env.UPSTASH_BLOB_TOKEN; + delete process.env.UPSTASH_EMAIL; + delete process.env.UPSTASH_API_KEY; + process.env.UPSTASH_CONFIG_HOME = directory; + process.env.UPSTASH_LEGACY_CONFIG_HOME = directory; + vi.spyOn(globalThis, "fetch").mockRejectedValue(new Error("Unexpected network request in offline test")); +}); + +afterEach(async () => { + vi.restoreAllMocks(); + process.env = { ...originalEnv }; + await rm(directory, { recursive: true, force: true }); +}); + +function token(id: string = randomUUID()): string { + const bucket = Buffer.from(id); + const password = Buffer.from("fixture-password"); + const hash = Buffer.from("fixture"); + return Buffer.concat([Buffer.from([2, 0, bucket.length, 0, password.length, hash.length]), bucket, password, hash]).toString("base64url"); +} + +function entry(rel: string, size = 1, mtime = 0): Entry { + return { rel, size, mtime, location: { type: "local", path: rel } }; +} + +describe("locations", () => { + it("parses blob URIs and local paths", () => { + expect(parseLocation("blob://bucket/a/b.txt")).toEqual({ type: "blob", bucket: "bucket", key: "a/b.txt" }); + expect(parseLocation("blob://bucket")).toEqual({ type: "blob", bucket: "bucket", key: "" }); + expect(parseLocation("./dir")).toEqual({ type: "local", path: "./dir" }); + expect(() => parseLocation("s3://bucket/key")).toThrow("unsupported location"); + expect(() => parseLocation("blob:///key")).toThrow("has no bucket"); + }); + + it("treats recursive sources as directories", () => { + expect(dirPrefix("")).toBe(""); + expect(dirPrefix("images")).toBe("images/"); + expect(dirPrefix("images/")).toBe("images/"); + }); + + it("maps entries to destinations like aws s3", () => { + const blob = { type: "blob" as const, bucket: "b", key: "site" }; + expect(destinationFor(entry("sub/a.txt"), blob, true)).toEqual({ type: "blob", bucket: "b", key: "site/sub/a.txt" }); + expect(destinationFor(entry("a.txt"), blob, false)).toEqual(blob); + expect(destinationFor(entry("a.txt"), { ...blob, key: "" }, true)).toMatchObject({ key: "a.txt" }); + expect(destinationFor(entry("sub/a.txt"), { type: "local", path: "out" }, true)).toEqual({ type: "local", path: join("out", "sub", "a.txt") }); + }); + + it("writes into a filesystem root", () => { + expect(destinationFor(entry("a/b.txt"), { type: "local", path: "/" }, true)).toEqual({ type: "local", path: join("/", "a", "b.txt") }); + }); + + it("refuses keys that would escape the local destination", () => { + for (const rel of ["../escape.txt", "a/../../escape.txt", ".."]) { + expect(() => destinationFor(entry(rel), { type: "local", path: directory }, true)).toThrow("outside"); + } + }); +}); + +describe("local collisions", () => { + it("normalizes keys the way the local listing reports them", () => { + expect(localRel("/a.txt")).toBe("a.txt"); + expect(localRel("a//b/./c")).toBe("a/b/c"); + expect(localRel("a/b")).toBe("a/b"); + }); + + it("lets only one key write a local file, ignoring case where the disk does", async () => { + const claimed = new Set(); + await claimLocal(claimed, { type: "local", path: join(directory, "README.md") }); + await expect(claimLocal(claimed, { type: "local", path: join(directory, "sub", "..", "README.md") })).rejects.toThrow("also written"); + const lower = claimLocal(claimed, { type: "local", path: join(directory, "readme.md") }); + if (process.platform === "linux") await expect(lower).resolves.toBeUndefined(); + else await expect(lower).rejects.toThrow("also written"); + await expect(claimLocal(claimed, { type: "blob", bucket: "b", key: "README.md" })).resolves.toBeUndefined(); + const folded = new Set(); + await claimLocal(folded, { type: "local", path: join(directory, "README.md") }, true); + await expect(claimLocal(folded, { type: "local", path: join(directory, "Readme.MD") }, true)).rejects.toThrow("also written"); + }); + + it("sees a file through a symbolic link as the same file", async () => { + await mkdir(join(directory, "real")); + await symlink(join(directory, "real"), join(directory, "link")); + const claimed = new Set(); + await claimLocal(claimed, { type: "local", path: join(directory, "real", "x") }); + await expect(claimLocal(claimed, { type: "local", path: join(directory, "link", "x") })).rejects.toThrow("also written"); + }); + + it("finds a destination prefix nested in the source within one bucket", async () => { + const resolver = new BucketResolver(new Command()); + const at = (key: string) => ({ type: "blob" as const, bucket: "b", key }); + await expect(nestedPrefix(at(""), at("archive"), resolver)).resolves.toBe("archive/"); + await expect(nestedPrefix(at("a/"), at("a/b/"), resolver)).resolves.toBe("a/b/"); + await expect(nestedPrefix(at("archive"), at(""), resolver)).resolves.toBeUndefined(); + await expect(nestedPrefix(at("a"), at("a"), resolver)).resolves.toBeUndefined(); + await expect(nestedPrefix(at("a"), at("ab"), resolver)).resolves.toBeUndefined(); + await expect(nestedPrefix({ type: "local", path: "." }, at("x"), resolver)).resolves.toBeUndefined(); + }); + + it("refuses copies onto keys below the source prefix", () => { + const at = (key: string) => ({ type: "blob" as const, bucket: "b", key }); + expect(() => checkOverwritesSource(at("x/1"), "x/")).toThrow("also a source"); + expect(() => checkOverwritesSource(at("1"), "x/")).not.toThrow(); + expect(() => checkOverwritesSource(at("x/1"), undefined)).not.toThrow(); + }); +}); + +describe("folder markers", () => { + const listing = { + list: async () => ({ + cursor: undefined, + blobs: [ + { path: "tmp/", size: 0, etag: "", uploadedAt: new Date(0) }, + { path: "tmp/sub/", size: 0, etag: "", uploadedAt: new Date(0) }, + { path: "tmp/x.txt", size: 1, etag: "", uploadedAt: new Date(0) }, + ], + }), + } as unknown as Parameters[0]; + const location = { type: "blob" as const, bucket: "b", key: "tmp" }; + + it("skips zero-byte markers except when deleting", async () => { + expect((await listBlobs(listing, location, "tmp/")).map((e) => e.rel)).toEqual(["x.txt"]); + expect((await listBlobs(listing, location, "tmp/", true)).map((e) => e.rel)).toEqual(["", "sub/", "x.txt"]); + }); +}); + +describe("filters", () => { + it("matches fnmatch patterns where * crosses slashes", () => { + expect(globToRegExp("*.txt").test("dir/a.txt")).toBe(true); + expect(globToRegExp("a?c").test("abc")).toBe(true); + expect(globToRegExp("[!a]b").test("ab")).toBe(false); + expect(globToRegExp("[!a]b").test("cb")).toBe(true); + expect(globToRegExp("[]]x").test("]x")).toBe(true); + expect(globToRegExp("a.b(c)").test("a.b(c)")).toBe(true); + expect(globToRegExp("a.b").test("axb")).toBe(false); + }); + + it("lets the last matching filter win", () => { + const filters = [ + { exclude: true, match: globToRegExp("*") }, + { exclude: false, match: globToRegExp("*.json") }, + ]; + expect(isIncluded("x/a.json", filters)).toBe(true); + expect(isIncluded("x/a.txt", filters)).toBe(false); + expect(isIncluded("anything", [])).toBe(true); + }); + + it("keeps --exclude and --include in command-line order", async () => { + await writeFile(join(directory, "a.json"), "{}"); + await writeFile(join(directory, "b.txt"), "b"); + const result = await runCommand(await createBlobProgram(), [ + "blob", "cp", directory, "blob://bucket/dest", "--recursive", + "--exclude", "*", "--include", "*.json", "--dryrun", "--quiet", + ]); + expect(result).toEqual({ + dry_run: true, + operations: [{ + action: "upload", + source: join(directory, "a.json"), + destination: "blob://bucket/dest/a.json", + size: 2, + }], + }); + expect(fetch).not.toHaveBeenCalled(); + }); +}); + +describe("local listing", () => { + it("follows symbolic links like aws s3, skipping loops and dangling links", async () => { + const source = join(directory, "src"); + await mkdir(join(directory, "real"), { recursive: true }); + await mkdir(source); + await writeFile(join(directory, "real", "a.txt"), "a"); + await symlink(join(directory, "real"), join(source, "linked")); + await symlink(source, join(source, "loop")); + await symlink(join(directory, "missing"), join(source, "dangling")); + const result = await runCommand(await createBlobProgram(), [ + "blob", "cp", source, "blob://bucket/dest", "--recursive", "--dryrun", "--quiet", + ]) as { operations: { destination: string }[] }; + expect(result.operations.map((operation) => operation.destination)).toEqual(["blob://bucket/dest/linked/a.txt"]); + }); +}); + +describe("sync comparison", () => { + const at = (seconds: number, size = 1): Entry => entry("a", size, seconds * 1000); + + it("copies uploads and copies when the size differs or the source is newer", () => { + expect(needsSync(at(10, 1), at(10, 2), "upload", {})).toBe(true); + expect(needsSync(at(11), at(10), "upload", {})).toBe(true); + expect(needsSync(at(10), at(11), "copy", {})).toBe(false); + expect(needsSync(at(10.9), at(10.1), "upload", {})).toBe(false); + expect(needsSync(at(11), at(10), "upload", { sizeOnly: true })).toBe(false); + }); + + it("skips same-sized downloads unless the local file is newer, as aws does", () => { + expect(needsSync(at(11), at(10), "download", {})).toBe(false); + expect(needsSync(at(10), at(11), "download", {})).toBe(true); + expect(needsSync(at(11), at(10), "download", { exactTimestamps: true })).toBe(true); + expect(needsSync(at(10), at(10), "download", { exactTimestamps: true })).toBe(false); + }); +}); + +describe("bucket resolution", () => { + it("reads the bucket id from a token", () => { + const id = randomUUID(); + expect(tokenBucketId(token(id))).toBe(id); + expect(tokenBucketId("x")).toBeUndefined(); + }); + + it("uses a token only for the bucket it was issued for", async () => { + const id = randomUUID(); + process.env.UPSTASH_BLOB_TOKEN = token(id); + const resolver = new BucketResolver(new Command()); + await expect(resolver.open(id)).resolves.toBeDefined(); + await expect(resolver.open("some-name")).rejects.toThrow(`The Blob token is for bucket ${id}, not "some-name"`); + expect(fetch).not.toHaveBeenCalled(); + }); + + it("prefers the explicit token over the environment", async () => { + const id = randomUUID(); + process.env.UPSTASH_BLOB_TOKEN = token(); + await expect(new BucketResolver(new Command(), token(id)).open(id)).resolves.toBeDefined(); + }); +}); + +describe("argument checks", () => { + it("rejects local-to-local copies and directories without --recursive", async () => { + await mkdir(join(directory, "d")); + const program = await createBlobProgram(); + await expect(runCommand(program, ["blob", "cp", "a", "b"])).rejects.toThrow("must be a blob://"); + await expect(runCommand(await createBlobProgram(), ["blob", "cp", join(directory, "d"), "blob://b/x"])) + .rejects.toThrow("is a directory; use --recursive"); + await expect(runCommand(await createBlobProgram(), ["blob", "rm", "blob://b"])).rejects.toThrow("names no object"); + await expect(runCommand(await createBlobProgram(), ["blob", "presign", "blob://b/k", "--expires-in", "0"])) + .rejects.toThrow(); + await expect(runCommand(await createBlobProgram(), ["blob", "mb", "blob://b/key"])).rejects.toThrow("includes a key"); + }); +}); diff --git a/tests/unit/blob-transfer.test.ts b/tests/unit/blob-transfer.test.ts new file mode 100644 index 0000000..75ef700 --- /dev/null +++ b/tests/unit/blob-transfer.test.ts @@ -0,0 +1,253 @@ +import { Bucket, BlobError } from "@upstash/blob"; +import { Command } from "commander"; +import { mkdir, mkdtemp, rm, symlink, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { randomUUID } from "node:crypto"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { BucketResolver } from "../../src/commands/blob/buckets.js"; +import { parseBlobLocation, parseBucket, putFile } from "../../src/commands/blob/transfer.js"; +import { createBlobProgram, runCommand } from "../helpers/program.js"; + +let directory: string; +const originalEnv = { ...process.env }; + +beforeEach(async () => { + directory = await mkdtemp(join(tmpdir(), "blob-transfer-test-")); + delete process.env.UPSTASH_BLOB_TOKEN; + delete process.env.UPSTASH_EMAIL; + delete process.env.UPSTASH_API_KEY; + process.env.UPSTASH_CONFIG_HOME = directory; + process.env.UPSTASH_LEGACY_CONFIG_HOME = directory; + vi.spyOn(globalThis, "fetch").mockRejectedValue(new Error("Unexpected network request in offline test")); +}); + +afterEach(async () => { + vi.restoreAllMocks(); + vi.useRealTimers(); + process.env = { ...originalEnv }; + await rm(directory, { recursive: true, force: true }); +}); + +function token(id: string = randomUUID()): string { + const bucket = Buffer.from(id); + const password = Buffer.from("fixture-password"); + const hash = Buffer.from("fixture"); + return Buffer.concat([Buffer.from([2, 0, bucket.length, 0, password.length, hash.length]), bucket, password, hash]).toString("base64url"); +} + +function storageCredentials(session = "session", expiresAt = Date.now() / 1000 + 600): Response { + return Response.json({ + accessKeyId: "key", secretAccessKey: "secret", sessionToken: session, expiresAt, + endpoint: "https://fixture.r2.cloudflarestorage.com", bucket: "fixture-bucket", region: "auto", + }); +} + +async function consume(body: Parameters[1]): Promise { + if (!(body instanceof ReadableStream)) throw new Error("expected a stream"); + const reader = body.getReader(); + while (!(await reader.read()).done) { /* drain the fixture stream */ } +} + +async function file(name = "file.txt", content: string | Buffer = "hello") { + const source = join(directory, name); + await writeFile(source, content); + return { source, key: name, size: Buffer.byteLength(content), contentType: "text/plain" }; +} + +describe("bucket paths", () => { + it("accepts a bucket path with or without blob://", () => { + expect(parseBlobLocation("my-bucket/a/b.txt")).toEqual({ type: "blob", bucket: "my-bucket", key: "a/b.txt" }); + expect(parseBlobLocation("my-bucket")).toEqual({ type: "blob", bucket: "my-bucket", key: "" }); + expect(parseBlobLocation("blob://my-bucket/a")).toEqual({ type: "blob", bucket: "my-bucket", key: "a" }); + expect(parseBlobLocation("my-bucket/a://b")).toEqual({ type: "blob", bucket: "my-bucket", key: "a://b" }); + expect(() => parseBlobLocation("s3://my-bucket/a")).toThrow("unsupported location"); + expect(() => parseBlobLocation("/a")).toThrow("has no bucket"); + expect(parseBucket("my-bucket/")).toBe("my-bucket"); + expect(() => parseBucket("my-bucket/key")).toThrow("includes a key"); + }); + + it("takes short flags", async () => { + await writeFile(join(directory, "a.txt"), "a"); + const result = await runCommand(await createBlobProgram(), ["blob", "cp", directory, "blob://bucket/dest", "-r", "-n", "-q"]); + expect(result).toMatchObject({ dry_run: true, operations: [{ destination: "blob://bucket/dest/a.txt" }] }); + const program = await createBlobProgram(); + const blob = program.commands.find((command) => command.name() === "blob")!; + const flags = (name: string) => blob.commands.find((command) => command.name() === name)!.options.map((option) => option.short); + expect(flags("sync")).toEqual(expect.arrayContaining(["-d", "-n", "-q"])); + expect(flags("rm")).toEqual(expect.arrayContaining(["-r", "-n", "-q"])); + expect(flags("rb")).toEqual(expect.arrayContaining(["-f", "-q"])); + expect(fetch).not.toHaveBeenCalled(); + }); + + it("checks arguments before any request", async () => { + await expect(runCommand(await createBlobProgram(), ["blob", "cp", directory, "blob://bucket/x", "-r", "--concurrency", "0"])) + .rejects.toThrow("1 to 16"); + await expect(runCommand(await createBlobProgram(), ["blob", "rm", "build", "-r"])).rejects.toThrow("write blob://build"); + expect(fetch).not.toHaveBeenCalled(); + }); +}); + +describe("sync --delete", () => { + it("keeps a local file reached through a symbolic link", async () => { + const dest = join(directory, "dest"); + await mkdir(join(dest, "a"), { recursive: true }); + await writeFile(join(dest, "a", "f.txt"), "hello"); + await writeFile(join(dest, "stale.txt"), "old"); + await symlink(join(dest, "a"), join(dest, "link")); + vi.spyOn(BucketResolver.prototype, "open").mockResolvedValue({ + list: async () => ({ blobs: [{ path: "a/f.txt", size: 5, etag: "", uploadedAt: new Date(Date.now() + 1e7) }] }), + } as unknown as Bucket); + const result = await runCommand(await createBlobProgram(), ["blob", "sync", "blob://b/", dest, "-d", "-n", "-q"]); + expect(result).toEqual({ + dry_run: true, + operations: [{ action: "delete", source: join(dest, "stale.txt"), size: 0 }], + unchanged: 1, + }); + }); +}); + +describe("putFile", () => { + it("reopens the stream when retrying a transient failure", async () => { + let attempts = 0; + const bodies: unknown[] = []; + const bucket = { + put: vi.fn(async (_key, body) => { + bodies.push(body); + await consume(body); + if (++attempts === 1) throw new BlobError("request_failed", { status: 503 }); + return {}; + }), + } as unknown as Pick; + await putFile(bucket, await file()); + expect(attempts).toBe(2); + expect(bodies[0]).not.toBe(bodies[1]); + }); + + it("retries a credential request that timed out, but not a rejected one", async () => { + let attempts = 0; + const flaky = { + put: vi.fn(async () => { + if (++attempts === 1) throw new DOMException("The operation was aborted due to timeout", "TimeoutError"); + return {}; + }), + } as unknown as Pick; + await putFile(flaky, await file()); + expect(attempts).toBe(2); + const denied = { put: vi.fn().mockRejectedValue(new BlobError("unauthorized")) }; + await expect(putFile(denied, await file())).rejects.toThrow(); + expect(denied.put).toHaveBeenCalledTimes(1); + }); + + it("refreshes credentials between multipart parts after the original credentials expire", async () => { + const size = 17 * 1024 * 1024; + const large = await file("large.bin", Buffer.alloc(size, 7)); + let now = Date.now(); + vi.spyOn(Date, "now").mockImplementation(() => now); + let mints = 0; + const parts: { session: string | null; size: number }[] = []; + let completed = false; + vi.mocked(fetch).mockImplementation(async (input, init) => { + const url = new URL(String(input)); + if (url.hostname === "blob.upstash.io") return storageCredentials(`session-${++mints}`, now / 1000 + 600); + if (url.hostname !== "fixture.r2.cloudflarestorage.com") throw new Error("Unexpected host"); + if (url.searchParams.has("uploads")) return new Response("fixture-upload"); + if (url.searchParams.has("partNumber")) { + expect(init?.body).toBeInstanceOf(Uint8Array); + parts.push({ session: new Headers(init?.headers).get("x-amz-security-token"), size: (init!.body as Uint8Array).byteLength }); + if (parts.length === 1) now += 11 * 60 * 1000; + return new Response(null, { headers: { etag: `"part-${parts.length}"` } }); + } + if (init?.method === "POST" && url.searchParams.has("uploadId")) { + completed = true; + expect(new Headers(init.headers).get("x-amz-security-token")).toBe("session-2"); + return new Response('"complete"'); + } + throw new Error(`Unexpected offline request: ${init?.method} ${url.pathname}`); + }); + await putFile(new Bucket({ token: token(), enableTelemetry: false }), large); + expect(mints).toBe(2); + expect(parts[0]!.session).toBe("session-1"); + expect(parts.slice(1).every((part) => part.session === "session-2")).toBe(true); + expect(parts.reduce((sum, part) => sum + part.size, 0)).toBe(size); + expect(completed).toBe(true); + }); + + it("aborts a failed multipart upload instead of reporting success", async () => { + const large = await file("large.bin", Buffer.alloc(17 * 1024 * 1024)); + let aborted = false; + vi.mocked(fetch).mockImplementation(async (input, init) => { + const url = new URL(String(input)); + if (url.hostname === "blob.upstash.io") return storageCredentials(); + if (url.searchParams.has("uploads")) return new Response("fixture-upload"); + if (init?.method === "DELETE") { aborted = true; return new Response(null, { status: 204 }); } + return new Response("AccessDenied", { status: 403 }); + }); + await expect(putFile(new Bucket({ token: token(), enableTelemetry: false }), large)).rejects.toThrow(); + expect(aborted).toBe(true); + }); +}); + +describe("cp with the real SDK and offline storage", () => { + it.each(["environment", "flag"])("uploads with a %s bucket token and no management credentials", async (source) => { + await writeFile(join(directory, "hello.txt"), "hello"); + const id = randomUUID(); + const bucketToken = token(id); + process.env.UPSTASH_BLOB_TOKEN = source === "environment" ? bucketToken : token(); + const uploaded: string[] = []; + vi.mocked(fetch).mockImplementation(async (input, init) => { + const url = new URL(String(input)); + if (url.hostname === "blob.upstash.io") { + expect(new Headers(init?.headers).get("authorization")).toBe(`Bearer ${bucketToken}`); + return storageCredentials(); + } + if (url.hostname !== "fixture.r2.cloudflarestorage.com") throw new Error("Unexpected host"); + expect(new Headers(init?.headers).get("content-type")).toBe("text/plain"); + expect(init?.method).toBe("PUT"); + expect(await new Response(init?.body).text()).toBe("hello"); + uploaded.push(url.pathname); + return new Response(null, { headers: { etag: '"hello"' } }); + }); + const flags = source === "flag" ? ["--token", bucketToken] : []; + const result = await runCommand(await createBlobProgram(), ["blob", "cp", directory, `blob://${id}/assets`, "-r", "-q", ...flags]); + expect(result).toEqual({ completed: 1, bytes: 5, failed: [], remaining: 0 }); + expect(uploaded).toEqual(["/fixture-bucket/assets/hello.txt"]); + }); + + it("waits for a freshly created bucket, addressed by name, to finish provisioning", async () => { + await writeFile(join(directory, "hello.txt"), "hello"); + process.env.UPSTASH_EMAIL = "user@example.com"; + process.env.UPSTASH_API_KEY = "api-key"; + const bucketToken = token(); + let mints = 0; + vi.mocked(fetch).mockImplementation(async (input) => { + const url = new URL(String(input)); + if (url.pathname.endsWith("/v2/blob/bucket")) return Response.json([{ id: "bucket_123", name: "fresh" }]); + if (url.pathname.endsWith("/v2/blob/bucket/bucket_123")) { + return Response.json({ id: "bucket_123", name: "fresh", token: bucketToken, creation_time: Date.now() / 1000 }); + } + if (url.hostname === "blob.upstash.io") { + if (++mints === 1) return new Response('{"error":"unauthorized"}', { status: 401 }); + return storageCredentials(); + } + if (url.hostname !== "fixture.r2.cloudflarestorage.com") throw new Error("Unexpected host"); + return new Response(null, { headers: { etag: '"hello"' } }); + }); + vi.useFakeTimers({ toFake: ["setTimeout"] }); + const pending = runCommand(await createBlobProgram(), ["blob", "cp", join(directory, "hello.txt"), "blob://fresh/", "-q"]); + await vi.waitFor(() => expect(mints).toBe(1)); + await vi.advanceTimersByTimeAsync(3000); + expect(await pending).toEqual({ completed: 1, bytes: 5, failed: [], remaining: 0 }); + expect(mints).toBeGreaterThan(1); + }); + + it("rejects an empty explicit token instead of falling back to the environment", async () => { + await writeFile(join(directory, "hello.txt"), "hello"); + const id = randomUUID(); + process.env.UPSTASH_BLOB_TOKEN = token(id); + await expect(runCommand(await createBlobProgram(), ["blob", "cp", join(directory, "hello.txt"), `blob://${id}/`, "--token", " "])) + .rejects.toThrow("1 of 1 operations failed"); + await expect(new BucketResolver(new Command(), " ").open(id)).rejects.toThrow("--token must be a non-empty"); + expect(fetch).not.toHaveBeenCalled(); + }); +}); diff --git a/tests/unit/blob-upload.test.ts b/tests/unit/blob-upload.test.ts deleted file mode 100644 index e43a435..0000000 --- a/tests/unit/blob-upload.test.ts +++ /dev/null @@ -1,309 +0,0 @@ -import { Bucket, BlobError } from "@upstash/blob"; -import { mkdtemp, mkdir, rm, symlink, writeFile } from "node:fs/promises"; -import { tmpdir } from "node:os"; -import { join } from "node:path"; -import { randomUUID } from "node:crypto"; -import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; -import { planUpload, uploadFiles } from "../../src/commands/blob/upload.js"; -import { createBlobProgram, runCommand } from "../helpers/program.js"; - -let directory: string; -const originalEnv = { ...process.env }; - -beforeEach(async () => { - directory = await mkdtemp(join(tmpdir(), "blob-upload-test-")); - delete process.env.UPSTASH_BLOB_TOKEN; - delete process.env.UPSTASH_EMAIL; - delete process.env.UPSTASH_API_KEY; - vi.spyOn(globalThis, "fetch").mockRejectedValue(new Error("Unexpected network request in offline test")); -}); - -afterEach(async () => { - vi.restoreAllMocks(); - vi.useRealTimers(); - process.env = { ...originalEnv }; - await rm(directory, { recursive: true, force: true }); -}); - -function token(): string { - const id = Buffer.from(randomUUID()); - const password = Buffer.from("fixture-password"); - const hash = Buffer.from("fixture"); - return Buffer.concat([Buffer.from([2, 0, id.length, 0, password.length, hash.length]), id, password, hash]).toString("base64url"); -} - -async function consume(body: Parameters[1]): Promise { - if (!(body instanceof ReadableStream)) throw new Error("expected a stream"); - const reader = body.getReader(); - while (!(await reader.read()).done) { /* drain the fixture stream */ } -} - -describe("upload planning", () => { - it("preserves nested paths and MIME types while skipping symlinks and empty folders", async () => { - await mkdir(join(directory, "images")); - await mkdir(join(directory, "empty")); - await writeFile(join(directory, "images", "photo.png"), "png"); - await writeFile(join(directory, "app.css"), "body {}"); - await symlink(join(directory, "images"), join(directory, "linked")); - const files = await planUpload(directory, "/assets/"); - expect(files.map(({ path, contentType }) => ({ path, contentType })).sort((a, b) => a.path.localeCompare(b.path))).toEqual([ - { path: "assets/app.css", contentType: "text/css" }, - { path: "assets/images/photo.png", contentType: "image/png" }, - ]); - }); - - it("accepts a single file and rejects a symlink source", async () => { - const file = join(directory, "one.txt"); - await writeFile(file, "hello"); - expect(await planUpload(file, "docs")).toEqual([ - { source: file, path: "docs/one.txt", size: 5, contentType: "text/plain" }, - ]); - await symlink(file, join(directory, "link")); - await expect(planUpload(join(directory, "link"), "")).rejects.toThrow("symbolic links"); - }); - - it("rejects unsafe prefixes before authentication", async () => { - for (const prefix of ["../outside", "assets/../secret", "a\\b", "a\nb"]) { - await expect(planUpload(directory, prefix)).rejects.toThrow("prefix"); - } - expect(fetch).not.toHaveBeenCalled(); - }); - - it("dry-run works without credentials or network requests", async () => { - await writeFile(join(directory, "file.txt"), "hello"); - const result = await runCommand(await createBlobProgram(), ["blob", "upload", directory, "--prefix", "assets", "--dry-run"]); - expect(result).toMatchObject({ dry_run: true, bytes: 5, files: [{ path: "assets/file.txt" }] }); - expect(fetch).not.toHaveBeenCalled(); - }); - - it("validates concurrency before authentication", async () => { - const program = await createBlobProgram(); - program.configureOutput({ writeErr: () => {} }); - await expect(runCommand(program, ["blob", "upload", directory, "--concurrency", "0"])).rejects.toThrow("concurrency"); - expect(fetch).not.toHaveBeenCalled(); - }); -}); - -describe("upload scheduling", () => { - it("streams bytes, respects concurrency, and reports confirmed uploads", async () => { - for (let i = 0; i < 6; i++) await writeFile(join(directory, `${i}.txt`), "hello"); - let active = 0; - let maximum = 0; - const bucket = { - exists: vi.fn(), - put: vi.fn(async (_path, body) => { - active++; - maximum = Math.max(maximum, active); - await consume(body); - active--; - return {}; - }), - } as unknown as Pick; - const summary = await uploadFiles(bucket, await planUpload(directory, ""), { concurrency: 2 }, () => {}); - expect(summary).toEqual({ uploaded: 6, skipped: 0, bytes: 30, failed: [], remaining: 0 }); - expect(maximum).toBeLessThanOrEqual(2); - expect(bucket.exists).not.toHaveBeenCalled(); - }); - - it("skip-existing does not open or overwrite existing objects", async () => { - const file = join(directory, "file.txt"); - await writeFile(file, "hello"); - const files = await planUpload(directory, ""); - await rm(file); - const bucket = { exists: vi.fn().mockResolvedValue(true), put: vi.fn() }; - expect(await uploadFiles(bucket, files, { concurrency: 1, skipExisting: true }, () => {})).toEqual({ - uploaded: 0, skipped: 1, bytes: 0, failed: [], remaining: 0, - }); - expect(bucket.put).not.toHaveBeenCalled(); - }); - - it("reports failures and stops scheduling remaining files", async () => { - await writeFile(join(directory, "a.txt"), "hello"); - await writeFile(join(directory, "b.txt"), "hello"); - const files = await planUpload(directory, ""); - const bucket = { exists: vi.fn(), put: vi.fn().mockRejectedValue(new BlobError("unauthorized")) }; - const summary = await uploadFiles(bucket, files, { concurrency: 1 }, () => {}); - expect(summary).toMatchObject({ uploaded: 0, remaining: 1, failed: [{ path: files[0]!.path }] }); - expect(bucket.put).toHaveBeenCalledTimes(1); - }); - - it("reopens a stream when retrying a transient failure", async () => { - await writeFile(join(directory, "file.txt"), "hello"); - let attempts = 0; - const bodies: unknown[] = []; - const bucket = { - exists: vi.fn(), - put: vi.fn(async (_path, body) => { - bodies.push(body); - await consume(body); - if (++attempts === 1) throw new BlobError("request_failed", { status: 503 }); - return {}; - }), - } as unknown as Pick; - const summary = await uploadFiles(bucket, await planUpload(directory, ""), { concurrency: 1 }, () => {}); - expect(summary.uploaded).toBe(1); - expect(attempts).toBe(2); - expect(bodies[0]).not.toBe(bodies[1]); - }); - - it("retries a credential request that timed out", async () => { - await writeFile(join(directory, "file.txt"), "hello"); - let attempts = 0; - const bucket = { - exists: vi.fn(), - put: vi.fn(async () => { - if (++attempts === 1) throw new DOMException("The operation was aborted due to timeout", "TimeoutError"); - return {}; - }), - } as unknown as Pick; - const summary = await uploadFiles(bucket, await planUpload(directory, ""), { concurrency: 1 }, () => {}); - expect(summary.uploaded).toBe(1); - expect(attempts).toBe(2); - }); - - it("does not start work after cancellation", async () => { - await writeFile(join(directory, "file.txt"), "hello"); - const controller = new AbortController(); - controller.abort(); - const bucket = { exists: vi.fn(), put: vi.fn() }; - const summary = await uploadFiles(bucket, await planUpload(directory, ""), { concurrency: 1 }, () => {}, controller.signal); - expect(summary.remaining).toBe(1); - expect(bucket.put).not.toHaveBeenCalled(); - }); -}); - -describe("real SDK with offline storage transport", () => { - it.each(["environment", "flag"])("uploads with a %s bucket token and no management credentials", async (source) => { - await writeFile(join(directory, "hello.txt"), "hello"); - const bucketToken = token(); - process.env.UPSTASH_BLOB_TOKEN = source === "environment" ? bucketToken : "ignored-ambient-token"; - const uploaded: string[] = []; - vi.mocked(fetch).mockImplementation(async (input, init) => { - const url = new URL(String(input)); - if (url.hostname === "blob.upstash.io") { - expect(new Headers(init?.headers).get("authorization")).toBe(`Bearer ${bucketToken}`); - return Response.json({ - accessKeyId: "key", secretAccessKey: "secret", sessionToken: "session", expiresAt: Date.now() / 1000 + 600, - endpoint: "https://fixture.r2.cloudflarestorage.com", bucket: "fixture-bucket", region: "auto", - }); - } - if (url.hostname !== "fixture.r2.cloudflarestorage.com") throw new Error("Unexpected host"); - expect(new Headers(init?.headers).get("content-type")).toBe("text/plain"); - expect(init?.method).toBe("PUT"); - expect(await new Response(init?.body).text()).toBe("hello"); - uploaded.push(url.pathname); - return new Response(null, { headers: { etag: '"hello"' } }); - }); - const flags = source === "flag" ? ["--token", bucketToken] : []; - const result = await runCommand(await createBlobProgram(), ["blob", "upload", directory, "--prefix", "assets", "--quiet", ...flags]); - expect(result).toEqual({ uploaded: 1, skipped: 0, bytes: 5, failed: [], remaining: 0 }); - expect(uploaded).toEqual(["/fixture-bucket/assets/hello.txt"]); - }); - - it("waits for a freshly created bucket to finish provisioning", async () => { - await writeFile(join(directory, "hello.txt"), "hello"); - process.env.UPSTASH_EMAIL = "user@example.com"; - process.env.UPSTASH_API_KEY = "api-key"; - const bucketToken = token(); - let mints = 0; - vi.mocked(fetch).mockImplementation(async (input) => { - const url = new URL(String(input)); - if (url.pathname.endsWith("/v2/blob/bucket/bucket_123")) { - return Response.json({ id: "bucket_123", token: bucketToken, creation_time: Date.now() / 1000 }); - } - if (url.hostname === "blob.upstash.io") { - if (++mints === 1) return new Response('{"error":"unauthorized"}', { status: 401 }); - return Response.json({ - accessKeyId: "key", secretAccessKey: "secret", sessionToken: "session", expiresAt: Date.now() / 1000 + 600, - endpoint: "https://fixture.r2.cloudflarestorage.com", bucket: "fixture-bucket", region: "auto", - }); - } - if (url.hostname !== "fixture.r2.cloudflarestorage.com") throw new Error("Unexpected host"); - return new Response(null, { headers: { etag: '"hello"' } }); - }); - vi.useFakeTimers({ toFake: ["setTimeout"] }); - const pending = runCommand(await createBlobProgram(), ["blob", "upload", directory, "--bucket-id", "bucket_123", "--quiet"]); - await vi.waitFor(() => expect(mints).toBe(1)); - await vi.advanceTimersByTimeAsync(3000); - expect(await pending).toEqual({ uploaded: 1, skipped: 0, bytes: 5, failed: [], remaining: 0 }); - expect(mints).toBeGreaterThan(1); - }); - - it("rejects an empty explicit token instead of falling back to ambient credentials", async () => { - await writeFile(join(directory, "hello.txt"), "hello"); - process.env.UPSTASH_BLOB_TOKEN = token(); - await expect(runCommand(await createBlobProgram(), ["blob", "upload", directory, "--token", " "])) - .rejects.toThrow("--token must be a non-empty"); - expect(fetch).not.toHaveBeenCalled(); - }); - - it("rejects conflicting explicit bucket selectors without making a request", async () => { - await writeFile(join(directory, "hello.txt"), "hello"); - await expect(runCommand(await createBlobProgram(), ["blob", "upload", directory, "--token", token(), "--bucket-id", "another-bucket"])) - .rejects.toThrow("Use either --token or --bucket-id"); - expect(fetch).not.toHaveBeenCalled(); - }); - - it("refreshes credentials between multipart parts after the original credentials expire", async () => { - const file = join(directory, "large.bin"); - const size = 17 * 1024 * 1024; - await writeFile(file, Buffer.alloc(size, 7)); - let now = Date.now(); - vi.spyOn(Date, "now").mockImplementation(() => now); - let mints = 0; - const parts: { session: string | null; size: number }[] = []; - let completed = false; - vi.mocked(fetch).mockImplementation(async (input, init) => { - const url = new URL(String(input)); - if (url.hostname === "blob.upstash.io") { - mints++; - return Response.json({ - accessKeyId: "fixture-key", secretAccessKey: "fixture-secret", - sessionToken: `session-${mints}`, expiresAt: now / 1000 + 600, - endpoint: "https://fixture.r2.cloudflarestorage.com", bucket: "fixture-bucket", region: "auto", - }); - } - if (url.hostname !== "fixture.r2.cloudflarestorage.com") throw new Error("Unexpected host"); - if (url.searchParams.has("uploads")) return new Response("fixture-upload"); - if (url.searchParams.has("partNumber")) { - expect(init?.body).toBeInstanceOf(Uint8Array); - parts.push({ session: new Headers(init?.headers).get("x-amz-security-token"), size: (init!.body as Uint8Array).byteLength }); - if (parts.length === 1) now += 11 * 60 * 1000; - return new Response(null, { headers: { etag: `"part-${parts.length}"` } }); - } - if (init?.method === "POST" && url.searchParams.has("uploadId")) { - completed = true; - expect(new Headers(init.headers).get("x-amz-security-token")).toBe("session-2"); - return new Response('"complete"'); - } - throw new Error(`Unexpected offline request: ${init?.method} ${url.pathname}`); - }); - const bucket = new Bucket({ token: token(), enableTelemetry: false }); - const summary = await uploadFiles(bucket, await planUpload(file, ""), { concurrency: 1 }, () => {}); - expect(summary).toEqual({ uploaded: 1, skipped: 0, bytes: size, failed: [], remaining: 0 }); - expect(mints).toBe(2); - expect(parts[0]!.session).toBe("session-1"); - expect(parts.slice(1).every((part) => part.session === "session-2")).toBe(true); - expect(parts.reduce((sum, part) => sum + part.size, 0)).toBe(size); - expect(completed).toBe(true); - }); - - it("aborts a failed multipart upload instead of reporting success", async () => { - await writeFile(join(directory, "large.bin"), Buffer.alloc(17 * 1024 * 1024)); - let aborted = false; - vi.mocked(fetch).mockImplementation(async (input, init) => { - const url = new URL(String(input)); - if (url.hostname === "blob.upstash.io") return Response.json({ - accessKeyId: "key", secretAccessKey: "secret", sessionToken: "session", expiresAt: Date.now() / 1000 + 600, - endpoint: "https://fixture.r2.cloudflarestorage.com", bucket: "fixture-bucket", region: "auto", - }); - if (url.searchParams.has("uploads")) return new Response("fixture-upload"); - if (init?.method === "DELETE") { aborted = true; return new Response(null, { status: 204 }); } - return new Response("AccessDenied", { status: 403 }); - }); - const summary = await uploadFiles(new Bucket({ token: token(), enableTelemetry: false }), await planUpload(directory, ""), { concurrency: 1 }, () => {}); - expect(summary.uploaded).toBe(0); - expect(summary.failed).toHaveLength(1); - expect(aborted).toBe(true); - }); -}); diff --git a/tests/unit/blob.test.ts b/tests/unit/blob.test.ts index e9c79a1..d095572 100644 --- a/tests/unit/blob.test.ts +++ b/tests/unit/blob.test.ts @@ -57,14 +57,20 @@ describe("blob command registration", () => { const blob = program.commands.find((command) => command.name() === "blob"); expect(blob).toBeDefined(); - expect(blob?.description()).toBe("Manage Blob buckets"); + expect(blob?.description()).toBe("Manage Blob buckets and objects"); expect(blob?.commands.map((command) => command.name())).toEqual([ "create", - "list", "get", "delete", "credentials", - "upload", + "ls", + "cp", + "mv", + "rm", + "sync", + "presign", + "mb", + "rb", ]); }); }); @@ -205,6 +211,60 @@ describe("blob CRUD commands", () => { }); }); +describe("bucket arguments", () => { + const listed = (): Response => new Response(JSON.stringify([ + makeBucket({ id: "bucket_other", name: "other", token: undefined, token_next: undefined }), + makeBucket({ token: undefined, token_next: undefined }), + ]), { status: 200 }); + + it("ls lists buckets, and list is the same command", async () => { + const buckets = [makeBucket({ token: undefined, token_next: undefined })]; + vi.spyOn(globalThis, "fetch").mockImplementation(async () => new Response(JSON.stringify(buckets), { status: 200 })); + expect(await runCommand(await createBlobProgram(), ["blob", "ls"])).toEqual(buckets); + expect(await runCommand(await createBlobProgram(), ["blob", "list"])).toEqual(buckets); + }); + + it.each([["my-bucket"], ["blob://my-bucket"], ["bucket_123"], ["my-bucket/"]])("get %s looks the bucket up by name or id", async (name) => { + const bucket = makeBucket(); + const fetchSpy = vi.spyOn(globalThis, "fetch") + .mockResolvedValueOnce(listed()) + .mockResolvedValueOnce(new Response(JSON.stringify(bucket), { status: 200 })); + expect(await runCommand(await createBlobProgram(), ["blob", "get", name])).toEqual(bucket); + expect(fetchSpy.mock.calls[1]?.[0]).toBe("https://api.upstash.com/v2/blob/bucket/bucket_123"); + }); + + it("delete by name deletes that bucket's id", async () => { + const fetchSpy = vi.spyOn(globalThis, "fetch") + .mockResolvedValueOnce(listed()) + .mockResolvedValueOnce(new Response('"OK"', { status: 200 })); + expect(await runCommand(await createBlobProgram(), ["blob", "delete", "my-bucket"])).toEqual({ deleted: true, bucket_id: "bucket_123" }); + expect(fetchSpy.mock.calls[1]?.[0]).toBe("https://api.upstash.com/v2/blob/bucket/bucket_123"); + expect((fetchSpy.mock.calls[1]?.[1] as RequestInit).method).toBe("DELETE"); + }); + + it("credentials by name exchanges that bucket's token", async () => { + const credentials = makeCredentials(); + const fetchSpy = vi.spyOn(globalThis, "fetch") + .mockResolvedValueOnce(listed()) + .mockResolvedValueOnce(new Response(JSON.stringify(makeBucket({ token: "bucket-token" })), { status: 200 })) + .mockResolvedValueOnce(new Response(JSON.stringify(credentials), { status: 200 })); + expect(await runCommand(await createBlobProgram(), ["blob", "credentials", "my-bucket"])).toEqual(credentials); + expect((fetchSpy.mock.calls[2]?.[1] as RequestInit).headers).toEqual({ Authorization: "Bearer bucket-token" }); + }); + + it("rejects a missing, unknown, doubled or keyed bucket", async () => { + const fetchSpy = vi.spyOn(globalThis, "fetch").mockImplementation(async () => listed()); + await expect(runCommand(await createBlobProgram(), ["blob", "get"])).rejects.toThrow("Name a bucket: upstash blob get "); + await expect(runCommand(await createBlobProgram(), ["blob", "delete", "nope"])).rejects.toThrow('Blob bucket "nope" not found'); + await expect(runCommand(await createBlobProgram(), ["blob", "get", "a", "--bucket-id", "b"])).rejects.toThrow("Name the bucket once"); + await expect(runCommand(await createBlobProgram(), ["blob", "delete", "my-bucket/key"])).rejects.toThrow("includes a key"); + fetchSpy.mockClear(); + await expect(runCommand(await createBlobProgram(), ["blob", "credentials", "my-bucket", "--token", "t"])) + .rejects.toThrow("Use either --token or a bucket"); + expect(fetchSpy).not.toHaveBeenCalled(); + }); +}); + describe("blob credentials command", () => { it("by bucket id fetches the bucket first, then exchanges its token", async () => { const bucket = makeBucket({ id: "bucket_456", token: "bucket-token" }); @@ -277,7 +337,7 @@ describe("blob credentials command", () => { const program = await createBlobProgram(); await expect(runCommand(program, ["blob", "credentials"])) - .rejects.toThrow(/--token.*UPSTASH_BLOB_TOKEN.*--bucket-id/); + .rejects.toThrow(/Name a bucket.*--token.*UPSTASH_BLOB_TOKEN/); }); it("prints successful credential responses unchanged", async () => {