Skip to content

Commit 6f9040c

Browse files
afrindmeta-codesync[bot]
authored andcommitted
Emit the subgroup End of Group bit
Summary: containsLastInGroup is plumbed from BeginSubgroupOptions through StreamPublisherImpl into the stream type, and the relay forwards it to downstream subscribers, but writeSubgroupHeader recomputed the type with endOfGroup hardcoded to false. The SG_HAS_END_OF_GROUP bit therefore never made it onto the wire and no subscriber could tell which subgroup ends a group. Pass the flag through to writeSubgroupHeader instead. It goes at the end of the parameter list so the existing positional callers, several of which pass another bool, keep their meaning. Reviewed By: sandarsh Differential Revision: D116502269 fbshipit-source-id: 2faebdf29a4899be7d1d919e1112fd80d6f7b354
1 parent 3ab33bf commit 6f9040c

9 files changed

Lines changed: 173 additions & 90 deletions

moxygen/MoQFramer.cpp

Lines changed: 6 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -4350,9 +4350,7 @@ WriteResult MoQFrameWriter::writeSubgroupHeader(
43504350
folly::IOBufQueue& writeBuf,
43514351
TrackAlias trackAlias,
43524352
const ObjectHeader& objectHeader,
4353-
SubgroupIDFormat format,
4354-
bool includeExtensions,
4355-
bool beginsWithFirstObject) const noexcept {
4353+
const SubgroupOptions& options) const noexcept {
43564354
size_t size = 0;
43574355
bool error = false;
43584356

@@ -4364,11 +4362,12 @@ WriteResult MoQFrameWriter::writeSubgroupHeader(
43644362

43654363
auto streamType = getSubgroupStreamType(
43664364
*version_,
4367-
objectHeader.subgroup == 0 ? SubgroupIDFormat::Zero : format,
4368-
includeExtensions,
4369-
/*endOfGroup=*/false,
4365+
objectHeader.subgroup == 0 ? SubgroupIDFormat::Zero
4366+
: options.subgroupIDFormat,
4367+
options.hasExtensions,
4368+
options.hasEndOfGroup,
43704369
priorityPresent,
4371-
beginsWithFirstObject);
4370+
options.beginsWithFirstObject);
43724371
auto streamTypeInt = folly::to_underlying(streamType);
43734372
writeVarint(writeBuf, streamTypeInt, size, error);
43744373
writeVarint(writeBuf, trackAlias.value, size, error);
@@ -4450,32 +4449,6 @@ void MoQFrameWriter::setFetchGroupOrder(GroupOrder groupOrder) noexcept {
44504449
groupOrder == GroupOrder::Default ? GroupOrder::OldestFirst : groupOrder;
44514450
}
44524451

4453-
WriteResult MoQFrameWriter::writeSingleObjectStream(
4454-
folly::IOBufQueue& writeBuf,
4455-
TrackAlias trackAlias,
4456-
const ObjectHeader& objectHeader,
4457-
std::unique_ptr<folly::IOBuf> objectPayload) const noexcept {
4458-
bool hasExtensions = objectHeader.extensions.size() > 0;
4459-
auto res = writeSubgroupHeader(
4460-
writeBuf,
4461-
trackAlias,
4462-
objectHeader,
4463-
objectHeader.subgroup == objectHeader.id ? SubgroupIDFormat::FirstObject
4464-
: SubgroupIDFormat::Present,
4465-
hasExtensions,
4466-
/*beginsWithFirstObject=*/true);
4467-
if (res) {
4468-
return writeStreamObject(
4469-
writeBuf,
4470-
hasExtensions ? StreamType::SUBGROUP_HEADER_SG_EXT
4471-
: StreamType::SUBGROUP_HEADER_SG,
4472-
objectHeader,
4473-
std::move(objectPayload));
4474-
} else {
4475-
return res;
4476-
}
4477-
}
4478-
44794452
void MoQFrameWriter::writeKeyValuePairs(
44804453
folly::IOBufQueue& writeBuf,
44814454
const std::vector<Extension>& extensions,

moxygen/MoQFramer.h

Lines changed: 3 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -566,13 +566,13 @@ WriteResult writeServerSetup(
566566
// are version-agnostic, so we are leaving them out of the MoQFrameWriter.
567567
class MoQFrameWriter {
568568
public:
569+
// options.priorityPresent is ignored; objectHeader.priority decides whether
570+
// the priority byte is written.
569571
WriteResult writeSubgroupHeader(
570572
folly::IOBufQueue& writeBuf,
571573
TrackAlias trackAlias,
572574
const ObjectHeader& objectHeader,
573-
SubgroupIDFormat format = SubgroupIDFormat::Present,
574-
bool includeExtensions = true,
575-
bool beginsWithFirstObject = false) const noexcept;
575+
const SubgroupOptions& options) const noexcept;
576576

577577
WriteResult writeFetchHeader(folly::IOBufQueue& writeBuf, RequestID requestID)
578578
const noexcept;
@@ -605,12 +605,6 @@ class MoQFrameWriter {
605605
std::unique_ptr<folly::IOBuf> objectPayload,
606606
bool forwardingPreferenceIsDatagram = false) const noexcept;
607607

608-
WriteResult writeSingleObjectStream(
609-
folly::IOBufQueue& writeBuf,
610-
TrackAlias trackAlias,
611-
const ObjectHeader& objectHeader,
612-
std::unique_ptr<folly::IOBuf> objectPayload) const noexcept;
613-
614608
WriteResult writeSubscribeRequest(
615609
folly::IOBufQueue& writeBuf,
616610
const SubscribeRequest& subscribeRequest) const noexcept;

moxygen/MoQSession.cpp

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -599,9 +599,11 @@ StreamPublisherImpl::StreamPublisherImpl(
599599
writeBuf_,
600600
trackAlias_,
601601
header_,
602-
format,
603-
includeExtensions,
604-
beginsWithFirstObject);
602+
SubgroupOptions{
603+
.hasExtensions = includeExtensions,
604+
.subgroupIDFormat = format,
605+
.hasEndOfGroup = endOfGroup,
606+
.beginsWithFirstObject = beginsWithFirstObject});
605607
}
606608

607609
// Private methods

moxygen/test/MoQCodecTest.cpp

Lines changed: 29 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -391,7 +391,8 @@ TEST_P(MoQCodecTest, UnderflowObjects) {
391391

392392
TEST_P(MoQCodecTest, ObjectStreamPayloadFin) {
393393
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
394-
moqFrameWriter_.writeSingleObjectStream(
394+
writeSingleObjectStream(
395+
moqFrameWriter_,
395396
writeBuf,
396397
TrackAlias(1),
397398
ObjectHeader(2, 3, 4, 5, 11),
@@ -410,7 +411,8 @@ TEST_P(MoQCodecTest, ObjectStreamPayloadFin) {
410411

411412
TEST_P(MoQCodecTest, ObjectStreamPayload) {
412413
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
413-
moqFrameWriter_.writeSingleObjectStream(
414+
writeSingleObjectStream(
415+
moqFrameWriter_,
414416
writeBuf,
415417
TrackAlias(1),
416418
ObjectHeader(2, 3, 4, 5, 11),
@@ -431,7 +433,8 @@ TEST_P(MoQCodecTest, ObjectStreamPayload) {
431433

432434
TEST_P(MoQCodecTest, EmptyObjectPayload) {
433435
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
434-
moqFrameWriter_.writeSingleObjectStream(
436+
writeSingleObjectStream(
437+
moqFrameWriter_,
435438
writeBuf,
436439
TrackAlias(1),
437440
ObjectHeader(2, 3, 4, 5, ObjectStatus::END_OF_GROUP),
@@ -454,7 +457,7 @@ TEST_P(MoQCodecTest, EmptyObjectPayload) {
454457
TEST_P(MoQCodecTest, TruncatedObject) {
455458
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
456459
auto res = moqFrameWriter_.writeSubgroupHeader(
457-
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5));
460+
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5), SubgroupOptions{});
458461
res = moqFrameWriter_.writeStreamObject(
459462
writeBuf,
460463
StreamType::SUBGROUP_HEADER_SG,
@@ -472,7 +475,10 @@ TEST_P(MoQCodecTest, TruncatedObject) {
472475
TEST_P(MoQCodecTest, TruncatedObjectPayload) {
473476
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
474477
auto res = moqFrameWriter_.writeSubgroupHeader(
475-
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5));
478+
writeBuf,
479+
TrackAlias(1),
480+
ObjectHeader(2, 3, 4, 5),
481+
SubgroupOptions{.hasExtensions = true});
476482
res = moqFrameWriter_.writeStreamObject(
477483
writeBuf,
478484
StreamType::SUBGROUP_HEADER_SG_EXT,
@@ -683,7 +689,7 @@ TEST_P(MoQCodecTest, SubgroupHeaderWithEOF) {
683689

684690
// Write a subgroup header with extensions
685691
auto res = moqFrameWriter_.writeSubgroupHeader(
686-
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5));
692+
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5), SubgroupOptions{});
687693
EXPECT_TRUE(res);
688694

689695
// Expect only onSubgroup and onEndOfStream
@@ -702,9 +708,7 @@ TEST_P(MoQCodecTest, FirstObjectSubgroupOption) {
702708
writeBuf,
703709
TrackAlias(1),
704710
ObjectHeader(2, 3, 4, 5),
705-
SubgroupIDFormat::Present,
706-
/*includeExtensions=*/true,
707-
/*beginsWithFirstObject=*/true);
711+
SubgroupOptions{.beginsWithFirstObject = true});
708712
ASSERT_TRUE(res);
709713

710714
const bool expectedBeginsWithFirstObject =
@@ -732,7 +736,10 @@ TEST_P(MoQCodecTest, FirstObjectSubgroupOption) {
732736
TEST_P(MoQCodecTest, CallbackReturnsErrorTerminateOnObjectBegin) {
733737
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
734738
auto res = moqFrameWriter_.writeSubgroupHeader(
735-
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5));
739+
writeBuf,
740+
TrackAlias(1),
741+
ObjectHeader(2, 3, 4, 5),
742+
SubgroupOptions{.hasExtensions = true});
736743
// First object - will trigger ERROR_TERMINATE
737744
res = moqFrameWriter_.writeStreamObject(
738745
writeBuf,
@@ -778,7 +785,10 @@ TEST_P(MoQCodecTest, CallbackReturnsErrorTerminateOnObjectBegin) {
778785
TEST_P(MoQCodecTest, CallbackReturnsErrorTerminateOnObjectPayload) {
779786
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
780787
auto res = moqFrameWriter_.writeSubgroupHeader(
781-
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5));
788+
writeBuf,
789+
TrackAlias(1),
790+
ObjectHeader(2, 3, 4, 5),
791+
SubgroupOptions{.hasExtensions = true});
782792
res = moqFrameWriter_.writeStreamObject(
783793
writeBuf,
784794
StreamType::SUBGROUP_HEADER_SG_EXT,
@@ -819,7 +829,10 @@ TEST_P(MoQCodecTest, CallbackReturnsErrorTerminateOnObjectPayload) {
819829
TEST_P(MoQCodecTest, CallbackReturnsErrorTerminateOnObjectStatus) {
820830
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
821831
auto res = moqFrameWriter_.writeSubgroupHeader(
822-
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5));
832+
writeBuf,
833+
TrackAlias(1),
834+
ObjectHeader(2, 3, 4, 5),
835+
SubgroupOptions{.hasExtensions = true});
823836
// First object with status - will trigger ERROR_TERMINATE
824837
ObjectHeader statusObj(2, 3, 4, 5);
825838
statusObj.status = ObjectStatus::END_OF_GROUP;
@@ -865,7 +878,7 @@ TEST_P(MoQCodecTest, CallbackReturnsErrorTerminateOnObjectStatus) {
865878
TEST_P(MoQCodecTest, CallbackReturnsErrorTerminateOnSubgroup) {
866879
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
867880
auto res = moqFrameWriter_.writeSubgroupHeader(
868-
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5));
881+
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 4, 5), SubgroupOptions{});
869882
res = moqFrameWriter_.writeStreamObject(
870883
writeBuf,
871884
StreamType::SUBGROUP_HEADER_SG,
@@ -932,7 +945,8 @@ TEST_P(MoQCodecTest, CallbackReturnsErrorTerminateOnFetchHeader) {
932945
// Test that callbacks returning CONTINUE work as expected
933946
TEST_P(MoQCodecTest, CallbackReturnsContinue) {
934947
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
935-
moqFrameWriter_.writeSingleObjectStream(
948+
writeSingleObjectStream(
949+
moqFrameWriter_,
936950
writeBuf,
937951
TrackAlias(1),
938952
ObjectHeader(2, 3, 4, 5, 11),
@@ -963,11 +977,7 @@ TEST_P(MoQCodecTest, CallbackReturnsContinue) {
963977
TEST_P(MoQCodecTest, ZeroLengthObjectFollowedByNormalObject) {
964978
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
965979
auto res = moqFrameWriter_.writeSubgroupHeader(
966-
writeBuf,
967-
TrackAlias(1),
968-
ObjectHeader(2, 3, 0, 5),
969-
SubgroupIDFormat::Present,
970-
false);
980+
writeBuf, TrackAlias(1), ObjectHeader(2, 3, 0, 5), SubgroupOptions{});
971981

972982
// Write a zero-length NORMAL object
973983
ObjectHeader zeroLenObj(0, 0, 4, 0);

moxygen/test/MoQFramerTest.cpp

Lines changed: 46 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -659,14 +659,14 @@ TEST_P(MoQFramerTest, SubgroupObjectHeaderIsNotMarkedAsDatagram) {
659659
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
660660
auto streamType =
661661
getSubgroupStreamType(GetParam(), SubgroupIDFormat::Zero, false, false);
662-
ASSERT_TRUE(writer_
663-
.writeSubgroupHeader(
664-
writeBuf,
665-
TrackAlias(22),
666-
objectHeader,
667-
SubgroupIDFormat::Zero,
668-
false)
669-
.hasValue());
662+
ASSERT_TRUE(
663+
writer_
664+
.writeSubgroupHeader(
665+
writeBuf,
666+
TrackAlias(22),
667+
objectHeader,
668+
SubgroupOptions{.subgroupIDFormat = SubgroupIDFormat::Zero})
669+
.hasValue());
670670
ASSERT_TRUE(writer_
671671
.writeStreamObject(
672672
writeBuf,
@@ -1056,8 +1056,7 @@ TEST_P(MoQFramerTest, ParseStreamHeader) {
10561056
writeBuf,
10571057
TrackAlias(22),
10581058
expectedObjectHeader,
1059-
SubgroupIDFormat::Zero,
1060-
false);
1059+
SubgroupOptions{.subgroupIDFormat = SubgroupIDFormat::Zero});
10611060
EXPECT_TRUE(result.hasValue());
10621061
result = writer_.writeStreamObject(
10631062
writeBuf,
@@ -1396,7 +1395,8 @@ TEST_P(MoQFramerTest, All) {
13961395

13971396
TEST_P(MoQFramerTest, SingleObjectStream) {
13981397
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
1399-
auto result = writer_.writeSingleObjectStream(
1398+
auto result = moxygen::test::writeSingleObjectStream(
1399+
writer_,
14001400
writeBuf,
14011401
TrackAlias(22), // trackAlias
14021402
ObjectHeader(
@@ -1439,6 +1439,31 @@ TEST_P(MoQFramerTest, SingleObjectStream) {
14391439
EXPECT_EQ(parseResult->value.status, ObjectStatus::NORMAL);
14401440
EXPECT_EQ(*parseResult->value.length, 4);
14411441
cursor.skip(*parseResult->value.length);
1442+
1443+
// The same stream ending its group differs only in the type byte.
1444+
MoQFrameWriter eogWriter;
1445+
eogWriter.initializeVersion(GetParam());
1446+
folly::IOBufQueue eogBuf{folly::IOBufQueue::cacheChainLength()};
1447+
EXPECT_TRUE(
1448+
moxygen::test::writeSingleObjectStream(
1449+
eogWriter,
1450+
eogBuf,
1451+
TrackAlias(22),
1452+
ObjectHeader(33, 44, 44, 55, 4),
1453+
folly::IOBuf::copyBuffer("abcd"),
1454+
/*endOfGroup=*/true)
1455+
.hasValue());
1456+
auto eogSerialized = eogBuf.move();
1457+
folly::io::Cursor eogCursor(eogSerialized.get());
1458+
EXPECT_EQ(
1459+
parseStreamType(eogCursor),
1460+
getSubgroupStreamType(
1461+
GetParam(),
1462+
SubgroupIDFormat::FirstObject,
1463+
/*includeExtensions=*/false,
1464+
/*endOfGroup=*/true,
1465+
/*priorityPresent=*/true,
1466+
/*beginsWithFirstObject=*/true));
14421467
}
14431468

14441469
TEST_P(MoQFramerTest, ParseTrackStatus) {
@@ -2054,7 +2079,7 @@ TEST_P(MoQFramerTest, OddExtensionLengthVarintBoundary) {
20542079

20552080
// Write subgroup header (includeExtensions=true) and the stream object
20562081
auto res = writer_.writeSubgroupHeader(
2057-
writeBuf, TrackAlias(1), obj, SubgroupIDFormat::Present, true);
2082+
writeBuf, TrackAlias(1), obj, SubgroupOptions{.hasExtensions = true});
20582083
EXPECT_TRUE(res.hasValue());
20592084
res = writer_.writeStreamObject(
20602085
writeBuf, StreamType::SUBGROUP_HEADER_SG_EXT, obj, nullptr);
@@ -2766,7 +2791,7 @@ TEST_P(MoQFramerTest, SubgroupObjectWithExtensionsAndNonNormalStatus) {
27662791
auto streamType =
27672792
getSubgroupStreamType(GetParam(), SubgroupIDFormat::Present, true, false);
27682793
auto headerResult = writer_.writeSubgroupHeader(
2769-
writeBuf, TrackAlias(1), obj, SubgroupIDFormat::Present, true);
2794+
writeBuf, TrackAlias(1), obj, SubgroupOptions{.hasExtensions = true});
27702795
EXPECT_TRUE(headerResult.hasValue());
27712796

27722797
auto objResult =
@@ -3505,7 +3530,7 @@ void testSubgroupPriorityRoundTrip(
35053530
100, 50, 200, priority, ObjectStatus::NORMAL, noExtensions(), 0};
35063531

35073532
auto result = writer.writeSubgroupHeader(
3508-
writeBuf, TrackAlias(25), objHeader, SubgroupIDFormat::Present, false);
3533+
writeBuf, TrackAlias(25), objHeader, SubgroupOptions{});
35093534
EXPECT_TRUE(result.hasValue());
35103535

35113536
auto serialized = writeBuf.move();
@@ -3552,9 +3577,7 @@ TEST(MoQFramerTest, FirstObjectSubgroupHeaderRoundTripDraft18) {
35523577
writeBuf,
35533578
TrackAlias(25),
35543579
objHeader,
3555-
SubgroupIDFormat::Present,
3556-
/*includeExtensions=*/false,
3557-
/*beginsWithFirstObject=*/true);
3580+
SubgroupOptions{.beginsWithFirstObject = true});
35583581
ASSERT_TRUE(result.hasValue());
35593582

35603583
auto serialized = writeBuf.move();
@@ -3586,9 +3609,7 @@ TEST(MoQFramerTest, FirstObjectSubgroupHeaderIgnoredBeforeDraft18) {
35863609
writeBuf,
35873610
TrackAlias(25),
35883611
objHeader,
3589-
SubgroupIDFormat::Present,
3590-
/*includeExtensions=*/false,
3591-
/*beginsWithFirstObject=*/true);
3612+
SubgroupOptions{.beginsWithFirstObject = true});
35923613
ASSERT_TRUE(result.hasValue());
35933614

35943615
auto serialized = writeBuf.move();
@@ -4637,7 +4658,7 @@ TEST_P(MoQFramerV16PlusTest, ExtensionBlockLengthWithDeltaEncoding) {
46374658
auto streamType =
46384659
getSubgroupStreamType(GetParam(), SubgroupIDFormat::Present, true, false);
46394660
auto res = writer_.writeSubgroupHeader(
4640-
writeBuf, TrackAlias(1), obj, SubgroupIDFormat::Present, true);
4661+
writeBuf, TrackAlias(1), obj, SubgroupOptions{.hasExtensions = true});
46414662
ASSERT_TRUE(res.hasValue());
46424663
res = writer_.writeStreamObject(
46434664
writeBuf, streamType, obj, folly::IOBuf::copyBuffer("AAAA"));
@@ -6558,7 +6579,10 @@ TEST_P(MoQFramerTest, SubgroupObjectUnderflowDoesNotCorruptDeltaState) {
65586579
ObjectHeader obj(1, 0, 0, 128, ObjectStatus::NORMAL, noExtensions(), 4);
65596580
folly::IOBufQueue writeBuf{folly::IOBufQueue::cacheChainLength()};
65606581
writer_.writeSubgroupHeader(
6561-
writeBuf, TrackAlias(1), obj, SubgroupIDFormat::Zero, false);
6582+
writeBuf,
6583+
TrackAlias(1),
6584+
obj,
6585+
SubgroupOptions{.subgroupIDFormat = SubgroupIDFormat::Zero});
65626586
writer_.writeStreamObject(
65636587
writeBuf, streamType, obj, folly::IOBuf::copyBuffer("AAAA"));
65646588
obj.id = 1;

0 commit comments

Comments
 (0)