跳到主要内容

接入引擎:SDK 与 worker 握手

30 秒导读: 前面几章讲了引擎里的注册表、消息分发、触发器。这一章换到 worker 一侧: 一个普通业务进程,如何用几行 SDK 代码开一条 WebSocket、连上引擎、领到一个会话身份, 再把自己的函数和触发器"挂"到引擎上,然后开始接活。这是所有协议真正被"插上电"的地方。

1. 这是什么(零基础也能懂)

一句话定义: worker = 一个跑你业务代码的进程;SDK = 让它"三行代码连上引擎并暴露能力"的客户端库。

在 iii 里,引擎(engine)自己不跑你的业务逻辑。真正的函数体活在 worker 进程里—— 可能是一段 Node 服务、一个 Python 脚本、一个 Rust 二进制。worker 靠 SDK 主动拨号连到引擎, 报到、注册能力,之后就像坐席一样等引擎派单。

用起来什么样。 一个最小的 Node worker 长这样:

// 示意,非源码(基于 sdk/packages/node/iii/src/iii.ts 的公开 API)
import { registerWorker } from 'iii-sdk'

// 1. 拨号:开一条 WebSocket 连到引擎(默认端口 49134)
const worker = registerWorker('ws://localhost:49134', { workerName: 'greeter' })

// 2. 暴露一个函数:给它一个全局唯一 id + 一个处理器
worker.registerFunction('hello::greet', async ({ name }) => {
return { message: `Hello, ${name}!` }
})
// 就这样。之后引擎收到对 hello::greet 的调用,就会把它派给这个进程。

注意:没有单独的 connect()registerWorker(...) 一返回,连接就已在后台建立 (registerWorker 内部直接 new Sdk(...) 并在构造器里 connect(),sdk/packages/node/iii/src/iii.ts:1097:156)。

一句话直觉: 把引擎当成一个电话总机,worker 是坐席。开机(WebSocket)→ 报工号(握手)→ 登记"我能接哪些业务"(注册函数/触发器)→ 之后总机把匹配的来电转给你。断线了自动重拨。

本章讲这条"拨号—报到—登记—接活"链路的两端怎么对上。不重复引擎内部怎么路由消息(见 02-engine-core-routing.md)、线上消息 schema 的字段细节(见 01-primitives-and-protocol.md)。

2. 顶层全景(它大概怎么转)

worker 和引擎之间只有一条 WebSocket,双向跑 JSON 消息。引擎这侧接客的不是什么特殊组件, 而是一个内置 worker——iii-worker-manager,它就是那个 WebSocket 监听器 (crate::register_worker!("iii-worker-manager", …, mandatory),engine/src/workers/worker/mod.rs:260)。

怎么读下面这张图: 左边是你的进程,右边是引擎进程;中间那条竖线是唯一的 WebSocket; 箭头是消息流向。

worker 进程(你的代码 + SDK) 引擎进程(iii)
┌───────────────────────────┐ ┌────────────────────────────────────┐
│ registerWorker(url) │ │ iii-worker-manager(内置 mandatory) │
│ └ new WebSocket ─────────┼──WS 升级──►│ ws_handler → engine.handle_worker │
│ │ │ ① handle_session(RBAC 认证) │
│ registerFunction(id, fn) ──┼──注册消息──►│ ② 登记进注册表(可加 prefix) │
│ registerTrigger(cfg) ──────┼──注册消息──►│ ③ 回发 WorkerRegistered │
│ │◄──worker_id─┼─ │
│ onInvokeFunction(fn) ◄─────┼──派单调用──┼─ 引擎路由把调用转过来 │
│ └ 回 invocationresult ───┼──结果─────►│ │
└───────────────────────────┘ └────────────────────────────────────┘

各部件一句话职责:

部件干什么在哪(文件:符号)
registerWorkerSDK 入口,建连接、返回客户端句柄sdk/.../iii.ts:1097 registerWorker
WorkerManager引擎侧的 WS 监听器(一个内置 worker)engine/src/workers/worker/mod.rs:84 WorkerManager
ws_handler把 HTTP 请求升级成 WebSocket,转交引擎engine/src/workers/worker/mod.rs:216 ws_handler
handle_worker握手 + 每连接的消息主循环engine/src/engine/mod.rs:1429 handle_worker
handle_session跑 RBAC 认证,产出一个 Sessionengine/src/workers/worker/rbac_session.rs:114 handle_session
EngineBuilder装配并启动引擎(含内置 + 外部 worker)engine/src/workers/config.rs:542 EngineBuilder

主线走一遍(高层):

  1. 引擎启动时,EngineBuilderiii-worker-manager 装配起来,它在 0.0.0.0:49134 开始监听。
  2. worker 进程调 registerWorker,SDK 向该端口发起 WebSocket 升级。
  3. 引擎 handle_worker 先做一次 RBAC 握手,得到一个 Session,再回发一条 WorkerRegistered 告知工号。
  4. SDK 在连接打开后,把已登记的触发器类型、函数、触发器一次性重放给引擎,并上报一次 worker 元数据。
  5. 之后进入稳态:引擎把匹配的调用(invokefunction)推给 worker,worker 执行完回 invocationresult

3. 核心原理(逐个机制,由浅入深)

3.1 拨号:一条 WebSocket 就是全部

worker 侧建连很朴素:构造器里 new WebSocket(this.address, { headers }),挂三个回调 (open/close/error),仅此而已(sdk/packages/node/iii/src/iii.ts:621 connect)。

引擎侧,iii-worker-manager 用 axum 起一个路由,把 / 这个路径交给 ws_handler:

// engine/src/workers/worker/mod.rs:127-131(节选)
let app = Router::new()
.route("/", get(ws_handler))
.route("/otel", get(otel_ws_handler))
.route("/ws/channels/{channel_id}", get(channel_ws_upgrade))
.with_state(state);

ws_handler 收到升级请求后,ws.on_upgrade(...) 把裸 socket 交给 engine.handle_worker(...) (engine/src/workers/worker/mod.rs:216-234)。监听端口默认 49134 (DEFAULT_PORT,engine/src/workers/worker/mod.rs:36),可由 WorkerManagerConfig.port 改 (engine/src/workers/worker/mod.rs:54)。

值得留意:遥测(OpenTelemetry)走的是另一条 /otel 连接,和业务 WS 分开,免得把注册表 塞进没有元数据的"幽灵 worker"(handle_otel 的注释,engine/src/engine/mod.rs:1552)。

3.2 报到:RBAC 握手产出一个 Session

连接升级后,handle_worker 做的第一件事就是握手认证。它调用 handle_session,后者跑一次 Session::authenticate(engine/src/workers/worker/rbac_session.rs:67):

  • 没有配 auth 函数 → 直接放行,给一个"全开"的默认 AuthResult (allow_trigger_type_registration: trueallow_function_registration: true, rbac_session.rs:76-86)。这是本地开发的常态。
  • 配了 auth_function_id → 引擎调用那个函数,把 headers/query_params/ip_address 作为输入喂进去(rbac_session.rs:88-97),函数返回一份 AuthResult 决定这条连接的权限。

握手成功后得到一个 Session,它就是这条连接的身份证,字段如下:

Session 字段作用落到协议的哪一步
allowed_functions / forbidden_functions允许/禁止调用的函数 id每次 invokefunction 前的准入判定(见 02)
allowed_trigger_types该连接可用的触发器类型触发器体系(见 04)
allow_trigger_type_registration是否允许注册新的触发器类型见 §3.4
allow_function_registration是否允许注册函数registerfunction 的准入
function_registration_prefix给该连接注册的函数 id 自动加命名空间前缀见 §3.3
contextauth 函数返回的自定义上下文(透传给中间件/hook)中间件与注册 hook

握手失败时,引擎回一条 {"type":"error", …} 再发 Close,连接就此结束 (engine/src/engine/mod.rs:1446-1456)。

紧接着引擎把 Session 绑进一个 WorkerConnection,登记进注册表,并回发一条 WorkerRegistered 把分配的 worker_id 告诉 worker(engine/src/engine/mod.rs:1480-1492)。SDK 收到后记下工号、 启动指标上报(onMessage 里的 WorkerRegistered 分支,sdk/packages/node/iii/src/iii.ts:1032)。

整条握手的时序(从 worker 视角):

worker engine.handle_worker
│ ── GET / (WebSocket Upgrade) ──► │
│ handle_session:
│ authenticate() 跑 auth 函数(或直接放行)
│ ◄── (若认证失败) {type:error} + Close │ 建 Session(allowed/prefix/context/…)
│ register_worker + 分配 worker_id
│ ◄──────── WorkerRegistered ──────────│ (带 worker_id)
│ ──────── registerfunction ─────────► │ ② 进入消息主循环:router_msg 逐条处理
│ ──────── registertrigger ──────────► │
│ ── engine::workers::register(void) ─►│ ③ 上报 worker 元数据
│ ◄──────── invokefunction ────────────│ 引擎路由派来的调用
│ ──────── invocationresult ─────────► │

3.3 命名空间:function_registration_prefix 怎么两头对上

它要解决的小问题: 多个租户/项目的 worker 连到同一个引擎,函数 id 会撞车。解法是给某条连接 注册的所有函数自动加前缀,但对 worker 自己透明——它注册和被调用时用的都是"短 id"。

这靠一个对称变换实现,注册时加前缀、调用时去前缀:

  • 注册时(加前缀): worker 发来 registerfunction { id: "greet" },引擎用 resolve_registration_id 把它变成 myproject::greet 再登记:

    // engine/src/engine/mod.rs:257-267 resolve_registration_id
    fn resolve_registration_id(worker: &WorkerConnection, id: &str) -> String {
    if let Some(prefix) = worker.session.as_ref()
    .and_then(|s| s.function_registration_prefix.as_ref())
    {
    format!("{prefix}::{id}") // greet → myproject::greet
    } else {
    id.to_string()
    }
    }

    (这正是 02 里注册路径调用的那个 resolve_registration_id; 本章只说明它由 worker 的 Session 前缀驱动,不复述注册表内部结构。)

  • 调用时(去前缀): 引擎要把对 myproject::greet 的调用推回 worker 时,先剥掉前缀, 让 worker 收到的仍是它认识的 greet(engine/src/worker_connections/traits.rs:105-117):

    // engine/src/worker_connections/traits.rs:106-114(节选)
    if let Some(prefix) = &session.function_registration_prefix {
    let needle = format!("{prefix}::");
    function_id.strip_prefix(&needle).map(String::from).unwrap_or(function_id)
    }

一句话:前缀只活在引擎的注册表里,worker 两头看到的都是短 id。

3.4 注册即协议:三个 register 调用对应三种线上消息

SDK 的 registerFunction / registerTrigger / registerTriggerType 不做别的,本质是把一条 协议消息发上线。三者一一对应 01 里的 Message 变体:

SDK 调用发出的 MessageType线上 typeschema(SDK 侧)
registerFunction(id, fn, opts)RegisterFunctionregisterfunctionRegisterFunctionMessage,iii-types.ts:84
registerTrigger(cfg)RegisterTriggerregistertriggerRegisterTriggerMessage,iii-types.ts:47
registerTriggerType(t, h)RegisterTriggerTyperegistertriggertypeRegisterTriggerTypeMessage,iii-types.ts:16

线上 type 的字符串值由 MessageType 枚举定死(全小写,sdk/packages/node/iii/src/iii-types.ts:3), 和引擎 Rust 侧 Message 的 serde 标签严格对齐——这是两端能互认的根基。

registerFunction 的骨架:它先做本地校验(id 非空、未重复),拼出 RegisterFunctionMessage, sendMessage(...) 发上线,再把 handler 存进本地 functions 表以备被调用 (sdk/packages/node/iii/src/iii.ts:292-402)。返回的 FunctionRef 带一个 unregister(), 它会发 unregisterfunction 并从本地表删掉。

触发器类型的注册还要过引擎的两道 RBAC 闸(worker 发 registertriggertype 后):

  1. 准入闸:session.allow_trigger_type_registration 为 false,引擎直接静默丢弃 (engine/src/engine/mod.rs:680-687)。
  2. 改写 hook: 若配了 on_trigger_type_registration_function_id,引擎调用该 hook, hook 可以改写 trigger_type_id/description,或否决注册 (返回非 object 即视为拒绝,engine/src/engine/mod.rs:689-718)。

这两个开关都定义在 RbacConfig(engine/src/workers/worker/rbac_config.rs:15),由握手时的 auth 函数通过 Session 传导下来。触发器类型本身的语义见 04

3.5 接活:worker 侧如何处理引擎推来的消息

稳态下,消息是双向的。上面几节是 worker→引擎;反过来,引擎会推两类消息给 worker, 都在 SDK 的 onMessage 分发(sdk/packages/node/iii/src/iii.ts:996-1038):

  • invokefunctiononInvokeFunction:按 function_id 找到本地 handler,执行, 把结果封成 invocationresult 回发;handler 抛错则回一条带 error 的结果 (iii.ts:871-930)。若消息没带 invocation_id,说明是 fire-and-forget,执行完不回结果。
  • registertriggeronRegisterTrigger:当这个 worker 注册过某个触发器类型时, 引擎会把"给这个类型绑定一个触发器"的请求推回来,worker 调用该类型的 registerTrigger 回调,再回一条 triggerregistrationresult(iii.ts:932-963)。

这就是"worker 既是被调用方、也是某些触发器类型的宿主"的双重身份的落地方式。

4. 深入实现:引擎怎么启动并装配 worker

前面都在讲"连上一个已经在跑的引擎"。这一节补上引擎自己是怎么起来、怎么把 iii-worker-manager(以及别的内置/外部 worker)装配上去的——这是 worker 有得可连的前提。

入口是 EngineBuilder,一个链式 builder(engine/src/workers/config.rs:542)。典型用法:

// engine/src/workers/config.rs:518-532 文档注释里的示例
EngineBuilder::new()
.config_file("config.yaml")? // 或 .default_config()
.build().await?
.serve().await?;

三层数据结构:

结构是什么关键方法
EngineConfig一份 YAML 配置:modules + workers 两张 worker 清单config.rs:30
WorkerRegistryworker 工厂表:name → 创建实例的闭包config.rs:258
EngineBuilder把配置 + 注册表 + 引擎捏在一起,建好并启动config.rs:542

① 内置 worker 从哪来:default_modules 与 inventory。 iii 用 inventory crate 做编译期收集: 每个内置 worker 用 register_worker! 宏声明(如 iii-worker-manager),启动时 default_worker_entries() 遍历所有标了 is_default 的登记项,自动放进 modules 清单 (config.rs:158-168)。WorkerRegistry::with_inventory() 则把每个登记项的工厂函数装进工厂表 (config.rs:269-280:375)。你不用手写"启用 worker-manager",它作为 mandatory 内置 worker 被自动装上。

② 内置守护进程兜底:ensure_builtin_daemons 有些能力(如 SDK 的 worker::* 触发器)靠一个 外部守护进程 iii-worker-ops 提供。ensure_builtin_daemons 会在能找到其二进制的前提下, 把它幂等地注入 workers 清单,免得用户还得在配置里手写一行(config.rs:131-155)。找不到二进制 就跳过——否则每台没装 iii-worker 的机器都会开机失败。

③ 编程式注册自定义 worker:register_worker / add_worker builder 上这两个方法配合使用: register_worker::<M>(name) 把一个实现了 Worker trait 的类型登记进工厂表;add_worker(name, cfg) 把一条 worker 条目加进配置清单(config.rs:614:620)。前者说"这个名字怎么造",后者说"造一个"。

build():逐个创建并初始化 worker。 buildworkersmodules 两张清单合并, 补齐所有 mandatory 内置 worker,给重名条目分配 #1/#2 实例号,然后逐个 create_workerinitializeregister_functions,并记录成 RunningWorker(config.rs:639-724)。其中 WorkerRegistry::create_worker 有一条四步解析链:内置工厂 → 传统外部 worker(iii.toml) → 委派给 iii-worker start(自动下载/拉起 OCI 镜像)(config.rs:306-369)。

serve():起后台任务并守住主循环。 serve 给每个 worker 起后台任务(iii-worker-manager 的后台任务就是那个 axum 监听器),接线全局 shutdown,监听配置文件变更做热重载,然后在 tokio::select! 循环里守着,直到收到 shutdown 或重载失败(config.rs:734-953)。

一句话串起来:inventory 收集内置 worker → EngineConfig 定清单 → WorkerRegistry 提供工厂 → EngineBuilder.build 把它们造出来 → serve 让 iii-worker-manager 开始监听 → worker 才连得上。

5. 巧妙之处(可借鉴的技术)

重放式重连:注册状态是幂等的、可回放的。 SDK 不维护"连接级"的复杂状态机。它把所有已注册的 触发器类型/函数/触发器都留在本地 Map 里;每次连接打开(含重连),onSocketOpen 就把它们 全部重新发一遍,再重发排队的消息(sdk/packages/node/iii/src/iii.ts:709-741)。这样"断线—重连— 恢复"不需要任何特殊的会话续期协议:重连 = 重新握手 + 重放全部注册。

指数退避 + 抖动的重连。 scheduleReconnectinitialDelayMs * backoffMultiplier ** attempt 算延迟,封顶到 maxDelayMs,再叠一层随机抖动,避免大批 worker 同时重连时的"惊群" (sdk/packages/node/iii/src/iii.ts:640-670)。maxRetries: -1 表示永久重试。

元数据上报走的是普通函数调用。 worker 报到时不是发个特殊的"hello"消息,而是调用一个引擎函数 engine::workers::register,用 fire-and-forget(void)动作发出去 (registerWorkerMetadata,sdk/packages/node/iii/src/iii.ts:526-550)。也就是说"注册元数据"这件事 本身也复用了统一的调用协议,没有额外的专用信道。

RBAC hook 复用调用协议。 auth 函数、on_trigger_type_registration hook——这些控制面动作, 引擎都用 self.call(fn_id, input) 发起,和数据面调用同一套机制。控制逻辑因此可以用任意语言的 worker 来实现,而不是写死在引擎里。

优雅关闭吞掉竞态错误。 shutdown 会拒绝所有 pending 调用、关掉 WS,并特意挂一个空的 error 监听器,吞掉"连接还在 CONNECTING 就 close"这类竞态异常(sdk/packages/node/iii/src/iii.ts:573-611)。

6. 边界与局限

  • worker 永远是发起方。 引擎不会主动去连 worker;是 worker 拨号进来。防火墙/NAT 后的 worker 也能连出,但引擎必须有一个可达的监听地址。
  • 握手失败即断连,不重试认证。 auth 函数报错或返回空,引擎回 error + Close,连接直接结束 (engine/src/engine/mod.rs:1446-1456);是否重连由 SDK 的重连策略决定,但认证仍会失败—— 换句话说,权限问题不会靠重连自愈。
  • 函数 id 本地去重。 同一个 worker 内重复 registerFunction 同一个 id 会抛错 (sdk/packages/node/iii/src/iii.ts:300);跨 worker 的 id 冲突则由引擎的所有权/前缀机制处理。
  • 前缀是纯字符串拼接。 resolve_registration_idformat!("{prefix}::{id}"),去前缀用 strip_prefix;它不校验 id 里是否已含 ::,靠约定而非强制。
  • /otel 是二进制专用信道。 文本帧发到 /otel 会被丢弃,那条连接只收遥测二进制帧 (engine/src/engine/mod.rs:1598 附近)。

7. 横向对比:一份协议,多语言绑定

sdk/packages/ 下有四套 SDK:nodepythonrustgo。它们不是四套协议,而是同一线上协议的 四种语言绑定——发的都是 {"type":"registerfunction", …} 这样的 JSON。

证据:Go SDK 的协议测试直接断言线上字节 {"type":"registerfunction","id":"hello::greet",…}(sdk/packages/go/iii/protocol_test.go:113), Python SDK 的测试同样按 m.get("type") == "registerfunction" 断言 (sdk/packages/python/iii/tests/test_register_function_args.py:83),四套 SDK 的公开入口都叫 register_worker / registerWorker(如 sdk/packages/python/iii-example/src/main.py:211sdk/packages/rust/iii-example/src/http_example.rs)。

结论:本章以 Node SDK 为主线讲的握手与注册流程,对其余三种语言等价成立;差异只在语言习惯 (命名、类型、async 风格),协议线格式完全一致。 想加第五种语言?实现同一套 MessageType 的 序列化/反序列化即可,引擎无需改动。

同 shelf 的其它 agent-runtime 项目多把 worker 逻辑和调度耦合在一个进程里;iii 的取舍是把引擎和 worker 用一条 WebSocket 彻底解耦,代价是多一跳网络,收益是 worker 可用任意语言、可独立部署与重启。

8. 代码地图(导航索引)

主题文件符号名
SDK 入口(建连)sdk/packages/node/iii/src/iii.tsregisterWorker / Sdk.connect
连接打开时重放注册sdk/packages/node/iii/src/iii.tsonSocketOpen
注册函数sdk/packages/node/iii/src/iii.tsregisterFunction
注册触发器 / 触发器类型sdk/packages/node/iii/src/iii.tsregisterTrigger / registerTriggerType
worker 元数据上报sdk/packages/node/iii/src/iii.tsregisterWorkerMetadata
处理引擎推来的调用sdk/packages/node/iii/src/iii.tsonInvokeFunction / onRegisterTrigger
重连 / 优雅关闭sdk/packages/node/iii/src/iii.tsscheduleReconnect / shutdown
线上消息类型枚举 + schemasdk/packages/node/iii/src/iii-types.tsMessageType / IIIMessage
协议导出面sdk/packages/node/iii/src/protocol.ts · index.ts(re-exports)
WS 监听器(内置 worker)engine/src/workers/worker/mod.rsWorkerManager / ws_handler / DEFAULT_PORT
握手主循环engine/src/engine/mod.rshandle_worker
RBAC 会话与认证engine/src/workers/worker/rbac_session.rshandle_session / Session / authenticate
RBAC 配置(前缀/hook/开关)engine/src/workers/worker/rbac_config.rsRbacConfig / is_function_allowed
注册 id 加/去前缀engine/src/engine/mod.rs · worker_connections/traits.rsresolve_registration_id
触发器类型注册的 RBAC 闸engine/src/engine/mod.rsrouter_msg(RegisterTriggerType 分支)
引擎装配与启动engine/src/workers/config.rsEngineConfig / WorkerRegistry / EngineBuilder
内置守护进程注入engine/src/workers/config.rsensure_builtin_daemons
多语言 SDK 协议一致性sdk/packages/go/iii/protocol_test.goTestRegisterFunctionMessage(线上字节断言)