|
| 1 | +// Flags: --experimental-stream-iter |
| 2 | +'use strict'; |
| 3 | + |
| 4 | +constcommon=require('../common'); |
| 5 | +constassert=require('assert'); |
| 6 | +const{ broadcast, text }=require('stream/iter'); |
| 7 | + |
| 8 | +// ============================================================================= |
| 9 | +// Basic broadcast |
| 10 | +// ============================================================================= |
| 11 | + |
| 12 | +asyncfunctiontestBasicBroadcast(){ |
| 13 | +const{ writer,broadcast: bc}=broadcast(); |
| 14 | + |
| 15 | +// Create two consumers |
| 16 | +constconsumer1=bc.push(); |
| 17 | +constconsumer2=bc.push(); |
| 18 | + |
| 19 | +assert.strictEqual(bc.consumerCount,2); |
| 20 | + |
| 21 | +awaitwriter.write('hello'); |
| 22 | +awaitwriter.end(); |
| 23 | + |
| 24 | +const[data1,data2]=awaitPromise.all([ |
| 25 | +text(consumer1), |
| 26 | +text(consumer2), |
| 27 | +]); |
| 28 | + |
| 29 | +assert.strictEqual(data1,'hello'); |
| 30 | +assert.strictEqual(data2,'hello'); |
| 31 | +} |
| 32 | + |
| 33 | +asyncfunctiontestMultipleWrites(){ |
| 34 | +const{ writer,broadcast: bc}=broadcast({highWaterMark: 10}); |
| 35 | + |
| 36 | +constconsumer=bc.push(); |
| 37 | + |
| 38 | +awaitwriter.write('a'); |
| 39 | +awaitwriter.write('b'); |
| 40 | +awaitwriter.write('c'); |
| 41 | +awaitwriter.end(); |
| 42 | + |
| 43 | +constdata=awaittext(consumer); |
| 44 | +assert.strictEqual(data,'abc'); |
| 45 | +} |
| 46 | + |
| 47 | +asyncfunctiontestConsumerCount(){ |
| 48 | +const{broadcast: bc}=broadcast(); |
| 49 | + |
| 50 | +assert.strictEqual(bc.consumerCount,0); |
| 51 | + |
| 52 | +constc1=bc.push(); |
| 53 | +assert.strictEqual(bc.consumerCount,1); |
| 54 | + |
| 55 | +bc.push(); |
| 56 | +assert.strictEqual(bc.consumerCount,2); |
| 57 | + |
| 58 | +bc.cancel(); |
| 59 | + |
| 60 | +// After cancel, consumer count drops to 0 |
| 61 | +assert.strictEqual(bc.consumerCount,0); |
| 62 | + |
| 63 | +// Consumers are detached and yield nothing |
| 64 | +constbatches=[]; |
| 65 | +forawait(constbatchofc1){ |
| 66 | +batches.push(batch); |
| 67 | +} |
| 68 | +assert.strictEqual(batches.length,0); |
| 69 | +} |
| 70 | + |
| 71 | +// ============================================================================= |
| 72 | +// Writer methods |
| 73 | +// ============================================================================= |
| 74 | + |
| 75 | +asyncfunctiontestWriteSync(){ |
| 76 | +const{ writer,broadcast: bc}=broadcast({highWaterMark: 2}); |
| 77 | +constconsumer=bc.push(); |
| 78 | + |
| 79 | +assert.strictEqual(writer.writeSync('a'),true); |
| 80 | +assert.strictEqual(writer.writeSync('b'),true); |
| 81 | +// Buffer full (highWaterMark=2, strict policy) |
| 82 | +assert.strictEqual(writer.writeSync('c'),false); |
| 83 | + |
| 84 | +writer.endSync(); |
| 85 | + |
| 86 | +constdata=awaittext(consumer); |
| 87 | +assert.strictEqual(data,'ab'); |
| 88 | +} |
| 89 | + |
| 90 | +asyncfunctiontestWritevSync(){ |
| 91 | +const{ writer,broadcast: bc}=broadcast({highWaterMark: 10}); |
| 92 | +constconsumer=bc.push(); |
| 93 | + |
| 94 | +assert.strictEqual(writer.writevSync(['hello',' ','world']),true); |
| 95 | +writer.endSync(); |
| 96 | + |
| 97 | +constdata=awaittext(consumer); |
| 98 | +assert.strictEqual(data,'hello world'); |
| 99 | +} |
| 100 | + |
| 101 | +asyncfunctiontestWriterEnd(){ |
| 102 | +const{ writer,broadcast: bc}=broadcast(); |
| 103 | +constconsumer=bc.push(); |
| 104 | + |
| 105 | +awaitwriter.write('data'); |
| 106 | +consttotalBytes=awaitwriter.end(); |
| 107 | +assert.strictEqual(totalBytes,4);// 'data' = 4 UTF-8 bytes |
| 108 | + |
| 109 | +constdata=awaittext(consumer); |
| 110 | +assert.strictEqual(data,'data'); |
| 111 | +} |
| 112 | + |
| 113 | +asyncfunctiontestWriterFail(){ |
| 114 | +const{ writer,broadcast: bc}=broadcast(); |
| 115 | +constconsumer=bc.push(); |
| 116 | + |
| 117 | +writer.fail(newError('test error')); |
| 118 | + |
| 119 | +awaitassert.rejects( |
| 120 | +async()=>{ |
| 121 | +// eslint-disable-next-line no-unused-vars |
| 122 | +forawait(const_ofconsumer){ |
| 123 | +assert.fail('Should not reach here'); |
| 124 | +} |
| 125 | +}, |
| 126 | +{message: 'test error'}, |
| 127 | +); |
| 128 | +} |
| 129 | + |
| 130 | +// ============================================================================= |
| 131 | +// Cancel |
| 132 | +// ============================================================================= |
| 133 | + |
| 134 | +asyncfunctiontestCancelWithoutReason(){ |
| 135 | +const{broadcast: bc}=broadcast(); |
| 136 | +constconsumer=bc.push(); |
| 137 | + |
| 138 | +bc.cancel(); |
| 139 | + |
| 140 | +constbatches=[]; |
| 141 | +forawait(constbatchofconsumer){ |
| 142 | +batches.push(batch); |
| 143 | +} |
| 144 | +assert.strictEqual(batches.length,0); |
| 145 | +} |
| 146 | + |
| 147 | +asyncfunctiontestCancelWithReason(){ |
| 148 | +const{broadcast: bc}=broadcast(); |
| 149 | + |
| 150 | +// Start a consumer that is waiting for data (promise pending) |
| 151 | +constconsumer=bc.push(); |
| 152 | +constresultPromise=text(consumer).catch((err)=>err); |
| 153 | + |
| 154 | +// Give the consumer time to enter the waiting state |
| 155 | +awaitnewPromise((resolve)=>setImmediate(resolve)); |
| 156 | + |
| 157 | +bc.cancel(newError('cancelled')); |
| 158 | + |
| 159 | +constresult=awaitresultPromise; |
| 160 | +assert.ok(resultinstanceofError); |
| 161 | +assert.strictEqual(result.message,'cancelled'); |
| 162 | +} |
| 163 | + |
| 164 | +// ============================================================================= |
| 165 | +// Writer fail detaches consumers |
| 166 | +// ============================================================================= |
| 167 | + |
| 168 | +asyncfunctiontestFailDetachesConsumers(){ |
| 169 | +const{ writer,broadcast: bc}=broadcast(); |
| 170 | +constconsumer1=bc.push(); |
| 171 | +constconsumer2=bc.push(); |
| 172 | + |
| 173 | +assert.strictEqual(bc.consumerCount,2); |
| 174 | + |
| 175 | +// Write some data, then fail the writer |
| 176 | +awaitwriter.write('data'); |
| 177 | +awaitwriter.fail(newError('writer failed')); |
| 178 | + |
| 179 | +// After fail, consumers are detached |
| 180 | +assert.strictEqual(bc.consumerCount,0); |
| 181 | + |
| 182 | +// Both consumers should see the error |
| 183 | +awaitassert.rejects( |
| 184 | +async()=>{ |
| 185 | +// eslint-disable-next-line no-unused-vars |
| 186 | +forawait(const_ofconsumer1){ |
| 187 | +assert.fail('Should not reach here'); |
| 188 | +} |
| 189 | +}, |
| 190 | +{message: 'writer failed'}, |
| 191 | +); |
| 192 | + |
| 193 | +awaitassert.rejects( |
| 194 | +async()=>{ |
| 195 | +// eslint-disable-next-line no-unused-vars |
| 196 | +forawait(const_ofconsumer2){ |
| 197 | +assert.fail('Should not reach here'); |
| 198 | +} |
| 199 | +}, |
| 200 | +{message: 'writer failed'}, |
| 201 | +); |
| 202 | +} |
| 203 | + |
| 204 | +// ============================================================================= |
| 205 | +// Writer fail idempotent |
| 206 | +// ============================================================================= |
| 207 | + |
| 208 | +asyncfunctiontestWriterFailIdempotent(){ |
| 209 | +const{ writer,broadcast: bc}=broadcast(); |
| 210 | +constconsumer=bc.push(); |
| 211 | +writer.writeSync('hello'); |
| 212 | +writer.fail(newError('fail!')); |
| 213 | +// Second call is a no-op (already errored) |
| 214 | +writer.fail(newError('fail2')); |
| 215 | +awaitassert.rejects(async()=>{ |
| 216 | +// eslint-disable-next-line no-unused-vars |
| 217 | +forawait(const_ofconsumer){/* consume */} |
| 218 | +},{message: 'fail!'}); |
| 219 | +} |
| 220 | + |
| 221 | +// cancel() with falsy reason (0, "", false) should still treat as error |
| 222 | +asyncfunctiontestCancelWithFalsyReason(){ |
| 223 | +const{broadcast: bc}=broadcast(); |
| 224 | +constconsumer=bc.push(); |
| 225 | +constresultPromise=text(consumer).catch((err)=>err); |
| 226 | +awaitnewPromise((resolve)=>setImmediate(resolve)); |
| 227 | +bc.cancel(0); |
| 228 | +constresult=awaitresultPromise; |
| 229 | +assert.strictEqual(result,0); |
| 230 | +} |
| 231 | + |
| 232 | +// Late-joining consumer should read from oldest buffered entry |
| 233 | +asyncfunctiontestLateJoinerSeesBufferedData(){ |
| 234 | +const{ writer,broadcast: bc}=broadcast({highWaterMark: 16}); |
| 235 | + |
| 236 | +// Write data before any consumer joins |
| 237 | +writer.writeSync('before-join'); |
| 238 | +writer.endSync(); |
| 239 | + |
| 240 | +// Consumer joins after data is written |
| 241 | +constconsumer=bc.push(); |
| 242 | +constresult=awaittext(consumer); |
| 243 | +assert.strictEqual(result,'before-join'); |
| 244 | +} |
| 245 | + |
| 246 | +Promise.all([ |
| 247 | +testBasicBroadcast(), |
| 248 | +testMultipleWrites(), |
| 249 | +testConsumerCount(), |
| 250 | +testWriteSync(), |
| 251 | +testWritevSync(), |
| 252 | +testWriterEnd(), |
| 253 | +testWriterFail(), |
| 254 | +testCancelWithoutReason(), |
| 255 | +testCancelWithReason(), |
| 256 | +testCancelWithFalsyReason(), |
| 257 | +testFailDetachesConsumers(), |
| 258 | +testWriterFailIdempotent(), |
| 259 | +testLateJoinerSeesBufferedData(), |
| 260 | +]).then(common.mustCall()); |
0 commit comments