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/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/auth.ts b/src/auth.ts index 8985196..96c895d 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,18 @@ 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) => { + // 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"); + } + 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 b6fdecb..5d677e9 100644 --- a/src/index.ts +++ b/src/index.ts @@ -11,8 +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 { createReleaseSyncRunner, scheduleReleaseSync } from "./release-sync"; import { initializeWebRTCSignaling } from "./webrtc-signaling"; declare global { @@ -47,6 +49,7 @@ declare global { ICE_SERVERS: string; ALLOWED_IDENTITIES?: string; + RELEASE_SYNC_TOKEN?: string; } } } @@ -114,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); @@ -218,3 +232,6 @@ const server = app.listen(PORT, () => { }); initializeWebRTCSignaling(server); + +// Register new R2 releases at the default rollout, now and every 30 minutes. +scheduleReleaseSync(releaseSync); diff --git a/src/release-sync.ts b/src/release-sync.ts new file mode 100644 index 0000000..b176040 --- /dev/null +++ b/src/release-sync.ts @@ -0,0 +1,431 @@ +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; + /** One version whose upload the caller vouches is complete: it skips the settle check. */ + settled?: { type: ReleaseType; version: string }; +} + +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"; + +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; + +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)); +} + +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" + ); +} + +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. + 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( + `[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, 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"; + } + 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: SyncStats = { + 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) { + // 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).catch((error: unknown) => { + throw new Error(`[sync-releases] ${type} ${version}: sync failed`, { cause: error }); + }); + 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 one unattended sync, or resolves to "busy" while another run is in progress. */ +export type ReleaseSyncRunner = ( + options?: Pick, +) => Promise; + +/** + * 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 createReleaseSyncRunner( + clients: SyncClients, + config: SyncConfig, +): ReleaseSyncRunner { + let running = false; + + return async (options = {}) => { + if (running) { + return "busy"; + } + running = true; + try { + return await syncReleases( + clients, + { uploadSettleMs: UPLOAD_SETTLE_MS, ...config, ...options }, + createAtDefaultRollout, + ); + } 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); + return () => clearInterval(timer); +} diff --git a/src/releases.ts b/src/releases.ts index d34ba83..30310fd 100644 --- a/src/releases.ts +++ b/src/releases.ts @@ -1,15 +1,12 @@ 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, - 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, @@ -24,6 +21,7 @@ import { isKnownSku, legacyCompatibleSkus, otaArtifacts, + OTA_PREFIXES, type OtaKind, type Artifact, } from "./skus"; @@ -75,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(", "); @@ -101,15 +103,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 +127,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 +136,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 +161,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 +199,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 +345,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; } @@ -471,16 +427,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 +563,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); }), ); @@ -685,7 +649,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}`); } @@ -742,7 +706,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}`); } @@ -752,3 +716,34 @@ 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. 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 settled = parseOrBadRequest(syncBodySchema, req.body ?? {}); + const stats = await runner({ settled }); + if (stats === "busy") { + throw new ConflictError("A release sync is already in progress"); + } + return res.json(stats); + }; +} 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/auth.test.ts b/test/auth.test.ts new file mode 100644 index 0000000..d324c3c --- /dev/null +++ b/test/auth.test.ts @@ -0,0 +1,35 @@ +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.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); + expect(next).not.toHaveBeenCalled(); + }); +}); 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( diff --git a/test/sync-releases.test.ts b/test/sync-releases.test.ts index c0f1a92..3d340d4 100644 --- a/test/sync-releases.test.ts +++ b/test/sync-releases.test.ts @@ -4,12 +4,24 @@ import { ListObjectsV2Command, S3Client, } from "@aws-sdk/client-s3"; -import { describe, expect, beforeEach, it } from "vitest"; +import { Prisma, PrismaClient } from "@prisma/client"; +import type { Request, Response } from "express"; +import { afterEach, describe, expect, beforeEach, it, vi } from "vitest"; -import { collectReleaseArtifacts, syncReleases } from "../scripts/sync-releases"; +import { + collectReleaseArtifacts, + createAtDefaultRollout, + createReleaseSyncRunner, + scheduleReleaseSync, + syncReleases, + UPLOAD_SETTLE_MS, + type ReleaseDecider, + type ReleaseType, +} from "../src/release-sync"; +import { Sync } from "../src/releases"; +import { BadRequestError, ConflictError } from "../src/errors"; import { otaFileForPrefix } from "../src/skus"; -type ReleaseType = string; import { createAsyncIterable, s3Mock, testPrisma } from "./setup"; const DEFAULT_SKU = "jetkvm-v2"; @@ -20,9 +32,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 +86,24 @@ 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: [] }); +}); + +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"); @@ -178,16 +218,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 +269,319 @@ 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 () => [], + findUnique: (args: unknown) => testPrisma.release.findUnique(args as any), + 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"); + }); + + 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", () => { + 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(newRunner(), 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(newRunner(), 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); + }); +}); + +describe("Sync handler", () => { + 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 only the vouched version and defers the other fresh one", async () => { + mockTwoFreshVersions(); + const res = response(); + + await Sync(newRunner())(request({ type: "app", version: FRESH_VERSION }), res); + + 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(), response())).rejects.toThrow(ConflictError); + }); });