iii 引擎协议详解:SDK Worker 与 Engine 之间的 WebSocket 线级协议
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
本文以 iii 仓库中docs/next/reference/engine-protocol.mdx的官方协议参考为主体,逐帧讲解 engine 与各语言 SDK worker 在 WebSocket 上交换的线级(wire-level)协议:端口分配、连接建立与命名空间挂起机制、全部消息帧的字段语义、触发器命名空间解析、注册冲突处理、engine::*发现函数与引擎采集指标。读完本文后,你可以直接读懂或手写一个不走语言 SDK 的原生协议客户端,并理解 engine/src/protocol.rs 中各字段在引擎内部的真实行为。
说明:绝大多数项目通过语言 SDK(Node、Python、Rust、Browser)接入,从不需要直接触碰协议;本文所列的 JSON 形态就是这些 SDK 序列化时的“事实来源”(source of truth)。另外,可观测性内省(traces、logs、metrics、采样规则、告警、rollups)由 iii-observability worker 端到端负责,该 worker 由引擎自动注入,不要将其声明进
config.yaml、engine.workers或containers。
连接端口
引擎绑定 worker 与 stream 两个 WebSocket;项目 HTTP 路由与可观测性则各有独立的 worker 承载:
| 端口 | 绑定方 | 表面 |
|---|---|---|
3111 | httpworker | 项目 HTTP 路由(可在worker-compose.yaml中配置)。 |
3112 | engine | Stream API(WebSocket,消费端的流订阅)。 |
49134 | engine | SDK WebSocket;即iii_sdk::register_worker打开的地址。 |
9464 | iii-observabilityworker | Prometheus 指标端点(通常与引擎暴露在同一个容器中)。 |
Console UI 运行在3113,由iii console单独启动。
从源码结构看,SDK WebSocket 的默认端口在 engine/src/workers/worker/mod.rs 中定义为pub const DEFAULT_PORT: u16 = 49134;;而 stream/queue worker 的 bridge 适配器也硬编码了ws://localhost:49134作为缺省bridge_url(见 engine/src/workers/queue/adapters/bridge.rs 与 engine/src/workers/stream/adapters/bridge.rs),因此引擎内部 worker 与外部 SDK worker 走的是同一条协议通道。
连接流程
worker 打开 SDK WebSocket(默认ws://127.0.0.1:49134)。连接建立后,引擎为该 worker 分配一个 UUID,并发送携带它的WorkerRegistered { worker_id }帧。worker 随后发送其在内存中持有的注册帧(每条RegisterFunction、RegisterTrigger与RegisterTriggerType),并调用engine::workers::register发布自身元数据(runtime、version、OS、PID、isolation、可选namespace、可选的单行description),引擎以RegisterWorkerResult确认。
此后连接完全双向:引擎向 worker 推送InvokeFunction帧,worker 回推InvocationResult、追加注册或反注册。
命名空间挂起(hold)机制
连接所属的命名空间来自engine::workers::register调用。客户端可以在该调用之前发送注册帧——引擎会把这些帧挂起,直到它知道命名空间为止;引擎不会先把它们注册到default再迁移。当engine::workers::register到达时,引擎注册 worker、设定连接命名空间,并按到达顺序处理被挂起的帧;如果注册超时先到期,引擎将连接命名空间设为default。
在 engine/src/engine/mod.rs 中可以看到对应的注释语义:被挂起的帧在“注册调用真正到达”与“超时真正触发”两种情形下有完全不同的归宿(源码注释明确区分了engine::workers::register到达与 timeout 真正 fired 两个分支)。
源码中的额外细节:Reattach 与 WorkerMetrics
Message枚举中还包含文档表格未单列、但对运维很关键的两处扩展:
Reattach帧(worker → engine):重连时的第一条消息,携带上一次连接中WorkerRegistered下发的previous_worker_id与reattach_token。引擎将旧连接按正常断开流程退役,使重连后的注册重放落在干净状态上,而不是与旧连接的清理竞速。token 是必需的——worker id 对外可发现,而 token 只通过旧连接自身的 socket 下发过。旧引擎对未知消息类型只记录警告日志并忽略,因此该设计对版本偏斜(version skew)安全。对应测试见 engine/src/protocol.rs 中的reattach_round_trips与worker_registered_token_is_optional_on_the_wire。WorkerMetrics结构:worker 上报的资源指标(heap/RSS 内存、CPU 微秒与百分比、事件循环滞后、运行时长、时间戳、runtime 名称),即engine::workers::list返回的latest_metrics。源码中附带了 JavaScript 精度说明:u64字段理论上可超过Number.MAX_SAFE_INTEGER,但内存值需超过约 9 PB 才会丢精度,CPU 微秒需约 285 年连续运行才会触及,因此多数场景无需 BigInt 解析。
消息类型
每个帧都是 JSON 对象,由type字段判别(取变体名的全小写形式,如registerfunction)。完整集合定义在 engine/src/protocol.rs 的Message枚举上(#[serde(tag = "type", rename_all = "lowercase")]):
| 帧 | 方向 | 用途 |
|---|---|---|
RegisterFunction | worker → engine | 使一个函数可按function_id被调用。 |
UnregisterFunction | worker → engine | 移除已注册的函数。 |
RegisterTrigger | worker → engine | 将函数绑定到一个触发器实例。 |
UnregisterTrigger | worker → engine | 移除触发器绑定。 |
TriggerRegistrationResult | engine → worker | 对RegisterTrigger的确认 / 错误。 |
RegisterTriggerType | worker → engine | 声明 worker 通告的新触发器类型。 |
RegisterService | worker → engine | 将相关函数归组到一个 service id 下。 |
InvokeFunction | engine → worker | 携带载荷调用一个已注册的函数。 |
InvocationResult | worker → engine | 回传函数结果或错误。 |
WorkerRegistered | engine → worker | 确认 worker,附带分配的worker_id。 |
RegistrationRejected | engine → worker | 拒绝与同命名空间内存活 worker 冲突的注册。 |
Ping/Pong | 双向 | 存活探测;防止空闲连接超时。 |
RegisterFunction
{ "type": "registerfunction", "id": "math::add", "description": "Add two numbers.", "request_format": { "type": "object", "properties": { "a": { "type": "number" }, "b": { "type": "number" } } }, "response_format": { "type": "object", "properties": { "c": { "type": "number" } } }, "metadata": { "owner": "math-team" }, "invocation": null }id必填。description、request_format、response_format、metadata均为可选,用于喂给 iii console 与 agent 可读的 skills。invocation保留给外部 HTTP 函数(HttpInvocationRef),进程内处理器应置null。
对照 engine/src/protocol.rs 的定义,HttpInvocationRef的完整字段为:
pub struct HttpInvocationRef { pub url: String, #[serde(default = "default_http_method")] pub method: HttpMethod, // 缺省为 POST #[serde(default)] pub timeout_ms: Option<u64>, #[serde(default)] pub headers: HashMap<String, String>, #[serde(default)] pub auth: Option<HttpAuthConfig>, }也就是说,一个声明了invocation的函数并不需要本 worker 持有处理器:引擎可直接按url/method/headers/auth/timeout_ms对外发起 HTTP 调用,这就是“无进程内 handler 也能注册函数”的机制。
RegisterTrigger
{ "type": "registertrigger", "id": "math::add@http", "trigger_type": "http", "function_id": "math::add", "config": { "api_path": "/math/add", "http_method": "POST" }, "metadata": null, "namespace": "orders", "trigger_namespace": null }config是逐触发器类型的配置,其形状由通告该trigger_type的 worker 定义(例如http触发器的http配置)。引擎以TriggerRegistrationResult应答,其中携带可选的error: ErrorBody(ErrorBody结构为code、message与可选stacktrace,见 engine/src/protocol.rs)。
namespace与trigger_namespace是两个不同的问题
namespace指定目标函数所在的命名空间,与 worker 注册使用同一套命名空间体系,通常与 worker 命名空间相同。触发器可以调用另一个命名空间里的函数,因此RegisterTrigger携带目标namespace。字段缺省时,function_id在default中解析;字段存在则必须是非空字符串——引擎拒绝null及其他任何非字符串值。trigger_namespace指定在哪里寻找触发器类型的 provider。一个定位目标函数,一个定位点燃它的 provider。
若trigger_namespace缺省,引擎按两步解析:先查注册连接自身的命名空间,再查default。这与发送"default"不同。两步式设计允许项目为引擎也提供的某个 trigger type id 注册自己的 provider,而未命名任何命名空间的 worker 仍能触达引擎的 provider。若trigger_namespace显式给出,则解析是严格的:要么该命名空间,要么什么都没有;显式指定了命名空间的绑定绝不会被迁移到其他 provider。
还有一个与启动顺序无关的规则:当某个 provider 在绑定已解析到default之后才在某个命名空间注册时,引擎会把该绑定迁移到新 provider——启动顺序不决定由哪个 provider 服务项目。
示例:这些字段指向default中的math::add:
{ "function_id": "math::add" }这些字段指向orders中的math::add:
{ "function_id": "math::add", "namespace": "orders" }RegisterTriggerType
{ "type": "registertriggertype", "id": "webhook", "description": "HTTP webhook trigger", "trigger_request_format": { "type": "object", ... }, "call_request_format": { "type": "object", ... }, "namespace": null }trigger_request_format是该触发器每个绑定config的 JSON Schema;call_request_format是触发器点火时投递给绑定函数的载荷的 JSON Schema。
namespace是该 provider 服务的命名空间。字段缺省时,引擎把 provider 记在注册连接的命名空间下。provider 以(namespace, trigger_type_id)为键,因此两个处于不同命名空间的 worker 可以通告相同的trigger_typeid 而互不覆盖。引擎自带的 provider(http、cron、state、stream)注册在default。
InvokeFunction
{ "type": "invokefunction", "function_id": "math::add", "data": { "a": 2, "b": 3 }, "metadata": { "tenant": "acme" }, "traceparent": "00-…", "baggage": "k=v,…", "action": { "type": "void" }, "namespace": "orders" }Void调用省略invocation_id(worker 没有可回复的结果通道)。- 触发器注册上的可选
metadata(上例中的null/None)是任意 JSON,与触发器一起存储,并在触发时作为独立参数(不并入data)与载荷一起投递给接收函数。它为接收函数提供关于触发器或执行上下文的上下文信息;被多个触发器共享的目标函数可用它判断是哪一次注册点火、带着什么上下文。metadata既可通过registerTrigger提供,也可在直接trigger()调用时提供。 traceparent与baggage携带 W3C trace context,用于分布式追踪。action是路由标志(见下文“触发器动作”),缺省/null表示同步。namespace可选,选择function_id在其中解析的命名空间。省略该字段即在default中解析;省略还能保证从不发送该字段的旧端对端保持线级兼容。字段存在时发送非空字符串;引擎拒绝null及其他非字符串值。解析规则见下文“命名空间”。
InvocationResult
成功:
{ "type": "invocationresult", "invocation_id": "9f3c…", "function_id": "math::add", "result": { "c": 5 }, "error": null, "traceparent": "00-…", "baggage": "k=v,…" }失败:
{ "type": "invocationresult", "invocation_id": "9f3c…", "function_id": "math::add", "result": null, "error": { "code": "invocation_failed", "message": "boom", "stacktrace": "TraceError: …" } }出现在InvocationResult.error中的ErrorBody.code值包括:invocation_failed(handler 抛出异常)、invocation_stopped(持有该函数的 worker 在飞行中断开,引擎取消进行中的调用并向调用方暴露此 code)、function_not_found、function_not_invokable、TIMEOUT(客户端侧超时)、FORBIDDEN(RBAC 拒绝)。
function_not_found消息会指明查找发生在哪个命名空间,并列出该 id 实际存在的命名空间,例如:
Function state::get not found in namespace default. It is registered in namespace(s): orders, analytics.
RegistrationRejected与冲突处理
{ "type": "registrationrejected", "code": "WORKER_NAMESPACE_CONFLICT", "namespace": "orders", "worker_name": "state", "owner_worker_id": "3f9c1a2e-…" }引擎在某注册与namespace内存活 worker 冲突时发送该消息。owner_worker_id标识持有该身份的 worker;code指明冲突的身份字段与严重程度。每条消息只含一个身份字段:
code | 身份字段 | 连接 | 严重程度 |
|---|---|---|---|
WORKER_NAMESPACE_CONFLICT | worker_name | 被引擎关闭 | 致命。SDK 停止 worker 且不再重连。 |
FUNCTION_NAMESPACE_CONFLICT | function_id | 保持打开 | 非致命。引擎拒绝单个函数注册,worker 仍服务其余函数。 |
函数冲突的帧形如:
{ "type": "registrationrejected", "code": "FUNCTION_NAMESPACE_CONFLICT", "namespace": "orders", "function_id": "state::get", "owner_worker_id": "3f9c1a2e-…" }函数冲突的行为
worker 注册与函数注册是两个独立操作:引擎可以接受一个 worker 但拒绝它其中的某个函数。对函数所有权冲突,引擎的处理顺序是:
- 保留当前函数持有者;
- 不注册新 handler;
- 向新 worker 发送
FUNCTION_NAMESPACE_CONFLICT; - 保持新 worker 连接打开;
- 继续注册并服务新 worker 的其余函数。
FUNCTION_NAMESPACE_CONFLICT是注册结果,不是调用结果。之后对同一命名空间、同一 function id 的调用会路由到当前持有者,而不会发给注册被拒的 worker。
注意:worker 已连接并不证明其所有函数都注册成功。SDK 把函数冲突作为警告上报并让 worker 继续活跃;如果 worker 要求其全部函数在线,应把该警告当作启动或部署错误处理。
一个值得注意的边界:worker 重启时撞上自己尚未清理完毕的旧连接不算冲突。引擎把已开始拆除的连接视为不存活,因此重启立即收回该名称。这一点在 engine/src/engine/mod.rs 中由 worker 名称租约(name lease)的 CAS 机制实现:只有当另一个仍存活的 worker 在该命名空间持有同名租约时才会失败,而持有者连接正在拆除的窗口则被新声明直接收回。
触发器动作(Trigger actions)
InvokeFunction.action以type标记,线级编码为小写(对应 engine/src/protocol.rs 中的TriggerAction枚举,仅两个变体):
| 线级形状 | 含义 |
|---|---|
省略 /null | 同步;worker 以InvocationResult回复。 |
{ "type": "void" } | 发后即忘;无invocation_id,无回复。 |
{ "type": "enqueue", "queue": "math" } | 路由到指定名称的队列(由queueworker 提供)。 |
调用生命周期
- 同步:引擎分配
invocation_id,把InvokeFunction转发给持有 worker,等待配对的InvocationResult。 Void:引擎不带invocation_id转发,且从不期待回复。Enqueue:引擎把调用交给 queue worker,由其持久化后按队列的重试策略在订阅者侧重放目标函数。
命名空间(Namespaces)
命名空间是随 function id 一起携带的路由值,不是函数名的一部分。例如state::get在每个命名空间中的 id 都是相同的。注册表以(namespace, function_id)与(namespace, worker_name)为键,因此同一个 id 或 worker 名在每个命名空间内至多出现一次。
worker 在engine::workers::register调用中声明自己的命名空间(对应 engine/src/workers/engine_fn/mod.rs 中RegisterWorkerInput.namespace字段,#[serde(default)]——缺省即default,保证旧 SDK 构建的 worker 无需改动继续注册)。未声明命名空间的连接落入default。
解析规则
调用路由是严格的,绝不回退到其他命名空间:
InvokeFunction.namespace | 解析于 |
|---|---|
| 缺省 | 仅default |
"orders" | 仅orders |
null或其他非字符串 | 引擎拒绝该帧 |
SDK 在调用方未设置时,从 worker 自身命名空间填充该字段,因此帧携带的是 worker 自己的命名空间(worker 未声明时即default)。
未命中返回function_not_found,并列出该 id 实际存在的命名空间。
内省解析则刻意宽松,使得一个所有 worker 都在非默认命名空间里的引擎仍能回答关于自身的问题。对未带显式namespace的engine::functions::info与engine::workers::info:default中的条目优先;否则,仅唯一存在于某一个非默认命名空间的 id 或名称可解析;同时出现在多个非默认命名空间的 id 或名称被报告为歧义并列出候选,绝不靠猜测解析。传入显式namespace即恢复严格解析。
保留 id
engine::*前缀保留给注册在default中的引擎基础设施。自定义 worker 试图在其他命名空间注册engine::*函数 id 时,该注册会被拒绝——engine/src/engine/mod.rs 中的测试即验证了“拒绝在default之外注册保留的engine::*函数 id”,且有断言确保这样的函数“绝不能泄漏进 default 命名空间”。
例外:引擎自带的 queue worker 将engine::queue::*函数注册在default,保留 id 检查不会拒绝它们。
线级兼容
所有 namespace 字段都是可选的,未设置时省略。针对旧 SDK 构建的 worker 不发送 namespace,落入default,行为与从前完全一致。显式null不同于字段缺省:引擎会拒绝它,如同拒绝任何非字符串值。
此外,engine/src/protocol.rs 中专门定义了空白命名空间的处理:is_blank_namespace把“缺省”与“空白”视为相反的意图——缺省是请求default;空白是命名了一个命名空间却没给出名字。若把空白读成缺省,worker 会被放进它从未要求的default,其全部调用与触发器随后跟随,而这恰恰是运维者读声明时看不到的地方;引擎作为“所有客户端都必须经过的唯一一方”在 SDK 之外再拒绝一次,防止遗忘(或手写客户端)创造出无人可寻址的匿名命名空间,对应错误码为INVALID_NAMESPACE。
引擎发现函数
引擎在engine::*前缀下注册一组用于内省与 worker 生命周期的函数,定义见 engine/src/workers/engine_fn/mod.rs:
| 函数 | 用途 |
|---|---|
engine::channels::create | 创建流式通道的读 / 写端对。 |
engine::functions::list | 列出全部已注册函数(可用include_internal过滤)。 |
engine::functions::info | 检查一个或多个函数(单个function_id,或至多 32 个function_ids):schema、持有者、已注册触发器。接受可选namespace。 |
engine::workers::list | 列出所有已连接 worker 及其指标。 |
engine::workers::info | 检查单个已连接 worker 的完整表面(函数、触发器类型、已注册触发器)。参数为name加可选namespace。 |
engine::triggers::list | 列出全部已注册触发器类型(可用include_internal过滤)。 |
engine::triggers::info | 检查单个触发器类型:schema、持有者、存活实例数。 |
engine::registered-triggers::list | 列出全部已注册触发器实例(可用include_internal过滤)。 |
engine::registered-triggers::info | 检查单个已注册触发器实例,含反规范化的触发器与函数详情。 |
engine::workers::register | 发布调用方 worker 的元数据(runtime、version、OS、PID、isolation、可选namespace、可选description)。 |
engine::register_trigger | 注册一个直接点燃function_id的触发器,可选metadata作为独立参数投递给 handler。返回触发器 id。 |
engine::unregister_trigger | 按 id 反注册触发器。幂等;报告其是否曾存在。 |
engine::workers::register的输入结构RegisterWorkerInput(见 engine/src/workers/engine_fn/mod.rs)还包含一个细节字段:_caller_worker_id(serde 重命名的worker_id),即调用方自身身份标识;engine::register_trigger的输入同理带有引擎从调用方 worker 注入的_caller_worker_id,用于将该触发器限定到该 worker 的范围内,使其在 worker 断连时被回收(进程内调用方则无此字段)。engine::unregister_trigger的返回结构UnregisterTriggerResult携带removed: bool,即“是否曾存在并被移除”。
engine::functions::list、engine::functions::info、engine::workers::list、engine::workers::info返回的每一行都带一个namespace字段,指明条目所注册的注册表键,用以区分两行共享同一function_id或同一 workername的情况。两个list函数都不接受namespace过滤——都返回所有命名空间的结果。
引擎发现触发器
| 触发器 | 点火时机 |
|---|---|
engine::functions-available | 函数被注册或反注册时。 |
engine::workers-available | worker 连接或断开时。 |
引擎采集指标
以下指标由引擎自身发出,与 worker 使用何种语言 SDK 无关。名称与单位来自 engine/src/workers/observability/metrics.rs。
调用(Invocations)
| 指标 | 仪器类型 | 单位 |
|---|---|---|
iii.invocations.total | counter | invocations |
iii.invocation.duration | histogram | s |
iii.invocation.errors.total | counter | errors |
Worker 全局(Workers)
| 指标 | 仪器类型 | 单位 |
|---|---|---|
iii.workers.active | gauge | workers |
iii.workers.spawns.total | counter | workers |
iii.workers.deaths.total | counter | workers |
iii.workers.by_status | gauge | workers |
逐 worker(Per-worker)
| 指标 | 仪器类型 | 单位 |
|---|---|---|
iii.worker.memory.heap.bytes | gauge | bytes |
iii.worker.memory.rss.bytes | gauge | bytes |
iii.worker.cpu.percent | gauge | % |
iii.worker.event_loop.lag.ms | gauge | ms |
iii.worker.uptime.seconds | gauge | s |
小结:读懂这份协议的三条主线
- 一切注册都是命名空间化的:函数、worker、触发器类型、触发器绑定全部以
(namespace, id)为键;路由严格不回落,内省宽松不猜解,空白命名空间被显式拒绝(INVALID_NAMESPACE)。 - 注册与连接解耦:worker 身份冲突是致命的(连接被关闭,SDK 不重连),函数冲突是非致命的(连接保持、其余函数照常服务),因此“worker 已连接”不等于“函数全部在线”。
- 线级兼容优先:所有新字段(
namespace、trigger_namespace、metadata、reattach_token等)一律Option+ 缺省省略,保证旧 SDK 构建的 worker 无感接入;协议测试(见 engine/src/protocol.rs 的tests模块)专门覆盖 token 缺省、未知类型容忍等版本偏斜场景。
掌握以上三帧(RegisterFunction/InvokeFunction/InvocationResult)的字段语义、触发器双命名空间解析与RegistrationRejected的两种 code,即可独立实现或排障任何直接对接49134端口的协议客户端。
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考