Skip to content

Commit 3ff89af

Browse files
authored
MoQForwarder: restore passive-subscriber guard in duplicate beginSubgroup (#326)
1 parent 0d9f3af commit 3ff89af

2 files changed

Lines changed: 108 additions & 2 deletions

File tree

moxygen/relay/MoQForwarder.cpp

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -630,10 +630,20 @@ MoQForwarder::beginSubgroup(
630630
if (it != sub->subgroups.end() && it->second) {
631631
it->second->reset(ResetStreamErrorCode::CANCELLED);
632632
sub->subgroups.erase(it);
633-
anyReset = true;
633+
// Passive subscribers (e.g. the relay's own top-N/cache observer chain)
634+
// are not real downstream consumers: they never stop_sending, so they
635+
// must not mask the "no active consumers" signal. Reset their stale
636+
// subgroup but do not count them toward anyReset, otherwise a duplicate
637+
// subgroup would never propagate CANCELLED back to the publisher once
638+
// all real consumers have stop_sent.
639+
if (!sub->passive) {
640+
anyReset = true;
641+
}
634642
} else if (
635-
it == sub->subgroups.end() && sub->trackConsumer &&
643+
!sub->passive && it == sub->subgroups.end() && sub->trackConsumer &&
636644
checkRange(*sub) && sub->checkShouldForward()) {
645+
// Passive subscribers can't renew interest either - only a real
646+
// downstream subscriber counts as a reopen candidate.
637647
anyReopenCandidate = true;
638648
}
639649
}

moxygen/relay/test/OpenMOQForwarderTest.cpp

Lines changed: 96 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -284,6 +284,102 @@ TEST_F(OpenMOQForwarderTest, DuplicateChannelSubscriberSameExecIsNoOp) {
284284
EXPECT_EQ(forwarder->numForwardingSubscribers(), 1u);
285285
}
286286

287+
// Test: Once all real subscribers have tombstoned a subgroup, duplicate
288+
// beginSubgroup returns CANCELLED - a passive observer's reset doesn't mask it.
289+
TEST_F(OpenMOQForwarderTest, PassiveResetDoesNotMaskDuplicateSubgroupCancel) {
290+
auto forwarder = std::make_shared<MoQForwarder>(kOpenFwdTestTrackName);
291+
auto realSession = createMockSession();
292+
auto passiveSession = createMockSession();
293+
auto realConsumer = createMockConsumer();
294+
auto passiveConsumer = createMockConsumer();
295+
296+
std::shared_ptr<MockSubgroupConsumer> realSg;
297+
EXPECT_CALL(*realConsumer, beginSubgroup(0, 0, _, _))
298+
.WillOnce(
299+
[this, &realSg](uint64_t, uint64_t, uint8_t, BeginSubgroupOptions) {
300+
realSg = createMockSubgroupConsumer();
301+
// Simulate stop_sending: soft error tombstones the subgroup.
302+
EXPECT_CALL(*realSg, object(0, _, _, false))
303+
.WillOnce(Return(
304+
folly::makeUnexpected(MoQPublishError(
305+
MoQPublishError::CANCELLED, "STOP_SENDING"))));
306+
return folly::makeExpected<
307+
MoQPublishError,
308+
std::shared_ptr<SubgroupConsumer>>(realSg);
309+
});
310+
311+
std::shared_ptr<MockSubgroupConsumer> passiveSg;
312+
EXPECT_CALL(*passiveConsumer, beginSubgroup(0, 0, _, _))
313+
.WillOnce([this, &passiveSg](
314+
uint64_t, uint64_t, uint8_t, BeginSubgroupOptions) {
315+
passiveSg = createMockSubgroupConsumer();
316+
EXPECT_CALL(*passiveSg, object(0, _, _, false))
317+
.WillOnce(Return(folly::unit));
318+
// The duplicate resets the passive observer's stale consumer...
319+
EXPECT_CALL(*passiveSg, reset(ResetStreamErrorCode::CANCELLED));
320+
return folly::
321+
makeExpected<MoQPublishError, std::shared_ptr<SubgroupConsumer>>(
322+
passiveSg);
323+
});
324+
325+
auto realHandle =
326+
addSubscriber(*forwarder, realSession, realConsumer, RequestID(1));
327+
ASSERT_NE(realHandle, nullptr);
328+
auto passiveHandle = forwarder->addSubscriber(
329+
passiveSession, /*forward=*/true, passiveConsumer, /*passive=*/true);
330+
ASSERT_NE(passiveHandle, nullptr);
331+
332+
auto pubSg = forwarder->beginSubgroup(0, 0, 0).value();
333+
EXPECT_TRUE(pubSg->object(0, test::makeBuf(4)).hasValue());
334+
335+
// ...but must not count toward anyReset: CANCELLED propagates.
336+
auto dupRes = forwarder->beginSubgroup(0, 0, 0);
337+
ASSERT_TRUE(dupRes.hasError());
338+
EXPECT_EQ(dupRes.error().code, MoQPublishError::CANCELLED);
339+
}
340+
341+
// Test: A passive subscriber with no subgroup state is not a reopen candidate.
342+
TEST_F(OpenMOQForwarderTest, PassiveDoesNotCountAsReopenCandidate) {
343+
auto forwarder = std::make_shared<MoQForwarder>(kOpenFwdTestTrackName);
344+
auto realSession = createMockSession();
345+
auto passiveSession = createMockSession();
346+
auto realConsumer = createMockConsumer();
347+
auto passiveConsumer = createMockConsumer();
348+
349+
std::shared_ptr<MockSubgroupConsumer> realSg;
350+
EXPECT_CALL(*realConsumer, beginSubgroup(0, 0, _, _))
351+
.WillOnce(
352+
[this, &realSg](uint64_t, uint64_t, uint8_t, BeginSubgroupOptions) {
353+
realSg = createMockSubgroupConsumer();
354+
EXPECT_CALL(*realSg, object(0, _, _, false))
355+
.WillOnce(Return(
356+
folly::makeUnexpected(MoQPublishError(
357+
MoQPublishError::CANCELLED, "STOP_SENDING"))));
358+
return folly::makeExpected<
359+
MoQPublishError,
360+
std::shared_ptr<SubgroupConsumer>>(realSg);
361+
});
362+
// The passive observer joins after the subgroup's objects were delivered
363+
// (object delivery lazily opens subgroups for earlier joiners), so it has
364+
// no subgroup entry; it must not be offered the duplicate subgroup.
365+
EXPECT_CALL(*passiveConsumer, beginSubgroup(_, _, _, _)).Times(0);
366+
367+
auto realHandle =
368+
addSubscriber(*forwarder, realSession, realConsumer, RequestID(1));
369+
ASSERT_NE(realHandle, nullptr);
370+
371+
auto pubSg = forwarder->beginSubgroup(0, 0, 0).value();
372+
EXPECT_TRUE(pubSg->object(0, test::makeBuf(4)).hasValue());
373+
374+
auto passiveHandle = forwarder->addSubscriber(
375+
passiveSession, /*forward=*/true, passiveConsumer, /*passive=*/true);
376+
ASSERT_NE(passiveHandle, nullptr);
377+
378+
auto dupRes = forwarder->beginSubgroup(0, 0, 0);
379+
ASSERT_TRUE(dupRes.hasError());
380+
EXPECT_EQ(dupRes.error().code, MoQPublishError::CANCELLED);
381+
}
382+
287383
TEST_F(OpenMOQForwarderTest, ChannelSubscriberDrainsWhenSubgroupsOpen) {
288384
auto forwarder = std::make_shared<MoQForwarder>(kOpenFwdTestTrackName);
289385
auto consumer = createMockConsumer();

0 commit comments

Comments
 (0)