
使用 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 SDKcrate 名为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.9edition 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 Schemareqwest用于 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(), Boxdyn 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的闭包返回ResultValue, Error异步签名async move由 SDK 自动包装输入输出均为serde_json::Value。register_trigger(http, hello::greet, config)把 HTTP 触发器绑定到函数配置里的api_path与http_method决定外部如何触发。同步调用返回Value即函数执行结果。核心 API 一览README 用一张表总结了 SDK 的全部核心操作这里完整保留并补充说明OperationSignatureDescriptionInitializeregister_worker(address, options)Create an SDK instance and auto-connectRegister functioniii.register_function(id, \|input: Value\| ...)Register a function that can be invoked by nameRegister triggeriii.register_trigger(type, fn_id, config)?Bind a trigger (HTTP, cron, queue, etc.) to a functionInvoke (await)iii.trigger(TriggerRequest { ... }).await?Invoke a function and wait for the resultInvoke (fire-and-forget)iii.trigger(TriggerRequest { action: Some(TriggerAction::Void), ... }).await?Fire-and-forget invocationInvoke (enqueue)iii.trigger(TriggerRequest { action: Some(TriggerAction::Enqueue { queue }), ... }).await?Route invocation through a named queueregister_worker()会在一个独立的后台线程中建立并维护 WebSocket 通信该线程自带一个 tokio 运行时同时负责自动重连与 OpenTelemetry 埋点见 src/lib.rs 与 src/iii.rs。也就是说SDK 的自动重连不依赖调用方运行时Worker 主线程可以专注于业务逻辑。初始化与连接生命周期register_worker 与 InitOptionsregister_worker(address, options)是推荐的入口它会按InitOptions配置好客户端后自动调用connect()。InitOptions的可配置字段如下见 src/lib.rsmetadata: OptionWorkerMetadata自定义 Worker 元数据运行时、版本、名称、描述、PID、隔离模式等。默认自动探测 hostname、PID、操作系统与项目名Managed 身份模式下进程级环境变量会覆盖其中对应字段。headers: OptionHashMapString, StringWebSocket 握手时携带的自定义 HTTP 头可用于认证等场景。otel: OptionOtelConfigOpenTelemetry 配置。namespace: OptionStringWorker 所属命名空间作用范围远超注册本身——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()返回IIIConnectionStateDisconnected/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 用于证明我们就是那个 WorkerWorker 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_functionuse 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 JsonSchemarequest_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) - ResultOutput, 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_typeSDK 还支持注册自定义触发器类型通过RegisterTriggerType::new(id, description, handler)创建handler 需实现TriggerHandlertrait再通过.trigger_request_format::T()与.call_request_format::T()绑定配置与调用请求的类型从而在TriggerTypeRef上获得编译期类型安全的register_function/register_triggerlet 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| - Resultserde_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::createiii.register_trigger(http, orders::create, json!({ api_path: /orders, http_method: POST }))?;底层实现中register_trigger接收RegisterTriggerInputtrigger_type、function_id、config为必填另有可选metadata、namespace、trigger_namespace内部自动生成 UUID 作为触发器 id返回一个可调用unregister()的Trigger句柄。命名空间语义值得注意见 src/iii.rs 与 src/protocol.rsnamespace目标函数在哪个命名空间解析。不填时默认继承 Worker 的命名空间因为触发器指向的函数是当前 Worker 注册的而函数落在 Worker 的命名空间想绑定到其他命名空间包括引擎的default需显式声明。trigger_namespace触发器类型的 provider在哪个命名空间查找。不填时引擎先查当前连接命名空间、再查引擎自身命名空间——这让尚未迁移到引擎 provider 的 Worker 无需声明即可继续工作。对于引擎内置的触发器类型如 HTTP、cron、queue推荐使用IIITrigger位于 src/builtin_triggers.rs它知道类型 id 与配置形状。调用函数同步、Fire-and-forget 与队列异步TriggerRequest的核心字段见 src/protocol.rsfunction_id: String要调用的函数 id。payload: Value传给函数的输入数据。action: OptionTriggerAction路由方式。None表示同步请求/响应Some(TriggerAction::Void)表示 fire-and-forgetSome(TriggerAction::Enqueue { queue })表示经命名队列异步处理。timeout_ms: Optionu64覆盖默认调用超时默认 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_namegroup_iditem_id组织的 KV 型数据结构支持原子更新操作。写入流条目use iii_sdk::{register_worker, InitOptions, TriggerRequest, UpdateOp}; use serde_json::json; #[tokio::main] async fn main() - Result(), Boxdyn 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的Loggeruse iii_helpers::observability::Logger; let logger Logger::new(Some(my-function.to_string())); logger.info(Processing started, None);Logger会发射 OpenTelemetryLogRecord当 OTel 未初始化时自动回退到tracingcrate 输出日志。也就是说同一段日志代码在有无 OTel Collector 的环境中都能正常工作。调用链跟踪TraceSDK 与引擎通过 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_PAYLOADS1关闭 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_NAMESPACENone此时由引擎套用其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.rsWORKER_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/Subtests/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),仅供参考