@@ -58,11 +58,12 @@ SubscriberState::SubscriberState(
5858 receiver_(
5959 std::make_shared<ObjectReceiver>(
6060 ObjectReceiver::SUBSCRIBE ,
61- std::shared_ptr<ObjectReceiverCallback>(
62- std::shared_ptr<void >(),
63- &callback_))) {}
61+ callback_)) {}
6462
6563SubscriberState::~SubscriberState () {
64+ // The session may still hold subgroup receivers referencing callback_;
65+ // detach first so any late callback can't touch this destroyed object.
66+ callback_->detach ();
6667 try {
6768 // Unsubscribe if we have a valid subHandle_
6869 if (subHandle_) {
@@ -193,9 +194,12 @@ ObjectReceiverCallback::FlowControlState SubscriberState::Callback::onObject(
193194 std::optional<TrackAlias> /* trackAlias */ ,
194195 const ObjectHeader& objHeader,
195196 Payload payload) {
196- state_.objectsReceived_ ++;
197+ if (!state_) {
198+ return FlowControlState::UNBLOCKED ;
199+ }
200+ state_->objectsReceived_ ++;
197201 if (payload) {
198- state_. bytesReceived_ += payload->computeChainDataLength ();
202+ state_-> bytesReceived_ += payload->computeChainDataLength ();
199203 }
200204
201205 if (auto sendTs =
@@ -205,15 +209,15 @@ ObjectReceiverCallback::FlowControlState SubscriberState::Callback::onObject(
205209 .count ();
206210 if (nowMs >= *sendTs) {
207211 uint64_t latencyMs = nowMs - *sendTs;
208- state_. totalLatencyMs_ += latencyMs;
209- state_. latencyObjects_ ++;
210- state_. testClient_ .recordLatency (latencyMs);
212+ state_-> totalLatencyMs_ += latencyMs;
213+ state_-> latencyObjects_ ++;
214+ state_-> testClient_ .recordLatency (latencyMs);
211215 }
212216 }
213217
214218 // Update largest object seen for track restart detection
215219 AbsoluteLocation location (objHeader.group , objHeader.id );
216- state_. testClient_ .updateLargestObjectSeen (location);
220+ state_-> testClient_ .updateLargestObjectSeen (location);
217221
218222 return FlowControlState::UNBLOCKED ;
219223}
@@ -229,37 +233,46 @@ void SubscriberState::Callback::onEndOfStream() {
229233}
230234
231235void SubscriberState::Callback::onError (ResetStreamErrorCode code) {
232- XLOG (ERR ) << " Subscriber " << state_.id_
236+ if (!state_) {
237+ return ;
238+ }
239+ XLOG (ERR ) << " Subscriber " << state_->id_
233240 << " received stream reset: " << static_cast <uint64_t >(code);
234241 // Don't unsubscribe immediately - let removeSubscriber handle it
235- state_. testClient_ .recordReset ();
242+ state_-> testClient_ .recordReset ();
236243}
237244
238245void SubscriberState::Callback::onPublishDone (PublishDone done) {
239- XLOG (DBG1 ) << " Subscriber " << state_.id_
246+ if (!state_) {
247+ return ;
248+ }
249+ XLOG (DBG1 ) << " Subscriber " << state_->id_
240250 << " received PublishDone - status: "
241251 << static_cast <uint32_t >(done.statusCode )
242252 << " , reason: " << done.reasonPhrase ;
243253
244254 // Library has already cleaned up subscription state; just release our handle
245- state_. subHandle_ .reset ();
255+ state_-> subHandle_ .reset ();
246256
247257 // Only signal completion if track ended naturally
248258 if (done.statusCode == PublishDoneStatusCode::TRACK_ENDED ) {
249- XLOG (DBG1 ) << " Subscriber " << state_. id_ << " - track ended naturally" ;
250- state_. testClient_ .completed ();
259+ XLOG (DBG1 ) << " Subscriber " << state_-> id_ << " - track ended naturally" ;
260+ state_-> testClient_ .completed ();
251261 } else {
252262 // Other status codes (errors, going away, etc.) don't end the test
253- XLOG (DBG1 ) << " Subscriber " << state_. id_
263+ XLOG (DBG1 ) << " Subscriber " << state_-> id_
254264 << " - PublishDone with non-TRACK_ENDED status, not ending test" ;
255265 }
256266}
257267
258268void SubscriberState::Callback::onAllDataReceived () {
259- XLOG (DBG1 ) << " Subscriber " << state_.id_
269+ if (!state_) {
270+ return ;
271+ }
272+ XLOG (DBG1 ) << " Subscriber " << state_->id_
260273 << " - all data received, removing subscriber" ;
261274 // Remove this subscriber from the client's map now that all streams are done
262- state_. testClient_ .removeSubscriber (state_. id_ );
275+ state_-> testClient_ .removeSubscriber (state_-> id_ );
263276}
264277
265278// ============================================================================
0 commit comments