From 211b705a91fc893a511a0eba76f40825e55392c9 Mon Sep 17 00:00:00 2001 From: Jason Penilla <11360596+jpenilla@users.noreply.github.com> Date: Sun, 16 Aug 2026 13:09:20 -0700 Subject: [PATCH] feat(downloads): add live update backend Add Durable Object coordination, regional WebSocket shards, and versioned download snapshots for live build updates. --- package.json | 1 + src/downloads-live.ts | 321 +++++++++++++++++++++++++++++++ src/utils/download.test.ts | 66 +++++++ src/utils/download.ts | 50 ++++- src/utils/downloads-live-path.ts | 1 + src/worker.ts | 3 + worker-configuration.d.ts | 5 +- wrangler.jsonc | 18 ++ 8 files changed, 460 insertions(+), 5 deletions(-) create mode 100644 src/downloads-live.ts create mode 100644 src/utils/download.test.ts create mode 100644 src/utils/downloads-live-path.ts diff --git a/package.json b/package.json index 1d4c5af7..91a09960 100644 --- a/package.json +++ b/package.json @@ -14,6 +14,7 @@ "lint": "bun run lint:eslint && bun run lint:prettier", "check": "bun run check:types", "check:all": "bun run check && bun run lint", + "test": "bun test", "data:update": "bun scripts/update-data.ts work", "postinstall": "wrangler types" }, diff --git a/src/downloads-live.ts b/src/downloads-live.ts new file mode 100644 index 00000000..3a1fd45a --- /dev/null +++ b/src/downloads-live.ts @@ -0,0 +1,321 @@ +import { DurableObject } from "cloudflare:workers"; +import { + type DownloadProjectId, + type DownloadRegion, + type DownloadsPageSnapshot, + DOWNLOAD_REGIONS, + downloadRegionForContinent, + downloadsPageDataKvKey, + downloadsPageDataRevision, + fetchDownloadsPageData, + isDownloadProjectId, + isValidDownloadsPageData, +} from "./utils/download"; +import { DOWNLOADS_LIVE_PATH } from "./utils/downloads-live-path"; + +type CoordinatorState = DownloadsPageSnapshot & { project: DownloadProjectId }; +type PublishRequest = CoordinatorState & { region: DownloadRegion }; +type ClientMessage = { type: "hello" | "resync"; streamId: string; generation: number; revision: string }; + +const COMMITTED_KEY = "committed"; +const PENDING_KEY = "pending"; + +export function downloadRegionForRequest(request: Request): DownloadRegion { + const continent = request.cf?.continent; + return downloadRegionForContinent(typeof continent === "string" ? continent : undefined); +} + +export function coordinatorStub(env: Env, project: DownloadProjectId): DurableObjectStub { + return env.DOWNLOAD_UPDATES.getByName(`updates:${project}`); +} + +export function downloadShardStub(env: Env, project: DownloadProjectId, region: DownloadRegion): DurableObjectStub { + return env.DOWNLOAD_CLIENTS.getByName(`downloads:${project}:${region}`, { locationHint: region }); +} + +export async function requestDownloadsRefresh(env: Env, project: DownloadProjectId, deliveryId?: string): Promise { + const response = await coordinatorStub(env, project).fetch("https://downloads.internal/refresh", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ project, deliveryId }), + }); + if (!response.ok) throw new Error(`Downloads refresh failed for ${project}: ${response.status} ${await response.text()}`); +} + +export class DownloadsUpdateCoordinator extends DurableObject { + private refreshQueue: Promise = Promise.resolve(); + + override async fetch(request: Request): Promise { + const url = new URL(request.url); + if (request.method === "GET" && url.pathname === "/snapshot") { + const snapshot = await this.latestState(); + return snapshot ? Response.json(snapshot) : new Response(null, { status: 404 }); + } + + if (request.method !== "POST" || url.pathname !== "/refresh") return new Response("Not found", { status: 404 }); + + const body = (await request.json()) as { project?: unknown; deliveryId?: unknown }; + if (!isDownloadProjectId(body.project)) return new Response("Invalid project", { status: 400 }); + const project = body.project; + const deliveryId = typeof body.deliveryId === "string" ? body.deliveryId : undefined; + + const refresh = this.refreshQueue.then(() => this.refresh(project, deliveryId)); + this.refreshQueue = refresh.catch(() => undefined); + try { + await refresh; + return new Response(null, { status: 204 }); + } catch (error) { + log("downloads_refresh_failed", { project, deliveryId, error: errorMessage(error) }); + return new Response("Refresh failed", { status: 502 }); + } + } + + private async refresh(project: DownloadProjectId, deliveryId?: string): Promise { + const startedAt = Date.now(); + log("downloads_refresh_started", { project, deliveryId }); + + await this.recoverPending(); + const committed = await this.ctx.storage.get(COMMITTED_KEY); + const data = await fetchDownloadsPageData(project); + if (!isValidDownloadsPageData(data)) throw new Error("Fill returned an incomplete downloads snapshot"); + + const revision = await downloadsPageDataRevision(data); + if (committed?.revision === revision) { + await this.publishToRegions(committed); + log("downloads_refresh_unchanged", { + project, + deliveryId, + generation: committed.generation, + revision, + durationMs: Date.now() - startedAt, + }); + return; + } + + const next: CoordinatorState = { + project, + streamId: committed?.streamId ?? crypto.randomUUID(), + generation: (committed?.generation ?? 0) + 1, + revision, + data, + }; + await this.ctx.storage.put(PENDING_KEY, next); + await this.env.WEBSITE_CACHE.put(downloadsPageDataKvKey(project), JSON.stringify(snapshotWithoutProject(next))); + await this.ctx.storage.transaction(async (transaction) => { + await transaction.put(COMMITTED_KEY, next); + await transaction.delete(PENDING_KEY); + }); + await this.publishToRegions(next); + log("downloads_refresh_succeeded", { + project, + deliveryId, + generation: next.generation, + revision, + durationMs: Date.now() - startedAt, + }); + } + + private async recoverPending(): Promise { + const pending = await this.ctx.storage.get(PENDING_KEY); + if (!pending) return; + + await this.env.WEBSITE_CACHE.put(downloadsPageDataKvKey(pending.project), JSON.stringify(snapshotWithoutProject(pending))); + await this.ctx.storage.transaction(async (transaction) => { + await transaction.put(COMMITTED_KEY, pending); + await transaction.delete(PENDING_KEY); + }); + await this.publishToRegions(pending); + log("downloads_refresh_recovered", { + project: pending.project, + generation: pending.generation, + revision: pending.revision, + }); + } + + private async latestState(): Promise { + return this.ctx.storage.get(COMMITTED_KEY); + } + + private async publishToRegions(snapshot: CoordinatorState): Promise { + await Promise.all( + DOWNLOAD_REGIONS.map(async (region) => { + const startedAt = Date.now(); + let lastError: unknown; + for (let attempt = 1; attempt <= 2; attempt++) { + try { + const response = await downloadShardStub(this.env, snapshot.project, region).fetch("https://downloads.internal/publish", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ ...snapshot, region } satisfies PublishRequest), + }); + if (!response.ok) throw new Error(`Shard returned ${response.status}`); + log("downloads_shard_published", { + project: snapshot.project, + region, + generation: snapshot.generation, + attempt, + durationMs: Date.now() - startedAt, + }); + return; + } catch (error) { + lastError = error; + } + } + log("downloads_shard_publish_failed", { + project: snapshot.project, + region, + generation: snapshot.generation, + error: errorMessage(lastError), + }); + }) + ); + } +} + +export class DownloadsWebSocketShard extends DurableObject { + override async fetch(request: Request): Promise { + const url = new URL(request.url); + if (request.method === "POST" && url.pathname === "/publish") { + const snapshot = (await request.json()) as PublishRequest; + if (!isDownloadProjectId(snapshot.project) || !DOWNLOAD_REGIONS.includes(snapshot.region)) { + return new Response("Invalid snapshot", { status: 400 }); + } + await this.acceptSnapshot(snapshot); + return new Response(null, { status: 204 }); + } + + if (url.pathname !== DOWNLOADS_LIVE_PATH || request.headers.get("upgrade")?.toLowerCase() !== "websocket") { + return new Response("Expected WebSocket", { status: 426 }); + } + + const project = url.searchParams.get("project"); + const region = url.searchParams.get("region"); + if (!isDownloadProjectId(project) || !isDownloadRegion(region)) return new Response("Invalid subscription", { status: 400 }); + + const pair = new WebSocketPair(); + const [client, server] = Object.values(pair); + server.serializeAttachment({ project, region }); + this.ctx.acceptWebSocket(server); + log("download_ws_open", { project, region, connections: this.ctx.getWebSockets().length }); + return new Response(null, { status: 101, webSocket: client }); + } + + override async webSocketMessage(webSocket: WebSocket, message: string | ArrayBuffer): Promise { + if (typeof message !== "string") return; + let hello: ClientMessage; + try { + hello = JSON.parse(message) as ClientMessage; + } catch { + return; + } + if ( + (hello.type !== "hello" && hello.type !== "resync") || + typeof hello.streamId !== "string" || + !Number.isSafeInteger(hello.generation) || + typeof hello.revision !== "string" + ) { + return; + } + + const attachment = webSocket.deserializeAttachment() as { project: DownloadProjectId; region: DownloadRegion }; + let local = await this.ctx.storage.get(COMMITTED_KEY); + if (hello.type === "resync" || !local || local.streamId !== hello.streamId || local.generation < hello.generation) { + const response = await coordinatorStub(this.env, attachment.project).fetch("https://downloads.internal/snapshot"); + if (response.ok) { + const latest = (await response.json()) as CoordinatorState; + await this.acceptSnapshot({ ...latest, region: attachment.region }); + local = latest; + } + } + + if (!local) return; + if (hello.type === "resync") { + webSocket.send(JSON.stringify({ type: "snapshot", shard: attachment.region, authoritative: true, ...local })); + return; + } + if ( + local.streamId !== hello.streamId || + local.generation > hello.generation || + (local.generation === hello.generation && local.revision !== hello.revision) + ) { + webSocket.send(JSON.stringify({ type: "snapshot", shard: attachment.region, ...local })); + } + } + + override webSocketClose(webSocket: WebSocket, code: number, reason: string): void { + const attachment = webSocket.deserializeAttachment() as { project?: string; region?: string } | null; + log("download_ws_close", { + project: attachment?.project, + region: attachment?.region, + code, + reason, + connections: this.ctx.getWebSockets().length, + }); + webSocket.close(code, reason); + } + + override webSocketError(webSocket: WebSocket, error: unknown): void { + const attachment = webSocket.deserializeAttachment() as { project?: string; region?: string } | null; + log("download_ws_error", { project: attachment?.project, region: attachment?.region, error: errorMessage(error) }); + } + + private async acceptSnapshot(snapshot: PublishRequest): Promise { + const current = await this.ctx.storage.get(COMMITTED_KEY); + if (current?.streamId === snapshot.streamId && current.generation > snapshot.generation) return; + if (current?.streamId === snapshot.streamId && current.generation === snapshot.generation && current.revision === snapshot.revision) + return; + if (current?.streamId === snapshot.streamId && current.generation === snapshot.generation && current.revision !== snapshot.revision) { + throw new Error(`Conflicting revision for generation ${snapshot.generation}`); + } + + const stored: CoordinatorState = { + project: snapshot.project, + streamId: snapshot.streamId, + generation: snapshot.generation, + revision: snapshot.revision, + data: snapshot.data, + }; + await this.ctx.storage.put(COMMITTED_KEY, stored); + + const message = JSON.stringify({ type: "snapshot", shard: snapshot.region, ...stored }); + const sockets = this.ctx.getWebSockets(); + const startedAt = Date.now(); + let sent = 0; + let failed = 0; + for (const socket of sockets) { + try { + socket.send(message); + sent++; + } catch { + failed++; + } + } + log("download_broadcast", { + project: snapshot.project, + region: snapshot.region, + generation: snapshot.generation, + revision: snapshot.revision, + payloadBytes: new TextEncoder().encode(message).byteLength, + recipients: sockets.length, + sent, + failed, + durationMs: Date.now() - startedAt, + }); + } +} + +function isDownloadRegion(value: string | null): value is DownloadRegion { + return value !== null && DOWNLOAD_REGIONS.some((region) => region === value); +} + +function snapshotWithoutProject(snapshot: CoordinatorState): DownloadsPageSnapshot { + return { streamId: snapshot.streamId, generation: snapshot.generation, revision: snapshot.revision, data: snapshot.data }; +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} + +function log(event: string, fields: Record): void { + console.log(JSON.stringify({ event, ...fields })); +} diff --git a/src/utils/download.test.ts b/src/utils/download.test.ts new file mode 100644 index 00000000..df9ba394 --- /dev/null +++ b/src/utils/download.test.ts @@ -0,0 +1,66 @@ +import assert from "node:assert/strict"; +import { describe, test } from "node:test"; +import { downloadRegionForContinent, type DownloadsPageData, downloadsPageDataRevision } from "./download"; + +const data: DownloadsPageData = { + projectResult: { + value: { + id: "paper", + name: "Paper", + latestStableVersion: "1.21.5", + latestExperimentalVersion: null, + latestVersionGroup: "1.21", + }, + }, + stableBuildsResult: { value: { builds: [] } }, + experimentalBuildsResult: null, +}; + +describe("downloads snapshot revisions", () => { + test("is independent of object key insertion order", async () => { + const reordered = { + experimentalBuildsResult: null, + stableBuildsResult: { value: { builds: [] } }, + projectResult: { + value: { + name: "Paper", + id: "paper", + latestVersionGroup: "1.21", + latestExperimentalVersion: null, + latestStableVersion: "1.21.5", + }, + }, + } satisfies DownloadsPageData; + + assert.equal(await downloadsPageDataRevision(reordered), await downloadsPageDataRevision(data)); + }); + + test("changes when snapshot content changes", async () => { + const changed = structuredClone(data); + if (changed.projectResult.value) changed.projectResult.value.latestStableVersion = "1.21.6"; + assert.notEqual(await downloadsPageDataRevision(changed), await downloadsPageDataRevision(data)); + }); + + test("matches JSON serialization semantics for undefined object properties", async () => { + const withUndefined = structuredClone(data); + if (withUndefined.stableBuildsResult.value) withUndefined.stableBuildsResult.value.latest = undefined; + assert.equal(await downloadsPageDataRevision(withUndefined), await downloadsPageDataRevision(data)); + }); +}); + +describe("downloads regional routing", () => { + const cases = [ + ["NA", "wnam"], + ["SA", "wnam"], + ["EU", "weur"], + ["AF", "weur"], + ["AS", "apac"], + ["OC", "apac"], + ] as const; + + for (const [continent, expected] of cases) { + test(`maps ${continent} to ${expected}`, () => { + assert.equal(downloadRegionForContinent(continent), expected); + }); + } +}); diff --git a/src/utils/download.ts b/src/utils/download.ts index c1b80914..0f82ad57 100644 --- a/src/utils/download.ts +++ b/src/utils/download.ts @@ -4,6 +4,8 @@ import { type ProjectDescriptor, type Build, type Project } from "@/utils/types" export const DOWNLOAD_PROJECT_IDS = ["paper", "velocity", "waterfall", "folia"] as const; export type DownloadProjectId = (typeof DOWNLOAD_PROJECT_IDS)[number]; +export const DOWNLOAD_REGIONS = ["wnam", "weur", "apac"] as const; +export type DownloadRegion = (typeof DOWNLOAD_REGIONS)[number]; const DOWNLOAD_PROJECT_ID_SET: ReadonlySet = new Set(DOWNLOAD_PROJECT_IDS); @@ -11,6 +13,12 @@ export function isDownloadProjectId(value: unknown): value is DownloadProjectId return typeof value === "string" && DOWNLOAD_PROJECT_ID_SET.has(value); } +export function downloadRegionForContinent(continent?: string): DownloadRegion { + if (continent === "EU" || continent === "AF") return "weur"; + if (continent === "AS" || continent === "OC") return "apac"; + return "wnam"; +} + export type ProjectDescriptorOrError = { error?: string; value?: ProjectDescriptor }; export type ProjectBuildsOrError = { error?: string; value?: { latest?: Build; builds: Build[] } }; export type DownloadsPageData = { @@ -19,21 +27,55 @@ export type DownloadsPageData = { experimentalBuildsResult: ProjectBuildsOrError | null; }; +export type DownloadsPageSnapshot = { + streamId: string; + generation: number; + revision: string; + data: DownloadsPageData; +}; + +export type DownloadsLiveStatus = "connecting" | "live" | "reconnecting" | "paused" | "offline"; + export function downloadsPageDataKvKey(projectId: string) { return `downloads:${projectId}`; } -export async function refreshDownloadsPageCache(projectId: string, kv: KVNamespace): Promise { - const data = await fetchDownloadsPageData(projectId); - if ( +export function isValidDownloadsPageData(data: DownloadsPageData): boolean { + return ( data.projectResult.error === undefined && data.stableBuildsResult.error === undefined && data.experimentalBuildsResult?.error === undefined - ) { + ); +} + +export async function downloadsPageDataRevision(data: DownloadsPageData): Promise { + const digest = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(canonicalJson(data))); + return `sha256:${Array.from(new Uint8Array(digest), (byte) => byte.toString(16).padStart(2, "0")).join("")}`; +} + +export async function refreshDownloadsPageCache(projectId: string, kv: KVNamespace): Promise { + const data = await fetchDownloadsPageData(projectId); + if (isValidDownloadsPageData(data)) { await kv.put(downloadsPageDataKvKey(projectId), JSON.stringify(data)); } } +function canonicalJson(value: unknown): string { + return JSON.stringify(sortJsonValue(value)); +} + +function sortJsonValue(value: unknown): unknown { + if (value === null || typeof value !== "object") return value; + if (Array.isArray(value)) return value.map(sortJsonValue); + + const record = value as Record; + return Object.fromEntries( + Object.keys(record) + .sort() + .map((key) => [key, sortJsonValue(record[key])]) + ); +} + export async function fetchDownloadsPageData(projectId: string, kv?: KVNamespace): Promise { if (kv) { const cachedString = await kv.get(downloadsPageDataKvKey(projectId)); diff --git a/src/utils/downloads-live-path.ts b/src/utils/downloads-live-path.ts new file mode 100644 index 00000000..8c880762 --- /dev/null +++ b/src/utils/downloads-live-path.ts @@ -0,0 +1 @@ +export const DOWNLOADS_LIVE_PATH = "/internal-api/downloads/live"; diff --git a/src/worker.ts b/src/worker.ts index 5ad40fc0..c5ed3b7b 100644 --- a/src/worker.ts +++ b/src/worker.ts @@ -1,6 +1,9 @@ import { handle } from "@astrojs/cloudflare/handler"; import { DOWNLOAD_PROJECT_IDS, refreshDownloadsPageCache } from "./utils/download"; import { PAPER_PLAYERCOUNT_KEY, fetchPaperBstatsPlayerCount } from "./utils/bstats"; +import { DownloadsUpdateCoordinator, DownloadsWebSocketShard } from "./downloads-live"; + +export { DownloadsUpdateCoordinator, DownloadsWebSocketShard }; const PLAYER_COUNT_CRON = "*/10 * * * *"; const DOWNLOADS_RECONCILIATION_CRON = "0 * * * *"; diff --git a/worker-configuration.d.ts b/worker-configuration.d.ts index 74082897..3f606d2a 100644 --- a/worker-configuration.d.ts +++ b/worker-configuration.d.ts @@ -1,14 +1,17 @@ /* eslint-disable */ -// Generated by Wrangler by running `wrangler types` (hash: 2ca44fc19c1df203233058883b303360) +// Generated by Wrangler by running `wrangler types` (hash: ec36602ee029116c34a130f2c0cffb96) // Runtime types generated with workerd@1.20260708.1 2026-02-19 disable_nodejs_process_v2,nodejs_compat interface __BaseEnv_Env { WEBSITE_CACHE: KVNamespace; ASSETS: Fetcher; FILL_WEBHOOK_SECRET: string; + DOWNLOAD_UPDATES: DurableObjectNamespace; + DOWNLOAD_CLIENTS: DurableObjectNamespace; } declare namespace Cloudflare { interface GlobalProps { mainModule: typeof import("./src/worker"); + durableNamespaces: "DownloadsUpdateCoordinator" | "DownloadsWebSocketShard"; } interface Env extends __BaseEnv_Env {} } diff --git a/wrangler.jsonc b/wrangler.jsonc index 1d1c0d63..563b2796 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -14,6 +14,24 @@ "id": "99e81ab8e2c74620b115ca861ebae59a", }, ], + "durable_objects": { + "bindings": [ + { + "name": "DOWNLOAD_UPDATES", + "class_name": "DownloadsUpdateCoordinator", + }, + { + "name": "DOWNLOAD_CLIENTS", + "class_name": "DownloadsWebSocketShard", + }, + ], + }, + "migrations": [ + { + "tag": "v1", + "new_sqlite_classes": ["DownloadsUpdateCoordinator", "DownloadsWebSocketShard"], + }, + ], "triggers": { "crons": ["*/10 * * * *", "0 * * * *"], },