iii-helpers Rust 助手库完全指南:HTTP 调用、观测、队列、Stream 与 RBAC 实战
iii-helpers Rust 助手库完全指南HTTP 调用、观测、队列、Stream 与 RBAC 实战【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iiiiii-helpers是 III 项目中 Rust SDKsdk/packages/rust/helpers内跨 SDK 共享的助手原语集合为 Worker 开发提供 HTTP 调用配置、OpenTelemetry 观测、队列投递、Stream 状态操作与 RBAC 鉴权等开箱即用的类型与工具。本文基于 docs/reference/helpers-rust.mdx 完整展开该 crate 的 API 参考并结合源码实现与可运行示例深入讲解读完你将掌握如何用iii_helpers编写被 HTTP 触发的函数、输出与分布式追踪关联的结构化日志、执行 Stream 原子更新以及为 Worker 配置细粒度 RBAC 权限。安装与模块总览在 Rust 项目中添加依赖cargo add iii-helpers从 sdk/packages/rust/helpers/Cargo.toml 可以看到该 crate 基于serde/serde_json启用unbounded_depth、tokio、tokio-tungsteniterustls 原生根证书、opentelemetry0.31 与opentelemetry_sdklogs/metrics/trace、reqwestJSON rustls、schemars构建并附带opentelemetry-http用于 reqwest 追踪插桩。模块布局定义在 sdk/packages/rust/helpers/src/lib.rs共五个公共模块httpHTTP 请求/响应类型、认证配置与调用配置observabilityLogger、OpenTelemetry 配置、Span 与 WebSocket 重连助手queue队列投递结果类型streamStream 触发器配置、变更事件、IO 输入与原子更新操作worker_connection_managerRBAC 鉴权与注册回调类型httpHTTP 调用型函数的请求、响应与认证HTTP 模块适用于被 HTTP 触发的函数文档明确提到 Lambda、Cloudflare Workers 等外部端点场景提供调用配置、三种认证方案以及处理器的输入输出结构。导入方式use iii_helpers::http;其序列化细节可在 sdk/packages/rust/helpers/src/http.rs 中核对HttpMethod通过#[serde(rename_all UPPERCASE)]序列化为大写字符串HttpAuthConfig使用#[serde(tag type, rename_all lowercase)]的标签式枚举HttpInvocationConfig::method在缺省时由default_http_method()提供POST。HttpAuthConfig三种认证方案HTTP 调用型函数的认证配置支持三个变体Hmac { secret_key: String }使用共享密钥进行 HMAC 签名校验Bearer { token_key: String }Bearer Token 认证ApiKey { header: String, value_key: String }通过自定义请求头发送 API Key注意ApiKey在 JSON 线格式中的类型标签为api_key源码通过#[serde(rename api_key)]显式指定。HttpInvocationConfig外部端点调用配置字段类型必填说明urlString是要调用的 URLmethodHttpMethod是HTTP 方法默认POSTtimeout_msOptionu64否超时时间毫秒headersHashMapString, String是随请求发送的自定义请求头authOptionHttpAuthConfig否认证配置HttpMethod调用方法枚举HttpMethod是HttpInvocationConfig接受的 HTTP 方法集合仅含Get、Post、Put、Patch、Delete五种。源码注释特别说明它与核心引擎builtin_triggers中的 HTTP 方法枚举不同——后者还覆盖HEAD/OPTIONS。HttpRequest 与 HttpResponse处理器输入输出函数处理器收到的缓冲 HTTP 请求泛型T默认Value字段类型说明query_paramsHashMapString, StringURL 中的查询字符串参数path_paramsHashMapString, String从匹配路由提取的路径参数headersHashMapString, String请求头pathString请求路径methodString请求的 HTTP 方法如GET、POSTbodyT解析后的请求体函数返回的缓冲 HTTP 响应同样泛型字段类型说明status_codeu16HTTP 状态码headersHashMapString, String响应头bodyT响应体实战函数内发起带追踪的外部 HTTP 调用sdk/packages/rust/iii-example/src/http_example.rs 给出了完整用法函数内用reqwest构造请求通过execute_traced_request(client, request)发起在fetch_instrumentation_enabled开启时该函数为出站请求创建 CLIENT span最后把上游结果包装成HttpResponse返回use iii_helpers::http::{HttpRequest, HttpResponse}; use iii_helpers::observability::{Logger, execute_traced_request}; use iii_sdk::builtin_triggers::{HttpMethod, HttpTriggerConfig}; use iii_sdk::trigger::IIITrigger; use iii_sdk::{Error, IIIClient, RegisterFunction}; use serde_json::json; iii.register_function( api::get::http::rust::fetch, RegisterFunction::new_async(move |_input: serde_json::Value| { let client client.clone(); let logger Logger::new(); async move { logger.info(Fetching todo from external API, None); let request client .get(https://jsonplaceholder.typicode.com/todos/1) .build() .map_err(|e| Error::Handler(e.to_string()))?; let response execute_traced_request(client, request) .await .map_err(|e| Error::Handler(e.to_string()))?; let status response.status().as_u16(); logger.info(Fetched todo successfully, Some(json!({ status: status }))); let data: serde_json::Value response .json::serde_json::Value() .await .map_err(|e| Error::Handler(e.to_string()))?; let api_response HttpResponse { status_code: 200, body: json!({ upstream_status: status, data: data }), headers: [(Content-Type.into(), application/json.into())].into(), }; Ok(serde_json::to_value(api_response)?) } }), );observabilityLogger、OpenTelemetry 与 Span 助手观测模块提供 Logger、OpenTelemetry 初始化配置、Span 处理器与 WebSocket 重连配置是 Worker 遥测能力的核心。导入方式use iii_helpers::observability;模块在 sdk/packages/rust/helpers/src/observability 下组织。除文档列出的类型外mod.rs 还公开了大量实用函数init_otel/shutdown_otel/flush_otel、run_in_span/with_span、上下文捕获与注入capture_otel_context、extract_traceparent、inject_baggage、get_baggage_entry等、Span 操作set_current_span_attribute、record_span_event、set_current_span_error、载荷脱敏redact、redact_and_truncate、REDACTED_PLACEHOLDER以及execute_traced_request。Logger输出 OTel LogRecord 的结构化日志Logger 把日志作为 OpenTelemetry LogRecord 发出每次日志调用自动捕获当前 trace 与 span 上下文无需手动接线即可把日志与分布式追踪关联当 OTel 未初始化时优雅回退到tracingcrate。文档建议把结构化数据作为第二个参数传入——使用serde_json::Value键值对对象而非字符串插值便于在观测后端过滤、聚合和构建看板。方法签名说明newfn() - Self创建新 Logger 实例debugfn(message: str, data: OptionValue)记录 debug 级别日志infofn(message: str, data: OptionValue)记录 info 级别日志warnfn(message: str, data: OptionValue)记录 warning 级别日志errorfn(message: str, data: OptionValue)记录 error 级别日志logger.rs 的实现会把serde_json::Value递归转换为 OTelAnyValue嵌套对象映射为kvlistValue、数组映射为arrayValue从而在 OTLP 属性中保留完整结构而不被字符串化。结合 logger_example.rs 的示例use iii_helpers::observability::Logger; use serde_json::{Value, json}; let logger Logger::new(); // 基础日志trace 上下文自动注入 logger.info(Processing request, Some(json!({ input: input }))); logger.debug(Validating input fields, Some(json!({ step: validation }))); // 结构化上下文用于看板与告警 logger.warn(Using default timeout, Some(json!({ timeout_ms: 5000, reason: not configured }))); logger.error(Payment failed, Some(json!({ order_id: ord_123, gateway: stripe, error_code: card_declined }))); logger.info(Request processed successfully, None);OtelConfigOpenTelemetry 初始化配置这是观测模块最核心的配置结构所有字段可选均带默认值与环境变量覆盖字段类型默认值说明enabledOptionbooltrue是否启用 OTel 导出设为false或环境变量OTEL_ENABLEDfalse/0/no/off可关闭service_nameOptionStringOTEL_SERVICE_NAME环境变量上报的服务名service_versionOptionStringSERVICE_VERSION或unknown上报的服务版本service_namespaceOptionStringSERVICE_NAMESPACE环境变量上报的服务命名空间service_instance_idOptionStringSERVICE_INSTANCE_ID或自动生成的 UUID服务实例 IDengine_ws_urlOptionStringIII_URL或ws://localhost:49134III 引擎 WebSocket URLmetrics_enabledOptionbooltrue是否启用指标导出OTEL_METRICS_ENABLEDfalse/0/no/off可关闭metrics_export_interval_msOptionu646000060 秒指标导出间隔毫秒reconnection_configOptionReconnectionConfig无WebSocket 重连配置shutdown_timeout_msOptionu6410000关闭序列超时毫秒channel_capacityOptionusize10000内部遥测消息通道容量即导出器与 WebSocket 连接循环之间的在途消息缓冲。有意大于ReconnectionConfig::max_pending_messages以便正常运行时吸收突发流量同时在重连时限制陈旧数据spans_flush_interval_msOptionu64100Span 处理器刷新延迟毫秒。OTel 默认 5000ms 正是导致 trace 在操作数秒后才出现的原因。环境变量覆盖OTEL_SPANS_FLUSH_INTERVAL_MSlogs_enabledOptionbooltrue是否启用日志导出器logs_flush_interval_msOptionu64100日志处理器刷新延迟毫秒logs_batch_sizeOptionusize1每批导出的日志记录最大条数fetch_instrumentation_enabledOptionboolSome(true)None视为true是否自动插桩出站 HTTP 调用开启后可用execute_traced_request()为 reqwest 请求创建 CLIENT spanlive_spansOptionbool开启向引擎发布零结束 OTLP 快照的 span 开始事件LiveSpanStartProcessor使实时 trace 视图能渲染进行中的工作每个 span 多一帧引擎以pending存储或在其实时 span 存储关闭时丢弃最终 span 原位替换。环境变量覆盖OTEL_LIVE_SPANSReconnectionConfigWebSocket 重连行为字段类型默认值说明initial_delay_msu641000起始延迟毫秒max_delay_msu6430000最大延迟上限毫秒backoff_multiplierf642指数退避乘数jitter_factorf640.3随机抖动因子取值 0-1max_retriesOptionu64None无限重试最大重试次数max_pending_messagesusize无重连期间最多保留的消息数超出即丢弃避免长断开后投递陈旧数据。有意小于OtelConfig::channel_capacityeffective_initial_delay_msfn() - u64无返回initial_delay_ms钳制到最小 1ms 以防除零其他类型BaggageSpanProcessornew() - Self。OpenTelemetry span 处理器把 OTel baggage 条目复制到每个已启动 span 的属性上实现跨请求的上下文传递。ConnectionState共享 WebSocket 的连接状态枚举取值Disconnected、Connecting、Connected、Reconnecting、Failed。WorkerGaugesOptions注册 Worker 指标gauges的选项。worker_idString必填为上报指标的 Worker 稳定标识worker_nameOptionString可选为 Worker 可读名称。queue队列投递结果队列模块目前只包含投递结果类型。导入方式use iii_helpers::queue;EnqueueResult当函数以TriggerAction.Enqueue方式被调用即消息进入队列时返回的结果见 sdk/packages/rust/helpers/src/queue.rs字段类型说明message_receipt_idString已入队消息的唯一回执 ID源码中该字段通过#[serde(rename messageReceiptId)]映射为驼峰式 JSON 字段名保证与 Node/Python SDK 的线格式一致。streamStream 触发器、变更事件与原子更新Stream 模块是类型最丰富的部分覆盖触发器配置、变更事件、IO 输入、鉴权与原子更新操作。导入方式use iii_helpers::stream;MergePath 与路径归一化关键实现细节MergePath是UpdateOp::Merge/UpdateOp::Append的路径目标接受单字符串传统/一级字段或字面量段数组嵌套路径。引擎施加的路径归一化规则缺省 /Single()/Segments(vec![])→ 根级合并Single(foo)等价于Segments(vec![foo.into()])Segments([a, b, c])依次走三个字面量键绝不把点号当特殊分隔符Segments(vec![a.b.into()])是名为a.b的单一字面量键变体顺序是承重设计load-bearing#[serde(untagged)]按声明顺序尝试变体Single必须排在Segments之前这样 JSON 字符串才会反序列化为Single而不是先让数组匹配失败。源码 sdk/packages/rust/helpers/src/stream.rs 中的注释明确指出重排会破坏线上兼容性字符串载荷会被反序列化成单元素Segments并有回归测试merge_path_single_variant_deserializes_string_first锁定该行为UpdateOp还提供了set/increment/decrement/append/append_root/append_at_path/remove/merge/merge_at/merge_at_path等构造函数简化调用。UpdateOp可原子应用的流值操作变体字段说明Set{ path: String, value: OptionValue }在路径上设值覆盖Merge{ path: OptionMergePath, value: Value }将对象合并进现有值仅限对象。path 可省略根合并、单一级键或字面量段数组Increment{ path: String, by: i64 }数值自增Decrement{ path: String, by: i64 }数值自减Append{ path: OptionMergePath, value: Value }向数组追加元素或在可选路径处拼接字符串。path 语义同 MergeRemove{ path: String }删除字段序列化采用#[serde(tag type, rename_all lowercase)]标签式枚举例如{type:append,path:chunks,value:{text:hello}}。测试update_append_serializes_as_tagged_operation、append_with_segments_path_round_trips_as_array与append_with_root_path_round_trips验证了单段路径、段数组路径与根路径的往返一致性根路径path: None时字段整体省略而非输出null保证跨 SDK 消费方Node/Python/浏览器都能按缺省字段解码。Stream 变更事件StreamChangeEventstream触发器由stream::set、stream::update或stream::delete触发的条目变更的处理输入。字段类型说明event_typeString恒为streamtimestampi64事件 Unix 时间戳毫秒stream_nameString发生变更的流group_idString发生变更的组idOptionString变更的条目 IDeventStreamChangeEventDetail含变更类型与数据的事件详情StreamChangeEventDetailevent_typeStreamEventType即Create/Update/DeleteJSON 中为create/update/deletedataValue。注意源码中StreamChangeEvent的stream_name/group_id通过 serde rename 映射为streamName/groupId驼峰字段。StreamJoinLeaveEventstream:join/stream:leave触发器的事件载荷含subscription_id唯一订阅标识、stream_name、group_id、id可选条目标识与context来自StreamAuthResult的鉴权上下文。Stream IO 输入/结果类型类型字段说明StreamGetInputstream_name、group_id、item_id读取单个条目StreamSetInputstream_name、group_id、item_id、data: Value写入条目StreamSetResultold_value: OptionValue、new_value: Value写入结果含旧值/新值StreamDeleteInputstream_name、group_id、item_id删除条目StreamDeleteResultold_value: OptionValue删除结果旧值若存在StreamListInputstream_name、group_id列出组内全部条目StreamListGroupsInputstream_name列出流内全部组StreamUpdateInputstream_name、group_id、item_id、ops: VecUpdateOp原子更新ops 为有序操作列表StreamUpdateResult原子更新的结果old_value更新前值若存在、new_value更新后值、errors: VecUpdateOpError应用操作时遇到的错误成功应用的操作仍反映在new_value中该字段为空时从 JSON 中省略以保证向后兼容——源码通过#[serde(skip_serializing_if Vec::is_empty)]实现测试update_result_without_errors_omits_field_from_json验证。UpdateOpError单操作错误op_index原ops数组中的索引、code稳定错误码如merge.path.too_deep、message含具体数字的可读描述、doc_url可选该错误类的文档链接。测试update_result_with_errors_serializes_field展示了深度超限场景Path depth 33 exceeds maximum of 32。Stream 鉴权与触发器配置StreamAuthInput流鉴权输入含headers请求头、path请求路径、query_params: HashMapString, VecString查询参数支持重复键、addr客户端地址。StreamAuthResult流鉴权结果context: OptionValue为鉴权后传给流处理器的任意上下文。StreamJoinResult加入流的结果unauthorized: bool标识是否未授权。StreamJoinLeaveTriggerConfigstream:join/stream:leave触发器配置。stream_name要监听的流、condition_function_id可选调用处理器前先评估的函数 ID。源码提供new()与链式 builderstream_name(name)、condition(function_id)并实现Default。StreamTriggerConfigstream触发器配置用于过滤哪些条目变更会触发处理器。stream_name监听的流、group_id组过滤、item_id条目过滤、condition_function_id可选前置条件函数。同样提供 builder 链式构造。worker_connection_managerRBAC 鉴权与注册回调该模块定义 Worker 通过 RBAC 端口连接时的鉴权输入/输出与三类注册钩子的回调类型与引擎侧rbac_session的默认值对齐sdk/packages/rust/helpers/src/worker_connection_manager.rs 中测试auth_result_defaults_match_engine校验默认值一致性。导入方式use iii_helpers::worker_connection_manager;AuthInput 与 AuthResultAuthInputWebSocket 升级期间传入 RBAC 鉴权函数的输入包含升级请求的 HTTP 头、查询参数与客户端 IP。query_params每个键映射到值数组以支持重复键如?a1a2。AuthResult鉴权函数返回值控制已认证 Worker 可调用哪些函数、注册哪些触发器以及转发给中间件的上下文字段类型说明namespacesHashMapString, VecString按命名空间授予的权限如{ orders: [svc::*] }值可为精确函数 ID 或通配符svc::*或match(svc::*)写法。键也是会话可在engine::workers::register上声明的命名空间留空表示不添加任何命名空间作用域授权allowed_functionsVecString除expose_functions配置外额外允许的函数 ID仅default命名空间命名空间授权请放namespacesforbidden_functionsVecString即使匹配expose_functions也拒绝的函数 ID优先级高于允许allowed_trigger_typesOptionVecString该 Worker 可注册触发器的触发器类型 IDNone表示全部允许allow_trigger_type_registrationbool是否允许注册新触发器类型默认falseallow_function_registrationbool是否允许注册新函数默认truecontextValue每次调用转发给中间件函数的任意上下文function_registration_prefixOptionString应用于该 Worker 注册的所有函数 ID 的可选前缀注册钩子类型三类钩子遵循同一模式输入结构携带被注册对象的元数据与会话鉴权上下文输出结构中的省略字段保持注册请求的原始值即只映射你显式给出的字段直接返回错误即可拒绝注册。每个输入结构还包含namespace字段源码default_namespace_field()默认default用于按目标命名空间授权——同名函数 ID 可存在于多个命名空间。OnFunctionRegistrationInput / OnFunctionRegistrationResultWorker 通过 RBAC 端口注册函数时触发on_function_registration_function_id钩子。输入function_id、description可选、metadata可选、namespace、context。输出function_id、description、metadata均可选映射值。OnTriggerRegistrationInput / OnTriggerRegistrationResult注册触发器时触发on_trigger_registration_function_id钩子。输入trigger_id、trigger_type、function_id、config、metadata可选、namespace、context。输出trigger_id、trigger_type、function_id、config均可选映射值。OnTriggerTypeRegistrationInput / OnTriggerTypeRegistrationResult注册新触发器类型时触发on_trigger_type_registration_function_id钩子。输入trigger_type_id、description、context。输出trigger_type_id、description均可选映射值。与 Rust SDK 的配合使用iii-helpers是 III Rust SDKsdk/packages/rust/iii的伴生库SDK 的InitOptions.otel字段直接接收iii_helpers::observability::OtelConfigregister_worker(ws://localhost:49134, InitOptions::default())建立与引擎的 WebSocket 连接专用后台线程 独立 tokio runtime随后即可注册函数与触发器参见 docs/reference/sdk-rust.mdx.skill.md。可运行的完整示例位于 sdk/packages/rust/iii-examplelogger_example.rs 展示各日志级别的结构化输出http_example.rs 展示外部 HTTP 调用与 CLIENT span 追踪custom_trigger_example.rs 展示自定义触发器类型注册。典型接入路径先以cargo add iii-helpers iii-sdk添加依赖用OtelConfig配置观测例如将spans_flush_interval_ms设为100让 trace 近乎实时可见按需通过OTEL_LIVE_SPANS/OTEL_SPANS_FLUSH_INTERVAL_MS等环境变量覆盖函数内用Logger输出结构化日志再以HttpInvocationConfig、UpdateOp与 RBAC 钩子类型完成业务集成与权限控制。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考