Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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=
15 changes: 15 additions & 0 deletions src/auth.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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];

@adamshiervani adamshiervani Sep 19, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in #85: the scheme is matched case-insensitively with one or more spaces, and the token is still compared exactly.

if (!presented || !timingSafeEqual(sha256(presented), expectedDigest)) {
throw new UnauthorizedError("Invalid bearer token");
}
next();
};
};
8 changes: 8 additions & 0 deletions src/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
18 changes: 15 additions & 3 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -49,6 +49,7 @@ declare global {
ICE_SERVERS: string;

ALLOWED_IDENTITIES?: string;
RELEASE_SYNC_TOKEN?: string;
}
}
}
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
54 changes: 37 additions & 17 deletions src/release-sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,8 @@ export type ReleaseOutcome =
| "skipped"
| "aborted";

export type SyncStats = Record<ReleaseOutcome, number>;

/** Settle window for unattended runs; shorter than the interval, so it costs at most one tick. */
export const UPLOAD_SETTLE_MS = 10 * 60 * 1000;

Expand Down Expand Up @@ -322,8 +324,8 @@ export async function syncReleases(
clients: SyncClients,
config: SyncConfig,
decide: ReleaseDecider,
): Promise<Record<ReleaseOutcome, number>> {
const stats: Record<ReleaseOutcome, number> = {
): Promise<SyncStats> {
const stats: SyncStats = {
created: 0,
"already-synced": 0,
uploading: 0,
Expand Down Expand Up @@ -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<SyncConfig, "uploadSettleMs">,
) => Promise<SyncStats | "busy">;

/**
* 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);
Expand Down
19 changes: 18 additions & 1 deletion src/releases.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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 });

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Keep settling enabled for unrelated uploads

When another firmware version or type is still being uploaded, a POST made after an independent upload completes performs a bucket-wide scan with the settle window disabled. It can therefore persist the other upload's partial SKU set or stale hash, and syncReleases deliberately never updates that row after the remaining objects arrive. Scope the bypass to the specific completed upload, or retain settling for all unrelated versions.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in #84: the body names the version the caller finished, and only that version skips the settle window. Every other version keeps it.

if (stats === "busy") {
throw new ConflictError("A release sync is already in progress");
}
return res.json(stats);
};
}
31 changes: 31 additions & 0 deletions test/auth.test.ts
Original file line number Diff line number Diff line change
@@ -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();
});
});
53 changes: 43 additions & 10 deletions test/sync-releases.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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");
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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);
Expand All @@ -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,
);
});
});
Loading