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
7 changes: 5 additions & 2 deletions src/release-sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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<SyncConfig, "uploadSettleMs">,
options?: Pick<SyncConfig, "settled">,
) => Promise<SyncStats | "busy">;

/**
Expand Down
30 changes: 25 additions & 5 deletions src/releases.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import {
isKnownSku,
legacyCompatibleSkus,
otaArtifacts,
OTA_PREFIXES,
type OtaKind,
type Artifact,
} from "./skus";
Expand Down Expand Up @@ -72,8 +73,12 @@ type RetrieveQuery = z.infer<typeof retrieveQuerySchema>;
* Parses query parameters and converts ZodError to BadRequestError.
*/
function parseQuery<T>(schema: z.ZodSchema<T>, req: Request): T {
return parseOrBadRequest(schema, req.query);
}

function parseOrBadRequest<T>(schema: z.ZodSchema<T>, 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(", ");
Expand Down Expand Up @@ -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");
}
Expand Down
75 changes: 56 additions & 19 deletions test/sync-releases.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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);
});
});
Loading