Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
19 kB · 545 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546import { 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, fingerprint, eq, 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 { verifyUpgradeIdentity } from "./artifact.ts";import type { Artifact } from "./artifact-types.ts";import { healthInstallation } from "./bridge.ts";import { runDomainWorkflow } from "./domain-workflow.ts";
export interface InstallationParams { kind?: "domain"; 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<Artifact> { 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<InstallationParams>, step: WorkflowStep) { const params = event.payload; if (params.kind === "domain") return runDomainWorkflow(this.env, params, step, this.network()); 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, progress: "preparing" | "deploying" | "provisioning" | "verifying", action: (context: { record: Installation; artifact: Artifact; api: DeploymentAPI; config: unknown; bootstrapSecret: string | null; operation: import("./operation.ts").InstallationOperation; checkpoint: (changes: InstallationChanges) => Promise<Installation>; }) => Promise<void>, ) => { 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"); if (record.progress !== progress) record = unwrap( await registry.update( record.ownerSubject, record.installationId, record.revision, { operationId: params.operationId, progress }, ), ); const artifact = await this.artifact(record.desiredRelease!); const api = new DeploymentAPI( grant.accessToken, record, this.network(), ); let config = name === "resolve customer origin" ? null : runtimeConfiguration( this.env, record, artifact, operation.workerIntent?.operationId ?? params.operationId, ); if (operation.upgrade && name !== "resolve customer origin") { const source = operation.upgrade.fromRelease; verifyUpgradeIdentity(source, artifact); const settings = await api.settings(); if (!settings) fail("resource_conflict"); const existing = api.configuration(settings); const expected = existing.release?.artifactDigest === source.artifactDigest ? operation.upgrade.configDigest : operation.workerIntent?.configDigest; if (!expected || (await fingerprint(existing)) !== expected) fail("resource_conflict"); config = { ...existing, release: { version: artifact.identity.version, artifactDigest: artifact.identity.artifactDigest, operationId: operation.workerIntent?.operationId ?? params.operationId, }, }; } await action({ record, operation, artifact, config, bootstrapSecret: authorization.bootstrapSecret, api, 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", "preparing", 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("reconcile model gateway", "preparing", async ({ api }) => { await api.ensureModelGateway(); }); await execute( "upload assets and Worker", "deploying", async ({ record, operation, artifact, api, config, bootstrapSecret, }) => { if (operation.upgrade) { const baseline = operation.upgrade; const settings = await api.settings(); if (!settings) fail("resource_conflict"); const existing = api.configuration(settings); if ( existing.release?.artifactDigest === artifact.identity.artifactDigest ) { if (!operation.workerIntent) fail("resource_conflict"); const adopted = await api.observe( artifact, operation.workerIntent.operationId, operation.workerIntent.configDigest, ); if (adopted.fingerprint !== baseline.fingerprint) fail("resource_conflict"); return; } const observed = await api.observe( { ...artifact, identity: baseline.fromRelease }, baseline.deployedOperationId, baseline.configDigest, baseline.sourceCodeHash, ); if ( !eq( { versionId: observed.versionId, deploymentId: observed.deploymentId, fingerprint: observed.fingerprint, containerFingerprint: observed.containerFingerprint, sourceCodeHash: observed.codeHash, }, { versionId: baseline.versionId, deploymentId: baseline.deploymentId, fingerprint: baseline.fingerprint, containerFingerprint: baseline.containerFingerprint, sourceCodeHash: baseline.sourceCodeHash, }, ) ) fail("resource_conflict"); if (operation.workerIntent) fail("recovery_required"); await api.upload( artifact, config, null, async () => { const checked = await api.observe( { ...artifact, identity: baseline.fromRelease }, baseline.deployedOperationId, baseline.configDigest, baseline.sourceCodeHash, ); if ( checked.fingerprint !== baseline.fingerprint || checked.containerFingerprint !== baseline.containerFingerprint || checked.deploymentId !== baseline.deploymentId || checked.versionId !== baseline.versionId ) fail("resource_conflict"); unwrap( await registry.intent( record.ownerSubject, record.installationId, params.operationId, record.revision, { workerIntent: { operationId: params.operationId, release: artifact.identity, configDigest: await fingerprint(config), }, }, ), ); if ( !eq(await api.activeDeployment(), { versionId: baseline.versionId, deploymentId: baseline.deploymentId, }) ) fail("resource_conflict"); }, { versionId: baseline.versionId, bindings: observed.settings.bindings, metadata: observed.metadata, }, ); return; } 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, configDigest: await fingerprint(config), }, }, ), ); }); }, ); await execute( "reconcile native namespaces", "provisioning", async ({ api, config, artifact, operation, checkpoint }) => { if (operation.upgrade) { if (!operation.workerIntent?.configDigest) fail("resource_conflict"); const observed = await api.observe( artifact, operation.workerIntent.operationId, operation.workerIntent.configDigest, ); if (observed.fingerprint !== operation.upgrade.fingerprint) fail("resource_conflict"); } await api.verifyContent(artifact); await checkpoint({ resources: await api.namespaces(config, artifact), }); }, ); await execute( "reconcile Containers application", "provisioning", async ({ record, operation, artifact, api, checkpoint }) => { let app = await api.findApplication(); if (!app && operation.upgrade) fail("resource_conflict"); 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 (operation.upgrade) { const observed = await api.containerFingerprint(app); if ( observed !== operation.upgrade.containerFingerprint && observed !== operation.upgrade.targetContainerFingerprint ) 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, operation.rolloutId, async (rolloutId) => { unwrap( await registry.intent( record.ownerSubject, record.installationId, params.operationId, record.revision, { rolloutId }, ), ); }, ); } if (!record.resources.sandboxApplicationId) await checkpoint({ resources: { sandboxApplicationId: app.id } }); }, ); await execute( "publish and verify runtime", "verifying", async ({ record, api, operation }) => { if (operation.upgrade) await api.endpoint(); else 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 }; } }}