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-manager):HTTP 流式响应封装的用法与底层控制帧机制、OTel 初始化配置全参数及 WebSocket 上报链路、Stream 触发器配置与原子更新操作(UpdateOp)语义,以及 Worker 连接 RBAC 鉴权类型。读完后你可以直接在 Worker 代码中使用这些 helpers 编写 HTTP 接口、接入分布式追踪与结构化日志,并理解每个参数在源码中的真实行为。
1. 安装与包结构
npm install @iii-dev/helpers@iii-dev/helpers是 iii SDK 体系中跨语言共享的辅助原语包,当前仓库中版本为0.23.0-rc.9(见 package.json)。它自身只依赖 OpenTelemetry SDK 系列包和ws(WebSocket 客户端),通过 5 个子路径对外暴露 API:
| 子路径 | 导入方式 | 职责 |
|---|---|---|
http | import { http } from '@iii-dev/helpers/http' | HTTP 风格 handler 包装器及请求/响应/鉴权类型 |
observability | import { initOtel, Logger } from '@iii-dev/helpers/observability' | 结构化日志、OTel 初始化、span/baggage 工具、worker 资源指标 |
queue | import { EnqueueResult } from '@iii-dev/helpers/queue' | 队列Enqueue触发动作的结果类型 |
stream | import { StreamTriggerConfig } from '@iii-dev/helpers/stream' | Stream 触发器配置、变更事件、IO 输入与更新操作类型 |
worker-connection-manager | import { AuthInput } from '@iii-dev/helpers/worker-connection-manager' | RBAC 鉴权与注册回调类型 |
从 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 格式,可直接传给registerFunction:
http(callback: (req: HttpStreamingRequest, res: HttpStreamingResponse) => Promise<void | HttpResponse<number, string | Buffer | Record<string, unknown>>>) => (req: HttpInternalRequest) => Promise<void | HttpResponse<number, string | Buffer | Record<string, 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.stream:Node.jsWritableStream,用于流式写响应体;res.close():关闭响应通道,结束请求。
2.2 源码级原理:控制帧协议
从 http/index.ts 的实现看,http()的核心是一个「协议适配层」。运行时(SDK 核心)实际传入的是HttpInternalRequest——它在流式请求字段之外多带了一个response: HttpStreamWriter(含sendMessage、stream、close)。helper 把这个内部 writer 拆包装成用户友好的HttpStreamingResponse:
status()会调用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 调用型函数的配置:
| 名称 | 类型 | 必填 | 说明 |
|---|---|---|---|
url | string | 是 | 要调用的 URL |
method | HttpMethod | 否 | HTTP 方法,默认POST |
timeout_ms | number | 否 | 超时时间(毫秒) |
headers | Record<string, string> | 否 | 随请求发送的自定义头 |
auth | HttpAuthConfig | 否 | 鉴权配置 |
HttpMethod—— 注意它与引擎核心builtin_triggers中的 HTTP 方法枚举不同:后者还覆盖HEAD/OPTIONS,此处仅 5 种:
type HttpMethod = "GET" | "POST" | "PUT" | "PATCH" | "DELETE"HttpRequest—— 函数 handler 接收到的缓冲式入站请求:
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
body | TBody | 是 | 已解析的请求体 |
headers | Record<string, string \| string[]> | 是 | 请求头 |
method | string | 是 | HTTP 方法(如GET、POST) |
path_params | Record<string, string> | 是 | 从匹配路由提取的路径参数 |
query_params | Record<string, string \| string[]> | 是 | 查询串参数 |
request_body | HttpStreamReader | 是 | 原始请求体的流式读取器 |
HttpResponse—— handler 返回的结构化缓冲式响应:
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
status_code | TStatus | 是 | HTTP 状态码 |
headers | Record<string, string> | 否 | 响应头 |
body | TBody | 否 | 响应体 |
源码中还有一个值得注意的设计细节:HttpStreamReader(含stream: ReadableStream、readAll()、onMessage()、close())是在 helpers 包内以「结构类型」本地声明的,注释明确说明这是为了让 helpers 包不产生对核心 SDK 的运行时依赖(见 http/index.ts)。
3. observability:OTel 初始化、日志、span 与 payload 工具
这是 helpers 包中体量最大的子模块,覆盖三大能力:OTel 生命周期管理、trace/baggage 上下文操作、worker 资源指标采集。
3.1 OTel 生命周期:initOtel / flushOtel / shutdownOtel
initOtel(config: OtelConfig): void—— 初始化 OpenTelemetry,应在应用启动时调用一次。未显式设置的字段由环境变量补齐(见 3.2)。flushOtel(): Promise<void>—— 强制刷新所有 OTel provider 但不拆掉它们,适用于短生命周期进程退出前想把 pending 的 span/metrics/logs 送出去、之后还要继续使用 OTel 的场景。shutdownOtel(): Promise<void>—— 关闭 OTel,关闭前尽力刷出 pending 数据。
从 telemetry-system/index.ts 的实现看,initOtel内部完成的工作包括:
- 禁用开关:
enabled ?? parseBoolEnv(process.env.OTEL_ENABLED, true),false/0/no/off均视为禁用; - 服务身份:
serviceName默认取OTEL_SERVICE_NAME否则为iii-node;serviceInstanceId未设置时自动生成 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 span,Node.js/Bun/Deno 均可工作。
3.2 OtelConfig 全参数与默认值
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
enabled | boolean | true | 是否启用 OTel 导出;设false或OTEL_ENABLED=false/0/no/off可禁用 |
serviceName | string | OTEL_SERVICE_NAME或iii-node | 上报的服务名 |
[Omitted for brevity]serviceVersion | string | SERVICE_VERSION或unknown | 服务版本 |
serviceNamespace | string | SERVICE_NAMESPACE环境变量 | 服务命名空间 |
serviceInstanceId | string | SERVICE_INSTANCE_ID或自动 UUID | 服务实例 ID |
engineWsUrl | string | III_URL或ws://localhost:49134 | III Engine WebSocket 地址 |
metricsEnabled | boolean | true | 是否启用指标导出,支持OTEL_METRICS_ENABLED覆盖 |
metricsExportIntervalMs | number | 60000 | 指标导出间隔(毫秒) |
spansFlushIntervalMs | number | 100 | span 批处理缓冲延迟;环境变量OTEL_SPANS_FLUSH_INTERVAL_MS覆盖 |
logsFlushIntervalMs | number | 100 | 日志处理器刷新延迟;OTEL_LOGS_FLUSH_INTERVAL_MS覆盖 |
logsBatchSize | number | 1 | 每批导出的最大日志记录数 |
fetchInstrumentationEnabled | boolean | true | 是否自动插桩globalThis.fetch |
instrumentations | Instrumentation[] | — | 额外注册的 OTel 插桩(如 PrismaInstrumentation) |
reconnectionConfig | Partial<ReconnectionConfig> | — | WebSocket 重连行为配置 |
上表与 types.ts 中的DEFAULT_OTEL_CONFIG常量一一对应。
ReconnectionConfig(types.ts 中定义了默认值):
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
initialDelayMs | number | 1000 | 起始重连延迟 |
maxDelayMs | number | 30000 | 延迟上限 |
backoffMultiplier | number | 2 | 指数退避倍数 |
jitterFactor | number | 0.3 | 随机抖动因子(0-1) |
maxRetries | number | -1 | 最大重试次数,-1表示无限 |
从 connection.ts 看,SharedEngineConnection的实现细节包括:WebSocket 握手超时 10 秒(与 iii-sdk/Rust SDK 对齐);断线后按指数退避 + 抖动重连;重连成功后自动 flush 断线期间积压的消息(pending 队列上限 1000 条);主动关闭(shuttingDown)时不会触发重连。三类 OTLP 帧通过消息前缀区分信号类型:OTLP(traces)、MTRC(metrics)、LOGS(logs)。
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.name;OTel 未初始化时优雅降级到console.*。第二个参数传入对象而非字符串插值,才能在可观测后端里按 key 过滤、聚合和建仪表盘。
3.3 span 与上下文操作函数族
| 函数 | 签名 | 说明 |
|---|---|---|
withSpan | (name, { kind?, traceparent? }, fn) => Promise<T> | 新建 span 并在其中执行回调,支持从 W3C traceparent 恢复父上下文 |
currentTraceId/currentSpanId | () => string \| undefined | 从活动 span 上下文提取 trace/span ID |
currentSpanIsRecording | () => boolean | 无活动 span 或采样器丢弃时返回false |
setCurrentSpanAttribute | (key, value) => void | 给活动 span 设属性(未 recording 时 no-op) |
setCurrentSpanError | (message) => void | 设置 span 的 ERROR 状态 |
recordSpanEvent | (name, attrs?) => void | 向活动 span 添加事件 |
executeTracedRequest | (input, init?: TracedFetchInit) => Promise<Response> | 在 OTel CLIENT span 中执行 fetch |
W3C 上下文传播:
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.ts:span 名采用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_number(OTEL 严重度 1-24:TRACE=1-4、DEBUG=5-8、INFO=9-12、WARN=13-16、ERROR=17-20、FATAL=21-24)与severity_text、可选的trace_id/span_id关联字段。
4. stream:触发器配置、变更事件与原子更新操作
@iii-dev/helpers/stream导出 stream 触发器的完整类型体系,可分为四组。
4.1 触发器配置与事件
- StreamTriggerConfig(
stream触发器,过滤哪些条目变更会触发 handler):
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
stream_name | string | 是 | 监听的 stream 名,只有该 stream 上的变更会触发 |
group_id | string | 否 | 设置后仅该 group 内的变更触发 |
item_id | string | 否 | 设置后仅该条目的变更触发 |
condition_function_id | string | 否 | 条件函数 ID,返回false时跳过 handler |
- StreamJoinLeaveTriggerConfig:
stream:join/stream:leave触发器配置,仅含可选condition_function_id。 - StreamChangeEvent:
stream触发器 handler 输入,在stream::set、stream::update、stream::delete发生时触发。字段:type: "stream"、streamName、groupId、id?、timestamp(Unix 时间戳)、event: StreamChangeEventDetail。其中StreamChangeEventDetail为{ type: "create" | "update" | "delete"; data: any }。 - StreamJoinLeaveEvent:join/leave 事件负载,含
stream_name、group_id、subscription_id(订阅唯一标识)、id?、context?(鉴权上下文)。 - 鉴权侧:
StreamAuthInput(addr/headers/path/query_params)、StreamAuthResult(可选context传递给后续 handler)、StreamContext(即StreamAuthResult["context"]的别名)、StreamJoinResult(unauthorized: boolean)。
4.2 IO 输入/结果类型
| 类型 | 字段 | 用途 |
|---|---|---|
StreamGetInput/StreamSetInput/StreamDeleteInput/StreamUpdateInput/StreamListInput/StreamListGroupsInput | stream_name+group_id+(条目操作)item_id,set 另带data、update 另带ops: UpdateOp[] | 各类 stream IO 的输入 |
StreamSetResult/StreamDeleteResult/StreamUpdateResult | new_value/old_value?/errors? | 对应操作结果 |
StreamUpdateResult.errors的语义值得强调:它由merge与append在校验拒绝时发出(路径深度/大小、值深度,或路径段/顶层 key 出现__proto__/constructor/prototype),以及append的type_mismatch与target_not_object场景;成功应用的 op 仍反映在new_value中;该字段为空时不出现在 JSON 线上格式里。
4.3 原子更新操作 UpdateOp
type UpdateOp = | UpdateSet | UpdateIncrement | UpdateDecrement | UpdateAppend | UpdateRemove | UpdateMerge- UpdateSet:
{ 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,均被拒绝。
- MergePath:
string | 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_deep"),op_index指向原ops数组中出错的位置。
4.4 queue:EnqueueResult
@iii-dev/helpers/queue只导出一个类型——函数以TriggerAction.Enqueue被调用时的返回值:
export type EnqueueResult = { /** 入队消息的唯一回执 ID。 */ messageReceiptId: string }见 queue/index.ts。
5. worker-connection-manager:RBAC 鉴权与注册回调类型
该子模块为 RBAC 代理 worker 提供 WebSocket 升级阶段的鉴权与注册钩子类型。
AuthInput—— WebSocket 升级时传给 RBAC 鉴权函数的输入:headers(升级请求的 HTTP 头)、ip_address(客户端 IP)、query_params(升级 URL 查询参数,value 为数组以支持重复 key)。
AuthResult—— 鉴权函数返回值,控制该 worker 能调用哪些函数及转发给中间件的上下文:
| 字段 | 类型 | 缺省 | 说明 |
|---|---|---|---|
allow_function_registration | boolean | true | 是否允许注册新函数 |
allow_trigger_type_registration | boolean | false | 是否允许注册新触发器类型 |
allowed_functions | string[] | [] | 在expose_functions配置之外额外允许的函数 ID |
forbidden_functions | string[] | [] | 即使匹配expose_functions也拒绝的函数 ID,优先级高于允许 |
allowed_trigger_types | string[] | 省略=全部允许 | 可注册触发器的类型 ID |
context | Record<string, unknown> | {} | 每次调用时转发给中间件函数的任意上下文 |
function_registration_prefix | string | — | 应用于该 worker 所有注册函数 ID 的前缀 |
三个注册钩子,均支持「映射」语义——返回可空结果改写注册字段,抛异常则拒绝注册:
OnFunctionRegistrationInput(function_id、context必填,description?、metadata?)→OnFunctionRegistrationResult(function_id/description/metadata均可映射);OnTriggerRegistrationInput(trigger_id、trigger_type、function_id、config、context必填,metadata?)→OnTriggerRegistrationResult(四个字段均可映射);OnTriggerTypeRegistrationInput(trigger_type_id、description、context必填)→OnTriggerTypeRegistrationResult(trigger_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 自动关联到分布式 trace;stream的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),仅供参考