@@ -1968,6 +1968,8 @@ class MoQSession::SubscribeTrackReceiveState
19681968 << " deliverPublishDoneAndRemove: Delivering PUBLISH_DONE to app; statusCode="
19691969 << folly::to_underlying (pendingPublishDone_->statusCode )
19701970 << " alias=" << alias_ << " requestID=" << requestID_;
1971+ MOQ_SUBSCRIBER_STATS (
1972+ session_->subscriberStatsCallback_ , onSubscriptionEnd);
19711973 auto token = cancelSource_.getToken ();
19721974 auto cb = std::exchange (callback_, nullptr );
19731975 cb->publishDone (std::move (*pendingPublishDone_));
@@ -2205,6 +2207,7 @@ void MoQSession::cleanup() {
22052207 auto requestID = it->first ;
22062208 auto pubTrack = std::move (it->second );
22072209 pubTracks_.erase (it);
2210+ MOQ_PUBLISHER_STATS (publisherStatsCallback_, onSubscriptionEnd);
22082211 pubTrack->terminatePublish (
22092212 PublishDone (
22102213 {requestID,
@@ -3668,6 +3671,7 @@ void MoQSession::onUnsubscribe(Unsubscribe unsubscribe) {
36683671 } else {
36693672 trackPublisher->unsubscribe ();
36703673 if (pubTracks_.erase (unsubscribe.requestID )) {
3674+ MOQ_PUBLISHER_STATS (publisherStatsCallback_, onSubscriptionEnd);
36713675 retireRequestID (/* signalWriteLoop=*/ true );
36723676 } // else, the caller invoked publishDone, which isn't needed but fine
36733677 }
@@ -3703,6 +3707,7 @@ void MoQSession::onPublishOk(PublishOk publishOk) {
37033707 std::static_pointer_cast<TrackPublisherImpl>(trackIt->second );
37043708 trackPublisher->onPublishOk (publishOk);
37053709 pendingPublishTracks_.erase (trackIt->second ->fullTrackName ());
3710+ MOQ_PUBLISHER_STATS (publisherStatsCallback_, onSubscriptionBegin);
37063711 }
37073712
37083713 publishPtr->setValue (std::move (publishOk));
@@ -4763,6 +4768,7 @@ Subscriber::PublishResult MoQSession::publish(
47634768void MoQSession::publishOk (const PublishOk& pubOk, ReplyContext& replyContext) {
47644769 XLOG (DBG1 ) << __func__ << " reqID=" << pubOk.requestID << " sess=" << this ;
47654770 MOQ_SUBSCRIBER_STATS (subscriberStatsCallback_, onPublishOk);
4771+ MOQ_SUBSCRIBER_STATS (subscriberStatsCallback_, onSubscriptionBegin);
47664772
47674773 if (logger_) {
47684774 logger_->logPublishOk (pubOk, ControlMessageType::CREATED );
@@ -4885,6 +4891,7 @@ folly::coro::Task<Publisher::SubscribeResult> MoQSession::subscribe(
48854891 co_return folly::makeUnexpected (subscribeResult.error ());
48864892 } else {
48874893 MOQ_SUBSCRIBER_STATS (subscriberStatsCallback_, onSubscribeSuccess);
4894+ MOQ_SUBSCRIBER_STATS (subscriberStatsCallback_, onSubscriptionBegin);
48884895 co_return std::make_shared<ReceiverSubscriptionHandle>(
48894896 std::move (subscribeResult.value ()), trackAlias, shared_from_this ());
48904897 }
@@ -4893,6 +4900,7 @@ folly::coro::Task<Publisher::SubscribeResult> MoQSession::subscribe(
48934900void MoQSession::sendSubscribeOk (const SubscribeOk& subOk, ReplyContext& ctx) {
48944901 XLOG (DBG1 ) << __func__ << " sess=" << this ;
48954902 MOQ_PUBLISHER_STATS (publisherStatsCallback_, onSubscribeSuccess);
4903+ MOQ_PUBLISHER_STATS (publisherStatsCallback_, onSubscriptionBegin);
48964904 auto res = moqFrameWriter_.writeSubscribeOk (ctx.writeBuf (), subOk);
48974905 if (!res) {
48984906 XLOG (ERR ) << " writeSubscribeOk failed sess=" << this ;
@@ -4951,6 +4959,9 @@ void MoQSession::unsubscribe(const Unsubscribe& unsubscribe) {
49514959 << " sess=" << this ;
49524960 // cancel() should send STOP_SENDING on any open streams for this
49534961 // subscription
4962+ if (trackIt->second ->getSubscribeCallback ()) {
4963+ MOQ_SUBSCRIBER_STATS (subscriberStatsCallback_, onSubscriptionEnd);
4964+ }
49544965 trackIt->second ->cancel ();
49554966 subTracks_.erase (trackIt);
49564967 reqIdToTrackAlias_.erase (trackAliasIt);
@@ -4979,6 +4990,7 @@ void MoQSession::sendPublishDone(const PublishDone& pubDone) {
49794990 << " sess=" << this ;
49804991 return ;
49814992 }
4993+ MOQ_PUBLISHER_STATS (publisherStatsCallback_, onSubscriptionEnd);
49824994 auto * ctx = it->second ->replyContext ();
49834995 SCOPE_EXIT {
49844996 pubTracks_.erase (it);
0 commit comments