diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 5ff911b..11bed82 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -39,6 +39,8 @@ jobs: - run: pnpm build:control-plane - run: pnpm test:oauth - run: pnpm test:ownership + - run: pnpm test:orchestrator + - run: pnpm test:bridge - run: pnpm test:web - run: pnpm test:browser - run: docker info diff --git a/control-plane/artifact-types.ts b/control-plane/artifact-types.ts new file mode 100644 index 0000000..520897c --- /dev/null +++ b/control-plane/artifact-types.ts @@ -0,0 +1,41 @@ +import type { ReleaseIdentity } from "./installation-metadata.ts"; + +export interface CatalogEntry { + identity: ReleaseIdentity; + development: boolean; + files: Record; +} +export interface ArtifactFile { + path: string; + size: number; + sha256: string; + mime: string; + assetHash?: string; +} +export interface Deployment { + main: string; + no_bundle: true; + compatibility_date: string; + compatibility_flags: string[]; + assets: { directory: string; binding: "ASSETS" }; + durable_objects: { bindings: { name: string; class_name: string }[] }; + exports: Record; + containers: { + name: string; + class_name: "Sandbox"; + image: string; + instance_type: "lite"; + max_instances: number; + }[]; + ai: { binding: "AI" }; + browser: { binding: "BROWSER" }; + worker_loaders: { binding: "LOADER" }[]; + observability: { enabled: false }; + keep_vars: true; +} +export interface Artifact { + identity: ReleaseIdentity; + files: ArtifactFile[]; + bytes: Record; + deployment: Deployment; +} diff --git a/control-plane/artifact.ts b/control-plane/artifact.ts new file mode 100644 index 0000000..0876ed3 --- /dev/null +++ b/control-plane/artifact.ts @@ -0,0 +1,171 @@ +import { z } from "zod"; +import { + releaseIdentity, + sameRelease, + type ReleaseIdentity, +} from "./installation-metadata.ts"; +import { fail } from "./deployment-errors.ts"; +import type { Artifact, CatalogEntry } from "./artifact-types.ts"; + +const fileSchema = z.strictObject({ + path: z + .string() + .regex(/^(?:deployment\.json|(?:worker|assets)\/[A-Za-z0-9_./-]+)$/), + size: z + .number() + .int() + .nonnegative() + .max(25 * 1024 * 1024), + sha256: z.string().regex(/^[a-f0-9]{64}$/), + mime: z.string().regex(/^[a-z0-9.+-]+\/[a-z0-9.+-]+$/), + assetHash: z + .string() + .regex(/^[a-f0-9]{32}$/) + .optional(), +}); +const deploymentSchema = z.strictObject({ + main: z.literal("./worker/index.js"), + no_bundle: z.literal(true), + compatibility_date: z.string().regex(/^\d{4}-\d{2}-\d{2}$/), + compatibility_flags: z.array(z.literal("nodejs_compat")).length(1), + assets: z.strictObject({ + directory: z.literal("./assets"), + binding: z.literal("ASSETS"), + }), + durable_objects: z.strictObject({ + bindings: z + .array( + z.strictObject({ + name: z.enum(["PersonalAgent", "Sandbox"]), + class_name: z.enum(["PersonalAgent", "Sandbox"]), + }), + ) + .length(2), + }), + exports: z.strictObject({ + PersonalAgent: z.strictObject({ + type: z.literal("durable-object"), + storage: z.literal("sqlite"), + }), + Sandbox: z.strictObject({ + type: z.literal("durable-object"), + storage: z.literal("sqlite"), + }), + }), + containers: z + .array( + z.strictObject({ + name: z.literal("flarebot-shell"), + class_name: z.literal("Sandbox"), + image: z + .string() + .regex( + /^docker\.io\/cloudflare\/sandbox:0\.12\.9@sha256:[a-f0-9]{64}$/, + ), + instance_type: z.literal("lite"), + max_instances: z.literal(4), + }), + ) + .length(1), + ai: z.strictObject({ binding: z.literal("AI") }), + browser: z.strictObject({ binding: z.literal("BROWSER") }), + worker_loaders: z + .array(z.strictObject({ binding: z.literal("LOADER") })) + .length(1), + observability: z.strictObject({ enabled: z.literal(false) }), + keep_vars: z.literal(true), +}); +export async function digest(bytes: ArrayBuffer) { + return [...new Uint8Array(await crypto.subtle.digest("SHA-256", bytes))] + .map((b) => b.toString(16).padStart(2, "0")) + .join(""); +} +const json = (bytes: ArrayBuffer) => + JSON.parse(new TextDecoder().decode(bytes)); +// Byte verification precedes every external effect. Data modules retain the +// original buffers; verification does not accumulate another copy of the release. +export async function loadArtifact( + entry: CatalogEntry, + pinned?: ReleaseIdentity, + development = false, +): Promise { + try { + const identity = releaseIdentity.parse(entry.identity); + if ( + (entry.development && !development) || + (pinned && !sameRelease(identity, pinned)) + ) + fail("artifact_unavailable"); + const manifestBytes = entry.files["manifest.json"]; + if ( + !manifestBytes || + (await digest(manifestBytes)) !== identity.artifactDigest + ) + fail("artifact_unavailable"); + const manifest = json(manifestBytes); + if ( + manifest.schemaVersion !== 1 || + manifest.kind !== "flarebot-customer-runtime" || + manifest.configurationVersion !== 1 || + manifest.release !== identity.version || + manifest.sourceRevision !== identity.sourceRevision || + (manifest.sourceDirty && !development) || + typeof manifest.sourceDirty !== "boolean" || + manifest.deployment !== "deployment.json" + ) + fail("artifact_unavailable"); + const files = z.array(fileSchema).min(3).max(256).parse(manifest.files); + const paths = new Set(); + let total = 0; + for (const file of files) { + if ( + paths.has(file.path) || + file.path.split("/").some((p) => !p || p === "." || p === "..") || + file.path.endsWith(".map") + ) + fail("artifact_unavailable"); + paths.add(file.path); + total += file.size; + const bytes = entry.files[file.path]; + if ( + !bytes || + bytes.byteLength !== file.size || + (await digest(bytes)) !== file.sha256 || + file.path.startsWith("assets/") !== !!file.assetHash + ) + fail("artifact_unavailable"); + if ( + file.path.startsWith("worker/") && + (!file.path.endsWith(".js") || file.mime !== "application/javascript") + ) + fail("artifact_unavailable"); + } + if ( + files + .filter((f) => f.path.startsWith("worker/")) + .reduce((sum, f) => sum + f.size, 0) > + 16 * 1024 * 1024 || + total > 24 * 1024 * 1024 || + Object.keys(entry.files).length !== paths.size + 1 || + !paths.has("worker/index.js") || + !paths.has("deployment.json") + ) + fail("artifact_unavailable"); + const deployment = deploymentSchema.parse( + json(entry.files["deployment.json"]), + ); + const bindings = deployment.durable_objects.bindings; + if ( + !bindings.some( + (b) => b.name === "PersonalAgent" && b.class_name === "PersonalAgent", + ) || + !bindings.some((b) => b.name === "Sandbox" && b.class_name === "Sandbox") + ) + fail("artifact_unavailable"); + if (deployment.containers[0].image !== manifest.shell?.image) + fail("artifact_unavailable"); + return { identity, files, bytes: entry.files, deployment }; + } catch { + return fail("artifact_unavailable"); + } +} diff --git a/control-plane/catalog.ts b/control-plane/catalog.ts new file mode 100644 index 0000000..90aede9 --- /dev/null +++ b/control-plane/catalog.ts @@ -0,0 +1,7 @@ +import entry from "../dist/catalog/catalog.ts"; +import { loadArtifact } from "./artifact.ts"; +import type { ReleaseIdentity } from "./installation-metadata.ts"; +export const customerArtifact = ( + pinned?: ReleaseIdentity, + development = false, +) => loadArtifact(entry, pinned, development); diff --git a/control-plane/cloudflare.ts b/control-plane/cloudflare.ts index 36bf881..4cdc620 100644 --- a/control-plane/cloudflare.ts +++ b/control-plane/cloudflare.ts @@ -21,7 +21,7 @@ async function requestJson( try { response = await network(url, { ...init, - redirect: "error", + redirect: "manual", signal: AbortSignal.timeout(10_000), }); if (!response.body) { @@ -252,7 +252,7 @@ export async function revoke( method: "POST", headers: clientAuthentication(config, form), body: form.toString(), - redirect: "error", + redirect: "manual", signal: AbortSignal.timeout(10_000), }); await response.body?.cancel(); diff --git a/control-plane/config.ts b/control-plane/config.ts index 40b315f..6a0a7eb 100644 --- a/control-plane/config.ts +++ b/control-plane/config.ts @@ -8,6 +8,9 @@ import type { AuthVault } from "./vault.ts"; import type { InstallationRegistry } from "./installation-registry.ts"; import release from "../package.json"; export interface Env extends ControlPlaneConfigBindings { + INSTALLATION_WORKFLOW: Workflow< + import("./installation-workflow.ts").InstallationParams + >; INSTALLATIONS: DurableObjectNamespace; AUTH_VAULT: DurableObjectNamespace; ASSETS: Fetcher; diff --git a/control-plane/data.d.ts b/control-plane/data.d.ts new file mode 100644 index 0000000..aeca78c --- /dev/null +++ b/control-plane/data.d.ts @@ -0,0 +1,8 @@ +declare module "*.bin" { + const bytes: ArrayBuffer; + export default bytes; +} +declare module "*catalog/catalog.ts" { + const entry: import("./artifact-types.ts").CatalogEntry; + export default entry; +} diff --git a/control-plane/deployment-api.ts b/control-plane/deployment-api.ts new file mode 100644 index 0000000..6a7fa80 --- /dev/null +++ b/control-plane/deployment-api.ts @@ -0,0 +1,629 @@ +import { digest } from "./artifact.ts"; +import type { Artifact } from "./artifact-types.ts"; +import type { Installation } from "./installation-metadata.ts"; +import { fail } from "./deployment-errors.ts"; + +export type DeploymentNetwork = typeof fetch; +type Binding = { + type: string; + name: string; + class_name?: string; + script_name?: string; + namespace_id?: string; + text?: string; + json?: string; +}; +type Settings = { bindings: Binding[] }; +export type Application = { + id: string; + name: string; + configuration: { image: string; instance_type: string }; + max_instances: number; + scheduling_policy: string; + durable_objects: { namespace_id: string }; +}; +export type Rollout = { + id: string; + description: string; + target_configuration: { image: string; instance_type: string }; + status?: string; +}; +const id = (value: unknown): value is string => + typeof value === "string" && + /^(?:[a-f0-9]{32}|[a-f0-9]{8}(?:-[a-f0-9]{4}){3}-[a-f0-9]{12})$/.test(value); +const object = (v: unknown): v is Record => + !!v && typeof v === "object" && !Array.isArray(v); +const token = (value: unknown): value is string => + typeof value === "string" && + value.length > 0 && + value.length <= 16_384 && + /^[A-Za-z0-9_.-]+$/.test(value); +function normalizedConfiguration(value: unknown) { + if (!object(value) || typeof value.image !== "string") + fail("resource_conflict"); + const config = { ...value }; + const lite = + config.vcpu === 0.0625 && + config.memory_mib === 256 && + object(config.disk) && + config.disk.size_mb === 2000; + if ( + lite && + Object.keys(config.disk as Record).some( + (key) => key !== "size_mb", + ) + ) + fail("resource_conflict"); + if (lite) { + config.instance_type = "lite"; + for (const key of ["vcpu", "memory", "memory_mib", "disk"]) + delete config[key]; + } + if (typeof config.instance_type !== "string") fail("resource_conflict"); + return config; +} +const eq = (a: unknown, b: unknown): boolean => { + if (a === b) return true; + if (Array.isArray(a) && Array.isArray(b)) + return a.length === b.length && a.every((v, i) => eq(v, b[i])); + if (!object(a) || !object(b)) return false; + const keys = Object.keys(a); + return ( + keys.length === Object.keys(b).length && keys.every((k) => eq(a[k], b[k])) + ); +}; +async function boundedJson(response: Response) { + const reader = response.body?.getReader(); + if (!reader) fail("temporarily_unavailable"); + const chunks: Uint8Array[] = []; + let size = 0; + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + size += value.length; + if (size > 2 * 1024 * 1024) { + await reader.cancel(); + fail("temporarily_unavailable"); + } + chunks.push(value); + } + const bytes = new Uint8Array(size); + let offset = 0; + for (const chunk of chunks) { + bytes.set(chunk, offset); + offset += chunk.length; + } + try { + return JSON.parse(new TextDecoder().decode(bytes)); + } catch { + fail("temporarily_unavailable"); + } +} +// Fixed native v4 endpoint adapter. No arbitrary host/path from browser input, +// automatic mutation retries, redirect following, or provider error propagation. +export class DeploymentAPI { + constructor( + private readonly accessToken: string, + private readonly record: Installation, + private readonly network: DeploymentNetwork = fetch, + ) {} + private async request( + path: string, + method = "GET", + body?: unknown, + notFound = false, + authorization = this.accessToken, + ): Promise { + let response: Response; + try { + response = await this.network( + `https://api.cloudflare.com/client/v4/accounts/${this.record.accountId}/${path}`, + { + method, + redirect: "manual", + signal: AbortSignal.timeout(60_000), + headers: { + Authorization: `Bearer ${authorization}`, + ...(body instanceof FormData + ? {} + : body === undefined + ? {} + : { "Content-Type": "application/json" }), + }, + ...(body === undefined + ? {} + : { body: body instanceof FormData ? body : JSON.stringify(body) }), + }, + ); + } catch { + return fail("temporarily_unavailable"); + } + if (response.status === 401) fail("reauthorization_required"); + if (response.status === 403) fail("account_denied"); + if (response.status === 404 && notFound) return null; + const envelope = await boundedJson(response); + if ( + !response.ok || + !object(envelope) || + envelope.success !== true || + !("result" in envelope) + ) + fail("temporarily_unavailable"); + return envelope.result; + } + private get script() { + return `workers/scripts/${this.record.resources.workerName}`; + } + async origin() { + const result = await this.request("workers/subdomain"); + if ( + !object(result) || + typeof result.subdomain !== "string" || + !/^[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?$/.test(result.subdomain) + ) + fail("setup_required"); + return `https://${this.record.resources.workerName}.${result.subdomain}.workers.dev`; + } + async settings(): Promise { + const result = await this.request( + `${this.script}/settings`, + "GET", + undefined, + true, + ); + if (result === null) return null; + if ( + !object(result) || + !Array.isArray(result.bindings) || + result.bindings.length > 256 || + result.bindings.some( + (b) => + !object(b) || + typeof b.type !== "string" || + typeof b.name !== "string", + ) + ) + fail("resource_conflict"); + return { bindings: result.bindings as Binding[] }; + } + verifyWorker(settings: Settings, installationConfig: unknown) { + const bindings = settings.bindings; + if (new Set(bindings.map((b) => b.name)).size !== bindings.length) + fail("resource_conflict"); + const byName = (name: string, type: string) => { + const binding = bindings.find((b) => b.name === name); + if (binding?.type !== type) fail("resource_conflict"); + return binding; + }; + let installed: unknown; + try { + installed = JSON.parse( + byName("FLAREBOT_INSTALLATION", "plain_text").text!, + ); + } catch { + fail("resource_conflict"); + } + if ( + !eq(installed, installationConfig) || + byName("FLAREBOT_MODE", "plain_text").text !== "customer-runtime" || + byName("FLAREBOT_ENV", "plain_text").text !== "production" + ) + fail("resource_conflict"); + for (const [name, type] of [ + ["ASSETS", "assets"], + ["AI", "ai"], + ["BROWSER", "browser"], + ["LOADER", "worker_loader"], + ["FLAREBOT_SESSION_SECRET", "secret_text"], + ]) + byName(name, type); + for (const name of ["PersonalAgent", "Sandbox"]) { + const b = byName(name, "durable_object_namespace"); + if ( + b.class_name !== name || + (b.script_name !== undefined && + b.script_name !== this.record.resources.workerName) + ) + fail("resource_conflict"); + } + } + async upload( + artifact: Artifact, + installationConfig: unknown, + bootstrapSecret: string, + beforeUpload: () => Promise, + ) { + const assets = artifact.files.filter((f) => f.assetHash); + const manifest = Object.fromEntries( + assets.map((f) => [ + `/${f.path.slice("assets/".length)}`, + { hash: f.assetHash!, size: f.size }, + ]), + ); + const session = await this.request( + `${this.script}/assets-upload-session`, + "POST", + { manifest }, + ); + if ( + !object(session) || + !token(session.jwt) || + !Array.isArray(session.buckets) || + session.buckets.length > assets.length + ) + fail("temporarily_unavailable"); + const byHash = new Map(assets.map((f) => [f.assetHash!, f])); + const seen = new Set(); + let completion = session.jwt; + let completed = session.buckets.length === 0; + for (const bucket of session.buckets) { + if ( + !Array.isArray(bucket) || + !bucket.length || + bucket.length > assets.length + ) + fail("temporarily_unavailable"); + const form = new FormData(); + for (const hash of bucket) { + const file = byHash.get(hash); + if (!file || seen.has(hash)) fail("temporarily_unavailable"); + seen.add(hash); + // Base64 encoding is transient per bucket; source bytes remain opaque. + const bytes = new Uint8Array(artifact.bytes[file.path]); + let binary = ""; + for (let i = 0; i < bytes.length; i += 8192) + binary += String.fromCharCode(...bytes.subarray(i, i + 8192)); + form.append(hash, new Blob([btoa(binary)], { type: file.mime }), hash); + } + const result = await this.request( + "workers/assets/upload?base64=true", + "POST", + form, + false, + session.jwt, + ); + if (!object(result)) fail("temporarily_unavailable"); + if (result.jwt !== undefined) { + if (!token(result.jwt)) fail("temporarily_unavailable"); + completion = result.jwt; + completed = true; + } + } + if (!completed) fail("temporarily_unavailable"); + const d = artifact.deployment; + const bindings = [ + { type: "assets", name: "ASSETS" }, + ...d.durable_objects.bindings.map((b) => ({ + type: "durable_object_namespace", + ...b, + })), + { type: "ai", name: "AI" }, + { type: "browser", name: "BROWSER" }, + { type: "worker_loader", name: "LOADER" }, + { type: "plain_text", name: "FLAREBOT_MODE", text: "customer-runtime" }, + { type: "plain_text", name: "FLAREBOT_ENV", text: "production" }, + { + type: "plain_text", + name: "FLAREBOT_INSTALLATION", + text: JSON.stringify(installationConfig), + }, + { + type: "secret_text", + name: "FLAREBOT_SESSION_SECRET", + text: bootstrapSecret, + }, + ]; + const form = new FormData(); + form.append( + "metadata", + new Blob( + [ + JSON.stringify({ + main_module: "index.js", + compatibility_date: d.compatibility_date, + compatibility_flags: d.compatibility_flags, + bindings, + exports: d.exports, + containers: [ + { + name: this.record.resources.sandboxApplicationName, + class_name: "Sandbox", + }, + ], + keep_bindings: ["secret_text", "secret_key", "plain_text", "json"], + observability: d.observability, + assets: { jwt: completion, config: {} }, + }), + ], + { type: "application/json" }, + ), + ); + for (const file of artifact.files.filter((f) => + f.path.startsWith("worker/"), + )) { + const name = file.path.slice("worker/".length); + form.append( + name, + new Blob([artifact.bytes[file.path]], { + type: "application/javascript+module", + }), + name, + ); + } + await beforeUpload(); + await this.request(this.script, "PUT", form); + } + async verifyContent(artifact: Artifact) { + const files = artifact.files.filter((file) => + file.path.startsWith("worker/"), + ); + const limit = files.reduce((sum, file) => sum + file.size, 0) + 65_536; + let response: Response; + try { + response = await this.network( + `https://api.cloudflare.com/client/v4/accounts/${this.record.accountId}/${this.script}/content/v2`, + { + redirect: "manual", + signal: AbortSignal.timeout(60_000), + headers: { Authorization: `Bearer ${this.accessToken}` }, + }, + ); + } catch { + fail("temporarily_unavailable"); + } + if (response.status === 401) fail("reauthorization_required"); + if (response.status === 403) fail("account_denied"); + if ( + !response.ok || + !response.headers.get("Content-Type")?.startsWith("multipart/") || + response.headers.get("cf-entrypoint") !== "index.js" || + !response.body + ) + fail("resource_conflict"); + const reader = response.body.getReader(); + const chunks: Uint8Array[] = []; + let size = 0; + for (;;) { + const { value, done } = await reader.read(); + if (done) break; + size += value.length; + if (size > limit) { + await reader.cancel(); + fail("resource_conflict"); + } + chunks.push(value); + } + const bytes = new Uint8Array(size); + let offset = 0; + for (const chunk of chunks) { + bytes.set(chunk, offset); + offset += chunk.length; + } + const form = await new Response(bytes, { + headers: { "Content-Type": response.headers.get("Content-Type")! }, + }).formData(); + if ([...form.keys()].length !== files.length) fail("resource_conflict"); + for (const file of files) { + const parts = form.getAll(file.path.slice("worker/".length)); + if (parts.length !== 1) fail("resource_conflict"); + const value = parts[0]; + const content = + typeof value === "string" + ? (new TextEncoder().encode(value).buffer as ArrayBuffer) + : await value.arrayBuffer(); + if ( + content.byteLength !== file.size || + (await digest(content)) !== file.sha256 + ) + fail("resource_conflict"); + } + } + async namespaces(installationConfig: unknown, artifact: Artifact) { + const deployments = await this.request(`${this.script}/deployments`); + const latest = deployments?.deployments?.[0]; + if ( + !Array.isArray(latest?.versions) || + latest.versions.length !== 1 || + latest.versions[0].percentage !== 100 || + !id(latest.versions[0].version_id) + ) + fail("resource_conflict"); + const version = await this.request( + `${this.script}/versions/${latest.versions[0].version_id}`, + ); + const runtime = version?.resources?.script_runtime; + if ( + !object(runtime) || + runtime.compatibility_date !== artifact.deployment.compatibility_date || + JSON.stringify(runtime.compatibility_flags) !== + JSON.stringify(artifact.deployment.compatibility_flags) || + !object(runtime.exports) + ) + fail("resource_conflict"); + for (const name of ["PersonalAgent", "Sandbox"]) { + const exported = runtime.exports[name]; + if ( + !object(exported) || + exported.type !== "durable-object" || + exported.storage !== "sqlite" || + (exported.state !== undefined && exported.state !== "created") || + (exported.container !== undefined && + (name !== "Sandbox" || + exported.container !== + this.record.resources.sandboxApplicationName)) + ) + fail("resource_conflict"); + } + if ( + Object.entries(runtime.exports).some( + ([name, value]) => + object(value) && + value.type === "durable-object" && + name !== "PersonalAgent" && + name !== "Sandbox", + ) + ) + fail("resource_conflict"); + const bindings = version?.resources?.bindings; + if (Array.isArray(bindings)) + this.verifyWorker({ bindings }, installationConfig); + if (!Array.isArray(bindings)) fail("temporarily_unavailable"); + const namespace = (name: string) => { + const matches = bindings.filter( + (b) => + b.type === "durable_object_namespace" && + b.name === name && + b.class_name === name && + (b.script_name === undefined || + b.script_name === this.record.resources.workerName), + ); + if (matches.length !== 1 || !id(matches[0].namespace_id)) + fail("resource_conflict"); + return matches[0].namespace_id as string; + }; + return { + personalAgentNamespaceId: namespace("PersonalAgent"), + sandboxNamespaceId: namespace("Sandbox"), + }; + } + private application(value: unknown): Application { + if ( + !object(value) || + !id(value.id) || + value.name !== this.record.resources.sandboxApplicationName || + !object(value.durable_objects) || + value.durable_objects.namespace_id !== + this.record.resources.sandboxNamespaceId || + !object(value.configuration) || + typeof value.configuration.image !== "string" || + false + ) + fail("resource_conflict"); + return { + ...value, + configuration: normalizedConfiguration(value.configuration), + } as unknown as Application; + } + async findApplication(): Promise { + const known = this.record.resources.sandboxApplicationId; + if (known) + return this.application( + await this.request(`containers/applications/${known}`), + ); + const list = await this.request( + `containers/applications?name=${this.record.resources.sandboxApplicationName}`, + ); + if (!Array.isArray(list) || list.length > 100) + fail("temporarily_unavailable"); + const matches = list.filter( + (app) => app?.name === this.record.resources.sandboxApplicationName, + ); + if (matches.length > 1) fail("resource_conflict"); + return matches.length ? this.application(matches[0]) : null; + } + desiredContainer(artifact: Artifact) { + const c = artifact.deployment.containers[0]; + return { + configuration: { image: c.image, instance_type: c.instance_type }, + max_instances: c.max_instances, + scheduling_policy: "default", + }; + } + matchesContainer(app: Application, artifact: Artifact) { + const d = this.desiredContainer(artifact); + return ( + app.configuration.image === d.configuration.image && + app.configuration.instance_type === d.configuration.instance_type && + app.max_instances === d.max_instances && + app.scheduling_policy === d.scheduling_policy + ); + } + async createApplication(artifact: Artifact) { + return this.application( + await this.request("containers/applications", "POST", { + name: this.record.resources.sandboxApplicationName, + instances: 0, + ...this.desiredContainer(artifact), + durable_objects: { + namespace_id: this.record.resources.sandboxNamespaceId, + }, + }), + ); + } + async patchApplication(app: Application, artifact: Artifact) { + await this.request(`containers/applications/${app.id}`, "PATCH", { + ...this.desiredContainer(artifact), + configuration: { + ...app.configuration, + ...this.desiredContainer(artifact).configuration, + }, + }); + } + async rollout( + app: Application, + artifact: Artifact, + operationId: string, + allowCreate: boolean, + ) { + const description = `Flarebot ${operationId}`; + const target = { + ...app.configuration, + ...this.desiredContainer(artifact).configuration, + }; + const list = await this.request( + `containers/applications/${app.id}/rollouts`, + ); + if (!Array.isArray(list) || list.length > 100) + fail("temporarily_unavailable"); + const matches = list.filter((r) => r?.description === description); + if (matches.length > 1) fail("resource_conflict"); + if (matches.length) { + if ( + !eq( + normalizedConfiguration(matches[0].target_configuration), + normalizedConfiguration(target), + ) || + !id(matches[0].id) + ) + fail("resource_conflict"); + if ( + matches[0].status === "pending" || + matches[0].status === "progressing" + ) + fail("temporarily_unavailable"); + if (matches[0].status !== "completed") fail("resource_conflict"); + return; + } + if (!allowCreate) fail("recovery_required"); + const created = await this.request( + `containers/applications/${app.id}/rollouts`, + "POST", + { + description, + strategy: "rolling", + kind: "full_auto", + step_percentage: 100, + target_configuration: target, + }, + ); + if ( + !id(created?.id) || + created?.status !== "completed" || + created?.description !== description || + !eq( + normalizedConfiguration(created?.target_configuration), + normalizedConfiguration(target), + ) + ) + fail("temporarily_unavailable"); + } + async publish() { + const current = await this.request(`${this.script}/subdomain`); + if (current?.enabled === true && current?.previews_enabled === false) + return; + await this.request(`${this.script}/subdomain`, "POST", { + enabled: true, + previews_enabled: false, + }); + } +} diff --git a/control-plane/deployment-config.ts b/control-plane/deployment-config.ts new file mode 100644 index 0000000..4381e78 --- /dev/null +++ b/control-plane/deployment-config.ts @@ -0,0 +1,27 @@ +import { loadControlPlaneConfig } from "../configuration/control-plane.ts"; +import type { Env } from "./config.ts"; +import type { Artifact } from "./artifact-types.ts"; +import type { Installation } from "./installation-metadata.ts"; +import { fail } from "./deployment-errors.ts"; +export function runtimeConfiguration( + env: Env, + record: Installation, + artifact: Artifact, + deployedOperationId: string, +) { + const cp = loadControlPlaneConfig(env).config; + if (!cp.bridge) fail("setup_required"); + return { + schemaVersion: 1, + installationId: record.installationId, + ownerSubject: record.ownerSubject, + runtimeOrigin: record.resources.runtimeOrigin, + controlPlaneOrigin: cp.publicOrigin, + bridge: { issuer: cp.publicOrigin, ...cp.bridge }, + release: { + version: artifact.identity.version, + artifactDigest: artifact.identity.artifactDigest, + operationId: deployedOperationId, + }, + }; +} diff --git a/control-plane/deployment-errors.ts b/control-plane/deployment-errors.ts new file mode 100644 index 0000000..8e4b629 --- /dev/null +++ b/control-plane/deployment-errors.ts @@ -0,0 +1,28 @@ +export const deploymentCodes = [ + "recovery_required", + "artifact_unavailable", + "setup_required", + "reauthorization_required", + "account_denied", + "resource_conflict", + "deployment_failed", + "health_failed", + "temporarily_unavailable", +] as const; +export type DeploymentCode = (typeof deploymentCodes)[number]; +export class DeploymentError extends Error { + readonly code: DeploymentCode; + constructor(code: DeploymentCode) { + super(code); + this.code = code; + } +} +export function deploymentCode(error: unknown): DeploymentCode { + if (error instanceof DeploymentError) return error.code; + // Never pass provider bodies, network URLs, credentials or arbitrary messages + // into Workflow error history or installation metadata. + return "temporarily_unavailable"; +} +export function fail(code: DeploymentCode): never { + throw new DeploymentError(code); +} diff --git a/control-plane/index.ts b/control-plane/index.ts index 6517ec4..e0487fb 100644 --- a/control-plane/index.ts +++ b/control-plane/index.ts @@ -3,6 +3,7 @@ import { loadControlPlaneConfig } from "../configuration/control-plane.ts"; import type { Env } from "./config.ts"; import { handleOAuth, privateResponse } from "./http.ts"; import { handleInstallations } from "./installations.ts"; +export { InstallationWorkflow } from "./installation-workflow.ts"; export { InstallationRegistry } from "./installation-registry.ts"; export { AuthVault } from "./vault.ts"; export default { diff --git a/control-plane/installation-metadata.ts b/control-plane/installation-metadata.ts index d02eb1e..06c748c 100644 --- a/control-plane/installation-metadata.ts +++ b/control-plane/installation-metadata.ts @@ -35,6 +35,9 @@ const status = z.enum([ ]); const errorCode = z .enum([ + "recovery_required", + "artifact_unavailable", + "setup_required", "reauthorization_required", "account_denied", "resource_conflict", diff --git a/control-plane/installation-registry.ts b/control-plane/installation-registry.ts index 475e872..7c8cfb0 100644 --- a/control-plane/installation-registry.ts +++ b/control-plane/installation-registry.ts @@ -10,16 +10,29 @@ import { registryResult, reject, sameRelease, + releaseIdentity, + type ReleaseIdentity, type Installation, type InstallationChanges, } from "./installation-metadata.ts"; +import { + operationSchema, + operationChanges, + type InstallationOperation, +} from "./operation.ts"; + const rowKey = (id: string) => `installation:${id}`; export const registryName = (subject: string) => `owner:${subject}`; const replaySchema = z.strictObject({ installationId, accountId: installationId, }); +const operationReplaySchema = z.strictObject({ + installationId, + operationId: installationId, + release: releaseIdentity, +}); export interface InstallationPage { installations: Installation[]; nextCursor: string | null; @@ -123,6 +136,234 @@ export class InstallationRegistry extends DurableObject { }; }); } + start( + subject: string, + id: string, + accountId: string, + requestId: string, + release: ReleaseIdentity, + deadline: number, + recovery: { + expectedRevision: number; + clearWorker: boolean; + clearContainer: boolean; + clearRollout: boolean; + } | null = null, + ) { + return registryResult(async () => { + this.#owner(subject); + parse(installationId, id); + parse(installationId, accountId); + parse(installationId, requestId); + const desired = parse(releaseIdentity, release); + if (recovery) + recovery = parse( + z.strictObject({ + expectedRevision: z.number().int().positive(), + clearWorker: z.boolean(), + clearContainer: z.boolean(), + clearRollout: z.boolean(), + }), + recovery, + ); + parse( + z + .number() + .int() + .min(Date.now() + 30_000) + .max(Date.now() + 30 * 60_000), + deadline, + ); + return this.ctx.storage.transaction(async (tx) => { + const current = this.#record(await tx.get(rowKey(id)), subject, id); + if (current.accountId !== accountId) reject("installation_conflict"); + const replayKey = `operation-request:${requestId}`; + const replayValue = await tx.get(replayKey); + const replay = + replayValue === undefined + ? null + : parse(operationReplaySchema, replayValue); + const previousValue = await tx.get(`operation:${id}`); + const previous = + previousValue === undefined + ? null + : parse(operationSchema, previousValue); + if (replay) { + if ( + replay.installationId !== id || + replay.operationId !== current.operationId || + !sameRelease(replay.release, desired) || + !previous + ) + reject("installation_conflict"); + return { + installation: current, + operation: parse(operationSchema, previous), + }; + } + if ( + current.status === "installing" || + current.status === "updating" || + current.installedRelease || + (current.desiredRelease && + !sameRelease(current.desiredRelease, desired)) + ) + reject("installation_conflict"); + if ( + recovery && + (current.status !== "failed" || + current.errorCode !== "recovery_required" || + current.revision !== recovery.expectedRevision) + ) + reject("installation_conflict"); + if ( + recovery?.clearWorker && + (current.resources.personalAgentNamespaceId || + current.resources.sandboxNamespaceId) + ) + reject("installation_conflict"); + if (recovery?.clearContainer && current.resources.sandboxApplicationId) + reject("installation_conflict"); + const operationId = crypto.randomUUID().replaceAll("-", ""); + const now = Math.max(Date.now(), current.updatedAt); + const operation = parse(operationSchema, { + operationId, + deadline, + workerIntent: recovery?.clearWorker + ? null + : (previous?.workerIntent ?? null), + containerIntent: recovery?.clearContainer + ? false + : (previous?.containerIntent ?? false), + containerUpdate: previous?.containerUpdate ?? false, + rolloutIntent: recovery?.clearRollout + ? null + : (previous?.rolloutIntent ?? null), + }); + const installation = parseInstallation({ + ...current, + desiredRelease: desired, + operationId, + status: "installing", + errorCode: null, + revision: current.revision + 1, + updatedAt: now, + }); + await tx.put({ + [rowKey(id)]: installation, + [`operation:${id}`]: operation, + [replayKey]: parse(operationReplaySchema, { + installationId: id, + operationId, + release: desired, + }), + }); + return { installation, operation }; + }); + }); + } + replay(subject: string, id: string, requestId: string) { + return registryResult(async () => { + this.#owner(subject); + parse(installationId, id); + parse(installationId, requestId); + return this.ctx.storage.transaction(async (tx) => { + const saved = await tx.get(`operation-request:${requestId}`); + if (saved === undefined) return null; + const replay = parse(operationReplaySchema, saved); + const current = this.#record(await tx.get(rowKey(id)), subject, id); + if ( + replay.installationId !== id || + replay.operationId !== current.operationId || + !sameRelease(replay.release, current.desiredRelease) + ) + reject("installation_conflict"); + return current; + }); + }); + } + operation(subject: string, id: string) { + return registryResult(async () => { + this.#owner(subject); + parse(installationId, id); + const value = await this.ctx.storage.get(`operation:${id}`); + return value === undefined ? null : parse(operationSchema, value); + }); + } + active(subject: string, id: string, operationId: string) { + return registryResult(async () => { + this.#owner(subject); + parse(installationId, id); + parse(installationId, operationId); + return this.ctx.storage.transaction(async (tx) => { + const installation = this.#record( + await tx.get(rowKey(id)), + subject, + id, + ); + const operation = parse( + operationSchema, + await tx.get(`operation:${id}`), + ); + if ( + installation.operationId !== operationId || + operation.operationId !== operationId || + installation.status !== "installing" + ) + reject("installation_conflict"); + return { installation, operation }; + }); + }); + } + intent( + subject: string, + id: string, + operationId: string, + expectedRevision: number, + changes: Partial< + Pick< + InstallationOperation, + "workerIntent" | "containerIntent" | "containerUpdate" | "rolloutIntent" + > + >, + ) { + return registryResult(async () => { + this.#owner(subject); + parse(installationId, id); + parse(installationId, operationId); + const input = parse(operationChanges, changes); + return this.ctx.storage.transaction(async (tx) => { + const installation = this.#record( + await tx.get(rowKey(id)), + subject, + id, + ); + const operation = parse( + operationSchema, + await tx.get(`operation:${id}`), + ); + if ( + installation.operationId !== operationId || + operation.operationId !== operationId || + installation.revision !== expectedRevision || + installation.status !== "installing" + ) + reject("installation_conflict"); + const next = parse(operationSchema, { ...operation, ...input }); + // Intent cannot be cleared or rewritten to conceal an ambiguous write. + for (const field of [ + "workerIntent", + "containerIntent", + "containerUpdate", + "rolloutIntent", + ] as const) + if (operation[field] && field in input) + reject("installation_conflict"); + await tx.put(`operation:${id}`, next); + return next; + }); + }); + } update( subject: string, id: string, diff --git a/control-plane/installation-workflow.ts b/control-plane/installation-workflow.ts new file mode 100644 index 0000000..e364e68 --- /dev/null +++ b/control-plane/installation-workflow.ts @@ -0,0 +1,358 @@ +import { + WorkflowEntrypoint, + type WorkflowEvent, + type WorkflowStep, +} from "cloudflare:workers"; +import { runtimeConfiguration } from "./deployment-config.ts"; +import type { Env } from "./config.ts"; +import { customerArtifact } from "./catalog.ts"; +import { DeploymentAPI, type DeploymentNetwork } from "./deployment-api.ts"; +import { deploymentCode, fail } from "./deployment-errors.ts"; +import { installationRegistry } from "./installations.ts"; +import { + unwrap, + sameRelease, + type Installation, + type InstallationChanges, + type ReleaseIdentity, +} from "./installation-metadata.ts"; +import { vault } from "./session.ts"; +import type { Artifact } from "./artifact-types.ts"; +import { healthInstallation } from "./bridge.ts"; + +export interface InstallationParams { + ownerSubject: string; + installationId: string; + operationId: string; +} +const retry = { + retries: { limit: 3, delay: "2 seconds", backoff: "exponential" as const }, + timeout: "2 minutes", +} as const; +export class InstallationWorkflow extends WorkflowEntrypoint< + Env, + InstallationParams +> { + // Only fixture subclasses replace these server seams. No configurable provider + // host, fixture flag, arbitrary artifact or health bypass exists in production. + protected artifact(pinned: ReleaseIdentity): Promise { + return customerArtifact(pinned); + } + protected network(): DeploymentNetwork { + return fetch; + } + protected health(record: Installation, deployedOperationId: string) { + return healthInstallation( + this.env, + record, + this.network(), + deployedOperationId, + ).catch(() => fail("health_failed")); + } + async run(event: WorkflowEvent, step: WorkflowStep) { + const params = event.payload; + const registry = installationRegistry(this.env, params.ownerSubject); + const binding = (record: Installation) => ({ + subject: record.ownerSubject, + accountId: record.accountId, + installationId: record.installationId, + operationId: params.operationId, + }); + const execute = async ( + name: string, + action: (context: { + record: Installation; + artifact: Artifact; + api: DeploymentAPI; + config: unknown; + bootstrapSecret: string | null; + operation: import("./operation.ts").InstallationOperation; + checkpoint: (changes: InstallationChanges) => Promise; + }) => Promise, + ) => { + const result = await step.do(name, retry, async () => { + try { + let { installation: record, operation } = unwrap( + await registry.active( + params.ownerSubject, + params.installationId, + params.operationId, + ), + ); + if (operation.deadline <= Date.now() + 60_000) + fail("reauthorization_required"); + const authorization = await vault( + this.env, + "operation", + params.operationId, + ).operation(binding(record)); + if (!authorization) fail("reauthorization_required"); + const grant = await vault( + this.env, + "grant", + authorization.grantRef, + ).grant(record.ownerSubject); + if (!grant || grant.expiresAt <= Date.now() + 60_000) + fail("reauthorization_required"); + const artifact = await this.artifact(record.desiredRelease!); + const config = runtimeConfiguration( + this.env, + record, + artifact, + operation.workerIntent?.operationId ?? params.operationId, + ); + await action({ + record, + operation, + artifact, + config, + bootstrapSecret: authorization.bootstrapSecret, + api: new DeploymentAPI(grant.accessToken, record, this.network()), + checkpoint: async (changes) => { + record = unwrap( + await registry.update( + record.ownerSubject, + record.installationId, + record.revision, + { ...changes, operationId: params.operationId }, + ), + ); + return record; + }, + }); + return { complete: true, errorCode: null }; + } catch (error) { + const code = deploymentCode(error); + if (code !== "temporarily_unavailable") + return { complete: false, errorCode: code }; + throw new Error(code); + } + }); + if (!result.complete) fail(result.errorCode!); + }; + try { + await execute( + "resolve customer origin", + async ({ record, api, checkpoint }) => { + const runtimeOrigin = await api.origin(); + if ( + record.resources.runtimeOrigin && + record.resources.runtimeOrigin !== runtimeOrigin + ) + fail("resource_conflict"); + if (!record.resources.runtimeOrigin) + await checkpoint({ resources: { runtimeOrigin } }); + }, + ); + await execute( + "upload assets and Worker", + async ({ + record, + operation, + artifact, + api, + config, + bootstrapSecret, + }) => { + const settings = await api.settings(); + if (settings) { + // The sole authorized upload intent must exist before adopting a Worker. + if (!operation.workerIntent) fail("resource_conflict"); + api.verifyWorker(settings, config); + await api.verifyContent(artifact); + return; + } + if (operation.workerIntent) fail("recovery_required"); + if (!bootstrapSecret) fail("reauthorization_required"); + await api.upload(artifact, config, bootstrapSecret, async () => { + unwrap( + await registry.intent( + record.ownerSubject, + record.installationId, + params.operationId, + record.revision, + { + workerIntent: { + operationId: params.operationId, + release: artifact.identity, + }, + }, + ), + ); + }); + }, + ); + await execute( + "reconcile native namespaces", + async ({ api, config, artifact, checkpoint }) => { + await api.verifyContent(artifact); + await checkpoint({ + resources: await api.namespaces(config, artifact), + }); + }, + ); + await execute( + "reconcile Containers application", + async ({ record, operation, artifact, api, checkpoint }) => { + let app = await api.findApplication(); + if (!app) { + if (operation.containerIntent) fail("recovery_required"); + unwrap( + await registry.intent( + record.ownerSubject, + record.installationId, + params.operationId, + record.revision, + { containerIntent: true }, + ), + ); + app = await api.createApplication(artifact); + } else if ( + !record.resources.sandboxApplicationId && + !operation.containerIntent + ) + fail("resource_conflict"); + if ( + !api.matchesContainer(app, artifact) || + operation.containerUpdate + ) { + // PATCH is declarative; its durable repair flag distinguishes a lost + // PATCH response from a rollout POST that was actually attempted. + if (!operation.containerUpdate) + unwrap( + await registry.intent( + record.ownerSubject, + record.installationId, + params.operationId, + record.revision, + { containerUpdate: true }, + ), + ); + if (!api.matchesContainer(app, artifact)) + await api.patchApplication(app, artifact); + const previous = operation.rolloutIntent; + if (!previous) + unwrap( + await registry.intent( + record.ownerSubject, + record.installationId, + params.operationId, + record.revision, + { rolloutIntent: params.operationId }, + ), + ); + await api.rollout( + app, + artifact, + previous ?? params.operationId, + !previous, + ); + } + if (!record.resources.sandboxApplicationId) + await checkpoint({ resources: { sandboxApplicationId: app.id } }); + }, + ); + await execute( + "publish and verify runtime", + async ({ record, api, operation }) => { + await api.publish(); + await this.health(record, operation.workerIntent!.operationId); + }, + ); + await step.do("record verified installation", retry, async () => { + try { + const current = unwrap( + await registry.get(params.ownerSubject, params.installationId), + ); + if ( + current?.operationId === params.operationId && + current.status === "ready" && + sameRelease(current.desiredRelease, current.installedRelease) + ) + return { complete: true }; + const active = unwrap( + await registry.active( + params.ownerSubject, + params.installationId, + params.operationId, + ), + ); + unwrap( + await registry.update( + params.ownerSubject, + params.installationId, + active.installation.revision, + { operationId: params.operationId, status: "ready" }, + ), + ); + return { complete: true }; + } catch { + throw new Error("temporarily_unavailable"); + } + }); + // A lost retirement response is harmless. The bounded protected record + // expires independently; no secret is copied into this Workflow's result. + await step.do("retire bootstrap material", retry, async () => { + try { + const record = unwrap( + await registry.get(params.ownerSubject, params.installationId), + ); + if ( + record?.status === "ready" && + record.operationId === params.operationId + ) + await vault( + this.env, + "operation", + params.operationId, + ).retireBootstrap(binding(record)); + return { complete: true }; + } catch { + return { complete: false }; + } + }); + return { status: "ready" }; + } catch (error) { + // Native Workflow errors carry our declared code only, never provider data. + const known = + typeof (error as Error)?.message === "string" + ? (error as Error).message + : ""; + const code = + ( + [ + "recovery_required", + "artifact_unavailable", + "setup_required", + "reauthorization_required", + "account_denied", + "resource_conflict", + "deployment_failed", + "health_failed", + ] as const + ).find((c) => known === c) ?? "temporarily_unavailable"; + await step.do("record safe failure", retry, async () => { + const state = await registry.active( + params.ownerSubject, + params.installationId, + params.operationId, + ); + if (state.ok) + unwrap( + await registry.update( + params.ownerSubject, + params.installationId, + state.value.installation.revision, + { + operationId: params.operationId, + status: "failed", + errorCode: code, + }, + ), + ); + return { errorCode: code }; + }); + return { status: "failed", errorCode: code }; + } + } +} diff --git a/control-plane/installations.ts b/control-plane/installations.ts index e19c16e..a851b81 100644 --- a/control-plane/installations.ts +++ b/control-plane/installations.ts @@ -13,6 +13,9 @@ import { import { registryName } from "./installation-registry.ts"; import { authenticatedPrincipal, selectedDeploymentGrant } from "./session.ts"; +import { startInstallation } from "./start-installation.ts"; +import { DeploymentError } from "./deployment-errors.ts"; + export const installationRegistry = (env: Env, subject: string) => env.INSTALLATIONS.get(env.INSTALLATIONS.idFromName(registryName(subject))); @@ -57,6 +60,34 @@ export async function handleInstallations( const origin = loadControlPlaneOrigin(env); if (url.origin !== origin) throw new OAuthError("forbidden"); if (request.method !== "GET") checkOrigin(request, origin); + const start = + /^\/api\/installations\/([a-f0-9]{32})\/(start|recover)$/.exec( + url.pathname, + ); + if (start && request.method === "POST") { + if (url.search) throw new InstallationError("invalid_metadata"); + const input = await form(request); + if ( + [...input.keys()].some((key) => key !== "requestId") || + input.getAll("requestId").length !== 1 + ) + throw new InstallationError("invalid_metadata"); + const requestId = parse(installationId, input.get("requestId")); + return json( + { + installation: await startInstallation( + request, + env, + start[1], + requestId, + network, + undefined, + start[2] === "recover", + ), + }, + 202, + ); + } if (url.pathname === "/api/installations" && request.method === "POST") { if (url.search) throw new InstallationError("invalid_metadata"); const input = await form(request); @@ -110,6 +141,15 @@ export async function handleInstallations( } return json({ error: "not_found" }, 404); } catch (error) { + if (error instanceof DeploymentError) + return json( + { error: error.code }, + error.code === "account_denied" + ? 403 + : error.code === "reauthorization_required" + ? 401 + : 503, + ); if (error instanceof InstallationError) return json( { error: error.code, message: metadataMessages[error.code] }, diff --git a/control-plane/operation.ts b/control-plane/operation.ts new file mode 100644 index 0000000..209ee19 --- /dev/null +++ b/control-plane/operation.ts @@ -0,0 +1,21 @@ +import { z } from "zod"; +import { installationId, releaseIdentity } from "./installation-metadata.ts"; +export const operationSchema = z.strictObject({ + operationId: installationId, + deadline: z.number().int().positive(), + // Intent is persisted BEFORE a mutation, and survives a new authorized attempt. + // An ambiguous absent resource must not cause a blind replay of that mutation. + workerIntent: z + .strictObject({ operationId: installationId, release: releaseIdentity }) + .nullable(), + containerIntent: z.boolean(), + containerUpdate: z.boolean(), + rolloutIntent: z + .string() + .regex(/^[a-f0-9]{32}$/) + .nullable(), +}); +export type InstallationOperation = z.infer; +export const operationChanges = operationSchema + .omit({ operationId: true, deadline: true }) + .partial(); diff --git a/control-plane/start-installation.ts b/control-plane/start-installation.ts new file mode 100644 index 0000000..55535eb --- /dev/null +++ b/control-plane/start-installation.ts @@ -0,0 +1,222 @@ +import type { Env } from "./config.ts"; +import type { CloudflareFetch } from "./cloudflare.ts"; +import { customerArtifact } from "./catalog.ts"; +import type { Artifact } from "./artifact-types.ts"; +import { loadControlPlaneConfig } from "../configuration/control-plane.ts"; +import { signBridgeAssertion } from "./bridge.ts"; +import { HEALTH_PURPOSE } from "../shared/bridge.ts"; +import { random } from "./crypto.ts"; +import { fail, DeploymentError } from "./deployment-errors.ts"; +import { DeploymentAPI } from "./deployment-api.ts"; +import { runtimeConfiguration } from "./deployment-config.ts"; +import { installationRegistry } from "./installations.ts"; +import { + unwrap, + InstallationError, + sameRelease, + type ReleaseIdentity, +} from "./installation-metadata.ts"; +import { selectedDeploymentGrant, vault } from "./session.ts"; + +// This request only starts/repairs native background execution. The browser +// never holds the deployment lease and never has to poll to keep it running. +export async function startInstallation( + request: Request, + env: Env, + id: string, + requestId: string, + network: CloudflareFetch = fetch, + artifactLoader: ( + pinned?: ReleaseIdentity, + ) => Promise = customerArtifact, + recover = false, +) { + const { principal, grant } = await selectedDeploymentGrant( + request, + env, + network, + ); + const registry = installationRegistry(env, principal.subject); + let current = unwrap(await registry.get(principal.subject, id)); + if (!current) throw new InstallationError("installation_not_found"); + if (current.accountId !== principal.selectedAccountId) fail("account_denied"); + if (recover && current.status === "ready") return current; + const cp = loadControlPlaneConfig(env).config; + if (!cp.bridge) fail("setup_required"); + const artifact = await artifactLoader(current.desiredRelease ?? undefined); + if ( + current.desiredRelease && + !sameRelease(current.desiredRelease, artifact.identity) + ) + fail("artifact_unavailable"); + try { + await signBridgeAssertion(env, { + aud: cp.publicOrigin, + sub: principal.subject, + installationId: id, + purpose: HEALTH_PURPOSE, + state: random(), + challenge: random(), + operationId: id, + artifactDigest: artifact.identity.artifactDigest, + version: artifact.identity.version, + }); + } catch { + fail("setup_required"); + } + if (grant.expiresAt <= Date.now() + 120_000) fail("reauthorization_required"); + const replay = unwrap( + await registry.replay(principal.subject, id, requestId), + ); + if (recover && !replay && current.status === "installing") { + const interruptedId = current.operationId!; + try { + const instance = await env.INSTALLATION_WORKFLOW.get(interruptedId); + await instance.terminate(); + if ((await instance.status()).status !== "terminated") + fail("temporarily_unavailable"); + } catch (error) { + const completed = unwrap(await registry.get(principal.subject, id)); + if ( + completed?.operationId === interruptedId && + completed.status === "ready" + ) + return completed; + const message = (error as Error)?.message; + // Native binding's documented absence signal; every other failure remains + // unavailable and cannot authorize recovery of an unknown execution. + if ( + message !== "instance.not_found" && + message !== "Error: instance.not_found" + ) + fail("temporarily_unavailable"); + } + // Termination cannot retract a remote write. Preserve all intent; only the + // native execution is stopped. A ready commit racing termination wins. + current = unwrap(await registry.get(principal.subject, id)); + if (!current || current.operationId !== interruptedId) + throw new InstallationError("installation_conflict"); + if (current.status === "ready") return current; + if (current.status === "installing") + current = unwrap( + await registry.update(principal.subject, id, current.revision, { + operationId: interruptedId, + status: "failed", + errorCode: "recovery_required", + }), + ); + } + const previousOperation = unwrap( + await registry.operation(principal.subject, id), + ); + const previousAuthorization = current.operationId + ? await vault(env, "operation", current.operationId).operation({ + subject: principal.subject, + accountId: current.accountId, + installationId: id, + operationId: current.operationId, + }) + : null; + let recovery: { + expectedRevision: number; + clearWorker: boolean; + clearContainer: boolean; + clearRollout: boolean; + } | null = null; + if (recover && current.status === "failed") { + if (current.errorCode !== "recovery_required" || !previousOperation) + fail("recovery_required"); + recovery = { + expectedRevision: current.revision, + clearWorker: false, + clearContainer: false, + clearRollout: false, + }; + const api = new DeploymentAPI(grant.accessToken, current, network); + const settings = await api.settings(); + if (settings) { + if (!previousOperation.workerIntent) fail("resource_conflict"); + api.verifyWorker( + settings, + runtimeConfiguration( + env, + current, + artifact, + previousOperation.workerIntent.operationId, + ), + ); + await api.verifyContent(artifact); + } else if (previousOperation.workerIntent) { + if ( + current.resources.personalAgentNamespaceId || + current.resources.sandboxNamespaceId + ) + fail("resource_conflict"); + recovery.clearWorker = true; + } + if (current.resources.sandboxNamespaceId) { + const app = await api.findApplication(); + if (!app && previousOperation.containerIntent) + recovery.clearContainer = true; + if (app && previousOperation.rolloutIntent) { + try { + await api.rollout( + app, + artifact, + previousOperation.rolloutIntent, + false, + ); + } catch (error) { + if ( + error instanceof DeploymentError && + error.code === "recovery_required" + ) + recovery.clearRollout = true; + else throw error; + } + } + } + } + const { installation, operation } = unwrap( + await registry.start( + principal.subject, + id, + principal.selectedAccountId!, + requestId, + artifact.identity, + Math.min(Date.now() + 20 * 60_000, grant.expiresAt - 60_000), + recovery, + ), + ); + if (installation.status !== "installing") return installation; + // Native services cannot share a transaction. The registry replay survives + // both missing vault creation and an ambiguous Workflow.create response. + await vault(env, "operation", operation.operationId).createOperation({ + subject: principal.subject, + accountId: installation.accountId, + installationId: id, + operationId: operation.operationId, + grantRef: principal.grantRef, + expiresAt: operation.deadline, + bootstrapSecret: previousAuthorization?.bootstrapSecret ?? random(), + }); + try { + await env.INSTALLATION_WORKFLOW.create({ + id: operation.operationId, + params: { + ownerSubject: principal.subject, + installationId: id, + operationId: operation.operationId, + }, + }); + } catch { + try { + await ( + await env.INSTALLATION_WORKFLOW.get(operation.operationId) + ).status(); + } catch { + fail("temporarily_unavailable"); + } + } + return installation; +} diff --git a/deployment/manifest.json b/deployment/manifest.json index 83b602b..4f92728 100644 --- a/deployment/manifest.json +++ b/deployment/manifest.json @@ -41,6 +41,14 @@ "required": true, "secret": true, "destination": "customer Worker FLAREBOT_SESSION_SECRET secret binding only; independently generated per installation, never manifest/vars/metadata" + }, + "bridgeVerification": { + "required": true, + "destination": "FLAREBOT_INSTALLATION.bridge; deliberate publisher Ed25519 public pin, never private signing key" + }, + "releaseIdentity": { + "required": true, + "destination": "FLAREBOT_INSTALLATION.release; trusted publisher catalog version/artifactDigest and original upload operation marker" } }, "storage": { @@ -80,21 +88,12 @@ "lifecycle": "declarative exports; never combine with migrations", "dataMigrations": "additive application migrations in customer storage, independent of class lifecycle", "destructiveChanges": "never automatic", - "keepBindings": [ - "secret_text", - "secret_key" - ] + "keepBindings": ["secret_text", "secret_key", "plain_text", "json"] }, "runtimeConfiguration": { "mode": "customer-runtime", - "variables": [ - "FLAREBOT_MODE", - "FLAREBOT_ENV", - "FLAREBOT_INSTALLATION" - ], - "secrets": [ - "FLAREBOT_SESSION_SECRET" - ], + "variables": ["FLAREBOT_MODE", "FLAREBOT_ENV", "FLAREBOT_INSTALLATION"], + "secrets": ["FLAREBOT_SESSION_SECRET"], "installationSchema": "configuration/customer.ts: InstallationConfig (schemaVersion 1)", "productionEnvironment": "production (default); development overrides forbidden", "developmentOnly": "FLAREBOT_DEV_OVERRIDES; origins only", diff --git a/docs/bug-lessons.md b/docs/bug-lessons.md index b6a3c5d..9b94762 100644 --- a/docs/bug-lessons.md +++ b/docs/bug-lessons.md @@ -318,3 +318,12 @@ Symptom-match new bug reports against these entries before theorising. logout against the real Worker and native Durable Objects. - **Prevention rule:** Test browser navigation forms as well as direct HTTP API calls. Never weaken Origin validation to accommodate a referrer-policy mistake. + +## 2026-09-06 — Native Workflow error messages are not a result protocol + +- **Affected area:** `control-plane/installation-workflow.ts` +- **Symptom signature:** Native `health_failed` and `resource_conflict` step failures became `temporarily_unavailable` in installation metadata, despite the correct callback error appearing in local logs. +- **Root cause:** The Workflow/RPC boundary changes nonretryable error messages; matching an exact application enum against the wrapped message loses the original classification. +- **Resolution:** Nonretryable callback failures return a strict safe result containing the enum. The Workflow interprets that persisted result outside the callback. Only transient failures throw a sanitized error for native retries. +- **Regression signal:** `pnpm test:orchestrator` checks exact durable failure codes through real local Workflow execution and confirms a failed boot never assigns an installed release. +- **Prevention rule:** Persist declared, nonsecret step results for domain failures. Do not depend on native exception message formatting as an application protocol. diff --git a/docs/installation-orchestrator.md b/docs/installation-orchestrator.md new file mode 100644 index 0000000..ee4f325 --- /dev/null +++ b/docs/installation-orchestrator.md @@ -0,0 +1,114 @@ +# Customer installation orchestration + +The publisher installs one bundled immutable customer release using a native +Cloudflare Workflow. The authenticated owner's `INSTALLATIONS` Durable Object +remains the sole installation authority. Browser requests can end immediately +after starting the operation; they do not execute or maintain the deployment. + +## Start and recovery + +After reserving an installation with `POST /api/installations`, send a +same-origin URL-encoded `POST /api/installations//start` containing only a +random 32-hex `requestId`. The server freshly verifies the selected account and +deployment grant, requires the immutable installation account to match, verifies +the publisher artifact and bridge signing key, and atomically records operation +intent and request replay. The response contains safe installation metadata. +Neither release URLs nor owner/account identifiers are accepted as deployment +inputs from the browser. FLA-10 owns the onboarding/status presentation; FLA-12 +owns selecting a newer release and upgrade policy. + +The Workflow ID is the server operation ID. A duplicate request repairs a missing +encrypted authorization record or missing Workflow creation without creating +another operation. Each callback rechecks the active operation, its deadline, +and the protected grant before loading credentials. Only owner/installation/ +operation identifiers enter Workflow parameters. Step results, errors and the +registry contain safe enums and resource identifiers, never credentials. + +A new authorized attempt after a failed installation retains the pinned release, +resource identities, prior upload marker and available encrypted bootstrap +secret. An existing owned Worker never receives a replacement session secret. +The last verified installed release remains distinct from the desired release; +only the final idempotent completion step assigns it after all checks pass. + +Worker upload, Containers creation and rollout intent are committed before their +respective mutations. Lost replies are reconciled through fixed native APIs, +exact Worker module hashes, configuration, namespace IDs, stable application +names and rollout operation descriptions. No failure automatically deletes +resources or adopts an unknown Worker/application. + +If a write outcome is ambiguous and no matching resource is observable, the +installation reports `recovery_required`. It does not blindly repeat the write. +An owner can deliberately send `POST /api/installations//recover` with a new +request ID after checking the customer account. This repeats fresh authorization, +rechecks the unchanged pinned artifact and actual remote absence, then clears +only the unresolved intent necessary for another attempt. It does not clear +recorded namespace/application identities. If bootstrap material has already +expired and the Worker remains absent, this explicit recovery creates a fresh +secret; any existing Worker must instead match the original recorded identity. +An explicit recover also handles an interrupted native Workflow: it terminates +the old execution, preserves its mutation intent, and starts a newly authorized +operation. A concurrent successful ready commit wins. Local Workflows currently +retain a running status after an abrupt process restart rather than automatically +resuming the interrupted callback; the acceptance gate exercises this concrete +owner recovery path. A missing Workflow at a persisted startup gap is recoverable +without retaining the original browser request ID. + +The provider offers no documented transaction spanning these resources or +create-only Worker CAS: a concurrent account administrator can still change +resources between observation and mutation. Conflicting observed identities +stop installation instead of being overwritten. + +## Native deployment sequence + +1. Resolve the account's existing workers.dev subdomain and fix the customer + origin to the reserved Worker name. No account-wide subdomain is renamed. +2. Start an assets upload session using native precomputed BLAKE3 asset hashes, + upload requested buckets with their MIME types, then upload actual Worker + modules with the completion JWT, declarative SQLite exports, app variables, + preserved customer variables/secrets and native Containers metadata. +3. Verify deployed content and bindings, resolve the PersonalAgent and Sandbox + namespace IDs from the active Worker version, and preserve those identities. +4. Reconcile the native Containers application against its Sandbox namespace. + An owned configuration repair uses PATCH followed by a separately reconciled + rolling rollout. Native expanded VM resources are normalized to the pinned + `lite` sizing. Unrelated configuration is retained during repair. +5. Enable the workers.dev endpoint and perform the separately signed metadata + health protocol. The customer validates identity, native parent readiness, + packaged assets, authentication and actual pinned Sandbox boot and destroy. + Health does not run AI inference or transfer conversations, files or model + credentials. See [owner login bridge](owner-login-bridge.md). +6. Atomically record ready and the verified installed release. Replaying this + step recognizes its own completed operation. Bootstrap retirement is bounded + best-effort cleanup backed by vault expiry and cannot undo readiness. + +## Publisher artifact and prerequisites + +Run `pnpm build:release`, then `pnpm build:control-plane` from committed source. +The latter generates a publisher-only catalog of opaque Data modules. Customer +JavaScript is never evaluated in the control-plane import graph, and these bytes +never enter the publisher's browser assets. Generated source maps are omitted +from the deployable inventory. Every inventory file has SHA-256 and size; the +trusted catalog pins the exact manifest SHA-256 plus version/source revision. +Asset API BLAKE3 addressing is separate from that inventory integrity. + +Publication rejects dirty customer artifacts, unknown/missing files, unsupported +bindings/exports, mismatched checksums and releases exceeding the bounded +publisher memory budget (16 MiB Worker modules, 24 MiB total inventory). Local +fixture work must explicitly use `pnpm build:control-plane:fixture`; production +artifact loading rejects a fixture catalog. Retain catalog artifacts needed by +active operations during control-plane rollouts. A missing pinned artifact fails +honestly instead of silently installing a newer build. + +Live installation requires a registered and verified Cloudflare OAuth client, +reviewed scope-catalog capability mapping covering Workers, Assets and Containers, +Workers Paid/Containers entitlement, an existing workers.dev subdomain, and a +configured bridge signing secret matching the deliberate public key pin. Local +fixtures use synthetic scopes and prove native execution and request contracts; +they do not certify live customer entitlement or third-party OAuth scope grants. +No live account deployment is performed by the local validation gates. + +The native endpoint contracts follow the installed Wrangler 4.128.0 deployment +implementation and [Workers multipart metadata](https://developers.cloudflare.com/workers/configuration/multipart-upload-metadata/), +[Static Assets direct upload](https://developers.cloudflare.com/workers/static-assets/direct-upload/), +[Worker content API](https://developers.cloudflare.com/api/resources/workers/subresources/scripts/subresources/content/methods/get/), +and the first-party [Containers application and rollout client](https://github.com/cloudflare/workers-sdk/tree/main/packages/containers-shared/src/client). diff --git a/docs/oauth-onboarding.md b/docs/oauth-onboarding.md index f5dad5b..2c104ce 100644 --- a/docs/oauth-onboarding.md +++ b/docs/oauth-onboarding.md @@ -152,7 +152,7 @@ An authorization error consumes its valid transaction too. All authorization, token, UserInfo, revoke and Accounts URLs are fixed to [Cloudflare's documented endpoints](https://developers.cloudflare.com/fundamentals/oauth/integrate-with-cloudflare/). -Credential-bearing requests use `redirect: error`, ten-second deadlines and +Credential-bearing requests use `redirect: manual` with explicit non-2xx rejection, ten-second deadlines and bounded response bodies. Token error categories are sanitized; no raw response, query, token, OAuth code, verifier or credential enters logs. The token response's scope is checked against the required manifest; when omitted, OAuth's diff --git a/package.json b/package.json index 1c4d98f..124a31e 100644 --- a/package.json +++ b/package.json @@ -30,9 +30,12 @@ "test:tasks-ui": "node --test tests/tasks-ui.test.mjs", "test:schedule-action": "node --test tests/schedule-action.test.mjs", "test:chat-ui": "node --test --test-reporter=tap tests/chat-ui.test.mjs", - "build:control-plane": "vite build --config control-plane/ui/vite.config.ts && wrangler deploy --config wrangler.control-plane.jsonc --dry-run --outdir dist/control-plane/worker", + "build:control-plane": "node scripts/build-catalog.mjs && vite build --config control-plane/ui/vite.config.ts && wrangler deploy --config wrangler.control-plane.jsonc --dry-run --outdir dist/control-plane/worker", "test:oauth": "node --test tests/oauth.test.mjs", - "test:ownership": "node --test --test-reporter=tap tests/ownership.test.mjs" + "test:ownership": "node --test --test-reporter=tap tests/ownership.test.mjs", + "build:control-plane:fixture": "node scripts/build-catalog.mjs --development-fixture && vite build --config control-plane/ui/vite.config.ts && wrangler deploy --config wrangler.control-plane.jsonc --dry-run --outdir dist/control-plane/worker", + "test:orchestrator": "node --test --test-reporter=tap tests/orchestrator.test.mjs", + "test:bridge": "node --test tests/bridge.test.mjs" }, "dependencies": { "@ai-sdk/anthropic": "4.0.49", @@ -44,13 +47,13 @@ "@octanejs/vite-plugin": "^0.1.52", "agents": "0.22.0", "ai": "7.0.93", + "cron-schedule": "6.0.0", "entities": "7.0.1", "marked": "18.0.11", "octane": "^0.2.2", "octane-kumo": "github:NathanBeddoeWebDev/octane-kumo#b86a46a4ed720eda9571d8940caa6f5ae6644fe3&path:/packages/octane-kumo", "workers-ai-provider": "4.0.0", - "zod": "4.4.3", - "cron-schedule": "6.0.0" + "zod": "4.4.3" }, "engines": { "node": "^24.12.0" @@ -60,6 +63,7 @@ "@tsrx/prettier-plugin": "^0.3.130", "@tsrx/typescript-plugin": "^0.3.130", "@types/node": "^26.4.1", + "blake3-wasm": "2.1.5", "jsonc-parser": "3.3.1", "playwright": "1.63.0", "prettier": "^3.9.6", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index e1c5f0e..c9b107f 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -173,6 +173,9 @@ importers: '@types/node': specifier: ^26.4.1 version: 26.4.1 + blake3-wasm: + specifier: 2.1.5 + version: 2.1.5 jsonc-parser: specifier: 3.3.1 version: 3.3.1 diff --git a/scripts/build-catalog.mjs b/scripts/build-catalog.mjs new file mode 100644 index 0000000..755cff0 --- /dev/null +++ b/scripts/build-catalog.mjs @@ -0,0 +1,57 @@ +import { readFile, writeFile, mkdir, rm } from "node:fs/promises"; +import { createHash } from "node:crypto"; + +// Publisher-controlled inputs only. This file is generated before CP packaging; +// customer JavaScript is imported as opaque Data, never executable CP modules. +const development = process.argv.includes("--development-fixture"); +const manifestBytes = await readFile("dist/release/manifest.json"); +const manifest = JSON.parse(manifestBytes); +if (manifest.sourceDirty && !development) + throw new Error( + "Refusing to publish a dirty customer artifact. Commit source and rebuild, or explicitly select --development-fixture for local tests.", + ); +if ( + manifest.files.reduce((sum, f) => sum + f.size, 0) > 24 * 1024 * 1024 || + manifest.files + .filter((f) => f.path.startsWith("worker/")) + .reduce((sum, f) => sum + f.size, 0) > + 16 * 1024 * 1024 +) + throw new Error("Release exceeds bounded publisher deployment memory budget"); +const artifactDigest = createHash("sha256").update(manifestBytes).digest("hex"); +if ( + artifactDigest !== + (await readFile("dist/release/manifest.sha256", "utf8")).trim() +) + throw new Error("Release manifest integrity mismatch"); +await rm("dist/catalog", { recursive: true, force: true }); +await mkdir("dist/catalog", { recursive: true }); +let source = + "// Generated immutable publisher catalog. Never shipped to browser assets.\n"; +const entries = []; +for (const [i, file] of [ + { path: "manifest.json", size: manifestBytes.length, sha256: artifactDigest }, + ...manifest.files, +].entries()) { + if ( + !/^(manifest\.json|deployment\.json|(?:worker|assets)\/[A-Za-z0-9_./-]+)$/.test( + file.path, + ) || + file.path.split("/").some((p) => p === ".." || !p) + ) + throw new Error("Unsafe release path"); + const bytes = await readFile(`dist/release/${file.path}`); + if ( + bytes.length !== file.size || + createHash("sha256").update(bytes).digest("hex") !== file.sha256 + ) + throw new Error("Release file integrity mismatch"); + await writeFile(`dist/catalog/${i}.bin`, bytes); + source += `import data${i} from "./${i}.bin";\n`; + entries.push(`${JSON.stringify(file.path)}: data${i}`); +} +source += `export default { identity: ${JSON.stringify({ version: manifest.release, sourceRevision: manifest.sourceRevision, artifactDigest })}, development: ${development}, files: {${entries.join(",")}} };\n`; +await writeFile("dist/catalog/catalog.ts", source); +console.log( + `Pinned publisher catalog: ${artifactDigest}${development ? " (development fixture)" : ""}`, +); diff --git a/scripts/build-release.mjs b/scripts/build-release.mjs index 6333f53..d436060 100644 --- a/scripts/build-release.mjs +++ b/scripts/build-release.mjs @@ -1,6 +1,7 @@ import { execFileSync } from "node:child_process"; import { cp, mkdir, readFile, readdir, rm, writeFile } from "node:fs/promises"; -import { join } from "node:path"; +import { join, extname } from "node:path"; +import { hash as blake3 } from "blake3-wasm"; import { fileURLToPath } from "node:url"; import { createHash } from "node:crypto"; import { parse } from "jsonc-parser"; @@ -30,6 +31,9 @@ execFileSync( }, ); await rm(join(output, "worker/README.md"), { force: true }); +// Debug maps remain in the normal build, outside the deployable inventory. +for (const file of await readdir(join(output, "worker"))) + if (file.endsWith(".map")) await rm(join(output, "worker", file)); await cp(join(root, "dist/client"), join(output, "assets"), { recursive: true, }); @@ -80,7 +84,32 @@ async function inventory(directory, prefix = "") { if (!entry.isFile()) throw new Error(`Unexpected artifact entry: ${path}`); const bytes = await readFile(join(directory, entry.name)); - files.push({ path, size: bytes.length, sha256: sha256(bytes) }); + const extension = extname(path).slice(1); + const mime = { + js: "application/javascript", + css: "text/css", + html: "text/html", + svg: "image/svg+xml", + png: "image/png", + ico: "image/x-icon", + woff2: "font/woff2", + txt: "text/plain", + json: "application/json", + }[extension]; + if (!mime) throw new Error(`Unsupported artifact MIME: ${path}`); + files.push({ + path, + size: bytes.length, + sha256: sha256(bytes), + mime, + ...(path.startsWith("assets/") + ? { + assetHash: blake3(bytes.toString("base64") + extension) + .toString("hex") + .slice(0, 32), + } + : {}), + }); } } return files; diff --git a/tests/fixtures/bridge-control-worker.ts b/tests/fixtures/bridge-control-worker.ts index 993c1b6..825b216 100644 --- a/tests/fixtures/bridge-control-worker.ts +++ b/tests/fixtures/bridge-control-worker.ts @@ -1,3 +1,4 @@ +export { InstallationWorkflow } from "../../control-plane/installation-workflow.ts"; // Explicit local test transport/admin seams. Never a production entry. import oauth from "./ownership-worker.ts"; import { InstallationRegistry } from "./ownership-worker.ts"; diff --git a/tests/fixtures/oauth-worker.ts b/tests/fixtures/oauth-worker.ts index 0aea56b..3faa13e 100644 --- a/tests/fixtures/oauth-worker.ts +++ b/tests/fixtures/oauth-worker.ts @@ -1,3 +1,4 @@ +export { InstallationWorkflow } from "../../control-plane/installation-workflow.ts"; // EXPLICIT TEST-ONLY network/admin seam. Never exported by a production entry. import app from "../../control-plane/index.ts"; import { AuthVault as ProductionVault } from "../../control-plane/vault.ts"; @@ -71,7 +72,7 @@ export const network: CloudflareFetch = async (input, init) => { ? input.href : input.url, ); - if (init?.redirect !== "error" || !init.signal) + if (init?.redirect !== "manual" || !init.signal) throw new Error("Missing redirect/timeout protection"); if (url.href === endpoints.token) { exchangeCount++; diff --git a/tests/fixtures/orchestrator-worker.ts b/tests/fixtures/orchestrator-worker.ts new file mode 100644 index 0000000..2307eb4 --- /dev/null +++ b/tests/fixtures/orchestrator-worker.ts @@ -0,0 +1,436 @@ +// Fixed stateful Cloudflare API fixture, native Workflow and native SQLite DOs. +// Never included in either production entrypoint. +import { DurableObject } from "cloudflare:workers"; +import oauth from "./ownership-worker.ts"; +import { network as oauthNetwork } from "./oauth-worker.ts"; +import { InstallationWorkflow as ProductionWorkflow } from "../../control-plane/installation-workflow.ts"; +import { customerArtifact } from "../../control-plane/catalog.ts"; +import { DeploymentError } from "../../control-plane/deployment-errors.ts"; +import { startInstallation } from "../../control-plane/start-installation.ts"; +import { installationRegistry } from "../../control-plane/installations.ts"; +import { + authenticatedPrincipal, + selectedDeploymentGrant, + vault, +} from "../../control-plane/session.ts"; +import type { Env } from "../../control-plane/config.ts"; +import type { + Installation, + ReleaseIdentity, +} from "../../control-plane/installation-metadata.ts"; +import { unwrap } from "../../control-plane/installation-metadata.ts"; +import { digest } from "../../control-plane/artifact.ts"; +import { InstallationRegistry as BaseRegistry } from "./ownership-worker.ts"; +import { AuthVault as BaseVault } from "./oauth-worker.ts"; +export class InstallationRegistry extends BaseRegistry { + async update(...args: Parameters) { + const result = await super.update(...args); + if ( + result.ok && + args[3].status === "ready" && + (await provider(this.env).fixtureFault("loseReady")) + ) + throw new Error("fixture lost ready reply"); + return result; + } +} +export class AuthVault extends BaseVault { + async retireBootstrap(...args: Parameters) { + const result = await super.retireBootstrap(...args); + if (await provider(this.env).fixtureFault("loseRetire")) + throw new Error("fixture lost retirement reply"); + return result; + } +} +interface TestEnv extends Env { + PROVIDER: DurableObjectNamespace; +} +const remoteId = async (value: string) => + (await digest(new TextEncoder().encode(value).buffer as ArrayBuffer)).slice( + 0, + 32, + ); +const safe = (result: unknown) => Response.json({ success: true, result }); +const unavailable = () => + Response.json( + { + success: false, + errors: [{ message: "provider-private-error-sentinel" }], + }, + { status: 503 }, + ); +const provider = (env: Env) => { + const e = env as TestEnv; + return e.PROVIDER.get(e.PROVIDER.idFromName("account")); +}; +const network = + (env: Env): typeof fetch => + async (input, init) => { + const request = new Request(input, init); + if ( + request.url.startsWith("https://api.cloudflare.com/client/v4/accounts/") + ) + return provider(env).fetch(request); + return oauthNetwork(input, init); + }; +export class Provider extends DurableObject { + async fixtureFault(name: string) { + const options = + (await this.ctx.storage.get>("options")) ?? {}; + if (!options[name] || (await this.ctx.storage.get(`fault:${name}`))) + return false; + await this.ctx.storage.put(`fault:${name}`, true); + return true; + } + async fixtureOptions(options: Record) { + await this.ctx.storage.put("options", options); + } + async fixtureInspect() { + return Object.fromEntries( + [...(await this.ctx.storage.list())].filter( + ([key]) => !key.startsWith("contents:"), + ), + ); + } + async fixtureHealth(record: Installation) { + const opts = await this.ctx.storage.get>("options"); + if (opts?.healthFailure) return false; + await this.ctx.storage.put(`health:${record.installationId}`, true); + return true; + } + async fetch(request: Request) { + const url = new URL(request.url); + const path = url.pathname.replace( + /^\/client\/v4\/accounts\/[a-f0-9]{32}\//, + "", + ); + const opts = + (await this.ctx.storage.get>("options")) ?? {}; + if ( + request.headers.get("Authorization") !== + "Bearer oauth-secret-sentinel-normal" && + !path.startsWith("workers/assets/upload") + ) + return new Response(null, { status: 401 }); + if (opts.revoked) + return new Response("provider-private-error-sentinel", { status: 401 }); + const trace = (await this.ctx.storage.get("trace")) ?? []; + trace.push(`${request.method} ${path}`); + await this.ctx.storage.put("trace", trace); + const once = async (name: string) => { + if (!opts[name]) return false; + const key = `fault:${name}`; + if (await this.ctx.storage.get(key)) return false; + await this.ctx.storage.put(key, true); + return true; + }; + if (path === "workers/subdomain") + return safe({ subdomain: "fixture-account" }); + const script = /^workers\/scripts\/([^/]+)(.*)$/.exec(path); + if (script) { + const [, name, suffix] = script; + const key = `worker:${name}`; + const worker = await this.ctx.storage.get(key); + if (suffix === "/settings") { + if (opts.collision && !worker) return safe({ bindings: [] }); + return worker + ? safe({ bindings: worker.bindings }) + : new Response(null, { status: 404 }); + } + if (suffix === "/assets-upload-session") { + const { manifest } = (await request.json()) as any; + await this.ctx.storage.put(`assets:${name}`, manifest); + const hashes = [ + ...new Set(Object.values(manifest).map((f: any) => f.hash)), + ]; + return safe({ + jwt: opts.emptyBuckets ? "completion.jwt" : "upload.jwt", + buckets: opts.emptyBuckets + ? [] + : opts.unknownAsset + ? [["f".repeat(32)]] + : [hashes], + }); + } + if (suffix === "" && request.method === "PUT") { + if (await once("rejectWorker")) return unavailable(); + const form = await request.formData(); + const metadata = JSON.parse( + await (form.get("metadata") as File).text(), + ); + const secret = metadata.bindings.find( + (b: any) => b.name === "FLAREBOT_SESSION_SECRET", + )?.text; + const modules = []; + const contents: Record = {}; + for (const [part, file] of form) { + if (part === "metadata") continue; + const f = file as File; + const bytes = await f.arrayBuffer(); + contents[part] = Math.ceil(bytes.byteLength / 1_048_576); + for (let offset = 0; offset < bytes.byteLength; offset += 1_048_576) + await this.ctx.storage.put( + `contents:${name}:${part}:${offset / 1_048_576}`, + bytes.slice(offset, offset + 1_048_576), + ); + modules.push({ + name: part, + mime: f.type, + size: f.size, + sha256: await digest(await f.arrayBuffer()), + }); + } + const bindings = await Promise.all( + metadata.bindings.map(async (b: any) => { + if (b.type === "secret_text") return { name: b.name, type: b.type }; + if (b.type === "durable_object_namespace") + return { ...b, namespace_id: await remoteId(name + b.name) }; + return b; + }), + ); + await this.ctx.storage.put(`contents:${name}`, contents); + await this.ctx.storage.put(key, { + bindings, + metadata: { ...metadata, bindings }, + modules, + secret, + versionId: "3".repeat(32), + }); + if (await once("holdWorker")) + await new Promise((resolve) => setTimeout(resolve, 90_000)); + if (await once("loseWorker")) return unavailable(); + return safe({ id: name, deployment_id: "3".repeat(32) }); + } + if (suffix === "/content/v2") { + const form = new FormData(); + const contents = await this.ctx.storage.get>( + `contents:${name}`, + ); + for (const [part, count] of Object.entries(contents ?? {})) { + const chunks: ArrayBuffer[] = []; + for (let i = 0; i < count; i++) + chunks.push( + (await this.ctx.storage.get( + `contents:${name}:${part}:${i}`, + ))!, + ); + form.append( + part, + new File(chunks, part, { type: "application/javascript+module" }), + ); + } + return new Response(form, { headers: { "cf-entrypoint": "index.js" } }); + } + if (suffix === "/deployments") + return safe({ + deployments: [ + { versions: [{ version_id: worker.versionId, percentage: 100 }] }, + ], + }); + if (suffix.startsWith("/versions/")) + return safe({ + resources: { + bindings: worker.bindings, + script_runtime: { + compatibility_date: worker.metadata.compatibility_date, + compatibility_flags: worker.metadata.compatibility_flags, + exports: worker.metadata.exports, + }, + }, + }); + if (suffix === "/subdomain") { + if (request.method === "POST") { + await this.ctx.storage.put(`published:${name}`, true); + return safe({ enabled: true, previews_enabled: false }); + } + return safe({ + enabled: !!(await this.ctx.storage.get(`published:${name}`)), + previews_enabled: false, + }); + } + } + if (path === "workers/assets/upload") { + if (await once("assetExpired")) return unavailable(); + const form = await request.formData(); + const files = []; + for (const [hash, value] of form) { + const f = value as File; + const bytes = Uint8Array.from(atob(await f.text()), (c) => + c.charCodeAt(0), + ); + files.push({ + hash, + mime: f.type, + size: bytes.length, + sha256: await digest(bytes.buffer), + }); + } + await this.ctx.storage.put("uploadedAssets", files); + return safe({ jwt: "completion.jwt" }); + } + if (path === "containers/applications") { + if (request.method === "GET") + return safe( + (await this.ctx.storage.list({ prefix: "application:" })) + .values() + .toArray() + .filter((app: any) => app.name === url.searchParams.get("name")), + ); + const body = (await request.json()) as any; + const app = { ...body, id: await remoteId(body.name) }; + if (opts.driftContainer) { + app.max_instances = 2; + app.configuration.environment_variables = [ + { name: "CUSTOM_CUSTOMER_SETTING", value: "retained-array-sentinel" }, + ]; + } + if (opts.expandedContainer) + app.configuration = { + ...app.configuration, + image: app.configuration.image, + vcpu: 0.0625, + memory_mib: 256, + disk: { size_mb: 2000 }, + }; + if (opts.expandedContainer) delete app.configuration.instance_type; + await this.ctx.storage.put(`application:${app.id}`, app); + if (await once("loseContainer")) return unavailable(); + return safe(app); + } + const application = /^containers\/applications\/([a-f0-9]{32})(.*)$/.exec( + path, + ); + if (application) { + const [, appId, suffix] = application; + const key = `application:${appId}`; + const app = await this.ctx.storage.get(key); + if (!suffix) { + if (request.method === "PATCH") { + const patch = await request.json(); + await this.ctx.storage.put(key, { ...app, ...(patch as object) }); + if (await once("losePatch")) return unavailable(); + return safe({ ...app, ...(patch as object) }); + } + return safe(app); + } + if (suffix === "/rollouts") { + if (request.method === "GET") + return safe((await this.ctx.storage.get(`rollouts:${appId}`)) ?? []); + const rollout = { + ...((await request.json()) as { target_configuration: any }), + id: "5".repeat(32), + status: "completed", + }; + if (opts.expandedContainer) + rollout.target_configuration = { + ...rollout.target_configuration, + image: rollout.target_configuration.image, + vcpu: 0.0625, + memory_mib: 256, + disk: { size_mb: 2000 }, + }; + if (opts.expandedContainer) + delete rollout.target_configuration.instance_type; + await this.ctx.storage.put(`rollouts:${appId}`, [rollout]); + if (await once("loseRollout")) return unavailable(); + return safe(rollout); + } + } + return new Response("Unexpected fixture endpoint", { status: 400 }); + } +} +export class InstallationWorkflow extends ProductionWorkflow { + protected artifact(pinned: ReleaseIdentity) { + return customerArtifact(pinned, true); + } + protected network() { + return network(this.env); + } + protected async health(record: Installation) { + if (!(await provider(this.env).fixtureHealth(record))) + throw new DeploymentError("health_failed"); + } +} +export default { + async fetch(request: Request, env: TestEnv, ctx: ExecutionContext) { + const url = new URL(request.url); + if (url.pathname === "/__test__/provider/options") { + await provider(env).fixtureOptions(await request.json()); + return Response.json({ ok: true }); + } + if (url.pathname === "/__test__/provider/inspect") + return Response.json(await provider(env).fixtureInspect()); + if (url.pathname === "/__test__/pending-start") { + const { id, requestId } = (await request.json()) as { + id: string; + requestId: string; + }; + const { principal, grant } = await selectedDeploymentGrant( + request, + env, + oauthNetwork, + ); + const artifact = await customerArtifact(undefined, true); + return Response.json( + unwrap( + await installationRegistry(env, principal.subject).start( + principal.subject, + id, + principal.selectedAccountId!, + requestId, + artifact.identity, + Math.min(Date.now() + 20 * 60_000, grant.expiresAt - 60_000), + ), + ), + ); + } + if (url.pathname === "/__test__/start") { + const { id, requestId, recover } = (await request.json()) as { + id: string; + requestId: string; + recover?: boolean; + }; + try { + return Response.json( + { + installation: await startInstallation( + request, + env, + id, + requestId, + network(env), + (pinned) => customerArtifact(pinned, true), + recover, + ), + }, + { status: 202 }, + ); + } catch (error) { + return Response.json( + { error: (error as any).code ?? "safe_fixture_error" }, + { status: 409 }, + ); + } + } + if (url.pathname === "/__test__/operation") { + const principal = await authenticatedPrincipal(request, env); + const record = unwrap( + await installationRegistry(env, principal.subject).get( + principal.subject, + url.searchParams.get("id")!, + ), + ); + const status = record?.operationId + ? await ( + await env.INSTALLATION_WORKFLOW.get(record.operationId) + ).status() + : null; + return Response.json({ record, status }); + } + return oauth.fetch( + request as Request, + env, + ctx, + ); + }, +} satisfies ExportedHandler; diff --git a/tests/fixtures/ownership-worker.ts b/tests/fixtures/ownership-worker.ts index 1644057..5af4b37 100644 --- a/tests/fixtures/ownership-worker.ts +++ b/tests/fixtures/ownership-worker.ts @@ -1,3 +1,4 @@ +export { InstallationWorkflow } from "../../control-plane/installation-workflow.ts"; // Explicit test-only storage/RPC inspection. Never part of a production entry. import oauth, { network } from "./oauth-worker.ts"; export { AuthVault } from "./oauth-worker.ts"; diff --git a/tests/orchestrator.test.mjs b/tests/orchestrator.test.mjs new file mode 100644 index 0000000..8aee3a5 --- /dev/null +++ b/tests/orchestrator.test.mjs @@ -0,0 +1,564 @@ +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]; +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 loadArtifact(entry, undefined, true); + assert.equal(loaded.identity.artifactDigest, entry.identity.artifactDigest); + await assert.rejects(() => loadArtifact(entry), { + message: "artifact_unavailable", + }); + await assert.rejects( + () => + loadArtifact( + entry, + { ...entry.identity, artifactDigest: "0".repeat(64) }, + true, + ), + { message: "artifact_unavailable" }, + ); + await assert.rejects( + () => + 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 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(origin), + FLAREBOT_CONTROL_PLANE: JSON.stringify({ + ...oauthConfig(origin), + 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: origin, ...options.headers }, + }); + const admin = async (path, body = {}, session = "") => + ( + await call("/__test__/" + path, { + method: "POST", + headers: { Cookie: session, "Content-Type": "application/json" }, + body: JSON.stringify(body), + }) + ).json(); + 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 session = await login(); + assert.ok(session); + assert.equal( + (await form("/api/account", session, { accountId: "a".repeat(32) })) + .status, + 200, + ); + const reserve = async () => { + const r = await form("/api/installations", session, { requestId: id() }); + assert.equal(r.status, 200); + return (await r.json()).installation; + }; + const inspect = () => admin("provider/inspect"); + const state = async (record) => + ( + await call(`/__test__/operation?id=${record.installationId}`, { + headers: { Cookie: session }, + }) + ).json(); + 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 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, + }); + 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(); + 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( + "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(); + await admin( + "pending-start", + { id: target.installationId, requestId: id() }, + session, + ); + const recovered = await startInstall(target, id(), true); + assert.ok(recovered.installation); + await wait(target); + }, + ); + 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(); + assert.equal( + (await form("/api/account", session, { accountId: "b".repeat(32) })) + .status, + 200, + ); + assert.equal((await startInstall(target)).error, "account_denied"); + assert.equal( + (await form("/api/account", session, { accountId: "a".repeat(32) })) + .status, + 200, + ); + 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(); + }, + ); + }, +); diff --git a/tsconfig.worker.json b/tsconfig.worker.json index 377e275..b97de27 100644 --- a/tsconfig.worker.json +++ b/tsconfig.worker.json @@ -17,6 +17,7 @@ "control-plane/*.ts", "tests/fixtures/oauth-worker.ts", "tests/fixtures/ownership-worker.ts", + "tests/fixtures/orchestrator-worker.ts", "tests/fixtures/bridge-control-worker.ts", "tests/fixtures/bridge-customer-worker.ts" ] diff --git a/wrangler.control-plane.jsonc b/wrangler.control-plane.jsonc index 1de0fca..aac2f62 100644 --- a/wrangler.control-plane.jsonc +++ b/wrangler.control-plane.jsonc @@ -5,6 +5,14 @@ "compatibility_date": "2026-09-04", "compatibility_flags": ["nodejs_compat"], "keep_names": true, + "rules": [{ "type": "Data", "globs": ["**/*.bin"], "fallthrough": true }], + "workflows": [ + { + "binding": "INSTALLATION_WORKFLOW", + "name": "flarebot-installation", + "class_name": "InstallationWorkflow", + }, + ], "keep_vars": true, "assets": { "directory": "./dist/control-plane/client", "binding": "ASSETS" }, "durable_objects": {