Skip to content

Commit 46ee83b

Browse files
afrindmeta-codesync[bot]
authored andcommitted
add onSubscriptionBegin/onSubscriptionEnd gauge stats; fix onPublishDone semantics (#159)
Summary: - onPublishDone fires only when a PUBLISH_DONE wire frame is received/sent, not for session-close synthesized cases - Add onSubscriptionBegin/onSubscriptionEnd to MoQStatsCallback (common base) for an active-subscription gauge that can never go negative; fires on both SUBSCRIBE and PUBLISH paths (publisher and subscriber sides) - onSubscriptionEnd fires at every subscription termination path: subscriber: deliverPublishDoneAndRemove, unsubscribe publisher: sendPublishDone, onUnsubscribe, cleanup - Switch stats test mocks to NiceMock to suppress uninteresting-call warnings Co-authored by: akash-a-n This is an alternative to #158 Pull Request resolved: #159 Reviewed By: sandarsh Differential Revision: D104433947 Pulled By: afrind fbshipit-source-id: cdd86716c2bc9634dd37ad55240978a0224605d1
1 parent 0878324 commit 46ee83b

5 files changed

Lines changed: 80 additions & 4 deletions

File tree

moxygen/MoQSession.cpp

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1958,6 +1958,8 @@ class MoQSession::SubscribeTrackReceiveState
19581958
<< "deliverPublishDoneAndRemove: Delivering PUBLISH_DONE to app; statusCode="
19591959
<< folly::to_underlying(pendingPublishDone_->statusCode)
19601960
<< " alias=" << alias_ << " requestID=" << requestID_;
1961+
MOQ_SUBSCRIBER_STATS(
1962+
session_->subscriberStatsCallback_, onSubscriptionEnd);
19611963
auto token = cancelSource_.getToken();
19621964
auto cb = std::exchange(callback_, nullptr);
19631965
cb->publishDone(std::move(*pendingPublishDone_));
@@ -2195,6 +2197,7 @@ void MoQSession::cleanup() {
21952197
auto requestID = it->first;
21962198
auto pubTrack = std::move(it->second);
21972199
pubTracks_.erase(it);
2200+
MOQ_PUBLISHER_STATS(publisherStatsCallback_, onSubscriptionEnd);
21982201
pubTrack->terminatePublish(
21992202
PublishDone(
22002203
{requestID,
@@ -3658,6 +3661,7 @@ void MoQSession::onUnsubscribe(Unsubscribe unsubscribe) {
36583661
} else {
36593662
trackPublisher->unsubscribe();
36603663
if (pubTracks_.erase(unsubscribe.requestID)) {
3664+
MOQ_PUBLISHER_STATS(publisherStatsCallback_, onSubscriptionEnd);
36613665
retireRequestID(/*signalWriteLoop=*/true);
36623666
} // else, the caller invoked publishDone, which isn't needed but fine
36633667
}
@@ -3693,6 +3697,7 @@ void MoQSession::onPublishOk(PublishOk publishOk) {
36933697
std::static_pointer_cast<TrackPublisherImpl>(trackIt->second);
36943698
trackPublisher->onPublishOk(publishOk);
36953699
pendingPublishTracks_.erase(trackIt->second->fullTrackName());
3700+
MOQ_PUBLISHER_STATS(publisherStatsCallback_, onSubscriptionBegin);
36963701
}
36973702

36983703
publishPtr->setValue(std::move(publishOk));
@@ -4753,6 +4758,7 @@ Subscriber::PublishResult MoQSession::publish(
47534758
void MoQSession::publishOk(const PublishOk& pubOk, ReplyContext& replyContext) {
47544759
XLOG(DBG1) << __func__ << " reqID=" << pubOk.requestID << " sess=" << this;
47554760
MOQ_SUBSCRIBER_STATS(subscriberStatsCallback_, onPublishOk);
4761+
MOQ_SUBSCRIBER_STATS(subscriberStatsCallback_, onSubscriptionBegin);
47564762

47574763
if (logger_) {
47584764
logger_->logPublishOk(pubOk, ControlMessageType::CREATED);
@@ -4875,6 +4881,7 @@ folly::coro::Task<Publisher::SubscribeResult> MoQSession::subscribe(
48754881
co_return folly::makeUnexpected(subscribeResult.error());
48764882
} else {
48774883
MOQ_SUBSCRIBER_STATS(subscriberStatsCallback_, onSubscribeSuccess);
4884+
MOQ_SUBSCRIBER_STATS(subscriberStatsCallback_, onSubscriptionBegin);
48784885
co_return std::make_shared<ReceiverSubscriptionHandle>(
48794886
std::move(subscribeResult.value()), trackAlias, shared_from_this());
48804887
}
@@ -4883,6 +4890,7 @@ folly::coro::Task<Publisher::SubscribeResult> MoQSession::subscribe(
48834890
void MoQSession::sendSubscribeOk(const SubscribeOk& subOk, ReplyContext& ctx) {
48844891
XLOG(DBG1) << __func__ << " sess=" << this;
48854892
MOQ_PUBLISHER_STATS(publisherStatsCallback_, onSubscribeSuccess);
4893+
MOQ_PUBLISHER_STATS(publisherStatsCallback_, onSubscriptionBegin);
48864894
auto res = moqFrameWriter_.writeSubscribeOk(ctx.writeBuf(), subOk);
48874895
if (!res) {
48884896
XLOG(ERR) << "writeSubscribeOk failed sess=" << this;
@@ -4941,6 +4949,9 @@ void MoQSession::unsubscribe(const Unsubscribe& unsubscribe) {
49414949
<< " sess=" << this;
49424950
// cancel() should send STOP_SENDING on any open streams for this
49434951
// subscription
4952+
if (trackIt->second->getSubscribeCallback()) {
4953+
MOQ_SUBSCRIBER_STATS(subscriberStatsCallback_, onSubscriptionEnd);
4954+
}
49444955
trackIt->second->cancel();
49454956
subTracks_.erase(trackIt);
49464957
reqIdToTrackAlias_.erase(trackAliasIt);
@@ -4969,6 +4980,7 @@ void MoQSession::sendPublishDone(const PublishDone& pubDone) {
49694980
<< " sess=" << this;
49704981
return;
49714982
}
4983+
MOQ_PUBLISHER_STATS(publisherStatsCallback_, onSubscriptionEnd);
49724984
auto* ctx = it->second->replyContext();
49734985
SCOPE_EXIT {
49744986
pubTracks_.erase(it);

moxygen/stats/MoQStats.h

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,25 @@ class MoQStatsCallback {
125125
* Subscriber: The publisher closed a subscription stream to us
126126
*/
127127
virtual void onSubscriptionStreamClosed() = 0;
128+
129+
/*
130+
* Publisher: A subscription became active (sent SUBSCRIBE_OK, or had a
131+
* PUBLISH accepted via PUBLISH_OK)
132+
* Subscriber: A subscription became active (received SUBSCRIBE_OK, or sent
133+
* PUBLISH_OK accepting a PUBLISH)
134+
* Pairs with onSubscriptionEnd to form a gauge of active subscriptions.
135+
*/
136+
virtual void onSubscriptionBegin() = 0;
137+
138+
/*
139+
* Publisher: An active subscription ended (sent PUBLISH_DONE, received
140+
* UNSUBSCRIBE, or session closed while subscription was active)
141+
* Subscriber: An active subscription ended (PUBLISH_DONE delivered to app,
142+
* unsubscribe called, or session closed while subscription was active)
143+
* Pairs with onSubscriptionBegin; this callback fires exactly once per
144+
* onSubscriptionBegin, so the difference is a non-negative gauge.
145+
*/
146+
virtual void onSubscriptionEnd() = 0;
128147
};
129148

130149
class MoQPublisherStatsCallback : public MoQStatsCallback {

moxygen/test/MoQSessionSubscribeTests.cpp

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -806,6 +806,10 @@ CO_TEST_P_X(MoQSessionTest, UnsubscribeWithinPublishDone) {
806806
auto subscribeHandler = res.value();
807807

808808
EXPECT_CALL(*clientSubscriberStatsCallback_, onUnsubscribe());
809+
EXPECT_CALL(
810+
*clientSubscriberStatsCallback_,
811+
onPublishDone(PublishDoneStatusCode::TRACK_ENDED));
812+
EXPECT_CALL(*clientSubscriberStatsCallback_, onSubscriptionEnd());
809813

810814
EXPECT_CALL(*subscribeCallback_, publishDone(_))
811815
.WillOnce(testing::Invoke([&](const auto&) {
@@ -1020,6 +1024,8 @@ CO_TEST_P_X(MoQSessionTest, UnsubscribeFromWithinPublishDoneHandler) {
10201024
// When publishDone is delivered, immediately call unsubscribe() from
10211025
// handler
10221026
folly::coro::Baton publishDoneInvoked;
1027+
// onPublishDone fires before this point (wire handler); NiceMock handles it.
1028+
EXPECT_CALL(*clientSubscriberStatsCallback_, onSubscriptionEnd());
10231029
EXPECT_CALL(*subscribeCallback_, publishDone(_))
10241030
.WillOnce(testing::Invoke([&](const auto&) {
10251031
subscribeHandle->unsubscribe();
@@ -1107,3 +1113,26 @@ CO_TEST_P_X(MoQSessionTest, BeginSubgroupCallbackError) {
11071113
co_await publishDone_;
11081114
clientSession_->close(SessionCloseErrorCode::NO_ERROR);
11091115
}
1116+
1117+
// Regression: session close fires onSubscriptionEnd for abandoned publishers;
1118+
// onPublishDone must NOT fire — no PUBLISH_DONE wire frame was ever received.
1119+
CO_TEST_P_X(MoQSessionTest, SubscriptionEndStatFiredOnAbandonedPublisher) {
1120+
co_await setupMoQSession();
1121+
expectSubscribe([](auto sub, auto) -> TaskSubscribeResult {
1122+
// Publisher accepts the subscription but intentionally never calls
1123+
// publishDone, simulating an abandoned / improperly-closed publisher.
1124+
co_return makeSubscribeOkResult(sub);
1125+
});
1126+
1127+
auto res = co_await clientSession_->subscribe(
1128+
getSubscribe(kTestTrackName), subscribeCallback_);
1129+
EXPECT_FALSE(res.hasError());
1130+
1131+
EXPECT_CALL(*clientSubscriberStatsCallback_, onPublishDone(_)).Times(0);
1132+
EXPECT_CALL(*clientSubscriberStatsCallback_, onSubscriptionEnd());
1133+
EXPECT_CALL(*subscribeCallback_, publishDone(_))
1134+
.WillOnce(testing::Return(folly::unit));
1135+
1136+
clientSession_->close(SessionCloseErrorCode::NO_ERROR);
1137+
serverSession_->close(SessionCloseErrorCode::NO_ERROR);
1138+
}

moxygen/test/MoQSessionTestCommon.cpp

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -194,16 +194,20 @@ void MoQSessionTest::SetUp() {
194194
folly::Expected<folly::Unit, MoQPublishError>(folly::unit)));
195195
EXPECT_CALL(*subscribeCallback_, setTrackAlias(_)).Times(testing::AtLeast(0));
196196

197-
clientSubscriberStatsCallback_ = std::make_shared<MockSubscriberStats>();
197+
clientSubscriberStatsCallback_ =
198+
std::make_shared<testing::NiceMock<MockSubscriberStats>>();
198199
clientSession_->setSubscriberStatsCallback(clientSubscriberStatsCallback_);
199200

200-
clientPublisherStatsCallback_ = std::make_shared<MockPublisherStats>();
201+
clientPublisherStatsCallback_ =
202+
std::make_shared<testing::NiceMock<MockPublisherStats>>();
201203
clientSession_->setPublisherStatsCallback(clientPublisherStatsCallback_);
202204

203-
serverSubscriberStatsCallback_ = std::make_shared<MockSubscriberStats>();
205+
serverSubscriberStatsCallback_ =
206+
std::make_shared<testing::NiceMock<MockSubscriberStats>>();
204207
serverSession_->setSubscriberStatsCallback(serverSubscriberStatsCallback_);
205208

206-
serverPublisherStatsCallback_ = std::make_shared<MockPublisherStats>();
209+
serverPublisherStatsCallback_ =
210+
std::make_shared<testing::NiceMock<MockPublisherStats>>();
207211
serverSession_->setPublisherStatsCallback(serverPublisherStatsCallback_);
208212

209213
// For Draft15+, initialize version via ALPN since it's required
@@ -436,7 +440,11 @@ void MoQSessionTest::expectPublishDone(MoQControlCodec::Direction recipient) {
436440
EXPECT_CALL(
437441
*getPublisherStatsCallback(oppositeDirection(recipient)),
438442
onPublishDone(_));
443+
EXPECT_CALL(
444+
*getPublisherStatsCallback(oppositeDirection(recipient)),
445+
onSubscriptionEnd());
439446
EXPECT_CALL(*getSubscriberStatsCallback(recipient), onPublishDone(_));
447+
EXPECT_CALL(*getSubscriberStatsCallback(recipient), onSubscriptionEnd());
440448
EXPECT_CALL(*subscribeCallback_, publishDone(_)).WillOnce([&] {
441449
publishDone_.post();
442450
return folly::unit;

moxygen/test/Mocks.h

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -422,6 +422,10 @@ class MockPublisherStats : public MoQPublisherStatsCallback {
422422

423423
MOCK_METHOD(void, onSubscriptionStreamClosed, (), (override));
424424

425+
MOCK_METHOD(void, onSubscriptionBegin, (), (override));
426+
427+
MOCK_METHOD(void, onSubscriptionEnd, (), (override));
428+
425429
MOCK_METHOD(void, recordPublishNamespaceLatency, (uint64_t), (override));
426430

427431
MOCK_METHOD(void, recordPublishLatency, (uint64_t), (override));
@@ -477,6 +481,10 @@ class MockSubscriberStats : public MoQSubscriberStatsCallback {
477481

478482
MOCK_METHOD(void, onSubscriptionStreamClosed, (), (override));
479483

484+
MOCK_METHOD(void, onSubscriptionBegin, (), (override));
485+
486+
MOCK_METHOD(void, onSubscriptionEnd, (), (override));
487+
480488
MOCK_METHOD(void, recordSubscribeLatency, (uint64_t), (override));
481489

482490
MOCK_METHOD(void, recordFetchLatency, (uint64_t), (override));

0 commit comments

Comments
 (0)