diff --git a/moxygen/MoQSession.cpp b/moxygen/MoQSession.cpp index 52b18ef27..6eac6d717 100644 --- a/moxygen/MoQSession.cpp +++ b/moxygen/MoQSession.cpp @@ -1968,6 +1968,10 @@ class MoQSession::SubscribeTrackReceiveState << "deliverPublishDoneAndRemove: Delivering PUBLISH_DONE to app; statusCode=" << folly::to_underlying(pendingPublishDone_->statusCode) << " alias=" << alias_ << " requestID=" << requestID_; + MOQ_SUBSCRIBER_STATS( + session_->subscriberStatsCallback_, + onPublishDone, + pendingPublishDone_->statusCode); auto token = cancelSource_.getToken(); auto cb = std::exchange(callback_, nullptr); cb->publishDone(std::move(*pendingPublishDone_)); @@ -4078,8 +4082,6 @@ void MoQSession::onPublishDone(PublishDone publishDone) { if (logger_) { logger_->logPublishDone(publishDone, ControlMessageType::PARSED); } - MOQ_SUBSCRIBER_STATS( - subscriberStatsCallback_, onPublishDone, publishDone.statusCode); // Handle regular subscription PUBLISH_DONE auto trackAliasIt = reqIdToTrackAlias_.find(publishDone.requestID); diff --git a/moxygen/test/MoQSessionSubscribeTests.cpp b/moxygen/test/MoQSessionSubscribeTests.cpp index c9b67a139..96922b1f3 100644 --- a/moxygen/test/MoQSessionSubscribeTests.cpp +++ b/moxygen/test/MoQSessionSubscribeTests.cpp @@ -806,6 +806,9 @@ CO_TEST_P_X(MoQSessionTest, UnsubscribeWithinPublishDone) { auto subscribeHandler = res.value(); EXPECT_CALL(*clientSubscriberStatsCallback_, onUnsubscribe()); + EXPECT_CALL( + *clientSubscriberStatsCallback_, + onPublishDone(PublishDoneStatusCode::TRACK_ENDED)); EXPECT_CALL(*subscribeCallback_, publishDone(_)) .WillOnce(testing::Invoke([&](const auto&) { @@ -1020,6 +1023,9 @@ CO_TEST_P_X(MoQSessionTest, UnsubscribeFromWithinPublishDoneHandler) { // When publishDone is delivered, immediately call unsubscribe() from // handler folly::coro::Baton publishDoneInvoked; + EXPECT_CALL( + *clientSubscriberStatsCallback_, + onPublishDone(PublishDoneStatusCode::TRACK_ENDED)); EXPECT_CALL(*subscribeCallback_, publishDone(_)) .WillOnce(testing::Invoke([&](const auto&) { subscribeHandle->unsubscribe(); @@ -1107,3 +1113,31 @@ CO_TEST_P_X(MoQSessionTest, BeginSubgroupCallbackError) { co_await publishDone_; clientSession_->close(SessionCloseErrorCode::NO_ERROR); } + +// Regression test: the onPublishDone subscriber stat must fire when the +// subscription ends due to session closure, even if the publisher never +// sent a PUBLISH_DONE frame (abandoned publisher). +CO_TEST_P_X(MoQSessionTest, OnPublishDoneStatFiredOnAbandonedPublisher) { + co_await setupMoQSession(); + expectSubscribe([](auto sub, auto) -> TaskSubscribeResult { + // Publisher accepts the subscription but intentionally never calls + // publishDone, simulating an abandoned / improperly-closed publisher. + co_return makeSubscribeOkResult(sub); + }); + + auto res = co_await clientSession_->subscribe( + getSubscribe(kTestTrackName), subscribeCallback_); + EXPECT_FALSE(res.hasError()); + + // When the session closes the client synthesizes a SESSION_CLOSED + // publishDone for every active subscription via the subscribeError path. + // The stat must fire here even though no PUBLISH_DONE frame was received. + EXPECT_CALL( + *clientSubscriberStatsCallback_, + onPublishDone(PublishDoneStatusCode::SESSION_CLOSED)); + EXPECT_CALL(*subscribeCallback_, publishDone(_)) + .WillOnce(testing::Return(folly::unit)); + + clientSession_->close(SessionCloseErrorCode::NO_ERROR); + serverSession_->close(SessionCloseErrorCode::NO_ERROR); +}