From 61e3d3de5a45c2cdecdac59c9371cd7ea86bce6d Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 08:23:45 -0400 Subject: [PATCH 01/10] feature: playcast edge relay worker and deploy script --- .github/workflows/cloudflare-workers.yml | 18 ++ cloudflare-workers/playcast-relay/.gitignore | 2 + .../playcast-relay/package.json | 4 + cloudflare-workers/playcast-relay/worker.js | 165 +++++++++++++ .../playcast-relay/worker.test.mjs | 225 ++++++++++++++++++ .../playcast-relay/wrangler.toml | 17 ++ playcast-relay.sh | 60 +++++ 7 files changed, 491 insertions(+) create mode 100644 .github/workflows/cloudflare-workers.yml create mode 100644 cloudflare-workers/playcast-relay/.gitignore create mode 100644 cloudflare-workers/playcast-relay/package.json create mode 100644 cloudflare-workers/playcast-relay/worker.js create mode 100644 cloudflare-workers/playcast-relay/worker.test.mjs create mode 100644 cloudflare-workers/playcast-relay/wrangler.toml create mode 100755 playcast-relay.sh diff --git a/.github/workflows/cloudflare-workers.yml b/.github/workflows/cloudflare-workers.yml new file mode 100644 index 0000000..a896e21 --- /dev/null +++ b/.github/workflows/cloudflare-workers.yml @@ -0,0 +1,18 @@ +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: node --test "cloudflare-workers/**/*.test.mjs" diff --git a/cloudflare-workers/playcast-relay/.gitignore b/cloudflare-workers/playcast-relay/.gitignore new file mode 100644 index 0000000..cc27ab8 --- /dev/null +++ b/cloudflare-workers/playcast-relay/.gitignore @@ -0,0 +1,2 @@ +.wrangler +node_modules diff --git a/cloudflare-workers/playcast-relay/package.json b/cloudflare-workers/playcast-relay/package.json new file mode 100644 index 0000000..094aa89 --- /dev/null +++ b/cloudflare-workers/playcast-relay/package.json @@ -0,0 +1,4 @@ +{ + "type": "module", + "private": true +} diff --git a/cloudflare-workers/playcast-relay/worker.js b/cloudflare-workers/playcast-relay/worker.js new file mode 100644 index 0000000..838745c --- /dev/null +++ b/cloudflare-workers/playcast-relay/worker.js @@ -0,0 +1,165 @@ +// Edge cache for a 5stack panel's Playcast relay, deployed as a route on the +// relay domain itself (tv.example.com/*). +// +// Game servers' posts pass straight through to the panel. Viewers' reads (CS2 +// clients and the game streamer) are cached here, so every fragment leaves the +// panel once per Cloudflare location rather than once per viewer. A route +// worker's fetch to its own hostname goes to the origin, never back into the +// worker. +// +// Fragment numbers start over when a new map starts a new broadcast, so a +// fragment is only cached under the broadcast token the panel reports on +// /sync. The same url can never be served data from an earlier broadcast. + +const NAME = "5stack-playcast-relay"; +const VERSION = "2"; + +const SYNC_CACHE_CONTROL = "public, max-age=3"; +const FRAGMENT_CACHE_CONTROL = "public, max-age=31536000, immutable"; +const TOKEN_HEADER = "X-Broadcast-Token"; +const CACHED_FIELDS = new Set(["full", "delta"]); + +export default { + async fetch(request, env, ctx) { + if (request.method !== "GET" && request.method !== "HEAD") { + return fetch(request); + } + + const response = await route(request, ctx); + + return request.method === "HEAD" ? new Response(null, response) : response; + }, +}; + +async function route(request, ctx) { + const url = new URL(request.url); + + // What the panel's settings page reads to tell the edge relay is live. Only + // the worker answers it: the panel's relay has no such path. + if (url.pathname === "/health") { + return Response.json( + { ok: true, worker: NAME, version: VERSION }, + { + headers: { + "Access-Control-Allow-Origin": "*", + "Cache-Control": "no-store", + }, + }, + ); + } + + const origin = url.origin; + const parts = url.pathname.split("/").filter(Boolean); + const [matchId] = parts; + + if (!matchId) { + return passThrough(`${origin}${url.pathname}${url.search}`); + } + + if (parts.length === 2 && parts[1] === "sync") { + return sync(url, origin, matchId, ctx); + } + + if ( + (parts.length === 3 || parts.length === 4) && + CACHED_FIELDS.has(parts[parts.length - 1]) + ) { + return fragment(url, origin, parts, ctx); + } + + return passThrough(`${origin}${url.pathname}${url.search}`); +} + +async function sync(url, origin, matchId, ctx) { + const key = new Request(url.toString()); + const cached = await caches.default.match(key); + if (cached) { + return cached; + } + + const response = await fetch(`${origin}/${matchId}/sync${url.search}`); + const result = await decoded(response, SYNC_CACHE_CONTROL); + + ctx.waitUntil(caches.default.put(key, result.clone())); + + return result; +} + +async function fragment(url, origin, parts, ctx) { + const [matchId] = parts; + const [fragmentIndex, field] = parts.slice(-2); + + const token = + parts.length === 4 + ? parts[1] + : await broadcastToken(url, matchId, origin, ctx); + + if (!token) { + return passThrough(`${origin}${url.pathname}${url.search}`); + } + + const path = `/${matchId}/${token}/${fragmentIndex}/${field}`; + const key = new Request(`${url.origin}${path}`); + + const scoped = parts.length === 4; + + const cached = await caches.default.match(key); + if (cached) { + return scoped ? cached : withCacheControl(cached, "no-store"); + } + + const response = await fetch(`${origin}${path}`); + + if (response.status !== 200) { + return decoded(response, "no-store"); + } + + const result = await decoded(response, FRAGMENT_CACHE_CONTROL); + + ctx.waitUntil(caches.default.put(key, result.clone())); + + // Only the token-scoped url names this data forever. The same fragment + // number without a token is a different fragment once a new map starts, so + // nothing past this worker may keep it. + return scoped ? result : withCacheControl(result, "no-store"); +} + +// Clients request fragments without the token, so it is read off the sync the +// panel serves (and this worker caches) for the match. +async function broadcastToken(url, matchId, origin, ctx) { + const response = await sync( + new URL(`${url.origin}/${matchId}/sync`), + origin, + matchId, + ctx, + ); + + return response.headers.get(TOKEN_HEADER); +} + +function withCacheControl(response, cacheControl) { + const result = new Response(response.body, response); + result.headers.set("Cache-Control", cacheControl); + return result; +} + +async function passThrough(target) { + return decoded(await fetch(target), "no-store"); +} + +// The panel gzips fragments, and fetch hands the worker the decoded body while +// keeping Content-Encoding. Passing that header on makes the runtime compress +// the body again on the way out of the cache, so viewers get it gzipped twice. +// Plain bytes are what Valve's relay serves and what CS2 reads. +async function decoded(response, cacheControl) { + const headers = new Headers(response.headers); + headers.delete("Content-Encoding"); + headers.delete("Content-Length"); + headers.set("Cache-Control", cacheControl); + + return new Response(await response.arrayBuffer(), { + status: response.status, + statusText: response.statusText, + headers, + }); +} diff --git a/cloudflare-workers/playcast-relay/worker.test.mjs b/cloudflare-workers/playcast-relay/worker.test.mjs new file mode 100644 index 0000000..073b069 --- /dev/null +++ b/cloudflare-workers/playcast-relay/worker.test.mjs @@ -0,0 +1,225 @@ +import { beforeEach, describe, it } from "node:test"; +import assert from "node:assert/strict"; +import worker from "./worker.js"; + +// The worker runs as a route on the relay domain, so the panel it fronts is +// the same host: its fetches to it go to the origin. +const RELAY = "https://tv.example.com"; + +let originRequests; +let originRoutes; +let cache; + +function respond(status, body = "", headers = {}) { + return new Response(status === 204 ? null : body, { status, headers }); +} + +beforeEach(() => { + originRequests = []; + originRoutes = new Map(); + cache = new Map(); + + globalThis.fetch = async (input) => { + const target = typeof input === "string" ? input : input.url; + originRequests.push(input); + const route = originRoutes.get(target); + return route ? route(input) : respond(404); + }; + + globalThis.caches = { + default: { + async match(request) { + const hit = cache.get(request.url); + return hit ? hit.clone() : undefined; + }, + async put(request, response) { + cache.set(request.url, response.clone()); + }, + }, + }; +}); + +const originUrls = () => + originRequests.map((input) => + typeof input === "string" ? input : input.url, + ); + +async function call(request) { + const pending = []; + const response = await worker.fetch( + request, + {}, + { waitUntil: (promise) => pending.push(promise) }, + ); + await Promise.all(pending); + return response; +} + +const get = (path) => call(new Request(`${RELAY}${path}`)); + +function broadcast(token, fragments = {}) { + originRoutes.set(`${RELAY}/match-1/sync`, () => + respond(200, JSON.stringify({ fragment: 42 }), { + "X-Broadcast-Token": token, + }), + ); + for (const [path, body] of Object.entries(fragments)) { + originRoutes.set(`${RELAY}/match-1/${token}/${path}`, () => + respond(200, body), + ); + } +} + +describe("playcast relay worker", () => { + it("tells the panel it is live, from any page", async () => { + const response = await get("/health"); + + assert.equal(response.status, 200); + assert.equal(response.headers.get("Access-Control-Allow-Origin"), "*"); + assert.deepEqual(await response.json(), { + ok: true, + worker: "5stack-playcast-relay", + version: "2", + }); + assert.equal(originRequests.length, 0); + }); + + it("passes a game server's post straight through to the panel", async () => { + originRoutes.set( + `${RELAY}/match-1/s1t1/45/full?tick=100`, + async (request) => + respond( + request.method === "POST" && + request.headers.get("x-origin-auth") === "match-1:secret" && + (await request.text()) === "fragment-bytes" + ? 200 + : 400, + ), + ); + + const response = await call( + new Request(`${RELAY}/match-1/s1t1/45/full?tick=100`, { + method: "POST", + headers: { "x-origin-auth": "match-1:secret" }, + body: "fragment-bytes", + }), + ); + + assert.equal(response.status, 200); + assert.equal(cache.size, 0); + }); + + it("serves sync from the panel and keeps it for a few seconds", async () => { + originRoutes.set(`${RELAY}/match-1/sync?fragment=0`, () => + respond(200, JSON.stringify({ fragment: 42 })), + ); + + const first = await get("/match-1/sync?fragment=0"); + assert.equal(first.status, 200); + assert.equal(first.headers.get("Cache-Control"), "public, max-age=3"); + assert.deepEqual(originUrls(), [`${RELAY}/match-1/sync?fragment=0`]); + + await get("/match-1/sync?fragment=0"); + assert.equal(originRequests.length, 1); + }); + + it("fetches a fragment once and serves every other viewer from the edge", async () => { + broadcast("s1t1", { "45/full": "full-45" }); + + assert.equal(await (await get("/match-1/45/full")).text(), "full-45"); + assert.equal(await (await get("/match-1/45/full")).text(), "full-45"); + assert.equal(await (await get("/match-1/s1t1/45/full")).text(), "full-45"); + + assert.equal( + originUrls().filter((url) => url.endsWith("/45/full")).length, + 1, + ); + }); + + it("lets only the token-scoped url be cached past the edge", async () => { + broadcast("s1t1", { "45/full": "full-45" }); + + const unscoped = await get("/match-1/45/full"); + assert.equal(unscoped.headers.get("Cache-Control"), "no-store"); + + const scoped = await get("/match-1/s1t1/45/full"); + assert.match(scoped.headers.get("Cache-Control"), /immutable/); + + const unscopedHit = await get("/match-1/45/full"); + assert.equal(unscopedHit.headers.get("Cache-Control"), "no-store"); + }); + + it("never passes on the panel's Content-Encoding", async () => { + broadcast("s1t1"); + originRoutes.set(`${RELAY}/match-1/s1t1/45/full`, () => + respond(200, "full-45", { "Content-Encoding": "gzip" }), + ); + + const response = await get("/match-1/s1t1/45/full"); + + assert.equal(response.headers.get("Content-Encoding"), null); + assert.equal(await response.text(), "full-45"); + }); + + it("answers HEAD without a body", async () => { + broadcast("s1t1", { "45/full": "full-45" }); + + const response = await call( + new Request(`${RELAY}/match-1/s1t1/45/full`, { method: "HEAD" }), + ); + + assert.equal(response.status, 200); + assert.equal(await response.text(), ""); + }); + + it("never serves a new map the fragment an earlier map had at that number", async () => { + broadcast("s1t1", { "3/full": "old-map" }); + assert.equal(await (await get("/match-1/3/full")).text(), "old-map"); + + cache.delete(`${RELAY}/match-1/sync`); + broadcast("s1t2", { "3/full": "new-map" }); + + assert.equal(await (await get("/match-1/3/full")).text(), "new-map"); + }); + + it("uses the token a client already has in its url", async () => { + broadcast("s1t1", { "45/delta": "delta-45" }); + + const response = await get("/match-1/s1t1/45/delta"); + + assert.equal(await response.text(), "delta-45"); + assert.ok(!originUrls().some((url) => url.endsWith("/sync"))); + }); + + it("does not cache a fragment the panel does not have yet", async () => { + broadcast("s1t1"); + + const missing = await get("/match-1/46/full"); + assert.equal(missing.status, 404); + assert.equal(missing.headers.get("Cache-Control"), "no-store"); + + originRoutes.set(`${RELAY}/match-1/s1t1/46/full`, () => + respond(200, "full-46"), + ); + assert.equal(await (await get("/match-1/46/full")).text(), "full-46"); + }); + + it("passes start through without caching it", async () => { + originRoutes.set(`${RELAY}/match-1/42/start`, () => + respond(200, "start-42"), + ); + + const response = await get("/match-1/42/start"); + + assert.equal(await response.text(), "start-42"); + assert.equal(response.headers.get("Cache-Control"), "no-store"); + assert.equal(cache.size, 0); + }); + + it("answers fragment requests from the panel when no broadcast is running", async () => { + const response = await get("/match-1/45/full"); + + assert.equal(response.status, 404); + assert.ok(originUrls().includes(`${RELAY}/match-1/45/full`)); + }); +}); diff --git a/cloudflare-workers/playcast-relay/wrangler.toml b/cloudflare-workers/playcast-relay/wrangler.toml new file mode 100644 index 0000000..222df8c --- /dev/null +++ b/cloudflare-workers/playcast-relay/wrangler.toml @@ -0,0 +1,17 @@ +# Edge cache for a 5stack panel's Playcast relay. Deploy it with +# ./playcast-relay.sh from the panel root, which puts it on the relay domain +# from overlays/config/api-config.env (RELAY_DOMAIN). By hand, from the panel +# root: +# +# npx wrangler deploy --config cloudflare-workers/playcast-relay/wrangler.toml \ +# --route "tv.example.com/*" +# +# The relay domain has to be proxied through Cloudflare (orange cloud). Set the +# route's request limit failure mode to "Fail open" in the Cloudflare dashboard, +# so broadcasts keep flowing straight to the panel if the Workers daily request +# limit is ever reached. + +name = "5stack-playcast-relay" +main = "worker.js" +compatibility_date = "2025-04-29" +workers_dev = false diff --git a/playcast-relay.sh b/playcast-relay.sh new file mode 100755 index 0000000..f9d2e21 --- /dev/null +++ b/playcast-relay.sh @@ -0,0 +1,60 @@ +#!/bin/bash + +# Deploys the Playcast edge relay: a Cloudflare Worker that runs as a route on +# the panel's relay domain (RELAY_DOMAIN). Game servers' uploads pass straight +# through it to the panel, and Playcast viewers are served from Cloudflare's +# cache. +# +# The relay domain has to be proxied through Cloudflare (orange cloud). +# Wrangler signs in to Cloudflare in a browser the first time; on a machine +# without one, export CLOUDFLARE_API_TOKEN (Workers Scripts: Edit and Workers +# Routes: Edit) first. + +PANEL_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +source "$PANEL_DIR/utils/colors.sh" +source "$PANEL_DIR/utils/print_domains_and_hosts.sh" + +load_domains_and_hosts + +if [ -z "$RELAY_DOMAIN" ]; then + err "RELAY_DOMAIN is not set in overlays/config/api-config.env. Run install.sh first." + exit 1 +fi + +if ! command -v npx >/dev/null 2>&1; then + err "npx was not found. Install Node.js (https://nodejs.org), or run this from a machine that has it." + exit 1 +fi + +step "Checking that $RELAY_DOMAIN goes through Cloudflare" +if ! curl -sI --max-time 10 "https://$RELAY_DOMAIN/" | grep -qi "^server: cloudflare"; then + err "$RELAY_DOMAIN is not proxied through Cloudflare." + err "Turn on the proxy (orange cloud) for its DNS record in Cloudflare, then run this again." + exit 1 +fi +ok "$RELAY_DOMAIN is proxied through Cloudflare" + +step "Deploying the Playcast edge relay to $RELAY_DOMAIN" +if ! npx --yes wrangler@4 deploy \ + --config "$PANEL_DIR/cloudflare-workers/playcast-relay/wrangler.toml" \ + --route "$RELAY_DOMAIN/*"; then + err "The deploy failed. See the wrangler output above." + exit 1 +fi + +# Only the worker answers /health; the panel's own relay has no such path. +step "Waiting for it to answer on https://$RELAY_DOMAIN/health" +for _ in $(seq 1 24); do + if curl -fsS --max-time 10 "https://$RELAY_DOMAIN/health" 2>/dev/null | grep -q '"worker":"5stack-playcast-relay"'; then + ok "The Playcast edge relay is active on $RELAY_DOMAIN" + warn "One last step in the Cloudflare dashboard: set this route's request limit" + warn "failure mode to \"Fail open\", so broadcasts keep reaching the panel if the" + warn "daily Workers limit is ever reached." + exit 0 + fi + sleep 5 +done + +err "https://$RELAY_DOMAIN/health is not answering from the worker yet." +err "Check the worker's route in the Cloudflare dashboard." +exit 1 From 0752e1137a6223eb8dd40f693c87a56aebddcb3f Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 11:22:54 -0400 Subject: [PATCH 02/10] feature: backblaze proxy worker and deploy script --- .github/workflows/cloudflare-workers.yml | 3 +- backblaze-proxy.sh | 119 ++++ cloudflare-workers/backblaze-proxy/.gitignore | 3 + .../backblaze-proxy/index.test.ts | 98 +++ cloudflare-workers/backblaze-proxy/index.ts | 581 ++++++++++++++++++ .../backblaze-proxy/package-lock.json | 19 + .../backblaze-proxy/package.json | 8 + .../backblaze-proxy/wrangler.toml | 25 + 8 files changed, 855 insertions(+), 1 deletion(-) create mode 100755 backblaze-proxy.sh create mode 100644 cloudflare-workers/backblaze-proxy/.gitignore create mode 100644 cloudflare-workers/backblaze-proxy/index.test.ts create mode 100644 cloudflare-workers/backblaze-proxy/index.ts create mode 100644 cloudflare-workers/backblaze-proxy/package-lock.json create mode 100644 cloudflare-workers/backblaze-proxy/package.json create mode 100644 cloudflare-workers/backblaze-proxy/wrangler.toml diff --git a/.github/workflows/cloudflare-workers.yml b/.github/workflows/cloudflare-workers.yml index a896e21..c505a48 100644 --- a/.github/workflows/cloudflare-workers.yml +++ b/.github/workflows/cloudflare-workers.yml @@ -15,4 +15,5 @@ jobs: - uses: actions/setup-node@v4 with: node-version: 24 - - run: node --test "cloudflare-workers/**/*.test.mjs" + - 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..5eacbe9 --- /dev/null +++ b/backblaze-proxy.sh @@ -0,0 +1,119 @@ +#!/bin/bash + +# Deploys the Backblaze proxy: a Cloudflare Worker in front of the S3 bucket +# (Backblaze B2) that serves demos, clips, news images and map assets through +# Cloudflare, so B2 egress is free and popular files come from the edge. The +# bucket, endpoint and keys come from the panel's config; the hostname to put +# it on is asked for (or given as the first argument). +# +# The hostname has to be on a domain proxied through Cloudflare (orange cloud). +# Wrangler signs in to Cloudflare in a browser the first time; on a machine +# without one, export CLOUDFLARE_API_TOKEN (Workers Scripts: Edit and Workers +# Routes: Edit) first. +# +# Panels keeping their secrets in Vault can export S3_ACCESS_KEY and S3_SECRET +# before running this instead. + +PANEL_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +source "$PANEL_DIR/utils/colors.sh" +source "$PANEL_DIR/utils/print_domains_and_hosts.sh" + +WORKER_DIR="$PANEL_DIR/cloudflare-workers/backblaze-proxy" + +read_env() { + grep -h "^$1=" "$2" 2>/dev/null | cut -d '=' -f2- +} + +load_domains_and_hosts + +S3_BUCKET="${S3_BUCKET:-$(read_env S3_BUCKET "$PANEL_DIR/overlays/config/s3-config.env")}" +S3_ENDPOINT="${S3_ENDPOINT:-$(read_env S3_ENDPOINT "$PANEL_DIR/overlays/config/s3-config.env")}" +S3_ENDPOINT="${S3_ENDPOINT#*://}" +S3_ENDPOINT="${S3_ENDPOINT%/}" +S3_ACCESS_KEY="${S3_ACCESS_KEY:-$(read_env S3_ACCESS_KEY "$PANEL_DIR/overlays/local-secrets/s3-secrets.env")}" +S3_SECRET="${S3_SECRET:-$(read_env S3_SECRET "$PANEL_DIR/overlays/local-secrets/s3-secrets.env")}" + +if [ -z "$S3_BUCKET" ] || [ -z "$S3_ENDPOINT" ]; then + err "S3_BUCKET and S3_ENDPOINT have to be set in overlays/config/s3-config.env." + exit 1 +fi + +# The worker reaches the bucket at https://., so the in-cluster +# storage the panel ships with cannot sit behind it. +if [[ "$S3_ENDPOINT" != *.* ]]; then + err "S3_ENDPOINT ($S3_ENDPOINT) is not a public S3 host. This worker is for a" + err "remote bucket such as Backblaze B2 (e.g. s3.us-east-005.backblazeb2.com)." + exit 1 +fi + +if [ -z "$S3_ACCESS_KEY" ] || [ -z "$S3_SECRET" ]; then + err "S3_ACCESS_KEY and S3_SECRET were not found in overlays/local-secrets/s3-secrets.env." + err "Export them before running this if your secrets live in Vault." + exit 1 +fi + +if ! command -v npx >/dev/null 2>&1; then + err "npx was not found. Install Node.js (https://nodejs.org), or run this from a machine that has it." + exit 1 +fi + +WORKER_HOST="$1" +if [ -z "$WORKER_HOST" ]; then + DEFAULT_HOST="demo-dl.${WEB_DOMAIN:-example.com}" + read -r -p "Hostname for the worker [$DEFAULT_HOST]: " WORKER_HOST + WORKER_HOST="${WORKER_HOST:-$DEFAULT_HOST}" +fi + +step "Checking that $WORKER_HOST goes through Cloudflare" +if ! curl -sI --max-time 10 "https://$WORKER_HOST/" | grep -qi "^server: cloudflare"; then + err "$WORKER_HOST is not proxied through Cloudflare." + err "Add a DNS record for it with the proxy on (orange cloud), then run this again." + exit 1 +fi +ok "$WORKER_HOST is proxied through Cloudflare" + +step "About to deploy the Backblaze proxy" +ok "Hostname: https://$WORKER_HOST" +ok "Bucket: $S3_BUCKET at $S3_ENDPOINT" +ok "API: https://$API_DOMAIN" +read -r -p "Deploy it? [y/N] " CONFIRM +if [[ ! "$CONFIRM" =~ ^[Yy]$ ]]; then + warn "Nothing deployed." + exit 0 +fi + +step "Installing the worker's dependencies" +if ! npm ci --silent --no-audit --no-fund --prefix "$WORKER_DIR"; then + err "npm ci failed. See the output above." + exit 1 +fi + +step "Deploying to $WORKER_HOST" +if ! npx --yes wrangler@4 deploy \ + --config "$WORKER_DIR/wrangler.toml" \ + --var "BUCKET_NAME:$S3_BUCKET" \ + --var "S3_ENDPOINT:$S3_ENDPOINT" \ + --var "API_URL:https://$API_DOMAIN" \ + --route "$WORKER_HOST/demo*" \ + --route "$WORKER_HOST/clips*" \ + --route "$WORKER_HOST/news*" \ + --route "$WORKER_HOST/maps*"; then + err "The deploy failed. See the wrangler output above." + exit 1 +fi + +step "Setting 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 ! npx --yes wrangler@4 secret bulk "$SECRETS_FILE" --config "$WORKER_DIR/wrangler.toml"; then + err "Setting the secrets failed. See the wrangler output above." + exit 1 +fi + +ok "The Backblaze proxy is deployed on https://$WORKER_HOST" +warn "Last step: in the panel, open Settings -> Application -> Demo settings and" +warn "set the Cloudflare Worker URL to https://$WORKER_HOST" 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..93a6edd --- /dev/null +++ b/cloudflare-workers/backblaze-proxy/index.test.ts @@ -0,0 +1,98 @@ +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"); + }); + } +}); diff --git a/cloudflare-workers/backblaze-proxy/index.ts b/cloudflare-workers/backblaze-proxy/index.ts new file mode 100644 index 0000000..beacc98 --- /dev/null +++ b/cloudflare-workers/backblaze-proxy/index.ts @@ -0,0 +1,581 @@ +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; + } + await response.body?.cancel(); + } + return response!; +} + +const VIEW_FRACTION = 0.5; +const BOT_UA = + /bot|crawl|spider|facebookexternalhit|slack|discord|telegram|whatsapp|preview|unfurl|embed|scrape|metainspector|skype|vkshare|redditbot|pinterest|googlebot|bingbot/i; + +// Whether a served response spans the file's VIEW_FRACTION mark, decided from +// Content-Range rather than by counting bytes through the stream. Counting +// meant piping every byte of an 18 MB clip through JS, which is what pushed +// the Worker over its CPU limit; this reads three numbers off a header. +// Tail seeks and end-of-file moov probes still do not qualify. Repeat beacons +// are deduped API-side per viewer. +function spansViewThreshold(status: number, headers: Headers): boolean { + if (status === 200) { + return true; + } + const contentRange = headers.get("content-range"); + if (!contentRange) { + return false; + } + const match = /^bytes\s+(\d+)-(\d+)\/(\d+)$/.exec(contentRange.trim()); + if (!match) { + return false; + } + const start = Number(match[1]); + const end = Number(match[2]); + const total = Number(match[3]); + if (!Number.isFinite(total) || total <= 0) { + return false; + } + const threshold = total * VIEW_FRACTION; + return start <= threshold && end >= 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); + 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 a route pattern in wrangler.toml 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