| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 2551f9e commit 575dc7d
3 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1280,19 +1280,14 @@ function waitForDrain(stream) { | |||
| 1280 | 1280 | ||
| 1281 | 1281 | // Writes a batch to the handle, awaiting drain if backpressured. | |
| 1282 | 1282 | // Returns true if the stream was destroyed during the wait. | |
| 1283 | - // Checks writeDesiredSize before writing to enforce backpressure | ||
| 1284 | - // against the outbound DataQueue's uncommitted bytes. | ||
| 1283 | + // Only waits when writeDesiredSize is 0 (no capacity at all). | ||
| 1284 | + // When there is any capacity, the write proceeds even if the batch | ||
| 1285 | + // is larger -- the C++ side buffers the data and writeDesiredSize | ||
| 1286 | + // drops toward 0, letting the normal drain mechanism take over. | ||
| 1285 | 1287 | async function writeBatchWithDrain(handle, stream, batch) { | |
| 1286 | 1288 | const state = getQuicStreamState(stream); | |
| 1287 | 1289 | ||
| 1288 | - // Calculate total batch size for the capacity check. | ||
| 1289 | - let len = 0; | ||
| 1290 | - for (const chunk of batch) len += TypedArrayPrototypeGetByteLength(chunk); | ||
| 1291 | - | ||
| 1292 | - // If insufficient capacity, wait for the C++ drain signal which | ||
| 1293 | - // fires when writeDesiredSize transitions from 0 to > 0 (i.e., | ||
| 1294 | - // ngtcp2 has consumed data from the outbound DataQueue). | ||
| 1295 | - if (len > state.writeDesiredSize) { | ||
| 1290 | + if (state.writeDesiredSize === 0) { | ||
| 1296 | 1291 | await waitForDrain(stream); | |
| 1297 | 1292 | if (stream.destroyed) return true; | |
| 1298 | 1293 | } | |
@@ -2030,9 +2025,15 @@ class QuicStream { | |||
| 2030 | 2025 | chunk = toUint8Array(chunk); | |
| 2031 | 2026 | const len = TypedArrayPrototypeGetByteLength(chunk); | |
| 2032 | 2027 | if (len === 0) return true; | |
| 2033 | - // Refuse the write if the chunk doesn't fit in the available | ||
| 2034 | - // buffer capacity. The caller should wait for drain and retry. | ||
| 2035 | - if (len > stream.#state.writeDesiredSize) return false; | ||
| 2028 | + // Refuse the write only when there is no available capacity at | ||
| 2029 | + // all. When writeDesiredSize > 0 we allow the write even if the | ||
| 2030 | + // chunk is larger than the remaining capacity -- the C++ side | ||
| 2031 | + // will accept the data into the DataQueue and | ||
| 2032 | + // UpdateWriteDesiredSize() will drop writeDesiredSize toward 0, | ||
| 2033 | + // at which point the standard drain mechanism takes over. | ||
| 2034 | + // This follows the Web Streams model where writes beyond the HWM | ||
| 2035 | + // succeed and backpressure applies to *subsequent* writes. | ||
| 2036 | + if (stream.#state.writeDesiredSize === 0) return false; | ||
| 2036 | 2037 | const result = handle.write([chunk]); | |
| 2037 | 2038 | if (result === undefined) return false; | |
| 2038 | 2039 | totalBytesWritten += len; | |
@@ -2071,7 +2072,7 @@ class QuicStream { | |||
| 2071 | 2072 | let len = 0; | |
| 2072 | 2073 | for (const c of chunks) len += TypedArrayPrototypeGetByteLength(c); | |
| 2073 | 2074 | if (len === 0) return true; | |
| 2074 | - if (len > stream.#state.writeDesiredSize) return false; | ||
| 2075 | + if (stream.#state.writeDesiredSize === 0) return false; | ||
| 2075 | 2076 | const result = handle.write(chunks); | |
| 2076 | 2077 | if (result === undefined) return false; | |
| 2077 | 2078 | totalBytesWritten += len; | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -1592,8 +1592,10 @@ void Stream::UpdateWriteDesiredSize() { | |||
| 1592 | 1592 | uint32_t old_size = state_->write_desired_size; | |
| 1593 | 1593 | state_->write_desired_size = clamped; | |
| 1594 | 1594 | ||
| 1595 | - // Fire drain when transitioning from 0 to non-zero | ||
| 1596 | - if (old_size == 0 && desired > 0) { | ||
| 1595 | + // Fire drain when transitioning from 0 to non-zero. | ||
| 1596 | + // writeDesiredSize == 0 means the buffer is full or flow control is | ||
| 1597 | + // exhausted, so the JS side may be waiting for capacity. | ||
| 1598 | + if (old_size == 0 && clamped > 0) { | ||
| 1597 | 1599 | EmitDrain(); | |
| 1598 | 1600 | } | |
| 1599 | 1601 | } | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -0,0 +1,91 @@ | |||
| 1 | + // Flags: --experimental-quic --experimental-stream-iter --no-warnings | ||
| 2 | + | ||
| 3 | + // Test: bidirectional data transfer with varying chunk sizes. | ||
| 4 | + // This is a regression test for a stall caused by a mismatch between | ||
| 5 | + // writeSync (which rejects when chunk > writeDesiredSize) and | ||
| 6 | + // drainableProtocol (which returned null when writeDesiredSize > 0). | ||
| 7 | + // When chunks don't evenly fill the high water mark, writeDesiredSize | ||
| 8 | + // can be positive but smaller than the next chunk, causing the | ||
| 9 | + // while(!writeSync) { dp(); await } loop to spin without yielding. | ||
| 10 | + // See: https://github.com/nodejs/node/issues/63216 | ||
| 11 | + | ||
| 12 | + import { hasQuic, skip, mustCall } from '../common/index.mjs'; | ||
| 13 | + import assert from 'node:assert'; | ||
| 14 | + | ||
| 15 | + const { strictEqual } = assert; | ||
| 16 | + | ||
| 17 | + if (!hasQuic) { | ||
| 18 | + skip('QUIC is not enabled'); | ||
| 19 | + } | ||
| 20 | + | ||
| 21 | + const { listen, connect } = await import('../common/quic.mjs'); | ||
| 22 | + const { bytes, drainableProtocol: dp } = await import('stream/iter'); | ||
| 23 | + | ||
| 24 | + // Varying chunk sizes — the pattern of alternating large and small | ||
| 25 | + // chunks is effective at triggering the writeDesiredSize gap. | ||
| 26 | + const chunkSizes = [60000, 12, 50000, 1600, 20000, 30000, 0, 100]; | ||
| 27 | + const numChunks = chunkSizes.length; | ||
| 28 | + const byteLength = chunkSizes.reduce((a, b) => a + b, 0); | ||
| 29 | + | ||
| 30 | + // Build a deterministic payload so we can verify integrity. | ||
| 31 | + function buildChunk(index) { | ||
| 32 | + const chunk = new Uint8Array(chunkSizes[index]); | ||
| 33 | + const val = index & 0xff; | ||
| 34 | + for (let i = 0; i < chunkSizes[index]; i++) { | ||
| 35 | + chunk[i] = (val + i) & 0xff; | ||
| 36 | + } | ||
| 37 | + return chunk; | ||
| 38 | + } | ||
| 39 | + | ||
| 40 | + function checksum(data) { | ||
| 41 | + let sum = 0; | ||
| 42 | + for (let i = 0; i < data.byteLength; i++) { | ||
| 43 | + sum = (sum + data[i]) | 0; | ||
| 44 | + } | ||
| 45 | + return sum; | ||
| 46 | + } | ||
| 47 | + | ||
| 48 | + // Compute expected checksum. | ||
| 49 | + let expectedChecksum = 0; | ||
| 50 | + for (let i = 0; i < numChunks; i++) { | ||
| 51 | + const chunk = buildChunk(i); | ||
| 52 | + expectedChecksum = (expectedChecksum + checksum(chunk)) | 0; | ||
| 53 | + } | ||
| 54 | + | ||
| 55 | + const done = Promise.withResolvers(); | ||
| 56 | + | ||
| 57 | + const serverEndpoint = await listen(mustCall((serverSession) => { | ||
| 58 | + serverSession.onstream = mustCall(async (stream) => { | ||
| 59 | + const received = await bytes(stream); | ||
| 60 | + strictEqual(received.byteLength, byteLength); | ||
| 61 | + strictEqual(checksum(received), expectedChecksum); | ||
| 62 | + | ||
| 63 | + stream.writer.endSync(); | ||
| 64 | + await stream.closed; | ||
| 65 | + serverSession.close(); | ||
| 66 | + done.resolve(); | ||
| 67 | + }); | ||
| 68 | + })); | ||
| 69 | + | ||
| 70 | + const clientSession = await connect(serverEndpoint.address); | ||
| 71 | + await clientSession.opened; | ||
| 72 | + | ||
| 73 | + const stream = await clientSession.createBidirectionalStream(); | ||
| 74 | + const w = stream.writer; | ||
| 75 | + | ||
| 76 | + // Write chunks, respecting backpressure via drainableProtocol. | ||
| 77 | + for (let i = 0; i < numChunks; i++) { | ||
| 78 | + const chunk = buildChunk(i); | ||
| 79 | + while (!w.writeSync(chunk)) { | ||
| 80 | + // Flow controlled — wait for drain before retrying. | ||
| 81 | + const drainable = w[dp](); | ||
| 82 | + if (drainable) await drainable; | ||
| 83 | + } | ||
| 84 | + } | ||
| 85 | + | ||
| 86 | + const totalWritten = w.endSync(); | ||
| 87 | + strictEqual(totalWritten, byteLength); | ||
| 88 | + | ||
| 89 | + await Promise.all([stream.closed, done.promise]); | ||
| 90 | + await clientSession.close(); | ||
| 91 | + await serverEndpoint.close(); | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments