Skip to content

Commit 05e2011

Browse files
afrindmeta-codesync[bot]
authored andcommitted
Frame datagrams on the way in via frameDatagram
Summary: Every transport framed its egress datagrams in a different place: H3 wrapped them in a `sendDatagram` override, while coro and qmux framed them in their write loops. Because framing happened that late, a datagram too large to ever fit on the wire was accepted by `sendDatagram` and then silently dropped by the write loop. This adds a `frameDatagram` hook on `WtSessionBase` that each transport overrides once, so what gets queued is already on-wire bytes. Coro and qmux check their size limits while framing, so callers now get an error back from `sendDatagram` instead of losing the datagram later. H2 has no datagram egress path yet, so its override has no caller until the capsule transports are wired up. Reviewed By: sharmafb Differential Revision: D115362262 fbshipit-source-id: c96e329f42aa3213ec193c2f1509681a4b229c9f
1 parent 33f24c6 commit 05e2011

12 files changed

Lines changed: 117 additions & 22 deletions

File tree

third-party/proxygen/src/proxygen/lib/http/coro/test/HttpWtUpstreamTests.cpp

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -777,6 +777,18 @@ CO_TEST_P_X(WtTest, MaxStreamsBidiUni) {
777777
wt->closeSession();
778778
}
779779

780+
// A datagram too large to ever fit in a write chunk is rejected up front
781+
// rather than left for the write loop to rediscover on every pass.
782+
CO_TEST_P_X(WtTest, SendDatagramExceedsMaxWriteSize) {
783+
XCHECK(wt);
784+
auto res = wt->sendDatagram(makeBuf(65'535));
785+
EXPECT_TRUE(res.hasError());
786+
co_await rescheduleN(2);
787+
EXPECT_TRUE(wtCodecCb.dgrams.empty());
788+
789+
wt->closeSession();
790+
}
791+
780792
CO_TEST_P_X(WtTest, Datagrams) {
781793
XCHECK(wt);
782794
// tx datagram to peer

third-party/proxygen/src/proxygen/lib/http/coro/util/CoroWtSession.cpp

Lines changed: 26 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -50,11 +50,33 @@ WtExpected<folly::Unit>::Type CoroWtSession::closeSession(
5050

5151
WtExpected<folly::Unit>::Type CoroWtSession::sendDatagram(
5252
IoBufPtr datagram) noexcept {
53-
if (!writeLoopDone_) {
54-
std::ignore = WtSessionBase::sendDatagram(std::move(datagram));
53+
if (writeLoopDone_) {
54+
return folly::unit;
55+
}
56+
auto res = WtSessionBase::sendDatagram(std::move(datagram));
57+
if (res.hasValue()) {
5558
wtSmEgressCb.waitForEvent.signal(); // wake up write loop
5659
}
57-
return folly::unit;
60+
return res;
61+
}
62+
63+
CoroWtSession::IoBufPtr CoroWtSession::frameDatagram(
64+
IoBufPtr datagram) noexcept {
65+
folly::IOBufQueue queue{folly::IOBufQueue::cacheChainLength()};
66+
if (!writeDatagram(
67+
queue,
68+
DatagramCapsule{.httpDatagramPayload = std::move(datagram)},
69+
FrameProtocol::WT_CAPSULE)) {
70+
return nullptr;
71+
}
72+
// A datagram larger than a whole write chunk could never fit, so reject it
73+
// here rather than letting the write loop discover it every pass.
74+
if (queue.chainLength() > kMaxWriteSize) {
75+
XLOG(ERR) << "datagram exceeds max write size; len=" << queue.chainLength()
76+
<< " kMaxWriteSize=" << kMaxWriteSize;
77+
return nullptr;
78+
}
79+
return queue.move();
5880
}
5981

6082
using WtCapsuleCallback = proxygen::detail::WtCapsuleCallback;
@@ -116,7 +138,7 @@ folly::coro::Task<void> CoroWtSession::writeLoop(Ptr self) {
116138
// TODO: currently datagrams can indefinitely preempt stream data
117139
auto datagrams = moveEgressDatagrams();
118140
for (auto& dgram : datagrams) {
119-
writeDatagram(egressBuf, DatagramCapsule{std::move(dgram)});
141+
egressBuf.append(std::move(dgram));
120142
}
121143

122144
auto* wh = sm.nextWritable();

third-party/proxygen/src/proxygen/lib/http/coro/util/CoroWtSession.h

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,8 @@ class CoroWtSession
122122
void start(Ptr self);
123123

124124
private:
125+
IoBufPtr frameDatagram(IoBufPtr datagram) noexcept override;
126+
125127
folly::coro::Task<void> readLoop(Ptr self);
126128
folly::coro::Task<void> writeLoop(Ptr self);
127129

third-party/proxygen/src/proxygen/lib/http/webtransport/QuicWtSession.cpp

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -201,7 +201,11 @@ QuicWtSessionBase::awaitBidiStreamCredit() noexcept {
201201
folly::Expected<folly::Unit, WebTransport::ErrorCode>
202202
QuicWtSessionBase::sendDatagram(IoBufPtr datagram) noexcept {
203203
XCHECK(quicSocket_);
204-
auto writeRes = quicSocket_->writeDatagram(std::move(datagram));
204+
auto framed = frameDatagram(std::move(datagram));
205+
if (!framed) {
206+
return folly::makeUnexpected(ErrorCode::GENERIC_ERROR);
207+
}
208+
auto writeRes = quicSocket_->writeDatagram(std::move(framed));
205209
if (writeRes.hasError()) {
206210
XLOG(ERR) << __func__ << "; err= " << writeRes.error();
207211
return folly::makeUnexpected(ErrorCode::GENERIC_ERROR);
@@ -511,8 +515,7 @@ folly::Expected<folly::Unit, WebTransport::ErrorCode> H3WtSession::closeSession(
511515
return folly::unit;
512516
}
513517

514-
folly::Expected<folly::Unit, WebTransport::ErrorCode> H3WtSession::sendDatagram(
515-
IoBufPtr datagram) noexcept {
518+
H3WtSession::IoBufPtr H3WtSession::frameDatagram(IoBufPtr datagram) noexcept {
516519
/**
517520
* RFC9297:
518521
* HTTP/3 Datagram {
@@ -528,7 +531,7 @@ folly::Expected<folly::Unit, WebTransport::ErrorCode> H3WtSession::sendDatagram(
528531
quarterStreamId, [&appender](auto val) { appender.writeBE(val); });
529532
XCHECK(encodeRes);
530533
queue.append(std::move(datagram));
531-
return QuicWtSessionBase::sendDatagram(queue.move());
534+
return queue.move();
532535
}
533536

534537
folly::Expected<WebTransport::StreamWriteHandle*, WebTransport::ErrorCode>

third-party/proxygen/src/proxygen/lib/http/webtransport/QuicWtSession.h

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -257,13 +257,12 @@ class H3WtSession final : public QuicWtSessionBase {
257257
override;
258258

259259
/**
260-
* Serializes the http/3 datagram, i.e. prefixes the datagram w/ quarter
261-
* stream id (connectStreamId_ / 4) per RFC9297. Similarly to
260+
* Prefixes the datagram w/ quarter stream id (connectStreamId_ / 4) per
261+
* RFC9297. Applies to every sendDatagram overload. Similarly to
262262
* ::create(Uni|Bidi)Stream functions above, ::sendDatagram writes bypass
263263
* the http session and directly write to the QuicSocket.
264264
*/
265-
folly::Expected<folly::Unit, ErrorCode> sendDatagram(
266-
IoBufPtr datagram) noexcept override;
265+
IoBufPtr frameDatagram(IoBufPtr datagram) noexcept override;
267266

268267
/**
269268
* Bidirectionally closes all associated webtransport streams. Notifies the

third-party/proxygen/src/proxygen/lib/http/webtransport/WebTransportSession.cpp

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -124,6 +124,17 @@ H2WtSession::H2WtSession(folly::EventBase* evb,
124124
wtHandler_(wtHandler) {
125125
}
126126

127+
H2WtSession::IoBufPtr H2WtSession::frameDatagram(IoBufPtr datagram) noexcept {
128+
folly::IOBufQueue queue{folly::IOBufQueue::cacheChainLength()};
129+
if (!writeDatagram(
130+
queue,
131+
DatagramCapsule{.httpDatagramPayload = std::move(datagram)},
132+
FrameProtocol::WT_CAPSULE)) {
133+
return nullptr;
134+
}
135+
return queue.move();
136+
}
137+
127138
H2WtSession::~H2WtSession() noexcept {
128139
// abort txn and detach handler if applicable
129140
if (auto* txn = std::exchange(txnHandler_.txn_, nullptr)) {

third-party/proxygen/src/proxygen/lib/http/webtransport/WebTransportSession.h

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -151,6 +151,8 @@ class H2WtSession
151151
WtStreamManager& sm,
152152
WebTransportHandler::Ptr& wtHandler) noexcept;
153153

154+
IoBufPtr frameDatagram(IoBufPtr datagram) noexcept override;
155+
154156
WtStreamManager& sm_;
155157
WebTransportHandler::Ptr& wtHandler_;
156158
};

third-party/proxygen/src/proxygen/lib/http/webtransport/WtUtils.cpp

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -501,7 +501,11 @@ WtExpected<folly::Unit>::Type WtSessionBase::stopSending(
501501

502502
WtExpected<folly::Unit>::Type WtSessionBase::sendDatagram(
503503
IoBufPtr datagram) noexcept {
504-
datagrams_.egress.push_back(std::move(datagram));
504+
auto framed = frameDatagram(std::move(datagram));
505+
if (!framed) {
506+
return folly::makeUnexpected(WebTransport::ErrorCode::GENERIC_ERROR);
507+
}
508+
datagrams_.egress.push_back(std::move(framed));
505509
return folly::unit;
506510
}
507511

third-party/proxygen/src/proxygen/lib/http/webtransport/WtUtils.h

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -225,10 +225,17 @@ class WtSessionBase : public WebTransport {
225225
std::vector<IoBufPtr> moveIngressDatagrams() noexcept {
226226
return std::move(datagrams_.ingress);
227227
}
228+
// Egress datagrams are framed on the way in, so what is queued is already
229+
// on-wire bytes.
228230
std::vector<IoBufPtr> moveEgressDatagrams() noexcept {
229231
return std::move(datagrams_.egress);
230232
}
231233

234+
// Applies this transport's per-datagram framing. Returns null to reject.
235+
virtual IoBufPtr frameDatagram(IoBufPtr datagram) noexcept {
236+
return datagram;
237+
}
238+
232239
private:
233240
folly::Executor* executor_;
234241
WtStreamManager& sm_;

third-party/proxygen/src/proxygen/lib/transport/qmux/QmuxSession.cpp

Lines changed: 28 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -174,11 +174,34 @@ folly::AsyncTransport* QmuxSession::getUnderlyingTransport() const noexcept {
174174

175175
proxygen::detail::WtExpected<folly::Unit>::Type QmuxSession::sendDatagram(
176176
std::unique_ptr<folly::IOBuf> datagram) noexcept {
177-
pendingDatagrams_.emplace_back(std::move(datagram));
177+
auto framed = frameDatagram(std::move(datagram));
178+
if (!framed) {
179+
return folly::makeUnexpected(WebTransport::ErrorCode::GENERIC_ERROR);
180+
}
181+
pendingDatagrams_.emplace_back(std::move(framed));
178182
wtSmEgressCb.waitForEvent.signal();
179183
return folly::unit;
180184
}
181185

186+
QmuxSession::IoBufPtr QmuxSession::frameDatagram(IoBufPtr datagram) noexcept {
187+
folly::IOBufQueue queue{folly::IOBufQueue::cacheChainLength()};
188+
if (!writeDatagram(
189+
queue,
190+
DatagramCapsule{.httpDatagramPayload = std::move(datagram)},
191+
FrameProtocol::QMUX)) {
192+
return nullptr;
193+
}
194+
// A datagram that cannot fit in a single record is unsendable, so reject it
195+
// here rather than letting the write loop discover it and drop it.
196+
if (queue.chainLength() > peerMaxRecordSize_) {
197+
XLOG(ERR)
198+
<< "datagram exceeds peer max_record_size; len=" << queue.chainLength()
199+
<< " peerMaxRecordSize=" << peerMaxRecordSize_;
200+
return nullptr;
201+
}
202+
return queue.move();
203+
}
204+
182205
void QmuxSession::start(Ptr self) {
183206
XLOG(DBG4) << "QmuxSession::start dir=" << (peerAddr_.describe());
184207
if (wtHandler_) {
@@ -366,15 +389,11 @@ folly::coro::Task<void> QmuxSession::writeLoop(Ptr self) {
366389

367390
while (!pendingDatagrams_.empty()) {
368391
folly::IOBufQueue frameBuf{folly::IOBufQueue::cacheChainLength()};
369-
writeDatagram(frameBuf,
370-
DatagramCapsule{.httpDatagramPayload =
371-
std::move(pendingDatagrams_.front())},
372-
FrameProtocol::QMUX);
392+
frameBuf.append(std::move(pendingDatagrams_.front()));
373393
pendingDatagrams_.pop_front();
374-
// Returns false only for a datagram too large for any record, which it
375-
// logs and drops.
376-
std::ignore = appendFrameToRecord(
377-
recordBuf, egressBuf, frameBuf, recordPayloadLimit);
394+
// Cannot fail; frameDatagram rejects anything too large for a record.
395+
XCHECK(appendFrameToRecord(
396+
recordBuf, egressBuf, frameBuf, recordPayloadLimit));
378397
}
379398

380399
// Flush all accumulated QMux records in one transport write.

0 commit comments

Comments
 (0)