Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
JavaScript
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126import * as Effect from "effect/Effect";import { upgradeAssetHash } from "./fixtures/upgrade-artifact.mjs";import assert from "node:assert/strict";import { randomBytes, generateKeyPairSync } from "node:crypto";import { mkdtemp, readFile, readdir, rm, writeFile } from "node:fs/promises";import { createServer } from "node:net";import { tmpdir } from "node:os";import { join, resolve } from "node:path";import { test } from "node:test";import { parse } from "jsonc-parser";import { unstable_dev } from "wrangler";import { oauthBindings, oauthConfig } from "./fixtures/oauth-config.mjs";import { loadArtifact } from "../control-plane/artifact.ts";import "./fixtures/config.mjs";const id = () => randomBytes(16).toString("hex");const cookie = (r, name) => r.headers .getSetCookie() .find((c) => c.startsWith(name + "=")) ?.split(";")[0];// Keep local transport failures attributable without exposing response bodies,// OAuth inputs, provider values or cookies in native CI diagnostics.async function responseJson(response, method) { try { return await response.json(); } catch { const path = new URL(response.url).pathname; const contentType = response.headers.get("Content-Type") ?? "missing"; assert.fail( `${method} ${path} returned invalid JSON (HTTP ${response.status}; Content-Type: ${contentType})`, ); }}async function freePort() { const s = createServer(); await new Promise((r) => s.listen(0, "127.0.0.1", r)); const port = s.address().port; await new Promise((r) => s.close(r)); return port;}async function artifactEntry() { const manifest = JSON.parse( await readFile("dist/release/manifest.json", "utf8"), ); const files = {}; for (const file of [{ path: "manifest.json" }, ...manifest.files]) { const buffer = await readFile(`dist/release/${file.path}`); files[file.path] = buffer.buffer.slice( buffer.byteOffset, buffer.byteOffset + buffer.byteLength, ); } return { development: true, identity: { version: manifest.release, sourceRevision: manifest.sourceRevision, artifactDigest: ( await readFile("dist/release/manifest.sha256", "utf8") ).trim(), }, files, };}test("immutable catalog refuses dirty production, wrong pin and mutated bytes before effects", async () => { const entry = await artifactEntry(); const loaded = await Effect.runPromise(loadArtifact(entry, undefined, true)); assert.equal(loaded.identity.artifactDigest, entry.identity.artifactDigest); await assert.rejects(() => Effect.runPromise(loadArtifact(entry)), { message: "artifact_unavailable", }); await assert.rejects( () => Effect.runPromise( loadArtifact( entry, { ...entry.identity, artifactDigest: "0".repeat(64) }, true, ), ), { message: "artifact_unavailable" }, ); await assert.rejects( () => Effect.runPromise( loadArtifact( { ...entry, files: { ...entry.files, "worker/index.js": new ArrayBuffer(0) }, }, undefined, true, ), ), { message: "artifact_unavailable" }, ); assert.ok(!loaded.files.some((f) => f.path.endsWith(".map")));});test( "native Workflow provisions a pinned customer artifact and reconciles lost provider replies", { timeout: 480_000 }, async (t) => { const logs = []; for (const method of ["log", "warn", "error"]) { const original = console[method]; t.mock.method(console, method, (...args) => { logs.push(args.map(String).join(" ")); original(...args); }); } const temporary = await mkdtemp(join(tmpdir(), "flarebot-orchestrator-")); const port = await freePort(); const origin = `http://127.0.0.1:${port}`; const publisherOrigin = "https://publisher.test"; const base = parse(await readFile("wrangler.control-plane.jsonc", "utf8")); const configPath = join(temporary, "wrangler.json"); await writeFile( configPath, JSON.stringify({ ...base, name: "flarebot-orchestrator-test", main: resolve("tests/fixtures/orchestrator-worker.ts"), assets: { directory: resolve("dist/control-plane/client"), binding: "ASSETS", }, durable_objects: { bindings: [ ...base.durable_objects.bindings, { name: "PROVIDER", class_name: "Provider" }, ], }, exports: { ...base.exports, Provider: { type: "durable-object", storage: "sqlite" }, }, }), ); const keys = generateKeyPairSync("ed25519"); const publicKey = keys.publicKey .export({ type: "spki", format: "der" }) .subarray(-32) .toString("base64url"); const vars = { ...oauthBindings(publisherOrigin), UPGRADE_ASSET_HASH: await upgradeAssetHash(), FLAREBOT_CONTROL_PLANE: JSON.stringify({ ...oauthConfig(publisherOrigin), bridge: { keyId: "fixture-key", publicKey }, }), FLAREBOT_BRIDGE_SIGNING_KEY: keys.privateKey .export({ type: "pkcs8", format: "der" }) .toString("base64url"), }; let worker; const start = () => unstable_dev("tests/fixtures/orchestrator-worker.ts", { config: configPath, vars, local: true, ip: "127.0.0.1", port, inspectorPort: 0, persist: true, persistTo: temporary, logLevel: "error", experimental: { disableExperimentalWarning: true, watch: false }, }); t.after(async () => { await worker?.stop(); await rm(temporary, { recursive: true, force: true }); }); const call = (path, options = {}) => fetch(origin + path, { redirect: "manual", ...options, headers: { Origin: publisherOrigin, ...options.headers }, }); const admin = async (path, body = {}, session = "") => responseJson( await call("/__test__/" + path, { method: "POST", headers: { Cookie: session, "Content-Type": "application/json" }, body: JSON.stringify(body), }), "POST", ); const form = (path, session, values) => call(path, { method: "POST", headers: { Cookie: session, "Content-Type": "application/x-www-form-urlencoded", }, body: new URLSearchParams(values).toString(), }); async function login() { const response = await form("/auth/start", "", {}); const url = new URL(response.headers.get("Location")); const code = id(); await admin("code", { code, challenge: url.searchParams.get("code_challenge"), mode: "normal", }); const r = await call( `/auth/callback?code=${code}&state=${url.searchParams.get("state")}`, { headers: { Cookie: cookie(response, "__Host-flarebot-oauth") } }, ); return cookie(r, "__Host-flarebot-control-session"); } worker = await start(); const publicRelease = await call("/api/releases/latest"); assert.equal(publicRelease.status, 200); assert.deepEqual(Object.keys(await publicRelease.json()).sort(), [ "artifactDigest", "sourceRevision", "version", ]); const session = await login(); assert.ok(session); const selectAccount = async (accountId) => { const response = await form("/api/account", session, { accountId }); const data = await responseJson(response, "POST"); assert.equal(response.status, 200); assert.ok( data.selectedAccountId === accountId, "Selected account must match the request", ); }; await selectAccount("a".repeat(32)); const reserve = async () => { const r = await form("/api/installations", session, { requestId: id() }); assert.equal(r.status, 200); return (await responseJson(r, "POST")).installation; }; const inspect = () => admin("provider/inspect"); const state = async (record) => responseJson( await call(`/__test__/operation?id=${record.installationId}`, { headers: { Cookie: session }, }), "GET", ); const wait = async (record, status = "ready") => { const deadline = Date.now() + 200_000; let latest; while (Date.now() < deadline) { latest = await state(record); if (latest.record.status === status) return latest; if (latest.record.status === "failed" && status !== "failed") assert.fail(JSON.stringify(latest)); await new Promise((r) => setTimeout(r, 200)); } assert.fail(JSON.stringify(latest)); }; const startInstall = (record, requestId = id(), recover = false) => admin( "start", { id: record.installationId, requestId, recover }, session, ); let upgradeTarget; const upgradeForm = async (path, ownerSession, values) => { if (!upgradeTarget) upgradeTarget = ( await responseJson( await call("/api/installations", { headers: { Cookie: ownerSession }, }), "GET", ) ).latestRelease; assert.ok( upgradeTarget?.artifactDigest, "fixture latest release unavailable", ); return form(path, ownerSession, { ...values, target: JSON.stringify(upgradeTarget), }); }; let record; await t.test( "duplicate start and ambiguous Worker/application responses preserve one operation and secrets", async () => { await admin("provider/options", { loseWorker: true, loseContainer: true, assetExpired: true, loseGateway: true, }); record = await reserve(); const requestId = id(); const results = await Promise.all([ startInstall(record, requestId), startInstall(record, requestId), ]); assert.equal( results[0].installation.operationId, results[1].installation.operationId, ); const result = await wait(record); const { installedAt, ...installed } = result.record.installedRelease; assert.deepEqual(result.record.desiredRelease, installed); const remote = await inspect(); assert.equal(remote.gateway.id, "default"); assert.equal(remote.modelProvider, undefined); assert.equal( remote.trace.filter((x) => x === "POST ai-gateway/custom-providers") .length, 0, ); assert.ok( remote.trace.indexOf("POST ai-gateway/gateways") < remote.trace.indexOf( `POST workers/scripts/${record.resources.workerName}/assets-upload-session`, ), ); const uploaded = remote[`worker:${record.resources.workerName}`]; assert.ok(uploaded.secret); const metadata = await admin("metadata/inspect", {}, session); assert.ok(!JSON.stringify(metadata).includes(uploaded.secret)); assert.ok(!JSON.stringify(metadata).includes("oauth-secret-sentinel")); assert.ok(!JSON.stringify(result.status).includes(uploaded.secret)); const writes = remote.trace.filter( (x) => x === `PUT workers/scripts/${record.resources.workerName}`, ); assert.equal(writes.length, 1); assert.equal( remote.trace.filter((x) => x === "POST containers/applications") .length, 1, ); const artifact = await artifactEntry(); const manifest = JSON.parse( new TextDecoder().decode(artifact.files["manifest.json"]), ); for (const module of uploaded.modules) { const file = manifest.files.find( (f) => f.path === `worker/${module.name}`, ); assert.equal(module.sha256, file.sha256); assert.equal(module.mime, "application/javascript+module"); } assert.equal(uploaded.metadata.main_module, "index.js"); assert.deepEqual(uploaded.metadata.keep_bindings, [ "secret_text", "secret_key", "plain_text", "json", ]); assert.equal(uploaded.metadata.migrations, undefined); assert.deepEqual(Object.keys(uploaded.metadata.exports).sort(), [ "PersonalAgent", "Sandbox", ]); for (const file of remote.uploadedAssets) { const source = manifest.files.find((f) => f.assetHash === file.hash); assert.equal(file.sha256, source.sha256); assert.equal(file.mime, source.mime); } assert.equal( result.record.resources.personalAgentNamespaceId, uploaded.bindings.find((b) => b.name === "PersonalAgent") .namespace_id, ); assert.equal( result.record.resources.sandboxNamespaceId, uploaded.bindings.find((b) => b.name === "Sandbox").namespace_id, ); }, ); await t.test( "versioned upgrade pins exact from/to, inherits every customer binding and reconciles lost PUT/PATCH/rollout", async () => { const beforeRecord = (await state(record)).record; await admin("provider/customize", { name: record.resources.workerName, drift: "legacy-attachments", }); const before = (await inspect())[ `worker:${record.resources.workerName}` ]; await admin("provider/options", { upgradeLatest: true, retireSource: true, loseUpgradeWorker: true, losePatch: true, loseRollout: true, removeModelGateway: true, }); const requestId = id(); const url = `/api/installations/${record.installationId}/upgrade`; const responses = await Promise.all([ upgradeForm(url, session, { requestId }), upgradeForm(url, session, { requestId }), ]); const started = await Promise.all( responses.map((r) => responseJson(r, "POST")), ); assert.ok(started[0].installation, JSON.stringify(started)); assert.equal( started[0].installation.operationId, started[1].installation.operationId, ); assert.equal(started[0].installation.status, "updating"); assert.deepEqual( await admin( "invalid-upgrade-intent", { id: record.installationId }, session, ), { ok: false, error: "invalid_metadata" }, ); assert.deepEqual( started[0].installation.installedRelease, beforeRecord.installedRelease, ); // Catalog changes do not redirect the pinned active target. await admin("provider/options", { upgradeLatest: false, loseUpgradeWorker: true, losePatch: true, loseRollout: true, }); const done = (await wait(record)).record; assert.equal(done.installedRelease.version, "0.1.0-fixture.2"); assert.deepEqual(done.resources, beforeRecord.resources); const remote = await inspect(); assert.equal(remote.modelProvider, undefined); assert.equal( remote.trace.filter((x) => x === "POST ai-gateway/gateways").length, 2, ); const after = remote[`worker:${record.resources.workerName}`]; assert.equal(after.secret, before.secret); assert.notEqual(after.versionId, before.versionId); assert.notDeepEqual(after.modules, before.modules); assert.ok( after.sentBindings .filter( (b) => !["ASSETS", "FLAREBOT_INSTALLATION", "ATTACHMENTS"].includes(b.name), ) .every((b) => b.type === "inherit" && b.version_id === "latest"), ); assert.equal(before.bindings.some((b) => b.name === "ATTACHMENTS"), false); assert.deepEqual(after.sentBindings.find((b) => b.name === "ATTACHMENTS"), { name: "ATTACHMENTS", type: "r2_bucket", bucket_name: `flarebot-attachments-${record.installationId}`, }); const script = `workers/scripts/${record.resources.workerName}`; const upgradePut = remote.trace.lastIndexOf(`PUT ${script}`); const assetStage = remote.trace.lastIndexOf( `POST ${script}/assets-upload-session`, ); assert.ok( remote.trace.slice(0, assetStage).includes(`GET ${script}/versions`), "the newest version is checked before asset staging", ); assert.ok( remote.trace .slice(assetStage + 1, upgradePut) .includes(`GET ${script}/versions`), "the newest version is rechecked after staging before inheriting latest", ); assert.equal( after.sentBindings.find((b) => b.name === "FLAREBOT_SESSION_SECRET") .text, undefined, ); for (const b of before.bindings.filter( (b) => b.name !== "FLAREBOT_INSTALLATION", )) assert.deepEqual( after.bindings.find((a) => a.name === b.name), b, ); assert.deepEqual(after.customerMetadata, before.customerMetadata); assert.deepEqual(after.customerMetadata.tags, [ "customer-owned", "production", ]); assert.deepEqual(after.customerMetadata.annotations, { "workers/message": "Customer deployment note", "workers/tag": "customer-release", }); assert.equal( after.metadata.annotations["workers/triggered_by"], undefined, ); assert.deepEqual(after.annotations, { "workers/triggered_by": "upload", }); assert.deepEqual( JSON.parse( after.bindings.find((b) => b.name === "FLAREBOT_INSTALLATION").text, ).bridge, JSON.parse( before.bindings.find((b) => b.name === "FLAREBOT_INSTALLATION") .text, ).bridge, ); const app = remote[`application:${done.resources.sandboxApplicationId}`]; assert.deepEqual(app.configuration.environment_variables, [ { name: "CUSTOM_ENV", value: "private-container-value" }, ]); assert.equal( remote.trace.filter( (x) => x === `PUT workers/scripts/${record.resources.workerName}`, ).length, 2, ); const replay = await responseJson( await upgradeForm(url, session, { requestId }), "POST", ); assert.ok(replay.installation, JSON.stringify(replay)); assert.equal(replay.installation.operationId, done.operationId); const metadata = JSON.stringify( await admin("metadata/inspect", {}, session), ); for (const sentinel of [ "private-custom-value", "private-container-value", before.secret, ]) assert.ok(!metadata.includes(sentinel)); await admin("provider/options", {}); }, ); await t.test( "upgrade rejects unsupported config, namespace, runtime and unknown settings before any customer mutation", async () => { for (const drift of [ "config", "namespace", "assets", "container", "setting", "annotation", "tags", "newer-version", ]) { await admin("provider/options", {}); const target = await reserve(); await startInstall(target); await wait(target); await admin("provider/customize", { name: target.resources.workerName, drift, }); await admin("provider/options", { upgradeLatest: true }); const before = (await inspect()).trace.length; const response = await upgradeForm( `/api/installations/${target.installationId}/upgrade`, session, { requestId: id() }, ); assert.equal( (await responseJson(response, "POST")).error, "resource_conflict", drift, ); const trace = (await inspect()).trace.slice(before); assert.ok( !trace.some((x) => /^(POST|PUT|PATCH) /.test(x)), JSON.stringify(trace), ); assert.equal((await state(target)).record.status, "ready"); } await admin("provider/options", {}); }, ); await t.test( "source code edits after baseline prevent upload without a retained bundle", async () => { await admin("provider/options", {}); const target = await reserve(); await startInstall(target); const before = (await wait(target)).record; const traceStart = (await inspect()).trace.length; await admin("provider/options", { upgradeLatest: true, retireSource: true, editCodeDuringAssets: true, }); const response = await upgradeForm( `/api/installations/${target.installationId}/upgrade`, session, { requestId: id() }, ); const accepted = await responseJson(response, "POST"); assert.equal(response.status, 202); assert.ok( accepted.installation?.installationId === target.installationId, "Accepted upgrade must identify the requested installation", ); const failed = (await wait(target, "failed")).record; assert.equal(failed.errorCode, "resource_conflict"); assert.deepEqual(failed.installedRelease, before.installedRelease); assert.ok( !(await inspect()).trace .slice(traceStart) .some((entry) => /^(PUT|PATCH) /.test(entry)), ); await admin("provider/options", {}); }, ); await t.test( "a newer staged version appearing during asset staging prevents the upgrade PUT", async () => { await admin("provider/options", {}); const target = await reserve(); await startInstall(target); const before = (await wait(target)).record; const beforeRemote = await inspect(); const workerKey = `worker:${target.resources.workerName}`; await admin("provider/options", { upgradeLatest: true, stageNewerVersionDuringAssets: true, }); const response = await upgradeForm( `/api/installations/${target.installationId}/upgrade`, session, { requestId: id() }, ); assert.equal(response.status, 202); await responseJson(response, "POST"); const failed = (await wait(target, "failed")).record; assert.equal(failed.errorCode, "resource_conflict"); assert.deepEqual(failed.installedRelease, before.installedRelease); const after = await inspect(); const trace = after.trace.slice(beforeRemote.trace.length); assert.ok( trace.includes( `POST workers/scripts/${target.resources.workerName}/assets-upload-session`, ), "the newer version appeared only after asset staging began", ); assert.ok( !trace.some((entry) => /^(PUT|PATCH) /.test(entry)), JSON.stringify(trace), ); assert.equal( after[workerKey].versionId, beforeRemote[workerKey].versionId, ); assert.notEqual( after[workerKey].newestVersionId, after[workerKey].versionId, ); assert.deepEqual( after[workerKey].bindings, beforeRemote[workerKey].bindings, ); assert.equal(after[workerKey].secret, beforeRemote[workerKey].secret); await admin("provider/options", {}); }, ); await t.test( "upgrade cannot become ready when an accepted upload drops customer tags", async () => { for (const loseUpgradeWorker of [false, true]) { await admin("provider/options", {}); const target = await reserve(); await startInstall(target); const before = (await wait(target)).record; await admin("provider/customize", { name: target.resources.workerName, }); await admin("provider/options", { upgradeLatest: true, dropUpgradeTags: true, loseUpgradeWorker, }); const response = await upgradeForm( `/api/installations/${target.installationId}/upgrade`, session, { requestId: id() }, ); assert.equal(response.status, 202); await responseJson(response, "POST"); const failed = (await wait(target, "failed")).record; assert.equal(failed.errorCode, "resource_conflict"); assert.deepEqual(failed.installedRelease, before.installedRelease); const after = (await inspect())[ `worker:${target.resources.workerName}` ]; assert.deepEqual(after.metadata.tags, [ "customer-owned", "production", ]); assert.deepEqual(after.customerMetadata.tags, []); } await admin("provider/options", {}); }, ); await t.test( "retired in-flight targets stop for manual recovery without customer writes", async () => { await admin("provider/options", { healthFailure: true }); const target = await reserve(); await startInstall(target); const before = (await wait(target, "failed")).record; const traceStart = (await inspect()).trace.length; await admin("provider/options", { upgradeLatest: true, retireSource: true, }); for (const action of ["start", "recover"]) { const response = await form( `/api/installations/${target.installationId}/${action}`, session, { requestId: id() }, ); assert.equal(response.status, 503); assert.equal( (await responseJson(response, "POST")).error, "artifact_unavailable", ); } assert.deepEqual((await state(target)).record, before); assert.ok( !(await inspect()).trace .slice(traceStart) .some((entry) => /^(POST|PUT|PATCH|DELETE) /.test(entry)), ); await admin("provider/options", {}); }, ); await t.test( "failed upgrade retains last verified release and retries immutable target with no new secret", async () => { await admin("provider/options", {}); const target = await reserve(); await startInstall(target); const before = (await wait(target)).record; await admin("provider/options", { upgradeLatest: true, healthFailure: true, }); const response = await upgradeForm( `/api/installations/${target.installationId}/upgrade`, session, { requestId: id() }, ); const accepted = await responseJson(response, "POST"); assert.equal(response.status, 202); assert.ok( accepted.installation?.installationId === target.installationId, "Accepted upgrade must identify the requested installation", ); const failed = (await wait(target, "failed")).record; assert.equal(failed.errorCode, "health_failed"); assert.deepEqual(failed.installedRelease, before.installedRelease); assert.deepEqual(failed.resources, before.resources); await admin("provider/options", {}); const retry = await upgradeForm( `/api/installations/${target.installationId}/upgrade`, session, { requestId: id() }, ); const retried = await responseJson(retry, "POST"); assert.equal(retry.status, 202); assert.ok( retried.installation?.installationId === target.installationId, "Accepted retry must identify the requested installation", ); assert.equal( (await wait(target)).record.installedRelease.version, "0.1.0-fixture.2", ); assert.equal( (await inspect()).trace.filter( (x) => x === `PUT workers/scripts/${target.resources.workerName}`, ).length, 2, ); }, ); await t.test( "owner recovery after process restart adopts the upgrade effect once", async () => { await admin("provider/options", {}); const target = await reserve(); await startInstall(target); const original = (await wait(target)).record; const originalWorker = (await inspect())[ `worker:${target.resources.workerName}` ]; await admin("provider/options", { upgradeLatest: true, holdUpgradeWorker: true, }); const response = await upgradeForm( `/api/installations/${target.installationId}/upgrade`, session, { requestId: id() }, ); assert.equal(response.status, 202); // Finish the accepted HTTP exchange before polling or stopping Wrangler. assert.equal( (await responseJson(response, "POST")).installation.status, "updating", ); let uploaded; for (let n = 0; n < 100; n++) { uploaded = (await inspect())[`worker:${target.resources.workerName}`]; if (uploaded.versionId !== originalWorker.versionId) break; await new Promise((r) => setTimeout(r, 100)); } assert.notEqual(uploaded.versionId, originalWorker.versionId); await worker.stop(); worker = await start(); await admin("provider/options", {}); const recovered = await startInstall(target, id(), true); assert.ok(recovered.installation, JSON.stringify(recovered)); const done = (await wait(target)).record; assert.equal(done.installedRelease.version, "0.1.0-fixture.2"); assert.deepEqual(done.resources, original.resources); const after = await inspect(); assert.equal( after[`worker:${target.resources.workerName}`].secret, originalWorker.secret, ); assert.equal( after.trace.filter( (x) => x === `PUT workers/scripts/${target.resources.workerName}`, ).length, 2, ); }, ); await t.test( "failed boot never assigns installed release; reauthorization retries preserve upload and secret", async () => { await admin("provider/options", { healthFailure: true, emptyBuckets: true, }); const target = await reserve(); await startInstall(target); const failed = await wait(target, "failed"); assert.equal(failed.record.installedRelease, null); assert.equal(failed.record.errorCode, "health_failed"); const before = await inspect(); const secret = before[`worker:${target.resources.workerName}`].secret; await admin("provider/options", { emptyBuckets: true }); await startInstall(target); const completed = await wait(target); assert.ok(completed.record.installedRelease); const after = await inspect(); assert.equal( after[`worker:${target.resources.workerName}`].secret, secret, ); assert.equal( after.trace.filter( (x) => x === `PUT workers/scripts/${target.resources.workerName}`, ).length, 1, ); }, ); await t.test( "unknown Worker collision and malformed asset bucket fail safely", async () => { await admin("provider/options", { collision: true }); const target = await reserve(); await startInstall(target); const result = await wait(target, "failed"); assert.equal(result.record.errorCode, "resource_conflict"); await admin("provider/options", { unknownAsset: true }); const other = await reserve(); await startInstall(other); const bad = await wait(other, "failed"); assert.equal(bad.record.installedRelease, null); assert.ok(!JSON.stringify(bad).includes("provider-private-error")); }, ); await t.test( "durable startup gap repairs the same deterministic Workflow", async () => { await admin("provider/options", {}); const target = await reserve(); const requestId = id(); const pending = await admin( "pending-start", { id: target.installationId, requestId }, session, ); assert.ok(pending.operation.operationId); const started = await startInstall(target, requestId); assert.equal( started.installation.operationId, pending.operation.operationId, ); await wait(target); }, ); await t.test( "fresh owner recovery also repairs a lost start without the original request ID", async () => { await admin("provider/options", {}); const target = await reserve(); const pending = await admin( "pending-start", { id: target.installationId, requestId: id() }, session, ); const path = `/api/installations/${target.installationId}`; const before = await responseJson( await call(path, { headers: { Cookie: session } }), "GET", ); const traceBefore = (await inspect()).trace; await admin("provider/options", { unknownWorkflowFailure: true }); const unavailable = await form(`${path}/recover`, session, { requestId: id(), }); assert.equal(unavailable.status, 503); assert.equal( (await responseJson(unavailable, "POST")).error, "temporarily_unavailable", ); assert.deepEqual( await responseJson( await call(path, { headers: { Cookie: session } }), "GET", ), before, ); assert.deepEqual((await inspect()).trace, traceBefore); assert.deepEqual( await admin("workflow-exists", { id: pending.operation.operationId }), { exists: false }, ); await admin("provider/options", { liveMissingWorkflow: true }); const response = await form(`${path}/recover`, session, { requestId: id(), }); const recovered = await responseJson(response, "POST"); assert.equal( response.status, 202, `owner recovery must accept deployed Workflow absence: ${JSON.stringify(recovered)}`, ); assert.ok(recovered.installation); await wait(target); await admin("provider/options", {}); }, ); await t.test( "ready commit replay and recovery racing completion never downgrade success", async () => { await admin("provider/options", { loseReady: true, loseRetire: true }); const target = await reserve(); await startInstall(target); const ready = await wait(target); const replay = await startInstall(target, id(), true); assert.equal(replay.installation.status, "ready"); assert.equal(replay.installation.operationId, ready.record.operationId); const until = Date.now() + 10_000; let final; do { final = await state(target); if (final.status.status === "complete") break; await new Promise((r) => setTimeout(r, 100)); } while (Date.now() < until); assert.equal(final.status.output.status, "ready"); assert.equal(final.record.status, "ready"); }, ); await t.test( "an ambiguous absent upload requires explicit fresh-authorized recovery", async () => { await admin("provider/options", { rejectWorker: true }); const target = await reserve(); await startInstall(target); const failed = await wait(target, "failed"); assert.equal(failed.record.errorCode, "recovery_required"); const before = await inspect(); assert.equal( before.trace.filter( (x) => x === `PUT workers/scripts/${target.resources.workerName}`, ).length, 1, ); await admin("provider/options", {}); const recovered = await startInstall(target, id(), true); assert.ok(recovered.installation); await wait(target); }, ); await t.test( "native expanded lite configuration and lost PATCH/rollout replies reconcile once", async () => { await admin("provider/options", { driftContainer: true, expandedContainer: true, losePatch: true, loseRollout: true, }); const target = await reserve(); await startInstall(target); const result = await wait(target); const remote = await inspect(); const appId = result.record.resources.sandboxApplicationId; assert.equal( remote.trace.filter( (x) => x === `POST containers/applications/${appId}/rollouts`, ).length, 1, ); assert.equal(remote[`rollouts:${appId}`][0].status, "completed"); assert.deepEqual( remote[`rollouts:${appId}`][0].target_configuration .environment_variables, [ { name: "CUSTOM_CUSTOMER_SETTING", value: "retained-array-sentinel", }, ], ); }, ); await t.test( "owner recovery resumes native execution after process restart following remote Worker effect", async () => { await admin("provider/options", { holdWorker: true }); const target = await reserve(); await startInstall(target); const deadline = Date.now() + 15_000; let remote; while (Date.now() < deadline) { remote = await inspect(); if (remote[`worker:${target.resources.workerName}`]) break; await new Promise((r) => setTimeout(r, 100)); } assert.ok(remote[`worker:${target.resources.workerName}`]); const secret = remote[`worker:${target.resources.workerName}`].secret; await worker.stop(); worker = await start(); const interrupted = await state(target); assert.equal(interrupted.record.status, "installing"); await startInstall(target, id(), true); await wait(target); const after = await inspect(); assert.equal( after[`worker:${target.resources.workerName}`].secret, secret, ); assert.equal( after.trace.filter( (x) => x === `PUT workers/scripts/${target.resources.workerName}`, ).length, 1, ); }, ); await t.test( "selected account cannot retarget an installation and provider revocation is durable", async () => { await admin("provider/options", {}); const target = await reserve(); await selectAccount("b".repeat(32)); assert.equal((await startInstall(target)).error, "account_denied"); await selectAccount("a".repeat(32)); await admin("provider/options", { revoked: true }); await startInstall(target); const failed = await wait(target, "failed"); assert.equal(failed.record.errorCode, "reauthorization_required"); assert.equal(failed.record.installedRelease, null); }, ); await t.test( "native metadata and Workflow result survive full Wrangler restart", async () => { const previous = await state(record); await worker.stop(); worker = await start(); const restarted = await state(record); assert.deepEqual(restarted.record, previous.record); assert.deepEqual(restarted.status.output, previous.status.output); const secrets = Object.values(await inspect()) .filter((value) => value?.secret) .map((value) => value.secret); await worker.stop(); worker = null; const workflowFiles = ( await readdir(temporary, { recursive: true }) ).filter( (path) => /workflow/i.test(path) && /\.sqlite(?:-wal)?$/.test(path), ); assert.ok( workflowFiles.length > 0, "inspect actual native Workflow SQLite files", ); let inspectedParameters = false; for (const path of workflowFiles) { const bytes = await readFile(join(temporary, path)); inspectedParameters ||= bytes.includes(previous.record.operationId); for (const value of [ "oauth-secret-sentinel", "provider-private-error-sentinel", "retained-array-sentinel", ...secrets, ]) { assert.ok(!bytes.includes(value)); assert.ok(!logs.join("\n").includes(value)); } } assert.ok( inspectedParameters, "native persisted parameters include the safe operation identifier", ); worker = await start(); }, ); },);