effect 4 新特性:为 Effect Cluster 提供 Deno 原生 Socket Runner 层
【免费下载链接】effectBuild production-ready applications in TypeScript项目地址: https://gitcode.com/GitHub_Trending/ef/effect
本篇文章介绍effect仓库中一项新近加入的特性:为 Effect Cluster 的 Runner(运行器)提供基于 Deno 原生 TCP socket 的传输层实现(@effect/platform-deno包中的DenoClusterSocket模块)。读者将掌握如何在 Deno 运行时下用原生 socket 搭建分布式分片(Sharding)集群、配置 runner 健康检查与消息/运行器存储,以及理解该实现与 Node 平台在连接超时、TLS 升级等细节上的差异。
特性来源与定位
该特性由仓库.changeset/pre/eff-155-deno-cluster-socket.md记录:
--- "@effect/platform-deno": patch --- Add native Deno socket layers for Effect Cluster runners.它属于@effect/platform-deno的一个 patch 级更新,为 Effect Cluster 的 runner 增加了原生 Deno socket 层。也就是说,在 Deno 运行时中,Cluster 的 runner 不再需要依赖 HTTP/WebSocket 传输,而是可以直接通过 TCP socket 建立点对点连接来接收与执行分片实体(Entity)的请求。相关代码位于 packages/platform/deno/src/DenoClusterSocket.ts,并在 packages/platform/deno/src/index.ts 中以DenoClusterSocket命名空间对外导出。
Effect Cluster 与 SocketRunner 背景
Effect Cluster 是 effect 的分片式分布式计算模型:一组 runner 进程各自承载若干分片(shard),实体(Entity)通过 RPC 调用被路由到持有对应分片的 runner 上执行。SocketRunner(effect/unstable/cluster/SocketRunner)是 Cluster 中基于原始 socket 的 runner 实现,与HttpRunner(基于 HTTP/WebSocket)并列。本次新增的 Deno socket 层,正是把SocketRunner所依赖的各种服务(RPC 客户端协议、socket 服务器、健康检查、存储等)用 Deno 的原生 API 组装起来,让 Deno 应用可以直接以 socket 传输方式加入 Cluster。
核心 API:DenoClusterSocket.layer
DenoClusterSocket.layer是一个一键式的分片(sharding)层构造器,它会根据传入的选项自动装配传输、序列化、存储、健康检查等全部依赖。函数签名与选项如下(来自 DenoClusterSocket.ts):
export const layer = < const ClientOnly extends boolean = false, const Storage extends "local" | "sql" | "byo" = never >( options?: { readonly serialization?: "binary" | "ndjson" | undefined readonly serializationMaxBufferSize?: number | "unbounded" | 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<ShardingConfig.ShardingConfig["Service"]> | undefined } ): Layer.Layer<...>选项说明
| 选项 | 取值 | 作用 |
|---|---|---|
serialization | "binary"/"ndjson" | RPC 消息序列化方式。缺省为"binary"(Schema 二进制编码);传"ndjson"则使用 NDJSON 行式编码 |
serializationMaxBufferSize | 数字 /"unbounded" | 序列化缓冲上限。二进制模式下作为maxFrameSize,NDJSON 模式下作为maxBufferSize |
clientOnly | true/false | 仅构建客户端(不启动 socket 服务器),用于只发起请求、不承载实体的进程 |
storage | "local"/"sql"/"byo" | 存储策略。"local"使用内存存储;"sql"使用 SQL 存储(需要SqlClient,并会自动装配DenoCrypto加密层);"byo"(bring your own)要求自行提供MessageStorage与RunnerStorage |
runnerHealth | "ping"/"k8s" | runner 健康检查方式。"ping"通过 RPC ping 检测存活;"k8s"通过 Kubernetes API 检测 |
runnerHealthK8s | { namespace?, labelSelector? } | "k8s"健康检查时的命名空间与标签选择器 |
shardingConfig | Partial<ShardingConfig> | 需要覆盖的分片配置项,最终会与ShardingConfig.layerFromEnv的环境变量配置合并 |
返回类型与依赖
layer的返回类型会随选项精确变化(从源码的类型重载可以清楚看到):
- 非
clientOnly时,产出Sharding | Runners | MessageStorage,错误通道为SocketServer.SocketServerError | Config.ConfigError; clientOnly: true时,不产出MessageStorage,错误通道退化为Config.ConfigError;storage为"local"时不再依赖外部存储;"sql"时需要SqlClient;"byo"时需要自行提供MessageStorage | RunnerStorage。
从实现看,layer内部实际是把SocketRunner.layer(或SocketRunner.layerClientOnly)与layerClientProtocol、layerSocketServer、健康检查层、存储层、序列化层逐层provide组合起来的组合子,最终返回一个开箱即用的分片层。
组合层的内部实现
layerClientProtocol:TCP 上的 RPC 客户端协议
layerClientProtocol提供 Cluster 所需的RpcClientProtocol,其核心是使用DenoSocket.makeTcp打开到目标 runner 地址(address.host/address.port)的 TCP 连接,并设置openTimeout: 1000(1 秒打开超时),随后通过RpcClient.makeProtocolSocket()把连接包装成 RPC 协议 socket(见 DenoClusterSocket.ts)。
值得注意的是源码中的注释:与 Node socket 不同,Deno 连接没有原生的 idle-timeout 选项,因此 peer 连接只使用这 1 秒的 open timeout,而不存在连接空闲自动断开的机制。这是迁移到 Deno 平台时需要考虑的行为差异。
layerSocketServer:runner 的 socket 服务器
layerSocketServer为 runner 提供SocketServer,监听地址取自ShardingConfig.runnerListenAddress,若该值为空则回退到runnerAddress;两者皆空时直接以Effect.die终止(见 DenoClusterSocket.ts)。它底层委托给DenoSocketServer.layer({ host, port })。
layerK8sHttpClient:Kubernetes 健康检查的 HTTP 客户端
当runnerHealth: "k8s"时,健康检查需要访问 Kubernetes API。layerK8sHttpClient提供了一个基于 Deno 原生Deno.createHttpClient的作用域化 HTTP 客户端:它会尝试读取 ServiceAccount 的 CA 证书(/var/run/secrets/kubernetes.io/serviceaccount/ca.crt),读取成功则用该 CA 创建专用客户端,否则回退到globalThis.fetch(见 DenoClusterSocket.ts)。这保证了在 Kubernetes 集群内与本地开发两种环境下都能工作。
底层 socket 适配:DenoSocket模块
DenoClusterSocket依赖的DenoSocket模块(packages/platform/deno/src/DenoSocket.ts)是 Deno 原生连接与 EffectSocket抽象之间的桥梁,它本身也是独立的可复用 API:
makeTcp(options):打开原生 Deno TCP/Unix 连接并返回Socket.Socket。支持noDelay、keepAlive、openTimeout选项;openTimeout会中断获取流程,但无法取消已经发出的Deno.connectpromise——超时后到达的连接会由作用域终结器(scope finalizer)负责关闭(见 DenoSocket.ts)。fromConn(open):把任意Deno.Conn适配为 Effect socket,返回的 socket 提供reader/writer,并支持就地 TLS 升级(见下文)。makeTcpChannel/layerTcp:将 TCP 连接包装为 Effect Channel 或Socket服务层。layerWebSocket/layerWebSocketConstructor:基于globalThis.WebSocket的 WebSocket socket 层(供 HTTP/WebSocket 传输的 runner 使用)。
TLS 升级与平台差异
fromConn的升级逻辑反映了 Deno 平台的两个关键限制(源码注释中有明确说明,见 DenoSocket.ts):
- TCP 连接可以通过
Deno.startTls就地升级为 TLS;由于 Deno 只暴露客户端侧的升级,Unix socket 升级会以SocketUpgradeError失败。 Deno.startTls不支持客户端证书,因此key、cert、passphrase、requestCert选项对 Deno 升级无效;rejectUnauthorized: false仅关闭主机名校验,证书链校验依然生效。
此外,Deno平台不支持keepAliveInitialDelay,noDelay与keepAlive对 Unix 连接无效;CloseEvent始终以优雅方式关闭,因为没有 Node 的 reset-on-close 等价物。
实战:用 Deno socket 层搭建 runner 与 client
仓库的测试 packages/platform/deno/test/cluster/SocketRunner.test.ts 完整演示了 runner 与 client 的装配方式,可直接作为模板。
装配一个 runner 层
const makeRunnerLayer = (port: number, entities: Layer.Layer<never, never, Sharding.Sharding>) => entities.pipe( Layer.provideMerge(SocketRunner.layer), Layer.provide(RunnerHealth.layerNoop), Layer.provide(DenoClusterSocket.layerSocketServer), Layer.provide(DenoClusterSocket.layerClientProtocol), Layer.provide(ShardingConfig.layer({ runnerAddress: Option.some(RunnerAddress.make("127.0.0.1", port)), entityTerminationTimeout: 0, entityMessagePollInterval: 5000, sendRetryInterval: 100 })), Layer.provide(RpcSerialization.layerSchemaBinary()) )要点:
RunnerAddress.make("127.0.0.1", port)指定 runner 的监听地址,layerSocketServer会根据runnerListenAddress(缺省回退到runnerAddress)启动监听;RunnerHealth.layerNoop或DenoClusterSocket.layer中的runnerHealth: "ping"/"k8s"决定健康检查策略;- 序列化层必须与对端一致,这里使用
RpcSerialization.layerSchemaBinary()。
装配一个 client-only 层
const makeClientLayer = (port: number) => SocketRunner.layerClientOnly.pipe( Layer.provide(DenoClusterSocket.layerClientProtocol), Layer.provide(ShardingConfig.layer({ runnerAddress: Option.some(RunnerAddress.make("127.0.0.1", port)), runnerListenAddress: Option.some(RunnerAddress.make("127.0.0.1", port)), entityTerminationTimeout: 0, entityMessagePollInterval: 5000, sendRetryInterval: 100 })), Layer.provide(RpcSerialization.layerSchemaBinary()) )client-only 模式不启动 socket 服务器,只负责把 RPC 请求发往指定 runner 并接收回复。测试中 runner 先以Effect.forkScoped启动,客户端稍后连接并调用实体的 RPC 方法(如client.Process(...))。
测试覆盖的行为
SocketRunner.test.ts中的两条用例(每条 30 秒超时)验证了关键语义:
- 丢弃(discard)请求语义:持久化(persisted)请求在
discard: true时仍需完成序列化(包括 BigDecimal 这类带循环引用的值),而易失(volatile)请求发送后即返回,不等待宿主 runner 处理完毕; - 错误隔离:某个请求的回复序列化失败(
MalformedMessage)只导致该请求失败,同一连接上的其他在途请求不受影响,也不会被错误地重新投递进实体的去重保护(不出现AlreadyProcessingMessage)。
安装与使用前提
@effect/platform-deno的安装方式(见 packages/platform/deno/README.md):
npm install effect@rc @effect/platform-deno@rc根据 packages/platform/deno/package.json,该包要求Deno >= 2.8.3,且以effect作为 peer 依赖。使用时从@effect/platform-deno导入DenoClusterSocket命名空间即可:
import { DenoClusterSocket } from "@effect/platform-deno"小结
本次 changeset 为 Effect Cluster 引入了完整的 Deno 原生 socket 传输链路:DenoClusterSocket.layer一键装配分片层,layerClientProtocol/layerSocketServer提供 TCP 上的 RPC 与监听能力,DenoSocket模块则给出 Deno 连接与 Effect Socket 抽象之间的底层适配(含 TLS 升级与平台差异处理)。对于希望避开 HTTP 开销、直接以原始 socket 在 Deno 上运行 Effect Cluster 的开发者而言,这提供了一条与 Node 平台对等的开箱即用路径。若需要 HTTP/WebSocket 传输,仓库中还提供了对应的 DenoClusterHttp.ts 模块,可满足不同网络环境下的部署选择。
【免费下载链接】effectBuild production-ready applications in TypeScript项目地址: https://gitcode.com/GitHub_Trending/ef/effect
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考