Skip to content

Commit cb2e8dd

Browse files
committed
quic: update writedesiredsize on maxdataframe
This commit will react on a callback from ngtcp2, when a maxdata frame arrives and update the scheduled streams accordingly. A PR to get this callback is pending. Fixes: #64835 Signed-off-by: Marten Richter <marten.richter@freenet.de>
1 parent 4551732 commit cb2e8dd

5 files changed

Lines changed: 118 additions & 1 deletion

File tree

‎src/quic/application.cc‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -419,7 +419,7 @@ class DefaultApplication final : public Session::Application {
419419
void BlockStream(stream_id id) override {
420420
if (auto stream = session().FindStream(id)) [[likely]] {
421421
// Remove the stream from the send queue. It will be re-scheduled
422-
// via ExtendMaxStreamData when the peer grants more flow control.
422+
// via ExtendMax(Stream)Data when the peer grants more flow control.
423423
// Without this, SendPendingData would repeatedly pop and retry
424424
// the same blocked stream in an infinite loop.
425425
stream->Unschedule();
@@ -434,6 +434,15 @@ class DefaultApplication final : public Session::Application {
434434
stream->Schedule(&stream_queue_);
435435
}
436436

437+
void ExtendMaxData(uint64_t max_data) override {
438+
// The peer granted more flow control for session. Re-schedule
439+
// all streams so SendPendingData will resume writing.
440+
for (auto& [id, stream] : session().streams()) {
441+
stream->Schedule(&stream_queue_);
442+
}
443+
}
444+
445+
437446
bool StreamCommit(Session::StreamData* stream_data, size_t datalen) override {
438447
DCHECK_NOT_NULL(stream_data);
439448
CHECK(stream_data->stream);

‎src/quic/application.h‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -149,6 +149,14 @@ class Session::Application : public MemoryRetainer {
149149
// By default do nothing.
150150
}
151151

152+
// Called when the Session determines that the flow control window for the
153+
// session has been expanded. Not all Application types will require
154+
// this notification so the default is to do nothing.
155+
virtual void ExtendMaxData(uint64_t max_data) {
156+
Debug(session_, "Application extending max data");
157+
// By default do nothing.
158+
}
159+
152160
// Different Applications may wish to set some application data in the
153161
// session ticket (e.g. http/3 would set server settings in the application
154162
// data). The first byte written MUST be the Application::Type enum value.

‎src/quic/http3.cc‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -378,6 +378,16 @@ class Http3ApplicationImpl final : public Session::Application {
378378
nghttp3_conn_unblock_stream(*this, stream->id());
379379
}
380380

381+
void ExtendMaxData(uint64_t max_data) override {
382+
Debug(&session(),
383+
"HTTP/3 application extending max data to %" PRIu64,
384+
max_data);
385+
for (auto& [id, stream] : session().streams()) {
386+
stream->UpdateWriteDesiredSize(); // the stream might be blocked on js side
387+
// is unblock stream also required?
388+
}
389+
}
390+
381391
void CollectSessionTicketAppData(
382392
SessionTicket::AppData* app_data) const override {
383393
uint8_t buf[kSessionTicketAppDataSize];

‎src/quic/session.cc‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1358,6 +1358,14 @@ struct Session::Impl final : public MemoryRetainer {
13581358
return NGTCP2_SUCCESS;
13591359
}
13601360

1361+
static int on_extend_max_data(ngtcp2_conn* conn,
1362+
uint64_t max_data,
1363+
void* user_data) {
1364+
NGTCP2_CALLBACK_SCOPE(session)
1365+
session->application().ExtendMaxData(max_data);
1366+
return NGTCP2_SUCCESS;
1367+
}
1368+
13611369
static int on_get_new_cid(ngtcp2_conn* conn,
13621370
ngtcp2_cid* cid,
13631371
ngtcp2_stateless_reset_token* token,
@@ -1700,6 +1708,9 @@ struct Session::Impl final : public MemoryRetainer {
17001708
on_receive_stream_stop_sending,
17011709
#ifdef NGTCP2_CALLBACKS_V5
17021710
nullptr,
1711+
#ifdef NGTCP2_CALLBACKS_V6
1712+
on_extend_max_data,
1713+
#endif
17031714
#endif // NGTCP2_CALLBACKS_V5
17041715
#endif // NGTCP2_CALLBACKS_V4
17051716
};
@@ -1754,6 +1765,9 @@ struct Session::Impl final : public MemoryRetainer {
17541765
on_receive_stream_stop_sending,
17551766
#ifdef NGTCP2_CALLBACKS_V5
17561767
nullptr,
1768+
#ifdef NGTCP2_CALLBACKS_V6
1769+
on_extend_max_data,
1770+
#endif
17571771
#endif // NGTCP2_CALLBACKS_V5
17581772
#endif // NGTCP2_CALLBACKS_V4
17591773
};
Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,76 @@
1+
// Flags: --experimental-quic --experimental-stream-iter --no-warnings
2+
3+
// Test: Quic maxdata updates on http/3
4+
// Client sends a body that precisely fills the session window size,
5+
// and verifies that it is data transfer is not stalled.
6+
7+
import { hasQuic, skip } from '../common/index.mjs';
8+
import { readFile } from 'node:fs/promises';
9+
import { setTimeout as sleep } from 'node:timers/promises';
10+
11+
if (!hasQuic) {
12+
skip('QUIC is not enabled');
13+
}
14+
const { listen, connect } = await import('node:quic');
15+
const { createPrivateKey } = await import('node:crypto');
16+
const { drainableProtocol } = await import('stream/iter');
17+
18+
const keys = 'test/fixtures/keys';
19+
const key = createPrivateKey(await readFile(`${keys}/agent1-key.pem`));
20+
const cert = await readFile(`${keys}/agent1-cert.pem`);
21+
22+
const WINDOW = 4096;
23+
// Fills the window exactly:
24+
// considers all framing including some initial session capsules
25+
const BODY = WINDOW - 38;
26+
27+
let letServerRead;
28+
const serverMayRead = new Promise((resolve) => { letServerRead = resolve; });
29+
30+
const endpoint = await listen((session) => {
31+
session.onstream = async (stream) => {
32+
await serverMayRead;
33+
// eslint-disable-next-line no-unused-vars
34+
for await (const _ of stream) { /* reading extends the window */ }
35+
};
36+
}, {
37+
sni: { '*': { keys: [key], certs: [cert] } },
38+
transportParams: {
39+
initialMaxStreamDataBidiRemote: 1024 * 1024, // Make sure maxstreamdata does not block
40+
initialMaxData: WINDOW,
41+
},
42+
onheaders() { this.sendHeaders({ ':status': '200' }); },
43+
});
44+
45+
const session = await connect(endpoint.address, {
46+
servername: 'localhost',
47+
verifyPeer: 'manual',
48+
});
49+
await session.opened;
50+
51+
// Budget well above the window, so the window is what stops the writer.
52+
const stream = await session.createBidirectionalStream({ budget: 1024 * 1024 });
53+
stream.sendHeaders({
54+
':method': 'POST',
55+
':path': '/',
56+
':scheme': 'https',
57+
':authority': 'localhost',
58+
}, { terminal: false });
59+
60+
const writer = stream.writer;
61+
writer.writeSync(new Uint8Array(BODY));
62+
63+
// Long enough for every byte to be acked. The peer acks as data arrives,
64+
// whether or not its application has read any of it, so by now the window is
65+
// exhausted, the send buffer is empty, and no further ACK can arrive.
66+
await sleep(500);
67+
68+
const watchdog = setTimeout(() => {
69+
console.error('STALLED: no drain after MAX_STREAM_DATA');
70+
process.exit(1);
71+
}, 5000);
72+
letServerRead(); // Extend the window, with no ack attached
73+
await writer[drainableProtocol]();
74+
75+
clearTimeout(watchdog);
76+
process.exit(0);

0 commit comments

Comments
 (0)