iii-sdk(Node.js/TypeScript)实战指南:连接 iii 引擎、注册函数与触发器、调用工作流
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
本指南以iii-sdk(Node.js/TypeScript 版)为核心,讲解如何将业务代码注册为可被 iii 引擎实时编排的函数(Function)、如何通过触发器(Trigger)把 HTTP、cron、队列等外部事件绑定到函数,以及如何使用trigger()以同步、异步(enqueue)与 fire-and-forget 三种方式发起调用。读完本文,你将能够独立完成一个 Worker 的注册、函数/触发器声明、调用与错误处理,并理解其背后基于 WebSocket 的协议与重连机制。
iii-sdk是 iii 项目的官方 Node.js SDK(Apache-2.0 许可,ESM 模块,详见 sdk/packages/node/iii/package.json),源码位于 sdk/packages/node/iii/src。它与 sdk/packages/python、sdk/packages/rust 等 SDK 保持接口与行为对齐。
安装
使用任一包管理器安装即可:
pnpm add iii-sdk # 或 npm install iii-sdk从 sdk/packages/node/iii/package.json 的exports字段可以看出,包除了主入口外,还按需导出了./stream、./state、./channel、./trigger、./runtime、./errors、./engine、./protocol、./helpers、./internal等子模块,均提供import/require双格式构建产物。其运行时依赖极简:@iii-dev/helpers(本仓库工作区包)、@opentelemetry/api(遥测)与ws(WebSocket 客户端)。
快速上手:Hello World
官方 README 给出的最小示例(见 sdk/packages/node/iii/README.md)覆盖了“注册函数 → 注册触发器 → 本地调用”的完整闭环:
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' } })流程可以拆解为四步:
registerWorker(url):建立与引擎的 WebSocket 连接,返回一个IIIClient实例(内部为Sdk类,见 sdk/packages/node/iii/src/iii.ts);registerFunction(id, handler):把本地异步函数注册为可被按名调用的函数;registerTrigger({ type, function_id, config }):把外部事件源(这里是 HTTP 路径/greet)绑定到目标函数,事件发生时引擎自动触发调用;await iii.trigger({ function_id, payload }):主动发起一次同步调用并等待结果。
注意:registerWorker的构造函数会立刻执行this.connect()(iii.ts),因此连接是自动建立、无需手动connect的;未连接成功前发送的消息会被缓存进messagesToSend队列,在onSocketOpen时统一冲刷(iii.ts)。
registerWorker:引擎地址解析与初始化选项
地址解析优先级
registerWorker(address?, options?)的第一个参数是引擎的 WebSocket 地址。从源码 resolveAddress 可以看到严格的解析顺序:
- 显式传入的
address(永远最高优先级); - 环境变量
III_URL(由拉起 Worker 的编排方设置,例如iii compose、容器运行时、systemd); - 兜底默认值
DEFAULT_ENGINE_URL = 'ws://127.0.0.1:49134'。
一个值得注意的细节是:默认地址刻意写成 IPv4 回环127.0.0.1而非localhost,因为localhost在部分主机上会解析为::1(IPv6),而引擎可能只监听 IPv4,导致连接失败(iii.ts)。因此本地联调时建议显式使用ws://127.0.0.1:49134或ws://localhost:49134均可,但以源码注释为准更稳妥。
InitOptions 详解
registerWorker的第二个参数options类型为InitOptions(iii.ts),常用字段如下:
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
workerName | string | hostname:pid | Worker 显示名;非空的III_WORKER_NAME环境变量会覆盖它(编排方分配的标识优先,见 resolveWorkerName) |
namespace | string | 无(引擎使用default) | Worker 所属命名空间;解析顺序为options.namespace→III_NAMESPACE→ undefined |
workerDescription | string | 无 | 一行人类/LLM 可读的 Worker 说明,会呈现在engine::workers::list/engine::workers::info中 |
enableMetricsReporting | boolean | true | 是否通过 OpenTelemetry 上报 Worker 指标 |
invocationTimeoutMs | number | 30000 | worker.trigger()默认调用超时(毫秒),单次调用可用timeoutMs覆盖 |
reconnectionConfig | Partial<IIIReconnectionConfig> | 见下表 | WebSocket 断线重连策略 |
otel | Omit<OtelConfig, 'engineWsUrl'> | 自动初始化 | OpenTelemetry 配置;{ enabled: false }或环境变量OTEL_ENABLED=false/0/no/off可关闭 |
headers | Record<string, string> | 无 | WebSocket 握手阶段发送的自定义 HTTP 头 |
官方文档示例给出了典型用法:
const worker = registerWorker('ws://localhost:49134', { workerName: 'my-worker', invocationTimeoutMs: 10000, reconnectionConfig: { maxRetries: 5 }, })命名空间(namespace)语义
命名空间是 iii 中隔离注册与调用范围的关键概念,SDK 对其做了非常严谨的处理(resolveNamespace):
- 解析顺序:
options.namespace→process.env.III_NAMESPACE→ undefined;为 undefined 时引擎应用其default命名空间; - 显式传入空白字符串的 namespace 会直接抛错——"声明了命名空间却不给它名字" 与 "不声明命名空间" 含义相反,前者会让整个项目静默注册到错误的命名空间且运维者无法从声明中察觉;
- 环境变量为空白(如
III_NAMESPACE=${NS}且 NS 未设置)则按"未设置"处理,这是 shell 表达未设置变量的惯用方式,SDK 尊重这一语义。
命名空间不仅作用于注册:Worker 及其函数注册在哪个命名空间,它后续的trigger调用与registerTrigger绑定就默认跟随哪个命名空间(除非调用时显式指定别的)。此外有一个专门的处理函数 invocationNamespace:引擎内建函数(engine::前缀)的隐式调用固定落在default命名空间,避免把引擎内建泄漏进 Worker 命名空间。
重连与心跳(连接可靠性)
SDK 内置了一套与 Rust SDK 对齐的连接可靠性机制,常量定义在 sdk/packages/node/iii/src/iii-constants.ts:
| 常量 | 默认值 | 含义 |
|---|---|---|
WS_HANDSHAKE_TIMEOUT_MS | 10000 | WebSocket 握手超时 |
WS_PING_INTERVAL_MS | 20000 | 客户端心跳 ping 间隔 |
WS_IDLE_TIMEOUT_MS | 60000 | 超过该时长无任何入站帧(message/ping/pong)则强制断开以触发重连 |
DEFAULT_INVOCATION_TIMEOUT_MS | 30000 | 调用默认超时 |
重连策略IIIReconnectionConfig(iii-constants.ts)采用指数退避 + 抖动:
| 字段 | 默认值 | 说明 |
|---|---|---|
initialDelayMs | 1000 | 起始延迟(毫秒) |
maxDelayMs | 30000 | 延迟上限 |
backoffMultiplier | 2 | 指数退避倍数 |
jitterFactor | 0.3 | 随机抖动因子 0-1 |
maxRetries | -1 | 最大重试次数,-1表示无限 |
连接状态通过getConnectionState()获取,可能的取值(iii-constants.ts)为:disconnected、connecting、connected、reconnecting、failed。其中failed是终态——它只出现在引擎致命拒绝注册之后,SDK 不会对该状态继续重连。
优雅关闭
进程退出时应当调用shutdown()释放资源(Sdk.shutdown):它会关闭 OpenTelemetry、清除重连与心跳定时器、以Error('iii is shutting down')拒绝所有挂起中的调用、最后关闭 WebSocket。官方建议配合信号处理:
process.on('SIGTERM', async () => { await worker.shutdown() process.exit(0) })注册函数:本地处理器与 HTTP 外部函数
registerFunction(functionId, handlerOrInvocation, options?)支持两种形态(Sdk.registerFunction)。
形态一:本地异步处理器
官方 README 示例:
iii.registerFunction('orders::create', async (input) => { return { status_code: 201, body: { id: '123', item: input.body.item } } })处理器类型RemoteFunctionHandler的签名是(data: TInput, metadata?: JsonValue) => Promise<TOutput>(sdk/packages/node/iii/src/types.ts):第一个参数是调用负载,第二个是可选的逐调用metadata(任意 JSON,随调用独立于 payload 传输,未附加时为undefined)。已有的单参数处理器完全兼容,多余参数会被忽略。
注意几个 SDK 强制的约束:
functionId为空字符串或纯空白会抛出id is required;- 同一个
functionId重复注册会抛出function id already registered: <id>(本地 Map 校验,iii.ts); - 返回值是一个
FunctionRef(含id与unregister()),可用于后续注销函数。
形态二:HTTP 外部函数(代理到远端)
除了本地处理器,还可以传入HttpInvocationConfig,把函数代理到外部 HTTP 服务(如 AWS Lambda、Cloudflare Workers 等),引擎侧负责转发调用。字段包括url、method(默认POST)、timeout_ms、headers、auth(如{ type: 'bearer', token_key: 'LAMBDA_AUTH_TOKEN' })。从 iii.ts 可以看到,注册消息中会携带完整的invocation配置块。
const lambdaRef = worker.registerFunction( 'external::my-lambda', { url: 'https://abc123.lambda-url.us-east-1.on.aws', method: 'POST', timeout_ms: 30_000, auth: { type: 'bearer', token_key: 'LAMBDA_AUTH_TOKEN' }, }, { description: 'Proxied Lambda function' }, )一个实现细节:对于 HTTP 形态注册的函数,SDK 不保存本地 handler。若引擎反向路由调用到该函数,SDK 会回传错误码function_not_invokable("Function is HTTP-invoked and cannot be invoked locally");函数根本不存在则回传function_not_found(onInvokeFunction)。
注册触发器:把外部事件绑定到函数
registerTrigger({ type, function_id, config })返回一个带unregister()的Trigger句柄(Sdk.registerTrigger)。官方 README 示例:
iii.registerTrigger({ type: 'http', function_id: 'orders::create', config: { api_path: '/orders', http_method: 'POST' }, })引擎内置的触发器类型包括http(HTTP 路由)、cron(定时调度,配置形如{ expression: '0 */5 * * * * *' })、queue(队列消费)以及durable:subscriber(可靠订阅)等,具体以引擎支持的触发器类型为准。触发器的type决定config的结构。
注册时的命名空间语义同样严格(iii.ts):未显式设置namespace时,触发器默认落在本 Worker 的命名空间,而非引擎的default。原因是触发器指向一个函数,而该函数注册在本 Worker 的命名空间——若默认落到default,触发器触发后就会解析不到目标函数。想绑定到其他命名空间(包括default)需要显式声明。
注销触发器:
const trigger = worker.registerTrigger({ type: 'cron', function_id: 'my-service::process-batch', config: { expression: '0 */5 * * * * *' }, }) // 稍后移除 trigger.unregister()自定义触发器类型:registerTriggerType
如果内置触发器类型不够用,SDK 允许向引擎注册自定义触发器类型,扩展"外部事件 → 函数调用"的映射。registerTriggerType(triggerType, handler)的 handler 需实现registerTrigger/unregisterTrigger两个回调(接口定义见 sdk/packages/node/iii/src/triggers.ts):
type CronConfig = { expression: string } worker.registerTriggerType<CronConfig>( { id: 'cron', description: 'Fires on a cron schedule' }, { async registerTrigger({ id, function_id, config }) { startCronJob(id, config.expression, () => worker.trigger({ function_id, payload: {} }), ) }, async unregisterTrigger({ id }) { stopCronJob(id) }, }, )TriggerConfig<TConfig>(triggers.ts)包含id、function_id、config、可选的metadata以及namespace——这个字段由 SDK 在注册时从 Worker 命名空间填充,触发方(provider)后续调用trigger()时必须透传,否则可能落到错误的命名空间。
registerTriggerType返回的TriggerTypeRef<TConfig>(types.ts)还提供了三个便捷方法:
registerTrigger(functionId, config, metadata?):绑定一个该类型的触发器;registerFunction(functionId, handler, config, metadata?):注册函数并立即绑定触发器(一次调用完成两件事,且自动把触发器命名空间默认到 Worker 命名空间,避免函数与触发器分离导致解析失败);unregister():注销整个触发器类型。
调用函数:同步、入队与 fire-and-forget
trigger(request)是唯一的调用入口,其路由行为由action字段决定。官方文档整理的三种模式如下:
action | 行为 | 返回类型 |
|---|---|---|
| (不传) | 同步:等待函数返回 | Promise<TOutput> |
TriggerAction.Enqueue({ queue }) | 经具名队列异步路由;引擎确认入队 | Promise<EnqueueResult>(含messageReceiptId) |
TriggerAction.Void() | fire-and-forget,无响应 | Promise<undefined> |
import { registerWorker, TriggerAction } from 'iii-sdk' const iii = registerWorker('ws://localhost:49134') // 同步调用,等待结果 const result = await iii.trigger({ function_id: 'orders::create', payload: { item: 'widget' } }) // fire-and-forget,不等待 iii.trigger({ function_id: 'analytics::track', payload: { event: 'page_view' }, action: TriggerAction.Void(), }) // 经命名队列异步处理 const { messageReceiptId } = await iii.trigger({ function_id: 'payments::charge', payload: { orderId: '123', amount: 49.99 }, action: TriggerAction.Enqueue({ queue: 'payment' }), })从 Sdk.trigger 的源码看,三种模式的实现差异明显:
- Void:直接发
InvokeFunction消息,不生成invocation_id、不注册 pending 调用、不等待响应,立即返回undefined; - 同步 / Enqueue:生成
invocation_id(crypto.randomUUID()),注册进invocationsMap 并用effectiveTimeout(timeoutMs ?? invocationTimeoutMs)启动超时定时器,超时则以InvocationError{ code: 'TIMEOUT' }拒绝; - 每次调用都会注入
traceparent/baggage(OpenTelemetry 上下文传播),供链路追踪使用。
TriggerAction构造器本身在 iii.ts 定义,并有一组专门的契约测试锁死其线上格式(sdk/packages/node/iii/tests/trigger-action.test.ts):TriggerAction.Enqueue({ queue: 'orders' })必须序列化为{ type: 'enqueue', queue: 'orders' },TriggerAction.Void()必须序列化为{ type: 'void' }——这是引擎侧TriggerAction反序列化所依赖的协议约定。
关于 Enqueue 的两个前提(源码注释明确说明):目标队列需要由worker-compose.yaml中的 queue worker 声明(queue_configs),否则触发会被引擎以enqueue_error(无队列提供者)拒绝。
错误处理与连接诊断
SDK 把调用失败统一收口为两个带类型码的Error子类(sdk/packages/node/iii/src/errors.ts):
InvocationError:trigger()失败时的统一错误类型,字段含code、message、function_id、stacktrace。覆盖三类失败:调用被拒绝(如 RBACFORBIDDEN)、handler 层失败(引擎回传invocation_failed,附带调用栈)、超时(TIMEOUT)。message 统一格式化为`${code}: ${message}`,避免旧版"拒绝值打印成[object Object]"的问题。未知形状的错误会被包装为code: 'UNKNOWN'(见 toInvocationError);RegistrationRejectedError:引擎拒绝 Worker 身份注册时的致命错误,字段含code、namespace、worker_name、function_id、owner_worker_id。触发后连接进入终态failed,不再重连。
注册拒绝(registrationrejected)的两种典型码(iii.ts):
| 错误码 | 场景 | 严重性 |
|---|---|---|
WORKER_NAMESPACE_CONFLICT | 另一个存活的 Worker 已持有同一(namespace, worker_name),引擎关闭连接 | 致命,SDK 停止且不重连 |
FUNCTION_NAMESPACE_CONFLICT | 同命名空间内另一 Worker 已导出同名函数,仅该注册被拒 | 非致命,仅告警,Worker 继续服务其他函数 |
诊断接口(与 Python/Rust SDK 对齐):
getConnectionState():当前连接状态;getAddress():实际解析到的引擎地址;getFatalError():致命注册拒绝的RegistrationRejectedError,健康时为undefined。
API 速查表
官方 README 整理的 API 一览(完整继承自 sdk/packages/node/iii/README.md):
| 操作 | 签名 | 说明 |
|---|---|---|
| 初始化 | registerWorker(url, options?) | 创建并连接引擎,返回ISdk实例 |
| 注册函数 | iii.registerFunction(id, handler, options?) | 注册一个可按名调用的函数 |
| 注册触发器 | iii.registerTrigger({ type, function_id, config }) | 把触发器(HTTP、cron、queue 等)绑定到函数 |
| 调用(等待) | await iii.trigger({ function_id, payload }) | 调用函数并等待结果 |
| 调用(fire-and-forget) | iii.trigger({ function_id, payload, action: TriggerAction.Void() }) | 调用但不等待 |
| 调用(入队) | iii.trigger({ function_id, payload, action: TriggerAction.Enqueue({ queue }) }) | 经具名队列路由调用 |
源码导读与延伸阅读
如果想深入理解 SDK 与引擎的交互,推荐按以下顺序阅读本仓库源码:
- sdk/packages/node/iii/src/iii.ts:
Sdk类全部实现——地址/命名空间解析、消息收发、重连与心跳、注册拒绝处理、TriggerAction构造器与registerWorker入口; - sdk/packages/node/iii/src/iii-constants.ts:引擎内建函数路径(
engine::functions::list、engine::workers::list、engine::workers::register等)与所有连接/重连/超时常量; - sdk/packages/node/iii/src/types.ts:
IIIClient公共接口、Trigger/FunctionRef/TriggerTypeRef句柄类型、流式请求响应类型; - sdk/packages/node/iii/src/errors.ts:
InvocationError与RegistrationRejectedError的定义与线上错误体识别; - sdk/packages/node/iii/tests/trigger-action.test.ts:
TriggerAction线上格式契约测试; - sdk/packages/node/iii/README.md:官方 SDK 文档(本文的基础)。
此外,SDK 测试目录还覆盖了连接握手超时、心跳、断线重连(reattach)、命名空间继承、RBAC Worker、流(stream)、状态(state)、发布订阅(pubsub)等大量行为测试(见 sdk/packages/node/iii/tests),是理解 SDK 边界行为的绝佳参考。包级配置(构建、测试、子模块导出)可查阅 sdk/packages/node/iii/package.json 与 sdk/packages/node/iii/vitest.config.ts。
【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考