使用 iii-sdk 在 Rust 中构建 III 引擎 Worker:函数注册、触发器绑定、流操作与可观测性实战
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
本指南以仓库中 sdk/packages/rust/iii/README.md 为骨架,系统讲解 III 引擎官方 Rust SDK(crate 名为iii-sdk)的完整使用方式:从 Cargo 依赖安装、Hello World 入门,到函数注册、触发器绑定、三种调用模式、流数据操作与日志观测。文章同时结合 src/lib.rs、src/iii.rs、src/protocol.rs 等源码实现,帮助你理解 SDK 底层的 WebSocket 连接管理、重连机制、命名空间解析与 OpenTelemetry 集成原理,读完即可写出可运行、可观测、可重连的 Rust Worker。
安装:将 iii-sdk 加入你的 Cargo 项目
iii-sdk是一个发布在 crates.io 的常规 Rust crate,只需在Cargo.toml中声明依赖即可。根据 Cargo.toml,该包当前版本为0.23.0-rc.9,edition = "2024",最低支持 Rust 1.85,协议为 Apache-2.0。
[dependencies] iii-sdk = "0.11" serde_json = "1" tokio = { version = "1", features = ["full"] }说明:README 中的
iii-sdk = "0.11"是发布在 crates.io 上的稳定版本号;当前仓库内工作区版本的 Cargo.toml 声明为0.23.0-rc.9(见 sdk/packages/rust/iii/Cargo.toml)。实际使用时请以 crates.io 上可解析的版本为准,或直接引用仓库内版本。SDK 的 lib 名称是iii_sdk(下划线),源码位于 src/lib.rs。
SDK 的依赖面体现了它的核心设计:tokio提供异步运行时,tokio-tungstenite(启用rustls-tls-native-roots)承载 WebSocket 通信,schemars负责从 Rust 类型自动推导 JSON Schema,reqwest用于 HTTP 调用型函数(Lambda、Cloudflare Workers 等)的注册配置,iii-helpers工作区 crate 提供可观测性与流操作的类型支持。
Hello World:注册函数并绑定 HTTP 触发器
以下是最小的完整 Worker 示例(源自 README 的 Hello World):连接引擎、注册一个hello::greet函数、绑定一个 HTTP POST 触发器,然后直接以编程方式调用并打印结果。
use iii_sdk::{register_worker, InitOptions, TriggerRequest}; use serde_json::{json, Value}; #[tokio::main] async fn main() -> Result<(), Box<dyn std::error::Error>> { let iii = register_worker("ws://localhost:49134", InitOptions::default()); iii.register_function("hello::greet", |input: Value| async move { let name = input.get("name").and_then(|v| v.as_str()).unwrap_or("world"); Ok(json!({ "message": format!("Hello, {name}!") })) }); iii.register_trigger("http", "hello::greet", json!({ "api_path": "/greet", "http_method": "POST" }))?; let result: Value = iii .trigger(TriggerRequest { function_id: "hello::greet".to_string(), payload: json!({ "name": "world" }), action: None, timeout_ms: None, }) .await?; println!("result: {result}"); Ok(()) }几点值得注意:
register_worker(address, options)建立与引擎的 WebSocket 连接,地址形如ws://localhost:49134。源码中定义了默认引擎地址常量DEFAULT_ENGINE_URL = "ws://127.0.0.1:49134",特意使用 IPv4 回环地址,因为localhost可能解析为::1,而引擎可能只监听 IPv4(见 src/lib.rs)。register_function的闭包返回Result<Value, Error>,异步签名(async move)由 SDK 自动包装,输入输出均为serde_json::Value。register_trigger("http", "hello::greet", config)把 HTTP 触发器绑定到函数,配置里的api_path与http_method决定外部如何触发。- 同步调用返回
Value,即函数执行结果。
核心 API 一览
README 用一张表总结了 SDK 的全部核心操作,这里完整保留并补充说明:
| Operation | Signature | Description |
|---|---|---|
| Initialize | register_worker(address, options) | Create an SDK instance and auto-connect |
| Register function | iii.register_function(id, \|input: Value\| ...) | Register a function that can be invoked by name |
| Register trigger | iii.register_trigger(type, fn_id, config)? | Bind a trigger (HTTP, cron, queue, etc.) to a function |
| Invoke (await) | iii.trigger(TriggerRequest { ... }).await? | Invoke a function and wait for the result |
| Invoke (fire-and-forget) | iii.trigger(TriggerRequest { action: Some(TriggerAction::Void), ... }).await? | Fire-and-forget invocation |
| Invoke (enqueue) | iii.trigger(TriggerRequest { action: Some(TriggerAction::Enqueue { queue }), ... }).await? | Route invocation through a named queue |
register_worker()会在一个独立的后台线程中建立并维护 WebSocket 通信,该线程自带一个 tokio 运行时,同时负责自动重连与 OpenTelemetry 埋点(见 src/lib.rs 与 src/iii.rs)。也就是说,SDK 的自动重连不依赖调用方运行时,Worker 主线程可以专注于业务逻辑。
初始化与连接生命周期
register_worker 与 InitOptions
register_worker(address, options)是推荐的入口,它会按InitOptions配置好客户端后自动调用connect()。InitOptions的可配置字段如下(见 src/lib.rs):
metadata: Option<WorkerMetadata>:自定义 Worker 元数据(运行时、版本、名称、描述、PID、隔离模式等)。默认自动探测 hostname、PID、操作系统与项目名;Managed 身份模式下,进程级环境变量会覆盖其中对应字段。headers: Option<HashMap<String, String>>:WebSocket 握手时携带的自定义 HTTP 头,可用于认证等场景。otel: Option<OtelConfig>:OpenTelemetry 配置。namespace: Option<String>:Worker 所属命名空间,作用范围远超注册本身——Worker 及其函数在此注册,之后该 Worker 的一切行为(trigger的目标解析、register_trigger的绑定位置)都默认跟随该命名空间。identity: WorkerIdentityMode:身份模式(见下文)。
从环境变量解析引擎地址
在iii compose、容器运行时或 systemd 等监督者场景下,引擎地址由监督者通过环境变量注入。SDK 提供了零参数形式:
use iii_sdk::{register_worker_from_env, InitOptions}; // III_URL when set, ws://127.0.0.1:49134 otherwise. let worker = register_worker_from_env(InitOptions::default()); worker.shutdown();engine_url_from_env()的解析顺序是:环境变量III_URL(非空时)→ 默认值ws://127.0.0.1:49134。监督者注入III_URL的方式与注入III_NAMESPACE、III_WORKER_NAME完全一致(见 src/lib.rs)。
优雅关闭:shutdown()
Rust 中进程在main返回时退出,所有线程随之终止。因此必须在main仍在运行时调用shutdown(),它负责:停止连接循环、发送关闭信号、join 后台连接线程,并在线程退出前冲刷 OpenTelemetry 数据:
worker.shutdown(); // cleanly stops the connection thread异步场景下也可以使用shutdown_async().await,它不 join 线程、不会阻塞执行器,但 OTel 冲刷可能在进程退出前来不及完成(见 src/iii.rs)。
连接状态与等待注册完成
get_connection_state()返回IIIConnectionState(Disconnected/Connecting/Connected/Reconnecting/Failed),可用来判断引擎是否可达。若需要等待引擎接受初始注册(例如启动时立即调用函数),可以使用wait_until_registered(timeout);注册被拒绝时返回Error::RegistrationRejected,超时返回Error::Timeout(见 src/iii.rs)。
底层自动重连机制
从源码看,连接循环具备完整的健康保障(见 src/iii.rs 与run_connection实现):
- 连接超时:单次 WS 连接(TCP + TLS + HTTP 升级)上限 10 秒,防止僵死 socket 卡住重连循环。
- 心跳保活:每 20 秒发送一次 WebSocket Ping,让空闲链路持续产生流量。
- 空闲超时:60 秒内未收到任何帧(含 Pong)即判定半开连接已死亡,强制重连——这是为了检测"引擎已把我们断开并注销了函数,但本地仍显示 Connected"的情况。
- 重连间隔:两次重连尝试之间等待 2 秒。
- 重连身份交接:重连时 SDK 会先发送
Reattach消息(携带上次引擎下发的worker_id与reattach_token密钥),让引擎先退役旧连接,再做注册重放,避免新旧连接竞态;token 用于证明"我们就是那个 Worker"(Worker id 本身是公开可列举的)。 - 注册重放与去重:重连成功后,SDK 通过
collect_registrations()重放 trigger type、function、trigger 的全部注册,并用dedupe_registrations/drain_pre_connect_duplicates丢弃握手前积压的重复注册消息,确保幂等。
连接线程的定时参数是私有旋钮(源码注释明确"留待运营有需求时再提升到InitOptions"),但单元测试会缩短这些参数来验证重连路径,例如connect_timeout_abandons_stalled_connect_and_retries与idle_timeout_reconnects_when_engine_goes_silent(见 src/iii.rs)。
注册函数
基础形式:直接传闭包
最简洁的注册方式是把异步闭包直接传给register_function:
use serde_json::{json, Value}; iii.register_function("orders::create", |input: Value| async move { let item = input["body"]["item"].as_str().unwrap_or(""); Ok(json!({ "status_code": 201, "body": { "id": "123", "item": item } })) });进阶:RegisterFunction 构建器
对于需要附加描述、元数据或 Schema 的场景,推荐使用RegisterFunction构建器(定义于 src/iii.rs)。它支持三种构造方式:
RegisterFunction::new(f):同步函数。RegisterFunction::new_async(f):异步函数(推荐)。RegisterFunction::http(config):HTTP 调用型函数(Lambda、Cloudflare Workers 等),不运行本地 handler,引擎侧通过配置的 URL 发起 HTTP 调用。
构建器方法(均为消费型链式调用):.description(desc)、.metadata(value)、.request_format(schema)、.response_format(schema)。其中new/new_async会自动通过schemars从参数与返回类型推导请求/响应 JSON Schema(类型约束T: Deserialize + JsonSchema、O: Serialize + JsonSchema),request_format/response_format可手动覆盖自动推导结果。
use iii_sdk::{register_worker, InitOptions, Error, RegisterFunction}; use serde::{Deserialize, Serialize}; use schemars::JsonSchema; #[derive(Deserialize, JsonSchema)] struct Input { name: String } #[derive(Serialize, JsonSchema)] struct Output { message: String } async fn greet(input: Input) -> Result<Output, Error> { Ok(Output { message: format!("Hello, {}!", input.name) }) } let worker = register_worker("ws://localhost:49134", InitOptions::default()); worker.register_function( "greetings::greet", RegisterFunction::new_async(greet).description("Greets a user"), );函数元数据与取消注册
register_function返回一个FunctionRef,可调用unregister()从引擎注销函数(内部会发送UnregisterFunction消息)。register_function要求函数 id 非空且不重复,否则直接 panic(见 src/iii.rs)。
自定义触发器类型(register_trigger_type)
SDK 还支持注册自定义触发器类型:通过RegisterTriggerType::new(id, description, handler)创建,handler 需实现TriggerHandlertrait;再通过.trigger_request_format::<T>()与.call_request_format::<T>()绑定配置与调用请求的类型,从而在TriggerTypeRef上获得编译期类型安全的register_function/register_trigger:
let my_trigger = worker.register_trigger_type( RegisterTriggerType::new("my-trigger", "My custom trigger", MyHandler) .trigger_request_format::<MyConfig>() .call_request_format::<MyRequest>(), ); // Compile-time safe: config must be MyConfig, function input must be MyRequest my_trigger.register_function("my::handler", |req: MyRequest| -> Result<serde_json::Value, iii_sdk::Error> { Ok(serde_json::json!({ "data": req.data })) }); my_trigger.register_trigger("my::handler", MyConfig { url: "/hook".into() });TriggerTypeRef::register_trigger_with_metadata还会默认把触发器命名空间设为当前 Worker 的命名空间——否则函数落在 Worker 命名空间、触发器却落在default,永远解析不到(见 src/iii.rs)。
注册触发器
READM 中的基础示例将 HTTP 触发器绑定到orders::create:
iii.register_trigger("http", "orders::create", json!({ "api_path": "/orders", "http_method": "POST" }))?;底层实现中,register_trigger接收RegisterTriggerInput(trigger_type、function_id、config为必填,另有可选metadata、namespace、trigger_namespace),内部自动生成 UUID 作为触发器 id,返回一个可调用unregister()的Trigger句柄。命名空间语义值得注意(见 src/iii.rs 与 src/protocol.rs):
namespace:目标函数在哪个命名空间解析。不填时默认继承 Worker 的命名空间(因为触发器指向的函数是当前 Worker 注册的,而函数落在 Worker 的命名空间);想绑定到其他命名空间(包括引擎的default)需显式声明。trigger_namespace:触发器类型的 provider在哪个命名空间查找。不填时引擎先查当前连接命名空间、再查引擎自身命名空间——这让尚未迁移到引擎 provider 的 Worker 无需声明即可继续工作。
对于引擎内置的触发器类型(如 HTTP、cron、queue),推荐使用IIITrigger(位于 src/builtin_triggers.rs),它知道类型 id 与配置形状。
调用函数:同步、Fire-and-forget 与队列异步
TriggerRequest的核心字段(见 src/protocol.rs):
function_id: String:要调用的函数 id。payload: Value:传给函数的输入数据。action: Option<TriggerAction>:路由方式。None表示同步请求/响应;Some(TriggerAction::Void)表示 fire-and-forget;Some(TriggerAction::Enqueue { queue })表示经命名队列异步处理。timeout_ms: Option<u64>:覆盖默认调用超时(默认 30 秒,源码常量DEFAULT_TIMEOUT_MS: u64 = 30_000,见 src/iii.rs)。
同步调用:等待结果
use iii_sdk::{TriggerRequest, TriggerAction}; use serde_json::json; // Synchronous -- waits for the result let result = iii.trigger(TriggerRequest { function_id: "orders::create".to_string(), payload: json!({ "body": { "item": "widget" } }), action: None, timeout_ms: None, }).await?;Fire-and-forget:只发不候
// Fire-and-forget iii.trigger(TriggerRequest { function_id: "analytics::track".to_string(), payload: json!({ "event": "page_view" }), action: Some(TriggerAction::Void), timeout_ms: None, }).await?;Void调用不生成invocation_id、不等待响应,SDK 立即返回Value::Null,适合日志、埋点等"发出即忘"场景。
队列异步:经命名队列路由
// Async via named queue iii.trigger(TriggerRequest { function_id: "orders::process".to_string(), payload: json!({ "order_id": "456" }), action: Some(TriggerAction::Enqueue { queue: "payments".to_string() }), timeout_ms: None, }).await?;Enqueue会把调用路由到指定队列(队列须在队列 Worker 的queue_configs中声明),由队列异步消费。TriggerAction枚举在 wire 协议上以type字段序列化为小写标签(enqueue/void,见 src/protocol.rs)。
附加元数据与目标命名空间
TriggerRequest提供两个非破坏性的扩展方法(不改变 struct 字面量的必需字段):
// 附加 per-invocation 元数据,handler 以独立参数接收 iii.trigger( TriggerRequest { function_id: "audit::write".to_string(), payload: json!({"event": "checkout"}), action: Some(TriggerAction::Void), timeout_ms: None, } .metadata(json!({"tenant": "acme"})), ).await?; // 指定本次调用的目标命名空间(不设则继承 Worker 命名空间;想从命名空间 Worker 打到引擎默认命名空间需显式写 "default") iii.trigger( TriggerRequest { function_id: "engine::some_fn".to_string(), payload: json!({}), action: None, timeout_ms: None, } .namespace("default"), ).await?;从实现看,trigger()对命名空间有一套精密的解析逻辑(invocation_namespace):显式指定的命名空间永远优先;未显式指定时,engine::前缀的引擎内置函数始终落在default,其余调用继承 Worker 命名空间(见 src/iii.rs)。
流数据操作:Stream Set 与原子更新
流(stream)是 III 中按stream_name+group_id+item_id组织的 KV 型数据结构,支持原子更新操作。
写入流条目
use iii_sdk::{register_worker, InitOptions, TriggerRequest, UpdateOp}; use serde_json::json; #[tokio::main] async fn main() -> Result<(), Box<dyn std::error::Error>> { let iii = register_worker("ws://localhost:49134", InitOptions::default()); // Set a stream item iii.trigger(TriggerRequest { function_id: "stream::set".into(), payload: json!({ "stream_name": "users", "group_id": "active", "item_id": "user-1", "data": { "status": "online" }, }), action: None, timeout_ms: None, }).await?; // Atomic update ops let ops = vec![ UpdateOp::increment("total", 100), UpdateOp::set("status", json!("processing")), ]; iii.trigger(TriggerRequest { function_id: "stream::update".into(), payload: json!({ "stream_name": "orders", "group_id": "user-123", "item_id": "order-456", "ops": ops, }), action: None, timeout_ms: None, }).await?; Ok(()) }UpdateOp 原子操作
UpdateOp定义于iii_helperscrate 的 sdk/packages/rust/helpers/src/stream.rs,是一组可以原子作用于流值的操作,按type字段序列化(小写标签),包括:
UpdateOp::set(path, value):在路径上覆盖写入值。UpdateOp::increment(path, amount):对路径上的数值做增量。UpdateOp::remove(path):删除路径上的值。
每个ops数组中的操作按顺序、原子地应用到同一流条目;UpdateOpError携带op_index指明出错的第几个操作,错误码如merge.path.too_deep用于定位深层路径合并问题。除stream::set/stream::update外,SDK 还支持stream::get、stream::delete、stream::list、stream::list_groups等引擎内置流函数,相关测试覆盖见 tests/stream.rs。
自定义流 Provider
如果需要把流的底层存储替换为自己的实现,SDK 提供create_stream(iii, stream_name, stream)辅助函数(见 src/helpers.rs):传入实现了IStreamtrait 的实例后,SDK 会自动在引擎上注册stream::get(名称)、stream::set(名称)、stream::delete(名称)、stream::list(名称)、stream::list_groups(名称)五个可调用函数,把它们接到你的实现上;注意update不会被注册,原子更新始终保留在引擎侧实现。
Logger 与可观测性
Logger 快速上手
SDK 提供基于iii_helpers::observability的Logger:
use iii_helpers::observability::Logger; let logger = Logger::new(Some("my-function".to_string())); logger.info("Processing started", None);Logger会发射 OpenTelemetryLogRecord;当 OTel 未初始化时,自动回退到tracingcrate 输出日志。也就是说,同一段日志代码在有无 OTel Collector 的环境中都能正常工作。
调用链跟踪(Trace)
SDK 与引擎通过 WebSocket 消息传递traceparent与baggage头,实现跨进程的分布式跟踪。处理端(handle_invoke_function)会:
- 从入站消息提取父级 trace 上下文,为每次调用创建一个名为
execute <function_id>的INTERNALspan(该命名特意与引擎发出的call/triggerspan 区分,避免重复),使 Worker handler 的 span 成为引擎 call span 的干净子节点。 - 以事件(event)形式记录调用输入输出:
iii.invocation.input与iii.invocation.output,并支持脱敏与截断(redact_and_truncate);可用环境变量III_DISABLE_TRACE_PAYLOADS=1关闭 payload 记录,payload 最大字节数也可通过环境变量调整。 - 根据结果设置 span 状态:成功
Ok,失败记录exception事件(含exception.type/exception.message/exception.stacktrace),错误信息会通过InvocationResult消息回传给调用方(见 src/iii.rs)。
命名空间与 Worker 身份
WorkerIdentityMode 两种模式
WorkerIdentityMode(见 src/iii.rs)决定 Worker 连接的引擎身份来源:
Managed(默认):采用监督者管理的III_WORKER_NAME与III_NAMESPACE环境变量(存在时覆盖 metadata 中对应字段)。适合iii compose、容器、systemd 等由外部编排注入身份的场景。Explicit:完全使用WorkerMetadata中的名称与显式选项/元数据中的命名空间,忽略进程级身份环境变量。适合一个 Worker 创建的辅助连接(auxiliary connections)。
命名空间解析顺序
Worker 有效命名空间的解析顺序是:InitOptions.namespace> 环境变量III_NAMESPACE>None(此时由引擎套用其default命名空间)。SDK 对"声明了但为空"的命名空间采取拒绝策略:reject_blank_namespace会直接 panic——因为"未设置"和"空白"含义相反,前者请求引擎默认命名空间,后者是想命名却拿不出名字,若被当作未设置,整个项目会悄悄在错误的命名空间运行(见 src/iii.rs)。
相关行为有专门的测试覆盖,例如 tests/namespace_inheritance.rs 验证命名空间继承语义,tests/env_contract.rs 验证环境变量契约。
注册冲突与错误处理
引擎在注册发生冲突时推送RegistrationRejected消息,SDK 按code区分严重程度(见 src/iii.rs 与 src/protocol.rs):
WORKER_NAMESPACE_CONFLICT:同命名空间下已有同名存活 Worker。致命错误——引擎关闭连接,SDK 置Failed状态、停止重连(避免再次撞进同一个冲突),并立即用RegistrationRejected失败所有在途调用。FUNCTION_NAMESPACE_CONFLICT:同命名空间下另一个 Worker 已导出该函数 id。非致命——只拒绝这一个函数注册,连接保持,其余函数继续服务,仅输出 warning 日志。- 未知 code:按致命错误处理(安全默认)。
同步调用的超时与连接错误会映射为Error::Timeout、Error::NotConnected等;远端函数失败返回Error::Remote(携带code、message、stacktrace)。相关 wire 行为测试见 tests/error_wire.rs。
小结
通过本指南,你已掌握 iii-sdk 的完整用法:安装依赖、初始化与优雅关闭、理解底层 WebSocket 自动重连与 Reattach 机制、注册同步/异步/HTTP 函数、绑定与注销触发器、以同步/发射即忘/队列三种模式调用函数、进行流数据原子操作,以及利用 Logger 与分布式跟踪实现可观测。SDK 源码中还有更多细节可供继续深入:
- 运行时与协议类型分组:
iii_sdk::runtime、iii_sdk::trigger、iii_sdk::channel、iii_sdk::errors、iii_sdk::protocol、iii_sdk::engine(见 src/lib.rs)。 - 内置触发器类型 src/builtin_triggers.rs、流 Provider trait src/stream_provider.rs、通道(Channel)实现 src/channels.rs。
- 集成测试覆盖了触发器动作(tests/trigger_action.rs)、重连(tests/reattach.rs)、注册去重(tests/registration_dedup.rs)、Pub/Sub(tests/pubsub.rs)与 HTTP 外部函数(tests/http_external_functions.rs)等场景,可作为你编写业务 Worker 时的参照。SDK 整体采用 Apache-2.0 协议(见 sdk/LICENSE)。
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考