File
Blob: src/node/internal/streams_transform.js
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | // |
| 5 | // Copyright Joyent, Inc. and other Node contributors. |
| 6 | // |
| 7 | // Permission is hereby granted, free of charge, to any person obtaining a |
| 8 | // copy of this software and associated documentation files (the |
| 9 | // "Software"), to deal in the Software without restriction, including |
| 10 | // without limitation the rights to use, copy, modify, merge, publish, |
| 11 | // distribute, sublicense, and/or sell copies of the Software, and to permit |
| 12 | // persons to whom the Software is furnished to do so, subject to the |
| 13 | // following conditions: |
| 14 | // |
| 15 | // The above copyright notice and this permission notice shall be included |
| 16 | // in all copies or substantial portions of the Software. |
| 17 | // |
| 18 | // THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS |
| 19 | // OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF |
| 20 | // MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN |
| 21 | // NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, |
| 22 | // DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR |
| 23 | // OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE |
| 24 | // USE OR OTHER DEALINGS IN THE SOFTWARE. |
| 25 | |
| 26 | // a transform stream is a readable/writable stream where you do |
| 27 | // something with the data. Sometimes it's called a "filter", |
| 28 | // but that's not a great name for it, since that implies a thing where |
| 29 | // some bits pass through, and others are simply ignored. (That would |
| 30 | // be a valid example of a transform, of course.) |
| 31 | // |
| 32 | // While the output is causally related to the input, it's not a |
| 33 | // necessarily symmetric or synchronous transformation. For example, |
| 34 | // a zlib stream might take multiple plain-text writes(), and then |
| 35 | // emit a single compressed chunk some time in the future. |
| 36 | // |
| 37 | // Here's how this works: |
| 38 | // |
| 39 | // The Transform stream has all the aspects of the readable and writable |
| 40 | // stream classes. When you write(chunk), that calls _write(chunk,cb) |
| 41 | // internally, and returns false if there's a lot of pending writes |
| 42 | // buffered up. When you call read(), that calls _read(n) until |
| 43 | // there's enough pending readable data buffered up. |
| 44 | // |
| 45 | // In a transform stream, the written data is placed in a buffer. When |
| 46 | // _read(n) is called, it transforms the queued up data, calling the |
| 47 | // buffered _write cb's as it consumes chunks. If consuming a single |
| 48 | // written chunk would result in multiple output chunks, then the first |
| 49 | // outputted bit calls the readcb, and subsequent chunks just go into |
| 50 | // the read buffer, and will cause it to emit 'readable' if necessary. |
| 51 | // |
| 52 | // This way, back-pressure is actually determined by the reading side, |
| 53 | // since _read has to be called to start processing a new chunk. However, |
| 54 | // a pathological inflate type of transform can cause excessive buffering |
| 55 | // here. For example, imagine a stream where every byte of input is |
| 56 | // interpreted as an integer from 0-255, and then results in that many |
| 57 | // bytes of output. Writing the 4 bytes {ff,ff,ff,ff} would result in |
| 58 | // 1kb of data being output. In this case, you could write a very small |
| 59 | // amount of input, and end up with a very large amount of output. In |
| 60 | // such a pathological inflating mechanism, there'd be no way to tell |
| 61 | // the system to stop doing the transform. A single 4MB write could |
| 62 | // cause the system to run out of memory. |
| 63 | // |
| 64 | // However, even in such a pathological case, only a single written chunk |
| 65 | // would be consumed, and then the rest would wait (un-transformed) until |
| 66 | // the results of the previous transformed chunk were consumed. |
| 67 | |
| 68 | 'use strict'; |
| 69 | |
| 70 | import { ERR_METHOD_NOT_IMPLEMENTED } from 'node-internal:internal_errors'; |
| 71 | import { nextTick } from 'node-internal:internal_process'; |
| 72 | |
| 73 | import { Duplex } from 'node-internal:streams_duplex'; |
| 74 | |
| 75 | import { getHighWaterMark } from 'node-internal:streams_state'; |
| 76 | |
| 77 | const streamsNodejsV24Compat = |
| 78 | Cloudflare.compatibilityFlags.enable_streams_nodejs_v24_compat; |
| 79 | |
| 80 | Object.setPrototypeOf(Transform.prototype, Duplex.prototype); |
| 81 | Object.setPrototypeOf(Transform, Duplex); |
| 82 | |
| 83 | const kCallback = Symbol('kCallback'); |
| 84 | |
| 85 | export function Transform(options) { |
| 86 | if (!(this instanceof Transform)) return new Transform(options); |
| 87 | |
| 88 | // TODO (ronag): This should preferably always be |
| 89 | // applied but would be semver-major. Or even better; |
| 90 | // make Transform a Readable with the Writable interface. |
| 91 | const readableHighWaterMark = options |
| 92 | ? getHighWaterMark(this, options, 'readableHighWaterMark', true) |
| 93 | : null; |
| 94 | if (readableHighWaterMark === 0) { |
| 95 | // A Duplex will buffer both on the writable and readable side while |
| 96 | // a Transform just wants to buffer hwm number of elements. To avoid |
| 97 | // buffering twice we disable buffering on the writable side. |
| 98 | options = { |
| 99 | ...options, |
| 100 | highWaterMark: null, |
| 101 | readableHighWaterMark, |
| 102 | writableHighWaterMark: options.writableHighWaterMark || 0, |
| 103 | }; |
| 104 | } |
| 105 | Duplex.call(this, options); |
| 106 | |
| 107 | // We have implemented the _read method, and done the other things |
| 108 | // that Readable wants before the first _read call, so unset the |
| 109 | // sync guard flag. |
| 110 | this._readableState.sync = false; |
| 111 | this[kCallback] = null; |
| 112 | if (options) { |
| 113 | if (typeof options.transform === 'function') |
| 114 | this._transform = options.transform; |
| 115 | if (typeof options.flush === 'function') this._flush = options.flush; |
| 116 | } |
| 117 | |
| 118 | // When the writable side finishes, then flush out anything remaining. |
| 119 | // Backwards compat. Some Transform streams incorrectly implement _final |
| 120 | // instead of or in addition to _flush. By using 'prefinish' instead of |
| 121 | // implementing _final we continue supporting this unfortunate use case. |
| 122 | this.on('prefinish', prefinish); |
| 123 | } |
| 124 | |
| 125 | function final(cb) { |
| 126 | if (typeof this._flush === 'function' && !this.destroyed) { |
| 127 | this._flush((er, data) => { |
| 128 | if (er) { |
| 129 | if (cb) { |
| 130 | cb(er); |
| 131 | } else { |
| 132 | this.destroy(er); |
| 133 | } |
| 134 | return; |
| 135 | } |
| 136 | if (data != null) { |
| 137 | this.push(data); |
| 138 | } |
| 139 | this.push(null); |
| 140 | if (cb) { |
| 141 | cb(); |
| 142 | } |
| 143 | }); |
| 144 | } else { |
| 145 | this.push(null); |
| 146 | if (cb) { |
| 147 | cb(); |
| 148 | } |
| 149 | } |
| 150 | } |
| 151 | |
| 152 | function prefinish() { |
| 153 | if (this._final !== final) { |
| 154 | final.call(this); |
| 155 | } |
| 156 | } |
| 157 | Transform.prototype._final = final; |
| 158 | |
| 159 | Transform.prototype._transform = function () { |
| 160 | throw new ERR_METHOD_NOT_IMPLEMENTED('_transform()'); |
| 161 | }; |
| 162 | |
| 163 | Transform.prototype._write = function (chunk, encoding, callback) { |
| 164 | const rState = this._readableState; |
| 165 | const wState = this._writableState; |
| 166 | const length = rState.length; |
| 167 | this._transform(chunk, encoding, (err, val) => { |
| 168 | if (err) { |
| 169 | callback(err); |
| 170 | return; |
| 171 | } |
| 172 | if (val != null) { |
| 173 | this.push(val); |
| 174 | } |
| 175 | // This is a semver-major change. Ref: https://github.com/nodejs/node/commit/557044af407376aff28a0a0800f3053bb58e9239 |
| 176 | if (streamsNodejsV24Compat && rState.ended) { |
| 177 | // If user has called this.push(null) we have to delay the callback to properly propagate the new |
| 178 | // state. |
| 179 | nextTick(callback); |
| 180 | return; |
| 181 | } else if ( |
| 182 | wState.ended || |
| 183 | // Backwards compat. |
| 184 | length === rState.length || |
| 185 | // Backwards compat. |
| 186 | rState.length < rState.highWaterMark |
| 187 | ) { |
| 188 | callback(); |
| 189 | } else { |
| 190 | this[kCallback] = callback; |
| 191 | } |
| 192 | }); |
| 193 | }; |
| 194 | |
| 195 | Transform.prototype._read = function (_size) { |
| 196 | if (this[kCallback]) { |
| 197 | const callback = this[kCallback]; |
| 198 | this[kCallback] = null; |
| 199 | callback(); |
| 200 | } |
| 201 | }; |
| 202 | |
| 203 | Object.setPrototypeOf(PassThrough.prototype, Transform.prototype); |
| 204 | Object.setPrototypeOf(PassThrough, Transform); |
| 205 | |
| 206 | export function PassThrough(options) { |
| 207 | if (!(this instanceof PassThrough)) return new PassThrough(options); |
| 208 | Transform.call(this, { |
| 209 | ...options, |
| 210 | transform: undefined, |
| 211 | flush: undefined, |
| 212 | }); |
| 213 | } |
| 214 | |
| 215 | PassThrough.prototype._transform = function (chunk, _, cb) { |
| 216 | cb(null, chunk); |
| 217 | }; |