使用 iii SDK 构建跨语言 Worker:Node、Python、Rust、Go 统一开发指南
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
iii 是一个实时编排引擎,围绕 Function、Trigger、Worker 三个核心原语提供跨语言执行、实时能力发现与可观测性。官方 SDK 将引擎能力封装为四套 API 完全对齐的客户端库(Node.js、Python、Rust、Go),让开发者可以用自己熟悉的语言编写 worker,通过 WebSocket 与引擎通信,注册函数与触发器并响应调用。本文以 sdk/README.md 为骨架,结合 sdk/packages 下各语言 SDK 的 README 与源码实现,系统讲解 SDK 的安装、Hello World、统一 API 面、触发器绑定、调用方式、连接重连机制以及开发测试流程。读完本文,你将能在任意一门受支持语言中独立完成"注册函数 -> 绑定触发器 -> 同步/异步调用 -> 接入可观测性"的完整开发闭环。
一、SDK 全景:一套 API 面,四种语言实现
iii SDK 是引擎的官方客户端集合,分布在仓库 sdk/packages 目录下。每个包都以iii-sdk为名发布到各自语言的生态仓库,但共享完全相同的编程模型:
| 包名 | 语言 | 安装方式 | 详细文档 |
|---|---|---|---|
iii-sdk | Node.js / TypeScript | pnpm add iii-sdk或npm install iii-sdk | README |
iii-sdk | Python | pip install iii-sdk | README |
iii-sdk | Rust | 添加依赖到Cargo.toml | README |
iii-sdk | Go | go get github.com/iii-hq/iii/sdk/packages/go/iii | README |
SDK 采用 Apache 2.0 许可(见 sdk/LICENSE)。从 engine/README.md 可以看到引擎的架构定位:引擎提供持久化编排与跨语言执行能力,一个进程通过打开一条 WebSocket 连接到引擎即成为一个 iiiworker,之后注册函数和触发器;引擎按需调用这些函数,worker 在同一条 socket 上回复结果(见 Go SDK README 的说明)。
端口速查(来自 engine/README.md):49134 为 WebSocket worker 连接端口,3111 为 HTTP API,3112 为 Stream API,9464 为 Prometheus metrics。
二、环境准备:先启动引擎,再安装 SDK
2.1 启动引擎
所有语言 SDK 都默认连接ws://localhost:49134。在本地运行 SDK 之前,需要先启动 iii 引擎:
# 安装引擎(含全部 CLI 命令) curl -fsSL https://install.iii.dev/iii/main/install.sh | sh # 验证安装 command -v iii && iii --version # 使用内置默认配置启动(自带 in-memory OpenTelemetry 配置) iii --use-default-config也可以创建config.yaml后以项目配置启动:iii --config /path/to/config.yaml。启动成功后可用iii console打开控制台。引擎默认 WebSocket 地址即为ws://localhost:49134。
2.2 安装各语言 SDK
# Node.js / TypeScript pnpm add iii-sdk # 或 npm install iii-sdk # Python pip install iii-sdk # Rust(Cargo.toml) [dependencies] iii-sdk = "0.11" serde_json = "1" tokio = { version = "1", features = ["full"] } # Go(要求 Go 1.24+) go get github.com/iii-hq/iii/sdk/packages/go/iii三、四语言 Hello World 全解
3.1 Node.js / TypeScript
import { registerWorker } from 'iii-sdk'; const iii = registerWorker('ws://localhost:49134'); iii.registerFunction('hello::greet', async (input) => { return { message: `Hello, ${input.name}!` }; }); iii.registerTrigger({ type: 'http', function_id: 'hello::greet', config: { api_path: '/greet', http_method: 'POST' }, }); const result = await iii.trigger({ function_id: 'hello::greet', payload: { name: 'world' } });3.2 Python
from iii import register_worker iii = register_worker("ws://localhost:49134") def greet(data): return {"message": f"Hello, {data['name']}!"} iii.register_function({"id": "hello::greet"}, greet) iii.register_trigger({ "type": "http", "function_id": "hello::greet", "config": {"api_path": "/greet", "http_method": "POST"} }) result = iii.trigger({"function_id": "hello::greet", "payload": {"name": "world"}})注意:Python 的register_function第一个参数是包含"id"键的字典,也可以直接传字符串"hello::greet"(见 Python SDK README 的写法)。
3.3 Rust
use iii_sdk::{register_worker, InitOptions, TriggerRequest, RegisterFunctionMessage, RegisterTriggerInput}; use serde_json::json; #[tokio::main] async fn main() -> Result<(), Box<dyn std::error::Error>> { let iii = register_worker("ws://127.0.0.1:49134", InitOptions::default())?; iii.register_function(RegisterFunctionMessage::with_id("hello::greet".into()), |input| 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(RegisterTriggerInput::new("http", "hello::greet", json!({ "api_path": "/greet", "http_method": "POST" })))?; let result: serde_json::Value = iii .trigger(TriggerRequest::new("hello::greet", json!({ "name": "world" }))) .await?; Ok(()) }RegisterFunctionMessage::with_id用于按 ID 注册函数,TriggerRequest::new组合了 function_id 与 payload,RegisterTriggerInput::new(type, fn_id, config)绑定触发器。InitOptions::default()负责默认的连接初始化参数。
3.4 Go
package main import ( "context" "encoding/json" "log" iii "github.com/iii-hq/iii/sdk/packages/go/iii" ) func main() { client := iii.RegisterWorker("ws://127.0.0.1:49134") client.RegisterFunction("hello::greet", func(ctx context.Context, data json.RawMessage) (any, error) { var req struct { Body struct { Name string `json:"name"` } `json:"body"` } _ = json.Unmarshal(data, &req) return map[string]any{ "status_code": 200, "body": map[string]string{"message": "Hello, " + req.Body.Name + "!"}, }, nil }) client.RegisterTrigger("hello-http", "http", "hello::greet", json.RawMessage(`{"api_path":"/greet","http_method":"POST"}`), nil) if err := client.Connect(context.Background()); err != nil { log.Fatal(err) } defer client.Close() }3.5 理解 HTTP 触发器信封(Envelope)
观察以上代码可以发现一个细节:Go 的 handler 返回值带有status_code与body两个字段,而 Node/Python/Rust 示例返回的是纯数据对象。原因在于 Go SDK README 明确说明的HTTP trigger envelope约定:
当一个函数经由
http触发器被调用时,引擎会把请求包装为{ path, method, body, headers, … },并期望函数返回{ status_code, body }。仅通过 socket 调用的函数可以使用任意 payload 形状。
也就是说,凡是绑定 HTTP 触发器的函数,入参和返回值都必须遵守 HTTP 信封格式;而通过 SDK 内部trigger()直接调用的函数则不受此约束。这是编写跨语言 worker 时最容易踩坑、也最需要统一的契约点。
四、统一 API 面:四语言操作对照
sdk/README.md 提供了一张核心操作对照表,四个 SDK 暴露的 API 面完全一致:注册函数、注册触发器、发起调用。下表完整继承原文档并补充各语言的实际签名:
| 操作 | Node.js | Python | Rust | Go | 说明 |
|---|---|---|---|---|---|
| 初始化 | registerWorker(url) | register_worker(url, options?) | register_worker(url, options) | iii.RegisterWorker(url) | 创建 SDK 实例并自动连接 |
| 注册函数 | iii.registerFunction(id, handler, options?) | iii.register_function(id, handler) | iii.register_function(id, \|input\| ...) | client.RegisterFunction(id, handler) | 注册可按名字调用的函数 |
| 注册触发器 | iii.registerTrigger({ type, function_id, config }) | iii.register_trigger({"type": ..., "function_id": ..., "config": ...}) | iii.register_trigger(type, fn_id, config)? | client.RegisterTrigger(id, type, fn, cfg, meta) | 将触发器(HTTP、cron、queue 等)绑定到函数 |
| 同步调用 | await iii.trigger({ function_id, payload }) | await iii.trigger({"function_id": id, "payload": data}) | iii.trigger(TriggerRequest::new(id, data)).await? | client.Trigger(ctx, iii.TriggerRequest{...}) | 调用函数并等待结果 |
| 异步调用(fire-and-forget) | iii.trigger({ function_id, payload, action: TriggerAction.Void() }) | iii.trigger({"function_id": id, ..., "action": TriggerAction.Void()}) | iii.trigger(TriggerRequest { action: Some(TriggerAction::Void), ... }) | client.Trigger(ctx, iii.TriggerRequest{Action: iii.VoidAction()}) | 不等待结果的调用 |
| 队列调用(enqueue) | iii.trigger({ function_id, payload, action: TriggerAction.Enqueue({ queue }) }) | iii.trigger({"function_id": id, ..., "action": TriggerAction.Enqueue(queue="name")}) | iii.trigger(TriggerRequest { action: Some(TriggerAction::Enqueue { queue }), ... }) | client.Trigger(ctx, iii.TriggerRequest{Action: iii.EnqueueAction("queue")}) | 通过命名队列路由调用 |
原文档中有一个重要的兼容性说明,写作本文时仍须强调:
call、callVoid、triggerVoid(以及 Python/Rust 的对应变体)已被移除,所有调用统一使用trigger()。需要 fire-and-forget 时,使用trigger({ function_id, payload, action: TriggerAction.Void() })。
registerWorker()/register_worker()创建 SDK 实例并自动连接引擎,内部处理 WebSocket 通信、自动重连与 OpenTelemetry 埋点。Node SDK 的导出集中在 sdk/packages/node/iii/src/index.ts:对外暴露registerWorker、InitOptions、TriggerAction、InvocationError等类型与函数。
4.1 注册函数:从简单到类型安全
Node 侧支持可选的 options 参数,Python/Rust 直接传入 id 与 handler。以"创建订单"为例:
// Node.js iii.registerFunction('orders::create', async (input) => { return { status_code: 201, body: { id: '123', item: input.body.item } } })# Python def create_order(data): return {"status_code": 201, "body": {"id": "123", "item": data["body"]["item"]}} iii.register_function("orders::create", create_order)// Rust 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 } })) });Go SDK 更进一步,提供RegisterFunctionTyped[Req, Resp]:通过类型参数推断请求/响应 JSON Schema,并作为request_format/response_format广播给引擎与仪表盘(对应 Rust SDK 的#[derive(JsonSchema)])。Go 没有编译期 derive,因此基于反射推断(依赖invopop/jsonschema),可用json与jsonschema结构体标签定制 schema:
type CreateOrderRequest struct { Item string `json:"item" jsonschema:"required"` Quantity int `json:"quantity" jsonschema:"minimum=1"` } type OrderResult struct { ID string `json:"id"` } iii.RegisterFunctionTypedCreateOrderRequest, OrderResult (OrderResult, error) { return OrderResult{ID: "ord_123"}, nil })4.2 注册触发器:把函数暴露给外部事件
触发器把函数与外部事件源(HTTP、cron、queue 等)绑定起来。四语言用法一致:
// Node.js iii.registerTrigger({ type: 'http', function_id: 'orders::create', config: { api_path: '/orders', http_method: 'POST' }, })# Python iii.register_trigger({ "type": "http", "function_id": "orders::create", "config": {"api_path": "/orders", "http_method": "POST"}, })// Rust iii.register_trigger("http", "orders::create", json!({ "api_path": "/orders", "http_method": "POST" }))?;// Go(注意 Go 的签名多一个触发器 ID 参数,并可选携带 metadata) client.RegisterTrigger("orders-http", "http", "orders::create", json.RawMessage(`{"api_path":"/orders","http_method":"POST"}`), nil)Go SDK 还支持RegisterTriggerType(id, description, handler)实现自定义触发器类型,这与 Node SDK 中registerTriggerType的能力对齐——四个 SDK 都允许把引擎的触发器能力延伸到自定义事件源。
4.3 调用函数:三种调用语义
统一调用入口trigger()支持三种语义:
// Node.js import { registerWorker, TriggerAction } from 'iii-sdk' const iii = registerWorker('ws://localhost:49134') // 1. 同步等待结果 const result = await iii.trigger({ function_id: 'orders::create', payload: { item: 'widget' } }) // 2. fire-and-forget(如埋点、日志) iii.trigger({ function_id: 'analytics::track', payload: { event: 'page_view' }, action: TriggerAction.Void() }) // 3. 通过命名队列异步路由 iii.trigger({ function_id: 'orders::process', payload: { order_id: '456' }, action: TriggerAction.Enqueue({ queue: 'payments' }) })Python 对应写法为iii.trigger({"function_id": ..., "action": TriggerAction.Enqueue(queue="name")});Rust 通过TriggerRequest { action: Some(TriggerAction::Enqueue { queue: "payments".to_string() }), .. }表达;Go 使用iii.EnqueueAction("jobs")。队列语义由引擎侧 queue 模块落地,SDK 只负责把调用请求路由进命名队列。
五、连接、重连与生命周期(源码视角)
以 Node SDK 为例,sdk/packages/node/iii/src/iii-constants.ts 定义了完整连接参数,这些默认值体现了 SDK 的健壮性设计:
| 常量 | 默认值 | 说明 |
|---|---|---|
DEFAULT_BRIDGE_RECONNECTION_CONFIG | 初始 1000ms、上限 30000ms、倍率 2、抖动 0.3、无限重试 | 指数退避 + 随机抖动 |
DEFAULT_INVOCATION_TIMEOUT_MS | 30000 | 默认调用超时 |
WS_HANDSHAKE_TIMEOUT_MS | 10000 | WebSocket 握手超时(与 Rust SDK 的 connect_timeout 对齐) |
WS_PING_INTERVAL_MS | 20000 | 保活 ping 间隔(与 Rust SDK 的 ping_interval 对齐) |
WS_IDLE_TIMEOUT_MS | 60000 | 长时间无入站帧强制重连(与 Rust SDK 的 idle_timeout 对齐) |
从 sdk/packages/node/iii/src/iii.ts 源码还可以看到两个关键设计:
- 地址解析优先级:显式传入的地址 > 环境变量
III_URL(由iii compose、容器运行时、systemd 等 supervisor 注入)> 默认值ws://127.0.0.1:49134。源码特意使用 IPv4 回环地址而非localhost,避免主机只监听 IPv4 时localhost解析到::1导致连不上。 - worker 身份与环境注入:
III_WORKER_NAME携带编排器分配的 worker 名(Compose 托管 worker 由 supervisor 注入),其优先级高于代码内嵌的显式名称,因为引擎按名字匹配在线注册;默认名形如${os.hostname()}:${process.pid}。
Go SDK 的连接行为与 Node 完全对齐(Go SDK README):指数退避重连(起始 1s、×2、封顶 30s、±30% 抖动、无限重试,可用iii.WithReconnectConfig覆盖);离线缓冲——断线期间发出的调用被缓冲并在重连后冲刷,注册信息则从内存注册表重放;每次连接最后注册 worker metadata(标记runtime: "go")。
Register*类调用可以在Connect之前或之后执行:注册保存在内存中,并在每次(重)连接时(重新)发送给引擎。这一点在 Node/Python/Rust/Go 四个 SDK 中行为一致,是断线重连后 worker 功能不丢失的基石。
六、进阶能力:数据通道、流操作与可观测性
sdk/README.md 明确指出:语言相关的进阶细节(modules、streams、OpenTelemetry)以各语言 SDK 的 README 为准。这里汇总仓库内可验证的四个进阶能力。
6.1 流式数据通道(Channels)
Go SDK 提供CreateChannel(ctx, bufferSize)打开双向数据通道(writer + reader 两端,各自独立 WebSocket)。通道上的约定是:文本帧即消息(SendMessage/OnMessage),二进制帧即流数据(Write/ReadAll),按 WebSocket opcode 区分,无需额外信封;writer 以Close()结束流,reader 将其视为 EOF。通道引用(WriterRef/ReaderRef)可以放进 trigger payload 传给另一个 worker,由对方打开对端——适合大体积或持续流式数据的传输,而非把所有数据塞进单次调用结果。Node/Python/Rust SDK 同样内置 channel 与 stream 模块(见 node/iii/src/channels.ts、rust/iii/src/channels.rs 等)。
6.2 流(Streams)与原子更新
Rust SDK README 展示了通过内置函数stream::set/stream::update操作流:以stream_name+group_id+item_id定位数据项,配合UpdateOp::increment("total", 100)、UpdateOp::set("status", json!("processing"))实现原子更新操作。这套能力在四个 SDK 中均有等价实现,用于跨 worker 共享实时状态。
6.3 可观测性:OpenTelemetry 一体化
SDK 内置 OpenTelemetry 埋点:registerWorker()自动完成 OTel 初始化,invocation 以 span 形式记录(Node 源码中通过withSpan、recordSpanEvent、injectTraceparent、injectBaggage等工具实现 W3C trace context 的注入与透传)。Rust 侧Logger结构体(iii_helpers::observability::Logger)发射 OTelLogRecord,未初始化 OTel 时回退到tracingcrate。引擎以--use-default-config启动时自带 in-memory OTel 配置,因此本地开发即可直接观察 traces、metrics、logs。
6.4 调用元数据与错误语义
- 元数据侧车:Go handler 通过
iii.MetadataFromContext(ctx)读取每次调用可选的元数据(如{"tenant":"acme"}),注册函数时也可用RegisterFunctionOptions{Metadata: ...}附加静态元数据。元数据随调用上下文传递,不改变 handler 签名。 - 错误类型:Go 中
errors.Is(err, iii.ErrTimeout)判断超时、errors.Is(err, iii.ErrNotConnected)判断关闭引起的取消、errors.As(err, &ie)拿到携带远端Code/Message/Stacktrace的*iii.InvocationError;handler 返回任意非InvocationError错误会被统一报告为invocation_failed,panic 也会被 recover 并同样上报,保证调用方得到错误而不是超时。Node SDK 对应导出InvocationError、RegistrationRejectedError(index.ts)。
七、开发与测试:从源码构建各语言 SDK
仓库内每个 SDK 包都自带构建与测试管线(sdk/packages/node/iii/package.json、sdk/packages/python/iii/pyproject.toml、sdk/packages/rust/iii/Cargo.toml)。
7.1 前提条件
- Node.js 20+ 与 pnpm(Node SDK)
- Python 3.10+ 与 uv(Python SDK)
- Rust 1.85+ 与 Cargo(Rust SDK)
- 运行在
ws://localhost:49134的 iii 引擎
7.2 构建
cd packages/node && pnpm install && pnpm build cd packages/python/iii && python -m build cd packages/rust/iii && cargo build --releasePython SDK 的开发模式安装、类型检查与 lint:
pip install -e . mypy src ruff check src7.3 测试
cd packages/node && pnpm test cd packages/python/iii && pytest cd packages/rust/iii && cargo test各 SDK 的测试目录覆盖了丰富的场景,是理解 API 语义的最佳教材:Node 侧 sdk/packages/node/iii/tests 包含连接握手超时、心跳、重连(connection-reattach.test.ts)、触发器动作(trigger-action.test.ts)、触发器注册错误(trigger-registration-error.test.ts)、流(stream.test.ts)、RBAC(rbac-workers.test.ts)等用例;Python 侧 sdk/packages/python/iii/tests 覆盖同步/异步 API、重连发送 reattach、环境契约、引擎常量等;Rust 侧 sdk/packages/rust/iii/tests 提供 mock engine 与集成测试;Go 侧 sdk/packages/go/iii/tests 覆盖触发器注册、数据通道、worker metadata、注册去重等。阅读这些测试可以帮助你理解 SDK 在异常与边界情况下的真实行为。
八、示例与下一步
仓库为每种语言都提供了可运行的示例工程:
- Node:
iii-example(含 HTTP、队列、DLQ、流、状态、中间件示例)——sdk/packages/node/iii-example - Python:
iii-example(函数、状态、流、触发器类型示例)——sdk/packages/python/iii-example - Rust:
iii-example(HTTP、cron、自定义触发器、Logger 示例)——sdk/packages/rust/iii-example - Go:
iii-example——sdk/packages/go/iii-example
仓库中还提供了 helpers 包(sdk/packages/node/helpers、sdk/packages/python/helpers、sdk/packages/rust/helpers),封装 HTTP、队列、流、worker 连接管理与可观测性工具,供 SDK 内部复用并在示例中展示。引擎架构细节可参阅 engine/README.md,SDK 与引擎的 WebSocket 协议交互可从 engine/src/protocol.rs 与各 SDK 的protocol.*源文件对照理解。
总结
iii SDK 的价值在于用一条 WebSocket 把任意语言进程接入统一编排引擎:四个语言包共享相同的 API 面(registerWorker->registerFunction->registerTrigger->trigger),统一处理连接、重连、离线缓冲、OpenTelemetry 埋点与错误语义;在此之上,channels 提供流式数据传输、streams 提供共享实时状态、队列触发器提供异步路由。无论团队技术栈是 TypeScript、Python、Rust 还是 Go,都可以用同样的心智模型开发 worker,并让不同语言的服务在同一引擎内相互调用、统一观测。
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考