Commit e2e2829

Browse files
pimterryaduh95
authored andcommitted
quic: extract transport logic from Application to Session
One small fix notably included: - Check is_destroyed() after StreamCommit, since it calls JS callbacks which could destroy the session. Signed-off-by: Tim Perry <pimterry@gmail.com> PR-URL: #64127 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 7ac97fc commit e2e2829

9 files changed

Lines changed: 608 additions & 562 deletions

File tree

β€Žsrc/quic/README.mdβ€Ž

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ data channels that carry application data.
7777

7878
Every entry point that may generate outbound data creates a
7979
`SendPendingDataScope`. Scopes nest β€” an internal depth counter ensures
80-
`Application::SendPendingData()` is called exactly once, when the outermost
80+
`Session::SendPendingData()` is called exactly once, when the outermost
8181
scope exits:
8282

8383
```cpp
@@ -218,13 +218,13 @@ Session::Receive()
218218

219219
```text
220220
SendPendingDataScope::~SendPendingDataScope()
221-
β†’ Application::SendPendingData()
221+
β†’ Session::SendPendingData()
222222
Loop (up to max_packet_count):
223-
β”œβ”€β”€ GetStreamData() // pull data from next stream
224-
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225-
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
223+
β”œβ”€β”€ application().GetStreamData() // pull data from next stream
224+
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225+
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
226226
β”‚ encrypts, frames, paces
227-
β”œβ”€β”€ if ndatalen > 0: StreamCommit()
227+
β”œβ”€β”€ if ndatalen > 0: application().StreamCommit()
228228
β”‚ stream->Commit(datalen, fin)
229229
β”œβ”€β”€ if nwrite > 0: Send() // uv_udp_send()
230230
β”œβ”€β”€ if WRITE_MORE: continue // room for more in this packet

β€Žsrc/quic/application.ccβ€Ž

Lines changed: 6 additions & 449 deletions
Large diffs are not rendered by default.

β€Žsrc/quic/application.hβ€Ž

Lines changed: 0 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -205,10 +205,6 @@ class Session::Application : public MemoryRetainer {
205205
returnfalse;
206206
}
207207

208-
// Signals to the Application that it should serialize and transmit any
209-
// pending session and stream packets it has accumulated.
210-
voidSendPendingData();
211-
212208
// Returns true if the application protocol supports sending and
213209
// receiving headers on streams (e.g. HTTP/3). Applications that
214210
// do not support headers should return false (the default).
@@ -243,10 +239,6 @@ class Session::Application : public MemoryRetainer {
243239
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
244240
}
245241

246-
// The StreamData struct is used by the application to pass pending stream
247-
// data to the session for transmission.
248-
structStreamData;
249-
250242
virtualintGetStreamData(StreamData* data) = 0;
251243
virtualboolStreamCommit(StreamData* data, size_t datalen) = 0;
252244

@@ -262,57 +254,9 @@ class Session::Application : public MemoryRetainer {
262254
}
263255

264256
private:
265-
Packet::Ptr CreateStreamDataPacket();
266-
267-
// Tries to pack a pending datagram into the current packet buffer.
268-
// If < 0 is returned, either NGTCP2_ERR_WRITE_MORE or a fatal error is
269-
// returned; the caller must check. If > 0 is returned, the packet is done
270-
// and the value is the size of the finalized packet. If 0 is returned,
271-
// the datagram is either congestion limited or was abandoned
272-
ssize_tTryWritePendingDatagram(PathStorage* path,
273-
uint8_t* dest,
274-
size_t destlen,
275-
uint64_t ts);
276-
277-
// Write the given stream_data into the buffer. The PacketInfo out-param
278-
// is populated by ngtcp2 with per-packet metadata (e.g., ECN codepoint)
279-
// that should be applied when sending the packet.
280-
ssize_tWriteVStream(PathStorage* path,
281-
PacketInfo* pi,
282-
uint8_t* buf,
283-
ssize_t* ndatalen,
284-
size_t max_packet_size,
285-
const StreamData& stream_data,
286-
uint64_t ts);
287-
288257
Session* session_ = nullptr;
289258
};
290259

291-
structSession::Application::StreamData final {
292-
// The actual number of vectors in the struct, up to kMaxVectorCount.
293-
size_t count = 0;
294-
// The stream identifier. If this is a negative value then no stream is
295-
// identified.
296-
stream_id id = -1;
297-
int fin = 0;
298-
ngtcp2_vec data[kMaxVectorCount]{};
299-
BaseObjectPtr<Stream> stream;
300-
301-
static_assert(sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
302-
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
303-
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
304-
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
305-
"ngtcp2_vec and nghttp3_vec must have identical layout");
306-
inlineoperator nghttp3_vec*() {
307-
returnreinterpret_cast<nghttp3_vec*>(data);
308-
}
309-
310-
inlineoperatorconst ngtcp2_vec*() const { return data; }
311-
inlineoperator ngtcp2_vec*() { return data; }
312-
313-
std::string ToString() const;
314-
};
315-
316260
// Create a DefaultApplication for the given session.
317261
std::unique_ptr<Session::Application> CreateDefaultApplication(
318262
Session* session, const Session::Application_Options& options);

β€Žsrc/quic/http3.ccβ€Ž

Lines changed: 40 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -209,24 +209,25 @@ class Http3ApplicationImpl final : public Session::Application {
209209
started_ = true;
210210
Debug(&session(), "Starting HTTP/3 application.");
211211

212-
auto params = ngtcp2_conn_get_remote_transport_params(session());
213-
if (params == nullptr) [[unlikely]] {
212+
constauto params = session().remote_transport_params();
213+
if (!params) [[unlikely]] {
214214
// The params are not available yet. Cannot start.
215215
Debug(&session(),
216216
"Cannot start HTTP/3 application yet. No remote transport params");
217217
returnfalse;
218218
}
219219

220-
if (params->initial_max_streams_uni < 3) {
220+
if (params.initial_max_streams_uni() < 3) {
221221
// HTTP3 requires 3 unidirectional control streams to be opened in each
222222
// direction in additional to the bidirectional streams that are used to
223223
// actually carry request and response payload back and forth.
224224
// See:
225225
// https://nghttp2.org/nghttp3/programmers-guide.html#binding-control-streams
226226
Debug(&session(),
227227
"Cannot start HTTP/3 application. Initial max "
228-
"unidirectional streams [%zu] is too low. Must be at least 3",
229-
params->initial_max_streams_uni);
228+
"unidirectional streams [%" PRIu64
229+
"] is too low. Must be at least 3",
230+
params.initial_max_streams_uni());
230231
returnfalse;
231232
}
232233

@@ -235,17 +236,14 @@ class Http3ApplicationImpl final : public Session::Application {
235236
// of requests that the client can actually created.
236237
if (session().is_server()) {
237238
nghttp3_conn_set_max_client_streams_bidi(
238-
*this, params->initial_max_streams_bidi);
239+
*this, params.initial_max_streams_bidi());
239240
}
240241

241242
Debug(&session(), "Creating and binding HTTP/3 control streams");
242243
bool ret =
243-
ngtcp2_conn_open_uni_stream(session(), &control_stream_id_, nullptr) ==
244-
0 &&
245-
ngtcp2_conn_open_uni_stream(
246-
session(), &qpack_enc_stream_id_, nullptr) == 0 &&
247-
ngtcp2_conn_open_uni_stream(
248-
session(), &qpack_dec_stream_id_, nullptr) == 0 &&
244+
session().OpenUnidirectionalStream(&control_stream_id_) &&
245+
session().OpenUnidirectionalStream(&qpack_enc_stream_id_) &&
246+
session().OpenUnidirectionalStream(&qpack_dec_stream_id_) &&
249247
nghttp3_conn_bind_control_stream(*this, control_stream_id_) == 0 &&
250248
nghttp3_conn_bind_qpack_streams(
251249
*this, qpack_enc_stream_id_, qpack_dec_stream_id_) == 0;
@@ -306,8 +304,7 @@ class Http3ApplicationImpl final : public Session::Application {
306304
Debug(&session(),
307305
"Extending stream and connection offset by %zd bytes",
308306
nread);
309-
session().ExtendStreamOffset(id, nread);
310-
session().ExtendOffset(nread);
307+
session().Consume(id, nread);
311308
}
312309

313310
// If this data arrived as 0-RTT, mark the stream. We set it after
@@ -365,24 +362,11 @@ class Http3ApplicationImpl final : public Session::Application {
365362
case EndpointLabel::LOCAL:
366363
return;
367364
case EndpointLabel::REMOTE: {
368-
switch (direction) {
369-
case Direction::BIDIRECTIONAL: {
370-
Debug(&session(),
371-
"HTTP/3 application extending max bidi streams by %" PRIu64,
372-
max_streams);
373-
ngtcp2_conn_extend_max_streams_bidi(
374-
session(), static_cast<size_t>(max_streams));
375-
break;
376-
}
377-
case Direction::UNIDIRECTIONAL: {
378-
Debug(&session(),
379-
"HTTP/3 application extending max uni streams by %" PRIu64,
380-
max_streams);
381-
ngtcp2_conn_extend_max_streams_uni(
382-
session(), static_cast<size_t>(max_streams));
383-
break;
384-
}
385-
}
365+
Debug(&session(),
366+
"HTTP/3 application extending max %s streams by %" PRIu64,
367+
direction == Direction::BIDIRECTIONAL ? "bidi" : "uni",
368+
max_streams);
369+
session().ExtendMaxStreams(direction, max_streams);
386370
}
387371
}
388372
}
@@ -530,8 +514,7 @@ class Http3ApplicationImpl final : public Session::Application {
530514
return;
531515
}
532516

533-
session().SetLastError(
534-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
517+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
535518
session().Close();
536519
}
537520

@@ -548,8 +531,7 @@ class Http3ApplicationImpl final : public Session::Application {
548531
return;
549532
}
550533

551-
session().SetLastError(
552-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
534+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
553535
session().Close();
554536
}
555537

@@ -687,17 +669,30 @@ class Http3ApplicationImpl final : public Session::Application {
687669
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
688670
}
689671

690-
intGetStreamData(StreamData* data) override {
672+
intGetStreamData(Session::StreamData* data) override {
673+
static_assert(
674+
sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
675+
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
676+
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
677+
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
678+
"ngtcp2_vec and nghttp3_vec must have identical layout");
691679
data->count = kMaxVectorCount;
692680
ssize_t ret = 0;
693681
Debug(&session(), "HTTP/3 application getting stream data");
694682
if (conn_ && session().max_data_left()) {
695-
ret = nghttp3_conn_writev_stream(
696-
*this, &data->id, &data->fin, *data, data->count);
683+
// nghttp3 reports fin through an int out-param; bridge it to the bool.
684+
int fin = 0;
685+
ret =
686+
nghttp3_conn_writev_stream(*this,
687+
&data->id,
688+
&fin,
689+
reinterpret_cast<nghttp3_vec*>(data->data),
690+
data->count);
697691
// A negative return value indicates an error.
698692
if (ret < 0) {
699693
returnstatic_cast<int>(ret);
700694
}
695+
data->fin = fin != 0;
701696

702697
data->count = static_cast<size_t>(ret);
703698
if (data->id >= 0 && data->id != control_stream_id_ &&
@@ -710,7 +705,7 @@ class Http3ApplicationImpl final : public Session::Application {
710705
return0;
711706
}
712707

713-
boolStreamCommit(StreamData* data, size_t datalen) override {
708+
boolStreamCommit(Session::StreamData* data, size_t datalen) override {
714709
Debug(&session(),
715710
"HTTP/3 application committing stream %" PRIi64 " data %zu",
716711
data->id,
@@ -720,8 +715,7 @@ class Http3ApplicationImpl final : public Session::Application {
720715
// nghttp3 tracks its own offset via add_write_offset.
721716
int err = nghttp3_conn_add_write_offset(*this, data->id, datalen);
722717
if (err != 0) {
723-
session().SetLastError(QuicError::ForApplication(
724-
nghttp3_err_infer_quic_app_error_code(err)));
718+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(err));
725719
returnfalse;
726720
}
727721
// Raw application bytes are committed to the stream's outbound
@@ -1212,10 +1206,10 @@ class Http3ApplicationImpl final : public Session::Application {
12121206
void* conn_user_data,
12131207
void* stream_user_data) {
12141208
NGHTTP3_CALLBACK_SCOPE(app);
1215-
auto& session = app.session();
1216-
Debug(&session, "HTTP/3 application deferred consume %zu bytes", consumed);
1217-
session.ExtendStreamOffset(id, consumed);
1218-
session.ExtendOffset(consumed);
1209+
Debug(&app.session(),
1210+
"HTTP/3 application deferred consume %zu bytes",
1211+
consumed);
1212+
app.session().Consume(id, consumed);
12191213
returnNGTCP2_SUCCESS;
12201214
}
12211215

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Commit e2e2829

Browse files
pimterryaduh95
authored andcommitted
quic: extract transport logic from Application to Session
One small fix notably included: - Check is_destroyed() after StreamCommit, since it calls JS callbacks which could destroy the session. Signed-off-by: Tim Perry <pimterry@gmail.com> PR-URL: #64127 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 7ac97fc commit e2e2829

9 files changed

Lines changed: 608 additions & 562 deletions

File tree

β€Žsrc/quic/README.mdβ€Ž

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ data channels that carry application data.
7777

7878
Every entry point that may generate outbound data creates a
7979
`SendPendingDataScope`. Scopes nest β€” an internal depth counter ensures
80-
`Application::SendPendingData()` is called exactly once, when the outermost
80+
`Session::SendPendingData()` is called exactly once, when the outermost
8181
scope exits:
8282

8383
```cpp
@@ -218,13 +218,13 @@ Session::Receive()
218218

219219
```text
220220
SendPendingDataScope::~SendPendingDataScope()
221-
β†’ Application::SendPendingData()
221+
β†’ Session::SendPendingData()
222222
Loop (up to max_packet_count):
223-
β”œβ”€β”€ GetStreamData() // pull data from next stream
224-
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225-
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
223+
β”œβ”€β”€ application().GetStreamData() // pull data from next stream
224+
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225+
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
226226
β”‚ encrypts, frames, paces
227-
β”œβ”€β”€ if ndatalen > 0: StreamCommit()
227+
β”œβ”€β”€ if ndatalen > 0: application().StreamCommit()
228228
β”‚ stream->Commit(datalen, fin)
229229
β”œβ”€β”€ if nwrite > 0: Send() // uv_udp_send()
230230
β”œβ”€β”€ if WRITE_MORE: continue // room for more in this packet

β€Žsrc/quic/application.ccβ€Ž

Lines changed: 6 additions & 449 deletions
Large diffs are not rendered by default.

β€Žsrc/quic/application.hβ€Ž

Lines changed: 0 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -205,10 +205,6 @@ class Session::Application : public MemoryRetainer {
205205
returnfalse;
206206
}
207207

208-
// Signals to the Application that it should serialize and transmit any
209-
// pending session and stream packets it has accumulated.
210-
voidSendPendingData();
211-
212208
// Returns true if the application protocol supports sending and
213209
// receiving headers on streams (e.g. HTTP/3). Applications that
214210
// do not support headers should return false (the default).
@@ -243,10 +239,6 @@ class Session::Application : public MemoryRetainer {
243239
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
244240
}
245241

246-
// The StreamData struct is used by the application to pass pending stream
247-
// data to the session for transmission.
248-
structStreamData;
249-
250242
virtualintGetStreamData(StreamData* data) = 0;
251243
virtualboolStreamCommit(StreamData* data, size_t datalen) = 0;
252244

@@ -262,57 +254,9 @@ class Session::Application : public MemoryRetainer {
262254
}
263255

264256
private:
265-
Packet::Ptr CreateStreamDataPacket();
266-
267-
// Tries to pack a pending datagram into the current packet buffer.
268-
// If < 0 is returned, either NGTCP2_ERR_WRITE_MORE or a fatal error is
269-
// returned; the caller must check. If > 0 is returned, the packet is done
270-
// and the value is the size of the finalized packet. If 0 is returned,
271-
// the datagram is either congestion limited or was abandoned
272-
ssize_tTryWritePendingDatagram(PathStorage* path,
273-
uint8_t* dest,
274-
size_t destlen,
275-
uint64_t ts);
276-
277-
// Write the given stream_data into the buffer. The PacketInfo out-param
278-
// is populated by ngtcp2 with per-packet metadata (e.g., ECN codepoint)
279-
// that should be applied when sending the packet.
280-
ssize_tWriteVStream(PathStorage* path,
281-
PacketInfo* pi,
282-
uint8_t* buf,
283-
ssize_t* ndatalen,
284-
size_t max_packet_size,
285-
const StreamData& stream_data,
286-
uint64_t ts);
287-
288257
Session* session_ = nullptr;
289258
};
290259

291-
structSession::Application::StreamData final {
292-
// The actual number of vectors in the struct, up to kMaxVectorCount.
293-
size_t count = 0;
294-
// The stream identifier. If this is a negative value then no stream is
295-
// identified.
296-
stream_id id = -1;
297-
int fin = 0;
298-
ngtcp2_vec data[kMaxVectorCount]{};
299-
BaseObjectPtr<Stream> stream;
300-
301-
static_assert(sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
302-
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
303-
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
304-
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
305-
"ngtcp2_vec and nghttp3_vec must have identical layout");
306-
inlineoperator nghttp3_vec*() {
307-
returnreinterpret_cast<nghttp3_vec*>(data);
308-
}
309-
310-
inlineoperatorconst ngtcp2_vec*() const { return data; }
311-
inlineoperator ngtcp2_vec*() { return data; }
312-
313-
std::string ToString() const;
314-
};
315-
316260
// Create a DefaultApplication for the given session.
317261
std::unique_ptr<Session::Application> CreateDefaultApplication(
318262
Session* session, const Session::Application_Options& options);

β€Žsrc/quic/http3.ccβ€Ž

Lines changed: 40 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -209,24 +209,25 @@ class Http3ApplicationImpl final : public Session::Application {
209209
started_ = true;
210210
Debug(&session(), "Starting HTTP/3 application.");
211211

212-
auto params = ngtcp2_conn_get_remote_transport_params(session());
213-
if (params == nullptr) [[unlikely]] {
212+
constauto params = session().remote_transport_params();
213+
if (!params) [[unlikely]] {
214214
// The params are not available yet. Cannot start.
215215
Debug(&session(),
216216
"Cannot start HTTP/3 application yet. No remote transport params");
217217
returnfalse;
218218
}
219219

220-
if (params->initial_max_streams_uni < 3) {
220+
if (params.initial_max_streams_uni() < 3) {
221221
// HTTP3 requires 3 unidirectional control streams to be opened in each
222222
// direction in additional to the bidirectional streams that are used to
223223
// actually carry request and response payload back and forth.
224224
// See:
225225
// https://nghttp2.org/nghttp3/programmers-guide.html#binding-control-streams
226226
Debug(&session(),
227227
"Cannot start HTTP/3 application. Initial max "
228-
"unidirectional streams [%zu] is too low. Must be at least 3",
229-
params->initial_max_streams_uni);
228+
"unidirectional streams [%" PRIu64
229+
"] is too low. Must be at least 3",
230+
params.initial_max_streams_uni());
230231
returnfalse;
231232
}
232233

@@ -235,17 +236,14 @@ class Http3ApplicationImpl final : public Session::Application {
235236
// of requests that the client can actually created.
236237
if (session().is_server()) {
237238
nghttp3_conn_set_max_client_streams_bidi(
238-
*this, params->initial_max_streams_bidi);
239+
*this, params.initial_max_streams_bidi());
239240
}
240241

241242
Debug(&session(), "Creating and binding HTTP/3 control streams");
242243
bool ret =
243-
ngtcp2_conn_open_uni_stream(session(), &control_stream_id_, nullptr) ==
244-
0 &&
245-
ngtcp2_conn_open_uni_stream(
246-
session(), &qpack_enc_stream_id_, nullptr) == 0 &&
247-
ngtcp2_conn_open_uni_stream(
248-
session(), &qpack_dec_stream_id_, nullptr) == 0 &&
244+
session().OpenUnidirectionalStream(&control_stream_id_) &&
245+
session().OpenUnidirectionalStream(&qpack_enc_stream_id_) &&
246+
session().OpenUnidirectionalStream(&qpack_dec_stream_id_) &&
249247
nghttp3_conn_bind_control_stream(*this, control_stream_id_) == 0 &&
250248
nghttp3_conn_bind_qpack_streams(
251249
*this, qpack_enc_stream_id_, qpack_dec_stream_id_) == 0;
@@ -306,8 +304,7 @@ class Http3ApplicationImpl final : public Session::Application {
306304
Debug(&session(),
307305
"Extending stream and connection offset by %zd bytes",
308306
nread);
309-
session().ExtendStreamOffset(id, nread);
310-
session().ExtendOffset(nread);
307+
session().Consume(id, nread);
311308
}
312309

313310
// If this data arrived as 0-RTT, mark the stream. We set it after
@@ -365,24 +362,11 @@ class Http3ApplicationImpl final : public Session::Application {
365362
case EndpointLabel::LOCAL:
366363
return;
367364
case EndpointLabel::REMOTE: {
368-
switch (direction) {
369-
case Direction::BIDIRECTIONAL: {
370-
Debug(&session(),
371-
"HTTP/3 application extending max bidi streams by %" PRIu64,
372-
max_streams);
373-
ngtcp2_conn_extend_max_streams_bidi(
374-
session(), static_cast<size_t>(max_streams));
375-
break;
376-
}
377-
case Direction::UNIDIRECTIONAL: {
378-
Debug(&session(),
379-
"HTTP/3 application extending max uni streams by %" PRIu64,
380-
max_streams);
381-
ngtcp2_conn_extend_max_streams_uni(
382-
session(), static_cast<size_t>(max_streams));
383-
break;
384-
}
385-
}
365+
Debug(&session(),
366+
"HTTP/3 application extending max %s streams by %" PRIu64,
367+
direction == Direction::BIDIRECTIONAL ? "bidi" : "uni",
368+
max_streams);
369+
session().ExtendMaxStreams(direction, max_streams);
386370
}
387371
}
388372
}
@@ -530,8 +514,7 @@ class Http3ApplicationImpl final : public Session::Application {
530514
return;
531515
}
532516

533-
session().SetLastError(
534-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
517+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
535518
session().Close();
536519
}
537520

@@ -548,8 +531,7 @@ class Http3ApplicationImpl final : public Session::Application {
548531
return;
549532
}
550533

551-
session().SetLastError(
552-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
534+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
553535
session().Close();
554536
}
555537

@@ -687,17 +669,30 @@ class Http3ApplicationImpl final : public Session::Application {
687669
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
688670
}
689671

690-
intGetStreamData(StreamData* data) override {
672+
intGetStreamData(Session::StreamData* data) override {
673+
static_assert(
674+
sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
675+
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
676+
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
677+
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
678+
"ngtcp2_vec and nghttp3_vec must have identical layout");
691679
data->count = kMaxVectorCount;
692680
ssize_t ret = 0;
693681
Debug(&session(), "HTTP/3 application getting stream data");
694682
if (conn_ && session().max_data_left()) {
695-
ret = nghttp3_conn_writev_stream(
696-
*this, &data->id, &data->fin, *data, data->count);
683+
// nghttp3 reports fin through an int out-param; bridge it to the bool.
684+
int fin = 0;
685+
ret =
686+
nghttp3_conn_writev_stream(*this,
687+
&data->id,
688+
&fin,
689+
reinterpret_cast<nghttp3_vec*>(data->data),
690+
data->count);
697691
// A negative return value indicates an error.
698692
if (ret < 0) {
699693
returnstatic_cast<int>(ret);
700694
}
695+
data->fin = fin != 0;
701696

702697
data->count = static_cast<size_t>(ret);
703698
if (data->id >= 0 && data->id != control_stream_id_ &&
@@ -710,7 +705,7 @@ class Http3ApplicationImpl final : public Session::Application {
710705
return0;
711706
}
712707

713-
boolStreamCommit(StreamData* data, size_t datalen) override {
708+
boolStreamCommit(Session::StreamData* data, size_t datalen) override {
714709
Debug(&session(),
715710
"HTTP/3 application committing stream %" PRIi64 " data %zu",
716711
data->id,
@@ -720,8 +715,7 @@ class Http3ApplicationImpl final : public Session::Application {
720715
// nghttp3 tracks its own offset via add_write_offset.
721716
int err = nghttp3_conn_add_write_offset(*this, data->id, datalen);
722717
if (err != 0) {
723-
session().SetLastError(QuicError::ForApplication(
724-
nghttp3_err_infer_quic_app_error_code(err)));
718+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(err));
725719
returnfalse;
726720
}
727721
// Raw application bytes are committed to the stream's outbound
@@ -1212,10 +1206,10 @@ class Http3ApplicationImpl final : public Session::Application {
12121206
void* conn_user_data,
12131207
void* stream_user_data) {
12141208
NGHTTP3_CALLBACK_SCOPE(app);
1215-
auto& session = app.session();
1216-
Debug(&session, "HTTP/3 application deferred consume %zu bytes", consumed);
1217-
session.ExtendStreamOffset(id, consumed);
1218-
session.ExtendOffset(consumed);
1209+
Debug(&app.session(),
1210+
"HTTP/3 application deferred consume %zu bytes",
1211+
consumed);
1212+
app.session().Consume(id, consumed);
12191213
returnNGTCP2_SUCCESS;
12201214
}
12211215

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit e2e2829

Browse files
pimterryaduh95
authored andcommitted
quic: extract transport logic from Application to Session
One small fix notably included: - Check is_destroyed() after StreamCommit, since it calls JS callbacks which could destroy the session. Signed-off-by: Tim Perry <pimterry@gmail.com> PR-URL: #64127 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 7ac97fc commit e2e2829

9 files changed

Lines changed: 608 additions & 562 deletions

File tree

β€Žsrc/quic/README.mdβ€Ž

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ data channels that carry application data.
7777

7878
Every entry point that may generate outbound data creates a
7979
`SendPendingDataScope`. Scopes nest β€” an internal depth counter ensures
80-
`Application::SendPendingData()` is called exactly once, when the outermost
80+
`Session::SendPendingData()` is called exactly once, when the outermost
8181
scope exits:
8282

8383
```cpp
@@ -218,13 +218,13 @@ Session::Receive()
218218

219219
```text
220220
SendPendingDataScope::~SendPendingDataScope()
221-
β†’ Application::SendPendingData()
221+
β†’ Session::SendPendingData()
222222
Loop (up to max_packet_count):
223-
β”œβ”€β”€ GetStreamData() // pull data from next stream
224-
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225-
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
223+
β”œβ”€β”€ application().GetStreamData() // pull data from next stream
224+
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225+
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
226226
β”‚ encrypts, frames, paces
227-
β”œβ”€β”€ if ndatalen > 0: StreamCommit()
227+
β”œβ”€β”€ if ndatalen > 0: application().StreamCommit()
228228
β”‚ stream->Commit(datalen, fin)
229229
β”œβ”€β”€ if nwrite > 0: Send() // uv_udp_send()
230230
β”œβ”€β”€ if WRITE_MORE: continue // room for more in this packet

β€Žsrc/quic/application.ccβ€Ž

Lines changed: 6 additions & 449 deletions
Large diffs are not rendered by default.

β€Žsrc/quic/application.hβ€Ž

Lines changed: 0 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -205,10 +205,6 @@ class Session::Application : public MemoryRetainer {
205205
returnfalse;
206206
}
207207

208-
// Signals to the Application that it should serialize and transmit any
209-
// pending session and stream packets it has accumulated.
210-
voidSendPendingData();
211-
212208
// Returns true if the application protocol supports sending and
213209
// receiving headers on streams (e.g. HTTP/3). Applications that
214210
// do not support headers should return false (the default).
@@ -243,10 +239,6 @@ class Session::Application : public MemoryRetainer {
243239
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
244240
}
245241

246-
// The StreamData struct is used by the application to pass pending stream
247-
// data to the session for transmission.
248-
structStreamData;
249-
250242
virtualintGetStreamData(StreamData* data) = 0;
251243
virtualboolStreamCommit(StreamData* data, size_t datalen) = 0;
252244

@@ -262,57 +254,9 @@ class Session::Application : public MemoryRetainer {
262254
}
263255

264256
private:
265-
Packet::Ptr CreateStreamDataPacket();
266-
267-
// Tries to pack a pending datagram into the current packet buffer.
268-
// If < 0 is returned, either NGTCP2_ERR_WRITE_MORE or a fatal error is
269-
// returned; the caller must check. If > 0 is returned, the packet is done
270-
// and the value is the size of the finalized packet. If 0 is returned,
271-
// the datagram is either congestion limited or was abandoned
272-
ssize_tTryWritePendingDatagram(PathStorage* path,
273-
uint8_t* dest,
274-
size_t destlen,
275-
uint64_t ts);
276-
277-
// Write the given stream_data into the buffer. The PacketInfo out-param
278-
// is populated by ngtcp2 with per-packet metadata (e.g., ECN codepoint)
279-
// that should be applied when sending the packet.
280-
ssize_tWriteVStream(PathStorage* path,
281-
PacketInfo* pi,
282-
uint8_t* buf,
283-
ssize_t* ndatalen,
284-
size_t max_packet_size,
285-
const StreamData& stream_data,
286-
uint64_t ts);
287-
288257
Session* session_ = nullptr;
289258
};
290259

291-
structSession::Application::StreamData final {
292-
// The actual number of vectors in the struct, up to kMaxVectorCount.
293-
size_t count = 0;
294-
// The stream identifier. If this is a negative value then no stream is
295-
// identified.
296-
stream_id id = -1;
297-
int fin = 0;
298-
ngtcp2_vec data[kMaxVectorCount]{};
299-
BaseObjectPtr<Stream> stream;
300-
301-
static_assert(sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
302-
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
303-
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
304-
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
305-
"ngtcp2_vec and nghttp3_vec must have identical layout");
306-
inlineoperator nghttp3_vec*() {
307-
returnreinterpret_cast<nghttp3_vec*>(data);
308-
}
309-
310-
inlineoperatorconst ngtcp2_vec*() const { return data; }
311-
inlineoperator ngtcp2_vec*() { return data; }
312-
313-
std::string ToString() const;
314-
};
315-
316260
// Create a DefaultApplication for the given session.
317261
std::unique_ptr<Session::Application> CreateDefaultApplication(
318262
Session* session, const Session::Application_Options& options);

β€Žsrc/quic/http3.ccβ€Ž

Lines changed: 40 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -209,24 +209,25 @@ class Http3ApplicationImpl final : public Session::Application {
209209
started_ = true;
210210
Debug(&session(), "Starting HTTP/3 application.");
211211

212-
auto params = ngtcp2_conn_get_remote_transport_params(session());
213-
if (params == nullptr) [[unlikely]] {
212+
constauto params = session().remote_transport_params();
213+
if (!params) [[unlikely]] {
214214
// The params are not available yet. Cannot start.
215215
Debug(&session(),
216216
"Cannot start HTTP/3 application yet. No remote transport params");
217217
returnfalse;
218218
}
219219

220-
if (params->initial_max_streams_uni < 3) {
220+
if (params.initial_max_streams_uni() < 3) {
221221
// HTTP3 requires 3 unidirectional control streams to be opened in each
222222
// direction in additional to the bidirectional streams that are used to
223223
// actually carry request and response payload back and forth.
224224
// See:
225225
// https://nghttp2.org/nghttp3/programmers-guide.html#binding-control-streams
226226
Debug(&session(),
227227
"Cannot start HTTP/3 application. Initial max "
228-
"unidirectional streams [%zu] is too low. Must be at least 3",
229-
params->initial_max_streams_uni);
228+
"unidirectional streams [%" PRIu64
229+
"] is too low. Must be at least 3",
230+
params.initial_max_streams_uni());
230231
returnfalse;
231232
}
232233

@@ -235,17 +236,14 @@ class Http3ApplicationImpl final : public Session::Application {
235236
// of requests that the client can actually created.
236237
if (session().is_server()) {
237238
nghttp3_conn_set_max_client_streams_bidi(
238-
*this, params->initial_max_streams_bidi);
239+
*this, params.initial_max_streams_bidi());
239240
}
240241

241242
Debug(&session(), "Creating and binding HTTP/3 control streams");
242243
bool ret =
243-
ngtcp2_conn_open_uni_stream(session(), &control_stream_id_, nullptr) ==
244-
0 &&
245-
ngtcp2_conn_open_uni_stream(
246-
session(), &qpack_enc_stream_id_, nullptr) == 0 &&
247-
ngtcp2_conn_open_uni_stream(
248-
session(), &qpack_dec_stream_id_, nullptr) == 0 &&
244+
session().OpenUnidirectionalStream(&control_stream_id_) &&
245+
session().OpenUnidirectionalStream(&qpack_enc_stream_id_) &&
246+
session().OpenUnidirectionalStream(&qpack_dec_stream_id_) &&
249247
nghttp3_conn_bind_control_stream(*this, control_stream_id_) == 0 &&
250248
nghttp3_conn_bind_qpack_streams(
251249
*this, qpack_enc_stream_id_, qpack_dec_stream_id_) == 0;
@@ -306,8 +304,7 @@ class Http3ApplicationImpl final : public Session::Application {
306304
Debug(&session(),
307305
"Extending stream and connection offset by %zd bytes",
308306
nread);
309-
session().ExtendStreamOffset(id, nread);
310-
session().ExtendOffset(nread);
307+
session().Consume(id, nread);
311308
}
312309

313310
// If this data arrived as 0-RTT, mark the stream. We set it after
@@ -365,24 +362,11 @@ class Http3ApplicationImpl final : public Session::Application {
365362
case EndpointLabel::LOCAL:
366363
return;
367364
case EndpointLabel::REMOTE: {
368-
switch (direction) {
369-
case Direction::BIDIRECTIONAL: {
370-
Debug(&session(),
371-
"HTTP/3 application extending max bidi streams by %" PRIu64,
372-
max_streams);
373-
ngtcp2_conn_extend_max_streams_bidi(
374-
session(), static_cast<size_t>(max_streams));
375-
break;
376-
}
377-
case Direction::UNIDIRECTIONAL: {
378-
Debug(&session(),
379-
"HTTP/3 application extending max uni streams by %" PRIu64,
380-
max_streams);
381-
ngtcp2_conn_extend_max_streams_uni(
382-
session(), static_cast<size_t>(max_streams));
383-
break;
384-
}
385-
}
365+
Debug(&session(),
366+
"HTTP/3 application extending max %s streams by %" PRIu64,
367+
direction == Direction::BIDIRECTIONAL ? "bidi" : "uni",
368+
max_streams);
369+
session().ExtendMaxStreams(direction, max_streams);
386370
}
387371
}
388372
}
@@ -530,8 +514,7 @@ class Http3ApplicationImpl final : public Session::Application {
530514
return;
531515
}
532516

533-
session().SetLastError(
534-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
517+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
535518
session().Close();
536519
}
537520

@@ -548,8 +531,7 @@ class Http3ApplicationImpl final : public Session::Application {
548531
return;
549532
}
550533

551-
session().SetLastError(
552-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
534+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
553535
session().Close();
554536
}
555537

@@ -687,17 +669,30 @@ class Http3ApplicationImpl final : public Session::Application {
687669
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
688670
}
689671

690-
intGetStreamData(StreamData* data) override {
672+
intGetStreamData(Session::StreamData* data) override {
673+
static_assert(
674+
sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
675+
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
676+
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
677+
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
678+
"ngtcp2_vec and nghttp3_vec must have identical layout");
691679
data->count = kMaxVectorCount;
692680
ssize_t ret = 0;
693681
Debug(&session(), "HTTP/3 application getting stream data");
694682
if (conn_ && session().max_data_left()) {
695-
ret = nghttp3_conn_writev_stream(
696-
*this, &data->id, &data->fin, *data, data->count);
683+
// nghttp3 reports fin through an int out-param; bridge it to the bool.
684+
int fin = 0;
685+
ret =
686+
nghttp3_conn_writev_stream(*this,
687+
&data->id,
688+
&fin,
689+
reinterpret_cast<nghttp3_vec*>(data->data),
690+
data->count);
697691
// A negative return value indicates an error.
698692
if (ret < 0) {
699693
returnstatic_cast<int>(ret);
700694
}
695+
data->fin = fin != 0;
701696

702697
data->count = static_cast<size_t>(ret);
703698
if (data->id >= 0 && data->id != control_stream_id_ &&
@@ -710,7 +705,7 @@ class Http3ApplicationImpl final : public Session::Application {
710705
return0;
711706
}
712707

713-
boolStreamCommit(StreamData* data, size_t datalen) override {
708+
boolStreamCommit(Session::StreamData* data, size_t datalen) override {
714709
Debug(&session(),
715710
"HTTP/3 application committing stream %" PRIi64 " data %zu",
716711
data->id,
@@ -720,8 +715,7 @@ class Http3ApplicationImpl final : public Session::Application {
720715
// nghttp3 tracks its own offset via add_write_offset.
721716
int err = nghttp3_conn_add_write_offset(*this, data->id, datalen);
722717
if (err != 0) {
723-
session().SetLastError(QuicError::ForApplication(
724-
nghttp3_err_infer_quic_app_error_code(err)));
718+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(err));
725719
returnfalse;
726720
}
727721
// Raw application bytes are committed to the stream's outbound
@@ -1212,10 +1206,10 @@ class Http3ApplicationImpl final : public Session::Application {
12121206
void* conn_user_data,
12131207
void* stream_user_data) {
12141208
NGHTTP3_CALLBACK_SCOPE(app);
1215-
auto& session = app.session();
1216-
Debug(&session, "HTTP/3 application deferred consume %zu bytes", consumed);
1217-
session.ExtendStreamOffset(id, consumed);
1218-
session.ExtendOffset(consumed);
1209+
Debug(&app.session(),
1210+
"HTTP/3 application deferred consume %zu bytes",
1211+
consumed);
1212+
app.session().Consume(id, consumed);
12191213
returnNGTCP2_SUCCESS;
12201214
}
12211215

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit e2e2829

Browse files
pimterryaduh95
authored andcommitted
quic: extract transport logic from Application to Session
One small fix notably included: - Check is_destroyed() after StreamCommit, since it calls JS callbacks which could destroy the session. Signed-off-by: Tim Perry <pimterry@gmail.com> PR-URL: #64127 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 7ac97fc commit e2e2829

9 files changed

Lines changed: 608 additions & 562 deletions

File tree

β€Žsrc/quic/README.mdβ€Ž

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ data channels that carry application data.
7777

7878
Every entry point that may generate outbound data creates a
7979
`SendPendingDataScope`. Scopes nest β€” an internal depth counter ensures
80-
`Application::SendPendingData()` is called exactly once, when the outermost
80+
`Session::SendPendingData()` is called exactly once, when the outermost
8181
scope exits:
8282

8383
```cpp
@@ -218,13 +218,13 @@ Session::Receive()
218218

219219
```text
220220
SendPendingDataScope::~SendPendingDataScope()
221-
β†’ Application::SendPendingData()
221+
β†’ Session::SendPendingData()
222222
Loop (up to max_packet_count):
223-
β”œβ”€β”€ GetStreamData() // pull data from next stream
224-
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225-
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
223+
β”œβ”€β”€ application().GetStreamData() // pull data from next stream
224+
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225+
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
226226
β”‚ encrypts, frames, paces
227-
β”œβ”€β”€ if ndatalen > 0: StreamCommit()
227+
β”œβ”€β”€ if ndatalen > 0: application().StreamCommit()
228228
β”‚ stream->Commit(datalen, fin)
229229
β”œβ”€β”€ if nwrite > 0: Send() // uv_udp_send()
230230
β”œβ”€β”€ if WRITE_MORE: continue // room for more in this packet

β€Žsrc/quic/application.ccβ€Ž

Lines changed: 6 additions & 449 deletions
Large diffs are not rendered by default.

β€Žsrc/quic/application.hβ€Ž

Lines changed: 0 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -205,10 +205,6 @@ class Session::Application : public MemoryRetainer {
205205
returnfalse;
206206
}
207207

208-
// Signals to the Application that it should serialize and transmit any
209-
// pending session and stream packets it has accumulated.
210-
voidSendPendingData();
211-
212208
// Returns true if the application protocol supports sending and
213209
// receiving headers on streams (e.g. HTTP/3). Applications that
214210
// do not support headers should return false (the default).
@@ -243,10 +239,6 @@ class Session::Application : public MemoryRetainer {
243239
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
244240
}
245241

246-
// The StreamData struct is used by the application to pass pending stream
247-
// data to the session for transmission.
248-
structStreamData;
249-
250242
virtualintGetStreamData(StreamData* data) = 0;
251243
virtualboolStreamCommit(StreamData* data, size_t datalen) = 0;
252244

@@ -262,57 +254,9 @@ class Session::Application : public MemoryRetainer {
262254
}
263255

264256
private:
265-
Packet::Ptr CreateStreamDataPacket();
266-
267-
// Tries to pack a pending datagram into the current packet buffer.
268-
// If < 0 is returned, either NGTCP2_ERR_WRITE_MORE or a fatal error is
269-
// returned; the caller must check. If > 0 is returned, the packet is done
270-
// and the value is the size of the finalized packet. If 0 is returned,
271-
// the datagram is either congestion limited or was abandoned
272-
ssize_tTryWritePendingDatagram(PathStorage* path,
273-
uint8_t* dest,
274-
size_t destlen,
275-
uint64_t ts);
276-
277-
// Write the given stream_data into the buffer. The PacketInfo out-param
278-
// is populated by ngtcp2 with per-packet metadata (e.g., ECN codepoint)
279-
// that should be applied when sending the packet.
280-
ssize_tWriteVStream(PathStorage* path,
281-
PacketInfo* pi,
282-
uint8_t* buf,
283-
ssize_t* ndatalen,
284-
size_t max_packet_size,
285-
const StreamData& stream_data,
286-
uint64_t ts);
287-
288257
Session* session_ = nullptr;
289258
};
290259

291-
structSession::Application::StreamData final {
292-
// The actual number of vectors in the struct, up to kMaxVectorCount.
293-
size_t count = 0;
294-
// The stream identifier. If this is a negative value then no stream is
295-
// identified.
296-
stream_id id = -1;
297-
int fin = 0;
298-
ngtcp2_vec data[kMaxVectorCount]{};
299-
BaseObjectPtr<Stream> stream;
300-
301-
static_assert(sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
302-
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
303-
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
304-
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
305-
"ngtcp2_vec and nghttp3_vec must have identical layout");
306-
inlineoperator nghttp3_vec*() {
307-
returnreinterpret_cast<nghttp3_vec*>(data);
308-
}
309-
310-
inlineoperatorconst ngtcp2_vec*() const { return data; }
311-
inlineoperator ngtcp2_vec*() { return data; }
312-
313-
std::string ToString() const;
314-
};
315-
316260
// Create a DefaultApplication for the given session.
317261
std::unique_ptr<Session::Application> CreateDefaultApplication(
318262
Session* session, const Session::Application_Options& options);

β€Žsrc/quic/http3.ccβ€Ž

Lines changed: 40 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -209,24 +209,25 @@ class Http3ApplicationImpl final : public Session::Application {
209209
started_ = true;
210210
Debug(&session(), "Starting HTTP/3 application.");
211211

212-
auto params = ngtcp2_conn_get_remote_transport_params(session());
213-
if (params == nullptr) [[unlikely]] {
212+
constauto params = session().remote_transport_params();
213+
if (!params) [[unlikely]] {
214214
// The params are not available yet. Cannot start.
215215
Debug(&session(),
216216
"Cannot start HTTP/3 application yet. No remote transport params");
217217
returnfalse;
218218
}
219219

220-
if (params->initial_max_streams_uni < 3) {
220+
if (params.initial_max_streams_uni() < 3) {
221221
// HTTP3 requires 3 unidirectional control streams to be opened in each
222222
// direction in additional to the bidirectional streams that are used to
223223
// actually carry request and response payload back and forth.
224224
// See:
225225
// https://nghttp2.org/nghttp3/programmers-guide.html#binding-control-streams
226226
Debug(&session(),
227227
"Cannot start HTTP/3 application. Initial max "
228-
"unidirectional streams [%zu] is too low. Must be at least 3",
229-
params->initial_max_streams_uni);
228+
"unidirectional streams [%" PRIu64
229+
"] is too low. Must be at least 3",
230+
params.initial_max_streams_uni());
230231
returnfalse;
231232
}
232233

@@ -235,17 +236,14 @@ class Http3ApplicationImpl final : public Session::Application {
235236
// of requests that the client can actually created.
236237
if (session().is_server()) {
237238
nghttp3_conn_set_max_client_streams_bidi(
238-
*this, params->initial_max_streams_bidi);
239+
*this, params.initial_max_streams_bidi());
239240
}
240241

241242
Debug(&session(), "Creating and binding HTTP/3 control streams");
242243
bool ret =
243-
ngtcp2_conn_open_uni_stream(session(), &control_stream_id_, nullptr) ==
244-
0 &&
245-
ngtcp2_conn_open_uni_stream(
246-
session(), &qpack_enc_stream_id_, nullptr) == 0 &&
247-
ngtcp2_conn_open_uni_stream(
248-
session(), &qpack_dec_stream_id_, nullptr) == 0 &&
244+
session().OpenUnidirectionalStream(&control_stream_id_) &&
245+
session().OpenUnidirectionalStream(&qpack_enc_stream_id_) &&
246+
session().OpenUnidirectionalStream(&qpack_dec_stream_id_) &&
249247
nghttp3_conn_bind_control_stream(*this, control_stream_id_) == 0 &&
250248
nghttp3_conn_bind_qpack_streams(
251249
*this, qpack_enc_stream_id_, qpack_dec_stream_id_) == 0;
@@ -306,8 +304,7 @@ class Http3ApplicationImpl final : public Session::Application {
306304
Debug(&session(),
307305
"Extending stream and connection offset by %zd bytes",
308306
nread);
309-
session().ExtendStreamOffset(id, nread);
310-
session().ExtendOffset(nread);
307+
session().Consume(id, nread);
311308
}
312309

313310
// If this data arrived as 0-RTT, mark the stream. We set it after
@@ -365,24 +362,11 @@ class Http3ApplicationImpl final : public Session::Application {
365362
case EndpointLabel::LOCAL:
366363
return;
367364
case EndpointLabel::REMOTE: {
368-
switch (direction) {
369-
case Direction::BIDIRECTIONAL: {
370-
Debug(&session(),
371-
"HTTP/3 application extending max bidi streams by %" PRIu64,
372-
max_streams);
373-
ngtcp2_conn_extend_max_streams_bidi(
374-
session(), static_cast<size_t>(max_streams));
375-
break;
376-
}
377-
case Direction::UNIDIRECTIONAL: {
378-
Debug(&session(),
379-
"HTTP/3 application extending max uni streams by %" PRIu64,
380-
max_streams);
381-
ngtcp2_conn_extend_max_streams_uni(
382-
session(), static_cast<size_t>(max_streams));
383-
break;
384-
}
385-
}
365+
Debug(&session(),
366+
"HTTP/3 application extending max %s streams by %" PRIu64,
367+
direction == Direction::BIDIRECTIONAL ? "bidi" : "uni",
368+
max_streams);
369+
session().ExtendMaxStreams(direction, max_streams);
386370
}
387371
}
388372
}
@@ -530,8 +514,7 @@ class Http3ApplicationImpl final : public Session::Application {
530514
return;
531515
}
532516

533-
session().SetLastError(
534-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
517+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
535518
session().Close();
536519
}
537520

@@ -548,8 +531,7 @@ class Http3ApplicationImpl final : public Session::Application {
548531
return;
549532
}
550533

551-
session().SetLastError(
552-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
534+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
553535
session().Close();
554536
}
555537

@@ -687,17 +669,30 @@ class Http3ApplicationImpl final : public Session::Application {
687669
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
688670
}
689671

690-
intGetStreamData(StreamData* data) override {
672+
intGetStreamData(Session::StreamData* data) override {
673+
static_assert(
674+
sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
675+
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
676+
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
677+
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
678+
"ngtcp2_vec and nghttp3_vec must have identical layout");
691679
data->count = kMaxVectorCount;
692680
ssize_t ret = 0;
693681
Debug(&session(), "HTTP/3 application getting stream data");
694682
if (conn_ && session().max_data_left()) {
695-
ret = nghttp3_conn_writev_stream(
696-
*this, &data->id, &data->fin, *data, data->count);
683+
// nghttp3 reports fin through an int out-param; bridge it to the bool.
684+
int fin = 0;
685+
ret =
686+
nghttp3_conn_writev_stream(*this,
687+
&data->id,
688+
&fin,
689+
reinterpret_cast<nghttp3_vec*>(data->data),
690+
data->count);
697691
// A negative return value indicates an error.
698692
if (ret < 0) {
699693
returnstatic_cast<int>(ret);
700694
}
695+
data->fin = fin != 0;
701696

702697
data->count = static_cast<size_t>(ret);
703698
if (data->id >= 0 && data->id != control_stream_id_ &&
@@ -710,7 +705,7 @@ class Http3ApplicationImpl final : public Session::Application {
710705
return0;
711706
}
712707

713-
boolStreamCommit(StreamData* data, size_t datalen) override {
708+
boolStreamCommit(Session::StreamData* data, size_t datalen) override {
714709
Debug(&session(),
715710
"HTTP/3 application committing stream %" PRIi64 " data %zu",
716711
data->id,
@@ -720,8 +715,7 @@ class Http3ApplicationImpl final : public Session::Application {
720715
// nghttp3 tracks its own offset via add_write_offset.
721716
int err = nghttp3_conn_add_write_offset(*this, data->id, datalen);
722717
if (err != 0) {
723-
session().SetLastError(QuicError::ForApplication(
724-
nghttp3_err_infer_quic_app_error_code(err)));
718+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(err));
725719
returnfalse;
726720
}
727721
// Raw application bytes are committed to the stream's outbound
@@ -1212,10 +1206,10 @@ class Http3ApplicationImpl final : public Session::Application {
12121206
void* conn_user_data,
12131207
void* stream_user_data) {
12141208
NGHTTP3_CALLBACK_SCOPE(app);
1215-
auto& session = app.session();
1216-
Debug(&session, "HTTP/3 application deferred consume %zu bytes", consumed);
1217-
session.ExtendStreamOffset(id, consumed);
1218-
session.ExtendOffset(consumed);
1209+
Debug(&app.session(),
1210+
"HTTP/3 application deferred consume %zu bytes",
1211+
consumed);
1212+
app.session().Consume(id, consumed);
12191213
returnNGTCP2_SUCCESS;
12201214
}
12211215

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Commit e2e2829

Browse files
pimterryaduh95
authored andcommitted
quic: extract transport logic from Application to Session
One small fix notably included: - Check is_destroyed() after StreamCommit, since it calls JS callbacks which could destroy the session. Signed-off-by: Tim Perry <pimterry@gmail.com> PR-URL: #64127 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 7ac97fc commit e2e2829

9 files changed

Lines changed: 608 additions & 562 deletions

File tree

β€Žsrc/quic/README.mdβ€Ž

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ data channels that carry application data.
7777

7878
Every entry point that may generate outbound data creates a
7979
`SendPendingDataScope`. Scopes nest β€” an internal depth counter ensures
80-
`Application::SendPendingData()` is called exactly once, when the outermost
80+
`Session::SendPendingData()` is called exactly once, when the outermost
8181
scope exits:
8282

8383
```cpp
@@ -218,13 +218,13 @@ Session::Receive()
218218

219219
```text
220220
SendPendingDataScope::~SendPendingDataScope()
221-
β†’ Application::SendPendingData()
221+
β†’ Session::SendPendingData()
222222
Loop (up to max_packet_count):
223-
β”œβ”€β”€ GetStreamData() // pull data from next stream
224-
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225-
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
223+
β”œβ”€β”€ application().GetStreamData() // pull data from next stream
224+
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225+
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
226226
β”‚ encrypts, frames, paces
227-
β”œβ”€β”€ if ndatalen > 0: StreamCommit()
227+
β”œβ”€β”€ if ndatalen > 0: application().StreamCommit()
228228
β”‚ stream->Commit(datalen, fin)
229229
β”œβ”€β”€ if nwrite > 0: Send() // uv_udp_send()
230230
β”œβ”€β”€ if WRITE_MORE: continue // room for more in this packet

β€Žsrc/quic/application.ccβ€Ž

Lines changed: 6 additions & 449 deletions
Large diffs are not rendered by default.

β€Žsrc/quic/application.hβ€Ž

Lines changed: 0 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -205,10 +205,6 @@ class Session::Application : public MemoryRetainer {
205205
returnfalse;
206206
}
207207

208-
// Signals to the Application that it should serialize and transmit any
209-
// pending session and stream packets it has accumulated.
210-
voidSendPendingData();
211-
212208
// Returns true if the application protocol supports sending and
213209
// receiving headers on streams (e.g. HTTP/3). Applications that
214210
// do not support headers should return false (the default).
@@ -243,10 +239,6 @@ class Session::Application : public MemoryRetainer {
243239
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
244240
}
245241

246-
// The StreamData struct is used by the application to pass pending stream
247-
// data to the session for transmission.
248-
structStreamData;
249-
250242
virtualintGetStreamData(StreamData* data) = 0;
251243
virtualboolStreamCommit(StreamData* data, size_t datalen) = 0;
252244

@@ -262,57 +254,9 @@ class Session::Application : public MemoryRetainer {
262254
}
263255

264256
private:
265-
Packet::Ptr CreateStreamDataPacket();
266-
267-
// Tries to pack a pending datagram into the current packet buffer.
268-
// If < 0 is returned, either NGTCP2_ERR_WRITE_MORE or a fatal error is
269-
// returned; the caller must check. If > 0 is returned, the packet is done
270-
// and the value is the size of the finalized packet. If 0 is returned,
271-
// the datagram is either congestion limited or was abandoned
272-
ssize_tTryWritePendingDatagram(PathStorage* path,
273-
uint8_t* dest,
274-
size_t destlen,
275-
uint64_t ts);
276-
277-
// Write the given stream_data into the buffer. The PacketInfo out-param
278-
// is populated by ngtcp2 with per-packet metadata (e.g., ECN codepoint)
279-
// that should be applied when sending the packet.
280-
ssize_tWriteVStream(PathStorage* path,
281-
PacketInfo* pi,
282-
uint8_t* buf,
283-
ssize_t* ndatalen,
284-
size_t max_packet_size,
285-
const StreamData& stream_data,
286-
uint64_t ts);
287-
288257
Session* session_ = nullptr;
289258
};
290259

291-
structSession::Application::StreamData final {
292-
// The actual number of vectors in the struct, up to kMaxVectorCount.
293-
size_t count = 0;
294-
// The stream identifier. If this is a negative value then no stream is
295-
// identified.
296-
stream_id id = -1;
297-
int fin = 0;
298-
ngtcp2_vec data[kMaxVectorCount]{};
299-
BaseObjectPtr<Stream> stream;
300-
301-
static_assert(sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
302-
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
303-
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
304-
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
305-
"ngtcp2_vec and nghttp3_vec must have identical layout");
306-
inlineoperator nghttp3_vec*() {
307-
returnreinterpret_cast<nghttp3_vec*>(data);
308-
}
309-
310-
inlineoperatorconst ngtcp2_vec*() const { return data; }
311-
inlineoperator ngtcp2_vec*() { return data; }
312-
313-
std::string ToString() const;
314-
};
315-
316260
// Create a DefaultApplication for the given session.
317261
std::unique_ptr<Session::Application> CreateDefaultApplication(
318262
Session* session, const Session::Application_Options& options);

β€Žsrc/quic/http3.ccβ€Ž

Lines changed: 40 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -209,24 +209,25 @@ class Http3ApplicationImpl final : public Session::Application {
209209
started_ = true;
210210
Debug(&session(), "Starting HTTP/3 application.");
211211

212-
auto params = ngtcp2_conn_get_remote_transport_params(session());
213-
if (params == nullptr) [[unlikely]] {
212+
constauto params = session().remote_transport_params();
213+
if (!params) [[unlikely]] {
214214
// The params are not available yet. Cannot start.
215215
Debug(&session(),
216216
"Cannot start HTTP/3 application yet. No remote transport params");
217217
returnfalse;
218218
}
219219

220-
if (params->initial_max_streams_uni < 3) {
220+
if (params.initial_max_streams_uni() < 3) {
221221
// HTTP3 requires 3 unidirectional control streams to be opened in each
222222
// direction in additional to the bidirectional streams that are used to
223223
// actually carry request and response payload back and forth.
224224
// See:
225225
// https://nghttp2.org/nghttp3/programmers-guide.html#binding-control-streams
226226
Debug(&session(),
227227
"Cannot start HTTP/3 application. Initial max "
228-
"unidirectional streams [%zu] is too low. Must be at least 3",
229-
params->initial_max_streams_uni);
228+
"unidirectional streams [%" PRIu64
229+
"] is too low. Must be at least 3",
230+
params.initial_max_streams_uni());
230231
returnfalse;
231232
}
232233

@@ -235,17 +236,14 @@ class Http3ApplicationImpl final : public Session::Application {
235236
// of requests that the client can actually created.
236237
if (session().is_server()) {
237238
nghttp3_conn_set_max_client_streams_bidi(
238-
*this, params->initial_max_streams_bidi);
239+
*this, params.initial_max_streams_bidi());
239240
}
240241

241242
Debug(&session(), "Creating and binding HTTP/3 control streams");
242243
bool ret =
243-
ngtcp2_conn_open_uni_stream(session(), &control_stream_id_, nullptr) ==
244-
0 &&
245-
ngtcp2_conn_open_uni_stream(
246-
session(), &qpack_enc_stream_id_, nullptr) == 0 &&
247-
ngtcp2_conn_open_uni_stream(
248-
session(), &qpack_dec_stream_id_, nullptr) == 0 &&
244+
session().OpenUnidirectionalStream(&control_stream_id_) &&
245+
session().OpenUnidirectionalStream(&qpack_enc_stream_id_) &&
246+
session().OpenUnidirectionalStream(&qpack_dec_stream_id_) &&
249247
nghttp3_conn_bind_control_stream(*this, control_stream_id_) == 0 &&
250248
nghttp3_conn_bind_qpack_streams(
251249
*this, qpack_enc_stream_id_, qpack_dec_stream_id_) == 0;
@@ -306,8 +304,7 @@ class Http3ApplicationImpl final : public Session::Application {
306304
Debug(&session(),
307305
"Extending stream and connection offset by %zd bytes",
308306
nread);
309-
session().ExtendStreamOffset(id, nread);
310-
session().ExtendOffset(nread);
307+
session().Consume(id, nread);
311308
}
312309

313310
// If this data arrived as 0-RTT, mark the stream. We set it after
@@ -365,24 +362,11 @@ class Http3ApplicationImpl final : public Session::Application {
365362
case EndpointLabel::LOCAL:
366363
return;
367364
case EndpointLabel::REMOTE: {
368-
switch (direction) {
369-
case Direction::BIDIRECTIONAL: {
370-
Debug(&session(),
371-
"HTTP/3 application extending max bidi streams by %" PRIu64,
372-
max_streams);
373-
ngtcp2_conn_extend_max_streams_bidi(
374-
session(), static_cast<size_t>(max_streams));
375-
break;
376-
}
377-
case Direction::UNIDIRECTIONAL: {
378-
Debug(&session(),
379-
"HTTP/3 application extending max uni streams by %" PRIu64,
380-
max_streams);
381-
ngtcp2_conn_extend_max_streams_uni(
382-
session(), static_cast<size_t>(max_streams));
383-
break;
384-
}
385-
}
365+
Debug(&session(),
366+
"HTTP/3 application extending max %s streams by %" PRIu64,
367+
direction == Direction::BIDIRECTIONAL ? "bidi" : "uni",
368+
max_streams);
369+
session().ExtendMaxStreams(direction, max_streams);
386370
}
387371
}
388372
}
@@ -530,8 +514,7 @@ class Http3ApplicationImpl final : public Session::Application {
530514
return;
531515
}
532516

533-
session().SetLastError(
534-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
517+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
535518
session().Close();
536519
}
537520

@@ -548,8 +531,7 @@ class Http3ApplicationImpl final : public Session::Application {
548531
return;
549532
}
550533

551-
session().SetLastError(
552-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
534+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
553535
session().Close();
554536
}
555537

@@ -687,17 +669,30 @@ class Http3ApplicationImpl final : public Session::Application {
687669
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
688670
}
689671

690-
intGetStreamData(StreamData* data) override {
672+
intGetStreamData(Session::StreamData* data) override {
673+
static_assert(
674+
sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
675+
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
676+
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
677+
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
678+
"ngtcp2_vec and nghttp3_vec must have identical layout");
691679
data->count = kMaxVectorCount;
692680
ssize_t ret = 0;
693681
Debug(&session(), "HTTP/3 application getting stream data");
694682
if (conn_ && session().max_data_left()) {
695-
ret = nghttp3_conn_writev_stream(
696-
*this, &data->id, &data->fin, *data, data->count);
683+
// nghttp3 reports fin through an int out-param; bridge it to the bool.
684+
int fin = 0;
685+
ret =
686+
nghttp3_conn_writev_stream(*this,
687+
&data->id,
688+
&fin,
689+
reinterpret_cast<nghttp3_vec*>(data->data),
690+
data->count);
697691
// A negative return value indicates an error.
698692
if (ret < 0) {
699693
returnstatic_cast<int>(ret);
700694
}
695+
data->fin = fin != 0;
701696

702697
data->count = static_cast<size_t>(ret);
703698
if (data->id >= 0 && data->id != control_stream_id_ &&
@@ -710,7 +705,7 @@ class Http3ApplicationImpl final : public Session::Application {
710705
return0;
711706
}
712707

713-
boolStreamCommit(StreamData* data, size_t datalen) override {
708+
boolStreamCommit(Session::StreamData* data, size_t datalen) override {
714709
Debug(&session(),
715710
"HTTP/3 application committing stream %" PRIi64 " data %zu",
716711
data->id,
@@ -720,8 +715,7 @@ class Http3ApplicationImpl final : public Session::Application {
720715
// nghttp3 tracks its own offset via add_write_offset.
721716
int err = nghttp3_conn_add_write_offset(*this, data->id, datalen);
722717
if (err != 0) {
723-
session().SetLastError(QuicError::ForApplication(
724-
nghttp3_err_infer_quic_app_error_code(err)));
718+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(err));
725719
returnfalse;
726720
}
727721
// Raw application bytes are committed to the stream's outbound
@@ -1212,10 +1206,10 @@ class Http3ApplicationImpl final : public Session::Application {
12121206
void* conn_user_data,
12131207
void* stream_user_data) {
12141208
NGHTTP3_CALLBACK_SCOPE(app);
1215-
auto& session = app.session();
1216-
Debug(&session, "HTTP/3 application deferred consume %zu bytes", consumed);
1217-
session.ExtendStreamOffset(id, consumed);
1218-
session.ExtendOffset(consumed);
1209+
Debug(&app.session(),
1210+
"HTTP/3 application deferred consume %zu bytes",
1211+
consumed);
1212+
app.session().Consume(id, consumed);
12191213
returnNGTCP2_SUCCESS;
12201214
}
12211215

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit e2e2829

Browse files
pimterryaduh95
authored andcommitted
quic: extract transport logic from Application to Session
One small fix notably included: - Check is_destroyed() after StreamCommit, since it calls JS callbacks which could destroy the session. Signed-off-by: Tim Perry <pimterry@gmail.com> PR-URL: #64127 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 7ac97fc commit e2e2829

9 files changed

Lines changed: 608 additions & 562 deletions

File tree

β€Žsrc/quic/README.mdβ€Ž

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ data channels that carry application data.
7777

7878
Every entry point that may generate outbound data creates a
7979
`SendPendingDataScope`. Scopes nest β€” an internal depth counter ensures
80-
`Application::SendPendingData()` is called exactly once, when the outermost
80+
`Session::SendPendingData()` is called exactly once, when the outermost
8181
scope exits:
8282

8383
```cpp
@@ -218,13 +218,13 @@ Session::Receive()
218218

219219
```text
220220
SendPendingDataScope::~SendPendingDataScope()
221-
β†’ Application::SendPendingData()
221+
β†’ Session::SendPendingData()
222222
Loop (up to max_packet_count):
223-
β”œβ”€β”€ GetStreamData() // pull data from next stream
224-
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225-
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
223+
β”œβ”€β”€ application().GetStreamData() // pull data from next stream
224+
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225+
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
226226
β”‚ encrypts, frames, paces
227-
β”œβ”€β”€ if ndatalen > 0: StreamCommit()
227+
β”œβ”€β”€ if ndatalen > 0: application().StreamCommit()
228228
β”‚ stream->Commit(datalen, fin)
229229
β”œβ”€β”€ if nwrite > 0: Send() // uv_udp_send()
230230
β”œβ”€β”€ if WRITE_MORE: continue // room for more in this packet

β€Žsrc/quic/application.ccβ€Ž

Lines changed: 6 additions & 449 deletions
Large diffs are not rendered by default.

β€Žsrc/quic/application.hβ€Ž

Lines changed: 0 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -205,10 +205,6 @@ class Session::Application : public MemoryRetainer {
205205
returnfalse;
206206
}
207207

208-
// Signals to the Application that it should serialize and transmit any
209-
// pending session and stream packets it has accumulated.
210-
voidSendPendingData();
211-
212208
// Returns true if the application protocol supports sending and
213209
// receiving headers on streams (e.g. HTTP/3). Applications that
214210
// do not support headers should return false (the default).
@@ -243,10 +239,6 @@ class Session::Application : public MemoryRetainer {
243239
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
244240
}
245241

246-
// The StreamData struct is used by the application to pass pending stream
247-
// data to the session for transmission.
248-
structStreamData;
249-
250242
virtualintGetStreamData(StreamData* data) = 0;
251243
virtualboolStreamCommit(StreamData* data, size_t datalen) = 0;
252244

@@ -262,57 +254,9 @@ class Session::Application : public MemoryRetainer {
262254
}
263255

264256
private:
265-
Packet::Ptr CreateStreamDataPacket();
266-
267-
// Tries to pack a pending datagram into the current packet buffer.
268-
// If < 0 is returned, either NGTCP2_ERR_WRITE_MORE or a fatal error is
269-
// returned; the caller must check. If > 0 is returned, the packet is done
270-
// and the value is the size of the finalized packet. If 0 is returned,
271-
// the datagram is either congestion limited or was abandoned
272-
ssize_tTryWritePendingDatagram(PathStorage* path,
273-
uint8_t* dest,
274-
size_t destlen,
275-
uint64_t ts);
276-
277-
// Write the given stream_data into the buffer. The PacketInfo out-param
278-
// is populated by ngtcp2 with per-packet metadata (e.g., ECN codepoint)
279-
// that should be applied when sending the packet.
280-
ssize_tWriteVStream(PathStorage* path,
281-
PacketInfo* pi,
282-
uint8_t* buf,
283-
ssize_t* ndatalen,
284-
size_t max_packet_size,
285-
const StreamData& stream_data,
286-
uint64_t ts);
287-
288257
Session* session_ = nullptr;
289258
};
290259

291-
structSession::Application::StreamData final {
292-
// The actual number of vectors in the struct, up to kMaxVectorCount.
293-
size_t count = 0;
294-
// The stream identifier. If this is a negative value then no stream is
295-
// identified.
296-
stream_id id = -1;
297-
int fin = 0;
298-
ngtcp2_vec data[kMaxVectorCount]{};
299-
BaseObjectPtr<Stream> stream;
300-
301-
static_assert(sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
302-
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
303-
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
304-
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
305-
"ngtcp2_vec and nghttp3_vec must have identical layout");
306-
inlineoperator nghttp3_vec*() {
307-
returnreinterpret_cast<nghttp3_vec*>(data);
308-
}
309-
310-
inlineoperatorconst ngtcp2_vec*() const { return data; }
311-
inlineoperator ngtcp2_vec*() { return data; }
312-
313-
std::string ToString() const;
314-
};
315-
316260
// Create a DefaultApplication for the given session.
317261
std::unique_ptr<Session::Application> CreateDefaultApplication(
318262
Session* session, const Session::Application_Options& options);

β€Žsrc/quic/http3.ccβ€Ž

Lines changed: 40 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -209,24 +209,25 @@ class Http3ApplicationImpl final : public Session::Application {
209209
started_ = true;
210210
Debug(&session(), "Starting HTTP/3 application.");
211211

212-
auto params = ngtcp2_conn_get_remote_transport_params(session());
213-
if (params == nullptr) [[unlikely]] {
212+
constauto params = session().remote_transport_params();
213+
if (!params) [[unlikely]] {
214214
// The params are not available yet. Cannot start.
215215
Debug(&session(),
216216
"Cannot start HTTP/3 application yet. No remote transport params");
217217
returnfalse;
218218
}
219219

220-
if (params->initial_max_streams_uni < 3) {
220+
if (params.initial_max_streams_uni() < 3) {
221221
// HTTP3 requires 3 unidirectional control streams to be opened in each
222222
// direction in additional to the bidirectional streams that are used to
223223
// actually carry request and response payload back and forth.
224224
// See:
225225
// https://nghttp2.org/nghttp3/programmers-guide.html#binding-control-streams
226226
Debug(&session(),
227227
"Cannot start HTTP/3 application. Initial max "
228-
"unidirectional streams [%zu] is too low. Must be at least 3",
229-
params->initial_max_streams_uni);
228+
"unidirectional streams [%" PRIu64
229+
"] is too low. Must be at least 3",
230+
params.initial_max_streams_uni());
230231
returnfalse;
231232
}
232233

@@ -235,17 +236,14 @@ class Http3ApplicationImpl final : public Session::Application {
235236
// of requests that the client can actually created.
236237
if (session().is_server()) {
237238
nghttp3_conn_set_max_client_streams_bidi(
238-
*this, params->initial_max_streams_bidi);
239+
*this, params.initial_max_streams_bidi());
239240
}
240241

241242
Debug(&session(), "Creating and binding HTTP/3 control streams");
242243
bool ret =
243-
ngtcp2_conn_open_uni_stream(session(), &control_stream_id_, nullptr) ==
244-
0 &&
245-
ngtcp2_conn_open_uni_stream(
246-
session(), &qpack_enc_stream_id_, nullptr) == 0 &&
247-
ngtcp2_conn_open_uni_stream(
248-
session(), &qpack_dec_stream_id_, nullptr) == 0 &&
244+
session().OpenUnidirectionalStream(&control_stream_id_) &&
245+
session().OpenUnidirectionalStream(&qpack_enc_stream_id_) &&
246+
session().OpenUnidirectionalStream(&qpack_dec_stream_id_) &&
249247
nghttp3_conn_bind_control_stream(*this, control_stream_id_) == 0 &&
250248
nghttp3_conn_bind_qpack_streams(
251249
*this, qpack_enc_stream_id_, qpack_dec_stream_id_) == 0;
@@ -306,8 +304,7 @@ class Http3ApplicationImpl final : public Session::Application {
306304
Debug(&session(),
307305
"Extending stream and connection offset by %zd bytes",
308306
nread);
309-
session().ExtendStreamOffset(id, nread);
310-
session().ExtendOffset(nread);
307+
session().Consume(id, nread);
311308
}
312309

313310
// If this data arrived as 0-RTT, mark the stream. We set it after
@@ -365,24 +362,11 @@ class Http3ApplicationImpl final : public Session::Application {
365362
case EndpointLabel::LOCAL:
366363
return;
367364
case EndpointLabel::REMOTE: {
368-
switch (direction) {
369-
case Direction::BIDIRECTIONAL: {
370-
Debug(&session(),
371-
"HTTP/3 application extending max bidi streams by %" PRIu64,
372-
max_streams);
373-
ngtcp2_conn_extend_max_streams_bidi(
374-
session(), static_cast<size_t>(max_streams));
375-
break;
376-
}
377-
case Direction::UNIDIRECTIONAL: {
378-
Debug(&session(),
379-
"HTTP/3 application extending max uni streams by %" PRIu64,
380-
max_streams);
381-
ngtcp2_conn_extend_max_streams_uni(
382-
session(), static_cast<size_t>(max_streams));
383-
break;
384-
}
385-
}
365+
Debug(&session(),
366+
"HTTP/3 application extending max %s streams by %" PRIu64,
367+
direction == Direction::BIDIRECTIONAL ? "bidi" : "uni",
368+
max_streams);
369+
session().ExtendMaxStreams(direction, max_streams);
386370
}
387371
}
388372
}
@@ -530,8 +514,7 @@ class Http3ApplicationImpl final : public Session::Application {
530514
return;
531515
}
532516

533-
session().SetLastError(
534-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
517+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
535518
session().Close();
536519
}
537520

@@ -548,8 +531,7 @@ class Http3ApplicationImpl final : public Session::Application {
548531
return;
549532
}
550533

551-
session().SetLastError(
552-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
534+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
553535
session().Close();
554536
}
555537

@@ -687,17 +669,30 @@ class Http3ApplicationImpl final : public Session::Application {
687669
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
688670
}
689671

690-
intGetStreamData(StreamData* data) override {
672+
intGetStreamData(Session::StreamData* data) override {
673+
static_assert(
674+
sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
675+
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
676+
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
677+
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
678+
"ngtcp2_vec and nghttp3_vec must have identical layout");
691679
data->count = kMaxVectorCount;
692680
ssize_t ret = 0;
693681
Debug(&session(), "HTTP/3 application getting stream data");
694682
if (conn_ && session().max_data_left()) {
695-
ret = nghttp3_conn_writev_stream(
696-
*this, &data->id, &data->fin, *data, data->count);
683+
// nghttp3 reports fin through an int out-param; bridge it to the bool.
684+
int fin = 0;
685+
ret =
686+
nghttp3_conn_writev_stream(*this,
687+
&data->id,
688+
&fin,
689+
reinterpret_cast<nghttp3_vec*>(data->data),
690+
data->count);
697691
// A negative return value indicates an error.
698692
if (ret < 0) {
699693
returnstatic_cast<int>(ret);
700694
}
695+
data->fin = fin != 0;
701696

702697
data->count = static_cast<size_t>(ret);
703698
if (data->id >= 0 && data->id != control_stream_id_ &&
@@ -710,7 +705,7 @@ class Http3ApplicationImpl final : public Session::Application {
710705
return0;
711706
}
712707

713-
boolStreamCommit(StreamData* data, size_t datalen) override {
708+
boolStreamCommit(Session::StreamData* data, size_t datalen) override {
714709
Debug(&session(),
715710
"HTTP/3 application committing stream %" PRIi64 " data %zu",
716711
data->id,
@@ -720,8 +715,7 @@ class Http3ApplicationImpl final : public Session::Application {
720715
// nghttp3 tracks its own offset via add_write_offset.
721716
int err = nghttp3_conn_add_write_offset(*this, data->id, datalen);
722717
if (err != 0) {
723-
session().SetLastError(QuicError::ForApplication(
724-
nghttp3_err_infer_quic_app_error_code(err)));
718+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(err));
725719
returnfalse;
726720
}
727721
// Raw application bytes are committed to the stream's outbound
@@ -1212,10 +1206,10 @@ class Http3ApplicationImpl final : public Session::Application {
12121206
void* conn_user_data,
12131207
void* stream_user_data) {
12141208
NGHTTP3_CALLBACK_SCOPE(app);
1215-
auto& session = app.session();
1216-
Debug(&session, "HTTP/3 application deferred consume %zu bytes", consumed);
1217-
session.ExtendStreamOffset(id, consumed);
1218-
session.ExtendOffset(consumed);
1209+
Debug(&app.session(),
1210+
"HTTP/3 application deferred consume %zu bytes",
1211+
consumed);
1212+
app.session().Consume(id, consumed);
12191213
returnNGTCP2_SUCCESS;
12201214
}
12211215

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Commit e2e2829

Browse files
pimterryaduh95
authored andcommitted
quic: extract transport logic from Application to Session
One small fix notably included: - Check is_destroyed() after StreamCommit, since it calls JS callbacks which could destroy the session. Signed-off-by: Tim Perry <pimterry@gmail.com> PR-URL: #64127 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 7ac97fc commit e2e2829

9 files changed

Lines changed: 608 additions & 562 deletions

File tree

β€Žsrc/quic/README.mdβ€Ž

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ data channels that carry application data.
7777

7878
Every entry point that may generate outbound data creates a
7979
`SendPendingDataScope`. Scopes nest β€” an internal depth counter ensures
80-
`Application::SendPendingData()` is called exactly once, when the outermost
80+
`Session::SendPendingData()` is called exactly once, when the outermost
8181
scope exits:
8282

8383
```cpp
@@ -218,13 +218,13 @@ Session::Receive()
218218

219219
```text
220220
SendPendingDataScope::~SendPendingDataScope()
221-
β†’ Application::SendPendingData()
221+
β†’ Session::SendPendingData()
222222
Loop (up to max_packet_count):
223-
β”œβ”€β”€ GetStreamData() // pull data from next stream
224-
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225-
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
223+
β”œβ”€β”€ application().GetStreamData() // pull data from next stream
224+
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225+
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
226226
β”‚ encrypts, frames, paces
227-
β”œβ”€β”€ if ndatalen > 0: StreamCommit()
227+
β”œβ”€β”€ if ndatalen > 0: application().StreamCommit()
228228
β”‚ stream->Commit(datalen, fin)
229229
β”œβ”€β”€ if nwrite > 0: Send() // uv_udp_send()
230230
β”œβ”€β”€ if WRITE_MORE: continue // room for more in this packet

β€Žsrc/quic/application.ccβ€Ž

Lines changed: 6 additions & 449 deletions
Large diffs are not rendered by default.

β€Žsrc/quic/application.hβ€Ž

Lines changed: 0 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -205,10 +205,6 @@ class Session::Application : public MemoryRetainer {
205205
returnfalse;
206206
}
207207

208-
// Signals to the Application that it should serialize and transmit any
209-
// pending session and stream packets it has accumulated.
210-
voidSendPendingData();
211-
212208
// Returns true if the application protocol supports sending and
213209
// receiving headers on streams (e.g. HTTP/3). Applications that
214210
// do not support headers should return false (the default).
@@ -243,10 +239,6 @@ class Session::Application : public MemoryRetainer {
243239
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
244240
}
245241

246-
// The StreamData struct is used by the application to pass pending stream
247-
// data to the session for transmission.
248-
structStreamData;
249-
250242
virtualintGetStreamData(StreamData* data) = 0;
251243
virtualboolStreamCommit(StreamData* data, size_t datalen) = 0;
252244

@@ -262,57 +254,9 @@ class Session::Application : public MemoryRetainer {
262254
}
263255

264256
private:
265-
Packet::Ptr CreateStreamDataPacket();
266-
267-
// Tries to pack a pending datagram into the current packet buffer.
268-
// If < 0 is returned, either NGTCP2_ERR_WRITE_MORE or a fatal error is
269-
// returned; the caller must check. If > 0 is returned, the packet is done
270-
// and the value is the size of the finalized packet. If 0 is returned,
271-
// the datagram is either congestion limited or was abandoned
272-
ssize_tTryWritePendingDatagram(PathStorage* path,
273-
uint8_t* dest,
274-
size_t destlen,
275-
uint64_t ts);
276-
277-
// Write the given stream_data into the buffer. The PacketInfo out-param
278-
// is populated by ngtcp2 with per-packet metadata (e.g., ECN codepoint)
279-
// that should be applied when sending the packet.
280-
ssize_tWriteVStream(PathStorage* path,
281-
PacketInfo* pi,
282-
uint8_t* buf,
283-
ssize_t* ndatalen,
284-
size_t max_packet_size,
285-
const StreamData& stream_data,
286-
uint64_t ts);
287-
288257
Session* session_ = nullptr;
289258
};
290259

291-
structSession::Application::StreamData final {
292-
// The actual number of vectors in the struct, up to kMaxVectorCount.
293-
size_t count = 0;
294-
// The stream identifier. If this is a negative value then no stream is
295-
// identified.
296-
stream_id id = -1;
297-
int fin = 0;
298-
ngtcp2_vec data[kMaxVectorCount]{};
299-
BaseObjectPtr<Stream> stream;
300-
301-
static_assert(sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
302-
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
303-
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
304-
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
305-
"ngtcp2_vec and nghttp3_vec must have identical layout");
306-
inlineoperator nghttp3_vec*() {
307-
returnreinterpret_cast<nghttp3_vec*>(data);
308-
}
309-
310-
inlineoperatorconst ngtcp2_vec*() const { return data; }
311-
inlineoperator ngtcp2_vec*() { return data; }
312-
313-
std::string ToString() const;
314-
};
315-
316260
// Create a DefaultApplication for the given session.
317261
std::unique_ptr<Session::Application> CreateDefaultApplication(
318262
Session* session, const Session::Application_Options& options);

β€Žsrc/quic/http3.ccβ€Ž

Lines changed: 40 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -209,24 +209,25 @@ class Http3ApplicationImpl final : public Session::Application {
209209
started_ = true;
210210
Debug(&session(), "Starting HTTP/3 application.");
211211

212-
auto params = ngtcp2_conn_get_remote_transport_params(session());
213-
if (params == nullptr) [[unlikely]] {
212+
constauto params = session().remote_transport_params();
213+
if (!params) [[unlikely]] {
214214
// The params are not available yet. Cannot start.
215215
Debug(&session(),
216216
"Cannot start HTTP/3 application yet. No remote transport params");
217217
returnfalse;
218218
}
219219

220-
if (params->initial_max_streams_uni < 3) {
220+
if (params.initial_max_streams_uni() < 3) {
221221
// HTTP3 requires 3 unidirectional control streams to be opened in each
222222
// direction in additional to the bidirectional streams that are used to
223223
// actually carry request and response payload back and forth.
224224
// See:
225225
// https://nghttp2.org/nghttp3/programmers-guide.html#binding-control-streams
226226
Debug(&session(),
227227
"Cannot start HTTP/3 application. Initial max "
228-
"unidirectional streams [%zu] is too low. Must be at least 3",
229-
params->initial_max_streams_uni);
228+
"unidirectional streams [%" PRIu64
229+
"] is too low. Must be at least 3",
230+
params.initial_max_streams_uni());
230231
returnfalse;
231232
}
232233

@@ -235,17 +236,14 @@ class Http3ApplicationImpl final : public Session::Application {
235236
// of requests that the client can actually created.
236237
if (session().is_server()) {
237238
nghttp3_conn_set_max_client_streams_bidi(
238-
*this, params->initial_max_streams_bidi);
239+
*this, params.initial_max_streams_bidi());
239240
}
240241

241242
Debug(&session(), "Creating and binding HTTP/3 control streams");
242243
bool ret =
243-
ngtcp2_conn_open_uni_stream(session(), &control_stream_id_, nullptr) ==
244-
0 &&
245-
ngtcp2_conn_open_uni_stream(
246-
session(), &qpack_enc_stream_id_, nullptr) == 0 &&
247-
ngtcp2_conn_open_uni_stream(
248-
session(), &qpack_dec_stream_id_, nullptr) == 0 &&
244+
session().OpenUnidirectionalStream(&control_stream_id_) &&
245+
session().OpenUnidirectionalStream(&qpack_enc_stream_id_) &&
246+
session().OpenUnidirectionalStream(&qpack_dec_stream_id_) &&
249247
nghttp3_conn_bind_control_stream(*this, control_stream_id_) == 0 &&
250248
nghttp3_conn_bind_qpack_streams(
251249
*this, qpack_enc_stream_id_, qpack_dec_stream_id_) == 0;
@@ -306,8 +304,7 @@ class Http3ApplicationImpl final : public Session::Application {
306304
Debug(&session(),
307305
"Extending stream and connection offset by %zd bytes",
308306
nread);
309-
session().ExtendStreamOffset(id, nread);
310-
session().ExtendOffset(nread);
307+
session().Consume(id, nread);
311308
}
312309

313310
// If this data arrived as 0-RTT, mark the stream. We set it after
@@ -365,24 +362,11 @@ class Http3ApplicationImpl final : public Session::Application {
365362
case EndpointLabel::LOCAL:
366363
return;
367364
case EndpointLabel::REMOTE: {
368-
switch (direction) {
369-
case Direction::BIDIRECTIONAL: {
370-
Debug(&session(),
371-
"HTTP/3 application extending max bidi streams by %" PRIu64,
372-
max_streams);
373-
ngtcp2_conn_extend_max_streams_bidi(
374-
session(), static_cast<size_t>(max_streams));
375-
break;
376-
}
377-
case Direction::UNIDIRECTIONAL: {
378-
Debug(&session(),
379-
"HTTP/3 application extending max uni streams by %" PRIu64,
380-
max_streams);
381-
ngtcp2_conn_extend_max_streams_uni(
382-
session(), static_cast<size_t>(max_streams));
383-
break;
384-
}
385-
}
365+
Debug(&session(),
366+
"HTTP/3 application extending max %s streams by %" PRIu64,
367+
direction == Direction::BIDIRECTIONAL ? "bidi" : "uni",
368+
max_streams);
369+
session().ExtendMaxStreams(direction, max_streams);
386370
}
387371
}
388372
}
@@ -530,8 +514,7 @@ class Http3ApplicationImpl final : public Session::Application {
530514
return;
531515
}
532516

533-
session().SetLastError(
534-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
517+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
535518
session().Close();
536519
}
537520

@@ -548,8 +531,7 @@ class Http3ApplicationImpl final : public Session::Application {
548531
return;
549532
}
550533

551-
session().SetLastError(
552-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
534+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
553535
session().Close();
554536
}
555537

@@ -687,17 +669,30 @@ class Http3ApplicationImpl final : public Session::Application {
687669
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
688670
}
689671

690-
intGetStreamData(StreamData* data) override {
672+
intGetStreamData(Session::StreamData* data) override {
673+
static_assert(
674+
sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
675+
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
676+
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
677+
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
678+
"ngtcp2_vec and nghttp3_vec must have identical layout");
691679
data->count = kMaxVectorCount;
692680
ssize_t ret = 0;
693681
Debug(&session(), "HTTP/3 application getting stream data");
694682
if (conn_ && session().max_data_left()) {
695-
ret = nghttp3_conn_writev_stream(
696-
*this, &data->id, &data->fin, *data, data->count);
683+
// nghttp3 reports fin through an int out-param; bridge it to the bool.
684+
int fin = 0;
685+
ret =
686+
nghttp3_conn_writev_stream(*this,
687+
&data->id,
688+
&fin,
689+
reinterpret_cast<nghttp3_vec*>(data->data),
690+
data->count);
697691
// A negative return value indicates an error.
698692
if (ret < 0) {
699693
returnstatic_cast<int>(ret);
700694
}
695+
data->fin = fin != 0;
701696

702697
data->count = static_cast<size_t>(ret);
703698
if (data->id >= 0 && data->id != control_stream_id_ &&
@@ -710,7 +705,7 @@ class Http3ApplicationImpl final : public Session::Application {
710705
return0;
711706
}
712707

713-
boolStreamCommit(StreamData* data, size_t datalen) override {
708+
boolStreamCommit(Session::StreamData* data, size_t datalen) override {
714709
Debug(&session(),
715710
"HTTP/3 application committing stream %" PRIi64 " data %zu",
716711
data->id,
@@ -720,8 +715,7 @@ class Http3ApplicationImpl final : public Session::Application {
720715
// nghttp3 tracks its own offset via add_write_offset.
721716
int err = nghttp3_conn_add_write_offset(*this, data->id, datalen);
722717
if (err != 0) {
723-
session().SetLastError(QuicError::ForApplication(
724-
nghttp3_err_infer_quic_app_error_code(err)));
718+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(err));
725719
returnfalse;
726720
}
727721
// Raw application bytes are committed to the stream's outbound
@@ -1212,10 +1206,10 @@ class Http3ApplicationImpl final : public Session::Application {
12121206
void* conn_user_data,
12131207
void* stream_user_data) {
12141208
NGHTTP3_CALLBACK_SCOPE(app);
1215-
auto& session = app.session();
1216-
Debug(&session, "HTTP/3 application deferred consume %zu bytes", consumed);
1217-
session.ExtendStreamOffset(id, consumed);
1218-
session.ExtendOffset(consumed);
1209+
Debug(&app.session(),
1210+
"HTTP/3 application deferred consume %zu bytes",
1211+
consumed);
1212+
app.session().Consume(id, consumed);
12191213
returnNGTCP2_SUCCESS;
12201214
}
12211215

0 commit comments

Comments
Β (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Commit e2e2829

Browse files
pimterryaduh95
authored andcommitted
quic: extract transport logic from Application to Session
One small fix notably included: - Check is_destroyed() after StreamCommit, since it calls JS callbacks which could destroy the session. Signed-off-by: Tim Perry <pimterry@gmail.com> PR-URL: #64127 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 7ac97fc commit e2e2829

9 files changed

Lines changed: 608 additions & 562 deletions

File tree

β€Žsrc/quic/README.mdβ€Ž

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ data channels that carry application data.
7777

7878
Every entry point that may generate outbound data creates a
7979
`SendPendingDataScope`. Scopes nest β€” an internal depth counter ensures
80-
`Application::SendPendingData()` is called exactly once, when the outermost
80+
`Session::SendPendingData()` is called exactly once, when the outermost
8181
scope exits:
8282

8383
```cpp
@@ -218,13 +218,13 @@ Session::Receive()
218218

219219
```text
220220
SendPendingDataScope::~SendPendingDataScope()
221-
β†’ Application::SendPendingData()
221+
β†’ Session::SendPendingData()
222222
Loop (up to max_packet_count):
223-
β”œβ”€β”€ GetStreamData() // pull data from next stream
224-
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225-
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
223+
β”œβ”€β”€ application().GetStreamData() // pull data from next stream
224+
β”‚ └── stream->Pull() // bob pull from Outboundβ†’DataQueue
225+
β”œβ”€β”€ WriteVStream() // ngtcp2_conn_writev_stream()
226226
β”‚ encrypts, frames, paces
227-
β”œβ”€β”€ if ndatalen > 0: StreamCommit()
227+
β”œβ”€β”€ if ndatalen > 0: application().StreamCommit()
228228
β”‚ stream->Commit(datalen, fin)
229229
β”œβ”€β”€ if nwrite > 0: Send() // uv_udp_send()
230230
β”œβ”€β”€ if WRITE_MORE: continue // room for more in this packet

β€Žsrc/quic/application.ccβ€Ž

Lines changed: 6 additions & 449 deletions
Large diffs are not rendered by default.

β€Žsrc/quic/application.hβ€Ž

Lines changed: 0 additions & 56 deletions
Original file line numberDiff line numberDiff line change
@@ -205,10 +205,6 @@ class Session::Application : public MemoryRetainer {
205205
returnfalse;
206206
}
207207

208-
// Signals to the Application that it should serialize and transmit any
209-
// pending session and stream packets it has accumulated.
210-
voidSendPendingData();
211-
212208
// Returns true if the application protocol supports sending and
213209
// receiving headers on streams (e.g. HTTP/3). Applications that
214210
// do not support headers should return false (the default).
@@ -243,10 +239,6 @@ class Session::Application : public MemoryRetainer {
243239
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
244240
}
245241

246-
// The StreamData struct is used by the application to pass pending stream
247-
// data to the session for transmission.
248-
structStreamData;
249-
250242
virtualintGetStreamData(StreamData* data) = 0;
251243
virtualboolStreamCommit(StreamData* data, size_t datalen) = 0;
252244

@@ -262,57 +254,9 @@ class Session::Application : public MemoryRetainer {
262254
}
263255

264256
private:
265-
Packet::Ptr CreateStreamDataPacket();
266-
267-
// Tries to pack a pending datagram into the current packet buffer.
268-
// If < 0 is returned, either NGTCP2_ERR_WRITE_MORE or a fatal error is
269-
// returned; the caller must check. If > 0 is returned, the packet is done
270-
// and the value is the size of the finalized packet. If 0 is returned,
271-
// the datagram is either congestion limited or was abandoned
272-
ssize_tTryWritePendingDatagram(PathStorage* path,
273-
uint8_t* dest,
274-
size_t destlen,
275-
uint64_t ts);
276-
277-
// Write the given stream_data into the buffer. The PacketInfo out-param
278-
// is populated by ngtcp2 with per-packet metadata (e.g., ECN codepoint)
279-
// that should be applied when sending the packet.
280-
ssize_tWriteVStream(PathStorage* path,
281-
PacketInfo* pi,
282-
uint8_t* buf,
283-
ssize_t* ndatalen,
284-
size_t max_packet_size,
285-
const StreamData& stream_data,
286-
uint64_t ts);
287-
288257
Session* session_ = nullptr;
289258
};
290259

291-
structSession::Application::StreamData final {
292-
// The actual number of vectors in the struct, up to kMaxVectorCount.
293-
size_t count = 0;
294-
// The stream identifier. If this is a negative value then no stream is
295-
// identified.
296-
stream_id id = -1;
297-
int fin = 0;
298-
ngtcp2_vec data[kMaxVectorCount]{};
299-
BaseObjectPtr<Stream> stream;
300-
301-
static_assert(sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
302-
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
303-
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
304-
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
305-
"ngtcp2_vec and nghttp3_vec must have identical layout");
306-
inlineoperator nghttp3_vec*() {
307-
returnreinterpret_cast<nghttp3_vec*>(data);
308-
}
309-
310-
inlineoperatorconst ngtcp2_vec*() const { return data; }
311-
inlineoperator ngtcp2_vec*() { return data; }
312-
313-
std::string ToString() const;
314-
};
315-
316260
// Create a DefaultApplication for the given session.
317261
std::unique_ptr<Session::Application> CreateDefaultApplication(
318262
Session* session, const Session::Application_Options& options);

β€Žsrc/quic/http3.ccβ€Ž

Lines changed: 40 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -209,24 +209,25 @@ class Http3ApplicationImpl final : public Session::Application {
209209
started_ = true;
210210
Debug(&session(), "Starting HTTP/3 application.");
211211

212-
auto params = ngtcp2_conn_get_remote_transport_params(session());
213-
if (params == nullptr) [[unlikely]] {
212+
constauto params = session().remote_transport_params();
213+
if (!params) [[unlikely]] {
214214
// The params are not available yet. Cannot start.
215215
Debug(&session(),
216216
"Cannot start HTTP/3 application yet. No remote transport params");
217217
returnfalse;
218218
}
219219

220-
if (params->initial_max_streams_uni < 3) {
220+
if (params.initial_max_streams_uni() < 3) {
221221
// HTTP3 requires 3 unidirectional control streams to be opened in each
222222
// direction in additional to the bidirectional streams that are used to
223223
// actually carry request and response payload back and forth.
224224
// See:
225225
// https://nghttp2.org/nghttp3/programmers-guide.html#binding-control-streams
226226
Debug(&session(),
227227
"Cannot start HTTP/3 application. Initial max "
228-
"unidirectional streams [%zu] is too low. Must be at least 3",
229-
params->initial_max_streams_uni);
228+
"unidirectional streams [%" PRIu64
229+
"] is too low. Must be at least 3",
230+
params.initial_max_streams_uni());
230231
returnfalse;
231232
}
232233

@@ -235,17 +236,14 @@ class Http3ApplicationImpl final : public Session::Application {
235236
// of requests that the client can actually created.
236237
if (session().is_server()) {
237238
nghttp3_conn_set_max_client_streams_bidi(
238-
*this, params->initial_max_streams_bidi);
239+
*this, params.initial_max_streams_bidi());
239240
}
240241

241242
Debug(&session(), "Creating and binding HTTP/3 control streams");
242243
bool ret =
243-
ngtcp2_conn_open_uni_stream(session(), &control_stream_id_, nullptr) ==
244-
0 &&
245-
ngtcp2_conn_open_uni_stream(
246-
session(), &qpack_enc_stream_id_, nullptr) == 0 &&
247-
ngtcp2_conn_open_uni_stream(
248-
session(), &qpack_dec_stream_id_, nullptr) == 0 &&
244+
session().OpenUnidirectionalStream(&control_stream_id_) &&
245+
session().OpenUnidirectionalStream(&qpack_enc_stream_id_) &&
246+
session().OpenUnidirectionalStream(&qpack_dec_stream_id_) &&
249247
nghttp3_conn_bind_control_stream(*this, control_stream_id_) == 0 &&
250248
nghttp3_conn_bind_qpack_streams(
251249
*this, qpack_enc_stream_id_, qpack_dec_stream_id_) == 0;
@@ -306,8 +304,7 @@ class Http3ApplicationImpl final : public Session::Application {
306304
Debug(&session(),
307305
"Extending stream and connection offset by %zd bytes",
308306
nread);
309-
session().ExtendStreamOffset(id, nread);
310-
session().ExtendOffset(nread);
307+
session().Consume(id, nread);
311308
}
312309

313310
// If this data arrived as 0-RTT, mark the stream. We set it after
@@ -365,24 +362,11 @@ class Http3ApplicationImpl final : public Session::Application {
365362
case EndpointLabel::LOCAL:
366363
return;
367364
case EndpointLabel::REMOTE: {
368-
switch (direction) {
369-
case Direction::BIDIRECTIONAL: {
370-
Debug(&session(),
371-
"HTTP/3 application extending max bidi streams by %" PRIu64,
372-
max_streams);
373-
ngtcp2_conn_extend_max_streams_bidi(
374-
session(), static_cast<size_t>(max_streams));
375-
break;
376-
}
377-
case Direction::UNIDIRECTIONAL: {
378-
Debug(&session(),
379-
"HTTP/3 application extending max uni streams by %" PRIu64,
380-
max_streams);
381-
ngtcp2_conn_extend_max_streams_uni(
382-
session(), static_cast<size_t>(max_streams));
383-
break;
384-
}
385-
}
365+
Debug(&session(),
366+
"HTTP/3 application extending max %s streams by %" PRIu64,
367+
direction == Direction::BIDIRECTIONAL ? "bidi" : "uni",
368+
max_streams);
369+
session().ExtendMaxStreams(direction, max_streams);
386370
}
387371
}
388372
}
@@ -530,8 +514,7 @@ class Http3ApplicationImpl final : public Session::Application {
530514
return;
531515
}
532516

533-
session().SetLastError(
534-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
517+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
535518
session().Close();
536519
}
537520

@@ -548,8 +531,7 @@ class Http3ApplicationImpl final : public Session::Application {
548531
return;
549532
}
550533

551-
session().SetLastError(
552-
QuicError::ForApplication(nghttp3_err_infer_quic_app_error_code(rv)));
534+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(rv));
553535
session().Close();
554536
}
555537

@@ -687,17 +669,30 @@ class Http3ApplicationImpl final : public Session::Application {
687669
return {StreamPriority::DEFAULT, StreamPriorityFlags::NON_INCREMENTAL};
688670
}
689671

690-
intGetStreamData(StreamData* data) override {
672+
intGetStreamData(Session::StreamData* data) override {
673+
static_assert(
674+
sizeof(ngtcp2_vec) == sizeof(nghttp3_vec) &&
675+
alignof(ngtcp2_vec) == alignof(nghttp3_vec) &&
676+
offsetof(ngtcp2_vec, base) == offsetof(nghttp3_vec, base) &&
677+
offsetof(ngtcp2_vec, len) == offsetof(nghttp3_vec, len),
678+
"ngtcp2_vec and nghttp3_vec must have identical layout");
691679
data->count = kMaxVectorCount;
692680
ssize_t ret = 0;
693681
Debug(&session(), "HTTP/3 application getting stream data");
694682
if (conn_ && session().max_data_left()) {
695-
ret = nghttp3_conn_writev_stream(
696-
*this, &data->id, &data->fin, *data, data->count);
683+
// nghttp3 reports fin through an int out-param; bridge it to the bool.
684+
int fin = 0;
685+
ret =
686+
nghttp3_conn_writev_stream(*this,
687+
&data->id,
688+
&fin,
689+
reinterpret_cast<nghttp3_vec*>(data->data),
690+
data->count);
697691
// A negative return value indicates an error.
698692
if (ret < 0) {
699693
returnstatic_cast<int>(ret);
700694
}
695+
data->fin = fin != 0;
701696

702697
data->count = static_cast<size_t>(ret);
703698
if (data->id >= 0 && data->id != control_stream_id_ &&
@@ -710,7 +705,7 @@ class Http3ApplicationImpl final : public Session::Application {
710705
return0;
711706
}
712707

713-
boolStreamCommit(StreamData* data, size_t datalen) override {
708+
boolStreamCommit(Session::StreamData* data, size_t datalen) override {
714709
Debug(&session(),
715710
"HTTP/3 application committing stream %" PRIi64 " data %zu",
716711
data->id,
@@ -720,8 +715,7 @@ class Http3ApplicationImpl final : public Session::Application {
720715
// nghttp3 tracks its own offset via add_write_offset.
721716
int err = nghttp3_conn_add_write_offset(*this, data->id, datalen);
722717
if (err != 0) {
723-
session().SetLastError(QuicError::ForApplication(
724-
nghttp3_err_infer_quic_app_error_code(err)));
718+
session().SetApplicationError(nghttp3_err_infer_quic_app_error_code(err));
725719
returnfalse;
726720
}
727721
// Raw application bytes are committed to the stream's outbound
@@ -1212,10 +1206,10 @@ class Http3ApplicationImpl final : public Session::Application {
12121206
void* conn_user_data,
12131207
void* stream_user_data) {
12141208
NGHTTP3_CALLBACK_SCOPE(app);
1215-
auto& session = app.session();
1216-
Debug(&session, "HTTP/3 application deferred consume %zu bytes", consumed);
1217-
session.ExtendStreamOffset(id, consumed);
1218-
session.ExtendOffset(consumed);
1209+
Debug(&app.session(),
1210+
"HTTP/3 application deferred consume %zu bytes",
1211+
consumed);
1212+
app.session().Consume(id, consumed);
12191213
returnNGTCP2_SUCCESS;
12201214
}
12211215

0 commit comments

Comments
Β (0)