iii (原 Motia) — 架构与原理
30 秒导读: iii 是一个用 Rust 写的后端运行时。你把任何进程(TypeScript API、Python 数据管道、Rust 微服务,甚至一个 AI agent)通过一条 WebSocket 连上它的引擎,注册几个"有名字的函数"和"什么时候该跑"的触发器。引擎当一张活注册表 + 消息路由器:此后每个进程都能发现别人注册了什么、直接调用它、并追踪整条链路——不需要两两之间写点对点集成。整套心智模型只有三个词:Worker · Function · Trigger。
1. 这是什么(零基础也能懂)
一句话定义。 iii 把"一个后端服务能做的每件事"归约成三种原语,交给一个中央 Rust 引擎去注册、路由和观测。
它解决谁的什么问题。 想象你有一堆服务:一个下单 API、一个发邮件的、一个跑定时任务的、一个 agent。传统做法是它们两两对接——A 要调 B 得知道 B 的地址、协议、重试、鉴权,换一个可观测性工具又要重新接一遍。iii 把这些整合进一个活的系统面:任何能力注册进来,就自动可被发现、可调用、可追踪。
三个原语(整个心智模型):
| 原语 | 白话 | 例子 |
|---|---|---|
| Worker(工作者) | 一个连上引擎的进程 | 一个 TS 服务、一个 Python 管道、一个 agent |
| Function(函数) | 有稳定 id 的一段工作 | content::classify、orders::validate |
| Trigger(触发器) | 让某个函数运行的"什么事" | 直接调用、HTTP 请求、cron、队列消息、state 变更、stream 事件 |
它能做什么。 注册可调用函数;声明式绑定各种触发器;跨语言互相调用并等结果;运行时动态加入新能力(iii worker add <anything>);全链路分布式追踪。触发器是声明式的——worker 只说"这个函数在这件事发生时跑",路由、序列化、投递都归引擎(README.md)。
用起来什么样。 一个 Node worker 注册一个函数,另一个进程直接按名字调它:
// 示意,基于 sdk/packages/node/iii —— 注册一个函数
const worker = registerWorker('ws://localhost:49134')
worker.registerFunction('orders::validate', async (order) => {
return { ok: order.total > 0 } // 返回值会被送回调用方
})
// 别处(任何语言的 worker)按名字调用,像调本地函数一样 await
const res = await worker.trigger({ function_id: 'orders::validate', payload: order })
一句话直觉/类比。 把引擎想成一部总机 + 电话簿:每个进程上线时在电话簿登记自己会接哪些"分机号"(function id),总机负责接线——你只报号码,不用知道对方在哪台机器、用什么语言。
2. 顶层全景(它大概怎么转)
怎么读这张图: 上面是各种进程(worker),都只通过 WebSocket 连到中间那个 Rust 引擎;引擎内部是"注册表 + 路由器 + 调用登记簿"三件套。所有调用都经过引擎中转,没有 worker↔worker 直连。
worker A (TS) worker B (Python) 内置 worker
注册函数/触发器 注册函数/触发器 http·cron·queue·state…
│ │ │
└────── WebSocket ──────┴────── WebSocket ────────┘
│
┌────────────▼─────────────┐
│ iii 引擎 (Rust) │
│ ┌────────────────────┐ │
│ │ 活注册表: │ │ 函数 / 触发器 / worker
│ │ 谁提供了哪个 id │ │
│ └────────────────────┘ │
│ router_msg 消息路由 │ 收发 WS 消息、分发
│ Invocation 调用登记簿 │ 把跨进程调用变成一次 await
└──────────────────────────┘
部件一句话职责:
| 部件 | 干什么 | 在哪 |
|---|---|---|
线上协议 Message | 定义 WS 上收发的所有消息种类(注册/调用/结果/心跳) | engine/src/protocol.rs:42 |
Engine | 持有全部注册表,是所有状态的中心 | engine/src/engine/mod.rs:229 |
router_msg | 收到一条 WS 消息后按类型分发处理 | engine/src/engine/mod.rs:589 |
InvocationHandler | 用 oneshot 通道登记"正在等结果"的调用 | engine/src/invocation/mod.rs:76 |
TriggerRegistry | 记录 trigger 类型与实例,把触发外包给拥有者 | engine/src/trigger.rs:199 |
| SDK(TS/Py/Rust/Go) | worker 侧:连接、注册、分发调用、自动重连 | sdk/packages/node/iii/src/iii.ts |
主线走一遍(高层,不进代码):
- 连上。 进程开一条 WebSocket 到引擎(默认
ws://localhost:49134),握手时建立一个带 RBAC 的会话(handle_worker,engine/src/engine/mod.rs:1429)。 - 注册。 进程发
RegisterFunction/RegisterTriggerType/RegisterTrigger,引擎把它们记进各自的注册表(谁提供了哪个 id)。 - 调用。 任何 worker 发
InvokeFunction { function_id, data };引擎查注册表找到"拥有该函数的 worker",把调用转过去。 - 回传。 目标 worker 执行完发
InvocationResult;引擎按invocation_id找回那个在等结果的调用方,把结果送回去——调用方那边只是一次await返回了。 - 触发。 当 cron 到点、HTTP 命中、state 变更时,拥有该触发器类型的 worker 调用引擎去跑绑定的函数,复用同一条调用路径扇出。
想看每一步的真实代码与数据结构,从 §3 的阅读地图按顺序进各章。
3. 阅读地图(建议顺序)
本子库拆成 6 章,由浅入深。人类建议顺序读;agent 可按 keyTopics 或下表"这章讲什么"直接跳到相关章。
| 顺序 | 章节 | 这章讲什么 |
|---|---|---|
| 0 | iii 是什么 · 全景 · 阅读地图(本页) | essence、顶层全景、导航;先判断相关性 |
| 1 | 三原语与线上协议 | Worker/Function/Trigger 精确定义;Message 枚举里每种 WS 消息的字段与语义 |
| 2 | 引擎中枢:注册表与消息分发 | Engine 里的几张 DashMap 注册表;router_msg 如何按消息类型处理;所有权与快重启竞态防护 |
| 3 | 一次调用的一生 | 从 InvokeFunction 到 spawn_invoke_function→handle_invocation;Deferred + oneshot 如何把跨进程调用折叠成一次 await;fire-and-forget 与 enqueue 变体 |
| 4 | 触发器体系与内置 worker | TriggerType 的 registrator 委派模型;内置 worker(http/cron/queue/state/stream/pubsub)如何提供触发器类型并扇出 |
| 5 | 活的系统:发现、运行时扩展与可观测性 | engine::* 自省函数;workers-available/functions-available 事件;运行时创建 worker;OTEL 遥测双通道 |
| 6 | 接入引擎:SDK 与 worker 握手 | SDK 侧的连接、onSocketOpen 重放注册、onInvokeFunction 分发、断线重连与消息缓冲 |
4. 巧妙之处(读完要带走的精华)
下面每条都是 iii 设计上"不显然但关键"的决定。想深看的进对应章。
① Deferred + oneshot:把跨 WebSocket 调用折叠成一次 await。 引擎调一个 WS worker 的函数时,并不能"就地拿到返回值"——结果得等 worker 那边异步回传。iii 的做法是让 worker 侧的处理器返回一个特殊结果 FunctionResult::Deferred(engine/src/function.rs:18),引擎把这次调用连同一个 oneshot::Sender 存进 DashMap(handle_invocation,engine/src/invocation/mod.rs:76),然后 await 那个 oneshot。等 worker 发回 InvocationResult,router_msg 按 invocation_id 取出 sender 一 send,await 就返回了。引擎全程不阻塞线程。详见第 3 章。
② 三原语统一一切入口。 HTTP、cron、队列、state 变更、直接调用——在引擎眼里全是"某个 trigger 让某个 function 跑"。引擎只认 function_id 和几种消息,不认具体协议。内置触发器类型是一张静态表 BUILTIN_TRIGGER_TYPES(engine/src/trigger.rs:16)。详见第 4 章。
③ registrator 委派:引擎自己不懂 HTTP/cron。 每个 trigger type 记着一个 registrator——拥有该类型的 worker。用户 worker 注册一个 http/cron 触发器时,引擎不自己去开端口或起定时器,而是把 RegisterTrigger 转发给拥有该类型的内置 worker,由它建立真实监听;触发时那个 worker 再回头 engine.call(function_id) 扇出(如 cron adapter 到点即调,engine/src/workers/cron/structs.rs:202)。引擎始终只当路由器。详见第 4 章。
④ 活注册表 + 推送式发现。 worker 一连上或一注册,引擎就 fire_triggers(workers-available, …)(engine/src/engine/mod.rs:1400),订阅了该事件的 worker 立刻收到通知——发现是推的,不是轮询。配合 engine::workers::list / functions::list 等自省函数,系统对自己"当前有什么能力"始终自知。详见第 5 章。
⑤ 快重启竞态防护。 一个 worker 崩溃重连、恰好复用同一个 function id 时,旧连接的 cleanup_worker(engine/src/engine/mod.rs:1626)可能误删新连接刚写下的注册。iii 用 function_owners 这张所有权表做 CAS 式校验:清理前先确认自己仍是登记的 owner,不是就跳过。这类细节撑起了"live 系统可安全热重启"。详见第 2 章。
⑥ 治理放在路由层。 会话可挂 RBAC/中间件 hook:注册函数、注册触发器、发起调用之前先跑一个钩子函数,可以改写参数或拒绝请求(见 router_msg 里对 on_*_registration_function_id / middleware_function_id 的处理,engine/src/engine/mod.rs:589 起)。鉴权与策略因此是运行时可插拔的,不写死在每个 worker 里。
⑦ 双通道遥测,跨语言拼一条 trace。 SDK 另开一条只走遥测的 WS,发 OTLP/MTRC/LOGS 前缀的二进制帧(handle_telemetry_frame,engine/src/engine/mod.rs:81);每条调用消息带 W3C traceparent/baggage 全程透传,于是"TS worker → 引擎 → Python worker"能在一条分布式 trace 里串起来。详见第 5 章。
5. 代码地图(导航索引)
按符号名可 grep 定位(行号会随上游漂移,符号名通常还在)。
| 主题 | 文件:行 | 关键符号 |
|---|---|---|
| 线上消息 schema(全部消息类型) | engine/src/protocol.rs:42 | enum Message(RegisterFunction/InvokeFunction/InvocationResult/RegisterTrigger…) |
| 引擎中心状态 | engine/src/engine/mod.rs:229 | struct Engine(functions/trigger_registry/worker_registry/invocations) |
| WS 消息分发 | engine/src/engine/mod.rs:589 | Engine::router_msg |
| 调用的后台执行与回传 | engine/src/engine/mod.rs:462 | Engine::spawn_invoke_function |
| WS 连接主循环 + 握手 | engine/src/engine/mod.rs:1429 | Engine::handle_worker |
| 断连清理(含所有权校验) | engine/src/engine/mod.rs:1626 | Engine::cleanup_worker / claim_function |
| 触发器扇出 | engine/src/engine/mod.rs:1400 | Engine::fire_triggers |
| 内部调用路径(hook/cron 复用) | engine/src/engine/mod.rs:1794 | impl EngineTrait for Engine::call |
| 调用登记簿(oneshot) | engine/src/invocation/mod.rs:76 | InvocationHandler::handle_invocation / halt_invocation |
| 函数注册表与结果类型 | engine/src/function.rs:18 | enum FunctionResult(Deferred)/FunctionsRegistry |
| worker 侧函数处理(返回 Deferred) | engine/src/worker_connections/traits.rs:77 | impl FunctionHandler for WorkerConnection::handle_function |
| 触发器注册表与内置类型表 | engine/src/trigger.rs:16 | BUILTIN_TRIGGER_TYPES / TriggerRegistry / TriggerType |
| worker 连接状态 | engine/src/worker_connections/mod.rs:225 | struct WorkerConnection / WorkerConnectionRegistry |
| 引擎自省 / 发现函数 | engine/src/workers/engine_fn/mod.rs:1287 | engine::workers::list / functions::list / triggers::list / register_worker |
| SDK:连接、重放注册、分发调用 | sdk/packages/node/iii/src/iii.ts:709 | Sdk.onSocketOpen / onInvokeFunction / onInvocationResult / trigger |