Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
93 changes: 47 additions & 46 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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://<bucket>/<key>` in place of
`s3://`. `<bucket>` 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

Expand Down
206 changes: 206 additions & 0 deletions src/commands/blob/buckets.ts
Original file line number Diff line number Diff line change
@@ -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<string> {
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()} <name-or-id>`);
const name = parseBucket(bucket);
return findAccountBucket(await request<BlobBucket[]>(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<string, Promise<Bucket>>();
private accountBuckets?: Promise<BlobBucket[]>;

constructor(private readonly command: Command, private readonly token?: string) {}

open(name: string): Promise<Bucket> {
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<BlobBucket[]> {
this.accountBuckets ??= request<BlobBucket[]>(this.auth(), "GET", "/v2/blob/bucket");
return this.accountBuckets;
}

private async resolve(name: string): Promise<Bucket> {
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<BlobBucket>(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(/&lt;/g, "<").replace(/&gt;/g, ">").replace(/&quot;/g, '"')
.replace(/&#39;|&apos;/g, "'").replace(/&amp;/g, "&");
}

function tag(xml: string, name: string): string | undefined {
const match = new RegExp(`<${name}>([\\s\\S]*?)</${name}>`).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]*?)</${name}>`, "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<DirectoryPage> {
const s3 = bucket.s3();
for (let attempt = 0; ; attempt++) {
const [{ url: endpoint }, credentials] = await Promise.all([s3.endpoint(), s3.credentials()]);
const query: Record<string, string> = { 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<string, string> = {
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,
};
}
}
Loading
Loading