Skip to content
10 changes: 10 additions & 0 deletions doc/api/quic.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down
3 changes: 2 additions & 1 deletion lib/internal/quic/quic.js
Original file line number Diff line number Diff line change
Expand Up @@ -3509,6 +3509,7 @@ class QuicSession {
incremental = false,
budget = kDefaultBudget,
headers,
waitUntilAvailable = true,
onheaders,
ontrailers,
oninfo,
Expand All @@ -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();
}
Expand Down
26 changes: 21 additions & 5 deletions src/quic/session.cc
Original file line number Diff line number Diff line change
Expand Up @@ -1168,17 +1168,26 @@ struct Session::Impl final : public MemoryRetainer {
}

DCHECK(args[0]->IsUint32());
DCHECK(args[1]->IsBoolean());

auto direction = FromV8Value<Direction>(args[0]);
if (!args[1].As<v8::Boolean>()->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<DataQueue> 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<Direction>(args[0]);
Local<Object> stream;
if (session->OpenStream(direction, std::move(data_source)).ToLocal(&stream))
[[likely]] {
Expand Down Expand Up @@ -3203,6 +3212,14 @@ BaseObjectPtr<Stream> 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<Object> Session::OpenStream(Direction direction,
std::shared_ptr<DataQueue> data_source) {
// If can_create_streams() returns false, we are not able to open a stream
Expand Down Expand Up @@ -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() {
Expand Down
3 changes: 3 additions & 0 deletions src/quic/session.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
59 changes: 45 additions & 14 deletions test/parallel/test-quic-stream-limits-pending.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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 },
});
Expand All @@ -50,27 +58,45 @@ 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({
body: encoder.encode('stream 3'),
});


// 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.
Expand All @@ -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();
1 change: 1 addition & 0 deletions test/parallel/test-quic-stream-limits-uni.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading