diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 2254c1b..3d18d2b 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -33,7 +33,7 @@ jobs: - name: Publish release candidate if: github.event.release.prerelease - run: npm publish --access public --tag next --provenance + run: npm publish --access public --tag=canary --provenance - name: Publish if: "!github.event.release.prerelease" diff --git a/README.md b/README.md index 93a3305..41f7e62 100644 --- a/README.md +++ b/README.md @@ -7,6 +7,8 @@ Agent-friendly CLI for managing & debugging Upstash resources from your terminal ## Installation +Requires Node.js 20 or newer. + ```bash npm i -g @upstash/cli ``` @@ -70,6 +72,7 @@ upstash qstash stats --qstash-id $QSTASH_ID --period 7d 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 # Team upstash team list @@ -78,6 +81,61 @@ 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 + +Set `UPSTASH_BLOB_TOKEN` in your environment or `.env` file, then run: + +```bash +upstash blob upload ./assets --prefix assets +``` + +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: + +```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" +``` + +`--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. + +```bash +upstash blob upload ./assets --prefix assets --dry-run +upstash blob upload ./assets --prefix assets --skip-existing +``` + +`--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`. + ## Telemetry The CLI identifies itself to the Upstash API on each request, so we can see which diff --git a/package-lock.json b/package-lock.json index cc883e3..86ceb97 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,16 +1,18 @@ { "name": "@upstash/cli", - "version": "1.0.0", + "version": "0.0.0", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@upstash/cli", - "version": "1.0.0", + "version": "0.0.0", "license": "MIT", "dependencies": { + "@upstash/blob": "0.0.5", "commander": "^13.0.0", - "dotenv": "^16.4.5" + "dotenv": "^16.4.5", + "mime": "4.1.0" }, "bin": { "upstash": "dist/cli.js" @@ -21,7 +23,7 @@ "vitest": "^2.0.0" }, "engines": { - "node": ">=18.0.0" + "node": ">=20.0.0" } }, "node_modules/@esbuild/aix-ppc64": { @@ -789,6 +791,23 @@ "undici-types": "~6.21.0" } }, + "node_modules/@upstash/blob": { + "version": "0.0.5", + "resolved": "https://registry.npmjs.org/@upstash/blob/-/blob-0.0.5.tgz", + "integrity": "sha512-5go9FxMd0yJU3g09+UyKDxGCGkwKJ9YG7kFRHpwEvTtHAuWuQtu4nkjuunRh6fuNiFMplVNLvTjosSD5UT/ZVA==", + "license": "MIT", + "engines": { + "node": ">=20" + }, + "peerDependencies": { + "react": ">=18" + }, + "peerDependenciesMeta": { + "react": { + "optional": true + } + } + }, "node_modules/@vitest/expect": { "version": "2.1.9", "resolved": "https://registry.npmjs.org/@vitest/expect/-/expect-2.1.9.tgz", @@ -1096,6 +1115,21 @@ "@jridgewell/sourcemap-codec": "^1.5.5" } }, + "node_modules/mime": { + "version": "4.1.0", + "resolved": "https://registry.npmjs.org/mime/-/mime-4.1.0.tgz", + "integrity": "sha512-X5ju04+cAzsojXKes0B/S4tcYtFAJ6tTMuSPBEn9CPGlrWr8Fiw7qYeLT0XyH80HSoAoqWCaz+MWKh22P7G1cw==", + "funding": [ + "https://github.com/sponsors/broofa" + ], + "license": "MIT", + "bin": { + "mime": "bin/cli.js" + }, + "engines": { + "node": ">=16" + } + }, "node_modules/ms": { "version": "2.1.3", "resolved": "https://registry.npmjs.org/ms/-/ms-2.1.3.tgz", diff --git a/package.json b/package.json index 4491786..2830931 100644 --- a/package.json +++ b/package.json @@ -25,8 +25,10 @@ "author": "Upstash", "license": "MIT", "dependencies": { + "@upstash/blob": "0.0.5", "commander": "^13.0.0", - "dotenv": "^16.4.5" + "dotenv": "^16.4.5", + "mime": "4.1.0" }, "devDependencies": { "@types/node": "^20.10.0", @@ -34,7 +36,7 @@ "vitest": "^2.0.0" }, "engines": { - "node": ">=18.0.0" + "node": ">=20.0.0" }, "publishConfig": { "access": "public" diff --git a/src/commands/blob/credentials.ts b/src/commands/blob/credentials.ts index 19bc1ac..20fea58 100644 --- a/src/commands/blob/credentials.ts +++ b/src/commands/blob/credentials.ts @@ -155,10 +155,18 @@ interface BucketTokenSource { unauthorizedRetries: number; } -function resolveBucketToken( - flags: { bucketId?: string }, +export function resolveBucketToken( + flags: { bucketId?: string; token?: string }, command: Command, ): Promise { + if (flags.token !== undefined) { + if (flags.bucketId !== undefined) { + return Promise.reject(new Error("Use either --token or --bucket-id, not both")); + } + const token = flags.token.trim(); + if (!token) return Promise.reject(new Error("--token must be a non-empty Blob bucket token")); + return Promise.resolve({ token, unauthorizedRetries: 0 }); + } if (flags.bucketId) { const auth = resolveAuth(command); return request(auth, "GET", `/v2/blob/bucket/${flags.bucketId}`).then((bucket) => { @@ -179,7 +187,7 @@ function resolveBucketToken( return Promise.reject( new Error( - "Blob credentials require either --bucket-id with Upstash account authentication or a non-empty UPSTASH_BLOB_TOKEN environment variable", + "Provide --token, UPSTASH_BLOB_TOKEN in the environment or .env, or --bucket-id with Upstash account authentication", ), ); } @@ -191,6 +199,7 @@ export function registerBlobCredentials(blob: Command): void { "Get temporary S3 credentials for a Blob bucket; expiresAt is the credential expiry", ) .option("--bucket-id ", "Blob bucket ID") + .option("--token ", "Blob bucket token; no management API key needed (overrides UPSTASH_BLOB_TOKEN)") .addHelpText( "after", ` @@ -198,7 +207,7 @@ With --bucket-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 }, command: Command) => { + .action(async (flags: { bucketId?: string; token?: string }, command: Command) => { const source = await resolveBucketToken(flags, command); const credentials = await fetchBlobCredentials(source.token, sleep, { unauthorizedRetries: source.unauthorizedRetries, diff --git a/src/commands/blob/index.ts b/src/commands/blob/index.ts index 8652c91..f2678c4 100644 --- a/src/commands/blob/index.ts +++ b/src/commands/blob/index.ts @@ -4,6 +4,7 @@ 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"; export function registerBlob(program: Command): void { const blob = program.command("blob").description("Manage Blob buckets"); @@ -13,4 +14,5 @@ export function registerBlob(program: Command): void { registerBlobGet(blob); registerBlobDelete(blob); registerBlobCredentials(blob); + registerBlobUpload(blob); } diff --git a/src/commands/blob/upload.ts b/src/commands/blob/upload.ts new file mode 100644 index 0000000..5218d36 --- /dev/null +++ b/src/commands/blob/upload.ts @@ -0,0 +1,191 @@ +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-upload.test.ts b/tests/unit/blob-upload.test.ts new file mode 100644 index 0000000..e43a435 --- /dev/null +++ b/tests/unit/blob-upload.test.ts @@ -0,0 +1,309 @@ +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 ca5edcf..e9c79a1 100644 --- a/tests/unit/blob.test.ts +++ b/tests/unit/blob.test.ts @@ -64,6 +64,7 @@ describe("blob command registration", () => { "get", "delete", "credentials", + "upload", ]); }); }); @@ -230,7 +231,9 @@ describe("blob credentials command", () => { }); }); - it("without bucket id uses UPSTASH_BLOB_TOKEN and skips Developer API auth", async () => { + it.each(["environment", "flag"])("uses the %s token and skips Developer API auth", async (source) => { + delete process.env.UPSTASH_EMAIL; + delete process.env.UPSTASH_API_KEY; process.env.UPSTASH_BLOB_TOKEN = "env-bucket-token"; const credentials = makeCredentials(); const fetchSpy = vi.spyOn(globalThis, "fetch").mockResolvedValue( @@ -238,11 +241,15 @@ describe("blob credentials command", () => { ); const program = await createBlobProgram(); - const result = await runCommand(program, ["blob", "credentials"]); + const flags = source === "flag" ? ["--token", "flag-bucket-token"] : []; + const result = await runCommand(program, ["blob", "credentials", ...flags]); expect(result).toEqual(credentials); expect(fetchSpy).toHaveBeenCalledTimes(1); expect(fetchSpy.mock.calls[0]?.[0]).toBe("https://blob.upstash.io/v1/credentials"); + expect((fetchSpy.mock.calls[0]?.[1] as RequestInit).headers).toEqual({ + Authorization: `Bearer ${source === "flag" ? "flag-bucket-token" : "env-bucket-token"}`, + }); }); it("explicit bucket id wins over an ambient UPSTASH_BLOB_TOKEN", async () => { @@ -270,7 +277,7 @@ describe("blob credentials command", () => { const program = await createBlobProgram(); await expect(runCommand(program, ["blob", "credentials"])) - .rejects.toThrow(/either --bucket-id.*UPSTASH_BLOB_TOKEN/); + .rejects.toThrow(/--token.*UPSTASH_BLOB_TOKEN.*--bucket-id/); }); it("prints successful credential responses unchanged", async () => {