| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent e6bfec9 commit 815424d
5 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -116,17 +116,32 @@ class BroadcastImpl { | |||
| 116 | 116 | ||
| 117 | 117 | push(...args) { | |
| 118 | 118 | const { transforms, options } = parsePullArgs(args); | |
| 119 | + const signal = options?.signal; | ||
| 120 | + validateAbortSignal(signal, 'options.signal'); | ||
| 121 | + | ||
| 122 | + // Avoid registering a consumer that the pre-aborted pipeline will never | ||
| 123 | + // read or detach. | ||
| 124 | + if (signal?.aborted) { | ||
| 125 | + return { | ||
| 126 | + __proto__: null, | ||
| 127 | + // eslint-disable-next-line require-yield | ||
| 128 | + async *[SymbolAsyncIterator]() { | ||
| 129 | + throw signal.reason; | ||
| 130 | + }, | ||
| 131 | + }; | ||
| 132 | + } | ||
| 133 | + | ||
| 119 | 134 | const rawConsumer = this.#createRawConsumer(); | |
| 120 | 135 | ||
| 121 | 136 | // When transforms are present, delegate to pull() which creates its | |
| 122 | 137 | // own internal AbortController that follows the external signal. | |
| 123 | 138 | // When no transforms, return rawConsumer directly (controller elided | |
| 124 | 139 | // per PULL-02 optimization -- no transforms means no signal recipient). | |
| 125 | - if (transforms.length > 0 || options?.signal) { | ||
| 140 | + if (transforms.length > 0 || signal) { | ||
| 126 | 141 | const pullArgs = [...transforms]; | |
| 127 | - if (options?.signal) { | ||
| 142 | + if (signal) { | ||
| 128 | 143 | ArrayPrototypePush(pullArgs, | |
| 129 | - { __proto__: null, signal: options.signal }); | ||
| 144 | + { __proto__: null, signal }); | ||
| 130 | 145 | } | |
| 131 | 146 | return pullWithTransforms(rawConsumer, ...pullArgs); | |
| 132 | 147 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -93,11 +93,29 @@ class ShareImpl { | |||
| 93 | 93 | ||
| 94 | 94 | pull(...args) { | |
| 95 | 95 | const { transforms, options } = parsePullArgs(args); | |
| 96 | + const signal = options?.signal; | ||
| 97 | + validateAbortSignal(signal, 'options.signal'); | ||
| 98 | + | ||
| 99 | + // Avoid registering a consumer that the pre-aborted pipeline will never | ||
| 100 | + // read or detach. | ||
| 101 | + if (signal?.aborted) { | ||
| 102 | + return { | ||
| 103 | + __proto__: null, | ||
| 104 | + // eslint-disable-next-line require-yield | ||
| 105 | + async *[SymbolAsyncIterator]() { | ||
| 106 | + throw signal.reason; | ||
| 107 | + }, | ||
| 108 | + }; | ||
| 109 | + } | ||
| 110 | + | ||
| 96 | 111 | const rawConsumer = this.#createRawConsumer(); | |
| 97 | 112 | ||
| 98 | - if (transforms.length > 0 || options?.signal) { | ||
| 99 | - if (options) { | ||
| 100 | - return pullWithTransforms(rawConsumer, ...transforms, options); | ||
| 113 | + if (transforms.length > 0 || signal) { | ||
| 114 | + if (signal) { | ||
| 115 | + return pullWithTransforms( | ||
| 116 | + rawConsumer, | ||
| 117 | + ...transforms, | ||
| 118 | + { __proto__: null, signal }); | ||
| 101 | 119 | } | |
| 102 | 120 | return pullWithTransforms(rawConsumer, ...transforms); | |
| 103 | 121 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -203,6 +203,17 @@ async function testPushAbortSignalRejectsPendingNext() { | |||
| 203 | 203 | await rejected; | |
| 204 | 204 | } | |
| 205 | 205 | ||
| 206 | + async function testPushPreAbortedSignalDoesNotAddConsumer() { | ||
| 207 | + const reason = new Error('already aborted'); | ||
| 208 | + const signal = AbortSignal.abort(reason); | ||
| 209 | + const { broadcast: bc } = broadcast(); | ||
| 210 | + const iter = bc.push({ signal })[Symbol.asyncIterator](); | ||
| 211 | + | ||
| 212 | + assert.strictEqual(bc.consumerCount, 0); | ||
| 213 | + await assert.rejects(iter.next(), (error) => error === reason); | ||
| 214 | + assert.strictEqual(bc.consumerCount, 0); | ||
| 215 | + } | ||
| 216 | + | ||
| 206 | 217 | // ============================================================================= | |
| 207 | 218 | // Writer fail detaches consumers | |
| 208 | 219 | // ============================================================================= | |
@@ -331,6 +342,7 @@ Promise.all([ | |||
| 331 | 342 | testCancelWithFalsyReason(), | |
| 332 | 343 | testPendingNextSettlesAfterReturn(), | |
| 333 | 344 | testPushAbortSignalRejectsPendingNext(), | |
| 345 | + testPushPreAbortedSignalDoesNotAddConsumer(), | ||
| 334 | 346 | testFailDetachesConsumers(), | |
| 335 | 347 | testWriterFailIdempotent(), | |
| 336 | 348 | testLateJoinerSeesBufferedData(), | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -226,6 +226,17 @@ async function testSharePullAbortSignalRejectsPendingNext() { | |||
| 226 | 226 | shared.cancel(); | |
| 227 | 227 | } | |
| 228 | 228 | ||
| 229 | + async function testSharePullPreAbortedSignalDoesNotAddConsumer() { | ||
| 230 | + const reason = new Error('already aborted'); | ||
| 231 | + const signal = AbortSignal.abort(reason); | ||
| 232 | + const shared = share(from('data')); | ||
| 233 | + const iter = shared.pull({ signal })[Symbol.asyncIterator](); | ||
| 234 | + | ||
| 235 | + assert.strictEqual(shared.consumerCount, 0); | ||
| 236 | + await assert.rejects(iter.next(), (error) => error === reason); | ||
| 237 | + assert.strictEqual(shared.consumerCount, 0); | ||
| 238 | + } | ||
| 239 | + | ||
| 229 | 240 | async function testShareAlreadyAborted() { | |
| 230 | 241 | const shared = share(from('data'), { signal: AbortSignal.abort() }); | |
| 231 | 242 | const consumer = shared.pull(); | |
@@ -372,6 +383,7 @@ Promise.all([ | |||
| 372 | 383 | testShareAbortSignal(), | |
| 373 | 384 | testShareAbortSignalWhileSourcePullPending(), | |
| 374 | 385 | testSharePullAbortSignalRejectsPendingNext(), | |
| 386 | + testSharePullPreAbortedSignalDoesNotAddConsumer(), | ||
| 375 | 387 | testShareAlreadyAborted(), | |
| 376 | 388 | testShareSourceError(), | |
| 377 | 389 | testShareLateJoiningConsumer(), | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -156,6 +156,15 @@ assert.throws(() => broadcast({ budget: 16383 }), { code: 'ERR_OUT_OF_RANGE' }); | |||
| 156 | 156 | assert.throws(() => broadcast({ signal: {} }), { code: 'ERR_INVALID_ARG_TYPE' }); | |
| 157 | 157 | assert.throws(() => broadcast({ backpressure: 'bad' }), { code: 'ERR_INVALID_ARG_VALUE' }); | |
| 158 | 158 | ||
| 159 | + // Broadcast consumer options.signal must be AbortSignal and validation must | ||
| 160 | + // not leave a consumer registered. | ||
| 161 | + { | ||
| 162 | + const { broadcast: bc } = broadcast(); | ||
| 163 | + assert.throws(() => bc.push({ signal: {} }), | ||
| 164 | + { code: 'ERR_INVALID_ARG_TYPE' }); | ||
| 165 | + assert.strictEqual(bc.consumerCount, 0); | ||
| 166 | + } | ||
| 167 | + | ||
| 159 | 168 | // BroadcastWriter options.signal must be AbortSignal | |
| 160 | 169 | { | |
| 161 | 170 | const { writer } = broadcast(); | |
@@ -212,6 +221,15 @@ assert.throws(() => share(from('a'), { budget: Number.MAX_SAFE_INTEGER + 1 }), | |||
| 212 | 221 | assert.throws(() => share(from('a'), { signal: {} }), { code: 'ERR_INVALID_ARG_TYPE' }); | |
| 213 | 222 | assert.throws(() => share(from('a'), { backpressure: 'bad' }), { code: 'ERR_INVALID_ARG_VALUE' }); | |
| 214 | 223 | ||
| 224 | + // Share consumer options.signal must be AbortSignal and validation must not | ||
| 225 | + // leave a consumer registered. | ||
| 226 | + { | ||
| 227 | + const shared = share(from('a')); | ||
| 228 | + assert.throws(() => shared.pull({ signal: {} }), | ||
| 229 | + { code: 'ERR_INVALID_ARG_TYPE' }); | ||
| 230 | + assert.strictEqual(shared.consumerCount, 0); | ||
| 231 | + } | ||
| 232 | + | ||
| 215 | 233 | // share() values < 16384 are rejected | |
| 216 | 234 | assert.throws(() => share(from('a'), { budget: 0 }), { code: 'ERR_OUT_OF_RANGE' }); | |
| 217 | 235 | assert.throws(() => share(from('a'), { budget: -1 }), { code: 'ERR_OUT_OF_RANGE' }); | |
| Back | FazBrowse Home | New Git URL |
0 commit comments