From 9f3bcf380eb898265953c342de8a74c4474b595c Mon Sep 17 00:00:00 2001 From: Adam Shiervani Date: Fri, 18 Sep 2026 22:46:05 +0200 Subject: [PATCH 1/6] fix(releases): serve staged releases before any release reaches 100% (#78) A prefix with no release at 100% made the default lookup throw a 500 before eligibility was checked, so no device got the staged build. That is every device for a new prefix until its first release is fully rolled out, and every JetKVM device the day no app or system row sits at 100%. The default lookup now returns null when nothing is at 100%. A device inside the rollout bucket gets the staged release as before; a device outside it gets a 404 saying no release is rolled out yet for its SKU, instead of a 500. Found by Bugbot on the release PR (#77). --- src/releases.ts | 18 +++++++++++++----- test/releases.test.ts | 37 +++++++++++++++++++++++++++++++------ 2 files changed, 44 insertions(+), 11 deletions(-) diff --git a/src/releases.ts b/src/releases.ts index d34ba83..fb0293f 100644 --- a/src/releases.ts +++ b/src/releases.ts @@ -471,16 +471,18 @@ function dbReleaseToMetadata(release: DbRelease, sku: string): ReleaseMetadata { }; } -async function getDefaultRelease(prefix: string, sku: string): Promise { +/** + * Newest fully rolled out release for the prefix, or null when no release + * has reached 100% yet (a prefix whose first release is still staged). + */ +async function getDefaultRelease(prefix: string, sku: string): Promise { const rolledOutReleases = await prisma.release.findMany({ where: { type: prefix, rolloutPercentage: 100 }, select: compatibleReleaseSelect(sku), }); if (rolledOutReleases.length === 0) { - throw new InternalServerError( - `No default release found for type ${prefix} and SKU "${sku}"`, - ); + return null; } // Only consider releases that ship a binary for this SKU. Without this, @@ -605,7 +607,13 @@ export async function Retrieve(req: Request, res: Response) { latest.artifacts.length > 0 && (await isDeviceEligibleForLatestRelease(latest.rolloutPercentage, deviceId)); - offered[kind] = dbReleaseToMetadata(useLatest ? latest : fallback, sku); + const chosen = useLatest ? latest : fallback; + if (!chosen) { + throw new NotFoundError( + `No ${prefix} release is rolled out yet for SKU "${sku}"`, + ); + } + offered[kind] = dbReleaseToMetadata(chosen, sku); }), ); diff --git a/test/releases.test.ts b/test/releases.test.ts index cc32725..9d672bc 100644 --- a/test/releases.test.ts +++ b/test/releases.test.ts @@ -587,15 +587,26 @@ describe("Retrieve handler", () => { expect(s3Mock.commandCalls(GetObjectCommand)).toHaveLength(0); }); - it("fails when no fully rolled out default exists for background checks", async () => { + it("serves a staged release to an in-bucket device even when nothing is at 100%", async () => { + // early-adopter hashes to rollout bucket 8, inside a 50% rollout. + await testPrisma.release.updateMany({ data: { rolloutPercentage: 50 } }); + + const res = createMockResponse(); + await Retrieve(createMockRequest({ deviceId: "early-adopter" }), res); + + expect(jsonBody(res)).toMatchObject({ appVersion: "1.2.0", systemVersion: "1.2.0" }); + }); + + it("answers 404, not 500, to an out-of-bucket device when nothing is at 100%", async () => { + // late-adopter hashes to rollout bucket 95, outside a 50% rollout. await testPrisma.release.updateMany({ data: { rolloutPercentage: 50 } }); await expect( - Retrieve( - createMockRequest({ deviceId: "no-default-device" }), - createMockResponse(), - ), - ).rejects.toThrow(InternalServerError); + Retrieve(createMockRequest({ deviceId: "late-adopter" }), createMockResponse()), + ).rejects.toThrow(NotFoundError); + await expect( + Retrieve(createMockRequest({ deviceId: "late-adopter" }), createMockResponse()), + ).rejects.toThrow(/No (app|system) release is rolled out yet for SKU "jetkvm-v2"/); }); }); @@ -789,6 +800,20 @@ describe("Retrieve handler", () => { ).rejects.toThrow('Version 1.0.0 has no artifact for SKU "jetkvm-mini-ethernet"'); }); + it("serves the first staged mini release to an in-bucket device and 404s the rest", async () => { + // No mini row at 100% yet: the first release is at 10%. + // early-adopter hashes to bucket 8, late-adopter to bucket 95. + await createDbRelease("mini", "1.0.0", 10, miniArtifacts("1.0.0")); + + const inBucket = createMockResponse(); + await Retrieve(createMockRequest({ deviceId: "early-adopter", sku: MINI_ETHERNET_SKU }), inBucket); + expect(jsonBody(inBucket)).toMatchObject({ systemVersion: "1.0.0" }); + + await expect( + Retrieve(createMockRequest({ deviceId: "late-adopter", sku: MINI_ETHERNET_SKU }), createMockResponse()), + ).rejects.toThrow('No mini release is rolled out yet for SKU "jetkvm-mini-ethernet"'); + }); + it("fails when no mini release exists for the SKU", async () => { const res = createMockResponse(); await expect( From a01f7a42f1dc0fe90b11e5ebbe784d9e263192dc Mon Sep 17 00:00:00 2001 From: Adam Shiervani Date: Sat, 19 Sep 2026 13:24:23 +0200 Subject: [PATCH 2/6] Releases: sync new R2 releases from the API every 30 minutes (#81) * feat(releases): sync new R2 releases from the API every 30 minutes The release sync only ran when an operator invoked scripts/sync-releases.ts by hand, so a firmware upload sat in R2 until someone remembered to run it. The API process now runs the same sync on start and every 30 minutes, registering each new stable version at the default 10% rollout. The non-interactive core (bucket listing, artifact collection, DB insert) moves from the script into src/release-sync.ts so the API can import it; the script keeps the GPG check and the confirmation prompt and passes them in as the release decider. Both entry points share one R2 client via src/s3.ts. Scheduled runs skip a tick while the previous run is still going, log and survive a failed run, and treat a unique-constraint hit on insert as "already synced" so several API instances can race without an error. The existence check now precedes the R2 artifact walk, so a tick over an already-synced bucket costs one list call per prefix plus one DB lookup per version. * fix(release-sync): defer versions whose upload is still settling A scheduled tick can land while a version is still being uploaded and register it with a partial SKU set or a hash that changes afterwards. Sync never rewrites a row, so that snapshot would be permanent. Before collecting artifacts, list every object under the version folder and skip the version when the newest object changed in the last 10 minutes. The window is shorter than the schedule interval, so a real upload costs at most one extra tick. * refactor(release-sync): share R2 config, paginate via SDK, probe SKUs concurrently * src/s3.ts now exports bucketName and baseUrl next to the client, so the request handlers, the scheduler and the script read R2 config from one place instead of three. * The upload settle window moves into SyncConfig.uploadSettleMs. The scheduler sets it; the one-shot operator script leaves it unset because an operator sees the artifact list before confirming and has no next run. * newestUploadTime uses the SDK's paginateListObjectsV2 instead of a hand-rolled continuation loop. * collectReleaseArtifacts probes SKUs concurrently and folds the results in SKU order, so artifact order and the primary artifact are unchanged. * Tests drop empty per-prefix listing stubs that the file-level default already covers. * fix(release-sync): take the settle listing after the artifact scan Listing the version folder before the scan left a gap: an upload that started between the listing and the scan passed the check and was captured partially. Listing after the scan closes it, because any upload active during the scan leaves an object newer than the window and the snapshot is discarded instead of registered. * fix(release-sync): paginate the version listing ListObjectsV2 returns at most 1,000 common prefixes per page. Once a release prefix grows past that, versions on later pages were never seen on any tick. listStableVersions now walks every page with the SDK paginator, as newestUploadTime already does, and the mock in the sync tests can serve a truncated listing. --- scripts/sync-releases.ts | 411 ++++++------------------------------- src/index.ts | 5 + src/release-sync.ts | 390 +++++++++++++++++++++++++++++++++++ src/releases.ts | 68 +----- src/s3.ts | 57 +++++ test/sync-releases.test.ts | 289 ++++++++++++++++++++++++-- 6 files changed, 794 insertions(+), 426 deletions(-) create mode 100644 src/release-sync.ts create mode 100644 src/s3.ts diff --git a/scripts/sync-releases.ts b/scripts/sync-releases.ts index c1fa91f..6a2edb4 100644 --- a/scripts/sync-releases.ts +++ b/scripts/sync-releases.ts @@ -8,54 +8,28 @@ import { Readable } from "node:stream"; import { pipeline } from "node:stream/promises"; import { createInterface } from "node:readline/promises"; -import { - GetObjectCommand, - HeadObjectCommand, - ListObjectsV2Command, - S3Client, -} from "@aws-sdk/client-s3"; +import { GetObjectCommand, S3Client } from "@aws-sdk/client-s3"; import { PrismaClient } from "@prisma/client"; import semver from "semver"; -import { objectKeyFromArtifactUrl, streamToString } from "../src/helpers"; -import { OTA_PREFIXES, legacyCompatibleSkus, otaFileForPrefix, skusForPrefix } from "../src/skus"; - -/** An R2 prefix, which is also the Release.type column value. */ -type ReleaseType = string; +import { objectKeyFromArtifactUrl } from "../src/helpers"; +import { baseUrl, bucketName, s3Client, s3ObjectExists } from "../src/s3"; +import { + DEFAULT_ROLLOUT_PERCENTAGE, + createAtDefaultRollout, + syncReleases, + type ReleaseArtifactInput, + type ReleaseDecider, + type ReleaseType, + type SyncClients, + type SyncConfig, +} from "../src/release-sync"; + +// Operator front end for src/release-sync.ts: signature verification output +// and a confirmation prompt before each production DB write. const OTA_ROOT_KEY_FPR = "AF5A36A993D828FEFE7C18C2D1B9856C26A79E95"; -interface SyncClients { - prisma: PrismaClient; - s3Client: S3Client; -} - -interface SyncConfig { - bucketName: string; - baseUrl: string; - skus?: string[]; -} - -interface ReleaseArtifactInput { - url: string; - hash: string; - compatibleSkus: string[]; -} - -const DEFAULT_ROLLOUT_PERCENTAGE = 10; - -type ReleaseOutcome = - | "created" - | "already-synced" - | "no-artifacts" - | "skipped" - | "aborted"; - -type ReleaseDecision = - | { kind: "create"; rolloutPercentage: number } - | { kind: "skip" } - | { kind: "abort" }; - interface LatestExistingRelease { version: string; rolloutPercentage: number; @@ -360,305 +334,41 @@ async function promptRolloutPercentage( } } -async function confirmProductionCreate( +function confirmProductionCreate( clients: SyncClients, config: SyncConfig, - type: ReleaseType, - version: string, - artifacts: ReleaseArtifactInput[], -): Promise { - if (process.env.NODE_ENV !== "production") { - return { kind: "create", rolloutPercentage: DEFAULT_ROLLOUT_PERCENTAGE }; - } - - if (!stdin.isTTY || !stdout.isTTY) { - throw new Error( - "Production release sync requires an interactive terminal for DB write confirmation.", - ); - } - - const [artifactInfos, latestExisting] = await Promise.all([ - loadArtifactDisplayInfo(clients, config, artifacts), - findLatestExistingRelease(clients.prisma, type), - ]); +): ReleaseDecider { + return async (type, version, artifacts) => { + const [artifactInfos, latestExisting] = await Promise.all([ + loadArtifactDisplayInfo(clients, config, artifacts), + findLatestExistingRelease(clients.prisma, type), + ]); - printArtifactSummary(type, version, artifactInfos, latestExisting); + printArtifactSummary(type, version, artifactInfos, latestExisting); - const readline = createInterface({ input: stdin, output: stdout }); - try { - const rolloutPercentage = await promptRolloutPercentage(readline); + const readline = createInterface({ input: stdin, output: stdout }); + try { + const rolloutPercentage = await promptRolloutPercentage(readline); - const confirmation = ( - await readline.question( - ` Create production ${type} release ${version} at ${rolloutPercentage}% rollout? [y/N/a (abort run)] `, + const confirmation = ( + await readline.question( + ` Create production ${type} release ${version} at ${rolloutPercentage}% rollout? [y/N/a (abort run)] `, + ) ) - ) - .trim() - .toLowerCase(); - - if (["a", "abort"].includes(confirmation)) { - return { kind: "abort" }; - } - if (!["y", "yes"].includes(confirmation)) { - return { kind: "skip" }; - } - return { kind: "create", rolloutPercentage }; - } finally { - readline.close(); - } -} - -function isS3NotFound(error: any): boolean { - return ( - error.name === "NotFound" || - error.name === "NoSuchKey" || - error.$metadata?.httpStatusCode === 404 - ); -} - -async function s3ObjectExists( - s3Client: S3Client, - bucketName: string, - key: string, -): Promise { - try { - await s3Client.send(new HeadObjectCommand({ Bucket: bucketName, Key: key })); - return true; - } catch (error: any) { - if (isS3NotFound(error)) { - return false; - } - throw error; - } -} - -async function versionHasSkuSupport( - s3Client: S3Client, - bucketName: string, - type: ReleaseType, - version: string, -): Promise { - const response = await s3Client.send( - new ListObjectsV2Command({ - Bucket: bucketName, - Prefix: `${type}/${version}/skus/`, - MaxKeys: 1, - }), - ); - return (response.Contents?.length ?? 0) > 0; -} - -async function readHash( - s3Client: S3Client, - bucketName: string, - artifactPath: string, -): Promise { - try { - const response = await s3Client.send( - new GetObjectCommand({ - Bucket: bucketName, - Key: `${artifactPath}.sha256`, - }), - ); - return streamToString(response.Body); - } catch (error: any) { - if (isS3NotFound(error)) { - return undefined; - } - throw error; - } -} - -function addArtifact( - artifactsByUrl: Map, - url: string, - hash: string, - sku: string, -): void { - const artifact = artifactsByUrl.get(url); - if (artifact) { - if (!artifact.compatibleSkus.includes(sku)) { - artifact.compatibleSkus.push(sku); - } - return; - } - - artifactsByUrl.set(url, { url, hash, compatibleSkus: [sku] }); -} - -export async function collectReleaseArtifacts( - clients: Pick, - config: SyncConfig, - type: ReleaseType, - version: string, -): Promise { - const skus = config.skus ?? skusForPrefix(type); - const artifactFileName = otaFileForPrefix(type); - - if (!(await versionHasSkuSupport(clients.s3Client, config.bucketName, type, version))) { - // Pre-SKU artifacts (no skus/ folder) are only safe on the SKUs that - // predate the layout. A type with no legacy form treats a version - // without skus/ as an upload mistake, not a release. - const compatibleSkus = legacyCompatibleSkus(type); - if (compatibleSkus.length === 0) { - return []; - } - - const artifactPath = `${type}/${version}/${artifactFileName}`; - const hash = await readHash(clients.s3Client, config.bucketName, artifactPath); - if (!hash) { - return []; - } - - return [ - { - url: `${config.baseUrl}/${artifactPath}`, - hash, - compatibleSkus, - }, - ]; - } - - const artifactsByUrl = new Map(); - for (const sku of skus) { - const artifactPath = `${type}/${version}/skus/${sku}/${artifactFileName}`; - if (!(await s3ObjectExists(clients.s3Client, config.bucketName, artifactPath))) { - continue; - } - - const hash = await readHash(clients.s3Client, config.bucketName, artifactPath); - if (!hash) { - continue; - } - addArtifact(artifactsByUrl, `${config.baseUrl}/${artifactPath}`, hash, sku); - } - - return Array.from(artifactsByUrl.values()); -} - -async function listStableVersions( - s3Client: S3Client, - bucketName: string, - type: ReleaseType, -): Promise { - const response = await s3Client.send( - new ListObjectsV2Command({ - Bucket: bucketName, - Prefix: `${type}/`, - Delimiter: "/", - }), - ); - - return (response.CommonPrefixes ?? []) - .map(cp => cp.Prefix?.split("/")[1]) - .filter((version): version is string => Boolean(version)) - .filter( - version => Boolean(semver.valid(version)) && semver.prerelease(version) === null, - ) - .sort(semver.compare); -} - -async function syncRelease( - clients: SyncClients, - config: SyncConfig, - type: ReleaseType, - version: string, - artifacts: ReleaseArtifactInput[], -): Promise { - if (artifacts.length === 0) { - console.log(`[sync-releases] ${type} ${version}: skipped, no compatible artifacts`); - return "no-artifacts"; - } - - // Sync only registers brand-new releases. Existing rows (rollout state, URLs, - // artifact compatibility) are left untouched — backfills/repairs are handled - // by one-off scripts so a routine sync run can never rewrite production data. - const existing = await clients.prisma.release.findUnique({ - where: { version_type: { version, type } }, - select: { id: true }, - }); - - if (existing) { - console.log(`[sync-releases] ${type} ${version}: already synced, skipping`); - return "already-synced"; - } - - const decision = await confirmProductionCreate( - clients, - config, - type, - version, - artifacts, - ); - if (decision.kind === "abort") { - console.log(`[sync-releases] ${type} ${version}: aborted by user`); - return "aborted"; - } - if (decision.kind === "skip") { - console.log(`[sync-releases] ${type} ${version}: skipped by user`); - return "skipped"; - } - - const primaryArtifact = artifacts[0]; - await clients.prisma.release.create({ - data: { - version, - type, - rolloutPercentage: decision.rolloutPercentage, - url: primaryArtifact.url, - hash: primaryArtifact.hash, - artifacts: { - create: artifacts.map(artifact => ({ - url: artifact.url, - hash: artifact.hash, - compatibleSkus: artifact.compatibleSkus, - })), - }, - }, - }); - - console.log( - `[sync-releases] ${type} ${version}: created with ${artifacts.length} artifact(s) at ${decision.rolloutPercentage}% rollout`, - ); - return "created"; -} - -export async function syncReleases( - clients: SyncClients, - config: SyncConfig, -): Promise { - const stats: Record = { - created: 0, - "already-synced": 0, - "no-artifacts": 0, - skipped: 0, - aborted: 0, - }; - let abortedAt: { type: ReleaseType; version: string } | null = null; - - outer: for (const type of OTA_PREFIXES) { - const versions = await listStableVersions(clients.s3Client, config.bucketName, type); - - for (const version of versions) { - const artifacts = await collectReleaseArtifacts(clients, config, type, version); - const outcome = await syncRelease(clients, config, type, version, artifacts); - stats[outcome]++; + .trim() + .toLowerCase(); - if (outcome === "aborted") { - abortedAt = { type, version }; - break outer; + if (["a", "abort"].includes(confirmation)) { + return { kind: "abort" }; } + if (!["y", "yes"].includes(confirmation)) { + return { kind: "skip" }; + } + return { kind: "create", rolloutPercentage }; + } finally { + readline.close(); } - } - - if (abortedAt) { - console.log( - `[sync-releases] aborted at ${abortedAt.type} ${abortedAt.version}; remaining versions in this run were not processed`, - ); - } - console.log( - `[sync-releases] done: created=${stats.created} skipped-by-user=${stats.skipped} already-synced=${stats["already-synced"]} no-artifacts=${stats["no-artifacts"]}`, - ); + }; } function describeDbTarget(): string { @@ -678,35 +388,32 @@ function describeDbTarget(): string { async function main(): Promise { console.log( - `[sync-releases] env=${process.env.NODE_ENV ?? "(unset)"} db=${describeDbTarget()} bucket=${process.env.R2_BUCKET ?? "(unset)"}`, + `[sync-releases] env=${process.env.NODE_ENV ?? "(unset)"} db=${describeDbTarget()} bucket=${bucketName ?? "(unset)"}`, ); + const isProduction = process.env.NODE_ENV === "production"; + if (isProduction && (!stdin.isTTY || !stdout.isTTY)) { + throw new Error( + "Production release sync requires an interactive terminal for DB write confirmation.", + ); + } + const prisma = new PrismaClient(); - const s3Client = new S3Client({ - endpoint: process.env.R2_ENDPOINT!, - credentials: { - accessKeyId: process.env.R2_ACCESS_KEY_ID!, - secretAccessKey: process.env.R2_SECRET_ACCESS_KEY!, - }, - region: "auto", - }); + const clients: SyncClients = { prisma, s3Client }; + const config: SyncConfig = { bucketName, baseUrl }; try { await syncReleases( - { prisma, s3Client }, - { - bucketName: process.env.R2_BUCKET!, - baseUrl: process.env.R2_CDN_URL!, - }, + clients, + config, + isProduction ? confirmProductionCreate(clients, config) : createAtDefaultRollout, ); } finally { await prisma.$disconnect(); } } -if (require.main === module) { - main().catch(error => { - console.error("[sync-releases] failed", error); - process.exit(1); - }); -} +main().catch(error => { + console.error("[sync-releases] failed", error); + process.exit(1); +}); diff --git a/src/index.ts b/src/index.ts index b6fdecb..2a235c1 100644 --- a/src/index.ts +++ b/src/index.ts @@ -13,6 +13,8 @@ import * as Releases from "./releases"; import { HttpError } from "./errors"; import { authenticated } from "./auth"; import { prisma } from "./db"; +import { baseUrl, bucketName, s3Client } from "./s3"; +import { scheduleReleaseSync } from "./release-sync"; import { initializeWebRTCSignaling } from "./webrtc-signaling"; declare global { @@ -218,3 +220,6 @@ const server = app.listen(PORT, () => { }); initializeWebRTCSignaling(server); + +// Register new R2 releases at the default rollout, now and every 30 minutes. +scheduleReleaseSync({ prisma, s3Client }, { bucketName, baseUrl }); diff --git a/src/release-sync.ts b/src/release-sync.ts new file mode 100644 index 0000000..2911427 --- /dev/null +++ b/src/release-sync.ts @@ -0,0 +1,390 @@ +import { GetObjectCommand, S3Client, paginateListObjectsV2 } from "@aws-sdk/client-s3"; +import { Prisma, PrismaClient } from "@prisma/client"; +import semver from "semver"; + +import { streamToString } from "./helpers"; +import { isS3NotFound, s3ObjectExists, versionHasSkuSupport } from "./s3"; +import { OTA_PREFIXES, legacyCompatibleSkus, otaFileForPrefix, skusForPrefix } from "./skus"; + +/** An R2 prefix, which is also the Release.type column value. */ +export type ReleaseType = string; + +export interface SyncClients { + prisma: PrismaClient; + s3Client: S3Client; +} + +export interface SyncConfig { + bucketName: string; + baseUrl: string; + skus?: string[]; + /** + * Defer a version whose newest object changed within this many ms, checked + * after the artifact scan: the upload script writes several files per SKU, + * and sync never rewrites a row, so a scan that overlapped an upload would + * freeze a partial SKU set or a stale hash. Unset or 0 disables the check + * (an operator at a terminal can judge the artifact list for themselves). + */ + uploadSettleMs?: number; +} + +export interface ReleaseArtifactInput { + url: string; + hash: string; + compatibleSkus: string[]; +} + +export const DEFAULT_ROLLOUT_PERCENTAGE = 10; + +export type ReleaseOutcome = + | "created" + | "already-synced" + | "uploading" + | "no-artifacts" + | "skipped" + | "aborted"; + +/** Settle window for unattended runs; shorter than the interval, so it costs at most one tick. */ +export const UPLOAD_SETTLE_MS = 10 * 60 * 1000; + +export type ReleaseDecision = + | { kind: "create"; rolloutPercentage: number } + | { kind: "skip" } + | { kind: "abort" }; + +/** + * Decides whether a release that exists in R2 but not in the DB gets created. + * The interactive script asks the operator; the in-process scheduler always + * creates at the default rollout. + */ +export type ReleaseDecider = ( + type: ReleaseType, + version: string, + artifacts: ReleaseArtifactInput[], +) => Promise; + +export const createAtDefaultRollout: ReleaseDecider = async () => ({ + kind: "create", + rolloutPercentage: DEFAULT_ROLLOUT_PERCENTAGE, +}); + +async function readHash( + s3Client: S3Client, + bucketName: string, + artifactPath: string, +): Promise { + try { + const response = await s3Client.send( + new GetObjectCommand({ + Bucket: bucketName, + Key: `${artifactPath}.sha256`, + }), + ); + return streamToString(response.Body); + } catch (error: any) { + if (isS3NotFound(error)) { + return undefined; + } + throw error; + } +} + +function addArtifact( + artifactsByUrl: Map, + url: string, + hash: string, + sku: string, +): void { + const artifact = artifactsByUrl.get(url); + if (artifact) { + if (!artifact.compatibleSkus.includes(sku)) { + artifact.compatibleSkus.push(sku); + } + return; + } + + artifactsByUrl.set(url, { url, hash, compatibleSkus: [sku] }); +} + +export async function collectReleaseArtifacts( + clients: Pick, + config: SyncConfig, + type: ReleaseType, + version: string, +): Promise { + const skus = config.skus ?? skusForPrefix(type); + const artifactFileName = otaFileForPrefix(type); + + if (!(await versionHasSkuSupport(clients.s3Client, config.bucketName, type, version))) { + // Pre-SKU artifacts (no skus/ folder) are only safe on the SKUs that + // predate the layout. A type with no legacy form treats a version + // without skus/ as an upload mistake, not a release. + const compatibleSkus = legacyCompatibleSkus(type); + if (compatibleSkus.length === 0) { + return []; + } + + const artifactPath = `${type}/${version}/${artifactFileName}`; + const hash = await readHash(clients.s3Client, config.bucketName, artifactPath); + if (!hash) { + return []; + } + + return [ + { + url: `${config.baseUrl}/${artifactPath}`, + hash, + compatibleSkus, + }, + ]; + } + + // SKUs are probed concurrently; folding afterwards in `skus` order keeps + // the primary artifact (artifacts[0]) and compatibleSkus order stable. + const found = await Promise.all( + skus.map(async sku => { + const artifactPath = `${type}/${version}/skus/${sku}/${artifactFileName}`; + if (!(await s3ObjectExists(clients.s3Client, config.bucketName, artifactPath))) { + return undefined; + } + const hash = await readHash(clients.s3Client, config.bucketName, artifactPath); + return hash ? { sku, url: `${config.baseUrl}/${artifactPath}`, hash } : undefined; + }), + ); + + const artifactsByUrl = new Map(); + for (const artifact of found) { + if (artifact) { + addArtifact(artifactsByUrl, artifact.url, artifact.hash, artifact.sku); + } + } + return Array.from(artifactsByUrl.values()); +} + +async function listStableVersions( + s3Client: S3Client, + bucketName: string, + type: ReleaseType, +): Promise { + const prefixes: string[] = []; + for await (const page of paginateListObjectsV2( + { client: s3Client }, + { Bucket: bucketName, Prefix: `${type}/`, Delimiter: "/" }, + )) { + for (const cp of page.CommonPrefixes ?? []) { + if (cp.Prefix) { + prefixes.push(cp.Prefix); + } + } + } + + return prefixes + .map(prefix => prefix.split("/")[1]) + .filter((version): version is string => Boolean(version)) + .filter( + version => Boolean(semver.valid(version)) && semver.prerelease(version) === null, + ) + .sort(semver.compare); +} + +/** Newest LastModified among every object under `${type}/${version}/`, or undefined when empty. */ +async function newestUploadTime( + s3Client: S3Client, + bucketName: string, + type: ReleaseType, + version: string, +): Promise { + let newest: Date | undefined; + for await (const page of paginateListObjectsV2( + { client: s3Client }, + { Bucket: bucketName, Prefix: `${type}/${version}/` }, + )) { + for (const object of page.Contents ?? []) { + if (object.LastModified && (!newest || object.LastModified > newest)) { + newest = object.LastModified; + } + } + } + return newest; +} + +async function listSyncedVersions(prisma: PrismaClient, type: ReleaseType): Promise> { + const releases = await prisma.release.findMany({ + where: { type }, + select: { version: true }, + }); + return new Set(releases.map(release => release.version)); +} + +function isUniqueViolation(error: unknown): boolean { + return ( + error instanceof Prisma.PrismaClientKnownRequestError && error.code === "P2002" + ); +} + +async function createRelease( + clients: SyncClients, + config: SyncConfig, + decide: ReleaseDecider, + type: ReleaseType, + version: string, +): Promise { + const artifacts = await collectReleaseArtifacts(clients, config, type, version); + + // Listed after the artifact scan on purpose: an upload that was active at + // any point during the scan leaves an object newer than the window, so the + // snapshot above is discarded rather than registered. + if (config.uploadSettleMs) { + const newest = await newestUploadTime(clients.s3Client, config.bucketName, type, version); + if (newest && Date.now() - newest.getTime() < config.uploadSettleMs) { + console.log( + `[sync-releases] ${type} ${version}: upload still settling (last object ${newest.toISOString()}), retrying next run`, + ); + return "uploading"; + } + } + + if (artifacts.length === 0) { + console.log(`[sync-releases] ${type} ${version}: skipped, no compatible artifacts`); + return "no-artifacts"; + } + + const decision = await decide(type, version, artifacts); + if (decision.kind === "abort") { + console.log(`[sync-releases] ${type} ${version}: aborted by user`); + return "aborted"; + } + if (decision.kind === "skip") { + console.log(`[sync-releases] ${type} ${version}: skipped by user`); + return "skipped"; + } + + const primaryArtifact = artifacts[0]; + try { + await clients.prisma.release.create({ + data: { + version, + type, + rolloutPercentage: decision.rolloutPercentage, + url: primaryArtifact.url, + hash: primaryArtifact.hash, + artifacts: { + create: artifacts.map(artifact => ({ + url: artifact.url, + hash: artifact.hash, + compatibleSkus: artifact.compatibleSkus, + })), + }, + }, + }); + } catch (error) { + // Another API instance can win the race between the version listing and + // this insert. The row it wrote is the one we wanted, so treat it as synced. + if (isUniqueViolation(error)) { + console.log(`[sync-releases] ${type} ${version}: created concurrently elsewhere, skipping`); + return "already-synced"; + } + throw error; + } + + console.log( + `[sync-releases] ${type} ${version}: created with ${artifacts.length} artifact(s) at ${decision.rolloutPercentage}% rollout`, + ); + return "created"; +} + +/** + * Registers every stable version in R2 that has no Release row yet. + * + * Sync only registers brand-new releases. Existing rows (rollout state, URLs, + * artifact compatibility) are left untouched — backfills/repairs are handled + * by one-off scripts so a routine sync run can never rewrite production data. + * Known versions are loaded once per prefix, so a run over an already-synced + * bucket costs one R2 list and one DB query per prefix. + * + * Returns the per-outcome counts so callers can log or assert on them. + */ +export async function syncReleases( + clients: SyncClients, + config: SyncConfig, + decide: ReleaseDecider, +): Promise> { + const stats: Record = { + created: 0, + "already-synced": 0, + uploading: 0, + "no-artifacts": 0, + skipped: 0, + aborted: 0, + }; + let abortedAt: { type: ReleaseType; version: string } | null = null; + + outer: for (const type of OTA_PREFIXES) { + const [versions, synced] = await Promise.all([ + listStableVersions(clients.s3Client, config.bucketName, type), + listSyncedVersions(clients.prisma, type), + ]); + + for (const version of versions) { + const outcome = synced.has(version) + ? "already-synced" + : await createRelease(clients, config, decide, type, version); + stats[outcome]++; + + if (outcome === "aborted") { + abortedAt = { type, version }; + break outer; + } + } + } + + if (abortedAt) { + console.log( + `[sync-releases] aborted at ${abortedAt.type} ${abortedAt.version}; remaining versions in this run were not processed`, + ); + } + console.log( + `[sync-releases] done: created=${stats.created} skipped-by-user=${stats.skipped} already-synced=${stats["already-synced"]} uploading=${stats.uploading} no-artifacts=${stats["no-artifacts"]}`, + ); + return stats; +} + +const RELEASE_SYNC_INTERVAL_MS = 30 * 60 * 1000; + +/** + * Runs the sync now and then every `intervalMs`, inside the API process. + * A tick that fires while the previous run is still going is skipped, not + * queued. A failed run is logged and the schedule continues, so one bad R2 + * or DB response never stops future syncs. + * Returns a function that stops the schedule. + */ +export function scheduleReleaseSync( + clients: SyncClients, + config: SyncConfig, + intervalMs: number = RELEASE_SYNC_INTERVAL_MS, +): () => void { + let running = false; + + const run = async () => { + if (running) { + console.warn("[sync-releases] previous run still in progress, skipping this tick"); + return; + } + running = true; + try { + await syncReleases( + clients, + { uploadSettleMs: UPLOAD_SETTLE_MS, ...config }, + createAtDefaultRollout, + ); + } catch (error) { + console.error("[sync-releases] scheduled run failed", error); + } finally { + running = false; + } + }; + + void run(); + const timer = setInterval(run, intervalMs); + return () => clearInterval(timer); +} diff --git a/src/releases.ts b/src/releases.ts index fb0293f..33d5323 100644 --- a/src/releases.ts +++ b/src/releases.ts @@ -3,13 +3,9 @@ import { prisma } from "./db"; import { BadRequestError, InternalServerError, NotFoundError } from "./errors"; import semver from "semver"; -import { - GetObjectCommand, - HeadObjectCommand, - ListObjectsV2Command, - S3Client, -} from "@aws-sdk/client-s3"; +import { GetObjectCommand, ListObjectsV2Command } from "@aws-sdk/client-s3"; import { LRUCache } from "lru-cache"; +import { baseUrl, bucketName, s3Client, s3ObjectExists, versionHasSkuSupport } from "./s3"; import { getDeviceRolloutBucket, @@ -101,15 +97,6 @@ interface DbRelease { }[]; } -const s3Client = new S3Client({ - endpoint: process.env.R2_ENDPOINT!, - credentials: { - accessKeyId: process.env.R2_ACCESS_KEY_ID!, - secretAccessKey: process.env.R2_SECRET_ACCESS_KEY!, - }, - region: "auto", -}); - const releaseCache = new LRUCache({ max: 1000, ttl: 5 * 60 * 1000, // 5 minutes @@ -134,9 +121,6 @@ export function clearCaches() { sigUrlCache.clear(); } -const bucketName = process.env.R2_BUCKET; -const baseUrl = process.env.R2_CDN_URL; - /** * The one error for "this version ships no artifact for this SKU", whichever * path detects it: a skus/ folder without this SKU, a pre-SKU version asked @@ -146,41 +130,7 @@ function noArtifactForSku(version: string, sku: string): NotFoundError { return new NotFoundError(`Version ${version} has no artifact for SKU "${sku}"`); } -/** - * Checks if an object exists in S3/R2 by attempting a HeadObjectCommand. - * Returns true if the object exists, false otherwise. - */ -async function s3ObjectExists(key: string): Promise { - try { - await s3Client.send(new HeadObjectCommand({ Bucket: bucketName, Key: key })); - return true; - } catch (error: any) { - // HeadObjectCommand throws NotFound, but some S3-compatible stores (like R2) may throw NoSuchKey - if ( - error.name === "NotFound" || - error.name === "NoSuchKey" || - error.$metadata?.httpStatusCode === 404 - ) { - return false; - } - throw error; - } -} - -/** - * Checks if a version was uploaded with SKU folder structure. - * Returns true if any skus/ subfolder exists for this version. - */ -async function versionHasSkuSupport(prefix: string, version: string): Promise { - const response = await s3Client.send( - new ListObjectsV2Command({ - Bucket: bucketName, - Prefix: `${prefix}/${version}/skus/`, - MaxKeys: 1, - }), - ); - return (response.Contents?.length ?? 0) > 0; -} +const objectExists = (key: string) => s3ObjectExists(s3Client, bucketName, key); /** * Resolves the artifact path for a given version and SKU. @@ -205,10 +155,10 @@ async function resolveArtifactPath( sku: string, file: string, ): Promise { - if (await versionHasSkuSupport(prefix, version)) { + if (await versionHasSkuSupport(s3Client, bucketName, prefix, version)) { const skuPath = `${prefix}/${version}/skus/${sku}/${file}`; - if (await s3ObjectExists(skuPath)) { + if (await objectExists(skuPath)) { return skuPath; } @@ -243,7 +193,7 @@ async function resolveSigUrl( try { const path = await resolveArtifactPath(artifact.prefix, version, sku, artifact.file); const sigKey = `${path}.sig`; - if (await s3ObjectExists(sigKey)) { + if (await objectExists(sigKey)) { const url = `${baseUrl}/${sigKey}`; sigUrlCache.set(cacheKey, url); return url; @@ -389,7 +339,7 @@ async function resolveSigUrlFromArtifactUrl( const sigUrl = `${artifactUrl}.sig`; try { const sigKey = `${objectKeyFromArtifactUrl(artifactUrl)}.sig`; - if (await s3ObjectExists(sigKey)) { + if (await objectExists(sigKey)) { sigUrlCache.set(cacheKey, sigUrl); return sigUrl; } @@ -693,7 +643,7 @@ export const RetrieveLatestSystemRecovery = cachedRedirect( recovery.file, ); - if (!(await s3ObjectExists(artifactPath))) { + if (!(await objectExists(artifactPath))) { throw new NotFoundError(`Recovery image not found for version ${latestVersion}`); } @@ -750,7 +700,7 @@ function latestArtifactRedirect(kind: OtaKind) { artifact.file, ); - if (!(await s3ObjectExists(artifactPath))) { + if (!(await objectExists(artifactPath))) { throw new NotFoundError(`${prefix} artifact not found for version ${latestVersion}`); } diff --git a/src/s3.ts b/src/s3.ts new file mode 100644 index 0000000..1c702bb --- /dev/null +++ b/src/s3.ts @@ -0,0 +1,57 @@ +import { HeadObjectCommand, ListObjectsV2Command, S3Client } from "@aws-sdk/client-s3"; + +/** The one R2 client the API process shares between request handlers and background jobs. */ +export const s3Client = new S3Client({ + endpoint: process.env.R2_ENDPOINT!, + credentials: { + accessKeyId: process.env.R2_ACCESS_KEY_ID!, + secretAccessKey: process.env.R2_SECRET_ACCESS_KEY!, + }, + region: "auto", +}); + +/** Bucket that holds release artifacts, and the public CDN origin they are served from. */ +export const bucketName = process.env.R2_BUCKET!; +export const baseUrl = process.env.R2_CDN_URL!; + +/** HeadObject throws NotFound, but some S3-compatible stores (like R2) may throw NoSuchKey. */ +export function isS3NotFound(error: any): boolean { + return ( + error.name === "NotFound" || + error.name === "NoSuchKey" || + error.$metadata?.httpStatusCode === 404 + ); +} + +export async function s3ObjectExists( + client: S3Client, + bucketName: string, + key: string, +): Promise { + try { + await client.send(new HeadObjectCommand({ Bucket: bucketName, Key: key })); + return true; + } catch (error: any) { + if (isS3NotFound(error)) { + return false; + } + throw error; + } +} + +/** True when the version was uploaded with the skus// folder layout. */ +export async function versionHasSkuSupport( + client: S3Client, + bucketName: string, + prefix: string, + version: string, +): Promise { + const response = await client.send( + new ListObjectsV2Command({ + Bucket: bucketName, + Prefix: `${prefix}/${version}/skus/`, + MaxKeys: 1, + }), + ); + return (response.Contents?.length ?? 0) > 0; +} diff --git a/test/sync-releases.test.ts b/test/sync-releases.test.ts index c0f1a92..51b3fbd 100644 --- a/test/sync-releases.test.ts +++ b/test/sync-releases.test.ts @@ -4,12 +4,20 @@ import { ListObjectsV2Command, S3Client, } from "@aws-sdk/client-s3"; -import { describe, expect, beforeEach, it } from "vitest"; +import { PrismaClient } from "@prisma/client"; +import { afterEach, describe, expect, beforeEach, it, vi } from "vitest"; -import { collectReleaseArtifacts, syncReleases } from "../scripts/sync-releases"; +import { + collectReleaseArtifacts, + createAtDefaultRollout, + scheduleReleaseSync, + syncReleases, + UPLOAD_SETTLE_MS, + type ReleaseDecider, + type ReleaseType, +} from "../src/release-sync"; import { otaFileForPrefix } from "../src/skus"; -type ReleaseType = string; import { createAsyncIterable, s3Mock, testPrisma } from "./setup"; const DEFAULT_SKU = "jetkvm-v2"; @@ -20,9 +28,27 @@ const SYNC_BUCKET = "test-bucket"; const SYNC_BASE_URL = "https://cdn.test.com"; const syncS3Client = new S3Client({}); -function mockS3ListVersions(prefix: ReleaseType, versions: string[]) { - s3Mock.on(ListObjectsV2Command, { Prefix: `${prefix}/` }).resolves({ - CommonPrefixes: versions.map(v => ({ Prefix: `${prefix}/${v}/` })), +/** Makes the settle check see one object under the version uploaded at `at`. */ +function mockS3UploadedAt(prefix: ReleaseType, version: string, at: Date) { + s3Mock.on(ListObjectsV2Command, { Prefix: `${prefix}/${version}/` }).resolves({ + Contents: [{ Key: `${prefix}/${version}/${otaFileForPrefix(prefix)}`, LastModified: at }], + }); +} + +/** Lists `versions` under `prefix`, one page per inner array, as R2 does past 1,000 keys. */ +function mockS3ListVersions(prefix: ReleaseType, ...pages: string[][]) { + pages.forEach((versions, index) => { + const last = index === pages.length - 1; + s3Mock + .on(ListObjectsV2Command, { + Prefix: `${prefix}/`, + ContinuationToken: index === 0 ? undefined : `page-${index}`, + }) + .resolves({ + CommonPrefixes: versions.map(v => ({ Prefix: `${prefix}/${v}/` })), + IsTruncated: !last, + NextContinuationToken: last ? undefined : `page-${index + 1}`, + }); }); } @@ -56,14 +82,17 @@ function mockS3SkuVersion( }); } -describe("sync-releases script", () => { - beforeEach(() => { - s3Mock.reset(); - s3Mock - .on(HeadObjectCommand) - .rejects({ name: "NotFound", $metadata: { httpStatusCode: 404 } }); - }); +beforeEach(() => { + s3Mock.reset(); + s3Mock + .on(HeadObjectCommand) + .rejects({ name: "NotFound", $metadata: { httpStatusCode: 404 } }); + // Listings a test does not stub (other prefixes, the upload settle check) see + // an empty folder. More specific .on(..., { Prefix }) stubs registered later win. + s3Mock.on(ListObjectsV2Command).resolves({ Contents: [] }); +}); +describe("syncReleases", () => { it("marks legacy app artifacts compatible with the default SKU only", async () => { mockS3HashFile("app", "9.9.1", "legacy-app-hash"); @@ -178,16 +207,25 @@ describe("sync-releases script", () => { mockS3ListVersions("app", [version, "10.0.0-beta.1"]); mockS3ListVersions("system", [version]); - mockS3ListVersions("mini", []); mockS3HashFile("app", version, "app-hash"); mockS3SkuVersion("system", version, DEFAULT_SKU, "system-hash-v2"); mockS3SkuVersion("system", version, SDMMC_SKU, "system-hash-sdmmc"); - await syncReleases( + const stats = await syncReleases( { prisma: testPrisma, s3Client: syncS3Client }, { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, + createAtDefaultRollout, ); + expect(stats).toEqual({ + created: 1, + "already-synced": 1, + uploading: 0, + "no-artifacts": 0, + skipped: 0, + aborted: 0, + }); + const appRelease = await testPrisma.release.findUniqueOrThrow({ where: { version_type: { version, type: "app" } }, include: { artifacts: true }, @@ -220,4 +258,225 @@ describe("sync-releases script", () => { // Prereleases are filtered out by listStableVersions. expect(prerelease).toBeNull(); }); + + it("registers versions from every page of a truncated listing", async () => { + const firstPage = "9.9.20"; + const secondPage = "9.9.21"; + mockS3ListVersions("app", [firstPage], [secondPage]); + mockS3HashFile("app", firstPage, "first-page-hash"); + mockS3HashFile("app", secondPage, "second-page-hash"); + + const stats = await syncReleases( + { prisma: testPrisma, s3Client: syncS3Client }, + { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, + createAtDefaultRollout, + ); + + expect(stats).toMatchObject({ created: 2 }); + const created = await testPrisma.release.findMany({ + where: { type: "app", version: { in: [firstPage, secondPage] } }, + orderBy: { version: "asc" }, + }); + expect(created.map(release => release.version)).toEqual([firstPage, secondPage]); + }); + + it("defers a version whose objects changed within the settle window", async () => { + const fresh = "9.9.10"; + const settled = "9.9.11"; + mockS3ListVersions("app", [fresh, settled]); + mockS3HashFile("app", fresh, "fresh-hash"); + mockS3HashFile("app", settled, "settled-hash"); + mockS3UploadedAt("app", fresh, new Date(Date.now() - 60 * 1000)); + mockS3UploadedAt("app", settled, new Date(Date.now() - UPLOAD_SETTLE_MS - 60 * 1000)); + + const stats = await syncReleases( + { prisma: testPrisma, s3Client: syncS3Client }, + { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL, uploadSettleMs: UPLOAD_SETTLE_MS }, + createAtDefaultRollout, + ); + + expect(stats).toMatchObject({ created: 1, uploading: 1 }); + expect( + await testPrisma.release.findUnique({ + where: { version_type: { version: fresh, type: "app" } }, + }), + ).toBeNull(); + // The settle listing must come after the artifact scan, so an upload that + // overlaps the scan is seen. The deferred version was therefore scanned. + const calls = s3Mock.calls().map(call => call.args[0].input as { Key?: string; Prefix?: string }); + const scanIndex = calls.findIndex(input => input.Key === `app/${fresh}/jetkvm_app.sha256`); + const settleIndex = calls.findIndex(input => input.Prefix === `app/${fresh}/`); + expect(scanIndex).toBeGreaterThanOrEqual(0); + expect(settleIndex).toBeGreaterThan(scanIndex); + expect( + await testPrisma.release.findUnique({ + where: { version_type: { version: settled, type: "app" } }, + }), + ).toMatchObject({ hash: "settled-hash" }); + }); + + it("honours the decider's rollout, skip and abort answers", async () => { + mockS3ListVersions("app", ["9.9.5", "9.9.6", "9.9.7"]); + mockS3ListVersions("system", ["9.9.5"]); + for (const version of ["9.9.5", "9.9.6", "9.9.7"]) { + mockS3HashFile("app", version, `app-hash-${version}`); + } + mockS3HashFile("system", "9.9.5", "system-hash"); + + const answers: Record>> = { + "app 9.9.5": { kind: "create", rolloutPercentage: 42 }, + "app 9.9.6": { kind: "skip" }, + "app 9.9.7": { kind: "abort" }, + }; + const decide: ReleaseDecider = async (type, version) => answers[`${type} ${version}`]; + + const stats = await syncReleases( + { prisma: testPrisma, s3Client: syncS3Client }, + { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, + decide, + ); + + expect(stats).toMatchObject({ created: 1, skipped: 1, aborted: 1 }); + + const created = await testPrisma.release.findUniqueOrThrow({ + where: { version_type: { version: "9.9.5", type: "app" } }, + }); + expect(created.rolloutPercentage).toBe(42); + + const notCreated = await testPrisma.release.findMany({ + where: { + OR: [ + { version: "9.9.6", type: "app" }, + { version: "9.9.7", type: "app" }, + // Abort stops the whole run, so system is never reached. + { version: "9.9.5", type: "system" }, + ], + }, + }); + expect(notCreated).toEqual([]); + }); + + it("treats a release created by another instance mid-run as already synced", async () => { + const version = "9.9.8"; + mockS3ListVersions("app", [version]); + mockS3HashFile("app", version, "app-hash"); + + // Simulate the race: the known-versions query sees nothing, but by the + // time this instance inserts, the row is there (written here up front so + // the real unique constraint fires on create). + await testPrisma.release.create({ + data: { + version, + type: "app", + rolloutPercentage: 10, + url: "https://cdn.test.com/other-instance", + hash: "other-instance-hash", + }, + }); + const racingPrisma = { + release: { + findMany: async () => [], + create: (args: unknown) => testPrisma.release.create(args as any), + }, + } as unknown as PrismaClient; + + const stats = await syncReleases( + { prisma: racingPrisma, s3Client: syncS3Client }, + { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, + createAtDefaultRollout, + ); + + expect(stats).toMatchObject({ created: 0, "already-synced": 1 }); + const release = await testPrisma.release.findUniqueOrThrow({ + where: { version_type: { version, type: "app" } }, + }); + expect(release.url).toBe("https://cdn.test.com/other-instance"); + }); +}); + +describe("scheduleReleaseSync", () => { + const INTERVAL_MS = 1000; + let stop: (() => void) | undefined; + + beforeEach(() => { + // Only the scheduler's own timer is faked; DB and S3 mock I/O stay real. + vi.useFakeTimers({ toFake: ["setInterval", "clearInterval"] }); + }); + + afterEach(() => { + stop?.(); + stop = undefined; + vi.useRealTimers(); + vi.restoreAllMocks(); + }); + + it("runs at start, keeps the schedule after a failed run and creates on the next tick", async () => { + const version = "9.9.9"; + s3Mock + .on(ListObjectsV2Command, { Prefix: "app/" }) + .rejectsOnce(new Error("R2 unavailable")) + .resolves({ CommonPrefixes: [{ Prefix: `app/${version}/` }] }); + mockS3HashFile("app", version, "app-hash"); + const errorLog = vi.spyOn(console, "error").mockImplementation(() => {}); + + stop = scheduleReleaseSync( + { prisma: testPrisma, s3Client: syncS3Client }, + { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, + INTERVAL_MS, + ); + + await vi.waitFor(() => expect(errorLog).toHaveBeenCalledOnce()); + expect( + await testPrisma.release.findUnique({ + where: { version_type: { version, type: "app" } }, + }), + ).toBeNull(); + + await vi.advanceTimersByTimeAsync(INTERVAL_MS); + + await vi.waitFor(async () => { + const release = await testPrisma.release.findUnique({ + where: { version_type: { version, type: "app" } }, + }); + expect(release?.rolloutPercentage).toBe(10); + }); + }); + + it("skips a tick while the previous run is still in progress", async () => { + let finishFirstRun!: () => void; + const firstListing = new Promise<{ CommonPrefixes: never[] }>(resolve => { + finishFirstRun = () => resolve({ CommonPrefixes: [] }); + }); + s3Mock + .on(ListObjectsV2Command, { Prefix: "app/" }) + .callsFakeOnce(() => firstListing) + .resolves({ CommonPrefixes: [] }); + const warnLog = vi.spyOn(console, "warn").mockImplementation(() => {}); + + stop = scheduleReleaseSync( + { prisma: testPrisma, s3Client: syncS3Client }, + { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, + INTERVAL_MS, + ); + + // Two ticks fire while the first run is still waiting on R2. + await vi.advanceTimersByTimeAsync(INTERVAL_MS * 2); + + expect(warnLog).toHaveBeenCalledTimes(2); + // The stalled run's list call is the only R2 traffic so far: the skipped + // ticks did not start a second walk of the bucket. + expect(s3Mock.commandCalls(ListObjectsV2Command)).toHaveLength(1); + + finishFirstRun(); + await vi.waitFor(() => + expect(s3Mock.commandCalls(ListObjectsV2Command)).toHaveLength(3), + ); + + // The next tick after the run completed starts a fresh run. + await vi.advanceTimersByTimeAsync(INTERVAL_MS); + await vi.waitFor(() => + expect(s3Mock.commandCalls(ListObjectsV2Command)).toHaveLength(6), + ); + expect(warnLog).toHaveBeenCalledTimes(2); + }); }); From 68c4dfcc8f47256c0b3f9a8d955b2841469da9cc Mon Sep 17 00:00:00 2001 From: Adam Shiervani Date: Sat, 19 Sep 2026 14:30:39 +0200 Subject: [PATCH 3/6] Release sync: a unique violation is only a race when the row exists (#82) * fix(release-sync): only treat a unique violation as a race when the row exists Sync caught every P2002 from the release insert as "created concurrently elsewhere". On staging the id sequences were behind the rows after a data import, so each insert failed on the primary key, was logged as a race, and left nothing in the table. After a unique violation, createRelease now looks the (version, type) row up. Present means another instance registered it first; absent means the insert really failed, and the error is rethrown with the type and version in its message so the scheduled run log names the release. * fix(release-sync): name the release in every per-version failure Wrapping only the insert error left the row lookup, and the S3 scan before it, free to escape without the type and version. syncReleases now wraps whatever createRelease throws for a version. --- src/release-sync.ts | 24 +++++++++++++++++++++--- test/sync-releases.test.ts | 37 ++++++++++++++++++++++++++++++++++++- 2 files changed, 57 insertions(+), 4 deletions(-) diff --git a/src/release-sync.ts b/src/release-sync.ts index 2911427..7a2f647 100644 --- a/src/release-sync.ts +++ b/src/release-sync.ts @@ -216,6 +216,18 @@ async function listSyncedVersions(prisma: PrismaClient, type: ReleaseType): Prom return new Set(releases.map(release => release.version)); } +async function releaseExists( + prisma: PrismaClient, + type: ReleaseType, + version: string, +): Promise { + const release = await prisma.release.findUnique({ + where: { version_type: { version, type } }, + select: { id: true }, + }); + return release !== null; +} + function isUniqueViolation(error: unknown): boolean { return ( error instanceof Prisma.PrismaClientKnownRequestError && error.code === "P2002" @@ -279,8 +291,10 @@ async function createRelease( }); } catch (error) { // Another API instance can win the race between the version listing and - // this insert. The row it wrote is the one we wanted, so treat it as synced. - if (isUniqueViolation(error)) { + // this insert, in which case the row it wrote is the one we wanted. Any + // other unique violation (a stale id sequence after a data import, say) + // leaves no row and is a real failure. + if (isUniqueViolation(error) && (await releaseExists(clients.prisma, type, version))) { console.log(`[sync-releases] ${type} ${version}: created concurrently elsewhere, skipping`); return "already-synced"; } @@ -326,9 +340,13 @@ export async function syncReleases( ]); for (const version of versions) { + // Name the release in any failure, whichever step raised it, so the + // scheduled run log does not need to be traced back to a version. const outcome = synced.has(version) ? "already-synced" - : await createRelease(clients, config, decide, type, version); + : await createRelease(clients, config, decide, type, version).catch((error: unknown) => { + throw new Error(`[sync-releases] ${type} ${version}: sync failed`, { cause: error }); + }); stats[outcome]++; if (outcome === "aborted") { diff --git a/test/sync-releases.test.ts b/test/sync-releases.test.ts index 51b3fbd..aa29e18 100644 --- a/test/sync-releases.test.ts +++ b/test/sync-releases.test.ts @@ -4,7 +4,7 @@ import { ListObjectsV2Command, S3Client, } from "@aws-sdk/client-s3"; -import { PrismaClient } from "@prisma/client"; +import { Prisma, PrismaClient } from "@prisma/client"; import { afterEach, describe, expect, beforeEach, it, vi } from "vitest"; import { @@ -376,6 +376,7 @@ describe("syncReleases", () => { const racingPrisma = { release: { findMany: async () => [], + findUnique: (args: unknown) => testPrisma.release.findUnique(args as any), create: (args: unknown) => testPrisma.release.create(args as any), }, } as unknown as PrismaClient; @@ -392,6 +393,40 @@ describe("syncReleases", () => { }); expect(release.url).toBe("https://cdn.test.com/other-instance"); }); + + it("fails on a unique violation that left no release row behind", async () => { + const version = "9.9.9"; + mockS3ListVersions("app", [version]); + mockS3HashFile("app", version, "app-hash"); + + // A stale id sequence (a database restored from a dump) makes the insert + // fail on the primary key with the same error code as the race above, + // but no row for this version exists afterwards. + const stalePrisma = { + release: { + findMany: async () => [], + findUnique: async () => null, + create: async () => { + throw new Prisma.PrismaClientKnownRequestError("Unique constraint failed", { + code: "P2002", + clientVersion: "test", + meta: { target: ["id"] }, + }); + }, + }, + } as unknown as PrismaClient; + + await expect( + syncReleases( + { prisma: stalePrisma, s3Client: syncS3Client }, + { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, + createAtDefaultRollout, + ), + ).rejects.toMatchObject({ + message: "[sync-releases] app 9.9.9: sync failed", + cause: expect.objectContaining({ code: "P2002" }), + }); + }); }); describe("scheduleReleaseSync", () => { From 02ee5a82cc8d141d791aa999c58565ddb0248750 Mon Sep 17 00:00:00 2001 From: Adam Shiervani Date: Sat, 19 Sep 2026 15:08:07 +0200 Subject: [PATCH 4/6] feat(releases): POST /releases/sync for the upload script (#83) The scheduled tick found a new version at most 30 minutes plus the settle window after upload. The upload script knows when its last object is written, so it can trigger the sync itself. One runner now owns the in-progress flag and serves both the timer and the endpoint. The endpoint is guarded by RELEASE_SYNC_TOKEN as a bearer token, compared in constant time, and is not registered without it. It runs with the settle window off and answers with the per-outcome counts, or 409 while a run is in progress. --- .env.example | 4 +++ src/auth.ts | 15 +++++++++++ src/errors.ts | 8 ++++++ src/index.ts | 18 ++++++++++--- src/release-sync.ts | 54 ++++++++++++++++++++++++++------------ src/releases.ts | 19 +++++++++++++- test/auth.test.ts | 31 ++++++++++++++++++++++ test/sync-releases.test.ts | 53 ++++++++++++++++++++++++++++++------- 8 files changed, 171 insertions(+), 31 deletions(-) create mode 100644 test/auth.test.ts diff --git a/.env.example b/.env.example index a7a246a..0ad2520 100644 --- a/.env.example +++ b/.env.example @@ -40,3 +40,7 @@ REAL_IP_HEADER= # ICE Servers for WebRTC, split by comma (e.g. stun:stun.l.google.com:19302,stun:stun1.l.google.com:19302) ICE_SERVERS= + +# Bearer token for POST /releases/sync (the upload script calls it after its last write). +# Leave empty to disable the endpoint. +RELEASE_SYNC_TOKEN= diff --git a/src/auth.ts b/src/auth.ts index 8985196..5856cca 100644 --- a/src/auth.ts +++ b/src/auth.ts @@ -1,3 +1,4 @@ +import { createHash, timingSafeEqual } from "crypto"; import { type NextFunction, type Request, type Response } from "express"; import * as jose from "jose"; import { UnauthorizedError } from "./errors"; @@ -56,3 +57,17 @@ export const authenticated = async (req: Request, res: Response, next: NextFunct next(); }; + +const sha256 = (value: string) => createHash("sha256").update(value).digest(); + +/** Guards a route with one static bearer token, compared in constant time. */ +export const bearerToken = (expected: string) => { + const expectedDigest = sha256(expected); + return (req: Request, res: Response, next: NextFunction) => { + const presented = req.headers.authorization?.match(/^Bearer (.+)$/)?.[1]; + if (!presented || !timingSafeEqual(sha256(presented), expectedDigest)) { + throw new UnauthorizedError("Invalid bearer token"); + } + next(); + }; +}; diff --git a/src/errors.ts b/src/errors.ts index 07f868b..936747d 100644 --- a/src/errors.ts +++ b/src/errors.ts @@ -31,6 +31,14 @@ export class ForbiddenError extends HttpError { } } +export class ConflictError extends HttpError { + constructor(message?: string, code?: string) { + super(409, message); + this.name = "ConflictError"; + this.code = code; + } +} + export class NotFoundError extends HttpError { constructor(message?: string, code?: string) { super(404, message); diff --git a/src/index.ts b/src/index.ts index 2a235c1..5d677e9 100644 --- a/src/index.ts +++ b/src/index.ts @@ -11,10 +11,10 @@ import * as Webrtc from "./webrtc"; import * as Releases from "./releases"; import { HttpError } from "./errors"; -import { authenticated } from "./auth"; +import { authenticated, bearerToken } from "./auth"; import { prisma } from "./db"; import { baseUrl, bucketName, s3Client } from "./s3"; -import { scheduleReleaseSync } from "./release-sync"; +import { createReleaseSyncRunner, scheduleReleaseSync } from "./release-sync"; import { initializeWebRTCSignaling } from "./webrtc-signaling"; declare global { @@ -49,6 +49,7 @@ declare global { ICE_SERVERS: string; ALLOWED_IDENTITIES?: string; + RELEASE_SYNC_TOKEN?: string; } } } @@ -116,12 +117,23 @@ app.get( }, ); +// One sync at a time, shared by the scheduled tick below and the upload +// script's POST /releases/sync. The route exists only with a token configured. +const releaseSync = createReleaseSyncRunner({ prisma, s3Client }, { bucketName, baseUrl }); + app.get("/releases", Releases.Retrieve); app.get( "/releases/system_recovery/latest", Releases.RetrieveLatestSystemRecovery, ); app.get("/releases/app/latest", Releases.RetrieveLatestApp); +if (process.env.RELEASE_SYNC_TOKEN) { + app.post( + "/releases/sync", + bearerToken(process.env.RELEASE_SYNC_TOKEN), + Releases.Sync(releaseSync), + ); +} app.get("/devices", authenticated, Devices.List); app.get("/devices/:id", authenticated, Devices.Retrieve); @@ -222,4 +234,4 @@ const server = app.listen(PORT, () => { initializeWebRTCSignaling(server); // Register new R2 releases at the default rollout, now and every 30 minutes. -scheduleReleaseSync({ prisma, s3Client }, { bucketName, baseUrl }); +scheduleReleaseSync(releaseSync); diff --git a/src/release-sync.ts b/src/release-sync.ts index 7a2f647..960c6c8 100644 --- a/src/release-sync.ts +++ b/src/release-sync.ts @@ -44,6 +44,8 @@ export type ReleaseOutcome = | "skipped" | "aborted"; +export type SyncStats = Record; + /** Settle window for unattended runs; shorter than the interval, so it costs at most one tick. */ export const UPLOAD_SETTLE_MS = 10 * 60 * 1000; @@ -322,8 +324,8 @@ export async function syncReleases( clients: SyncClients, config: SyncConfig, decide: ReleaseDecider, -): Promise> { - const stats: Record = { +): Promise { + const stats: SyncStats = { created: 0, "already-synced": 0, uploading: 0, @@ -369,38 +371,56 @@ export async function syncReleases( const RELEASE_SYNC_INTERVAL_MS = 30 * 60 * 1000; +/** Runs one unattended sync, or resolves to "busy" while another run is in progress. */ +export type ReleaseSyncRunner = ( + options?: Pick, +) => Promise; + /** - * Runs the sync now and then every `intervalMs`, inside the API process. - * A tick that fires while the previous run is still going is skipped, not - * queued. A failed run is logged and the schedule continues, so one bad R2 - * or DB response never stops future syncs. - * Returns a function that stops the schedule. + * One in-process sync at a time, shared by every trigger (the timer and the + * HTTP endpoint): a second run would only walk the bucket again and lose the + * (version, type) race to the first. */ -export function scheduleReleaseSync( +export function createReleaseSyncRunner( clients: SyncClients, config: SyncConfig, - intervalMs: number = RELEASE_SYNC_INTERVAL_MS, -): () => void { +): ReleaseSyncRunner { let running = false; - const run = async () => { + return async (options = {}) => { if (running) { - console.warn("[sync-releases] previous run still in progress, skipping this tick"); - return; + return "busy"; } running = true; try { - await syncReleases( + return await syncReleases( clients, - { uploadSettleMs: UPLOAD_SETTLE_MS, ...config }, + { uploadSettleMs: UPLOAD_SETTLE_MS, ...config, ...options }, createAtDefaultRollout, ); - } catch (error) { - console.error("[sync-releases] scheduled run failed", error); } finally { running = false; } }; +} + +/** + * Runs the sync now and then every `intervalMs`. A failed run is logged and + * the schedule continues. Returns a function that stops the schedule. + */ +export function scheduleReleaseSync( + runner: ReleaseSyncRunner, + intervalMs: number = RELEASE_SYNC_INTERVAL_MS, +): () => void { + const run = async () => { + try { + if ((await runner()) === "busy") { + console.warn("[sync-releases] previous run still in progress, skipping this tick"); + } + } catch (error) { + console.error("[sync-releases] scheduled run failed", error); + } + }; void run(); const timer = setInterval(run, intervalMs); diff --git a/src/releases.ts b/src/releases.ts index 33d5323..1763513 100644 --- a/src/releases.ts +++ b/src/releases.ts @@ -1,6 +1,7 @@ import { Request, Response } from "express"; import { prisma } from "./db"; -import { BadRequestError, InternalServerError, NotFoundError } from "./errors"; +import { BadRequestError, ConflictError, InternalServerError, NotFoundError } from "./errors"; +import type { ReleaseSyncRunner } from "./release-sync"; import semver from "semver"; import { GetObjectCommand, ListObjectsV2Command } from "@aws-sdk/client-s3"; @@ -710,3 +711,19 @@ function latestArtifactRedirect(kind: OtaKind) { } export const RetrieveLatestApp = latestArtifactRedirect("app"); + +/** + * POST /releases/sync: register every stable R2 version missing from the DB + * and answer with the per-outcome counts. Made for the upload script's last + * step, so the settle window is off: the caller vouches its last object is + * written. + */ +export function Sync(runner: ReleaseSyncRunner) { + return async (req: Request, res: Response) => { + const stats = await runner({ uploadSettleMs: 0 }); + if (stats === "busy") { + throw new ConflictError("A release sync is already in progress"); + } + return res.json(stats); + }; +} diff --git a/test/auth.test.ts b/test/auth.test.ts new file mode 100644 index 0000000..2adcc51 --- /dev/null +++ b/test/auth.test.ts @@ -0,0 +1,31 @@ +import type { Request, Response } from "express"; +import { describe, expect, it, vi } from "vitest"; + +import { bearerToken } from "../src/auth"; +import { UnauthorizedError } from "../src/errors"; + +describe("bearerToken", () => { + const guard = bearerToken("release-sync-secret"); + const res = {} as Response; + + function request(authorization?: string): Request { + return { headers: { authorization } } as unknown as Request; + } + + it("calls next for the configured token", () => { + const next = vi.fn(); + guard(request("Bearer release-sync-secret"), res, next); + expect(next).toHaveBeenCalledOnce(); + }); + + it.each([ + ["no header", undefined], + ["wrong token", "Bearer nope"], + ["missing scheme", "release-sync-secret"], + ["prefix of the token", "Bearer release-sync"], + ])("rejects %s without calling next", (_label, authorization) => { + const next = vi.fn(); + expect(() => guard(request(authorization), res, next)).toThrow(UnauthorizedError); + expect(next).not.toHaveBeenCalled(); + }); +}); diff --git a/test/sync-releases.test.ts b/test/sync-releases.test.ts index aa29e18..fd55f52 100644 --- a/test/sync-releases.test.ts +++ b/test/sync-releases.test.ts @@ -5,17 +5,21 @@ import { S3Client, } from "@aws-sdk/client-s3"; import { Prisma, PrismaClient } from "@prisma/client"; +import type { Request, Response } from "express"; import { afterEach, describe, expect, beforeEach, it, vi } from "vitest"; import { collectReleaseArtifacts, createAtDefaultRollout, + createReleaseSyncRunner, scheduleReleaseSync, syncReleases, UPLOAD_SETTLE_MS, type ReleaseDecider, type ReleaseType, } from "../src/release-sync"; +import { Sync } from "../src/releases"; +import { ConflictError } from "../src/errors"; import { otaFileForPrefix } from "../src/skus"; import { createAsyncIterable, s3Mock, testPrisma } from "./setup"; @@ -92,6 +96,13 @@ beforeEach(() => { s3Mock.on(ListObjectsV2Command).resolves({ Contents: [] }); }); +function newRunner() { + return createReleaseSyncRunner( + { prisma: testPrisma, s3Client: syncS3Client }, + { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, + ); +} + describe("syncReleases", () => { it("marks legacy app artifacts compatible with the default SKU only", async () => { mockS3HashFile("app", "9.9.1", "legacy-app-hash"); @@ -454,11 +465,7 @@ describe("scheduleReleaseSync", () => { mockS3HashFile("app", version, "app-hash"); const errorLog = vi.spyOn(console, "error").mockImplementation(() => {}); - stop = scheduleReleaseSync( - { prisma: testPrisma, s3Client: syncS3Client }, - { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, - INTERVAL_MS, - ); + stop = scheduleReleaseSync(newRunner(), INTERVAL_MS); await vi.waitFor(() => expect(errorLog).toHaveBeenCalledOnce()); expect( @@ -488,11 +495,7 @@ describe("scheduleReleaseSync", () => { .resolves({ CommonPrefixes: [] }); const warnLog = vi.spyOn(console, "warn").mockImplementation(() => {}); - stop = scheduleReleaseSync( - { prisma: testPrisma, s3Client: syncS3Client }, - { bucketName: SYNC_BUCKET, baseUrl: SYNC_BASE_URL }, - INTERVAL_MS, - ); + stop = scheduleReleaseSync(newRunner(), INTERVAL_MS); // Two ticks fire while the first run is still waiting on R2. await vi.advanceTimersByTimeAsync(INTERVAL_MS * 2); @@ -515,3 +518,33 @@ describe("scheduleReleaseSync", () => { expect(warnLog).toHaveBeenCalledTimes(2); }); }); + +describe("Sync handler", () => { + const request = {} as Request; + + it("registers a version uploaded moments ago and answers with the counts", async () => { + const version = "9.9.12"; + mockS3ListVersions("app", [version]); + mockS3HashFile("app", version, "fresh-hash"); + // Inside the settle window: the scheduled tick would defer this version, + // the endpoint trusts the caller and registers it. + mockS3UploadedAt("app", version, new Date(Date.now() - 60 * 1000)); + const res = { json: vi.fn() } as unknown as Response; + + await Sync(newRunner())(request, res); + + expect(res.json).toHaveBeenCalledWith(expect.objectContaining({ created: 1, uploading: 0 })); + const release = await testPrisma.release.findUniqueOrThrow({ + where: { version_type: { version, type: "app" } }, + }); + expect(release.rolloutPercentage).toBe(10); + }); + + it("answers 409 while a run is already in progress", async () => { + const busyRunner = vi.fn().mockResolvedValue("busy"); + + await expect(Sync(busyRunner)(request, { json: vi.fn() } as unknown as Response)).rejects.toThrow( + ConflictError, + ); + }); +}); From 2abf19dd0b3c2aa469f9716e89d457bda4e79b95 Mon Sep 17 00:00:00 2001 From: Adam Shiervani Date: Sat, 19 Sep 2026 15:19:55 +0200 Subject: [PATCH 5/6] fix(auth): accept the bearer scheme in any case (#85) Scheme names are case-insensitive (RFC 9110). The token is still compared exactly. --- src/auth.ts | 3 ++- test/auth.test.ts | 14 +++++++++----- 2 files changed, 11 insertions(+), 6 deletions(-) diff --git a/src/auth.ts b/src/auth.ts index 5856cca..96c895d 100644 --- a/src/auth.ts +++ b/src/auth.ts @@ -64,7 +64,8 @@ const sha256 = (value: string) => createHash("sha256").update(value).digest(); export const bearerToken = (expected: string) => { const expectedDigest = sha256(expected); return (req: Request, res: Response, next: NextFunction) => { - const presented = req.headers.authorization?.match(/^Bearer (.+)$/)?.[1]; + // The scheme name is case-insensitive (RFC 9110); the token is not. + const presented = req.headers.authorization?.match(/^Bearer +(.+)$/i)?.[1]; if (!presented || !timingSafeEqual(sha256(presented), expectedDigest)) { throw new UnauthorizedError("Invalid bearer token"); } diff --git a/test/auth.test.ts b/test/auth.test.ts index 2adcc51..d324c3c 100644 --- a/test/auth.test.ts +++ b/test/auth.test.ts @@ -12,17 +12,21 @@ describe("bearerToken", () => { return { headers: { authorization } } as unknown as Request; } - it("calls next for the configured token", () => { - const next = vi.fn(); - guard(request("Bearer release-sync-secret"), res, next); - expect(next).toHaveBeenCalledOnce(); - }); + it.each(["Bearer release-sync-secret", "bearer release-sync-secret", "BEARER release-sync-secret"])( + "calls next for %j", + authorization => { + const next = vi.fn(); + guard(request(authorization), res, next); + expect(next).toHaveBeenCalledOnce(); + }, + ); it.each([ ["no header", undefined], ["wrong token", "Bearer nope"], ["missing scheme", "release-sync-secret"], ["prefix of the token", "Bearer release-sync"], + ["token in a different case", "Bearer RELEASE-SYNC-SECRET"], ])("rejects %s without calling next", (_label, authorization) => { const next = vi.fn(); expect(() => guard(request(authorization), res, next)).toThrow(UnauthorizedError); From 6d5da4fa7a28d4c0fbd0134379761d1f5813d325 Mon Sep 17 00:00:00 2001 From: Adam Shiervani Date: Sat, 19 Sep 2026 15:21:02 +0200 Subject: [PATCH 6/6] fix(releases): scope the sync endpoint's settle bypass to one version (#84) The endpoint turned the settle window off for every version in the bucket. A call made while a different upload was still running registered that upload half done, and sync never revisits a row. The body now names the version the caller finished, and only that version skips the window. Without a body the call is a normal tick. --- src/release-sync.ts | 7 +++- src/releases.ts | 30 ++++++++++++--- test/sync-releases.test.ts | 75 ++++++++++++++++++++++++++++---------- 3 files changed, 86 insertions(+), 26 deletions(-) diff --git a/src/release-sync.ts b/src/release-sync.ts index 960c6c8..b176040 100644 --- a/src/release-sync.ts +++ b/src/release-sync.ts @@ -26,6 +26,8 @@ export interface SyncConfig { * (an operator at a terminal can judge the artifact list for themselves). */ uploadSettleMs?: number; + /** One version whose upload the caller vouches is complete: it skips the settle check. */ + settled?: { type: ReleaseType; version: string }; } export interface ReleaseArtifactInput { @@ -248,7 +250,8 @@ async function createRelease( // Listed after the artifact scan on purpose: an upload that was active at // any point during the scan leaves an object newer than the window, so the // snapshot above is discarded rather than registered. - if (config.uploadSettleMs) { + const vouched = config.settled?.type === type && config.settled.version === version; + if (config.uploadSettleMs && !vouched) { const newest = await newestUploadTime(clients.s3Client, config.bucketName, type, version); if (newest && Date.now() - newest.getTime() < config.uploadSettleMs) { console.log( @@ -373,7 +376,7 @@ const RELEASE_SYNC_INTERVAL_MS = 30 * 60 * 1000; /** Runs one unattended sync, or resolves to "busy" while another run is in progress. */ export type ReleaseSyncRunner = ( - options?: Pick, + options?: Pick, ) => Promise; /** diff --git a/src/releases.ts b/src/releases.ts index 1763513..30310fd 100644 --- a/src/releases.ts +++ b/src/releases.ts @@ -21,6 +21,7 @@ import { isKnownSku, legacyCompatibleSkus, otaArtifacts, + OTA_PREFIXES, type OtaKind, type Artifact, } from "./skus"; @@ -72,8 +73,12 @@ type RetrieveQuery = z.infer; * Parses query parameters and converts ZodError to BadRequestError. */ function parseQuery(schema: z.ZodSchema, req: Request): T { + return parseOrBadRequest(schema, req.query); +} + +function parseOrBadRequest(schema: z.ZodSchema, input: unknown): T { try { - return schema.parse(req.query); + return schema.parse(input); } catch (error) { if (error instanceof ZodError) { const message = error.issues.map((e: z.ZodIssue) => e.message).join(", "); @@ -712,15 +717,30 @@ function latestArtifactRedirect(kind: OtaKind) { export const RetrieveLatestApp = latestArtifactRedirect("app"); +const syncBodySchema = z + .object({ + type: z.string().refine(type => OTA_PREFIXES.includes(type), "Unknown release type"), + version: z.string().min(1), + }) + .partial() + .refine( + body => (body.type === undefined) === (body.version === undefined), + "type and version go together", + ) + .transform(body => + body.type && body.version ? { type: body.type, version: body.version } : undefined, + ); + /** * POST /releases/sync: register every stable R2 version missing from the DB - * and answer with the per-outcome counts. Made for the upload script's last - * step, so the settle window is off: the caller vouches its last object is - * written. + * and answer with the per-outcome counts. With `{ type, version }` in the + * body, that one version skips the settle window: the caller vouches its + * last object is written. Every other version keeps it. */ export function Sync(runner: ReleaseSyncRunner) { return async (req: Request, res: Response) => { - const stats = await runner({ uploadSettleMs: 0 }); + const settled = parseOrBadRequest(syncBodySchema, req.body ?? {}); + const stats = await runner({ settled }); if (stats === "busy") { throw new ConflictError("A release sync is already in progress"); } diff --git a/test/sync-releases.test.ts b/test/sync-releases.test.ts index fd55f52..3d340d4 100644 --- a/test/sync-releases.test.ts +++ b/test/sync-releases.test.ts @@ -19,7 +19,7 @@ import { type ReleaseType, } from "../src/release-sync"; import { Sync } from "../src/releases"; -import { ConflictError } from "../src/errors"; +import { BadRequestError, ConflictError } from "../src/errors"; import { otaFileForPrefix } from "../src/skus"; import { createAsyncIterable, s3Mock, testPrisma } from "./setup"; @@ -520,31 +520,68 @@ describe("scheduleReleaseSync", () => { }); describe("Sync handler", () => { - const request = {} as Request; + const FRESH_VERSION = "9.9.12"; + const OTHER_FRESH_VERSION = "9.9.13"; + + function request(body?: unknown): Request { + return { body } as Request; + } + + function response(): Response { + return { json: vi.fn() } as unknown as Response; + } + + /** Two versions inside the settle window; the scheduled tick would defer both. */ + function mockTwoFreshVersions() { + mockS3ListVersions("app", [FRESH_VERSION, OTHER_FRESH_VERSION]); + for (const version of [FRESH_VERSION, OTHER_FRESH_VERSION]) { + mockS3HashFile("app", version, `${version}-hash`); + mockS3UploadedAt("app", version, new Date(Date.now() - 60 * 1000)); + } + } - it("registers a version uploaded moments ago and answers with the counts", async () => { - const version = "9.9.12"; - mockS3ListVersions("app", [version]); - mockS3HashFile("app", version, "fresh-hash"); - // Inside the settle window: the scheduled tick would defer this version, - // the endpoint trusts the caller and registers it. - mockS3UploadedAt("app", version, new Date(Date.now() - 60 * 1000)); - const res = { json: vi.fn() } as unknown as Response; + it("registers only the vouched version and defers the other fresh one", async () => { + mockTwoFreshVersions(); + const res = response(); - await Sync(newRunner())(request, res); + await Sync(newRunner())(request({ type: "app", version: FRESH_VERSION }), res); - expect(res.json).toHaveBeenCalledWith(expect.objectContaining({ created: 1, uploading: 0 })); - const release = await testPrisma.release.findUniqueOrThrow({ - where: { version_type: { version, type: "app" } }, - }); - expect(release.rolloutPercentage).toBe(10); + expect(res.json).toHaveBeenCalledWith(expect.objectContaining({ created: 1, uploading: 1 })); + expect( + await testPrisma.release.findUnique({ + where: { version_type: { version: FRESH_VERSION, type: "app" } }, + }), + ).toMatchObject({ rolloutPercentage: 10 }); + expect( + await testPrisma.release.findUnique({ + where: { version_type: { version: OTHER_FRESH_VERSION, type: "app" } }, + }), + ).toBeNull(); + }); + + it("keeps the settle window for every version without a body", async () => { + mockTwoFreshVersions(); + const res = response(); + + await Sync(newRunner())(request(), res); + + expect(res.json).toHaveBeenCalledWith(expect.objectContaining({ created: 0, uploading: 2 })); + }); + + it.each([ + ["type without version", { type: "app" }], + ["version without type", { version: "1.0.0" }], + ["unknown type", { type: "firmware", version: "1.0.0" }], + ])("rejects %s without running", async (_label, body) => { + const run = vi.fn(); + + await expect(Sync(run)(request(body), response())).rejects.toThrow(BadRequestError); + expect(run).not.toHaveBeenCalled(); }); it("answers 409 while a run is already in progress", async () => { const busyRunner = vi.fn().mockResolvedValue("busy"); - await expect(Sync(busyRunner)(request, { json: vi.fn() } as unknown as Response)).rejects.toThrow( - ConflictError, - ); + await expect(Sync(busyRunner)(request(), response())).rejects.toThrow(ConflictError); }); });