Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion lib/internal/streams/iter/broadcast.js
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ class BroadcastImpl {
// own internal AbortController that follows the external signal.
// When no transforms, return rawConsumer directly (controller elided
// per PULL-02 optimization -- no transforms means no signal recipient).
if (transforms.length > 0) {
if (transforms.length > 0 || options?.signal) {
const pullArgs = [...transforms];
if (options?.signal) {
ArrayPrototypePush(pullArgs,
Expand Down
2 changes: 1 addition & 1 deletion lib/internal/streams/iter/share.js
Original file line number Diff line number Diff line change
Expand Up @@ -97,7 +97,7 @@ class ShareImpl {
const { transforms, options } = parsePullArgs(args);
const rawConsumer = this.#createRawConsumer();

if (transforms.length > 0) {
if (transforms.length > 0 || options?.signal) {
if (options) {
return pullWithTransforms(rawConsumer, ...transforms, options);
}
Expand Down
14 changes: 14 additions & 0 deletions test/parallel/test-stream-iter-broadcast-basic.js
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,19 @@ async function testPendingNextSettlesAfterReturn() {
assert.strictEqual(result.value, undefined);
}

async function testPushAbortSignalRejectsPendingNext() {
const ac = new AbortController();
const reason = new Error('push aborted');
const { broadcast: bc } = broadcast();
const iter = bc.push({ signal: ac.signal })[Symbol.asyncIterator]();

const pendingNext = iter.next();
const rejected = assert.rejects(pendingNext, (error) => error === reason);
ac.abort(reason);

await rejected;
}

// =============================================================================
// Writer fail detaches consumers
// =============================================================================
Expand Down Expand Up @@ -267,6 +280,7 @@ Promise.all([
testCancelWithReason(),
testCancelWithFalsyReason(),
testPendingNextSettlesAfterReturn(),
testPushAbortSignalRejectsPendingNext(),
testFailDetachesConsumers(),
testWriterFailIdempotent(),
testLateJoinerSeesBufferedData(),
Expand Down
20 changes: 20 additions & 0 deletions test/parallel/test-stream-iter-share-async.js
Original file line number Diff line number Diff line change
Expand Up @@ -196,6 +196,25 @@ async function testShareAbortSignalWhileSourcePullPending() {
await Promise.all([rejected1, rejected2]);
}

async function testSharePullAbortSignalRejectsPendingNext() {
const ac = new AbortController();
const reason = new Error('pull aborted');
const shared = share(
// eslint-disable-next-line require-yield
(async function* never() {
await new Promise(() => {});
})(),
);
const iter = shared.pull({ signal: ac.signal })[Symbol.asyncIterator]();

const pendingNext = iter.next();
const rejected = assert.rejects(pendingNext, (error) => error === reason);
ac.abort(reason);

await rejected;
shared.cancel();
}

async function testShareAlreadyAborted() {
const shared = share(from('data'), { signal: AbortSignal.abort() });
const consumer = shared.pull();
Expand Down Expand Up @@ -340,6 +359,7 @@ Promise.all([
testShareCancelWithReason(),
testShareAbortSignal(),
testShareAbortSignalWhileSourcePullPending(),
testSharePullAbortSignalRejectsPendingNext(),
testShareAlreadyAborted(),
testShareSourceError(),
testShareLateJoiningConsumer(),
Expand Down
Loading