Skip to content

Commit e096c6b

Browse files
Use weak session in cancel token callback (eclipse-zenoh#2315)
* use weak session in cancellation tokens callback * add test
1 parent 329c17d commit e096c6b

4 files changed

Lines changed: 32 additions & 6 deletions

File tree

zenoh/src/api/builders/liveliness.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -530,9 +530,9 @@ where
530530
.map(|qid| {
531531
#[cfg(feature = "unstable")]
532532
if let Some(cancellation_token) = cancellation_token {
533-
let session_clone = self.session.clone();
533+
let weak_session = self.session.downgrade();
534534
let on_cancel = move || {
535-
let _ = session_clone.0.cancel_liveliness_query(qid); // fails only if no associated query exists - likely because it was already finalized
535+
let _ = weak_session.cancel_liveliness_query(qid); // fails only if no associated query exists - likely because it was already finalized
536536
Ok(())
537537
};
538538
cancellation_token.add_on_cancel_handler(Box::new(on_cancel));

zenoh/src/api/builders/querier.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -514,9 +514,9 @@ where
514514
.map(|qid| {
515515
#[cfg(feature = "unstable")]
516516
if let Some(cancellation_token) = cancellation_token {
517-
let session_clone = self.querier.session.clone();
517+
let weak_session = self.querier.session.clone();
518518
let on_cancel = move || {
519-
let _ = session_clone.cancel_query(qid); // fails only if no associated query exists - likely because it was already finalized
519+
let _ = weak_session.cancel_query(qid); // fails only if no associated query exists - likely because it was already finalized
520520
Ok(())
521521
};
522522
cancellation_token.add_on_cancel_handler(Box::new(on_cancel));

zenoh/src/api/builders/query.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -411,9 +411,9 @@ where
411411
.map(|qid| {
412412
#[cfg(feature = "unstable")]
413413
if let Some(cancellation_token) = cancellation_token {
414-
let session_clone = self.session.clone();
414+
let weak_session = self.session.downgrade();
415415
let on_cancel = move || {
416-
let _ = session_clone.0.cancel_query(qid); // fails only if no associated query exists - likely because it was already finalized
416+
let _ = weak_session.cancel_query(qid); // fails only if no associated query exists - likely because it was already finalized
417417
Ok(())
418418
};
419419
cancellation_token.add_on_cancel_handler(Box::new(on_cancel));

zenoh/tests/cancellation.rs

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -184,3 +184,29 @@ async fn test_cancellation_querier_get() {
184184
let replies = ztimeout!(querier.get().cancellation_token(cancellation_token.clone())).unwrap();
185185
assert!(replies.is_disconnected());
186186
}
187+
188+
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
189+
async fn test_cancellation_does_not_prevent_session_from_close() {
190+
zenoh::init_log_from_env_or("error");
191+
let (session1, session2) = ztimeout!(create_peer_client_pair("tcp/127.0.0.1:50004"));
192+
let cancellation_token = zenoh::cancellation::CancellationToken::default();
193+
194+
let ke = "test/query_cancellation_does_not_prevent_session_from_close";
195+
let _queryable = ztimeout!(session1.declare_queryable(ke)).unwrap();
196+
197+
let querier = ztimeout!(session2.declare_querier(ke)).unwrap();
198+
tokio::time::sleep(Duration::from_secs(1)).await;
199+
200+
let replies = ztimeout!(session2
201+
.get(ke)
202+
.cancellation_token(cancellation_token.clone()))
203+
.unwrap();
204+
205+
let replies2 = ztimeout!(querier.get().cancellation_token(cancellation_token.clone())).unwrap();
206+
207+
std::mem::drop(session2);
208+
assert!(replies.is_disconnected());
209+
assert!(replies2.is_disconnected());
210+
211+
ztimeout!(cancellation_token.cancel()).unwrap();
212+
}

0 commit comments

Comments
 (0)