
workerd Streams 实现指南ReadableStream / WritableStream / TransformStream 的双实现架构与源码剖析【免费下载链接】workerdThe JavaScript / Wasm runtime that powers Cloudflare Workers项目地址: https://gitcode.com/GitHub_Trending/wo/workerd导读本文基于 docs/streams.md 展开系统讲解 Cloudflare Workers 运行时 workerd 中 Streams API 的完整实现——同一套ReadableStream/WritableStream/TransformStreamJavaScript 接口背后同时运行着Internalkj 后端与StandardWHATWG 规范两套实现。读完本文你将掌握两种流的内核差异、队列与背压机制、tee()与pipeTo()的底层分流/选路逻辑以及兼容性标志如何决定new TransformStream()的行为并能在源码中找到每一处关键实现的落点。一、总览一套 API两套实现workerd 的 Streams 子系统拥有两套互相独立的实现但对外暴露的是同一组 JavaScript APIReadableStream、WritableStream、TransformStream。1.1 Internal Streams内部流定位workerd 最早期的实现专为运行时自身需求设计——读取请求体request.body、写入响应体等。本质对 kj 异步 I/O 原语kj::AsyncInputStream、kj::AsyncOutputStream的薄封装。规范贴合度与 WHATWG Streams 规范仅有表面联系属于非标准实现。数据形态仅面向字节只处理TypedArray与ArrayBuffer。从源码看其控制器定义在 internal.hclass ReadableStreamInternalController: public ReadableStreamController, public kj::PtrTarget { public: using Readable IoOwnReadableStreamSource; // 状态机Readable / Closed / Errored 三态 explicit ReadableStreamInternalController(StreamStates::Closed closed) ... explicit ReadableStreamInternalController(StreamStates::Errored errored) ... explicit ReadableStreamInternalController(Readable readable) ...注释明确说明Every stream implementation that originates fromwithinthe Workers runtime will use these凡源自 Workers 运行时内部的流都使用这套控制器且其行为not entirely compliant with the streams specification。1.2 Standard Streams标准流定位符合 WHATWG Streams 规范的独立实现完全由 JavaScript Promise 与用户提供的回调函数驱动。数据形态既可以字节导向也可以值导向可处理任意 JavaScript 值包括undefined、null。核心代码主要在 standard.h 与 standard.c其中ReadableStreamJsController与WritableStreamJsController合计约 5400 行是整个子系统的复杂度中心见 AGENTS.md。1.3 为什么两套实现并存因为删除内部实现会破坏向后兼容。runtime 通过兼容性标志系统决定行为归属request.body等各类 API 始终返回Internal 流带用户回调的new ReadableStream(...)始终创建Standard 流new TransformStream()具体走哪套由兼容性标志transformstream_enable_standard_constructor决定详见下文第四节。两类实现的完整对照引用 README.md 的分类矩阵维度InternalStandard规范贴合度非标准kj 后端WHATWG Streams数据类型仅字节TypedArray/ArrayBuffer字节或任意 JS 值队列模型无队列仅单个 pending read双队列数据队列 pending read 队列异步模型kj::Promise/ kj 事件循环JS Promise / 微任务Isolate 锁数据流在锁外流动数据流在锁内流动背压隐式kj 流控显式highwater mark size 算法Reader 类型Default BYOBDefault BYOB仅字节流Readable 控制器ReadableStreamInternalControllerReadableStreamJsControllerWritable 控制器WritableStreamInternalControllerWritableStreamJsControllerReadable 后端ReadableStreamSource包装kj::AsyncInputStreamJS pull/cancel 算法Writable 后端WritableStreamSink包装kj::AsyncOutputStreamJS write/abort/close 算法创建方式request.body、内部 APInew ReadableStream({...})注意ReadableStream/WritableStream的公开 API不提供任何运行时检查手段来区分一个流是字节导向还是值导向这也直接影响了后续pipeTo桥接层的设计见第五节。二、核心术语2.1 控制器Controllers每个ReadableStream和WritableStream都有一个底层控制器由控制器提供真正的实现。Internal 与 Standard 各有专属控制器类流类型Internal 控制器Standard 控制器定义文件ReadableStreamReadableStreamInternalControllerReadableStreamJsControllerinternal.h / standard.hWritableStreamWritableStreamInternalControllerWritableStreamJsControllerinternal.h / standard.h控制器模式在 AGENTS.md 中被概括为一条清晰的依赖链ReadableStream → ReadableStreamController → 具体实现控制器Internal / Standard2.2 字节导向 vs 值导向字节导向byte-oriented只处理TypedArray和ArrayBuffer数据。值导向value-oriented可处理任意 JavaScript 值包括undefined与null。一个TypedArray或ArrayBuffer可以穿过值导向流但会被当作不透明的 JavaScript 值——流不会将其解释为字节数据。规则很简单Internal 流永远是字节导向Standard 流两种皆可。2.3 BYOB 读取器 vs Default 读取器Reader从ReadableStream中消费数据共有两种BYOB ReaderBring Your Own Buffer仅适用于字节导向流。调用方提供一个目标TypedArray由流来填充数据——数据直接写入调用方的缓冲区减少拷贝。Default Reader字节导向与值导向流都可用。数据以流自身产生的形态交付给调用方。三、ReadableStream 工作原理先看最基础的用法const readableStream getReadableSomehow(); const reader readableStream.getReader(); const chunk await reader.read(); console.log(chunk.value); // 读到的数据 console.log(chunk.done); // 流完成时为 true用户代码反复调用read()直到chunk.done true。幕后发生了什么完全取决于这是 Internal 还是 Standard 流。3.1 Standard ReadableStream双队列 四算法StandardReadableStream维护两个内部队列可用数据队列queue of available data与pending read 队列queue of pending reads。控制器依赖四个回调函数规范中称为 algorithmsstartReadableStream创建时立即调用用于初始化pull请求源提供更多数据cancel流被显式取消时调用size计算每个 chunk 的大小用于背压计算。创建与填充流程流创建后 start 算法立即执行一旦完成流检查highwater mark高水位线——即用 size 算法计算出的、数据队列中应持有的最大数据量。若当前队列大小低于 highwater mark则调用 pull 算法。pull 算法向流内推入数据pull 却实际在 push原文也承认这种反讽若没有 pending read、或队列中已有数据新数据进入队列若存在 pending read则立即用新数据满足该读请求超出该次读取所需的多余部分进入队列。---------------- | pull algorithm | ------------------------------------------ ---------------- | | | v | --------------- | | enqueue(data) | | --------------- | | | v | -------------------- ------------------- | | has a pending read | ---- | has data in queue | | -------------------- yes ------------------- | | | no| | no| | v | v | ---------------------- | ------------------- yes | | fulfill pending read | | | add data to queue | -------- ---------------------- | ------------------- | | | | | | v | | ----------------------------- no | -- | is queue at highwater mark? | ----------- ----------------------------- yes | | v (done)reader.read()的处理队列有数据 → 立即满足该读请求若此举使队列降至 highwater mark 以下则调用 pull 算法队列为空 → 将读请求加入 pending read 队列若队列大小低于 highwater mark则调用 pull 算法。--------------- | reader.read() | --------------- | v ----------------- ------------------ --------------------- | queue has data? | ---- | add pending read | ---- | call pull algorithm | ----------------- no ------------------ --------------------- | yes | v -------------- | fulfill read | -------------- | v (done)重要警告源自源码与文档的反复强调源可以在没有任何活跃 reader、甚至在背压已被触发之后继续向队列推数据——这些数据会无界地堆积在内存中单次 push 的数据量可能超过一次 read 能消费的量多出的部分留在队列里一旦用户代码拿到控制器引用就可以在任何时候独立于 pull 算法向队列 enqueue 数据。标准流背压的核心实现是 queue.h 中的ValueQueue与ByteQueue两类队列queue.h它们分别服务于值导向与字节导向的标准流配合 highwater mark 完成背压计算。3.2 Internal ReadableStream单 pending read无队列Internal 流的运行方式截然不同关键差异有三同时只允许一个 pending read、没有内部数据队列、没有 pull 算法。const readable new ReadableStream(); // Standard 流 const reader readable.getReader(); reader.read(); reader.read(); // 没问题 —— 作为 pending read 排队 const readable request.body; // Internal 流 const reader readable.getReader(); reader.read(); reader.read(); // 报错Internal 流由ReadableStreamSource支撑其内部包装了一个kj::AsyncInputStream。调用read()直接转化为对源上的tryRead()调用返回一个kj::Promise承载数据同一时间只能有一个 read 在飞行中in flight。--------------- | reader.read() | --------------- | v ------------------- ----------------------------------- | has pending read? | --- | tryRead() on ReadableStreamSource | ------------------- no ----------------------------------- yes | v ------- | error | -------这也意味着数据只有在存在活跃 reader 时才会流过 Internal 流且每次 read 返回的数据量绝不超过调用中指定的最大字节数。3.3 TeeReadableStream.tee()tee()将数据流拆分为两个独立的ReadableStream实例分支。WHATWG 规范中的行为tee()创建两个分支共享原流trunk上的一个 reader。当一个分支 pull 时数据满足该分支的读取并复制一份推入另一个分支的队列。这意味着一个分支读得比另一个快会导致慢分支出现无界的内存增长。workerd 的修改不再复制数据分支之间持有数据的refcounted 引用kj::RcEntry。向 trunk 的背压信号基于未消费数据最多的分支来计算。---------------- | pull algorithm | ---------------- | v .......................................................... --------------- . --------------------- ------------------- | enqueue(data) | --- | push data to branch | --- | has pending read? | --------------- . --------------------- ------------------- | . no | yes | | . ------------------- | -------------- | . | add data to queue | ----- | fulfill read | | . ------------------- -------------- | ............................................................ | . --------------------- ------------------- -------- | push data to branch | --- | has pending read? | . --------------------- ------------------- . no | yes | . ------------------- | -------------- . | add data to queue | ----- | fulfill read | . ------------------- -------------- ............................................................这项优化并不能彻底阻止某个分支读得远慢于另一个分支但只要底层源尊重背压信号就能避免原本会发生的内存堆积。对应安全模式见 README.md 中的RcEntry模式class Entry: public kj::Refcounted { kj::RcEntry clone(jsg::Lock js); };四、WritableStream 工作原理4.1 Standard WritableStream五算法 双背压通道StandardWritableStream使用五个算法start准备接收数据write每个写入的 chunk 都会调用abort异常终止时调用close最后一个 chunk 之后调用size计算 chunk 大小以服务于 highwater mark。const writable getWritableSomehow(); const writer writable.getWriter(); await writer.write(chunk of data); await writer.write(another chunk of data);write 算法是异步的每次write(data)都会调用它返回的 promise 在写入完成时 resolve。这个 promise 是主要的背压机制。背压通过两种方式呈现writer.desiredSize——在队列填满之前还可写入的数据量writer.ready——背压解除时 resolve 的 promise每次触发背压时都会替换成一个新的 promise。重要highwater mark 只是建议性的advisory。write 算法应当始终做好被调用的准备无论当前 highwater mark 取值如何——即不能因为背压就假定不会被调用。4.2 Internal WritableStream直通 SinkInternalWritableStream由WritableStreamSink支撑它是kj::AsyncOutputStream的薄封装。每次write()直接透传给 sink 的write()。此后的行为取决于具体 sink 的实现例如 HTTP 响应体、TLS 连接等不同 sink 各有其语义。相关代码见 internal.h 中的WritableStreamInternalController。五、TransformStreamTransformStream连接一个ReadableStream与一个WritableStream——写入 writable 侧的数据可以被读取于 readable 侧且可能经过变换。5.1 IdentityTransformStreamInternal旧实现最初的TransformStream是一个恒等变换数据从 writable 侧到 readable 侧原样通过。它用一个类同时实现ReadableStreamSource与WritableStreamSink两个接口。该实现不满足规范。transform.h 的注释记录了这段历史Workers 中最初的 TransformStream 不过是仅处理字节数据的恒等直通identity passthrough不执行任何真正的变换也不符合流规范这个旧版本被迁移到了IdentityTransformStream类中。兼容性标志transformstream_enable_standard_constructor未启用时TransformStream就是IdentityTransformStream的别名启用后TransformStream才实现标准化行为。旧行为仍可通过new IdentityTransformStream()获得。const { readable, writable } new IdentityTransformStream(); const enc new TextEncoder(); const writer writable.getWriter(); const reader readable.getReader(); // 注意write 的 promise 直到 read() 被调用才会 resolve await Promise.all([writer.write(enc.encode(hello)), reader.read()]);IdentityTransformStream是一个简单的状态机同一时间只允许一个 in-flight 的 read 或 write。write()的 promise 不会在对应的read()发生前 resolve反之亦然。它只支持字节数据ArrayBuffer/TypedArray。其状态机来自 README.md一目了然当前状态事件下一状态动作Idlewrite()Write Pending持有数据返回 pending promiseIdleread()Read Pending返回 pending promiseRead Pendingwrite()Idle用数据满足 readresolve writeWrite Pendingread()Idle用持有的数据满足 read源码层面该类定义于 identity-transform-stream.h并通过newIdentityPipe()identity-transform-stream.h构造底层管道。值得一提的还有同文件中的FixedLengthStream——它与恒等流几乎相同但在 readable 侧带有已知字节长度目前并不强制校验该上限其作用是说服 kj-http 层输出Content-Length头只要数据未被 gzip 等处理。5.2 Standard TransformStream新实现StandardTransformStream使用三个算法start初始化变换transform接收一个 chunk修改它并将结果 enqueueflush完成变换。const { writable, readable } new TransformStream({ transform(chunk, controller) { controller.enqueue(${chunk}!.toUpperCase()); }, }); const writer writable.getWriter(); const reader readable.getReader(); // write 的 promise 不等待 read。 await writer.write(hello); await reader.read(); // { value: HELLO!, done: false }两侧都是完整的 Standard 流各自拥有独立的队列与背压。与IdentityTransformStream不同write 不会被 read 阻塞除非背压信号表明队列已满。任意 JavaScript 值都可以流过。启用标准构造器的测试配置示例如 htmlrewriter-test.wd-testcompatibilityFlags [nodejs_compat, experimental, streams_enable_constructors, transformstream_enable_standard_constructor, ...]六、PipingpipeTo()的四种管道回路Piping 建立从ReadableStream到WritableStream的数据流动。由目标 writable 决定管道如何实现——经由其控制器的tryPipeFrom()方法完成。该入口定义在 common.h// The tryPipeFrom attempts to establish a data pipe where sources data // is delivered to this WritableStreamController as efficiently as possible. virtual kj::Maybejsg::Promisevoid tryPipeFrom( jsg::Lock js, jsg::RefReadableStream source, PipeToOptions options) 0;根据源与目标的流类型组合共有四种 pipe loop 变体--------------------------- | readable.pipeTo(writable) | --------------------------- | v ------------------------------------------ | writableController.tryPipeFrom(readable) | ------------------------------------------ | v ----------------------- ----------------------- | is internal writable? | ---- | is internal readable? | ----------------------- yes ----------------------- | no | yes | no | | | v ---- --------------- ----------------------- | | kj-to-kj pipe | | is internal readable? | | | loop | ----------------------- | --------------- no | yes | | | | | v | | --------------- | v | JS-to-JS pipe | | --------------- | loop | | | JS-to-kj pipe | --------------- | | loop | v --------------- --------------- | kj-to-JS pipe | | loop | ---------------四种组合的对照表来自 README.mdReadableWritableLoop 类型Isolate 锁数据限制InternalInternalkj-to-kj不持有仅字节InternalStandardkj-to-JS持有仅字节StandardInternalJS-to-kj持有仅字节StandardStandardJS-to-JS持有任意值kj-to-kj最优化路径。kj 直接在AsyncInputStream与AsyncOutputStream之间搬运数据完全在 JavaScript isolate 锁之外进行整个过程不运行任何 JavaScript 代码数据也从不进入 JS 堆。kj-to-JS / JS-to-kj在 kj 与 JavaScript Promise 之间搭桥。整个数据流必须在 isolate 锁内进行。数据被限制为字节——因为 Standard 流可能是值导向的而 API 没有提供运行时检查手段所以桥接层只能强制按字节处理以保证安全。JS-to-JS纯 JavaScript promise 链式调用概念上等价于async function pipe(reader, writer) { for await (const chunk of reader) { await writer.write(chunk); } }两套控制器的tryPipeFrom分别实现在 standard.c 与 internal.c注释均明确指出源 can be either a JavaScript-backed ReadableStream or ReadableStreamSource-backed可以是 JS 后端或 ReadableStreamSource 后端。internal-test.c中的测试 internal-test.c 还验证了pipe 期间对同一 sink 再次发起tryPipeFrom会被拒绝currently being piped to。6.1new Response(standardReadable)标准流到内部 API 的桥接将 StandardReadableStream传给new Response()这类 API 是件复杂的事因为这些 API 是为 Internal 流构建的内部使用 kj 异步 I/O。为打通二者StandardReadableStream可以通过ReadableStreamSourceAPI被消费与 Internal 流所用的是同一套 API。当适配器上的pumpTo()被调用时获取 isolate 锁运行读 JS 流 → 写 kj 输出的 promise 循环直到数据耗尽或发生错误才结束。pumpTo的定义见 common.h其中特别说明ReadableStreamSource版本的pumpTo()没有amount参数因为 Streams 规范只定义了泵出全部数据而 common.c 中对应实现特意注明不使用KJ_CO_MAGIC BEGIN_DEFERRED_PROXYING因为底层内存可能与 V8 堆绑定如ArrayBuffer、Blob 数据。七、复杂度预算为什么这块代码很难改workerd 的 streams 实现需要同时平衡多个来源的复杂度同一规范的两套实现Internal 与 Standard两种数据导向字节与值两块内存堆JavaScript 堆与 kj 堆两种异步模型JavaScript Promise 与 kj Promise通过特性标志保持严格向后兼容。最关键的分界线是 isolate 锁Internal 流的数据流动发生在 isolate 锁之外由 kj 事件循环驱动Standard 流的数据流动发生在 isolate 锁之内由 JavaScript Promise 驱动当两个世界交互时跨类型 pipe、把标准流传给内部 API桥接代码必须小心翼翼地同时管理两种异步模型。7.1 跨请求模型workerd 中多个请求通过绿色线程green threads共享一个 isolate。当某个请求因 I/O 让出yield时另一个请求可能开始执行。SetPromiseCrossContextResolveCallback机制负责把 promise 反应reactions延迟调度到正确的请求上下文。streams 代码通过以下手段与之协作详见 README.mdioContext.addFunctor()——把 continuation 绑定到正确的IoContextIoOwn——确保对象只在正确的上下文中被访问Promise 上下文标记——所有 promise 都标记其来源IoContext若标记与当前上下文不匹配反应会被延迟执行。这也解释了为什么ReadableStreamInternalController中的Readable类型被定义为IoOwnReadableStreamSourceinternal.h。八、安全模式目录给维护者的纪律清单README.md 归纳了一组 When/Why/How 安全模式是修改这段代码时的强制纪律AGENTS.md 则将其压缩为 7 条不可违背的不变量INVARIANTS。这里择要说明永远使用deferControllerStateChange()当调用可能在 read/write 操作中途触发 JS 回调的代码时。JS 回调可能在操作进行中触发 close/error状态必须等到操作完成后再变更。实现原理是计数器beginOperation()递增计数器endOperation()在计数器归零时应用待定状态转换。controller.state.beginOperation(); // 递增计数器 auto result readCallback(); // 可能触发 JS 调用 close() controller.state.endOperation(); // 计数器为 0 时应用待定状态迭代消费者时永远使用snapshot()若循环体可能触发 JS迭代过程中消费者可能被增删导致迭代器失效。先拷贝再遍历auto consumers ready.consumers.snapshot(); for (auto consumer: consumers) { consumer-push(js, entry-clone(js)); // 可能触发 JS 修改 consumers }StateListener 回调之后永远不要再访问this回调可能经由owner.doClose()销毁this。lambda continuation 中永远重新检查锁状态绑定在 promise continuation 上的 lambda 可能在锁释放后才执行捕获与执行之间引用的对象可能已被销毁。规则是绝不捕获可能悬垂的裸引用改用addRef()或重新获取。用户可能持有超过底层对象生命周期的句柄时使用 WeakRef如ByobRequest使用前检查存活impl.controller-runIfAlive( [](ReadableByteStreamController controller) { controller.maybeByobRequest kj::none; });先转状态、后 resolve promisecontinuation 必须看到一致的状态。void doClose(jsg::Lock js) { state.transitionToStreamStates::Closed(); // 状态此刻变更 maybeResolvePromise(js, locked.getClosedFulfiller()); // 调度微任务 }异步写引用 JS 堆数据时用V8Ref持有缓冲kj 异步写进行期间 GC 可能回收缓冲。struct Write { jsg::V8Refv8::ArrayBuffer ownBytes; // 阻止 GC kj::ArrayPtrkj::byte bytes; // 指向 ownBytes 的裸指针 };此外还有 Pipe-Lock 生命周期管理当ReadableStream被pipeTo/pipeThrough时源的锁状态机进入PipeLocked目标端的 pipe 机制持有指向该状态的kj::PtrPipeController只有目标端的 pipe 机制才能释放锁且必须在丢弃所有kj::Ptr之后通过ReadableStreamController::releasePipeLock()完成——Pipe::releaseSource()与WritableLockImpl::PipeLocked::releaseSource()封装了这一顺序。源侧的状态机依赖 state-machine.hStateMachine、PendingStates、ActiveState。九、兼容性标志速查标志作用streams_enable_constructors启用标准 Streams 构造器new ReadableStream(...)等是各测试中常见的组合标志之一见 fs-readstream-test.wd-testtransformstream_enable_standard_constructor启用后new TransformStream()创建 Standard TransformStream未启用时等价于new IdentityTransformStream()见 transform.horiginal-transform-stream-backpressure出现在 htmlrewriter-test.wd-test 等测试中用于控制旧版 TransformStream 的背压行为十、进一步阅读src/workerd/api/streams/README.md——速查参考分类矩阵、状态机、管道回路选择、安全模式目录src/workerd/api/streams/AGENTS.md——文件地图、架构摘要、编码不变量与代码评审规则关键实现文件internal.h / internal.c、standard.h / standard.c、queue.h、identity-transform-stream.h、transform.h、common.h、readable.h、writable.h测试参考internal-test.c、standard-test.c、streams-test.js、src/tests/streams 目录下的 340 个测试文件规范原文WHATWG Streams specstreams.spec.whatwg.org【免费下载链接】workerdThe JavaScript / Wasm runtime that powers Cloudflare Workers项目地址: https://gitcode.com/GitHub_Trending/wo/workerd创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考