ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

iii Rust SDK iii-helpers 全解:http、stream、observability、queue 与 RBAC 五大类型模块源码级参考

2026/9/14 17:18:46 拓冰建站 浏览量
iii Rust SDK iii-helpers 全解:http、stream、observability、queue 与 RBAC 五大类型模块源码级参考 iii Rust SDK iii-helpers 全解http、stream、observability、queue 与 RBAC 五大类型模块源码级参考【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii本文围绕 iii 项目的iii-helperscrateRust 版 Helper 原语库展开覆盖其http、observability、queue、stream、worker_connection_manager五个模块的全部公开类型与字段定义、默认值与环境变量覆盖规则、wire 格式细节并结合 helpers 源码 与测试用例说明这些类型在跨 SDK 场景下的兼容性设计帮助 Rust Worker 开发者可直接复制运行的 API 参考。安装与模块总览iii-helpers是 iii 多语言 SDK 共享的 Helper 原语 crate定义见 Cargo.tomlcrate 名iii-helpers当前版本 0.23.0-rc.9要求 Rust 1.85、edition 2024许可证 Apache-2.0cargo add iii-helpers从 lib.rs 看crate 只导出五个顶层模块各自职责清晰模块导入职责httpuse iii_helpers::http;HTTP 请求/响应类型、认证配置、HttpInvocationConfigobservabilityuse iii_helpers::observability;OpenTelemetry 初始化、OtelConfig、结构化Logger、span 辅助queueuse iii_helpers::queue;队列入队结果类型streamuse iii_helpers::stream;stream 触发器配置、变更事件、IO 输入、原子更新操作worker_connection_manageruse iii_helpers::worker_connection_manager;RBAC 认证与注册回调类型该 crate 的底层依赖见 Cargo.toml包括opentelemetry/opentelemetry_sdk0.31、tokio、tokio-tungstenite、reqwestrustls、schemars为触发器配置生成 JSON Schema供 codegen 使用。需要注意官方 API 参考页 helpers-rust 是自动生成的由 generate-api-docs.mts 从sdk/packages/rust/helpers/src下的 doc-comments 渲染因此阅读源码注释与阅读文档是同一份事实来源。http 模块HTTP 触发函数与外部 HTTP 调用http模块http.rs提供 HTTP 请求/响应类型、认证配置以及用于调用外部 HTTP 函数Lambda、Cloudflare Workers 等的HttpInvocationConfig。HttpMethodHttpInvocationConfig接受的 HTTP 方法枚举。源码注释特别指出它与核心builtin_triggers的 HTTP 方法枚举不同后者还覆盖 HEAD/OPTIONS。serde 序列化为全大写字符串#[serde(rename_all UPPERCASE)]变体wire 值GetGETPostPOSTPutPUTPatchPATCHDeleteDELETEHttpAuthConfigHTTP 被调函数的认证配置按type字段做 tagged 序列化小写Hmac使用共享密钥的 HMAC 签名校验wire 形如{type: hmac, secret_key: ...}BearerBearer token 认证{type: bearer, token_key: ...}ApiKey通过自定义 header 发送 API keyserde rename 为api_key{type: api_key, header: X-Api-Key, value_key: ...}。变体字段说明Hmacsecret_key: StringHMAC 共享密钥Bearertoken_key: StringBearer tokenApiKeyheader: String, value_key: String自定义 header 名与 token 来源HttpInvocationConfig配置一个被 HTTP 调用的函数端点字段类型必填说明urlString是要调用的 URLmethodHttpMethod是HTTP 方法缺省时 serde default 为POST见default_http_method()timeout_msOptionu64否超时毫秒数headersHashMapString, String是随请求发送的自定义 headerauthOptionHttpAuthConfig否认证配置timeout_ms与auth在序列化为None时会被跳过skip_serializing_if保证与旧版 wire 载荷兼容。HttpRequest 与 HttpResponse函数处理器收到/返回的缓冲式 HTTP 消息均为泛型body 默认serde_json::ValueHttpRequestT Value字段类型说明query_paramsHashMapString, String请求 URL 的 query 参数path_paramsHashMapString, String从匹配路由提取的路径参数headersHashMapString, String请求头pathString请求路径methodStringHTTP 方法如GET、POSTbodyT解析后的请求体所有字段都带#[serde(default)]旧版引擎省略字段时可安全反序列化。HttpResponseT Value字段类型说明status_codeu16HTTP 状态码headersHashMapString, String响应头bodyT响应体queue 模块入队结果queue模块queue.rs只有一个类型EnqueueResult当函数以TriggerAction.Enqueue方式被调用时返回字段类型wire 名说明message_receipt_idStringmessageReceiptId入队消息的唯一回执 ID注意其 serde rename 为 camelCase 的messageReceiptId——与 Node/Python SDK 共享同一 wire 契约跨语言消费时必须按 camelCase 解析。stream 模块stream 触发器、变更事件与原子更新stream模块stream.rs是五个模块中类型最多、wire 细节最重的一个覆盖 stream 触发器配置、变更事件、IO 输入类型与原子更新操作。IO 输入类型类型字段说明StreamGetInputstream_name,group_id,item_id读取单个 stream 项StreamSetInputstream_name,group_id,item_id,data: Value设置覆盖单个 stream 项StreamDeleteInputstream_name,group_id,item_id删除单个 stream 项StreamListInputstream_name,group_id列出组内所有项StreamListGroupsInputstream_name列出流内所有组StreamUpdateInputstream_name,group_id,item_id,ops: VecUpdateOp原子地按序应用一组更新操作更新操作 UpdateOp 与 MergePathUpdateOp是以type为 tag小写的 tagged 枚举表示可对 stream 值原子执行的操作操作JSON 形如说明Set{type: set, path: a.b, value: ...}在 path 处覆盖写入Merge{type: merge, path?: MergePath, value: {}}对象合并进已有值仅对象。path 可省略根合并、为单个一级 key、或字面量段数组嵌套合并Increment{type: increment, path: n, by: 1}数值自增Decrement{type: decrement, path: n, by: 1}数值自减Append{type: append, path?: MergePath, value: ...}向数组追加元素或对字符串拼接path 语义同 MergeRemove{type: remove, path: x}删除字段Rust 侧提供了一组构造器辅助方法stream.rs 中impl UpdateOpset、increment、decrement、append、append_root、append_at_path、remove、merge、merge_at、merge_at_path例如use iii_helpers::stream::UpdateOp; use serde_json::json; let ops vec![ UpdateOp::increment(count, 1), UpdateOp::append_at_path([entityId, buffer], json!(chunk)), UpdateOp::merge_at_path([sessions, abc], json!({ts: now})), UpdateOp::set(status, Some(json!(done))), UpdateOp::remove(temp), ];MergePath的 wire 兼容陷阱MergePath接受单个字符串legacy/一级字段或字面量段数组嵌套路径。引擎侧的路径归一化规则缺省 /Single()/Segments(vec![])→ 根合并Single(foo)等价于Segments(vec![foo])Segments([a, b, c])走三级字面 key从不把点号解释为分隔符Segments(vec![a.b])是一个名为a.b的单一 key。源码中有一行显眼的注释“VARIANT-ORDER-LOAD-BEARING”#[serde(untagged)]按声明顺序尝试变体Single必须位于Segments之前否则裸 JSON 字符串会被反序列化成单元素Segments破坏跨 SDK 的 wire 兼容。这个约束由单元测试merge_path_single_variant_deserializes_string_first显式锁定若变体顺序被工具链或格式化改动测试即失败。更新结果与逐操作错误类型字段说明StreamSetResultold_value: OptionValue,new_value: Valueset 前后的值StreamDeleteResultold_value: OptionValue删除前的值若存在StreamUpdateResultold_value,new_value,errors: VecUpdateOpError原子更新结果成功应用的操作仍反映在new_value中errors为空时从 JSON 中省略向后兼容UpdateOpError报告单个操作的失败细节字段类型说明op_indexusize出错操作在原ops数组中的下标codeString稳定错误码例如merge.path.too_deepmessageString人类可读描述适用时含具体数值doc_urlOptionString该错误类的文档 URL可选测试update_result_with_errors_serializes_field给出了一个真实错误样例code merge.path.too_deep、message Path depth 33 exceeds maximum of 32——即引擎对合并路径深度设了 32 的上限。变更事件与触发器配置StreamChangeEvent是stream触发器的处理器输入由stream::set、stream::update、stream::delete引发的变更触发。注意其 wire 字段名是 camelCasestreamName、groupId字段类型说明event_typeString恒为streamwire 字段名typetimestampi64事件 Unix 时间戳毫秒stream_nameString发生变更的 streamgroup_idString发生变更的 groupidOptionString变更的项 IDeventStreamChangeEventDetail变更详情StreamChangeEventDetail包含event_type: StreamEventTypewire 字段名type取值为小写create/update/delete与data: Value。触发器配置类型均带 builder 方法与Default实现供 codegen 与运行时共用并 derive 了JsonSchema类型可选过滤字段StreamTriggerConfigstream_name、group_id、item_id、condition_function_idStreamJoinLeaveTriggerConfigstream_name、condition_function_id用于stream:join/stream:leave触发器condition_function_id指定一个在调用 handler 前先求值的前置条件函数。stream 认证与加入结果类型字段说明StreamAuthInputheaders,path,query_params: HashMapString, VecString,addrstream 认证函数的输入请求头、路径、query 参数、客户端地址StreamAuthResultcontext: OptionValue认证成功后传递给 stream 处理器的任意上下文StreamJoinResultunauthorized: booljoin 请求是否被拒绝StreamJoinLeaveEventsubscription_id,stream_name,group_id,id: OptionString,context: OptionValuejoin/leave 触发器的事件载荷context来自StreamAuthResultobservability 模块OTel 初始化、结构化日志与重连策略observability模块observability/mod.rs是iii-helpers中最重的模块负责把 Worker 的 trace、metrics、logs 通过共享 WebSocket 上报到 III Engine。除下述配置类型外它还 re-export 了一组运行时辅助函数init_otel、flush_otel、shutdown_otel、run_in_span、with_span以及 span/baggage 操作record_span_event、set_current_span_attribute、inject_traceparent、run_with_baggage等与execute_traced_request为 reqwest 出站请求创建 CLIENT span并转导出原始opentelemetrycrate让使用方无需直接依赖 OTel。OtelConfigOpenTelemetry 初始化的全量配置。所有字段都是OptionderiveDefault未设置时按下列默认值与环境变量回退字段类型默认值 / 环境变量覆盖enabledOptionbool默认 true设为false或OTEL_ENABLEDfalse/0/no/off可禁用导出service_nameOptionString默认取OTEL_SERVICE_NAME环境变量service_versionOptionString默认取SERVICE_VERSION环境变量否则unknownservice_namespaceOptionString默认取SERVICE_NAMESPACE环境变量service_instance_idOptionString默认取SERVICE_INSTANCE_ID环境变量或自动生成的 UUIDengine_ws_urlOptionString默认取III_URL环境变量或ws://localhost:49134metrics_enabledOptionbool默认 trueOTEL_METRICS_ENABLEDfalse/0/no/off可禁用metrics_export_interval_msOptionu64默认 6000060 秒reconnection_configOptionReconnectionConfigWebSocket 重连配置shutdown_timeout_msOptionu64关闭序列超时默认 10,000mschannel_capacityOptionusize内部遥测消息通道容量默认 10,000。这是 exporter 与 WebSocket 连接循环之间的在途缓冲刻意大于max_pending_messages以在正常运行时吸收突发同时限制重连时的陈旧数据spans_flush_interval_msOptionu64span 处理器 flush 延迟默认 100msOTel 默认 5000ms 是 trace 延迟数秒才出现的原因。环境变量OTEL_SPANS_FLUSH_INTERVAL_MSlogs_enabledOptionbool是否启用 log exporter默认 truelogs_flush_interval_msOptionu64log flush 延迟默认 100mslogs_batch_sizeOptionusize每批导出的最大 log 记录数默认 1fetch_instrumentation_enabledOptionbool是否自动插桩出站 HTTP 调用。默认含None为 true此时可用execute_traced_request()为 reqwest 请求创建 CLIENT spanSome(false)表示退出live_spansOptionbool是否把 span 的开始作为零时长 OTLP 快照上报给 engineLiveSpanStartProcessor让 live trace 视图能渲染进行中的工作。每个 span 多一帧engine 将其存为pending若其 live-span 存储关闭则丢弃最终 span 原地替换它。默认启用环境变量OTEL_LIVE_SPANS与 engine 自身的 start mirroring 使用同一个开关init_otel(config: OtelConfig) - bool的实现在 telemetry/mod.rs全局 OTel 状态只初始化一次flush_otel/shutdown_otel分别用于进程退出前强制 flush 与带超时的优雅关闭超时即shutdown_timeout_ms。底层传输帧使用魔数前缀区分三类遥测OTLPtraces、MTRCmetrics、LOGSlogs定义在 telemetry/types.rs。ReconnectionConfigWebSocket 重连行为配置types.rsDefault实现与测试test_reconnection_config_defaults均可核对字段类型默认值说明initial_delay_msu641000起始延迟毫秒max_delay_msu6430000延迟上限毫秒backoff_multiplierf642.0指数退避倍数jitter_factorf640.3随机抖动因子0–1max_retriesOptionu64None无限重试最大重试次数max_pending_messagesusize1000跨重连保留的最大消息数超出部分被丢弃避免长时间断连后投递陈旧数据。刻意小于OtelConfig::channel_capacityeffective_initial_delay_ms()返回被钳制到最小 1ms 的初始延迟防止零值参与退避计算导致除零测试test_reconnection_config_zero_delay_clamped验证。ConnectionState 与 BaggageSpanProcessorConnectionState共享 WebSocket 的五态机——Disconnected、Connecting、Connected、Reconnecting、Failed。BaggageSpanProcessorfn new() - Self构造的 OTel span processor把当前 OTel baggage 条目复制为每个新 span 的属性使同一请求链上的 span 携带相同的业务上下文。Logger结构化 logger把日志作为 OTel LogRecord 导出use iii_helpers::observability::Logger; use serde_json::json; let logger Logger::new(); logger.info(Worker connected, None); logger.info(Order processed, Some(json!({ order_id: ord_123, amount: 49.99, currency: USD }))); logger.warn(Retry attempt, Some(json!({ attempt: 3, max_retries: 5 }))); logger.error(Payment failed, Some(json!({ order_id: ord_123, error_code: card_declined })));四个日志方法debug/info/warn/error均为fn(str, OptionValue)。两个关键行为见 logger.rs 的文档注释与实现每次日志调用自动捕获当前 trace 与 span 上下文把日志与分布式 trace 关联起来无需手动接线OTel 未初始化时优雅回退到tracingcrate第二个参数传serde_json::Value键值对而非字符串插值时json_value_to_anyvalue会把嵌套对象/数组原样转换为 OTLP 的kvlistValue/arrayValue结构而不是被字符串化——这样可在可观测性后端做过滤、聚合和仪表盘。WorkerGaugesOptions注册 worker gauge 指标的选项otel_worker_gauges.rs字段类型说明worker_idString上报 gauge 的 worker 的稳定标识worker_nameOptionStringworker 的人类可读名称该模块还有配套的集成测试tests/observability_init_test.rs、tests/observability_logger_test.rs 等覆盖初始化、日志发射与上下文捕获路径。worker_connection_manager 模块RBAC 认证与注册回调worker_connection_manager模块worker_connection_manager.rs定义 RBAC 代理 Worker 使用的认证与注册钩子类型与 engine 侧的 rbac_session 实现 对应测试注释中明确指明默认值必须与引擎的规范AuthResult保持一致。认证AuthInput 与 AuthResultWorker 的 WebSocket upgrade 请求经过 RBAC 端口时认证函数收到AuthInput字段类型说明headersHashMapString, Stringupgrade 请求的 HTTP 头query_paramsHashMapString, VecStringupgrade URL 的 query 参数值用 Vec 支持重复 key如?a1a2ip_addressString连接客户端的 IP认证函数返回AuthResult控制已认证 worker 可调用哪些函数、以及转发给中间件什么上下文字段类型说明namespacesHashMapString, VecString按命名空间划分的授权如{ orders: [svc::*] }值是精确函数 ID 或通配符裸写法svc::*或match(svc::*)均可。键同时也是该会话在engine::workers::register时可声明的命名空间留空表示不加命名空间级授权——此时会话可声明任意命名空间仅allowed_functions与expose_functions生效allowed_functionsVecString在default命名空间中、超出expose_functions配置额外允许的函数 ID按命名空间的授权应放入namespacesforbidden_functionsVecString即使匹配expose_functions也要拒绝的函数 ID优先级高于允许项allowed_trigger_typesOptionVecString允许注册触发器的触发器类型 IDNone表示全部允许allow_trigger_type_registrationbool是否允许注册新触发器类型默认falseallow_function_registrationbool是否允许注册新函数默认truecontextValue每次调用时转发给中间件函数的任意上下文默认空对象{}function_registration_prefixOptionString应用于该 worker 注册的所有函数 ID 的前缀可选serde 默认值{}载荷可完整反序列化由测试auth_result_defaults_match_engine锁定防止与引擎侧规范默认值漂移。注册钩子函数、触发器、触发器类型Worker 通过 RBAC 端口注册资源时三类钩子各有一组 Input/Result 类型。统一约定返回 Result省略的字段保留注册请求的原值或返回错误以拒绝注册Result 结构均 deriveDefault即“全省略 不修改”钩子Input 关键字段Result 可映射字段on_function_registration_function_idfunction_id、description?、metadata?、namespace、contextfunction_id?、description?、metadata?on_trigger_registration_function_idtrigger_id、trigger_type、function_id、config、metadata?、namespace、contexttrigger_id?、trigger_type?、function_id?、config?on_trigger_type_registration_function_idtrigger_type_id、description、contexttrigger_type_id?、description?两个 Input 中的namespace字段带 wire 兼容默认值旧版引擎省略该字段时按default解析default_namespace_field()保证新旧引擎与 helpers 版本混跑时反序列化不失败。namespace之所以必须传入钩子是因为同一个函数/触发器 ID 可以存在于多个命名空间钩子需要按目标命名空间分别授权。兼容性保障从测试与生成管线看 wire 契约这个 crate 的定位决定了它对 wire 格式极为敏感——它的类型是 Rust 与 Node/Python/浏览器 SDK 之间的共享契约。仓库中有几类证据可以佐证其契约约束变体顺序回归测试merge_path_single_variant_deserializes_string_first专门防止MergePath的 untagged 变体顺序被工具链改动跨 SDK 字段省略语义测试append_with_root_path_round_trips等测试确认path: None序列化时整体省略字段而非显式null并注释说明 Node/Python/浏览器侧把path?解析为“缺省”update_result_without_errors_omits_field_from_json确认errors为空时省略默认值对拍测试auth_result_defaults_match_engine把 helpers 侧AuthResult的 serde 默认值与引擎规范实现逐一比对文档即契约参考页由 generate-api-docs.mts 从源码 doc-comments 生成注释中写明任何文本修正都应改源码注释后重新生成而不是手改生成文件。主要路径索引资源路径API 参考生成docs/next/reference/helpers-rust.mdxcrate 入口sdk/packages/rust/helpers/src/lib.rshttp 类型sdk/packages/rust/helpers/src/http.rsstream 类型与 wire 测试sdk/packages/rust/helpers/src/stream.rsqueue 类型sdk/packages/rust/helpers/src/queue.rsobservability 模块入口sdk/packages/rust/helpers/src/observability/mod.rsOtelConfig / ReconnectionConfigsdk/packages/rust/helpers/src/observability/telemetry/types.rsinit_otel 实现sdk/packages/rust/helpers/src/observability/telemetry/mod.rsRBAC 类型sdk/packages/rust/helpers/src/worker_connection_manager.rsRust Worker 示例sdk/packages/rust/iii-example含logger_example.rs等适用前提与限制以上内容以仓库当前iii-helpers0.23.0-rc.9 为准OtelConfig的engine_ws_url默认指向本机ws://localhost:49134部署到远程 Engine 时须通过III_URL或显式配置覆盖spans_flush_interval_ms、OTEL_LIVE_SPANS等行为的最终呈现还依赖 engine 侧的对应开关如 live-span 存储两侧需版本配套。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考