news 2026/9/14 17:30:52

iii 引擎协议详解:SDK Worker 与 Engine 之间的 WebSocket 线级协议

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
iii 引擎协议详解:SDK Worker 与 Engine 之间的 WebSocket 线级协议

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.yamlengine.workerscontainers

连接端口

引擎绑定 worker 与 stream 两个 WebSocket;项目 HTTP 路由与可观测性则各有独立的 worker 承载:

端口绑定方表面
3111httpworker项目 HTTP 路由(可在worker-compose.yaml中配置)。
3112engineStream API(WebSocket,消费端的流订阅)。
49134engineSDK WebSocket;即iii_sdk::register_worker打开的地址。
9464iii-observabilityworkerPrometheus 指标端点(通常与引擎暴露在同一个容器中)。

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 随后发送其在内存中持有的注册帧(每条RegisterFunctionRegisterTriggerRegisterTriggerType),并调用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_idreattach_token。引擎将旧连接按正常断开流程退役,使重连后的注册重放落在干净状态上,而不是与旧连接的清理竞速。token 是必需的——worker id 对外可发现,而 token 只通过旧连接自身的 socket 下发过。旧引擎对未知消息类型只记录警告日志并忽略,因此该设计对版本偏斜(version skew)安全。对应测试见 engine/src/protocol.rs 中的reattach_round_tripsworker_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")]):

方向用途
RegisterFunctionworker → engine使一个函数可按function_id被调用。
UnregisterFunctionworker → engine移除已注册的函数。
RegisterTriggerworker → engine将函数绑定到一个触发器实例。
UnregisterTriggerworker → engine移除触发器绑定。
TriggerRegistrationResultengine → workerRegisterTrigger的确认 / 错误。
RegisterTriggerTypeworker → engine声明 worker 通告的新触发器类型。
RegisterServiceworker → engine将相关函数归组到一个 service id 下。
InvokeFunctionengine → worker携带载荷调用一个已注册的函数。
InvocationResultworker → engine回传函数结果或错误。
WorkerRegisteredengine → worker确认 worker,附带分配的worker_id
RegistrationRejectedengine → 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必填。descriptionrequest_formatresponse_formatmetadata均为可选,用于喂给 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: ErrorBodyErrorBody结构为codemessage与可选stacktrace,见 engine/src/protocol.rs)。

namespacetrigger_namespace是两个不同的问题

  • namespace指定目标函数所在的命名空间,与 worker 注册使用同一套命名空间体系,通常与 worker 命名空间相同。触发器可以调用另一个命名空间里的函数,因此RegisterTrigger携带目标namespace。字段缺省时,function_iddefault中解析;字段存在则必须是非空字符串——引擎拒绝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(httpcronstatestream)注册在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()调用时提供。
  • traceparentbaggage携带 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_foundfunction_not_invokableTIMEOUT(客户端侧超时)、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_CONFLICTworker_name被引擎关闭致命。SDK 停止 worker 且不再重连。
FUNCTION_NAMESPACE_CONFLICTfunction_id保持打开非致命。引擎拒绝单个函数注册,worker 仍服务其余函数。

函数冲突的帧形如:

{ "type": "registrationrejected", "code": "FUNCTION_NAMESPACE_CONFLICT", "namespace": "orders", "function_id": "state::get", "owner_worker_id": "3f9c1a2e-…" }

函数冲突的行为

worker 注册与函数注册是两个独立操作:引擎可以接受一个 worker 但拒绝它其中的某个函数。对函数所有权冲突,引擎的处理顺序是:

  1. 保留当前函数持有者;
  2. 不注册新 handler;
  3. 向新 worker 发送FUNCTION_NAMESPACE_CONFLICT
  4. 保持新 worker 连接打开;
  5. 继续注册并服务新 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.actiontype标记,线级编码为小写(对应 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 都在非默认命名空间里的引擎仍能回答关于自身的问题。对未带显式namespaceengine::functions::infoengine::workers::infodefault中的条目优先;否则,仅唯一存在于某一个非默认命名空间的 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::listengine::functions::infoengine::workers::listengine::workers::info返回的每一行都带一个namespace字段,指明条目所注册的注册表键,用以区分两行共享同一function_id或同一 workername的情况。两个list函数都不接受namespace过滤——都返回所有命名空间的结果。

引擎发现触发器

触发器点火时机
engine::functions-available函数被注册或反注册时。
engine::workers-availableworker 连接或断开时。

引擎采集指标

以下指标由引擎自身发出,与 worker 使用何种语言 SDK 无关。名称与单位来自 engine/src/workers/observability/metrics.rs。

调用(Invocations)

指标仪器类型单位
iii.invocations.totalcounterinvocations
iii.invocation.durationhistograms
iii.invocation.errors.totalcountererrors

Worker 全局(Workers)

指标仪器类型单位
iii.workers.activegaugeworkers
iii.workers.spawns.totalcounterworkers
iii.workers.deaths.totalcounterworkers
iii.workers.by_statusgaugeworkers

逐 worker(Per-worker)

指标仪器类型单位
iii.worker.memory.heap.bytesgaugebytes
iii.worker.memory.rss.bytesgaugebytes
iii.worker.cpu.percentgauge%
iii.worker.event_loop.lag.msgaugems
iii.worker.uptime.secondsgauges

小结:读懂这份协议的三条主线

  1. 一切注册都是命名空间化的:函数、worker、触发器类型、触发器绑定全部以(namespace, id)为键;路由严格不回落,内省宽松不猜解,空白命名空间被显式拒绝(INVALID_NAMESPACE)。
  2. 注册与连接解耦:worker 身份冲突是致命的(连接被关闭,SDK 不重连),函数冲突是非致命的(连接保持、其余函数照常服务),因此“worker 已连接”不等于“函数全部在线”。
  3. 线级兼容优先:所有新字段(namespacetrigger_namespacemetadatareattach_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/14 17:29:05

Hadoop源码深度剖析:从核心模块到RPC与HDFS读写链路实战

我决定把这一年多啃Hadoop源码的笔记整理成一篇可以直接照着读的索引式分享。不是那种罗列类名的源码导读&#xff0c;也不是贴一堆注释的代码复述&#xff0c;而是从“我为什么会去读源码、读了哪些模块、怎么搭建调试环境、核心流程到底怎么跑通、踩了哪些坑”这几个角度&…

作者头像 李华
网站建设 2026/9/14 17:23:25

【C语言】 数组

目录 1&#xff0c;数组的概念 2&#xff0c;数组的创建和初始化 3&#xff0c;数组的使用 4&#xff0c;数组的内存存储情况 5&#xff0c;sizeof 计算数组的元素个数 6&#xff0c;二维数组 7&#xff0c;二维数组的初始化和创建 8&#xff0c;二维数组的使用 9&…

作者头像 李华
网站建设 2026/9/14 17:22:31

供配电实训仿真软件:倒闸操作与故障处理全流程解析

干电气培训这些年&#xff0c;我见过太多学员第一次面对高压柜时的表情——手放在断路器分闸按钮上&#xff0c;迟迟不敢按下去。不是不知道步骤&#xff0c;而是怕按错了出大事。这种怕是对的&#xff0c;10kV开关柜一旦带负荷拉隔离开关&#xff0c;电弧瞬间就能把人灼伤&…

作者头像 李华
网站建设 2026/9/14 17:19:13

C++序列输出题全攻略:从读题到OJ提交的完整避坑指南

1. 一道短得不像话的题&#xff0c;凭什么让我交了三版才过东华OJ的基础题里有一类题属于“看着简单、做着崩溃”&#xff0c;第50题“按要求输出序列”就是典型。题面可能短到只有一句话&#xff0c;给一个整数N&#xff0c;让你按某种规则输出一串数。很多人的第一反应是&…

作者头像 李华