import type { WorkflowStep } from "cloudflare:workers"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import { healthDomain } from "./bridge.ts"; import type { Env } from "./config.ts"; import { DomainAPI } from "./domain-api.ts"; import { domainErrorCode, domainFailure, DomainHttpsPending, DomainReauthorizationRequired, DomainResourceConflict, type DomainFailure, } from "./domain-errors.ts"; import { attachDomain, checkDomainHttps, verifyDomainMapping, type SaveDomain, } from "./domain-operations.ts"; import { configureRuntimeDomain, requireDomainCapability, } from "./domain-setup.ts"; import type { RegistryFailure } from "./installation-errors.ts"; import type { InstallationParams } from "./installation-workflow.ts"; import { installationRegistry } from "./registry-client.ts"; import { vault } from "./session.ts"; import { domainExitCode, runDomainStep } from "./workflow-boundary.ts"; export async function runDomainWorkflow( env: Env, params: InstallationParams, step: WorkflowStep, network: typeof fetch = fetch, ) { const registry = installationRegistry(env, params.ownerSubject); // This Effect is rerun inside every native step attempt. No token or mutable // registry snapshot survives a retry, checkpoint or durable sleep. const context = Effect.gen(function* () { const record = yield* registry.get( params.ownerSubject, params.installationId, ); const domain = yield* registry.getDomain( params.ownerSubject, params.installationId, ); if ( !record || !domain || domain.operationId !== params.operationId || !["connecting", "removing"].includes(domain.status) ) return yield* new DomainResourceConflict(); const authorization = yield* vault( env, "operation", params.operationId, ).operation({ subject: params.ownerSubject, accountId: record.accountId, installationId: params.installationId, operationId: params.operationId, }); if (!authorization || authorization.expiresAt <= Date.now() + 15_000) return yield* new DomainReauthorizationRequired(); const grant = yield* vault(env, "grant", authorization.grantRef).grant( params.ownerSubject, ); if (!grant) return yield* new DomainReauthorizationRequired(); yield* requireDomainCapability(env, grant.scopes); return { record, domain, api: new DomainAPI(grant.accessToken, record, network), }; }); const save: SaveDomain = (domain, changes) => registry.updateDomain( params.ownerSubject, params.installationId, params.operationId, domain.revision, { status: domain.status, errorCode: domain.errorCode, domainId: domain.domainId, writeIntent: domain.writeIntent, ...changes, }, ); async function execute( name: string, action: ( value: Effect.Success, ) => Effect.Effect, ) { const result = await step.do( name, { retries: { limit: 3, delay: "5 seconds", backoff: "exponential" }, timeout: "2 minutes", }, () => runDomainStep(context.pipe(Effect.flatMap(action))), ); if (result.error) throw domainFailure(result.error); } try { const initial = await Effect.runPromise( registry.getDomain(params.ownerSubject, params.installationId), ); if (!initial || initial.operationId !== params.operationId) return; if (["active", "removed"].includes(initial.status)) return; if (initial.action === "attach") { await execute("attach custom domain", ({ domain, api }) => attachDomain(api, domain, save), ); await execute("authorize domain in runtime", ({ record, domain }) => configureRuntimeDomain(env, record, domain, network), ); // TLS readiness is observed, not inferred from attachment or a DNS record. let ready = false; for (let attempt = 0; attempt < 30; attempt++) { const observation = await step.do( `check domain HTTPS ${attempt}`, { retries: { limit: 0, delay: "1 second" }, timeout: "30 seconds" }, async () => { const exit = await Effect.runPromiseExit( context.pipe( Effect.flatMap(({ record, domain, api }) => checkDomainHttps( api, domain, healthDomain(env, record, domain.origin!, network), ), ), ), ); return Exit.isSuccess(exit) ? { ready: exit.value, error: null } : { ready: false, error: domainExitCode(exit) }; }, ); if (observation.error) throw domainFailure(observation.error); ready = observation.ready; if (ready) break; await step.sleep(`wait for domain HTTPS ${attempt}`, "30 seconds"); } if (!ready) throw new DomainHttpsPending(); await execute("activate custom domain", ({ domain, api }) => verifyDomainMapping(api, domain).pipe( Effect.andThen(save(domain, { status: "active" })), Effect.asVoid, ), ); } else { await execute("revoke runtime domain", ({ record, domain }) => configureRuntimeDomain(env, record, domain, network), ); await execute("detach custom domain", ({ domain, api }) => api.remove(domain), ); await execute("finish domain removal", ({ domain }) => save(domain, { status: "removed" }).pipe(Effect.asVoid), ); } } catch (error) { const code = domainErrorCode(error); await step.do("record domain failure", () => Effect.runPromise( Effect.gen(function* () { const domain = yield* registry.getDomain( params.ownerSubject, params.installationId, ); if ( domain?.operationId === params.operationId && ["connecting", "removing"].includes(domain.status) ) yield* save(domain, { status: "failed", errorCode: code }); }), ), ); } await step.do("retire domain authorization", async () => { await Effect.runPromise( vault(env, "operation", params.operationId).destroy(), ); }); }