Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
6.5 kB · 191 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192import 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<typeof context>, ) => Effect.Effect<void, DomainFailure | RegistryFailure>, ) { 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(), ); });}