Skip to content

Commit c8a1eab

Browse files
sandarshmeta-codesync[bot]
authored andcommitted
Fix one subgroup per object validation
Summary: Fixes the MoQTest client validation for ONE_SUBGROUP_PER_OBJECT forwarding preference, which was incorrectly expecting objects to arrive in order. Since each object is sent on its own stream, they can arrive out-of-order. Reviewed By: sharmafb Differential Revision: D89417533 fbshipit-source-id: 4c01d43df0b8f05373fd59f794adc75d8a5c9b76
1 parent d262107 commit c8a1eab

4 files changed

Lines changed: 139 additions & 27 deletions

File tree

moxygen/moqtest/MoQTestClient.cpp

Lines changed: 122 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -263,6 +263,14 @@ void MoQTestClient::onObjectStatus(
263263
return;
264264
}
265265

266+
// Remove the end-of-group marker from the scoreboard
267+
// End-of-group markers don't go through validateSubscribedData(), so we need
268+
// to erase them here to avoid false "objects still expected" errors
269+
if (params_.forwardingPreference != ForwardingPreference::DATAGRAM) {
270+
auto key = std::make_pair(header.group, header.id);
271+
expectedObjects_.erase(key);
272+
}
273+
266274
// Adjust the expected data
267275
if (adjustExpected(params_, &objHeader) ==
268276
AdjustedExpectedResult::RECEIVED_ALL_DATA) {
@@ -285,22 +293,42 @@ void MoQTestClient::onAllDataReceived() {
285293
auto subHandleResetGuard = folly::makeGuard([this] { subHandle_.reset(); });
286294

287295
if (params_.forwardingPreference == ForwardingPreference::DATAGRAM) {
296+
// For datagrams, some drops are allowed based on datagramDropPercentage
297+
uint64_t totalExpected = (((params_.lastGroupInTrack - params_.startGroup) /
298+
params_.groupIncrement) +
299+
1) *
300+
(((params_.lastObjectInTrack - params_.startObject) /
301+
params_.objectIncrement) +
302+
1);
303+
// Allow configured percentage of drops, with minimum of 1
304+
uint64_t dropsAllowed = std::max(
305+
uint64_t{1}, totalExpected * params_.datagramDropPercentage / 100);
288306
if (datagramObjects_ == 0) {
289307
XLOG(ERR)
290308
<< "MoQTest verification result: FAILURE! reason: Datagram Failed - 0 Objects Recieved";
291309
subHandle_->unsubscribe();
292310
return;
311+
} else if (expectedObjects_.size() > dropsAllowed) {
312+
XLOG(ERR)
313+
<< "MoQTest verification result: FAILURE! reason: Datagram had too many drops: "
314+
<< expectedObjects_.size() << " missing, allowed " << dropsAllowed;
315+
subHandle_->unsubscribe();
316+
return;
293317
} else {
294318
XLOG(INFO) << "MoQTest verification result: SUCCESS! Datagram Recieved "
295319
<< datagramObjects_ << " objects";
296320
return;
297321
}
298322
}
299-
if (params_.forwardingPreference != ForwardingPreference::DATAGRAM &&
300-
adjustExpected(params_, nullptr) ==
301-
AdjustedExpectedResult::STILL_RECEIVING_DATA) {
323+
324+
// For non-datagram: success == scoreboard.empty()
325+
if (!expectedObjects_.empty()) {
302326
XLOG(ERR)
303-
<< "MoQTest verification result: FAILURE! reason: SubscribeDone recieved while objects are still expected";
327+
<< "MoQTest verification result: FAILURE! reason: SubscribeDone recieved while "
328+
<< expectedObjects_.size() << " objects are still expected";
329+
for (const auto& [group, objId] : expectedObjects_) {
330+
XLOG(ERR) << " Missing object: group=" << group << " id=" << objId;
331+
}
304332
subHandle_->unsubscribe();
305333
return;
306334
}
@@ -319,12 +347,32 @@ bool MoQTestClient::validateSubscribedData(
319347
XLOG(DBG1) << "MoQTest DEBUGGING: Object Group=" << header.group
320348
<< " end of group markers=" << params_.sendEndOfGroupMarkers
321349
<< " expected end of group markers=" << expectEndOfGroup_;
322-
if (params_.forwardingPreference != ForwardingPreference::DATAGRAM &&
323-
header.group != expectedGroup_) {
324-
XLOG(ERR)
325-
<< "MoQTest verification result: FAILURE! reason: Group Mismatch: Actual="
326-
<< header.group << " Expected=" << expectedGroup_;
327-
return false;
350+
if (params_.forwardingPreference != ForwardingPreference::DATAGRAM) {
351+
if (params_.forwardingPreference ==
352+
ForwardingPreference::ONE_SUBGROUP_PER_OBJECT) {
353+
// Allow out-of-order groups, just validate range
354+
if (header.group < params_.startGroup ||
355+
header.group > params_.lastGroupInTrack) {
356+
XLOG(ERR)
357+
<< "MoQTest verification result: FAILURE! reason: Group out of range: "
358+
<< header.group << " not in [" << params_.startGroup << ", "
359+
<< params_.lastGroupInTrack << "]";
360+
return false;
361+
}
362+
// Validate group increment
363+
if ((header.group - params_.startGroup) % params_.groupIncrement != 0) {
364+
XLOG(ERR)
365+
<< "MoQTest verification result: FAILURE! reason: Group not on increment boundary: "
366+
<< header.group << " with startGroup=" << params_.startGroup
367+
<< " and groupIncrement=" << params_.groupIncrement;
368+
return false;
369+
}
370+
} else if (header.group != expectedGroup_) {
371+
XLOG(ERR)
372+
<< "MoQTest verification result: FAILURE! reason: Group Mismatch: Actual="
373+
<< header.group << " Expected=" << expectedGroup_;
374+
return false;
375+
}
328376
}
329377

330378
if (params_.forwardingPreference ==
@@ -363,6 +411,43 @@ bool MoQTestClient::validateSubscribedData(
363411
return false;
364412
}
365413

414+
// Validate ONE_SUBGROUP_PER_OBJECT constraints
415+
if (params_.forwardingPreference ==
416+
ForwardingPreference::ONE_SUBGROUP_PER_OBJECT) {
417+
// Subgroup must equal object ID
418+
if (header.subgroup != header.id) {
419+
XLOG(ERR)
420+
<< "MoQTest verification result: FAILURE! reason: SubGroup must equal "
421+
<< "Object ID for ONE_SUBGROUP_PER_OBJECT: subgroup="
422+
<< header.subgroup << " id=" << header.id;
423+
return false;
424+
}
425+
// Object ID must be in valid range
426+
if (header.id < params_.startObject ||
427+
header.id > params_.lastObjectInTrack) {
428+
XLOG(ERR)
429+
<< "MoQTest verification result: FAILURE! reason: Object ID out of range: "
430+
<< header.id << " not in [" << params_.startObject << ", "
431+
<< params_.lastObjectInTrack << "]";
432+
return false;
433+
}
434+
}
435+
436+
// Scoreboard-based duplicate detection for non-DATAGRAM forwarding
437+
// preferences
438+
if (params_.forwardingPreference != ForwardingPreference::DATAGRAM) {
439+
auto key = std::make_pair(header.group, header.id);
440+
auto it = expectedObjects_.find(key);
441+
if (it == expectedObjects_.end()) {
442+
XLOG(ERR)
443+
<< "MoQTest verification result: FAILURE! reason: Duplicate or unexpected object: "
444+
<< "group=" << header.group << " id=" << header.id;
445+
return false;
446+
}
447+
// Erase from scoreboard - object received successfully
448+
expectedObjects_.erase(it);
449+
}
450+
366451
if (params_.forwardingPreference != ForwardingPreference::DATAGRAM &&
367452
params_.forwardingPreference !=
368453
ForwardingPreference::ONE_SUBGROUP_PER_OBJECT &&
@@ -423,20 +508,10 @@ AdjustedExpectedResult MoQTestClient::adjustExpectedForOneSubgroupPerGroup(
423508
return AdjustedExpectedResult::STILL_RECEIVING_DATA;
424509
}
425510

426-
AdjustedExpectedResult MoQTestClient::adjustExpectedForOneSubgroupPerObject(
427-
MoQTestParameters& params) {
428-
// Adjust Expected Group, ObjectId and Subgroup
429-
if (expectedGroup_ < params.lastGroupInTrack &&
430-
subgroupToExpectedObjId_[0] == params.lastObjectInTrack) {
431-
// Increment Group, Reset ObjectId and Subgroup
432-
expectedGroup_ += params.groupIncrement;
433-
subgroupToExpectedObjId_[0] = params.startObject;
434-
expectedSubgroup_ = 0;
435-
} else if (subgroupToExpectedObjId_[0] < params.lastObjectInTrack) {
436-
// Increment ObjectId and Subgroup
437-
subgroupToExpectedObjId_[0] += params.objectIncrement;
438-
expectedSubgroup_ += params.objectIncrement;
439-
} else {
511+
AdjustedExpectedResult MoQTestClient::adjustExpectedForOneSubgroupPerObject() {
512+
// With scoreboard approach, we check if all expected objects have been
513+
// received
514+
if (expectedObjects_.empty()) {
440515
return AdjustedExpectedResult::RECEIVED_ALL_DATA;
441516
}
442517
return AdjustedExpectedResult::STILL_RECEIVING_DATA;
@@ -552,6 +627,16 @@ void MoQTestClient::initializeExpecteds(MoQTestParameters& params) {
552627
expectedSubgroup_ = 0;
553628
expectEndOfGroup_ = params.sendEndOfGroupMarkers;
554629

630+
// Initialize scoreboard with all expected (group, objectId) pairs
631+
expectedObjects_.clear();
632+
for (uint64_t group = params.startGroup; group <= params.lastGroupInTrack;
633+
group += params.groupIncrement) {
634+
for (uint64_t obj = params.startObject; obj <= params.lastObjectInTrack;
635+
obj += params.objectIncrement) {
636+
expectedObjects_.insert({group, obj});
637+
}
638+
}
639+
555640
// Only relevant for Datagram Forwarding Preference
556641
datagramObjects_ = 0;
557642
}
@@ -564,7 +649,7 @@ AdjustedExpectedResult MoQTestClient::adjustExpected(
564649
return adjustExpectedForOneSubgroupPerGroup(params);
565650
}
566651
case (ForwardingPreference::ONE_SUBGROUP_PER_OBJECT): {
567-
return adjustExpectedForOneSubgroupPerObject(params);
652+
return adjustExpectedForOneSubgroupPerObject();
568653
}
569654
case (ForwardingPreference::TWO_SUBGROUPS_PER_GROUP): {
570655
return adjustExpectedForTwoSubgroupsPerGroup(header, params);
@@ -622,6 +707,18 @@ bool MoQTestClient::validateDatagramObjects(const ObjectHeader& header) {
622707
return false;
623708
}
624709

710+
// Check for duplicates - if not in scoreboard, we already received this
711+
// object
712+
auto key = std::make_pair(header.group, header.id);
713+
auto it = expectedObjects_.find(key);
714+
if (it == expectedObjects_.end()) {
715+
XLOG(ERR)
716+
<< "MoQTest verification result: FAILURE! reason: Duplicate datagram object: "
717+
<< "group=" << header.group << " id=" << header.id;
718+
return false;
719+
}
720+
expectedObjects_.erase(it);
721+
625722
return true;
626723
}
627724

moxygen/moqtest/MoQTestClient.h

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -138,6 +138,11 @@ class MoQTestClient {
138138
uint64_t expectedSubgroup_{};
139139
std::array<uint64_t, 2> subgroupToExpectedObjId_{};
140140

141+
// Scoreboard of expected (group, objectId) pairs
142+
// When receiving: if present, erase; if absent, it's a duplicate
143+
// At end: success == scoreboard.empty() (or within drop limit for datagrams)
144+
std::set<std::pair<uint64_t, uint64_t>> expectedObjects_;
145+
141146
// Holds if current request expects end of group markers
142147
bool expectEndOfGroup_{};
143148

@@ -163,8 +168,7 @@ class MoQTestClient {
163168
const ObjectHeader* header);
164169
AdjustedExpectedResult adjustExpectedForOneSubgroupPerGroup(
165170
MoQTestParameters& params);
166-
AdjustedExpectedResult adjustExpectedForOneSubgroupPerObject(
167-
MoQTestParameters& params);
171+
AdjustedExpectedResult adjustExpectedForOneSubgroupPerObject();
168172
AdjustedExpectedResult adjustExpectedForTwoSubgroupsPerGroup(
169173
const ObjectHeader* header,
170174
MoQTestParameters& params);

moxygen/moqtest/MoQTestClientMain.cpp

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,10 @@ DEFINE_bool(
7373
quic_transport,
7474
false,
7575
"Use QUIC transport instead of WebTransport");
76+
DEFINE_uint64(
77+
datagram_drops_allowed_percentage,
78+
moxygen::kDefaultDatagramDropPercentage,
79+
"Allowed datagram drop percentage for DATAGRAM forwarding (default 1%)");
7680

7781
int main(int argc, char** argv) {
7882
gflags::ParseCommandLineFlags(&argc, &argv, false);
@@ -100,6 +104,8 @@ int main(int argc, char** argv) {
100104
defaultMoqParams.testVariableExtension = FLAGS_test_variable_extension;
101105
defaultMoqParams.publisherDeliveryTimeout = FLAGS_publisher_delivery_timeout;
102106
defaultMoqParams.deliveryTimeout = FLAGS_delivery_timeout;
107+
defaultMoqParams.datagramDropPercentage =
108+
FLAGS_datagram_drops_allowed_percentage;
103109
defaultMoqParams.lastObjectInTrack =
104110
FLAGS_last_object_in_track == moxygen::kLocationMax.object
105111
? FLAGS_object_increment *

moxygen/moqtest/Types.h

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ const uint64_t kDefaultLastGroupInTrack = (1ULL << 62) - 1;
2626
const uint64_t kDefaultStart = 0;
2727
const uint64_t kDefaultIncrement = 1;
2828
const uint64_t kDefaultPublisherDeliveryTimeout = 0;
29+
const uint64_t kDefaultDatagramDropPercentage = 1; // Allow up to 1% drops
2930

3031
struct MoQTestParameters {
3132
ForwardingPreference forwardingPreference =
@@ -49,6 +50,10 @@ struct MoQTestParameters {
4950
uint64_t deliveryTimeout =
5051
0; // Tuple Field 16 - Delivery timeout in milliseconds (0 = disabled)
5152

53+
// Client-side parameter (not part of track namespace)
54+
uint64_t datagramDropPercentage =
55+
kDefaultDatagramDropPercentage; // Allowed datagram drop percentage
56+
5257
uint64_t lastObjectInTrack = this->objectsPerGroup +
5358
this->sendEndOfGroupMarkers; // Tuple Field 5 (Out of Order to Ensure
5459
// this->sendEndOfGroupMarkers is

0 commit comments

Comments
 (0)