From 19234a0ed5264bd49dfc11d9120a7b6131e65d35 Mon Sep 17 00:00:00 2001 From: Tim Smart Date: Wed, 29 Jul 2026 16:45:09 +1200 Subject: [PATCH] Add Deno cluster HTTP support --- .changeset/eff-154-deno-cluster-http.md | 5 + packages/platform-deno/src/DenoClusterHttp.ts | 153 ++++++++++++++++++ packages/platform-deno/src/index.ts | 5 + 3 files changed, 163 insertions(+) create mode 100644 .changeset/eff-154-deno-cluster-http.md create mode 100644 packages/platform-deno/src/DenoClusterHttp.ts diff --git a/.changeset/eff-154-deno-cluster-http.md b/.changeset/eff-154-deno-cluster-http.md new file mode 100644 index 00000000000..36d03289586 --- /dev/null +++ b/.changeset/eff-154-deno-cluster-http.md @@ -0,0 +1,5 @@ +--- +"@effect/platform-deno": patch +--- + +Add native Deno HTTP and WebSocket layers for Effect Cluster runners. diff --git a/packages/platform-deno/src/DenoClusterHttp.ts b/packages/platform-deno/src/DenoClusterHttp.ts new file mode 100644 index 00000000000..d72e78eed0d --- /dev/null +++ b/packages/platform-deno/src/DenoClusterHttp.ts @@ -0,0 +1,153 @@ +/** + * Native Deno HTTP and WebSocket layers for Effect Cluster runners. + * + * `layerHttpServer` provides the Deno HTTP server used by cluster runners. The + * main `layer` builds a sharding layer for HTTP or WebSocket transport, + * choosing serialization, runner health checks, runner storage, message + * storage, and optional client-only mode from the supplied options. + * + * @since 4.0.0 + */ +import type * as Config from "effect/Config" +import * as Effect from "effect/Effect" +import * as Layer from "effect/Layer" +import * as Option from "effect/Option" +import * as HttpRunner from "effect/unstable/cluster/HttpRunner" +import * as MessageStorage from "effect/unstable/cluster/MessageStorage" +import * as RunnerHealth from "effect/unstable/cluster/RunnerHealth" +import * as Runners from "effect/unstable/cluster/Runners" +import * as RunnerStorage from "effect/unstable/cluster/RunnerStorage" +import type { Sharding } from "effect/unstable/cluster/Sharding" +import * as ShardingConfig from "effect/unstable/cluster/ShardingConfig" +import * as SqlMessageStorage from "effect/unstable/cluster/SqlMessageStorage" +import * as SqlRunnerStorage from "effect/unstable/cluster/SqlRunnerStorage" +import type * as Etag from "effect/unstable/http/Etag" +import type { HttpPlatform } from "effect/unstable/http/HttpPlatform" +import type { HttpServer } from "effect/unstable/http/HttpServer" +import type { ServeError } from "effect/unstable/http/HttpServerError" +import * as RpcSerialization from "effect/unstable/rpc/RpcSerialization" +import type { SqlClient } from "effect/unstable/sql/SqlClient" +import { layerK8sHttpClient } from "./DenoClusterSocket.ts" +import * as DenoCrypto from "./DenoCrypto.ts" +import * as DenoHttpClient from "./DenoHttpClient.ts" +import * as DenoHttpServer from "./DenoHttpServer.ts" +import type { DenoServices } from "./DenoServices.ts" +import * as DenoSocket from "./DenoSocket.ts" + +export { + /** + * Layer that provides a Kubernetes HTTP client for runner health checks. + * + * @category re-exports + * @since 4.0.0 + */ + layerK8sHttpClient +} + +/** + * Layer that provides a native Deno HTTP server for cluster runners. + * + * @category layers + * @since 4.0.0 + */ +export const layerHttpServer: Layer.Layer< + | HttpPlatform + | Etag.Generator + | DenoServices + | HttpServer, + ServeError, + ShardingConfig.ShardingConfig +> = Effect.gen(function*() { + const config = yield* ShardingConfig.ShardingConfig + const listenAddress = Option.orElse(config.runnerListenAddress, () => config.runnerAddress) + if (Option.isNone(listenAddress)) { + return yield* Effect.die("DenoClusterHttp.layerHttpServer: ShardingConfig.runnerAddress is None") + } + return DenoHttpServer.layer({ + hostname: listenAddress.value.host, + port: listenAddress.value.port, + onListen: () => {} + }) +}).pipe(Layer.unwrap) + +/** + * Creates Deno cluster layers for HTTP or WebSocket transport, configuring + * serialization, storage, runner health, and optional client-only mode. + * + * @category layers + * @since 4.0.0 + */ +export const layer = < + const ClientOnly extends boolean = false, + const Storage extends "local" | "sql" | "byo" = never +>(options: { + readonly transport: "http" | "websocket" + readonly serialization?: "msgpack" | "ndjson" | undefined + readonly clientOnly?: ClientOnly | undefined + readonly storage?: Storage | undefined + readonly runnerHealth?: "ping" | "k8s" | undefined + readonly runnerHealthK8s?: { + readonly namespace?: string | undefined + readonly labelSelector?: string | undefined + } | undefined + readonly shardingConfig?: Partial | undefined +}): ClientOnly extends true ? Layer.Layer< + Sharding | Runners.Runners | ("byo" extends Storage ? never : MessageStorage.MessageStorage), + Config.ConfigError, + "local" extends Storage ? never + : "byo" extends Storage ? (MessageStorage.MessageStorage | RunnerStorage.RunnerStorage) + : SqlClient + > : + Layer.Layer< + Sharding | Runners.Runners | ("byo" extends Storage ? never : MessageStorage.MessageStorage), + ServeError | Config.ConfigError, + "local" extends Storage ? never + : "byo" extends Storage ? (MessageStorage.MessageStorage | RunnerStorage.RunnerStorage) + : SqlClient + > => +{ + const layer: Layer.Layer = options.clientOnly + ? options.transport === "http" + ? Layer.provide(HttpRunner.layerHttpClientOnly, DenoHttpClient.layer) + : Layer.provide(HttpRunner.layerWebsocketClientOnly, DenoSocket.layerWebSocketConstructor) + : options.transport === "http" + ? Layer.provide(HttpRunner.layerHttp, [layerHttpServer, DenoHttpClient.layer]) + : Layer.provide(HttpRunner.layerWebsocket, [layerHttpServer, DenoSocket.layerWebSocketConstructor]) + + const runnerHealth: Layer.Layer = options.clientOnly + ? Layer.empty as any + : options.runnerHealth === "k8s" + ? RunnerHealth.layerK8s(options.runnerHealthK8s).pipe( + Layer.provide(layerK8sHttpClient) + ) + : RunnerHealth.layerPing.pipe( + Layer.provide(Runners.layerRpc), + Layer.provide( + options.transport === "http" + ? HttpRunner.layerClientProtocolHttpDefault.pipe(Layer.provide(DenoHttpClient.layer)) + : HttpRunner.layerClientProtocolWebsocketDefault.pipe(Layer.provide(DenoSocket.layerWebSocketConstructor)) + ) + ) + + return layer.pipe( + Layer.provide(runnerHealth), + Layer.provideMerge( + options.storage === "local" + ? MessageStorage.layerNoop + : options.storage === "byo" + ? Layer.empty + : Layer.orDie(SqlMessageStorage.layer).pipe(Layer.provide(DenoCrypto.layer)) + ), + Layer.provide( + options.storage === "local" + ? RunnerStorage.layerMemory + : options.storage === "byo" + ? Layer.empty + : Layer.orDie(SqlRunnerStorage.layer) + ), + Layer.provide(ShardingConfig.layerFromEnv(options.shardingConfig)), + Layer.provide( + options.serialization === "ndjson" ? RpcSerialization.layerNdjson : RpcSerialization.layerMsgPack + ) + ) as any +} diff --git a/packages/platform-deno/src/index.ts b/packages/platform-deno/src/index.ts index 68be2f9ddc7..11134127627 100644 --- a/packages/platform-deno/src/index.ts +++ b/packages/platform-deno/src/index.ts @@ -9,6 +9,11 @@ */ export * as DenoChildProcessSpawner from "./DenoChildProcessSpawner.ts" +/** + * @since 4.0.0 + */ +export * as DenoClusterHttp from "./DenoClusterHttp.ts" + /** * @since 4.0.0 */