Skip to content

Commit 457b5cb

Browse files
Y1D7NGPuffinRun
authored andcommitted
stream: destroy source when map iterator closes early
Fixes: #64261 Signed-off-by: y1d7ng <y1d7ng@yeah.net>
1 parent 2dfdb6a commit 457b5cb

2 files changed

Lines changed: 33 additions & 1 deletion

File tree

lib/internal/streams/operators.js

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ const {
2929
validateFunction,
3030
} = require('internal/validators');
3131
const { kWeakHandler, kResistStopPropagation } = require('internal/event_target');
32+
const destroyImpl = require('internal/streams/destroy');
3233
const { finished } = require('internal/streams/end-of-stream');
3334

3435
const kEmpty = Symbol('kEmpty');
@@ -177,6 +178,7 @@ function map(fn, options) {
177178
resume();
178179
resume = null;
179180
}
181+
destroyImpl.destroyer(stream, null);
180182
}
181183
}.call(this);
182184
}

test/parallel/test-stream-some-find-every.mjs

Lines changed: 31 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import * as common from '../common/index.mjs';
22
import { setTimeout } from 'timers/promises';
3-
import { Readable } from 'stream';
3+
import { PassThrough, Readable } from 'stream';
44
import assert from 'assert';
55

66

@@ -77,6 +77,36 @@ function oneTo5Async() {
7777
await findStream.find(common.mustCall((x) => x > 1, 2));
7878
await checkDestroyed(findStream);
7979

80+
const openSomeStream = new PassThrough({ objectMode: true });
81+
openSomeStream.write(1);
82+
openSomeStream.write(2);
83+
openSomeStream.write(3);
84+
assert.strictEqual(
85+
await openSomeStream.some(common.mustCall((x) => x > 2, 3)),
86+
true,
87+
);
88+
await checkDestroyed(openSomeStream);
89+
90+
const openEveryStream = new PassThrough({ objectMode: true });
91+
openEveryStream.write(1);
92+
openEveryStream.write(2);
93+
openEveryStream.write(3);
94+
assert.strictEqual(
95+
await openEveryStream.every(common.mustCall((x) => x < 3, 3)),
96+
false,
97+
);
98+
await checkDestroyed(openEveryStream);
99+
100+
const openFindStream = new PassThrough({ objectMode: true });
101+
openFindStream.write(1);
102+
openFindStream.write(2);
103+
openFindStream.write(3);
104+
assert.strictEqual(
105+
await openFindStream.find(common.mustCall((x) => x > 1, 2)),
106+
2,
107+
);
108+
await checkDestroyed(openFindStream);
109+
80110
// When short circuit isn't possible the whole stream is iterated
81111
await oneTo5().some(common.mustCall(() => false, 5));
82112
await oneTo5().every(common.mustCall(() => true, 5));

0 commit comments

Comments
 (0)