ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

effect 4 新特性:为 Effect Cluster 提供 Deno 原生 Socket Runner 层

effect 4 新特性:为 Effect Cluster 提供 Deno 原生 Socket Runner 层 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 上执行。SocketRunnereffect/unstable/cluster/SocketRunner是 Cluster 中基于原始 socket 的 runner 实现与HttpRunner基于 HTTP/WebSocket并列。本次新增的 Deno socket 层正是把SocketRunner所依赖的各种服务RPC 客户端协议、socket 服务器、健康检查、存储等用 Deno 的原生 API 组装起来让 Deno 应用可以直接以 socket 传输方式加入 Cluster。核心 APIDenoClusterSocket.layerDenoClusterSocket.layer是一个一键式的分片sharding层构造器它会根据传入的选项自动装配传输、序列化、存储、健康检查等全部依赖。函数签名与选项如下来自 DenoClusterSocket.tsexport 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?: PartialShardingConfig.ShardingConfig[Service] | undefined } ): Layer.Layer...选项说明选项取值作用serializationbinary/ndjsonRPC 消息序列化方式。缺省为binarySchema 二进制编码传ndjson则使用 NDJSON 行式编码serializationMaxBufferSize数字 /unbounded序列化缓冲上限。二进制模式下作为maxFrameSizeNDJSON 模式下作为maxBufferSizeclientOnlytrue/false仅构建客户端不启动 socket 服务器用于只发起请求、不承载实体的进程storagelocal/sql/byo存储策略。local使用内存存储sql使用 SQL 存储需要SqlClient并会自动装配DenoCrypto加密层byobring your own要求自行提供MessageStorage与RunnerStoragerunnerHealthping/k8srunner 健康检查方式。ping通过 RPC ping 检测存活k8s通过 Kubernetes API 检测runnerHealthK8s{ namespace?, labelSelector? }k8s健康检查时的命名空间与标签选择器shardingConfigPartialShardingConfig需要覆盖的分片配置项最终会与ShardingConfig.layerFromEnv的环境变量配置合并返回类型与依赖layer的返回类型会随选项精确变化从源码的类型重载可以清楚看到非clientOnly时产出Sharding | Runners | MessageStorage错误通道为SocketServer.SocketServerError | Config.ConfigErrorclientOnly: true时不产出MessageStorage错误通道退化为Config.ConfigErrorstorage为local时不再依赖外部存储sql时需要SqlClientbyo时需要自行提供MessageStorage | RunnerStorage。从实现看layer内部实际是把SocketRunner.layer或SocketRunner.layerClientOnly与layerClientProtocol、layerSocketServer、健康检查层、存储层、序列化层逐层provide组合起来的组合子最终返回一个开箱即用的分片层。组合层的内部实现layerClientProtocolTCP 上的 RPC 客户端协议layerClientProtocol提供 Cluster 所需的RpcClientProtocol其核心是使用DenoSocket.makeTcp打开到目标 runner 地址address.host/address.port的 TCP 连接并设置openTimeout: 10001 秒打开超时随后通过RpcClient.makeProtocolSocket()把连接包装成 RPC 协议 socket见 DenoClusterSocket.ts。值得注意的是源码中的注释与 Node socket 不同Deno 连接没有原生的 idle-timeout 选项因此 peer 连接只使用这 1 秒的 open timeout而不存在连接空闲自动断开的机制。这是迁移到 Deno 平台时需要考虑的行为差异。layerSocketServerrunner 的 socket 服务器layerSocketServer为 runner 提供SocketServer监听地址取自ShardingConfig.runnerListenAddress若该值为空则回退到runnerAddress两者皆空时直接以Effect.die终止见 DenoClusterSocket.ts。它底层委托给DenoSocketServer.layer({ host, port })。layerK8sHttpClientKubernetes 健康检查的 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抽象之间的桥梁它本身也是独立的可复用 APImakeTcp(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.tsTCP 连接可以通过Deno.startTls就地升级为 TLS由于 Deno 只暴露客户端侧的升级Unix socket 升级会以SocketUpgradeError失败。Deno.startTls不支持客户端证书因此key、cert、passphrase、requestCert选项对 Deno 升级无效rejectUnauthorized: false仅关闭主机名校验证书链校验依然生效。此外Deno平台不支持keepAliveInitialDelaynoDelay与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.Layernever, 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.mdnpm install effectrc effect/platform-denorc根据 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),仅供参考
返回列表