Skip to content

Commit 446780b

Browse files
afrindmeta-codesync[bot]
authored andcommitted
Add lastByteStreamOffset to WtBufferedStreamData::DequeueResult
Summary: Populate lastByteStreamOffset in dequeue() from the completed PendingWrite's offset field, and add test coverage in WtStreamManager::DeliveryCallback verifying the correct offset is reported for each dequeue scenario. X-link: facebook/proxygen#600 Reviewed By: hanidamlaj Differential Revision: D95707692 Pulled By: afrind fbshipit-source-id: 40ad112bdcc0bf0584c0f00c960e9890ea35cc82
1 parent 50d6b99 commit 446780b

3 files changed

Lines changed: 12 additions & 0 deletions

File tree

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,8 @@ WtBufferedStreamData::DequeueResult WtBufferedStreamData::dequeue(
7676
}
7777

7878
window_.commit(resQueue.chainLength());
79+
res.lastByteStreamOffset =
80+
std::max<uint64_t>(window_.getCurrentOffset(), 1) - 1;
7981
res.data = resQueue.move();
8082
return res;
8183
}

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -106,6 +106,7 @@ class WtBufferedStreamData {
106106
std::unique_ptr<folly::IOBuf> data;
107107
bool fin{false};
108108
WebTransport::ByteEventCallback* deliveryCallback{nullptr};
109+
uint64_t lastByteStreamOffset{0};
109110
};
110111

111112
/**

third-party/proxygen/src/proxygen/lib/http/webtransport/test/WtStreamManagerTest.cpp

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1139,6 +1139,7 @@ TEST(WtStreamManager, DeliveryCallback) {
11391139
EXPECT_EQ(deq1.data->computeChainDataLength(), 100);
11401140
EXPECT_TRUE(deq1.fin);
11411141
EXPECT_EQ(deq1.deliveryCallback, &cb1);
1142+
EXPECT_EQ(deq1.lastByteStreamOffset, 99);
11421143

11431144
// multiple writes with different callbacks
11441145
auto* h2 = CHECK_NOTNULL(streamManager.createEgressHandle());
@@ -1149,15 +1150,18 @@ TEST(WtStreamManager, DeliveryCallback) {
11491150
auto d2a = streamManager.dequeue(*h2, 100);
11501151
EXPECT_EQ(d2a.data->computeChainDataLength(), 100);
11511152
EXPECT_EQ(d2a.deliveryCallback, &cb1);
1153+
EXPECT_EQ(d2a.lastByteStreamOffset, 99);
11521154

11531155
auto d2b = streamManager.dequeue(*h2, 100);
11541156
EXPECT_EQ(d2b.data->computeChainDataLength(), 50);
11551157
EXPECT_EQ(d2b.deliveryCallback, &cb2);
1158+
EXPECT_EQ(d2b.lastByteStreamOffset, 149);
11561159

11571160
auto d2c = streamManager.dequeue(*h2, 100);
11581161
EXPECT_EQ(d2c.data->computeChainDataLength(), 75);
11591162
EXPECT_TRUE(d2c.fin);
11601163
EXPECT_EQ(d2c.deliveryCallback, &cb3);
1164+
EXPECT_EQ(d2c.lastByteStreamOffset, 224);
11611165

11621166
// writes without callbacks coalesce with final callback write
11631167
auto* h3 = CHECK_NOTNULL(streamManager.createEgressHandle());
@@ -1169,6 +1173,7 @@ TEST(WtStreamManager, DeliveryCallback) {
11691173
EXPECT_EQ(deq3.data->computeChainDataLength(), 250);
11701174
EXPECT_TRUE(deq3.fin);
11711175
EXPECT_EQ(deq3.deliveryCallback, &cb4);
1176+
EXPECT_EQ(deq3.lastByteStreamOffset, 249);
11721177

11731178
// callback delivered only when write completes
11741179
auto* h4 = CHECK_NOTNULL(streamManager.createEgressHandle());
@@ -1177,24 +1182,28 @@ TEST(WtStreamManager, DeliveryCallback) {
11771182
auto d4a = streamManager.dequeue(*h4, 60);
11781183
EXPECT_EQ(d4a.data->computeChainDataLength(), 60);
11791184
EXPECT_EQ(d4a.deliveryCallback, nullptr);
1185+
EXPECT_EQ(d4a.lastByteStreamOffset, 59);
11801186

11811187
auto d4b = streamManager.dequeue(*h4, 60);
11821188
EXPECT_EQ(d4b.data->computeChainDataLength(), 40);
11831189
EXPECT_TRUE(d4b.fin);
11841190
EXPECT_EQ(d4b.deliveryCallback, &cb5);
1191+
EXPECT_EQ(d4b.lastByteStreamOffset, 99);
11851192

11861193
// fin-only write with callback
11871194
auto* h5 = CHECK_NOTNULL(streamManager.createEgressHandle());
11881195
h5->writeStreamData(makeBuf(50), /*fin=*/false, nullptr);
11891196
auto d5a = streamManager.dequeue(*h5, 100);
11901197
EXPECT_EQ(d5a.data->computeChainDataLength(), 50);
11911198
EXPECT_EQ(d5a.deliveryCallback, nullptr);
1199+
EXPECT_EQ(d5a.lastByteStreamOffset, 49);
11921200

11931201
h5->writeStreamData(nullptr, /*fin=*/true, &cb1);
11941202
auto d5b = streamManager.dequeue(*h5, 100);
11951203
EXPECT_EQ(d5b.data->computeChainDataLength(), 0);
11961204
EXPECT_TRUE(d5b.fin);
11971205
EXPECT_EQ(d5b.deliveryCallback, &cb1);
1206+
EXPECT_EQ(d5b.lastByteStreamOffset, 49);
11981207
}
11991208

12001209
TEST(WtStreamManager, ByteEventCancellation) {

0 commit comments

Comments
 (0)