| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 217b495 commit 45780df
3 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -74,9 +74,7 @@ function duplex(options = { __proto__: null }) { | |||
| 74 | 74 | if (aClosed) return; | |
| 75 | 75 | aClosed = true; | |
| 76 | 76 | // End the writer (signals end-of-stream to B's readable) | |
| 77 | - if (aWriter.endSync() < 0) { | ||
| 78 | - await aWriter.end(); | ||
| 79 | - } | ||
| 77 | + aWriter.endSync(); | ||
| 80 | 78 | // Stop iteration of this channel's readable | |
| 81 | 79 | if (aReadableIterator?.return) { | |
| 82 | 80 | await aReadableIterator.return(); | |
@@ -104,9 +102,7 @@ function duplex(options = { __proto__: null }) { | |||
| 104 | 102 | async close() { | |
| 105 | 103 | if (bClosed) return; | |
| 106 | 104 | bClosed = true; | |
| 107 | - if (bWriter.endSync() < 0) { | ||
| 108 | - await bWriter.end(); | ||
| 109 | - } | ||
| 105 | + bWriter.endSync(); | ||
| 110 | 106 | if (bReadableIterator?.return) { | |
| 111 | 107 | await bReadableIterator.return(); | |
| 112 | 108 | bReadableIterator = null; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -274,7 +274,10 @@ class PushQueue { | |||
| 274 | 274 | if (this.#writerState === 'errored') { | |
| 275 | 275 | return -2; // Signal to reject with stored error | |
| 276 | 276 | } | |
| 277 | - if (this.#writerState === 'closing' || this.#writerState === 'closed') { | ||
| 277 | + if (this.#writerState === 'closing') { | ||
| 278 | + return -3; // Signal to PushWriter: wait for drain to complete | ||
| 279 | + } | ||
| 280 | + if (this.#writerState === 'closed') { | ||
| 278 | 281 | return this.#bytesWritten; // Idempotent | |
| 279 | 282 | } | |
| 280 | 283 | ||
@@ -636,6 +639,10 @@ class PushWriter { | |||
| 636 | 639 | if (result === -3) { | |
| 637 | 640 | // Closing: buffer has data, create deferred promise that resolves | |
| 638 | 641 | // when consumer drains past the end sentinel | |
| 642 | + const pendingEndPromise = this.#queue.pendingEndPromise; | ||
| 643 | + if (pendingEndPromise !== null) { | ||
| 644 | + return pendingEndPromise; | ||
| 645 | + } | ||
| 639 | 646 | const { promise, resolve, reject } = PromiseWithResolvers(); | |
| 640 | 647 | this.#queue.setPendingEnd({ __proto__: null, promise, resolve, reject }); | |
| 641 | 648 | return promise; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -232,6 +232,25 @@ async function testEndAsyncReturnValue() { | |||
| 232 | 232 | await consume; | |
| 233 | 233 | } | |
| 234 | 234 | ||
| 235 | + async function testEndAfterEndSyncWaitsForDrain() { | ||
| 236 | + const { writer, readable } = push(); | ||
| 237 | + writer.writeSync('hello'); | ||
| 238 | + assert.strictEqual(writer.endSync(), -1); | ||
| 239 | + | ||
| 240 | + let ended = false; | ||
| 241 | + const end = writer.end().then((n) => { | ||
| 242 | + ended = true; | ||
| 243 | + return n; | ||
| 244 | + }); | ||
| 245 | + | ||
| 246 | + await Promise.resolve(); | ||
| 247 | + assert.strictEqual(ended, false); | ||
| 248 | + | ||
| 249 | + // eslint-disable-next-line no-unused-vars | ||
| 250 | + for await (const _ of readable) { /* drain */ } | ||
| 251 | + assert.strictEqual(await end, 5); | ||
| 252 | + } | ||
| 253 | + | ||
| 235 | 254 | async function testWriteUint8Array() { | |
| 236 | 255 | const { writer, readable } = push(); | |
| 237 | 256 | writer.write(new Uint8Array([72, 73])); // 'HI' | |
@@ -413,6 +432,7 @@ Promise.all([ | |||
| 413 | 432 | testOndrainProtocolErrorPropagates(), | |
| 414 | 433 | testFail(), | |
| 415 | 434 | testEndAsyncReturnValue(), | |
| 435 | + testEndAfterEndSyncWaitsForDrain(), | ||
| 416 | 436 | testWriteUint8Array(), | |
| 417 | 437 | testOndrainWaitsForDrain(), | |
| 418 | 438 | testConsumerThrowRejectsWrites(), | |
| Back | FazBrowse Home | New Git URL |
0 commit comments