import {
Credentials,
oauthCredentials,
Retry,
type CloudflareOpContext,
} from "@distilled.cloud/cloudflare";
import * as Effect from "effect/Effect";
import * as FetchHttpClient from "effect/unstable/http/FetchHttpClient";
import { bodyBytes, withResponse } from "../server/http.ts";
// No process credentials, refresh cache or nested retries. The calling request
// or native Workflow attempt has already resolved its own authorized grant.
export function cloudflare(
operation: Effect.Effect,
accessToken: string,
network: typeof fetch,
timeout = 60_000,
limit = 2 * 1024 * 1024,
) {
const boundedFetch: typeof fetch = (input, init) =>
Effect.runPromise(
withResponse(
network,
input,
init ?? {},
timeout,
new Error("Cloudflare transport unavailable"),
(response) =>
Effect.gen(function* () {
if (response.status >= 300 && response.status < 400)
return yield* Effect.die(
new Error("Cloudflare redirect rejected"),
);
if (!response.body) return response;
const bytes = yield* bodyBytes(
response,
limit,
new Error("Cloudflare response unavailable"),
);
// The SDK also accepts bare JSON for other services. These control-plane
// operations require Cloudflare's successful v4 envelope explicitly.
if (response.ok) {
yield* Effect.try({
try: () => {
const envelope: unknown = JSON.parse(
new TextDecoder().decode(bytes),
);
if (
!envelope ||
typeof envelope !== "object" ||
!("success" in envelope) ||
envelope.success !== true ||
!("result" in envelope)
)
throw new Error("Invalid Cloudflare envelope");
},
catch: () => new Error("Cloudflare response unavailable"),
});
}
return new Response(bytes, response);
}),
),
{ signal: init?.signal ?? undefined },
);
return operation.pipe(
Retry.none,
Effect.provideService(
Credentials,
Effect.succeed(oauthCredentials({ accessToken })),
),
Effect.provide(FetchHttpClient.layer),
Effect.provideService(FetchHttpClient.Fetch, boundedFetch),
Effect.provideService(FetchHttpClient.RequestInit, { redirect: "manual" }),
);
}