import { WorkflowEntrypoint, type WorkflowEvent, type WorkflowStep, } from "cloudflare:workers"; import * as Effect from "effect/Effect"; import type { InstallationConfig } from "../configuration/customer.ts"; import type { Artifact } from "./artifact-types.ts"; import { healthInstallation } from "./bridge.ts"; import { customerArtifact } from "./catalog.ts"; import type { Env } from "./config.ts"; import { deploymentConfiguration } from "./deployment-config.ts"; import type { DeploymentFailure } from "./deployment-errors.ts"; import { deploymentFailure, HealthFailed, ReauthorizationRequired, RecoveryRequired, ResourceConflict, } from "./deployment-errors.ts"; import { fingerprint } from "./deployment-values.ts"; import { Deployment } from "./deployment.ts"; import { runDomainWorkflow } from "./domain-workflow.ts"; import type { RegistryFailure } from "./installation-errors.ts"; import { sameRelease, type Installation, type InstallationChanges, type ReleaseIdentity, } from "./installation-metadata.ts"; import { ensureModelGateway } from "./model-gateway.ts"; import type { InstallationOperation } from "./operation.ts"; import { installationRegistry } from "./registry-client.ts"; import { vault } from "./session.ts"; import { uploadWorker, verifyWorkerNamespaces } from "./worker-deployment.ts"; import { runDeploymentStep } from "./workflow-boundary.ts"; export interface InstallationParams { kind?: "domain"; ownerSubject: string; installationId: string; operationId: string; } type DeploymentProgress = "preparing" | "deploying" | "provisioning" | "verifying"; interface DeploymentAttempt { record: Installation; artifact: Artifact; api: Deployment; bootstrapSecret: string | null; operation: InstallationOperation; checkpoint: ( changes: InstallationChanges, ) => Effect.Effect; } 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) { return customerArtifact(pinned); } protected network(): typeof fetch { return fetch; } protected health( record: Installation, deployedOperationId: string, ): Effect.Effect { return healthInstallation( this.env, record, this.network(), deployedOperationId, ).pipe( Effect.mapError(() => new HealthFailed()), Effect.catchDefect(() => new HealthFailed()), ); } async run(event: WorkflowEvent, 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 runStep = async ( name: string, progress: DeploymentProgress, action: ( context: DeploymentAttempt, ) => Effect.Effect, ) => { const result = await step.do(name, retry, () => runDeploymentStep( Effect.gen({ self: this }, function* () { let { installation: record, operation } = yield* registry.active( params.ownerSubject, params.installationId, params.operationId, ); if (operation.deadline <= Date.now() + 60_000) return yield* new ReauthorizationRequired(); const authorization = yield* vault( this.env, "operation", params.operationId, ).operation(binding(record)); if (!authorization) return yield* new ReauthorizationRequired(); const grant = yield* vault( this.env, "grant", authorization.grantRef, ).grant(record.ownerSubject); if (!grant || grant.expiresAt <= Date.now() + 60_000) return yield* new ReauthorizationRequired(); if (record.progress !== progress) record = yield* registry.update( record.ownerSubject, record.installationId, record.revision, { operationId: params.operationId, progress }, ); const artifact = yield* this.artifact(record.desiredRelease!); const api = new Deployment( grant.accessToken, record, this.network(), ); yield* action({ record, operation, artifact, bootstrapSecret: authorization.bootstrapSecret, api, checkpoint: (changes) => registry .update( record.ownerSubject, record.installationId, record.revision, { ...changes, operationId: params.operationId }, ) .pipe( Effect.tap((updated) => Effect.sync(() => { record = updated; }), ), ), }); return { complete: true, errorCode: null }; }), ), ); if (!result.complete) throw deploymentFailure(result.errorCode); }; const execute = ( name: string, progress: DeploymentProgress, action: ( context: DeploymentAttempt & { config: InstallationConfig }, ) => Effect.Effect, ) => runStep(name, progress, (context) => deploymentConfiguration( this.env, context.record, context.artifact, context.operation, context.api.worker, ).pipe(Effect.flatMap((config) => action({ ...context, config }))), ); try { await runStep( "resolve customer origin", "preparing", ({ record, api, checkpoint }) => Effect.gen({ self: this }, function* () { const runtimeOrigin = yield* api.worker.origin(); if ( record.resources.runtimeOrigin && record.resources.runtimeOrigin !== runtimeOrigin ) return yield* new ResourceConflict(); if (!record.resources.runtimeOrigin) yield* checkpoint({ resources: { runtimeOrigin } }); }), ); await execute("reconcile model gateway", "preparing", ({ api }) => ensureModelGateway(api.account), ); await execute( "upload assets and Worker", "deploying", ({ record, operation, artifact, api, config, bootstrapSecret }) => Effect.gen({ self: this }, function* () { const recordIntent = Effect.gen(function* () { const configDigest = yield* fingerprint(config); yield* registry.intent( record.ownerSubject, record.installationId, params.operationId, record.revision, { workerIntent: { operationId: params.operationId, release: artifact.identity, configDigest, }, }, ); }); yield* uploadWorker(api, { artifact, config, operation, bootstrapSecret, recordIntent, }); }), ); await execute( "reconcile native namespaces", "provisioning", ({ api, config, artifact, operation, checkpoint }) => Effect.gen({ self: this }, function* () { const resources = yield* verifyWorkerNamespaces( api, artifact, config, operation, ); yield* checkpoint({ resources }); }), ); await execute( "reconcile Containers application", "provisioning", ({ record, operation, artifact, api, checkpoint }) => Effect.gen({ self: this }, function* () { let app = yield* api.containers.findApplication(); if (!app && operation.upgrade) return yield* new ResourceConflict(); if (!app) { if (operation.containerIntent) return yield* new RecoveryRequired(); yield* registry.intent( record.ownerSubject, record.installationId, params.operationId, record.revision, { containerIntent: true }, ); app = yield* api.containers.createApplication(artifact); } else if ( !record.resources.sandboxApplicationId && !operation.containerIntent ) return yield* new ResourceConflict(); if (operation.upgrade) { const observed = yield* api.containers.containerFingerprint(app); if ( observed !== operation.upgrade.containerFingerprint && observed !== operation.upgrade.targetContainerFingerprint ) return yield* new ResourceConflict(); } if ( !api.containers.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) yield* registry.intent( record.ownerSubject, record.installationId, params.operationId, record.revision, { containerUpdate: true }, ); if (!api.containers.matchesContainer(app, artifact)) yield* api.containers.patchApplication(app, artifact); const previous = operation.rolloutIntent; if (!previous) yield* registry.intent( record.ownerSubject, record.installationId, params.operationId, record.revision, { rolloutIntent: params.operationId }, ); yield* api.containers.rollout( app, artifact, previous ?? params.operationId, !previous, operation.rolloutId, (rolloutId) => registry .intent( record.ownerSubject, record.installationId, params.operationId, record.revision, { rolloutId }, ) .pipe(Effect.asVoid), ); } if (!record.resources.sandboxApplicationId) yield* checkpoint({ resources: { sandboxApplicationId: app.id }, }); }), ); await execute( "publish and verify runtime", "verifying", ({ record, api, operation }) => Effect.gen({ self: this }, function* () { if (operation.upgrade) yield* api.worker.endpoint(); else yield* api.worker.publish(); yield* this.health(record, operation.workerIntent!.operationId); }), ); await step.do("record verified installation", retry, () => Effect.runPromise( Effect.gen({ self: this }, function* () { const current = yield* registry.get( params.ownerSubject, params.installationId, ); if ( current?.operationId === params.operationId && current.status === "ready" && sameRelease(current.desiredRelease, current.installedRelease) ) return { complete: true }; const active = yield* registry.active( params.ownerSubject, params.installationId, params.operationId, ); yield* registry.update( params.ownerSubject, params.installationId, active.installation.revision, { operationId: params.operationId, status: "ready" }, ); return { complete: true }; }).pipe( Effect.catch(() => Effect.fail(new Error("temporarily_unavailable")), ), Effect.catchDefect(() => Effect.fail(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, () => Effect.runPromise( Effect.gen({ self: this }, function* () { const record = yield* registry.get( params.ownerSubject, params.installationId, ); if ( record?.status === "ready" && record.operationId === params.operationId ) yield* vault( this.env, "operation", params.operationId, ).retireBootstrap(binding(record)); return { complete: true }; }).pipe( Effect.catch(() => Effect.succeed({ complete: false })), Effect.catchDefect(() => Effect.succeed({ 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 () => { await Effect.runPromise( Effect.gen(function* () { // A stale or absent operation cannot be downgraded by this attempt. // Only this lookup suppresses expected registry failures, as before. const state = yield* registry .active( params.ownerSubject, params.installationId, params.operationId, ) .pipe(Effect.catch(() => Effect.succeed(null))); if (state) yield* registry.update( params.ownerSubject, params.installationId, state.installation.revision, { operationId: params.operationId, status: "failed", errorCode: code, }, ); }), ); return { errorCode: code }; }); return { status: "failed", errorCode: code }; } } }