diff --git a/doc/api/quic.md b/doc/api/quic.md index 5bdac8189fac..b163ffc1fd5e 100644 --- a/doc/api/quic.md +++ b/doc/api/quic.md @@ -1356,6 +1356,11 @@ added: v23.8.0 will buffer before `writeSync()` returns `false`. When the buffered data exceeds this limit, the caller should wait for drain before writing more. **Default:** `65536` (64 KB). + * `waitUntilAvailable` {boolean} When true the promise will wait until flow + control will allow to open the stream. If set to false, the function + will return a rejected promise, if flow control will not allow to + open the stream immediately. + **Default:** `true` * `onheaders` {Function} Callback for received initial response headers. Called with `(headers)`. * `ontrailers` {Function} Callback for received trailing headers. @@ -1397,6 +1402,11 @@ added: v23.8.0 will buffer before `writeSync()` returns `false`. When the buffered data exceeds this limit, the caller should wait for drain before writing more. **Default:** `65536` (64 KB). + * `waitUntilAvailable` {boolean} When true the promise will wait until flow + control will allow to open the stream. If set to false, the function + will return a rejected promise, if flow control will not allow + to open the stream immediately. + **Default:** `false` * `onheaders` {Function} Callback for received initial response headers. Called with `(headers)`. * `ontrailers` {Function} Callback for received trailing headers. diff --git a/lib/internal/quic/quic.js b/lib/internal/quic/quic.js index 4922ce562751..c5dbb7484c08 100644 --- a/lib/internal/quic/quic.js +++ b/lib/internal/quic/quic.js @@ -3509,6 +3509,7 @@ class QuicSession { incremental = false, budget = kDefaultBudget, headers, + waitUntilAvailable = true, onheaders, ontrailers, oninfo, @@ -3520,7 +3521,7 @@ class QuicSession { const validatedBody = validateBody(body); - const handle = this.#handle.openStream(direction, validatedBody); + const handle = this.#handle.openStream(direction, waitUntilAvailable, validatedBody); if (handle === undefined) { throw new ERR_QUIC_OPEN_STREAM_FAILED(); } diff --git a/src/quic/session.cc b/src/quic/session.cc index 90d59487a2a7..4985efd4cf5b 100644 --- a/src/quic/session.cc +++ b/src/quic/session.cc @@ -1168,17 +1168,26 @@ struct Session::Impl final : public MemoryRetainer { } DCHECK(args[0]->IsUint32()); + DCHECK(args[1]->IsBoolean()); + + auto direction = FromV8Value(args[0]); + if (!args[1].As()->Value()) { + // This is waitUntilAvailable + if (!session->CanImmediatelyOpenStream(direction)) { + return THROW_ERR_INVALID_STATE( + env, "No new stream available within flow control"); + } + } // GetDataQueueFromSource handles type validation. std::shared_ptr data_source; - if (!Stream::GetDataQueueFromSource(env, args[1]).To(&data_source)) + if (!Stream::GetDataQueueFromSource(env, args[2]).To(&data_source)) [[unlikely]] { return THROW_ERR_INVALID_ARG_VALUE(env, "Invalid data source"); } session->impl_->handshake_deferred_ = false; SendPendingDataScope send_scope(session); - auto direction = FromV8Value(args[0]); Local stream; if (session->OpenStream(direction, std::move(data_source)).ToLocal(&stream)) [[likely]] { @@ -3203,6 +3212,14 @@ BaseObjectPtr Session::CreateStream( return {}; } +bool Session::CanImmediatelyOpenStream(Direction direction) { + if (direction == Direction::BIDIRECTIONAL) { + return max_local_streams_bidi() > 0; + } else { + return max_local_streams_uni() > 0; + } +} + MaybeLocal Session::OpenStream(Direction direction, std::shared_ptr data_source) { // If can_create_streams() returns false, we are not able to open a stream @@ -3508,13 +3525,12 @@ void Session::SetApplicationError(error_code app_error_code) { uint64_t Session::max_local_streams_uni() const { DCHECK(!is_destroyed()); - return ngtcp2_conn_get_streams_uni_left(*this); + return ngtcp2_conn_get_streams_uni_left2(*this); } uint64_t Session::max_local_streams_bidi() const { DCHECK(!is_destroyed()); - return ngtcp2_conn_get_local_transport_params(*this) - ->initial_max_streams_bidi; + return ngtcp2_conn_get_streams_bidi_left2(*this); } void Session::set_wrapped() { diff --git a/src/quic/session.h b/src/quic/session.h index 9834aa7ec130..cc67f0ce3a67 100644 --- a/src/quic/session.h +++ b/src/quic/session.h @@ -535,6 +535,9 @@ class Session final : public AsyncWrap, private SessionTicket::AppData::Source { size_t max_packet_size() const; void set_priority_supported(bool on = true); + // Check whether flow control permits opening another stream + bool CanImmediatelyOpenStream(Direction direction); + // Open a new locally-initialized stream with the specified directionality. // If the session is not yet in a state where the stream can be opened -- // such as when the handshake is not yet sufficiently far along and ORTT diff --git a/test/parallel/test-quic-stream-limits-pending.mjs b/test/parallel/test-quic-stream-limits-pending.mjs index 2d81b0d79096..76220df8ad5b 100644 --- a/test/parallel/test-quic-stream-limits-pending.mjs +++ b/test/parallel/test-quic-stream-limits-pending.mjs @@ -15,9 +15,11 @@ if (!hasQuic) { const { listen, connect } = await import('../common/quic.mjs'); const { bytes } = await import('stream/iter'); +const { setTimeout: sleep } = await import('timers/promises'); const encoder = new TextEncoder(); const allDone = Promise.withResolvers(); +const twoDone = Promise.withResolvers(); let serverStreamCount = 0; // Server allows only 1 bidi stream at a time. @@ -26,11 +28,17 @@ const serverEndpoint = await listen(mustCall((serverSession) => { await bytes(stream); stream.writer.endSync(); await stream.closed; - if (++serverStreamCount === 2) { - serverSession.close(); + ++serverStreamCount; + if (serverStreamCount === 2) { + twoDone.resolve(); + } + if (serverStreamCount === 3) { allDone.resolve(); } - }, 2); + if (serverStreamCount === 4) { + serverSession.close(); + } + }, 3); }), { transportParams: { initialMaxStreamsBidi: 1 }, }); @@ -50,17 +58,34 @@ s1.opened.then(() => { opened++; }); -// Second stream is created but queued as pending because the -// server only allows 1 concurrent bidi stream. -const s2 = await clientSession.createBidirectionalStream({ - body: encoder.encode('stream 2'), +let s2; +await assert.rejects( + async () => { + // Second stream should not open, but throw. + s2 = await clientSession.createBidirectionalStream({ + body: encoder.encode('stream 2a'), + waitUntilAvailable: false, + }); + // eslint-disable-next-line node-core/must-call-assert + s2.opened.then(() => { + opened++; + }); + }, + { + name: 'Error', + message: 'No new stream available within flow control', + }, +); +// Ok try again a second second stream, that patiently waits +s2 = await clientSession.createBidirectionalStream({ + body: encoder.encode('stream 2b') }); - // eslint-disable-next-line node-core/must-call-assert s2.opened.then(() => { opened++; }); + // Third stream is created but queued as pending because the // server only allows 1 concurrent bidi stream. const s3 = await clientSession.createBidirectionalStream({ @@ -68,9 +93,10 @@ const s3 = await clientSession.createBidirectionalStream({ }); -// s2 should be pending until s1 closes and the server grants +// s2 and s3 should be pending until s1 closes and the server grants // more stream credits. assert.strictEqual(s2.pending, true); +assert.strictEqual(s3.pending, true); assert.strictEqual(opened, 1); // Drain and close the first stream. @@ -82,15 +108,20 @@ s3.destroy(err); await Promise.all([assert.rejects(s3.opened, err), assert.rejects(s3.closed, err)]); - -// After s1 closes, the server sends MAX_STREAMS which opens s2. +// After s1 closes, the server sends MAX_STREAMS which opens s3. // Wait for the server to receive both streams. -await allDone.promise; -assert.strictEqual(opened, 2); - +await twoDone.promise; // s2 should no longer be pending. for await (const _ of s2) { /* drain */ } // eslint-disable-line no-unused-vars await s2.closed; +await sleep(10); // We wait a bit, as we do not have a callback exposed to js +// fourth stream should open immediately and not throw +const s4 = await clientSession.createBidirectionalStream({ + body: encoder.encode('stream 4'), + waitUntilAvailable: false +}); +await Promise.all([s4.closed, allDone.promise]); + await clientSession.close(); await serverEndpoint.close(); diff --git a/test/parallel/test-quic-stream-limits-uni.mjs b/test/parallel/test-quic-stream-limits-uni.mjs index f3686ca45e3c..1c7fbd974386 100644 --- a/test/parallel/test-quic-stream-limits-uni.mjs +++ b/test/parallel/test-quic-stream-limits-uni.mjs @@ -39,6 +39,7 @@ let opened = 0; // First uni stream opens immediately. const s1 = await clientSession.createUnidirectionalStream({ body: encoder.encode('uni 1'), + waitUntilAvailable: false, }); // eslint-disable-next-line node-core/must-call-assert