| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -620,7 +620,11 @@ class ReadableStream { | |||
| 620 | 620 | const transfer = lazyTransfer(); | |
| 621 | 621 | setupReadableStreamDefaultControllerFromSource( | |
| 622 | 622 | this, | |
| 623 | - new transfer.CrossRealmTransformReadableSource(port), | ||
| 623 | + // The MessagePort is set to be referenced when reading. | ||
| 624 | + // After two MessagePorts are closed, there is a problem with | ||
| 625 | + // lingering promise not being properly resolved. | ||
| 626 | + // https://github.com/nodejs/node/issues/51486 | ||
| 627 | + new transfer.CrossRealmTransformReadableSource(port, true), | ||
| 624 | 628 | 0, () => 1); | |
| 625 | 629 | } | |
| 626 | 630 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -104,10 +104,11 @@ function InternalCloneableDOMException() { | |||
| 104 | 104 | InternalCloneableDOMException[kDeserialize] = () => {}; | |
| 105 | 105 | ||
| 106 | 106 | class CrossRealmTransformReadableSource { | |
| 107 | - constructor(port) { | ||
| 107 | + constructor(port, unref) { | ||
| 108 | 108 | this[kState] = { | |
| 109 | 109 | port, | |
| 110 | 110 | controller: undefined, | |
| 111 | + unref, | ||
| 111 | 112 | }; | |
| 112 | 113 | ||
| 113 | 114 | port.onmessage = ({ data }) => { | |
@@ -145,13 +146,19 @@ class CrossRealmTransformReadableSource { | |||
| 145 | 146 | error); | |
| 146 | 147 | port.close(); | |
| 147 | 148 | }; | |
| 149 | + | ||
| 150 | + port.unref(); | ||
| 148 | 151 | } | |
| 149 | 152 | ||
| 150 | 153 | start(controller) { | |
| 151 | 154 | this[kState].controller = controller; | |
| 152 | 155 | } | |
| 153 | 156 | ||
| 154 | 157 | async pull() { | |
| 158 | + if (this[kState].unref) { | ||
| 159 | + this[kState].unref = false; | ||
| 160 | + this[kState].port.ref(); | ||
| 161 | + } | ||
| 155 | 162 | this[kState].port.postMessage({ type: 'pull' }); | |
| 156 | 163 | } | |
| 157 | 164 | ||
@@ -172,11 +179,12 @@ class CrossRealmTransformReadableSource { | |||
| 172 | 179 | } | |
| 173 | 180 | ||
| 174 | 181 | class CrossRealmTransformWritableSink { | |
| 175 | - constructor(port) { | ||
| 182 | + constructor(port, unref) { | ||
| 176 | 183 | this[kState] = { | |
| 177 | 184 | port, | |
| 178 | 185 | controller: undefined, | |
| 179 | 186 | backpressurePromise: createDeferredPromise(), | |
| 187 | + unref, | ||
| 180 | 188 | }; | |
| 181 | 189 | ||
| 182 | 190 | port.onmessage = ({ data }) => { | |
@@ -213,13 +221,18 @@ class CrossRealmTransformWritableSink { | |||
| 213 | 221 | port.close(); | |
| 214 | 222 | }; | |
| 215 | 223 | ||
| 224 | + port.unref(); | ||
| 216 | 225 | } | |
| 217 | 226 | ||
| 218 | 227 | start(controller) { | |
| 219 | 228 | this[kState].controller = controller; | |
| 220 | 229 | } | |
| 221 | 230 | ||
| 222 | 231 | async write(chunk) { | |
| 232 | + if (this[kState].unref) { | ||
| 233 | + this[kState].unref = false; | ||
| 234 | + this[kState].port.ref(); | ||
| 235 | + } | ||
| 223 | 236 | if (this[kState].backpressurePromise === undefined) { | |
| 224 | 237 | this[kState].backpressurePromise = { | |
| 225 | 238 | promise: PromiseResolve(), | |
@@ -264,12 +277,12 @@ class CrossRealmTransformWritableSink { | |||
| 264 | 277 | } | |
| 265 | 278 | ||
| 266 | 279 | function newCrossRealmReadableStream(writable, port) { | |
| 267 | - const readable = | ||
| 268 | - new ReadableStream( | ||
| 269 | - new CrossRealmTransformReadableSource(port)); | ||
| 280 | + // MessagePort should always be unref. | ||
| 281 | + // There is a problem with the process not terminating. | ||
| 282 | + // https://github.com/nodejs/node/issues/44985 | ||
| 283 | + const readable = new ReadableStream(new CrossRealmTransformReadableSource(port, false)); | ||
| 270 | 284 | ||
| 271 | - const promise = | ||
| 272 | - readableStreamPipeTo(readable, writable, false, false, false); | ||
| 285 | + const promise = readableStreamPipeTo(readable, writable, false, false, false); | ||
| 273 | 286 | ||
| 274 | 287 | setPromiseHandled(promise); | |
| 275 | 288 | ||
@@ -280,12 +293,15 @@ function newCrossRealmReadableStream(writable, port) { | |||
| 280 | 293 | } | |
| 281 | 294 | ||
| 282 | 295 | function newCrossRealmWritableSink(readable, port) { | |
| 283 | - const writable = | ||
| 284 | - new WritableStream( | ||
| 285 | - new CrossRealmTransformWritableSink(port)); | ||
| 296 | + // MessagePort should always be unref. | ||
| 297 | + // There is a problem with the process not terminating. | ||
| 298 | + // https://github.com/nodejs/node/issues/44985 | ||
| 299 | + const writable = new WritableStream(new CrossRealmTransformWritableSink(port, false)); | ||
| 286 | 300 | ||
| 287 | 301 | const promise = readableStreamPipeTo(readable, writable, false, false, false); | |
| 302 | + | ||
| 288 | 303 | setPromiseHandled(promise); | |
| 304 | + | ||
| 289 | 305 | return { | |
| 290 | 306 | writable, | |
| 291 | 307 | promise, | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -294,8 +294,6 @@ class WritableStream { | |||
| 294 | 294 | this[kState].transfer.readable = readable; | |
| 295 | 295 | this[kState].transfer.promise = promise; | |
| 296 | 296 | ||
| 297 | - setPromiseHandled(this[kState].transfer.promise); | ||
| 298 | - | ||
| 299 | 297 | return { | |
| 300 | 298 | data: { port: this[kState].transfer.port2 }, | |
| 301 | 299 | deserializeInfo: | |
@@ -314,7 +312,11 @@ class WritableStream { | |||
| 314 | 312 | const transfer = lazyTransfer(); | |
| 315 | 313 | setupWritableStreamDefaultControllerFromSink( | |
| 316 | 314 | this, | |
| 317 | - new transfer.CrossRealmTransformWritableSink(port), | ||
| 315 | + // The MessagePort is set to be referenced when reading. | ||
| 316 | + // After two MessagePorts are closed, there is a problem with | ||
| 317 | + // lingering promise not being properly resolved. | ||
| 318 | + // https://github.com/nodejs/node/issues/51486 | ||
| 319 | + new transfer.CrossRealmTransformWritableSink(port, true), | ||
| 318 | 320 | 1, | |
| 319 | 321 | () => 1); | |
| 320 | 322 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -0,0 +1,16 @@ | |||
| 1 | + 'use strict'; | ||
| 2 | + | ||
| 3 | + require('../common'); | ||
| 4 | + const { ok } = require('node:assert'); | ||
| 5 | + | ||
| 6 | + // This test verifies that cloned ReadableStream and WritableStream instances | ||
| 7 | + // do not keep the process alive. The test fails if it timesout (it should just | ||
| 8 | + // exit immediately) | ||
| 9 | + | ||
| 10 | + const rs1 = new ReadableStream(); | ||
| 11 | + const ws1 = new WritableStream(); | ||
| 12 | + | ||
| 13 | + const [rs2, ws2] = structuredClone([rs1, ws1], { transfer: [rs1, ws1] }); | ||
| 14 | + | ||
| 15 | + ok(rs2 instanceof ReadableStream); | ||
| 16 | + ok(ws2 instanceof WritableStream); | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -454,12 +454,23 @@ const theData = 'hello'; | |||
| 454 | 454 | tracker.verify(); | |
| 455 | 455 | }); | |
| 456 | 456 | ||
| 457 | + // We create an interval to keep the event loop alive while | ||
| 458 | + // we wait for the stream read to complete. The reason this is needed is because there's | ||
| 459 | + // otherwise nothing to keep the worker thread event loop alive long enough to actually | ||
| 460 | + // complete the read from the stream. Under the covers the ReadableStream uses an | ||
| 461 | + // unref'd MessagePort to communicate with the main thread. Because the MessagePort | ||
| 462 | + // is unref'd, it's existence would not keep the thread alive on its own. There was previously | ||
| 463 | + // a bug where this MessagePort was ref'd which would block the thread and main thread | ||
| 464 | + // from terminating at all unless the stream was consumed/closed. | ||
| 465 | + const i = setInterval(() => {}, 1000); | ||
| 466 | + | ||
| 457 | 467 | parentPort.onmessage = tracker.calls(({ data }) => { | |
| 458 | 468 | assert(isReadableStream(data)); | |
| 459 | 469 | const reader = data.getReader(); | |
| 460 | 470 | reader.read().then(tracker.calls((result) => { | |
| 461 | 471 | assert(!result.done); | |
| 462 | 472 | assert(result.value instanceof Uint8Array); | |
| 473 | + clearInterval(i); | ||
| 463 | 474 | })); | |
| 464 | 475 | parentPort.close(); | |
| 465 | 476 | }); | |
| Back | FazBrowse Home | New Git URL |
0 commit comments