ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

iii 项目 Node.js SDK 完整参考:从 Worker 注册到流式通道与可观测性

iii 项目 Node.js SDK 完整参考:从 Worker 注册到流式通道与可观测性 iii 项目 Node.js SDK 完整参考从 Worker 注册到流式通道与可观测性【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii本指南以 iii 项目 Node.js SDK发布名iii-sdk源码位于 sdk/packages/node/iii为核心系统讲解如何将一个 Node.js Worker 连接到运行中的 iii 引擎注册可调用函数、绑定触发器、定义自定义触发器类型、跨 Worker 触发调用、使用流式通道与结构化日志并理解错误码、连接状态与底层 WebSocket 协议。读完本文你将能够基于这套 API 独立编写、注册、触发和优雅关闭一个生产可用的 Node.js Worker并具备排查InvocationError与连接故障的源码级认知。安装Node.js SDK 通过 npm 安装npm install iii-sdk包名iii-sdk在 sdk/packages/node/iii/package.json 中声明为 ESM 优先type: module同时通过exports字段提供.,./helpers,./stream,./state,./channel,./trigger,./errors等多个子路径入口支持import与require两种加载方式。运行环境需要 Node.js 与ws库WebSocket 实现运行时依赖仅ws与opentelemetry/apiOpenTelemetry 相关能力由同仓库的iii-dev/helpers工作区包提供。连接引擎registerWorkerregisterWorker将一个 Worker 连接到正在运行的 iii 引擎并返回其句柄function registerWorker(address: string, options?: InitOptions): ISdk;address传入引擎的 SDK WebSocket URL例如环境变量process.env.III_URLoptions配置 Worker 身份、超时、重连与 OpenTelemetry。返回的ISdk携带本文介绍的全部方法。从源码 sdk/packages/node/iii/src/iii.ts 看地址解析遵循明确的优先级resolveAddress显式传入的address参数始终优先否则读取环境变量III_URL由iii compose、容器运行时、systemd 等 supervisor 在拉起进程时设置与III_NAMESPACE、III_WORKER_NAME一套兜底常量DEFAULT_ENGINE_URL ws://127.0.0.1:49134。源码特意把 IPv4 回环地址写死而非使用localhost原因在注释中说明某些主机上localhost会解析为::1而引擎只监听 IPv4会导致连接失败。这一点在排查“明明引擎在跑却连不上”时非常关键。InitOptions 关键配置配置项类型默认值说明workerNamestringhostname:pidWorker 显示名非空III_WORKER_NAME环境变量会覆盖它Compose 托管 Worker 由 supervisor 赋名引擎按名字匹配在线注册故托管名优先namespacestring见解析顺序解析顺序options.namespace→III_NAMESPACE→ undefined引擎归入default命名空间空字符串会被拒绝并抛错workerDescriptionstring无一行人类/LLM 可读的 Worker 摘要出现在engine::workers::list/engine::workers::infoenableMetricsReportingbooleantrue通过 OpenTelemetry 上报 Worker 指标invocationTimeoutMsnumber30000worker.trigger()的默认调用超时毫秒reconnectionConfigPartialIIIReconnectionConfig见下文WebSocket 重连行为otelOmitOtelConfig, engineWsUrl自动初始化OpenTelemetry 配置设{ enabled: false }或环境变量OTEL_ENABLEDfalse/0/no/off可关闭headersRecordstring, string无WebSocket 握手时携带的自定义 HTTP 头超时、心跳与重连的默认常量全部集中在 sdk/packages/node/iii/src/iii-constants.tsexport const DEFAULT_BRIDGE_RECONNECTION_CONFIG { initialDelayMs: 1000, // 起始重连延迟 1s maxDelayMs: 30000, // 最大延迟上限 30s backoffMultiplier: 2, // 指数退避倍数 jitterFactor: 0.3, // 抖动因子 0-1避免惊群 maxRetries: -1, // -1 表示无限重试 } export const DEFAULT_INVOCATION_TIMEOUT_MS 30000 export const WS_HANDSHAKE_TIMEOUT_MS 10000 // 握手超时 export const WS_PING_INTERVAL_MS 20000 // 客户端心跳 ping 周期 export const WS_IDLE_TIMEOUT_MS 60000 // 60s 无入站帧则强制重连注意scheduleReconnect的实现细节重连延迟按initialDelayMs * backoffMultiplier^attempt指数增长、以maxDelayMs封顶并叠加jitterFactor的随机抖动maxRetries达到上限后连接状态置为failed并放弃。此外 SDK 启动后会以WS_PING_INTERVAL_MS周期发送 ping并在WS_IDLE_TIMEOUT_MS内未收到任何入站帧时主动terminate()连接以强制重连——这些行为与 Rust SDK 保持一致性源码注释明确标注 “parity with Rust SDK”。registerWorker构造完成即自动建立 WebSocket 连接new Sdk(resolveAddress(address), options)中直接调用this.connect()无需额外手动连线的步骤。注册可调用函数registerFunctionregisterFunction在当前 Worker 上注册一个可调用函数worker.registerFunction( functionId: string, handlerOrInvocation: RemoteFunctionHandler | HttpInvocationConfig, options?: RegisterFunctionOptions, ): FunctionRef;options接受description、metadata以及可选的request_format/response_formatJSON Schema随函数一起存储供 iii 控制台与 Agent 可读的 skills 使用。RegisterFunctionOptions对应消息类型RegisterFunctionMessagesdk/packages/node/iii/src/iii-types.ts字段包括description、request_format、response_format、metadata、invocation。其中request_format/response_format是RegisterFunctionFormat类型支持name、description、typestring | number | boolean | object | array | null | map | integer、properties、items、required等描述字段。本地处理器处理器签名在 sdk/packages/node/iii/src/types.ts 中定义为RemoteFunctionHandlertype RemoteFunctionHandlerTInput any, TOutput any ( data: TInput, metadata?: JsonValue, ) PromiseTOutputmetadata是独立于 payload 的任意 JSON由调用方在trigger请求中附带未附带时为undefined既有单参数处理器不受影响忽略多余参数即可。const fn worker.registerFunction( greet, async (input: { name: string }) { return { message: Hello, ${input.name}! } }, { description: Greets a user }, )返回的FunctionRef携带id与unregister()。源码对注册做了两处防御functionId为空字符串直接抛错同一functionId重复注册抛function id already registered。unregister()会向引擎发送UnregisterFunction消息并从本地 Map 移除。HTTP 远程函数第二个参数也可以传HttpInvocationConfig而非本地函数将函数代理到外部 HTTP 端点Lambda、Cloudflare Workers 等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 }, )源码中isHandler判断区分两种注册函数类型时携带本地handlerHTTP 配置时构造invocation字段url、method默认POST、timeout_ms、headers、auth并仅存储消息。这类函数被调用时会收到function_not_invokable错误因为引擎会直接走 HTTP 代理而非本地执行。执行跟踪与负载观测从源码看本地处理器的实际执行被包了一层内部包装iii.ts中this.functions.set处的handler包装默认会把输入/输出 payload 脱敏截断后记录为 span 事件iii.invocation.input/iii.invocation.output含iii.payload.truncated标记并包裹在命名为execute functionId的SpanKind.INTERNALspan 中。引擎会抑制自己的call fn包装 span因此execute命名可读为干净的内部子 span。相关测试见 sdk/packages/node/iii/tests/payload.test.ts 与 sdk/packages/node/iii/tests/span-ops.test.ts。绑定触发器registerTriggerregisterTrigger将一个已注册函数绑定到已配置的触发器实例worker.registerTrigger(trigger: RegisterTriggerInput): Trigger;const trigger worker.registerTrigger({ type: http, function_id: greet, config: { api_path: /greet, http_method: GET }, }) // 稍后移除 trigger.unregister()返回的Trigger携带运行时句柄通过Trigger.unregister()移除SDK 没有顶层unregisterTrigger函数。RegisterTriggerMessage见 iii-types.ts的字段包括type触发器类型标识如http、cron、durable:subscriber、function_id、config触发器类型期望的配置结构、可选的metadata、namespace目标函数解析的命名空间与trigger_namespace触发器类型 provider 所在命名空间。源码实现有两点值得注意触发器id由 SDK 用crypto.randomUUID()生成不在调用方传入命名空间默认语义registerTrigger会把未显式指定的namespace填充为当前 Worker 的命名空间而不是引擎的default——因为触发器指向的函数是当前 Worker 注册的落在 Worker 自己的命名空间若默认到default会“触发器触发了却解析不到函数”。想绑定其他命名空间含default必须显式声明。触发器的生命周期消息底层通过MessageType.RegisterTrigger/MessageType.UnregisterTrigger消息与引擎交互toWireFormat还会把 SDK 内部的type字段转换为线格式的trigger_type保证与引擎协议一致。声明自定义触发器类型registerTriggerType / unregisterTriggerTyperegisterTriggerType声明一个新的触发器类型由当前 Worker 对外提供其他 Worker 即可把函数绑定到该类型上worker.registerTriggerTypeTConfig( triggerType: RegisterTriggerTypeInput, handler: TriggerHandlerTConfig, ): TriggerTypeRefTConfig;TriggerHandler需要实现registerTrigger/unregisterTrigger两个回调对应引擎下发绑定与解绑指令type CronConfig { expression: string } worker.registerTriggerTypeCronConfig( { 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) }, }, )返回的TriggerTypeRefTConfig提供三个便捷方法registerTrigger(functionId, config, metadata?)注册绑定到该类型的触发器默认把触发器命名空间设为当前 Worker 的命名空间源码注释说明否则函数落在 Worker 命名空间、触发器落在default永远无法相互解析registerFunction(functionId, handler, config, metadata?)注册函数并一步完成触发器绑定unregister()撤销该触发器类型。unregisterTriggerType(triggerType)移除先前注册的触发器类型向引擎发送UnregisterTriggerType消息并清理本地 Map。当引擎侧把某个触发器的绑定/解绑指令路由到本 Worker 时SDK 的onRegisterTrigger/onUnregisterTrigger会调用对应 handler 回调并把结果以TriggerRegistrationResult回执给引擎成功或携带trigger_registration_failed/trigger_type_not_found错误码。这部分测试覆盖见 sdk/packages/node/iii/tests/trigger-type-lifecycle.test.ts 与 sdk/packages/node/iii/tests/service-triggers.test.ts。触发调用worker.triggertrigger调用一个已注册函数worker.triggerTInput, TOutput(request: TriggerRequestTInput): PromiseTOutput;返回语义取决于action字段源码 JSDoc 中的对照表action行为返回类型省略同步等待函数返回PromiseTOutputTriggerAction.Enqueue(...)经命名队列异步路由引擎确认入队PromiseEnqueueResultTriggerAction.Void()即发即忘无响应PromiseundefinedTriggerRequest的字段见 iii-types.tsfunction_id、payload、可选action、timeoutMs覆盖默认调用超时、metadata任意 JSON作为第二个参数透传给目标 handler、namespace。import { TriggerAction } from iii-sdk // 同步调用 const result await worker.trigger({ function_id: get-order, payload: { id: 123 } }) // 入队需要 engine 侧 queue 声明 const { messageReceiptId } await worker.trigger({ function_id: payments::charge, payload: { orderId: 123, amount: 49.99 }, action: TriggerAction.Enqueue({ queue: payment }), }) // 即发即忘 worker.trigger({ function_id: notifications::send, payload: { userId: 123 }, action: TriggerAction.Void(), })源码实现细节iii.ts 中trigger方法Void动作直接发送InvokeFunction消息且不带invocation_id不等待任何响应立即 resolveundefined同步与 Enqueue生成crypto.randomUUID()作为invocation_id注册到this.invocationsMap用setTimeout实现超时引擎返回InvocationResult时按invocation_id匹配 resolve/reject超时错误码固定为TIMEOUT消息形如invocation timed out after ${effectiveTimeout}ms命名空间解析由invocationNamespace完成显式request.namespace优先未显式时以engine::开头的内置函数解析到default其余解析到当前 Worker 命名空间避免把引擎内置函数泄漏进 Worker 命名空间每次调用自动注入traceparent与baggage用于跨进程链路追踪。优雅关闭shutdownworker.shutdown(): Promisevoid;shutdown断开与引擎的连接并释放资源同时冲刷待发送的可观测性数据。源码中的关闭顺序shutdown方法置isShuttingDown true停止指标上报stopMetricsReporting关闭 OpenTelemetryawait shutdownOtel()确保 span/日志导出完成清理重连定时器与心跳clearReconnectTimeout/stopHeartbeat以iii is shutting down拒绝所有在途调用并清空invocations移除 WebSocket 全部监听器、挂一个空的error监听兜底后close()状态置为disconnected。典型用法配合进程信号process.on(SIGTERM, async () { await worker.shutdown() process.exit(0) })触发器动作TriggerActionTriggerAction是一个运行时常量对象iii.ts 末尾导出产生传给trigger的action字段值TriggerAction.Void(); // fire-and-forget TriggerAction.Enqueue({ queue: math }); // route through iii-queue底层类型是判别联合TriggerActionTypetype TriggerActionType { type: enqueue; queue: string } | { type: void }源码定义export const TriggerAction { Enqueue: (opts: { queue: string }) ({ type: enqueue as const, ...opts }), Void: () ({ type: void as const }), } as const注意Enqueue的 JSDoc 提醒入队路由要求worker-compose.yaml中声明了queueworker否则trigger会以enqueue_error无 queue provider拒绝。EnqueueResult类型从iii-dev/helpers/queue重导出包含messageReceiptId等字段。错误类型IIIInvocationErrorIIIInvocationError extends Error。所有跨越 SDK 边界的失败都以该类的实例抛出携带code: string、message: string、可选function_id与可选stacktraceclass IIIInvocationError extends Error { code: string; message: string; function_id?: string; stacktrace?: string; }实际源码中的类名为InvocationErrorsdk/packages/node/iii/src/errors.ts构造时把${code}: ${message}作为 Error message并保留code、function_id、stacktrace字段。此外该文件还导出RegistrationRejectedError引擎拒绝注册时的致命错误含code、namespace、worker_name、function_id、owner_worker_idisErrorBody类型守卫判断引擎线格式{ code, message, stacktrace? }。常见错误码文档列出的错误码与源码行为对应如下code触发场景来源invocation_failed目标 handler 抛异常服务端回包onInvokeFunction捕获后回传invocation_stopped引擎超时终止调用引擎function_not_found目标函数不存在本地回包fn为 undefinedfunction_not_invokable函数是 HTTP 远程函数、无法本地执行本地回包TIMEOUT客户端侧超时SDK 本地setTimeoutFORBIDDENRBAC 拒绝引擎另外onRegistrationRejected区分两种注册冲突WORKER_NAMESPACE_CONFLICT同一(namespace, worker_name)已被其他存活 Worker 占用引擎断开连接致命进入failed状态且不重连FUNCTION_NAMESPACE_CONFLICT同命名空间下函数 id 已被占用非致命仅拒绝该函数注册、连接保持。可通过getFatalError()轮询获取终止原因。toInvocationError包装逻辑保证所有拒绝路径都返回可读的 Error已是 Error 的直接透传符合线格式ErrorBody的包装为InvocationError其余含非字符串对象兜底用UNKNOWN码避免String(err) [object Object]。相关测试见 sdk/packages/node/iii/tests/errors.test.ts。流式通道ChannelReader 与 ChannelWriterChannelReader与ChannelWriter是运行时类包装引擎的流式 WebSocket用于 Worker 间数据流。StreamChannelRef是在 SDK 调用间传递以标识通道的类型type StreamChannelRef { channel_id: string; access_key: string; direction: read | write; };每个端用引擎的 WS base URL 与一个StreamChannelRef构造ChannelReader暴露 NodeReadable流另有.sendMessage()与.onMessage()ChannelWriter暴露Writable流另有.sendMessage()、.sendChunked()与.close()。源码 sdk/packages/node/iii/src/channels.ts 的关键实现ChannelDirection常量镜像 Rust SDK 的ChannelDirection枚举read/writeChannelItem提供Text(string)/Binary(Uint8Array)工厂与判别式ChannelWriter用FRAME_SIZE 64 * 1024分块发送二进制sendChunked内部维护pendingMessages队列在 WebSocket 尚未open时缓冲、open 后按序冲刷Writable.final关闭时会延迟约 10ms再发 close 帧注释解释若不延迟close 帧可能先于已缓冲的数据帧到达引擎造成数据截断对应测试 sdk/packages/node/iii/tests/channel-close-delay.test.tsSDK 内部的__helpers_create_channel通过调用引擎内置函数engine::channels::create创建通道对返回含writer/reader/writerRef/readerRef的Channel——ref可序列化后作为调用数据字段传给其他 Worker实现跨进程通道传递。高级封装createChannel位于iii-sdk/helpers子模块import { createChannel } from iii-sdk/helpers const channel await createChannel(iii) // 写二进制流 channel.writer.stream.write(Buffer.from(hello)) channel.writer.stream.end() // 或发送结构化文本消息 channel.writer.sendMessage(JSON.stringify({ type: event, data: test })) channel.writer.close()日志LoggerLogger是运行时类提供info、warn、error、debug四个方法签名均为(message: string, data?: unknown) void。输出与 SDK 的 OpenTelemetry 集成——即通过 OTel 日志管道SeverityNumber分级导出而非直接打 console。SDK 内部自身也复用这套日志设施logError/logWarn在 OTel logger 可用时走otelLogger.emit(...)否则回退console.error/console.warn前缀[iii]。导出端由同仓库sdk/packages/node/observabilityiii-observability负责包含 span exporter、log exporter、metrics exporter、fetch instrumentation 等实现。信息类型Info typesSDK 重导出引擎在列出系统状态时返回的结构化类型FunctionInfofunction_id可选description、request_format/response_format、metadataTriggerInfoid、trigger_type、function_id可选config、metadataWorkerInfoid、name、runtime/version/OS 字段、IP、status、connected_at_ms、function_count、已注册functions、active_invocations、可选isolation。引擎侧对应的内置函数路径定义在 iii-constants.ts 的EngineFunctionsengine::functions::list/info、engine::workers::list/info、engine::triggers::list/info覆盖触发器类型模板、engine::registered-triggers::list/info覆盖触发器实例行。这些内置函数同样通过worker.trigger()调用。WorkerMetadata不属于本 SDKWorker 侧元数据请使用WorkerInfo。线协议MessageTypeMessageType是运行时枚举命名 SDK 与引擎交换的每一种线帧。完整枚举值iii-types.tsexport enum MessageType { RegisterFunction registerfunction, UnregisterFunction unregisterfunction, InvokeFunction invokefunction, InvocationResult invocationresult, RegisterTriggerType registertriggertype, RegisterTrigger registertrigger, UnregisterTrigger unregistertrigger, UnregisterTriggerType unregistertriggertype, TriggerRegistrationResult triggerregistrationresult, WorkerRegistered workerregistered, RegistrationRejected registrationrejected, Reattach reattach, }调用方很少直接使用它它主要出现在中间件钩子与协议级自定义代码中。例如收到WorkerRegistered时 SDK 保存引擎分配的worker_id与reattach_token收到RegistrationRejected时走onRegistrationRejected分支见上文错误码一节。连接状态机连接状态字面量联合disconnected | connecting | connected | reconnecting | failed在当前构建中属于内部实现连接建立以registerWorker返回为准失败会在首次 SDK 调用时抛出。对应类型为IIIConnectionStateiii-constants.ts。不过从源码看SDK 实例仍暴露了三个只读诊断方法IIIClient接口types.tsgetConnectionState()当前连接状态failed为终态紧随致命的注册拒绝之后getAddress()解析到的引擎地址显式参数 →III_URL→ 默认地址与 Rust SDK 的address()、Python SDK 的get_address()对应getFatalError()导致连接终止的致命注册拒绝如WORKER_NAMESPACE_CONFLICT健康时为undefined与 Python/Rust SDK 的_fatal_error/fatal_error()对应。重连与 Reattach重连时onSocketOpen分支SDK 会先发送Reattach消息携带previous_worker_id与reattach_token让引擎先退役旧连接、再回放注册——否则旧连接仍持有(namespace, worker_name)元数据声明会撞上自身的WORKER_NAMESPACE_CONFLICT。之后依次声明 Worker 元数据registerWorkerMetadata经engine::workers::register上报 runtime/version/name/os/pid/namespace/telemetry 信息对应测试 sdk/packages/node/iii/tests/register-worker-metadata.test.ts→ 回放触发器类型 → 函数 → 触发器 → 冲刷messagesToSend中积压的调用消息其中已超时被移除的invocation_id会被跳过。完整示例组装一个 Worker综合上述 API一个最小但完整的 Node.js Worker 如下import { registerWorker, TriggerAction } from iii-sdk // 地址优先取显式参数其次 III_URL兜底 ws://127.0.0.1:49134 const worker registerWorker(process.env.III_URL, { workerName: math-worker, namespace: demo, invocationTimeoutMs: 10_000, reconnectionConfig: { maxRetries: 10, initialDelayMs: 500 }, }) // 1. 注册本地函数 worker.registerFunction( math::add, async (input: { a: number; b: number }) ({ sum: input.a input.b }), { description: Adds two numbers, request_format: { type: object, properties: { a: { type: number }, b: { type: number } }, required: [a, b] }, response_format: { type: object, properties: { sum: { type: number } } }, }, ) // 2. 绑定 HTTP 触发器 const trigger worker.registerTrigger({ type: http, function_id: math::add, config: { api_path: /add, http_method: POST }, }) // 3. 调用其他 Worker 的函数同步 / 入队 / 即发即忘 const result await worker.trigger{ id: string }, { name: string }({ function_id: orders::get, payload: { id: 123 }, }) await worker.trigger({ function_id: notify::send, payload: { userId: 123 }, action: TriggerAction.Void(), }) // 4. 优雅关闭 process.on(SIGTERM, async () { await worker.shutdown() process.exit(0) })运行前请确保 iii 引擎在线默认监听ws://127.0.0.1:49134且 Worker 与目标函数位于同一命名空间或调用时显式指定namespace。入队动作还需在worker-compose.yaml中声明对应 queue worker。进一步阅读Node SDK 源码sdk/packages/node/iii/src/iii.ts、sdk/packages/node/iii/src/iii-types.ts、sdk/packages/node/iii/src/iii-constants.ts、sdk/packages/node/iii/src/errors.ts、sdk/packages/node/iii/src/channels.ts、sdk/packages/node/iii/src/types.tsNode SDK 测试连接、触发器、通道、错误、OTel 全覆盖sdk/packages/node/iii/tests可观测性 helpers 与导出端sdk/packages/node/observability引擎内置函数与协议说明engine/src/protocol.rs控制台对函数元数据的呈现console/packages/console-frontend/src本文档对应的原始参考页位于 docs/0-17-0/sdk-reference/node-sdk.mdx其中标注为“手写快照、最终参考将从 SDK 源码生成”本文的源码级细节即与该生成方向保持一致。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表