diff --git a/.github/workflows/cloudflare-workers.yml b/.github/workflows/cloudflare-workers.yml new file mode 100644 index 0000000..c505a48 --- /dev/null +++ b/.github/workflows/cloudflare-workers.yml @@ -0,0 +1,19 @@ +name: Cloudflare workers + +on: + push: + branches: [main] + paths: ["cloudflare-workers/**"] + pull_request: + paths: ["cloudflare-workers/**"] + +jobs: + test: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-node@v4 + with: + node-version: 24 + - run: npm ci --prefix cloudflare-workers/backblaze-proxy + - run: node --test "cloudflare-workers/**/*.test.mjs" "cloudflare-workers/**/*.test.ts" diff --git a/backblaze-proxy.sh b/backblaze-proxy.sh new file mode 100755 index 0000000..bcd2b36 --- /dev/null +++ b/backblaze-proxy.sh @@ -0,0 +1,380 @@ +#!/bin/bash + +# Deploys the Backblaze proxy: a Cloudflare Worker in front of the S3 bucket +# (Backblaze B2) that serves demos, clips, news images, event media and map +# assets through Cloudflare, so B2 egress is free and popular files come from +# the edge. +# +# Walks through everything it needs: the bucket keys, the hostname (default +# cf.), signing in to Cloudflare, the hostname's DNS record, the +# deploy and the routes, then saves the hostname as CLOUDFLARE_WORKER_DOMAIN +# in overlays/config/api-config.env, which the panel builds its download URLs +# from, and offers Smart Tiered Cache. +# +# The bucket comes from the panel's config. The keys come from the cluster on +# Vault installs and from overlays/local-secrets otherwise, and Backblaze has +# to accept them before anything is deployed: one worker serves every panel +# pointed at it, so bad keys would break downloads for all of them. +# +# Safe to run again; it updates the worker in place and keeps any routes it +# already has on other hostnames, so older links keep working. +# +# ./backblaze-proxy.sh [hostname] + +PANEL_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +source "$PANEL_DIR/utils/colors.sh" +source "$PANEL_DIR/utils/interactive_select.sh" +source "$PANEL_DIR/utils/print_domains_and_hosts.sh" +source "$PANEL_DIR/utils/update_env_var.sh" +source "$PANEL_DIR/utils/cloudflare_workers.sh" + +# Existing Cloudflare routes point at this name; renaming spawns a second +# worker while the old one keeps serving. +WORKER_NAME="5stack" +WORKER_DIR="$PANEL_DIR/cloudflare-workers/backblaze-proxy" +DOCS_URL="https://docs.5stack.gg/advanced/s3/backblaze" +ROUTE_PATHS=("/demo*" "/clips*" "/news*" "/maps*" "/events*") + +read_env() { + grep -h "^$1=" "$2" 2>/dev/null | tail -n 1 | cut -d '=' -f2- +} + +# With Vault the files in overlays/local-secrets are placeholders; the real +# values are only in the cluster, synced from Vault under these names. +read_cluster_secret() { + kubectl --kubeconfig="$PANEL_KUBECONFIG" -n 5stack get secret "$1" -o "jsonpath={.data.$2}" 2>/dev/null \ + | base64 --decode 2>/dev/null +} + +normalize_host() { + local host="$1" + host="${host#*://}" + host="${host%%/*}" + host="$(echo "$host" | tr '[:upper:]' '[:lower:]')" + if [ -n "$host" ] && [[ "$host" != *.* ]] && [ -n "$WEB_DOMAIN" ]; then + host="$host.$WEB_DOMAIN" + fi + echo "$host" +} + +panel_host_name() { + local var + for var in WEB_DOMAIN API_DOMAIN WS_DOMAIN RELAY_DOMAIN DEMOS_DOMAIN GAME_STREAM_DOMAIN S3_CONSOLE_HOST TYPESENSE_HOST; do + if [ -n "${!var}" ] && [ "${!var}" = "$1" ]; then + echo "$var" + return + fi + done +} + +# panel_graphql QUERY VARIABLES -- runs QUERY against the panel's Hasura as +# admin, the same way the web's settings pages write settings. +panel_graphql() { + local body + body="$(node -p 'JSON.stringify({ query: process.argv[1], variables: JSON.parse(process.argv[2]) })' "$1" "$2")" + printf 'header = "x-hasura-admin-secret: %s"\n' "$HASURA_ADMIN_SECRET" \ + | curl -sS -K - --max-time 15 -H "Content-Type: application/json" \ + --data "$body" "https://$API_DOMAIN/v1/graphql" +} + +# The api copies CLOUDFLARE_WORKER_DOMAIN into this setting when it boots; +# writing it here as well means it applies without an ./update.sh. +set_live_worker_url() { + local variables + variables="$(node -p 'JSON.stringify({ value: process.argv[1] })' "$1")" + [ -n "$(panel_graphql \ + 'mutation ($value: String!) { insert_settings(objects: [{ name: "cloudflare_worker_url", value: $value }], on_conflict: { constraint: settings_pkey, update_columns: [value] }) { affected_rows } }' \ + "$variables" | cf_json 'j.data?.insert_settings ? "ok" : undefined')" ] +} + +# The api copies CLOUDFLARE_WORKER_DOMAIN into its setting when it restarts, +# which ./update.sh does whenever the config changed. +offer_update() { + local answer + if [ ! -f "$PANEL_KUBECONFIG" ]; then + warn "Could not apply it to the running panel from here. Run ./update.sh on your" + warn "panel's server to apply it." + return + fi + warn "Could not apply it to the running panel directly, so ./update.sh has to apply it." + read -r -p " Run ./update.sh now, against the cluster in $PANEL_KUBECONFIG? [Y/n] " answer + if [[ "$answer" =~ ^[Nn] ]]; then + warn "Run ./update.sh to apply it." + return + fi + if ! "$PANEL_DIR/update.sh"; then + die "./update.sh failed. See the output above." + fi + ok "The panel now serves demos, clips and media through $WORKER_URL" +} + +# Signs a read of an object that does not exist and prints Backblaze's error +# code. A missing object and a key without list access both come back as +# AccessDenied or NoSuchKey; only a bad key ID or secret is rejected outright. +bucket_key_check() { + (cd "$WORKER_DIR" && S3_ACCESS_KEY="$S3_ACCESS_KEY" S3_SECRET="$S3_SECRET" node --input-type=module -e ' + import { AwsClient } from "aws4fetch"; + const client = new AwsClient({ + accessKeyId: process.env.S3_ACCESS_KEY, + secretAccessKey: process.env.S3_SECRET, + service: "s3", + }); + const url = `https://${process.argv[1]}.${process.argv[2]}/.5stack-key-check-${crypto.randomUUID()}`; + try { + const signed = await client.sign(url, { method: "GET", headers: new Headers() }); + const response = await fetch(signed.url, { method: "GET", headers: signed.headers }); + const body = await response.text(); + process.stdout.write((body.match(/([^<]+)., so the in-cluster +# storage the panel ships with cannot sit behind it. +if [[ "$S3_ENDPOINT" != *.* ]]; then + err "S3_ENDPOINT ($S3_ENDPOINT) is the panel's own storage, not a public S3 host." + err "This worker is for a remote bucket such as Backblaze B2" + err "(e.g. s3.us-east-005.backblazeb2.com):" + cf_link "$DOCS_URL" + exit 1 +fi + +cf_require_node + +step "Installing the worker's dependencies" +if ! npm ci --silent --no-audit --no-fund --prefix "$WORKER_DIR"; then + die "npm ci failed. See the output above." +fi + +# The worker signs every read with these keys and checks upload tokens against +# S3_SECRET, which the api signs them with, so they have to be the api's keys. +step "Checking the bucket keys with Backblaze" +if [ -z "$S3_ACCESS_KEY" ] || [ -z "$S3_SECRET" ]; then + warn "Could not read S3_ACCESS_KEY and S3_SECRET from $KEYS_FROM." + S3_ACCESS_KEY="" + S3_SECRET="" +fi +while true; do + if [ -n "$S3_ACCESS_KEY" ] && [ -n "$S3_SECRET" ]; then + KEY_CHECK="$(bucket_key_check)" + case "$KEY_CHECK" in + unreachable*) + die "Could not reach https://$S3_BUCKET.$S3_ENDPOINT: ${KEY_CHECK#unreachable }" + ;; + InvalidAccessKeyId|SignatureDoesNotMatch|InvalidSecurity|InvalidToken|InvalidArgument) + err "Backblaze rejected the keys from $KEYS_FROM ($KEY_CHECK)." + ;; + *) + ok "Backblaze accepted the keys from $KEYS_FROM" + break + ;; + esac + fi + warn "Enter the application key the panel's API uses for $S3_BUCKET. It is stored in" + warn "Cloudflare as worker secrets and nowhere else." + S3_ACCESS_KEY="" + S3_SECRET="" + while [ -z "$S3_ACCESS_KEY" ]; do + read -r -p " S3 access key ID: " S3_ACCESS_KEY || die "Stopped." + done + while [ -z "$S3_SECRET" ]; do + read_masked " S3 secret key: " S3_SECRET + done + KEYS_FROM="what you entered" +done + +step "Hostname" +echo " The worker gets a hostname of its own, on a domain you have on Cloudflare." +DEFAULT_HOST="$CLOUDFLARE_WORKER_DOMAIN" +if [ -z "$DEFAULT_HOST" ] && [ -n "$WEB_DOMAIN" ]; then + DEFAULT_HOST="cf.$WEB_DOMAIN" +fi +WORKER_HOST="$(normalize_host "$1")" +while true; do + if [ -z "$WORKER_HOST" ]; then + read -r -p " Hostname for the worker${DEFAULT_HOST:+ [$DEFAULT_HOST]}: " WORKER_HOST + WORKER_HOST="$(normalize_host "${WORKER_HOST:-$DEFAULT_HOST}")" + fi + CONFLICT="$(panel_host_name "$WORKER_HOST")" + if [ -z "$WORKER_HOST" ]; then + continue + elif [[ "$WORKER_HOST" != *.* ]]; then + err "Enter the full hostname, like cf.example.com." + elif [ -n "$CONFLICT" ]; then + err "$WORKER_HOST is the panel's own $CONFLICT. The worker takes over" + err "${ROUTE_PATHS[*]} on its hostname, which would break the panel there." + err "Use a hostname only the worker answers, like ${DEFAULT_HOST:-cf.}." + else + break + fi + WORKER_HOST="" +done + +cf_sign_in "$WORKER_HOST" + +step "Finding $WORKER_HOST in Cloudflare" +cf_find_zone "$WORKER_HOST" + +cf_wait_for_proxied "$WORKER_HOST" placeholder + +ROUTES=() +for path in "${ROUTE_PATHS[@]}"; do + ROUTES+=("$WORKER_HOST$path") +done + +VARS=(--var "BUCKET_NAME:$S3_BUCKET" --var "S3_ENDPOINT:$S3_ENDPOINT") +if [ -n "$API_DOMAIN" ]; then + VARS+=(--var "API_URL:https://$API_DOMAIN") +fi + +step "Ready to deploy" +ok "Worker: $WORKER_NAME" +ok "Hostname: https://$WORKER_HOST (${ROUTE_PATHS[*]})" +ok "Bucket: $S3_BUCKET at $S3_ENDPOINT" +if [ -n "$API_DOMAIN" ]; then + ok "API: https://$API_DOMAIN (clip views are counted there)" +fi +read -r -p " Deploy it? [Y/n] " CONFIRM +if [[ "$CONFIRM" =~ ^[Nn] ]]; then + warn "Nothing deployed." + exit 0 +fi + +step "Deploying the worker" +if ! wrangler deploy --config "$WORKER_DIR/wrangler.toml" "${VARS[@]}"; then + die "The deploy failed. See the wrangler output above." +fi + +step "Storing the bucket keys as worker secrets" +SECRETS_FILE="$(mktemp)" +trap 'rm -f "$SECRETS_FILE"' EXIT +chmod 600 "$SECRETS_FILE" +S3_ACCESS_KEY="$S3_ACCESS_KEY" S3_SECRET="$S3_SECRET" node -e \ + 'process.stdout.write(JSON.stringify({ S3_ACCESS_KEY: process.env.S3_ACCESS_KEY, S3_SECRET: process.env.S3_SECRET }))' \ + > "$SECRETS_FILE" +if ! wrangler secret bulk "$SECRETS_FILE" --config "$WORKER_DIR/wrangler.toml"; then + die "Storing the secrets failed. See the wrangler output above." +fi + +step "Routing $WORKER_HOST through it" +cf_ensure_routes "$WORKER_NAME" false "${ROUTES[@]}" + +step "Waiting for it to answer on https://$WORKER_HOST" +# The worker's health check signs a read against the bucket itself, so this +# fails on keys Backblaze rejects, not only on a route that is missing. New +# secrets take a few seconds to reach every Cloudflare location, so a rejection +# right after storing them is the previous keys still answering: keep asking. +ANSWERING=false +BUCKET_STATE="" +for _ in $(seq 1 24); do + HEALTH="$(cf_curl "$WORKER_HOST" -sS --max-time 20 "https://$WORKER_HOST/demo/_health" 2>/dev/null)" + if [ "$(echo "$HEALTH" | cf_json 'j.worker')" = "5stack-backblaze-proxy" ]; then + ANSWERING=true + BUCKET_STATE="$(echo "$HEALTH" | cf_json 'j.bucket')" + if [ "$BUCKET_STATE" = "ok" ]; then + break + fi + fi + sleep 5 +done +if [ "$ANSWERING" != true ]; then + err "https://$WORKER_HOST/demo/_health is not answering from the worker yet." + err "Check the routes in $CF_ZONE_NAME's Workers Routes, or run this again in a minute." + cf_link "https://dash.cloudflare.com/$CLOUDFLARE_ACCOUNT_ID/$CF_ZONE_NAME/workers" + exit 1 +fi +if [ "$BUCKET_STATE" != "ok" ]; then + err "The worker is up, but it still cannot read $S3_BUCKET after two minutes" + err "($BUCKET_STATE: $(echo "$HEALTH" | cf_json 'j.code')). Run this again and enter the application" + err "key the panel's API uses." + exit 1 +fi +ok "The Backblaze proxy is live on https://$WORKER_HOST" + +if [ -n "$CF_OTHER_ROUTES" ]; then + ok "Its earlier routes were kept, so links that still use them keep working:" + while IFS= read -r route; do + ok " $route" + done <<< "$CF_OTHER_ROUTES" +fi + +step "Pointing the panel at it" +WORKER_URL="https://$WORKER_HOST" +LIVE_URL="" +if [ -n "$API_DOMAIN" ] && [ -n "$HASURA_ADMIN_SECRET" ]; then + LIVE_URL="$(panel_graphql 'query { settings_by_pk(name: "cloudflare_worker_url") { value } }' '{}' \ + | cf_json 'j.data ? (j.data.settings_by_pk?.value ?? "") : undefined')" +fi + +CURRENT_URL="$LIVE_URL" +if [ -n "$CLOUDFLARE_WORKER_DOMAIN" ]; then + CURRENT_URL="https://$CLOUDFLARE_WORKER_DOMAIN" +fi +SWITCH=y +if [ -n "$CURRENT_URL" ] && [ "$CURRENT_URL" != "$WORKER_URL" ]; then + ok "The panel serves files through $CURRENT_URL now." + read -r -p " Switch it to $WORKER_URL? [Y/n] " SWITCH + SWITCH="${SWITCH:-y}" +fi + +if [[ ! "$SWITCH" =~ ^[Yy] ]]; then + warn "Left it on $CURRENT_URL. Run this again to switch later." +else + if [ "$CLOUDFLARE_WORKER_DOMAIN" != "$WORKER_HOST" ]; then + update_env_var "$PANEL_DIR/overlays/config/api-config.env" CLOUDFLARE_WORKER_DOMAIN "$WORKER_HOST" + fi + ok "CLOUDFLARE_WORKER_DOMAIN=$WORKER_HOST is in overlays/config/api-config.env" + + if [ "$LIVE_URL" = "$WORKER_URL" ]; then + ok "The panel already serves files through $WORKER_URL" + elif [ -n "$API_DOMAIN" ] && [ -n "$HASURA_ADMIN_SECRET" ] && set_live_worker_url "$WORKER_URL"; then + ok "The panel now serves demos, clips and media through $WORKER_URL" + else + offer_update + fi +fi + +cf_offer_tiered_cache diff --git a/cloudflare-workers/backblaze-proxy/.dev.vars.example b/cloudflare-workers/backblaze-proxy/.dev.vars.example new file mode 100644 index 0000000..ab6903a --- /dev/null +++ b/cloudflare-workers/backblaze-proxy/.dev.vars.example @@ -0,0 +1,13 @@ +# Copy to `.dev.vars` (which is gitignored) and fill in. Used by +# `wrangler dev` for the local proxy. Production is deployed with +# ./backblaze-proxy.sh, which sets these as vars and worker secrets. + +BUCKET_NAME=5stack +S3_ENDPOINT=s3.us-east-005.backblazeb2.com + +# Base URL of the 5stack API; when set, the worker beacons clip playback +# starts to {API_URL}/clip-views/play to count views. Leave unset to disable. +API_URL="" + +S3_ACCESS_KEY="" +S3_SECRET="" diff --git a/cloudflare-workers/backblaze-proxy/.gitignore b/cloudflare-workers/backblaze-proxy/.gitignore new file mode 100644 index 0000000..153b333 --- /dev/null +++ b/cloudflare-workers/backblaze-proxy/.gitignore @@ -0,0 +1,3 @@ +.wrangler +node_modules +.dev.vars diff --git a/cloudflare-workers/backblaze-proxy/index.test.ts b/cloudflare-workers/backblaze-proxy/index.test.ts new file mode 100644 index 0000000..08e36ad --- /dev/null +++ b/cloudflare-workers/backblaze-proxy/index.test.ts @@ -0,0 +1,168 @@ +import { afterEach, beforeEach, describe, it, mock } from "node:test"; +import assert from "node:assert/strict"; +import worker from "./index.ts"; + +const env = { + S3_ACCESS_KEY: "key", + S3_SECRET: "secret", + BUCKET_NAME: "5stack", + S3_ENDPOINT: "s3.example.test", +}; +const ORIGIN = "https://5stack.gg"; + +const realFetch = globalThis.fetch; +let upstream: ReturnType; + +beforeEach(() => { + upstream = mock.fn( + async () => + new Response("{}", { + status: 200, + headers: { "Content-Type": "application/json" }, + }), + ); + globalThis.fetch = upstream as unknown as typeof fetch; + (globalThis as any).caches = { + default: { match: async () => undefined, put: async () => {} }, + }; +}); + +afterEach(() => { + globalThis.fetch = realFetch; + delete (globalThis as any).caches; +}); + +async function get(path: string) { + const ctx = { waitUntil: () => {}, passThroughOnException: () => {} }; + const response = await worker.fetch( + new Request(`https://demo-dl.5stack.gg/${path}`, { + headers: { Origin: ORIGIN }, + }), + env, + ctx as any, + ); + const init = upstream.mock.calls[0].arguments[1] as { + cf: { cacheTtlByStatus: Record }; + }; + return { response, edgeTtl: init.cf.cacheTtlByStatus["200-299"] }; +} + +describe("backblaze-proxy map assets", () => { + it("lets latest.json go stale within a minute, at the edge and in the browser", async () => { + const { response, edgeTtl } = await get("maps/latest.json"); + + assert.equal(response.status, 200); + assert.equal(response.headers.get("Cache-Control"), "public, max-age=60"); + assert.equal(edgeTtl, 60); + assert.match( + String(upstream.mock.calls[0].arguments[0]), + /\/maps\/latest\.json$/, + ); + }); + + for (const path of [ + "maps/25537370/manifest.json", + "maps/25537370/de_mirage.view.bin.gz", + ]) { + it(`keeps ${path} immutable`, async () => { + const { response, edgeTtl } = await get(path); + + assert.equal( + response.headers.get("Cache-Control"), + "public, max-age=2592000, immutable", + ); + assert.equal(edgeTtl, 2592000); + }); + } + + for (const path of [ + "maps/latest.json", + "maps/25537370/manifest.json", + "maps/25537370/de_mirage.tri.gz", + ]) { + it(`answers ${path} with the same CORS headers`, async () => { + const { response } = await get(path); + + assert.equal(response.headers.get("Access-Control-Allow-Origin"), ORIGIN); + assert.equal( + response.headers.get("Access-Control-Allow-Credentials"), + "true", + ); + assert.match( + response.headers.get("Access-Control-Expose-Headers") ?? "", + /Content-Length/, + ); + assert.equal(response.headers.get("Vary"), "Origin"); + }); + } +}); + +describe("backblaze-proxy upstream failures", () => { + it("answers with B2's 403 once every retry is spent", async () => { + upstream = mock.fn( + async () => + new Response("AccessDenied", { + status: 403, + }), + ); + globalThis.fetch = upstream as unknown as typeof fetch; + + const { response } = await get("clips/missing.mp4"); + + assert.equal(upstream.mock.callCount(), 3); + assert.equal(response.status, 403); + assert.match(await response.text(), /AccessDenied/); + }); +}); + +describe("backblaze-proxy health", () => { + async function health(body: string, overrides: Partial = {}) { + upstream = mock.fn(async () => new Response(body, { status: 403 })); + globalThis.fetch = upstream as unknown as typeof fetch; + const ctx = { waitUntil: () => {}, passThroughOnException: () => {} }; + const response = await worker.fetch( + new Request("https://cf.5stack.gg/demo/_health"), + { ...env, ...overrides }, + ctx as any, + ); + return { response, json: await response.json() }; + } + + it("reports the bucket as reachable when Backblaze accepts the keys", async () => { + const { response, json } = await health( + "AccessDenied", + ); + + assert.equal(response.headers.get("Access-Control-Allow-Origin"), "*"); + assert.equal(response.headers.get("Cache-Control"), "no-store"); + assert.deepEqual(json, { + ok: true, + worker: "5stack-backblaze-proxy", + version: "1", + bucket: "ok", + code: "AccessDenied", + }); + assert.match( + String(upstream.mock.calls[0].arguments[0]), + /^https:\/\/5stack\.s3\.example\.test\/\.5stack-health$/, + ); + }); + + it("reports keys Backblaze rejects", async () => { + const { json } = await health( + "InvalidAccessKeyId", + ); + + assert.equal(json.ok, false); + assert.equal(json.bucket, "rejected"); + assert.equal(json.code, "InvalidAccessKeyId"); + }); + + it("reports a worker deployed without its bucket keys", async () => { + const { json } = await health("", { S3_SECRET: "" }); + + assert.equal(json.bucket, "misconfigured"); + assert.equal(upstream.mock.callCount(), 0); + }); +}); + diff --git a/cloudflare-workers/backblaze-proxy/index.ts b/cloudflare-workers/backblaze-proxy/index.ts new file mode 100644 index 0000000..4dc1cfd --- /dev/null +++ b/cloudflare-workers/backblaze-proxy/index.ts @@ -0,0 +1,662 @@ +import { AwsClient } from "aws4fetch"; + +// How many times to re-sign and retry an upstream request that came back +// 403/5xx. B2 intermittently rejects otherwise-valid signed reads; without a +// retry a single bad roll surfaces in the UI as a permanently broken clip, +// because the failure is never cached and every cold request rolls again. +const UPSTREAM_ATTEMPTS = 3; + +const IMMUTABLE_TTL = 2592000; + +// Every object is immutable except the map-asset pointer, which moves each +// time a new CS2 build is published; the edge and the browser must both let +// go of it within a minute. +const MUTABLE_TTL: Record = { "maps/latest.json": 60 }; + +// Nothing from the client request is forwarded to B2. Every header we sign is +// a header Cloudflare may rewrite between sign() and fetch() — which B2 then +// reads as SignatureDoesNotMatch and answers 403. B2 needs none of them: +// Range is served out of the cached full 200, and conditionals would only +// defeat the edge cache. It also keeps the viewer's session cookie, UA and +// referer from being shipped to Backblaze on every clip view. +async function signedFetch( + method: string, + url: string, + env: { S3_ACCESS_KEY: string; S3_SECRET: string }, + ttl: number, +): Promise { + const client = new AwsClient({ + accessKeyId: env.S3_ACCESS_KEY, + secretAccessKey: env.S3_SECRET, + service: "s3", + }); + + let response: Response | null = null; + for (let attempt = 0; attempt < UPSTREAM_ATTEMPTS; attempt++) { + // Retries must carry a unique query param or they are not retries at all: + // every attempt shares one Cloudflare cache key, so a cached 403 would be + // replayed from the edge three times without B2 ever being asked again. + // S3 ignores unrecognized query params, and signing the busted URL keeps + // the signature valid. + const target = + attempt === 0 + ? url + : `${url}?x-5stack-retry=${attempt}-${crypto.randomUUID()}`; + const signed = await client.sign(target, { + method, + headers: new Headers(), + }); + response = await fetch(signed.url, { + method: signed.method, + headers: signed.headers, + cf: { + cacheEverything: true, + // 4xx must not be pinned: B2 has no ListBucket grant on our key, so a + // transient denial and a genuinely missing object both arrive as 403, + // and caching either one would outlive the condition that caused it. + cacheTtlByStatus: { "200-299": ttl, "400-499": 0, "500-599": 0 }, + }, + }); + if (response.status !== 403 && response.status < 500) { + return response; + } + // The last attempt is what the client gets, so its body has to survive: + // handing a cancelled body to a new Response throws, and the client sees + // Cloudflare's 1101 instead of the 403. + if (attempt < UPSTREAM_ATTEMPTS - 1) { + await response.body?.cancel(); + } + } + return response!; +} + +// Under /demo* because that is the one route every deployment of this worker +// has had; the panel's settings page and ./backblaze-proxy.sh read it. +const HEALTH_PATH = "/demo/_health"; + +// What Backblaze answers when the key ID or secret itself is wrong, as opposed +// to a key that works but may not read or list a particular object. +const REJECTED_KEY_CODES = new Set([ + "InvalidAccessKeyId", + "SignatureDoesNotMatch", + "InvalidSecurity", + "InvalidToken", +]); + +async function health(env: { + S3_ACCESS_KEY: string; + S3_SECRET: string; + BUCKET_NAME: string; + S3_ENDPOINT: string; +}): Promise { + let bucket: "ok" | "rejected" | "misconfigured" | "unreachable" = "ok"; + let code: string | null = null; + + if ( + !env.BUCKET_NAME || + !env.S3_ENDPOINT || + !env.S3_ACCESS_KEY || + !env.S3_SECRET + ) { + bucket = "misconfigured"; + } else { + try { + const client = new AwsClient({ + accessKeyId: env.S3_ACCESS_KEY, + secretAccessKey: env.S3_SECRET, + service: "s3", + }); + const signed = await client.sign( + `https://${env.BUCKET_NAME}.${env.S3_ENDPOINT}/.5stack-health`, + { method: "GET", headers: new Headers() }, + ); + const response = await fetch(signed.url, { + method: signed.method, + headers: signed.headers, + }); + const body = await response.text(); + code = /([^<]+)= threshold; +} + +async function viewerKey(request: Request): Promise { + const ip = request.headers.get("cf-connecting-ip") ?? ""; + const ua = request.headers.get("user-agent") ?? ""; + const digest = await crypto.subtle.digest( + "SHA-256", + new TextEncoder().encode(`${ip}\n${ua}`), + ); + return Array.from(new Uint8Array(digest).slice(0, 16)) + .map((b) => b.toString(16).padStart(2, "0")) + .join(""); +} + +// Whether this request is a viewable clip stream we should track. The +// actual increment is gated on what the response spans (see +// spansViewThreshold). +function shouldTrackView( + request: Request, + url: URL, + key: string | null, + env: { API_URL?: string }, +): boolean { + if (!env.API_URL || request.method !== "GET" || !key) { + return false; + } + if (!/^clips\/.+\.mp4$/.test(key)) { + return false; + } + if ( + url.searchParams.get("dl") === "1" || + url.searchParams.get("download") === "1" || + url.searchParams.get("noview") === "1" + ) { + return false; + } + if (BOT_UA.test(request.headers.get("user-agent") ?? "")) { + return false; + } + return true; +} + +function registerView( + request: Request, + key: string, + env: { S3_SECRET: string; API_URL?: string }, + ctx: ExecutionContext, +) { + ctx.waitUntil( + viewerKey(request) + .then((clientKey) => + fetch(`${env.API_URL!.replace(/\/+$/, "")}/clip-views/play`, { + method: "POST", + headers: { + "content-type": "application/json", + authorization: `Bearer ${env.S3_SECRET}`, + }, + body: JSON.stringify({ file: key, clientKey }), + }), + ) + .catch(() => {}), + ); +} + +// Decorates a response on its way to the client: mp4 Content-Type fix and the +// view beacon. The body is never read or rewritten here. +function clientResponse( + source: Response, + opts: { + isHead: boolean; + key: string; + track: boolean; + request: Request; + env: { S3_SECRET: string; API_URL?: string }; + ctx: ExecutionContext; + }, +): Response { + const headers = new Headers(source.headers); + if (/\.mp4$/i.test(opts.key)) { + // B2 stored old clips as binary/octet-stream; players refuse non-video/*. + headers.set("Content-Type", "video/mp4"); + } + if (!headers.has("Accept-Ranges")) { + headers.set("Accept-Ranges", "bytes"); + } + if (opts.track && spansViewThreshold(source.status, headers)) { + registerView(opts.request, opts.key, opts.env, opts.ctx); + } + // A HEAD keeps the headers (Content-Length included, per spec) and drops the + // body without reading it. + if (opts.isHead) { + opts.ctx.waitUntil(source.body?.cancel() ?? Promise.resolve()); + return new Response(null, { + status: source.status, + statusText: source.statusText, + headers, + }); + } + // The body is handed straight through. Nothing inspects, slices or copies + // it: range slicing is the cache layer's job now (see the cache.match on the + // hot path), which is native and costs the Worker no CPU. + return new Response(source.body, { + status: source.status, + statusText: source.statusText, + headers, + }); +} + +// Path-style (/clips//.mp4) or query-style (?file=). +function resolveKey(url: URL): string | null { + const fromQuery = url.searchParams.get("file"); + if (fromQuery) { + // Older callers built `?file=foo.mp4?dl=1` which parses as + // file=foo.mp4?dl=1; drop anything past the first `?`. + const cleaned = fromQuery.split("?")[0]; + return cleaned ? decodeURIComponent(cleaned) : null; + } + const fromPath = url.pathname.replace(/^\/+/, ""); + return fromPath ? decodeURIComponent(fromPath) : null; +} + +export default { + async fetch( + request: Request, + env: { + S3_ACCESS_KEY: string; + S3_SECRET: string; + BUCKET_NAME: string; + S3_ENDPOINT: string; + API_URL?: string; + }, + ctx: ExecutionContext, + ) { + const reqOrigin = request.headers.get("Origin"); + + if (request.method === "OPTIONS") { + return new Response(null, { + status: 204, + headers: corsHeaders(reqOrigin), + }); + } + + if (request.method === "PUT") { + return handleUpload(request, env, reqOrigin); + } + + if (!["GET", "HEAD"].includes(request.method)) { + return new Response(null, { + status: 405, + statusText: "Method Not Allowed", + }); + } + + const url = new URL(request.url); + if (url.pathname === HEALTH_PATH) { + return health(env); + } + + const key = resolveKey(url); + const track = shouldTrackView(request, url, key, env); + const rangeHeader = + request.method === "GET" ? request.headers.get("range") : null; + + // The edge cache stores one full 200 per object. Range requests are sliced + // out of it by cache.match itself — Cloudflare honours Range on a Cache API + // lookup and builds the 206 natively, so the Worker never touches a byte of + // video. Doing that slicing in JS is what burned 170-680ms of CPU per clip + // request and got the Worker killed mid-stream ("exceeded CPU time limit"), + // which browsers report as ERR_HTTP2_PROTOCOL_ERROR / ERR_QUIC_PROTOCOL_ERROR. + // + // Origin is in the cache key so each caller gets its own + // Access-Control-Allow-Origin (credentialed fetches can't take `*`). + const cache = caches.default; + const cacheUrl = `${request.url}#origin=${reqOrigin ?? ""}`; + const cacheKey = new Request(cacheUrl, { method: "GET" }); + const isHead = request.method === "HEAD"; + const lookupKey = rangeHeader + ? new Request(cacheUrl, { + method: "GET", + headers: { Range: rangeHeader }, + }) + : cacheKey; + const cached = await cache.match(lookupKey); + if (cached) { + return clientResponse(cached, { + isHead, + key: key ?? "", + track, + request, + env, + ctx, + }); + } + + if (!key) { + return new Response("No file provided", { status: 400 }); + } + + if (!env.BUCKET_NAME || !env.S3_ENDPOINT) { + return new Response( + "Worker misconfigured: BUCKET_NAME / S3_ENDPOINT not set", + { status: 500 }, + ); + } + if (!env.S3_ACCESS_KEY || !env.S3_SECRET) { + return new Response( + "Worker misconfigured: S3_ACCESS_KEY / S3_SECRET not set", + { status: 500 }, + ); + } + + // Always the full object, never a Range: one 200 per object is fetched from + // B2 and cached, and every later range is served out of that one copy. This + // is what keeps B2 egress and transactions flat no matter how many range + // requests a player makes. + const ttl = MUTABLE_TTL[key] ?? IMMUTABLE_TTL; + const upstream = await signedFetch( + "GET", + `https://${env.BUCKET_NAME}.${env.S3_ENDPOINT}/${key}`, + env, + ttl, + ); + + const headers = new Headers(upstream.headers); + + const requestedNameRaw = url.searchParams.get("name") ?? ""; + let requestedName = requestedNameRaw; + try { + requestedName = decodeURIComponent(requestedNameRaw); + } catch { + requestedName = requestedNameRaw; + } + const fallbackName = key.split("/").pop() ?? key; + const sanitize = (s: string) => s.replace(/[\\\/"\r\n]/g, "").trim(); + const filename = + sanitize(requestedName) || sanitize(fallbackName) || "clip.mp4"; + + // Default to inline so pasting a clip URL plays in the browser; + // ?dl=1 opts into attachment. + const wantDownload = + url.searchParams.get("dl") === "1" || + url.searchParams.get("download") === "1"; + headers.set( + "Content-Disposition", + `${wantDownload ? "attachment" : "inline"}; filename="${filename}"`, + ); + + if (!headers.has("Accept-Ranges")) { + headers.set("Accept-Ranges", "bytes"); + } + // Force long browser cache regardless of what B2 returned — second view + // of a clip serves from the user's disk cache and never hits the Worker. + headers.set( + "Cache-Control", + ttl === IMMUTABLE_TTL + ? `public, max-age=${IMMUTABLE_TTL}, immutable` + : `public, max-age=${ttl}`, + ); + for (const [k, v] of Object.entries(corsHeaders(reqOrigin))) { + headers.set(k, v); + } + + if (/\.mp4$/i.test(key)) { + // Stored on the cached copy, not patched on the way out, so range hits + // served straight from cache carry it too. + headers.set("Content-Type", "video/mp4"); + } + + const stored = new Response(upstream.body, { + headers, + status: upstream.status, + statusText: upstream.statusText, + }); + + if (upstream.status !== 200) { + return clientResponse(stored, { isHead, key, track, request, env, ctx }); + } + + // clone() rather than a manual tee: one branch fills the cache, the other + // answers this request. A miss costs one full pass; every subsequent range + // is served natively out of the stored object. + ctx.waitUntil(cache.put(cacheKey, stored.clone()).catch(() => {})); + + // This request is answered with the whole object even if it asked for a + // range. Serving a 200 to a Range request is explicitly allowed, and it + // avoids slicing in JS on the one path where the cache can't do it yet. + // Players re-request ranges afterwards and those hit the warm cache. + return clientResponse(stored, { isHead, key, track, request, env, ctx }); + }, +}; + +// Every write is still individually authorized by the API-minted HMAC token +// below; the prefix list only bounds which trees are writable at all. +// NOTE: a prefix here is necessary but NOT sufficient — the request only +// reaches the worker if one of its routes (ROUTE_PATHS in backblaze-proxy.sh) +// matches its path. +const UPLOAD_PREFIXES = ["demo-uploads/", "events/", "news/"]; +// B2 multipart part ceiling we accept; the API chunks at 64MiB so anything +// larger is a malformed/abusive request. +const MAX_PART_BYTES = 64 * 1024 * 1024; + +async function handleUpload( + request: Request, + env: { + S3_ACCESS_KEY: string; + S3_SECRET: string; + BUCKET_NAME: string; + S3_ENDPOINT: string; + }, + reqOrigin: string | null, +): Promise { + const url = new URL(request.url); + const key = url.pathname.replace(/^\/+/, ""); + const partNumber = url.searchParams.get("partNumber"); + const uploadId = url.searchParams.get("uploadId"); + const token = url.searchParams.get("token"); + + if ( + !UPLOAD_PREFIXES.some((prefix) => key.startsWith(prefix)) || + !partNumber || + !uploadId || + !token + ) { + return new Response("forbidden", { + status: 403, + headers: corsHeaders(reqOrigin), + }); + } + if ( + !env.BUCKET_NAME || + !env.S3_ENDPOINT || + !env.S3_ACCESS_KEY || + !env.S3_SECRET + ) { + return new Response("Worker misconfigured", { + status: 500, + headers: corsHeaders(reqOrigin), + }); + } + + // Authorize the part write: the API mints this token (HMAC over key+uploadId + // keyed on the shared S3_SECRET) only for authenticated admins. Without this + // check the worker would sign arbitrary writes for anyone with a valid + // uploadId. + if (!(await verifyUploadToken(env.S3_SECRET, token, key, uploadId))) { + return new Response("forbidden", { + status: 403, + headers: corsHeaders(reqOrigin), + }); + } + + // The streamed body needs an explicit, bounded Content-Length — we forward + // the client's value but never trust it blindly. + const contentLengthRaw = request.headers.get("content-length"); + if ( + !contentLengthRaw || + !/^\d+$/.test(contentLengthRaw) || + Number(contentLengthRaw) > MAX_PART_BYTES + ) { + return new Response("invalid content-length", { + status: 413, + headers: corsHeaders(reqOrigin), + }); + } + + const target = `https://${env.BUCKET_NAME}.${env.S3_ENDPOINT}/${key}?partNumber=${encodeURIComponent( + partNumber, + )}&uploadId=${encodeURIComponent(uploadId)}`; + + const signed = await new AwsClient({ + accessKeyId: env.S3_ACCESS_KEY, + secretAccessKey: env.S3_SECRET, + service: "s3", + }).sign(target, { + method: "PUT", + headers: { "x-amz-content-sha256": "UNSIGNED-PAYLOAD" }, + }); + + const headers = new Headers(signed.headers); + headers.set("content-length", contentLengthRaw); + + const upstream = await fetch(target, { + method: "PUT", + headers, + body: request.body, + }); + + const responseHeaders = new Headers(corsHeaders(reqOrigin)); + const etag = upstream.headers.get("etag"); + if (etag) { + responseHeaders.set("ETag", etag); + } + + if (upstream.ok) { + return new Response(null, { + status: upstream.status, + headers: responseHeaders, + }); + } + return new Response(await upstream.text(), { + status: upstream.status, + statusText: upstream.statusText, + headers: responseHeaders, + }); +} + +function base64urlToBytes(value: string): Uint8Array { + const b64 = value.replace(/-/g, "+").replace(/_/g, "/"); + const padded = b64.padEnd(Math.ceil(b64.length / 4) * 4, "="); + const binary = atob(padded); + const bytes = new Uint8Array(binary.length); + for (let i = 0; i < binary.length; i++) { + bytes[i] = binary.charCodeAt(i); + } + return bytes; +} + +// Verifies the API-minted token: HMAC-SHA256 over the `payload` segment, then +// confirms the payload is bound to this exact key+uploadId and not expired. +// crypto.subtle.verify is constant-time, avoiding signature-timing leaks. +async function verifyUploadToken( + secret: string, + token: string, + key: string, + uploadId: string, +): Promise { + const dot = token.lastIndexOf("."); + if (dot <= 0) return false; + const message = token.slice(0, dot); + const signature = token.slice(dot + 1); + + let signatureBytes: Uint8Array; + let payloadBytes: Uint8Array; + try { + signatureBytes = base64urlToBytes(signature); + payloadBytes = base64urlToBytes(message); + } catch { + return false; + } + + const cryptoKey = await crypto.subtle.importKey( + "raw", + new TextEncoder().encode(secret), + { name: "HMAC", hash: "SHA-256" }, + false, + ["verify"], + ); + const valid = await crypto.subtle.verify( + "HMAC", + cryptoKey, + signatureBytes, + new TextEncoder().encode(message), + ); + if (!valid) return false; + + let payload: { k?: unknown; u?: unknown; exp?: unknown }; + try { + payload = JSON.parse(new TextDecoder().decode(payloadBytes)); + } catch { + return false; + } + if (payload.k !== key || payload.u !== uploadId) return false; + if (typeof payload.exp !== "number" || payload.exp * 1000 < Date.now()) { + return false; + } + return true; +} + +function corsHeaders(reqOrigin: string | null): Record { + // Reflect the request's Origin so credentialed fetches (which can't + // accept Allow-Origin: *) work. Non-CORS requests (no Origin header) + // fall back to `*` — same as before, keeps anon