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/iii
iii-helpers是 III 项目中 Rust SDK(sdk/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-tungstenite(rustls 原生根证书)、opentelemetry0.31 与opentelemetry_sdk(logs/metrics/trace)、reqwest(JSON + rustls)、schemars构建,并附带opentelemetry-http用于 reqwest 追踪插桩。模块布局定义在 sdk/packages/rust/helpers/src/lib.rs,共五个公共模块:
http:HTTP 请求/响应类型、认证配置与调用配置observability:Logger、OpenTelemetry 配置、Span 与 WebSocket 重连助手queue:队列投递结果类型stream:Stream 触发器配置、变更事件、IO 输入与原子更新操作worker_connection_manager:RBAC 鉴权与注册回调类型
http:HTTP 调用型函数的请求、响应与认证
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:外部端点调用配置
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
url | String | 是 | 要调用的 URL |
method | HttpMethod | 是 | HTTP 方法,默认POST |
timeout_ms | Option<u64> | 否 | 超时时间(毫秒) |
headers | HashMap<String, String> | 是 | 随请求发送的自定义请求头 |
auth | Option<HttpAuthConfig> | 否 | 认证配置 |
HttpMethod:调用方法枚举
HttpMethod是HttpInvocationConfig接受的 HTTP 方法集合,仅含Get、Post、Put、Patch、Delete五种。源码注释特别说明:它与核心引擎builtin_triggers中的 HTTP 方法枚举不同——后者还覆盖HEAD/OPTIONS。
HttpRequest 与 HttpResponse:处理器输入输出
函数处理器收到的缓冲 HTTP 请求(泛型T默认Value):
| 字段 | 类型 | 说明 |
|---|---|---|
query_params | HashMap<String, String> | URL 中的查询字符串参数 |
path_params | HashMap<String, String> | 从匹配路由提取的路径参数 |
headers | HashMap<String, String> | 请求头 |
path | String | 请求路径 |
method | String | 请求的 HTTP 方法(如GET、POST) |
body | T | 解析后的请求体 |
函数返回的缓冲 HTTP 响应(同样泛型):
| 字段 | 类型 | 说明 |
|---|---|---|
status_code | u16 | HTTP 状态码 |
headers | HashMap<String, String> | 响应头 |
body | T | 响应体 |
实战:函数内发起带追踪的外部 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)?) } }), );observability:Logger、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键值对对象(而非字符串插值),便于在观测后端过滤、聚合和构建看板。
| 方法 | 签名 | 说明 |
|---|---|---|
new | fn() -> Self | 创建新 Logger 实例 |
debug | fn(message: &str, data: Option<Value>) | 记录 debug 级别日志 |
info | fn(message: &str, data: Option<Value>) | 记录 info 级别日志 |
warn | fn(message: &str, data: Option<Value>) | 记录 warning 级别日志 |
error | fn(message: &str, data: Option<Value>) | 记录 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);OtelConfig:OpenTelemetry 初始化配置
这是观测模块最核心的配置结构,所有字段可选,均带默认值与环境变量覆盖:
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
enabled | Option<bool> | true | 是否启用 OTel 导出;设为false或环境变量OTEL_ENABLED=false/0/no/off可关闭 |
service_name | Option<String> | OTEL_SERVICE_NAME环境变量 | 上报的服务名 |
service_version | Option<String> | SERVICE_VERSION或"unknown" | 上报的服务版本 |
service_namespace | Option<String> | SERVICE_NAMESPACE环境变量 | 上报的服务命名空间 |
service_instance_id | Option<String> | SERVICE_INSTANCE_ID或自动生成的 UUID | 服务实例 ID |
engine_ws_url | Option<String> | III_URL或"ws://localhost:49134" | III 引擎 WebSocket URL |
metrics_enabled | Option<bool> | true | 是否启用指标导出;OTEL_METRICS_ENABLED=false/0/no/off可关闭 |
metrics_export_interval_ms | Option<u64> | 60000(60 秒) | 指标导出间隔(毫秒) |
reconnection_config | Option<ReconnectionConfig> | 无 | WebSocket 重连配置 |
shutdown_timeout_ms | Option<u64> | 10000 | 关闭序列超时(毫秒) |
channel_capacity | Option<usize> | 10000 | 内部遥测消息通道容量,即导出器与 WebSocket 连接循环之间的在途消息缓冲。有意大于ReconnectionConfig::max_pending_messages,以便正常运行时吸收突发流量,同时在重连时限制陈旧数据 |
spans_flush_interval_ms | Option<u64> | 100 | Span 处理器刷新延迟(毫秒)。OTel 默认 5000ms 正是导致 trace 在操作数秒后才出现的原因。环境变量覆盖:OTEL_SPANS_FLUSH_INTERVAL_MS |
logs_enabled | Option<bool> | true | 是否启用日志导出器 |
logs_flush_interval_ms | Option<u64> | 100 | 日志处理器刷新延迟(毫秒) |
logs_batch_size | Option<usize> | 1 | 每批导出的日志记录最大条数 |
fetch_instrumentation_enabled | Option<bool> | Some(true)(None视为true) | 是否自动插桩出站 HTTP 调用;开启后可用execute_traced_request()为 reqwest 请求创建 CLIENT span |
live_spans | Option<bool> | 开启 | 向引擎发布零结束 OTLP 快照的 span 开始事件(LiveSpanStartProcessor),使实时 trace 视图能渲染进行中的工作;每个 span 多一帧,引擎以pending存储(或在其实时 span 存储关闭时丢弃),最终 span 原位替换。环境变量覆盖:OTEL_LIVE_SPANS |
ReconnectionConfig:WebSocket 重连行为
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
initial_delay_ms | u64 | 1000 | 起始延迟(毫秒) |
max_delay_ms | u64 | 30000 | 最大延迟上限(毫秒) |
backoff_multiplier | f64 | 2 | 指数退避乘数 |
jitter_factor | f64 | 0.3 | 随机抖动因子,取值 0-1 |
max_retries | Option<u64> | None(无限重试) | 最大重试次数 |
max_pending_messages | usize | 无 | 重连期间最多保留的消息数,超出即丢弃,避免长断开后投递陈旧数据。有意小于OtelConfig::channel_capacity |
effective_initial_delay_ms | fn() -> u64 | 无 | 返回initial_delay_ms,钳制到最小 1ms 以防除零 |
其他类型
- BaggageSpanProcessor:
new() -> Self。OpenTelemetry span 处理器,把 OTel baggage 条目复制到每个已启动 span 的属性上,实现跨请求的上下文传递。 - ConnectionState:共享 WebSocket 的连接状态枚举,取值
Disconnected、Connecting、Connected、Reconnecting、Failed。 - WorkerGaugesOptions:注册 Worker 指标(gauges)的选项。
worker_id(String,必填)为上报指标的 Worker 稳定标识;worker_name(Option<String>,可选)为 Worker 可读名称。
queue:队列投递结果
队列模块目前只包含投递结果类型。导入方式:
use iii_helpers::queue;EnqueueResult
当函数以TriggerAction.Enqueue方式被调用(即消息进入队列)时返回的结果,见 sdk/packages/rust/helpers/src/queue.rs:
| 字段 | 类型 | 说明 |
|---|---|---|
message_receipt_id | String | 已入队消息的唯一回执 ID |
源码中该字段通过#[serde(rename = "messageReceiptId")]映射为驼峰式 JSON 字段名,保证与 Node/Python SDK 的线格式一致。
stream:Stream 触发器、变更事件与原子更新
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: Option<Value> } | 在路径上设值(覆盖) |
Merge | { path: Option<MergePath>, value: Value } | 将对象合并进现有值(仅限对象)。path 可省略(根合并)、单一级键或字面量段数组 |
Increment | { path: String, by: i64 } | 数值自增 |
Decrement | { path: String, by: i64 } | 数值自减 |
Append | { path: Option<MergePath>, value: Value } | 向数组追加元素或在可选路径处拼接字符串。path 语义同 Merge |
Remove | { 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 变更事件
StreamChangeEvent:stream触发器(由stream::set、stream::update或stream::delete触发的条目变更)的处理输入。
| 字段 | 类型 | 说明 |
|---|---|---|
event_type | String | 恒为"stream" |
timestamp | i64 | 事件 Unix 时间戳(毫秒) |
stream_name | String | 发生变更的流 |
group_id | String | 发生变更的组 |
id | Option<String> | 变更的条目 ID |
event | StreamChangeEventDetail | 含变更类型与数据的事件详情 |
StreamChangeEventDetail:event_type(StreamEventType,即Create/Update/Delete,JSON 中为create/update/delete)+data(Value)。注意源码中StreamChangeEvent的stream_name/group_id通过 serde rename 映射为streamName/groupId驼峰字段。
StreamJoinLeaveEvent:stream:join/stream:leave触发器的事件载荷,含subscription_id(唯一订阅标识)、stream_name、group_id、id(可选条目标识)与context(来自StreamAuthResult的鉴权上下文)。
Stream IO 输入/结果类型
| 类型 | 字段 | 说明 |
|---|---|---|
StreamGetInput | stream_name、group_id、item_id | 读取单个条目 |
StreamSetInput | stream_name、group_id、item_id、data: Value | 写入条目 |
StreamSetResult | old_value: Option<Value>、new_value: Value | 写入结果(含旧值/新值) |
StreamDeleteInput | stream_name、group_id、item_id | 删除条目 |
StreamDeleteResult | old_value: Option<Value> | 删除结果(旧值若存在) |
StreamListInput | stream_name、group_id | 列出组内全部条目 |
StreamListGroupsInput | stream_name | 列出流内全部组 |
StreamUpdateInput | stream_name、group_id、item_id、ops: Vec<UpdateOp> | 原子更新,ops 为有序操作列表 |
StreamUpdateResult(原子更新的结果):old_value(更新前值,若存在)、new_value(更新后值)、errors: Vec<UpdateOpError>(应用操作时遇到的错误;成功应用的操作仍反映在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: HashMap<String, Vec<String>>(查询参数,支持重复键)、addr(客户端地址)。 - StreamAuthResult:流鉴权结果,
context: Option<Value>为鉴权后传给流处理器的任意上下文。 - StreamJoinResult:加入流的结果,
unauthorized: bool标识是否未授权。 - StreamJoinLeaveTriggerConfig:
stream:join/stream:leave触发器配置。stream_name(要监听的流)、condition_function_id(可选,调用处理器前先评估的函数 ID)。源码提供new()与链式 builderstream_name(name)、condition(function_id),并实现Default。 - StreamTriggerConfig:
stream触发器配置,用于过滤哪些条目变更会触发处理器。stream_name(监听的流)、group_id(组过滤)、item_id(条目过滤)、condition_function_id(可选前置条件函数)。同样提供 builder 链式构造。
worker_connection_manager:RBAC 鉴权与注册回调
该模块定义 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 与 AuthResult
AuthInput:WebSocket 升级期间传入 RBAC 鉴权函数的输入,包含升级请求的 HTTP 头、查询参数与客户端 IP。query_params每个键映射到值数组以支持重复键(如?a=1&a=2)。
AuthResult:鉴权函数返回值,控制已认证 Worker 可调用哪些函数、注册哪些触发器,以及转发给中间件的上下文:
| 字段 | 类型 | 说明 |
|---|---|---|
namespaces | HashMap<String, Vec<String>> | 按命名空间授予的权限,如{ "orders": ["svc::*"] };值可为精确函数 ID 或通配符(svc::*或match("svc::*")写法)。键也是会话可在engine::workers::register上声明的命名空间;留空表示不添加任何命名空间作用域授权 |
allowed_functions | Vec<String> | 除expose_functions配置外额外允许的函数 ID(仅default命名空间;命名空间授权请放namespaces) |
forbidden_functions | Vec<String> | 即使匹配expose_functions也拒绝的函数 ID,优先级高于允许 |
allowed_trigger_types | Option<Vec<String>> | 该 Worker 可注册触发器的触发器类型 ID;None表示全部允许 |
allow_trigger_type_registration | bool | 是否允许注册新触发器类型,默认false |
allow_function_registration | bool | 是否允许注册新函数,默认true |
context | Value | 每次调用转发给中间件函数的任意上下文 |
function_registration_prefix | Option<String> | 应用于该 Worker 注册的所有函数 ID 的可选前缀 |
注册钩子类型
三类钩子遵循同一模式:输入结构携带被注册对象的元数据与会话鉴权上下文,输出结构中的省略字段保持注册请求的原始值(即只映射你显式给出的字段);直接返回错误即可拒绝注册。每个输入结构还包含namespace字段(源码default_namespace_field()默认"default"),用于按目标命名空间授权——同名函数 ID 可存在于多个命名空间。
- OnFunctionRegistrationInput / OnFunctionRegistrationResult:Worker 通过 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 SDK(sdk/packages/rust/iii)的伴生库:SDK 的InitOptions.otel字段直接接收iii_helpers::observability::OtelConfig,register_worker("ws://localhost:49134", InitOptions::default())建立与引擎的 WebSocket 连接(专用后台线程 + 独立 tokio runtime),随后即可注册函数与触发器(参见 docs/reference/sdk-rust.mdx.skill.md)。可运行的完整示例位于 sdk/packages/rust/iii-example:logger_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),仅供参考