Checkpoint web parity through Chromium input tasks
This commit is contained in:
57
tools/web/server-job-idempotency.mjs
Normal file
57
tools/web/server-job-idempotency.mjs
Normal file
@@ -0,0 +1,57 @@
|
||||
import crypto from "node:crypto";
|
||||
import fs from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { commitServerJobResult, createServerJobResultIdentity, verifyServerJobResultReceipt } from "./server-job-result-binding.mjs";
|
||||
|
||||
export const SERVER_JOB_IDEMPOTENCY_SCHEMA = 1;
|
||||
const inFlight = new Map();
|
||||
|
||||
function invalid(message) { throw new Error(`SERVER_JOB_IDEMPOTENCY_INVALID: ${message}`); }
|
||||
function outputHash(bytes) { return crypto.createHash("sha256").update(bytes).digest("hex"); }
|
||||
function keyFor(identity) { return `${identity.projectId}:${identity.requestId}`; }
|
||||
function receiptPath(directory, identity) { return path.join(directory, `${identity.requestId}.receipt.json`); }
|
||||
|
||||
async function readReceipt(file) {
|
||||
try { return JSON.parse(await fs.readFile(file, "utf8")); }
|
||||
catch (error) { if (error?.code === "ENOENT") return null; throw error; }
|
||||
}
|
||||
|
||||
async function writeReceipt(file, receipt) {
|
||||
const temporary = `${file}.${crypto.randomUUID()}.tmp`;
|
||||
await fs.writeFile(temporary, `${JSON.stringify(receipt, null, 2)}\n`, { flag: "wx", mode: 0o600 });
|
||||
await fs.rename(temporary, file);
|
||||
}
|
||||
|
||||
async function submit(identityValue, outputBytes, options) {
|
||||
const identity = createServerJobResultIdentity(identityValue);
|
||||
if (!(outputBytes instanceof Uint8Array) || outputBytes.byteLength < 1) invalid("output bytes are empty");
|
||||
if (typeof options?.receiptDirectory !== "string" || !path.isAbsolute(options.receiptDirectory)) invalid("receipt directory is invalid");
|
||||
if (typeof options?.outputDirectory !== "string" || !path.isAbsolute(options.outputDirectory)) invalid("output directory is invalid");
|
||||
await fs.mkdir(options.receiptDirectory, { recursive: true, mode: 0o700 });
|
||||
await fs.mkdir(options.outputDirectory, { recursive: true, mode: 0o700 });
|
||||
const file = receiptPath(options.receiptDirectory, identity);
|
||||
const expectedOutputSha256 = outputHash(outputBytes);
|
||||
const existing = await readReceipt(file);
|
||||
if (existing) {
|
||||
if (existing.schemaVersion !== SERVER_JOB_IDEMPOTENCY_SCHEMA || existing.outputSha256 !== expectedOutputSha256 || existing.requestId !== identity.requestId || existing.projectId !== identity.projectId || existing.sourceSha256 !== identity.sourceSha256 || existing.settingsSha256 !== identity.settingsSha256 || existing.buildSha256 !== identity.buildSha256 || existing.baseRevision !== identity.baseRevision) {
|
||||
throw new Error("SERVER_JOB_IDEMPOTENCY_CONFLICT: request is already bound to a different result");
|
||||
}
|
||||
await verifyServerJobResultReceipt(existing.result, identity, options.outputDirectory);
|
||||
return Object.freeze({ ...existing.result, reused: true, idempotencyKey: keyFor(identity), execution: "DISABLED" });
|
||||
}
|
||||
const result = await commitServerJobResult(identity, outputBytes, { outputDirectory: options.outputDirectory, expectedIdentity: identity });
|
||||
const stored = Object.freeze({ schemaVersion: SERVER_JOB_IDEMPOTENCY_SCHEMA, requestId: identity.requestId, projectId: identity.projectId, baseRevision: identity.baseRevision, sourceSha256: identity.sourceSha256, settingsSha256: identity.settingsSha256, buildSha256: identity.buildSha256, outputSha256: result.outputSha256, result });
|
||||
await writeReceipt(file, stored);
|
||||
return Object.freeze({ ...result, reused: false, idempotencyKey: keyFor(identity), execution: "DISABLED" });
|
||||
}
|
||||
|
||||
export async function submitIdempotentServerJobResult(identityValue, outputBytes, options) {
|
||||
const identity = createServerJobResultIdentity(identityValue);
|
||||
const key = `${options?.receiptDirectory ?? ""}:${keyFor(identity)}`;
|
||||
const running = inFlight.get(key);
|
||||
if (running) return Object.freeze({ ...(await running), reused: true });
|
||||
const current = Promise.resolve().then(() => submit(identity, outputBytes, options));
|
||||
inFlight.set(key, current);
|
||||
try { return await current; }
|
||||
finally { if (inFlight.get(key) === current) inFlight.delete(key); }
|
||||
}
|
||||
Reference in New Issue
Block a user