iii Node.js Helpers 包详解http、stream、queue 与 OpenTelemetry 可观测性 API 全解析【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii本文基于 iii 仓库中iii-dev/helpers包的 API 参考文档系统讲解这个 Node.js/TypeScript 辅助包的五个子路径导出http、observability、queue、stream、worker-connection-managerHTTP 流式响应封装的用法与底层控制帧机制、OTel 初始化配置全参数及 WebSocket 上报链路、Stream 触发器配置与原子更新操作UpdateOp语义以及 Worker 连接 RBAC 鉴权类型。读完后你可以直接在 Worker 代码中使用这些 helpers 编写 HTTP 接口、接入分布式追踪与结构化日志并理解每个参数在源码中的真实行为。1. 安装与包结构npm install iii-dev/helpersiii-dev/helpers是 iii SDK 体系中跨语言共享的辅助原语包当前仓库中版本为0.23.0-rc.9见 package.json。它自身只依赖 OpenTelemetry SDK 系列包和wsWebSocket 客户端通过 5 个子路径对外暴露 API子路径导入方式职责httpimport { http } from iii-dev/helpers/httpHTTP 风格 handler 包装器及请求/响应/鉴权类型observabilityimport { initOtel, Logger } from iii-dev/helpers/observability结构化日志、OTel 初始化、span/baggage 工具、worker 资源指标queueimport { EnqueueResult } from iii-dev/helpers/queue队列Enqueue触发动作的结果类型streamimport { StreamTriggerConfig } from iii-dev/helpers/streamStream 触发器配置、变更事件、IO 输入与更新操作类型worker-connection-managerimport { AuthInput } from iii-dev/helpers/worker-connection-managerRBAC 鉴权与注册回调类型从 package.json 的exports字段看每个子路径都同时提供 ESM.mjs与 CJS.cjs产物TypeScript 类型文件独立发布此外还有一个./observability/internal内部子路径供 SDK 内部模块复用普通用户代码一般使用./observability即可。2. http把 Express 风格 handler 变成 iii 函数2.1 http() 函数签名与用法http是一个高阶函数把「分离req/res两个参数」的 HTTP 风格 handler 适配成 SDK 期望的函数 handler 格式可直接传给registerFunctionhttp(callback: (req: HttpStreamingRequest, res: HttpStreamingResponse) Promisevoid | HttpResponsenumber, string | Buffer | Recordstring, unknown) (req: HttpInternalRequest) Promisevoid | HttpResponsenumber, string | Buffer | Recordstring, unknown完整使用示例import { http } from iii-dev/helpers/http worker.registerFunction( my-api, http(async (req, res) { res.status(200) res.headers({ content-type: application/json }) res.stream.end(JSON.stringify({ hello: world })) res.close() }), )res提供四个能力res.status(statusCode)设置 HTTP 状态码res.headers(headers)设置响应头res.streamNode.jsWritableStream用于流式写响应体res.close()关闭响应通道结束请求。2.2 源码级原理控制帧协议从 http/index.ts 的实现看http()的核心是一个「协议适配层」。运行时SDK 核心实际传入的是HttpInternalRequest——它在流式请求字段之外多带了一个response: HttpStreamWriter含sendMessage、stream、close。helper 把这个内部 writer 拆包装成用户友好的HttpStreamingResponsestatus()会调用response.sendMessage(JSON.stringify({ type: set_status, status_code }))即通过 WebSocket 向引擎发送一个set_status控制帧headers()同理发送{ type: set_headers, headers }控制帧stream则直接暴露底层WritableStream供流式响应。也就是说iii 中 Worker 对外的 HTTP 响应并不是进程内直接写 socket而是「状态码/响应头走消息帧、响应体走流」的两段式协议。这也是为什么示例中先status、再headers、写stream.end(...)、最后close()的顺序是自然的调用模式。2.3 http 类型定义HttpAuthConfig—— HTTP 调用型函数的鉴权配置三种模式type HttpAuthConfig | { secret_key: string; type: hmac } // HMAC 签名验证共享密钥 | { token_key: string; type: bearer } // Bearer token | { header: string; type: api_key; value_key: string } // 自定义 header 中的 API keyHttpInvocationConfig—— 面向 Lambda、Cloudflare Workers 等 HTTP 调用型函数的配置名称类型必填说明urlstring是要调用的 URLmethodHttpMethod否HTTP 方法默认POSTtimeout_msnumber否超时时间毫秒headersRecordstring, string否随请求发送的自定义头authHttpAuthConfig否鉴权配置HttpMethod—— 注意它与引擎核心builtin_triggers中的 HTTP 方法枚举不同后者还覆盖HEAD/OPTIONS此处仅 5 种type HttpMethod GET | POST | PUT | PATCH | DELETEHttpRequest—— 函数 handler 接收到的缓冲式入站请求字段类型必填说明bodyTBody是已解析的请求体headersRecordstring, string \| string[]是请求头methodstring是HTTP 方法如GET、POSTpath_paramsRecordstring, string是从匹配路由提取的路径参数query_paramsRecordstring, string \| string[]是查询串参数request_bodyHttpStreamReader是原始请求体的流式读取器HttpResponse—— handler 返回的结构化缓冲式响应字段类型必填说明status_codeTStatus是HTTP 状态码headersRecordstring, string否响应头bodyTBody否响应体源码中还有一个值得注意的设计细节HttpStreamReader含stream: ReadableStream、readAll()、onMessage()、close()是在 helpers 包内以「结构类型」本地声明的注释明确说明这是为了让 helpers 包不产生对核心 SDK 的运行时依赖见 http/index.ts。3. observabilityOTel 初始化、日志、span 与 payload 工具这是 helpers 包中体量最大的子模块覆盖三大能力OTel 生命周期管理、trace/baggage 上下文操作、worker 资源指标采集。3.1 OTel 生命周期initOtel / flushOtel / shutdownOtelinitOtel(config: OtelConfig): void—— 初始化 OpenTelemetry应在应用启动时调用一次。未显式设置的字段由环境变量补齐见 3.2。flushOtel(): Promisevoid—— 强制刷新所有 OTel provider 但不拆掉它们适用于短生命周期进程退出前想把 pending 的 span/metrics/logs 送出去、之后还要继续使用 OTel 的场景。shutdownOtel(): Promisevoid—— 关闭 OTel关闭前尽力刷出 pending 数据。从 telemetry-system/index.ts 的实现看initOtel内部完成的工作包括禁用开关enabled ?? parseBoolEnv(process.env.OTEL_ENABLED, true)false/0/no/off均视为禁用服务身份serviceName默认取OTEL_SERVICE_NAME否则为iii-nodeserviceInstanceId未设置时自动生成 UUID共享 WebSocket 连接所有信号traces/metrics/logs复用一条连接且连接地址会被自动改写为引擎的/otel专用端点——源码注释解释这是为了避免遥测 socket 被误注册进worker_registry而显示为「幽灵 null-metadata worker」Span 管道先注册BaggageSpanProcessor把 baggage 条目物化为 span 属性再注册BatchSpanProcessor并把scheduledDelayMillis覆盖为 100ms——OpenTelemetry 默认是 5000ms这正是「操作完成后 trace 要好几秒才出现」的主因指标导出默认启用PeriodicExportingMetricReader按 60s 间隔导出fetch 自动插桩默认patchGlobalFetch为每次出站 fetch 创建 CLIENT spanNode.js/Bun/Deno 均可工作。3.2 OtelConfig 全参数与默认值配置项类型默认值说明enabledbooleantrue是否启用 OTel 导出设false或OTEL_ENABLEDfalse/0/no/off可禁用serviceNamestringOTEL_SERVICE_NAME或iii-node上报的服务名[Omitted for brevity]serviceVersionstringSERVICE_VERSION或unknown服务版本serviceNamespacestringSERVICE_NAMESPACE环境变量服务命名空间serviceInstanceIdstringSERVICE_INSTANCE_ID或自动 UUID服务实例 IDengineWsUrlstringIII_URL或ws://localhost:49134III Engine WebSocket 地址metricsEnabledbooleantrue是否启用指标导出支持OTEL_METRICS_ENABLED覆盖metricsExportIntervalMsnumber60000指标导出间隔毫秒spansFlushIntervalMsnumber100span 批处理缓冲延迟环境变量OTEL_SPANS_FLUSH_INTERVAL_MS覆盖logsFlushIntervalMsnumber100日志处理器刷新延迟OTEL_LOGS_FLUSH_INTERVAL_MS覆盖logsBatchSizenumber1每批导出的最大日志记录数fetchInstrumentationEnabledbooleantrue是否自动插桩globalThis.fetchinstrumentationsInstrumentation[]—额外注册的 OTel 插桩如 PrismaInstrumentationreconnectionConfigPartialReconnectionConfig—WebSocket 重连行为配置上表与 types.ts 中的DEFAULT_OTEL_CONFIG常量一一对应。ReconnectionConfigtypes.ts 中定义了默认值字段类型默认值说明initialDelayMsnumber1000起始重连延迟maxDelayMsnumber30000延迟上限backoffMultipliernumber2指数退避倍数jitterFactornumber0.3随机抖动因子0-1maxRetriesnumber-1最大重试次数-1表示无限从 connection.ts 看SharedEngineConnection的实现细节包括WebSocket 握手超时 10 秒与 iii-sdk/Rust SDK 对齐断线后按指数退避 抖动重连重连成功后自动 flush 断线期间积压的消息pending 队列上限 1000 条主动关闭shuttingDown时不会触发重连。三类 OTLP 帧通过消息前缀区分信号类型OTLPtraces、MTRCmetrics、LOGSlogs。3.2.1 Logger自动关联 trace 的结构化日志import { Logger } from iii-dev/helpers/observability const logger new Logger() logger.info(Worker connected) logger.info(Order processed, { orderId: ord_123, amount: 49.99, currency: USD }) logger.warn(Retry attempt, { attempt: 3, maxRetries: 5, endpoint: /api/charge }) logger.error(Payment failed, { orderId: ord_123, gateway: stripe, errorCode: card_declined })Logger把每条日志发射为 OpenTelemetry LogRecord。从 logger.ts 的emit()实现看每条日志会自动附加当前trace_id与span_id属性构造时可显式传入traceId/spanId/serviceName覆盖并自动注入service.nameOTel 未初始化时优雅降级到console.*。第二个参数传入对象而非字符串插值才能在可观测后端里按 key 过滤、聚合和建仪表盘。3.3 span 与上下文操作函数族函数签名说明withSpan(name, { kind?, traceparent? }, fn) PromiseT新建 span 并在其中执行回调支持从 W3C traceparent 恢复父上下文currentTraceId/currentSpanId() string \| undefined从活动 span 上下文提取 trace/span IDcurrentSpanIsRecording() boolean无活动 span 或采样器丢弃时返回falsesetCurrentSpanAttribute(key, value) void给活动 span 设属性未 recording 时 no-opsetCurrentSpanError(message) void设置 span 的 ERROR 状态recordSpanEvent(name, attrs?) void向活动 span 添加事件executeTracedRequest(input, init?: TracedFetchInit) PromiseResponse在 OTel CLIENT span 中执行 fetchW3C 上下文传播extractTraceparent(traceparent) Context解析 W3Ctraceparent头extractBaggage(baggage) Context解析 W3Cbaggage头extractContext(traceparent?, baggage?) Context一次性恢复两种上下文injectTraceparent() string | undefined/injectBaggage() string | undefined把当前上下文序列化回 W3C 头baggage 增删查setBaggageEntry(key, value)、removeBaggageEntry(key)、getBaggageEntry(key)、getAllBaggage()。executeTracedRequest的行为细节来自 http-instrumentation.tsspan 名采用METHOD /path形式自动注入http.request.method、url.full、server.address、url.path/query等语义约定属性出站请求头自动注入 W3Ctraceparent响应状态码 400 或网络异常时把 span 置为 ERROR 并记录error.type。TracedFetchInit在标准RequestInit基础上多一个可选tracer字段。另有patchGlobalFetch(tracer)/unpatchGlobalFetch()用于全局接管/还原globalThis.fetch。payload 安全工具redact(value: unknown) unknown—— 递归脱敏敏感 key返回新值redactAndTruncate(value, maxBytes) { json: string; truncated: boolean }—— 先脱敏再 JSON 序列化maxBytes为null时不限制resolveMaxBytesFromEnv() number | null—— 从环境变量解析 payload 字节上限safeStringify(value) string—— 安全序列化处理循环引用、BigInt 等边界情况失败时返回[unserializable]。3.4 worker 资源指标WorkerMetricsCollector采集 CPU、内存、事件循环延迟。源码注释说明它使用 Node.js 的monitorEventLoopDelayAPI 做高精度测量而非手工setImmediate计时WorkerMetricsCollectorOptions.eventLoopResolutionMs控制直方图分辨率越低越精确但开销越大。registerWorkerGauges(meter, options)/stopWorkerGauges()注册/注销 OTel 可观察 gauge每次采集周期上报 worker 的 CPU、内存、事件循环指标。WorkerGaugesOptions必填workerId稳定标识可选workerName可读名称。3.5 OtelLogEvent引擎回传日志事件结构OtelLogEvent描述引擎侧 OTEL 日志事件的线上格式主要字段attributes结构化属性、body日志正文、observed_timestamp_unix_nano与timestamp_unix_nano纳秒 Unix 时间戳、resource资源属性、service_name、severity_numberOTEL 严重度 1-24TRACE1-4、DEBUG5-8、INFO9-12、WARN13-16、ERROR17-20、FATAL21-24与severity_text、可选的trace_id/span_id关联字段。4. stream触发器配置、变更事件与原子更新操作iii-dev/helpers/stream导出 stream 触发器的完整类型体系可分为四组。4.1 触发器配置与事件StreamTriggerConfigstream触发器过滤哪些条目变更会触发 handler字段类型必填说明stream_namestring是监听的 stream 名只有该 stream 上的变更会触发group_idstring否设置后仅该 group 内的变更触发item_idstring否设置后仅该条目的变更触发condition_function_idstring否条件函数 ID返回false时跳过 handlerStreamJoinLeaveTriggerConfigstream:join/stream:leave触发器配置仅含可选condition_function_id。StreamChangeEventstream触发器 handler 输入在stream::set、stream::update、stream::delete发生时触发。字段type: stream、streamName、groupId、id?、timestampUnix 时间戳、event: StreamChangeEventDetail。其中StreamChangeEventDetail为{ type: create | update | delete; data: any }。StreamJoinLeaveEventjoin/leave 事件负载含stream_name、group_id、subscription_id订阅唯一标识、id?、context?鉴权上下文。鉴权侧StreamAuthInputaddr/headers/path/query_params、StreamAuthResult可选context传递给后续 handler、StreamContext即StreamAuthResult[context]的别名、StreamJoinResultunauthorized: boolean。4.2 IO 输入/结果类型类型字段用途StreamGetInput/StreamSetInput/StreamDeleteInput/StreamUpdateInput/StreamListInput/StreamListGroupsInputstream_namegroup_id条目操作item_idset 另带data、update 另带ops: UpdateOp[]各类 stream IO 的输入StreamSetResult/StreamDeleteResult/StreamUpdateResultnew_value/old_value?/errors?对应操作结果StreamUpdateResult.errors的语义值得强调它由merge与append在校验拒绝时发出路径深度/大小、值深度或路径段/顶层 key 出现__proto__/constructor/prototype以及append的type_mismatch与target_not_object场景成功应用的 op 仍反映在new_value中该字段为空时不出现在 JSON 线上格式里。4.3 原子更新操作 UpdateOptype UpdateOp | UpdateSet | UpdateIncrement | UpdateDecrement | UpdateAppend | UpdateRemove | UpdateMergeUpdateSet{ type: set; path: string; value: any }path为空字符串时指向根值。UpdateIncrement / UpdateDecrement{ type; path: string; by: number }对数值字段增减。UpdateRemove{ type: remove; path: string }。UpdateAppend向数组追加、字符串拼接或在嵌套路径 push 新值。引擎语义文档明确列出嵌套路径上缺失或非对象的中间节点会被自动替换为{}叶子缺失时——嵌套路径恒生成[value]单字符串路径在「字符串拼接层」生成字符串否则生成[value]已存在数组则 push已存在字符串且值为字符串则拼接已存在对象/标量则报append.type_mismatch。路径校验深度 32 段、单段 256 字节、或出现__proto__/constructor/prototype段均被拒绝并以结构化错误写入响应的errors字段该 op 不应用。UpdateMerge把对象浅合并进目标目标为根或数组字面量段指定的嵌套位置。校验上限深度 32 段、段 256 字节、值深度 16、顶层 key 超过 1024 个或出现受保护 key均被拒绝。MergePathstring | string[]。单个字符串是 legacy 一级字段数组是嵌套路径每个元素是一个字面 key点号不解释为分隔符[a.b]指向名为a.b的单个 key而非a → b省略、或[]指向根值。UpdateOpError{ code: string; message: string; op_index: number; doc_url?: string }code是稳定错误码如merge.path.too_deepop_index指向原ops数组中出错的位置。4.4 queueEnqueueResultiii-dev/helpers/queue只导出一个类型——函数以TriggerAction.Enqueue被调用时的返回值export type EnqueueResult { /** 入队消息的唯一回执 ID。 */ messageReceiptId: string }见 queue/index.ts。5. worker-connection-managerRBAC 鉴权与注册回调类型该子模块为 RBAC 代理 worker 提供 WebSocket 升级阶段的鉴权与注册钩子类型。AuthInput—— WebSocket 升级时传给 RBAC 鉴权函数的输入headers升级请求的 HTTP 头、ip_address客户端 IP、query_params升级 URL 查询参数value 为数组以支持重复 key。AuthResult—— 鉴权函数返回值控制该 worker 能调用哪些函数及转发给中间件的上下文字段类型缺省说明allow_function_registrationbooleantrue是否允许注册新函数allow_trigger_type_registrationbooleanfalse是否允许注册新触发器类型allowed_functionsstring[][]在expose_functions配置之外额外允许的函数 IDforbidden_functionsstring[][]即使匹配expose_functions也拒绝的函数 ID优先级高于允许allowed_trigger_typesstring[]省略全部允许可注册触发器的类型 IDcontextRecordstring, unknown{}每次调用时转发给中间件函数的任意上下文function_registration_prefixstring—应用于该 worker 所有注册函数 ID 的前缀三个注册钩子均支持「映射」语义——返回可空结果改写注册字段抛异常则拒绝注册OnFunctionRegistrationInputfunction_id、context必填description?、metadata?→OnFunctionRegistrationResultfunction_id/description/metadata均可映射OnTriggerRegistrationInputtrigger_id、trigger_type、function_id、config、context必填metadata?→OnTriggerRegistrationResult四个字段均可映射OnTriggerTypeRegistrationInputtrigger_type_id、description、context必填→OnTriggerTypeRegistrationResulttrigger_type_id/description可映射。6. 小结iii-dev/helpers用五个聚焦的子路径为 iii 的 Node.js Worker 提供了「协议适配 可观测性 状态原语」三层能力http通过控制帧适配层让你用熟悉的req/res心智写流式接口observability以一条自动重连的/otelWebSocket 同时承载 traces、metrics、logs 三类信号100ms 的默认 flush 间隔解决了 OTel 默认的 5 秒延迟问题Logger与executeTracedRequest让日志和出站 HTTP 自动关联到分布式 tracestream的UpdateOp类型把引擎侧的原子更新语义含路径安全校验上限完整暴露给 TypeScript 用户queue与worker-connection-manager则分别补齐了入队回执与 RBAC 注册治理的类型定义。所有接口的完整源码可进一步对照 sdk/packages/node/helpers/src 阅读API 参考文档由 docs/next/scripts/generate-api-docs.mts 基于源码 doc-comments 生成。【免费下载链接】iiiEffortlessly compose, extend, and observe every service in real-time for the first time ever.项目地址: https://gitcode.com/GitHub_Trending/mo/iii创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考