From faee5a32f752f77fa7dc6200004c7c1f5cec2429 Mon Sep 17 00:00:00 2001 From: yuanyuyuan Date: Tue, 28 Jul 2026 02:34:52 +0800 Subject: [PATCH 1/4] fix(rmw): stop invoking executor callbacks under a lock All five remaining callback-under-lock sites in rmw-zenoh-rs invoked an rclcpp executor callback with one or two mutex guards still live: rmw_subscription_set_on_new_message_callback callback + unread_count subscription delivery notifier callback + callback_user_data service delivery notifier callback + callback_user_data client delivery notifier (send_request) callback + callback_user_data rmw_{service,client}_set_on_new_*_callback unread_count `std::sync::Mutex` is not reentrant, so a callback that re-enters the rmw API for the same entity blocks on a lock its own thread already holds. Installing and clearing on_new_message callbacks is how rclcpp executors attach and detach, so this is reachable, and it needs no race: the worst of the five is hit by the ordinary startup path where messages arrive before the executor attaches and the backlog is replayed under two guards. Replace the per-entity trio of mutexes with one `ExecCallback` slot whose two operations collect what they need under the lock, drop the guard, and only then call into user code. Collapsing three mutexes into one is part of the fix rather than tidying: with three, notification had to hold two guards at once to read a callback and its user-data together, and correctness rested on every site agreeing on a lock order. The slot uses hiroz's `TrackedMutex` and dispatches through `invoke_user_callback!`, so a reintroduction panics in debug with the site name instead of hanging. This crate had no coverage for this behaviour at all. Eight unit tests are added with it, including `no_guard_is_live_when_the_callback_runs` (asserts the live guard count is zero at the moment of dispatch), `delivery_notification_survives_reentry_into_itself` and `installing_a_callback_over_a_backlog_survives_reentry`, which reproduce the self-deadlock directly. Reverting the guard-drop and keeping them turns each into a hang. These five were not found by the manual sweep that produced the earlier fixes in this series. They were found by a mechanical pass over every lock acquisition site in the workspace: enumerate acquisitions, classify each guard's lifetime, then look inside that lifetime for a call into user code, plus a lock-order graph for ABBA inversions. --- crates/rmw-zenoh-rs/src/exec_callback.rs | 391 +++++++++++++++++++++++ crates/rmw-zenoh-rs/src/lib.rs | 1 + crates/rmw-zenoh-rs/src/pubsub.rs | 5 +- crates/rmw-zenoh-rs/src/rmw.rs | 163 +--------- crates/rmw-zenoh-rs/src/service.rs | 28 +- 5 files changed, 413 insertions(+), 175 deletions(-) create mode 100644 crates/rmw-zenoh-rs/src/exec_callback.rs diff --git a/crates/rmw-zenoh-rs/src/exec_callback.rs b/crates/rmw-zenoh-rs/src/exec_callback.rs new file mode 100644 index 000000000..8fcf7068c --- /dev/null +++ b/crates/rmw-zenoh-rs/src/exec_callback.rs @@ -0,0 +1,391 @@ +//! The executor-notification slot shared by every rmw entity that can wake an +//! rclcpp executor: subscriptions, services, and clients. +//! +//! # The defect this type exists to prevent +//! +//! Each of those three entities used to carry the same trio of independent +//! mutexes — `callback`, `callback_user_data`, `unread_count` — and every site +//! that notified the executor locked one or two of them *and then invoked the +//! callback while still holding the guard*: +//! +//! ```ignore +//! if let Ok(cb) = callback_holder.lock() { +//! if let Some(callback_fn) = *cb { +//! if let Ok(user_data) = user_data_holder.lock() { +//! unsafe { callback_fn(user_data_ptr, 1) }; // <-- executor, under two guards +//! ``` +//! +//! `std::sync::Mutex` is not reentrant. A callback that re-enters the rmw API +//! for the same entity — which rclcpp does, because installing and clearing +//! `on_new_message` callbacks is how executors attach and detach — blocks on a +//! lock its own thread already holds. That is a deterministic hang, not a race. +//! +//! # The shape of the fix +//! +//! One mutex, and every operation *collects what it needs under the lock, drops +//! the guard, and only then calls into user code*. This is the same pattern the +//! hiroz core fixes established (`GraphEventManager::trigger_graph_change`, +//! `ParameterState::validate_and_apply`, `update_shared_event_status`). +//! +//! Collapsing three mutexes into one is part of the fix, not incidental +//! tidying: with three, "notify" had to hold two guards at once to read a +//! callback and its user-data together, and the correctness of the whole thing +//! rested on every site agreeing on a lock order. With one, the invariant is +//! local and there is no order to get wrong. +//! +//! # Enforcement +//! +//! The mutex is a [`hiroz::reentrancy::TrackedMutex`], so its guards are counted +//! on the current thread, and both call sites go through +//! [`hiroz::invoke_user_callback!`], which asserts the count is zero before +//! dispatching. In debug builds a reintroduction of the defect panics with the +//! site name instead of hanging; in release both the counter and the assertion +//! compile to nothing. + +use std::sync::Arc; + +use hiroz::invoke_user_callback; +use hiroz::reentrancy::TrackedMutex; + +/// The one callback signature rmw uses for all three "new item arrived" +/// notifications. `rmw_subscription_new_message_callback_t`, +/// `rmw_service_new_request_callback_t` and `rmw_client_new_response_callback_t` +/// are all aliases of `Option`. +pub type ExecCallbackFn = unsafe extern "C" fn(user_data: *const std::ffi::c_void, count: usize); + +#[derive(Default)] +struct State { + /// The executor callback, or `None` when no executor is attached. + callback: Option, + /// The executor's opaque pointer, held as a `usize` so `State` stays `Send`. + user_data: usize, + /// Items that arrived while `callback` was `None`, replayed on install. + unread: usize, +} + +/// Shared executor-notification state for one rmw entity. +/// +/// Cloning shares the underlying slot; the delivery thread and the rmw API +/// entry points hold clones of the same `ExecCallback`. +#[derive(Clone)] +pub struct ExecCallback { + state: Arc>, + /// Entity kind, reproduced in the re-entrancy panic message. + site: &'static str, +} + +impl ExecCallback { + /// `site` names the entity kind ("subscription", "service", "client") and + /// appears in the re-entrancy assertion message. + pub fn new(site: &'static str) -> Self { + Self { + state: Arc::new(TrackedMutex::new(State::default())), + site, + } + } + + /// One new item arrived. Invokes the executor callback if one is installed, + /// otherwise records the item as unread so it can be replayed by [`set`]. + /// + /// Runs on the zenoh delivery thread. + /// + /// [`set`]: Self::set + pub fn notify_one(&self) { + // Collect under the lock... + let armed = { + let Ok(mut state) = self.state.lock() else { + return; + }; + match state.callback { + Some(callback_fn) => { + Some((callback_fn, state.user_data as *const std::ffi::c_void)) + } + None => { + state.unread += 1; + None + } + } + }; + // ...guard is dropped, and only now do we call into the executor. + if let Some((callback_fn, user_data)) = armed { + invoke_user_callback!(self.site, unsafe { callback_fn(user_data, 1) }); + } + } + + /// Install (or clear) the executor callback. + /// + /// Backs the `rmw_subscription_set_on_new_message_callback`, + /// `rmw_service_set_on_new_request_callback` and + /// `rmw_client_set_on_new_response_callback` entry points. + /// + /// Installing a callback replays any backlog accumulated by [`notify_one`] + /// in a single call, which is what lets an executor that attaches after + /// messages have already arrived — the common startup race — see them. + /// + /// [`notify_one`]: Self::notify_one + pub fn set(&self, callback: Option, user_data: *mut crate::c_void) { + // Collect under the lock... + let backlog = { + let Ok(mut state) = self.state.lock() else { + return; + }; + state.callback = callback; + state.user_data = user_data as usize; + match callback { + Some(callback_fn) if state.unread > 0 => { + Some((callback_fn, std::mem::take(&mut state.unread))) + } + _ => None, + } + }; + // ...guard is dropped, and only now do we call into the executor. + if let Some((callback_fn, count)) = backlog { + tracing::debug!( + "[{}] replaying {} unread item(s) to a newly installed callback", + self.site, + count + ); + invoke_user_callback!(self.site, unsafe { + callback_fn(user_data as usize as *const std::ffi::c_void, count) + }); + } + } + + /// Items that arrived with no callback installed. Test/introspection only. + #[cfg(test)] + fn unread(&self) -> usize { + self.state.lock().map(|s| s.unread).unwrap_or(0) + } +} + +impl std::fmt::Debug for ExecCallback { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ExecCallback") + .field("site", &self.site) + .finish_non_exhaustive() + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::mpsc; + use std::time::Duration; + + /// How long a re-entrant call is given before we declare it deadlocked. + /// The fixed code returns in microseconds; the pre-fix code never returns. + const DEADLOCK_TIMEOUT: Duration = Duration::from_secs(5); + + /// Stands in for the rclcpp executor: the object the `user_data` pointer + /// actually points at, holding both the entity's slot and whatever + /// bookkeeping the callback needs. Per-test, so the suite stays parallel. + struct Executor { + slot: ExecCallback, + calls: AtomicUsize, + last_count: AtomicUsize, + /// How many more times the callback is allowed to re-enter. Bounded so + /// a *correct* implementation terminates: with the fix, re-entry is + /// legal and unbounded re-entry is unbounded recursion, which is a + /// property of this callback and not of the code under test. + reentries_left: AtomicUsize, + } + + impl Executor { + fn new(site: &'static str, reentries: usize) -> Box { + Box::new(Self { + slot: ExecCallback::new(site), + calls: AtomicUsize::new(0), + last_count: AtomicUsize::new(0), + reentries_left: AtomicUsize::new(reentries), + }) + } + + fn user_data(&self) -> *mut crate::c_void { + self as *const Self as *mut crate::c_void + } + + fn enter(&self, count: usize) -> bool { + self.calls.fetch_add(1, Ordering::SeqCst); + self.last_count.store(count, Ordering::SeqCst); + self.reentries_left + .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |n| n.checked_sub(1)) + .is_ok() + } + } + + /// An executor callback that re-enters the rmw API on the same entity by + /// clearing itself — what rclcpp does when an executor detaches from + /// inside a notification. + unsafe extern "C" fn reenters_by_clearing(user_data: *const std::ffi::c_void, count: usize) { + let exec = unsafe { &*(user_data as *const Executor) }; + if exec.enter(count) { + // Pre-fix, this blocks on a guard this very thread already holds. + exec.slot.set(None, std::ptr::null_mut()); + } + } + + /// An executor callback that re-enters by asking for a fresh notification, + /// covering the delivery-thread entry point rather than the install one. + unsafe extern "C" fn reenters_by_notifying(user_data: *const std::ffi::c_void, count: usize) { + let exec = unsafe { &*(user_data as *const Executor) }; + if exec.enter(count) { + exec.slot.notify_one(); + } + } + + /// A callback that does not re-enter. + unsafe extern "C" fn passive(user_data: *const std::ffi::c_void, count: usize) { + let exec = unsafe { &*(user_data as *const Executor) }; + exec.calls.fetch_add(1, Ordering::SeqCst); + exec.last_count.store(count, Ordering::SeqCst); + } + + /// Run `body` on a worker thread and fail if it has not finished within + /// [`DEADLOCK_TIMEOUT`]. + /// + /// A deadlocked worker stays blocked for the lifetime of the test binary. + /// That is deliberate: it cannot be unblocked, and a non-main thread does + /// not keep the process alive, so the run still terminates and reports. + fn assert_completes( + what: &str, + body: impl FnOnce() -> T + Send + 'static, + ) -> T { + let (tx, rx) = mpsc::channel(); + std::thread::spawn(move || { + let _ = tx.send(body()); + }); + match rx.recv_timeout(DEADLOCK_TIMEOUT) { + Ok(v) => v, + Err(_) => panic!( + "{what} did not return within {DEADLOCK_TIMEOUT:?} — the executor \ + callback was invoked while a guard on the same mutex was live, so \ + the re-entrant call is blocked on this thread's own lock" + ), + } + } + + // --- the deadlock detectors --- + + /// Site 1 (the worst): `rmw_subscription_set_on_new_message_callback` + /// replaying a backlog. Pre-fix this held the `callback` guard *and* the + /// `unread_count` guard across the invocation, so a callback that touched + /// either self-deadlocked. This is the common startup race: messages + /// arrive before the executor attaches. + #[test] + fn installing_a_callback_over_a_backlog_survives_reentry() { + let (calls, count) = assert_completes("set() replaying a backlog", || { + let exec = Executor::new("subscription", 1); + exec.slot.notify_one(); + exec.slot.notify_one(); + exec.slot.set(Some(reenters_by_clearing), exec.user_data()); + ( + exec.calls.load(Ordering::SeqCst), + exec.last_count.load(Ordering::SeqCst), + ) + }); + assert_eq!(calls, 1, "backlog replays in exactly one call"); + assert_eq!(count, 2, "and reports both unread items"); + } + + /// Sites 2/3/4: the delivery-thread notifier, which pre-fix held the + /// `callback` guard and the `callback_user_data` guard across the + /// invocation. + #[test] + fn delivery_notification_survives_reentry_into_set() { + let calls = assert_completes("notify_one() dispatching to the executor", || { + let exec = Executor::new("service", 1); + exec.slot.set(Some(reenters_by_clearing), exec.user_data()); + exec.slot.notify_one(); + exec.calls.load(Ordering::SeqCst) + }); + assert_eq!(calls, 1); + } + + /// A callback that re-enters the *same* entry point it was dispatched + /// from. Pre-fix this is a self-deadlock on `callback`. + #[test] + fn delivery_notification_survives_reentry_into_itself() { + let calls = assert_completes("notify_one() re-entered from its own callback", || { + let exec = Executor::new("client", 1); + exec.slot.set(Some(reenters_by_notifying), exec.user_data()); + exec.slot.notify_one(); + exec.calls.load(Ordering::SeqCst) + }); + assert_eq!(calls, 2, "outer dispatch plus one re-entrant dispatch"); + } + + /// The tripwire is only meaningful if the guard really is released before + /// dispatch. Assert that directly, rather than trusting an assertion that + /// passes vacuously whenever the guard count is zero for the wrong reason. + #[test] + fn no_guard_is_live_when_the_callback_runs() { + unsafe extern "C" fn record_live_guards(user_data: *const std::ffi::c_void, _c: usize) { + let exec = unsafe { &*(user_data as *const Executor) }; + exec.last_count + .store(hiroz::reentrancy::live_guards(), Ordering::SeqCst); + } + + let exec = Executor::new("subscription", 0); + exec.last_count.store(usize::MAX, Ordering::SeqCst); + exec.slot.notify_one(); + exec.slot.set(Some(record_live_guards), exec.user_data()); + assert_eq!( + exec.last_count.load(Ordering::SeqCst), + 0, + "set() must dispatch the backlog with no guard live" + ); + + exec.last_count.store(usize::MAX, Ordering::SeqCst); + exec.slot.notify_one(); + assert_eq!( + exec.last_count.load(Ordering::SeqCst), + 0, + "notify_one() must dispatch with no guard live" + ); + } + + // --- semantics preserved from the pre-fix code --- + + #[test] + fn items_arriving_with_no_callback_are_counted_not_dropped() { + let exec = Executor::new("subscription", 0); + exec.slot.notify_one(); + exec.slot.notify_one(); + exec.slot.notify_one(); + assert_eq!(exec.slot.unread(), 3); + assert_eq!(exec.calls.load(Ordering::SeqCst), 0); + } + + #[test] + fn installing_a_callback_drains_the_backlog() { + let exec = Executor::new("subscription", 0); + exec.slot.notify_one(); + exec.slot.set(Some(passive), exec.user_data()); + assert_eq!(exec.slot.unread(), 0, "backlog is reset once replayed"); + assert_eq!(exec.calls.load(Ordering::SeqCst), 1); + } + + #[test] + fn installing_a_callback_with_no_backlog_does_not_dispatch() { + let exec = Executor::new("subscription", 0); + exec.slot.set(Some(passive), exec.user_data()); + assert_eq!(exec.calls.load(Ordering::SeqCst), 0); + } + + #[test] + fn a_cleared_callback_goes_back_to_counting() { + let exec = Executor::new("subscription", 0); + exec.slot.set(Some(passive), exec.user_data()); + exec.slot.notify_one(); + assert_eq!(exec.calls.load(Ordering::SeqCst), 1); + exec.slot.set(None, std::ptr::null_mut()); + exec.slot.notify_one(); + assert_eq!( + exec.calls.load(Ordering::SeqCst), + 1, + "no dispatch once cleared" + ); + assert_eq!(exec.slot.unread(), 1, "and the item is counted instead"); + } +} diff --git a/crates/rmw-zenoh-rs/src/lib.rs b/crates/rmw-zenoh-rs/src/lib.rs index 6c754e1a6..fbe94bba2 100644 --- a/crates/rmw-zenoh-rs/src/lib.rs +++ b/crates/rmw-zenoh-rs/src/lib.rs @@ -13,6 +13,7 @@ macro_rules! cfile { pub mod common; pub mod context; +pub mod exec_callback; pub mod guard_condition; pub mod msg; pub mod node; diff --git a/crates/rmw-zenoh-rs/src/pubsub.rs b/crates/rmw-zenoh-rs/src/pubsub.rs index a9c4e88e2..e1dbae39d 100644 --- a/crates/rmw-zenoh-rs/src/pubsub.rs +++ b/crates/rmw-zenoh-rs/src/pubsub.rs @@ -59,10 +59,7 @@ pub struct SubscriptionImpl { pub topic: CString, pub options: rmw_subscription_options_t, pub qos: rmw_qos_profile_t, - pub callback: - std::sync::Arc>, - pub callback_user_data: std::sync::Arc>, // Store pointer as usize for thread safety - pub unread_count: std::sync::Arc>, // Track messages arrived before callback was set + pub exec_callback: crate::exec_callback::ExecCallback, pub graph: std::sync::Arc, pub entity: hiroz::entity::EndpointEntity, pub notifier: std::sync::Arc, diff --git a/crates/rmw-zenoh-rs/src/rmw.rs b/crates/rmw-zenoh-rs/src/rmw.rs index d9fc35fc1..0ad4eaeba 100644 --- a/crates/rmw-zenoh-rs/src/rmw.rs +++ b/crates/rmw-zenoh-rs/src/rmw.rs @@ -539,36 +539,15 @@ pub extern "C" fn rmw_create_subscription( let zsub_builder = zsub_builder.with_type_info(ts.get_type_info()); // Set type_info BEFORE build() so liveliness token has correct type - // Create shared callback and user_data holders that will be populated after SubscriptionImpl is created - let callback_holder: std::sync::Arc< - std::sync::Mutex, - > = std::sync::Arc::new(std::sync::Mutex::new(None)); - let user_data_holder = std::sync::Arc::new(std::sync::Mutex::new(0usize)); // Store pointer as usize for thread safety - let unread_count_holder = std::sync::Arc::new(std::sync::Mutex::new(0usize)); // Track unread messages + // Executor-notification slot, populated after SubscriptionImpl is created. + let exec_callback = crate::exec_callback::ExecCallback::new("subscription"); // Create notification callback that will wake up wait sets and invoke user callback let notifier_clone = notifier.clone(); - let callback_holder_clone = callback_holder.clone(); - let user_data_holder_clone = user_data_holder.clone(); - let unread_count_clone = unread_count_holder.clone(); + let exec_callback_clone = exec_callback.clone(); let notify_callback = move || { notifier_clone.notify_all(); - // Invoke the user callback if set, otherwise increment unread count - if let Ok(cb) = callback_holder_clone.lock() { - if let Some(callback_fn) = *cb { - if let Ok(user_data_usize) = user_data_holder_clone.lock() { - unsafe { - let user_data_ptr = *user_data_usize as *const std::ffi::c_void; - callback_fn(user_data_ptr, 1); // 1 new message - } - } - } else { - // No callback set, increment unread count - if let Ok(mut unread) = unread_count_clone.lock() { - *unread += 1; - } - } - } + exec_callback_clone.notify_one(); }; let zsub = match zsub_builder.build_with_notifier(notify_callback) { @@ -713,9 +692,7 @@ pub extern "C" fn rmw_create_subscription( topic: topic_cstr.clone(), options: unsafe { *subscription_options }, qos: crate::qos::normalize_rmw_qos(unsafe { &*qos_policies }), - callback: callback_holder, - callback_user_data: user_data_holder, - unread_count: unread_count_holder, + exec_callback, graph: graph.clone(), entity: entity.clone(), notifier, @@ -1041,11 +1018,8 @@ pub extern "C" fn rmw_create_client( .with_qos(qos); let entity = zclient_builder.entity().clone(); - // Create shared callback and user_data holders - let callback_holder: std::sync::Arc< - std::sync::Mutex, - > = std::sync::Arc::new(std::sync::Mutex::new(None)); - let user_data_holder = std::sync::Arc::new(std::sync::Mutex::new(0usize)); + // Executor-notification slot, shared with every per-request notifier. + let exec_callback = crate::exec_callback::ExecCallback::new("client"); // Build the client (notification callback will be set per-request in send_request) let zclient = match zclient_builder.build() { @@ -1069,10 +1043,8 @@ pub extern "C" fn rmw_create_client( }, request_ts: service_type_support, response_ts: service_type_support, - callback: callback_holder, - callback_user_data: user_data_holder, + exec_callback, notifier, - unread_count: std::sync::Arc::new(std::sync::Mutex::new(0)), graph: node_impl.inner.graph().clone(), entity: entity.clone(), }; @@ -1237,36 +1209,15 @@ pub extern "C" fn rmw_create_service( .with_qos(qos); let entity = zserver_builder.entity().clone(); - // Create shared callback and user_data holders that will be populated after ServiceImpl is created - let callback_holder: std::sync::Arc< - std::sync::Mutex, - > = std::sync::Arc::new(std::sync::Mutex::new(None)); - let user_data_holder = std::sync::Arc::new(std::sync::Mutex::new(0usize)); // Store pointer as usize for thread safety - let unread_count_holder = std::sync::Arc::new(std::sync::Mutex::new(0usize)); // Track unread requests + // Executor-notification slot, populated after ServiceImpl is created. + let exec_callback = crate::exec_callback::ExecCallback::new("service"); // Create notification callback that will wake up wait sets and invoke user callback let notifier_clone = notifier.clone(); - let callback_holder_clone = callback_holder.clone(); - let user_data_holder_clone = user_data_holder.clone(); - let unread_count_clone = unread_count_holder.clone(); + let exec_callback_clone = exec_callback.clone(); let notify_callback = move || { notifier_clone.notify_all(); - // Invoke user callback if set, otherwise increment unread count - if let Ok(cb) = callback_holder_clone.lock() { - if let Some(callback_fn) = *cb { - if let Ok(user_data_usize) = user_data_holder_clone.lock() { - unsafe { - let user_data_ptr = *user_data_usize as *const std::ffi::c_void; - callback_fn(user_data_ptr, 1); // 1 new request - } - } - } else { - // No callback set, increment unread count - if let Ok(mut unread) = unread_count_clone.lock() { - *unread += 1; - } - } - } + exec_callback_clone.notify_one(); }; let zserver = match zserver_builder.build_with_notifier(notify_callback) { @@ -1289,9 +1240,7 @@ pub extern "C" fn rmw_create_service( request_ts: service_type_support, response_ts: service_type_support, qos: crate::qos::normalize_rmw_qos(unsafe { &*qos_profile }), - callback: callback_holder, - callback_user_data: user_data_holder, - unread_count: unread_count_holder, // Use the same Arc created earlier for notify_callback + exec_callback, graph: node_impl.inner.graph().clone(), entity: entity.clone(), }; @@ -2473,29 +2422,7 @@ pub extern "C" fn rmw_subscription_set_on_new_message_callback( Err(_) => return RMW_RET_INVALID_ARGUMENT as _, }; - // Set user_data first - if let Ok(mut ud) = subscription_impl.callback_user_data.lock() { - *ud = user_data as usize; // Store pointer as usize for thread safety - } - - // Then set callback and check for unread messages - if let Ok(mut cb) = subscription_impl.callback.lock() { - if callback.is_some() { - // Check if there are unread messages and invoke callback if needed - if let Ok(mut unread) = subscription_impl.unread_count.lock() { - if *unread > 0 { - // Invoke callback with unread count - unsafe { - if let Some(callback_fn) = callback { - callback_fn(user_data as *const std::ffi::c_void, *unread); - } - } - *unread = 0; // Reset unread count after notifying - } - } - } - *cb = callback; - } + subscription_impl.exec_callback.set(callback, user_data); RMW_RET_OK as _ } @@ -2515,36 +2442,7 @@ pub extern "C" fn rmw_service_set_on_new_request_callback( Err(_) => return RMW_RET_INVALID_ARGUMENT as _, }; - if let Some(callback_fn) = callback { - // Push events arrived before setting the executor callback (retroactive notification) - if let Ok(mut unread) = service_impl.unread_count.lock() { - if *unread > 0 { - tracing::debug!( - "[rmw_service_set_on_new_request_callback] Invoking callback retroactively for {} unread requests", - *unread - ); - unsafe { - callback_fn(user_data as *const std::ffi::c_void, *unread); - } - *unread = 0; // Reset unread count after notification - } - } - // Store the new callback and user_data - if let Ok(mut cb) = service_impl.callback.lock() { - *cb = callback; - } - if let Ok(mut ud) = service_impl.callback_user_data.lock() { - *ud = user_data as usize; - } - } else { - // Callback is being cleared (set to None) - if let Ok(mut cb) = service_impl.callback.lock() { - *cb = None; - } - if let Ok(mut ud) = service_impl.callback_user_data.lock() { - *ud = 0; - } - } + service_impl.exec_callback.set(callback, user_data); RMW_RET_OK as _ } @@ -2564,36 +2462,7 @@ pub extern "C" fn rmw_client_set_on_new_response_callback( Err(_) => return RMW_RET_INVALID_ARGUMENT as _, }; - if let Some(callback_fn) = callback { - // Push events arrived before setting the executor callback (retroactive notification) - if let Ok(mut unread) = client_impl.unread_count.lock() { - if *unread > 0 { - tracing::debug!( - "[rmw_client_set_on_new_response_callback] Invoking callback retroactively for {} unread responses", - *unread - ); - unsafe { - callback_fn(user_data as *const std::ffi::c_void, *unread); - } - *unread = 0; // Reset unread count after notification - } - } - // Store the new callback and user_data - if let Ok(mut cb) = client_impl.callback.lock() { - *cb = callback; - } - if let Ok(mut ud) = client_impl.callback_user_data.lock() { - *ud = user_data as usize; - } - } else { - // Callback is being cleared (set to None) - if let Ok(mut cb) = client_impl.callback.lock() { - *cb = None; - } - if let Ok(mut ud) = client_impl.callback_user_data.lock() { - *ud = 0; - } - } + client_impl.exec_callback.set(callback, user_data); RMW_RET_OK as _ } diff --git a/crates/rmw-zenoh-rs/src/service.rs b/crates/rmw-zenoh-rs/src/service.rs index effb92596..a4288407f 100644 --- a/crates/rmw-zenoh-rs/src/service.rs +++ b/crates/rmw-zenoh-rs/src/service.rs @@ -1,6 +1,5 @@ use std::collections::HashMap; use std::ffi::CString; -use std::sync::Mutex; use crate::c_void; use crate::rmw_impl_has_data_ptr; @@ -15,11 +14,8 @@ pub struct ClientImpl { pub options: rmw_client_options_t, pub request_ts: crate::type_support::ServiceTypeSupport, pub response_ts: crate::type_support::ServiceTypeSupport, - pub callback: std::sync::Arc>, - pub callback_user_data: std::sync::Arc>, + pub exec_callback: crate::exec_callback::ExecCallback, pub notifier: std::sync::Arc, - /// Tracks responses that arrived while no callback was set - pub unread_count: std::sync::Arc>, pub graph: std::sync::Arc, pub entity: hiroz::entity::EndpointEntity, } @@ -29,23 +25,10 @@ impl ClientImpl { let req = crate::msg::RosMessage::new(request, self.request_ts.request); let notifier = self.notifier.clone(); - let callback_holder = self.callback.clone(); - let user_data_holder = self.callback_user_data.clone(); - let unread_count_holder = self.unread_count.clone(); + let exec_callback = self.exec_callback.clone(); let notify_callback = move || { notifier.notify_all(); - if let Ok(cb) = callback_holder.lock() { - if let Some(callback_fn) = *cb { - if let Ok(user_data_usize) = user_data_holder.lock() { - unsafe { - let user_data_ptr = *user_data_usize as *const std::ffi::c_void; - callback_fn(user_data_ptr, 1); - } - } - } else if let Ok(mut unread) = unread_count_holder.lock() { - *unread += 1; - } - } + exec_callback.notify_one(); }; // rmw_send_request returns the sequence number stamped into the attachment, @@ -161,10 +144,7 @@ pub struct ServiceImpl { pub request_ts: crate::type_support::ServiceTypeSupport, pub response_ts: crate::type_support::ServiceTypeSupport, pub qos: rmw_qos_profile_t, - pub callback: std::sync::Arc>, - pub callback_user_data: std::sync::Arc>, - /// Tracks requests that arrived while no callback was set - pub unread_count: std::sync::Arc>, + pub exec_callback: crate::exec_callback::ExecCallback, pub graph: std::sync::Arc, pub entity: hiroz::entity::EndpointEntity, } From e87a722774df987401c958044501472654e7360e Mon Sep 17 00:00:00 2001 From: yuanyuyuan Date: Thu, 30 Jul 2026 09:46:22 +0800 Subject: [PATCH 2/4] fix(rmw): make set() wait for in-flight dispatches on other threads Dropping the state guard before dispatch fixed the deadlock but removed a lifetime guarantee rmw_zenoh_cpp provides deliberately: it invokes the callback under event_mutex_, which event_set_callback also takes, so once set_callback(nullptr) returns no callback is in flight and the caller may free what user_data pointed at. Without that, a delivery thread could snapshot (callback, user_data), drop the guard, and then be overtaken by a set(None) that clears the slot and returns -- after which the entity is destroyed and the snapshot dangles. Restoring the C++ shape would reinstate the deadlock, so exclusion and lifetime are separated: each dispatch registers its thread, and set() waits for registrations on *other* threads. Scoping the wait that way is what keeps a callback re-entering set() on its own thread from waiting on itself. --- crates/rmw-zenoh-rs/src/exec_callback.rs | 315 ++++++++++++++++++++++- 1 file changed, 304 insertions(+), 11 deletions(-) diff --git a/crates/rmw-zenoh-rs/src/exec_callback.rs b/crates/rmw-zenoh-rs/src/exec_callback.rs index 8fcf7068c..f1f2727ab 100644 --- a/crates/rmw-zenoh-rs/src/exec_callback.rs +++ b/crates/rmw-zenoh-rs/src/exec_callback.rs @@ -33,6 +33,35 @@ //! rested on every site agreeing on a lock order. With one, the invariant is //! local and there is no order to get wrong. //! +//! # Why dropping the guard is not sufficient on its own +//! +//! `rmw_zenoh_cpp` dispatches this callback *under* its `event_mutex_` +//! (`DataCallbackManager::trigger_callback`), and `event_set_callback` takes the +//! same mutex. That is not only exclusion — it is a **lifetime guarantee**: once +//! `set_callback(nullptr)` returns, no callback is in flight, so the caller may +//! free whatever `user_data` pointed at. +//! +//! Collecting the callback and its `user_data` under the lock and dispatching +//! after the guard drops removes that guarantee. The window: +//! +//! 1. A delivery thread enters [`ExecCallback::notify_one`], snapshots +//! `(callback, user_data)`, and drops the guard. +//! 2. The executor thread calls [`ExecCallback::set`] with `None`. It takes the +//! lock, clears the slot, and returns — without waiting. +//! 3. The entity is destroyed and `user_data` is freed. +//! 4. The delivery thread calls the callback with its stale snapshot — +//! use-after-free. +//! +//! Restoring the C++ shape is not an option: dispatching under the lock is +//! exactly the deadlock above. So the exclusion and the lifetime guarantee are +//! separated. Each dispatch registers the thread performing it, and `set` waits +//! for registered dispatches **on other threads** to finish before it returns. +//! +//! The thread distinction is the whole point, and it is what makes this +//! different from the original bug. A callback that re-enters `set` on its own +//! thread must not wait — waiting for itself is the deadlock. A `set` on an +//! unrelated thread must wait, because it is the one about to free the pointer. +//! //! # Enforcement //! //! The mutex is a [`hiroz::reentrancy::TrackedMutex`], so its guards are counted @@ -41,8 +70,12 @@ //! dispatching. In debug builds a reintroduction of the defect panics with the //! site name instead of hanging; in release both the counter and the assertion //! compile to nothing. +//! +//! Neither the state guard nor the in-flight guard is ever held across user +//! code. Where both are taken, the order is always state then in-flight. -use std::sync::Arc; +use std::sync::{Arc, Condvar, Mutex}; +use std::thread::ThreadId; use hiroz::invoke_user_callback; use hiroz::reentrancy::TrackedMutex; @@ -63,6 +96,44 @@ struct State { unread: usize, } +/// Threads currently inside a dispatch. +/// +/// One entry per active dispatch, not per thread: a callback may re-enter +/// `notify_one` on the same thread, so the same `ThreadId` can appear more than +/// once and each occurrence must be removed independently. +#[derive(Default)] +struct InFlight { + threads: Vec, +} + +/// The in-flight registry and the condvar `set` waits on. +type InFlightSlot = (Mutex, Condvar); + +/// Deregisters this thread's dispatch on scope exit, **including on unwind**. +/// +/// Without the unwind path, a panicking executor callback would leave its entry +/// behind forever and every later `set` would block on a dispatch that has +/// already finished — trading a use-after-free for a permanent hang. +struct DispatchToken<'a> { + slot: &'a InFlightSlot, +} + +impl Drop for DispatchToken<'_> { + fn drop(&mut self) { + let (lock, condvar) = self.slot; + // Poison-tolerant: this lock is only ever held for a push or a remove, + // never across user code, so a poisoned state carries no broken + // invariant — but failing to deregister here would hang every `set`. + let mut in_flight = lock.lock().unwrap_or_else(|e| e.into_inner()); + let me = std::thread::current().id(); + if let Some(i) = in_flight.threads.iter().position(|t| *t == me) { + in_flight.threads.swap_remove(i); + } + drop(in_flight); + condvar.notify_all(); + } +} + /// Shared executor-notification state for one rmw entity. /// /// Cloning shares the underlying slot; the delivery thread and the rmw API @@ -70,6 +141,11 @@ struct State { #[derive(Clone)] pub struct ExecCallback { state: Arc>, + /// Dispatches currently running, so [`set`] can wait for the ones that + /// captured the outgoing `user_data`. + /// + /// [`set`]: Self::set + in_flight: Arc, /// Entity kind, reproduced in the re-entrancy panic message. site: &'static str, } @@ -80,10 +156,43 @@ impl ExecCallback { pub fn new(site: &'static str) -> Self { Self { state: Arc::new(TrackedMutex::new(State::default())), + in_flight: Arc::new((Mutex::new(InFlight::default()), Condvar::new())), site, } } + /// Register this thread as dispatching, returning the token that removes it. + /// + /// **Must be called while the state guard is still held.** `set` swaps the + /// slot under that same guard, so registering before releasing it is what + /// guarantees that every dispatch holding the outgoing `user_data` is + /// already visible to `set`'s wait. Register after, and step 2 of the race + /// in the module docs reopens. + fn enter_dispatch(&self) -> DispatchToken<'_> { + let (lock, _) = &*self.in_flight; + lock.lock() + .unwrap_or_else(|e| e.into_inner()) + .threads + .push(std::thread::current().id()); + DispatchToken { + slot: &self.in_flight, + } + } + + /// Block until no *other* thread is inside a dispatch. + /// + /// Entries for the current thread are ignored on purpose: a callback that + /// re-enters `set` is the deadlock this type exists to prevent, and it + /// cannot be waiting for a pointer it is itself about to stop using. + fn wait_for_other_dispatches(&self) { + let (lock, condvar) = &*self.in_flight; + let me = std::thread::current().id(); + let mut in_flight = lock.lock().unwrap_or_else(|e| e.into_inner()); + while in_flight.threads.iter().any(|t| *t != me) { + in_flight = condvar.wait(in_flight).unwrap_or_else(|e| e.into_inner()); + } + } + /// One new item arrived. Invokes the executor callback if one is installed, /// otherwise records the item as unread so it can be replayed by [`set`]. /// @@ -91,14 +200,20 @@ impl ExecCallback { /// /// [`set`]: Self::set pub fn notify_one(&self) { - // Collect under the lock... + // Collect under the lock, and register the dispatch before releasing it + // so `set` cannot conclude that no callback is in flight. let armed = { let Ok(mut state) = self.state.lock() else { return; }; match state.callback { Some(callback_fn) => { - Some((callback_fn, state.user_data as *const std::ffi::c_void)) + let token = self.enter_dispatch(); + Some(( + callback_fn, + state.user_data as *const std::ffi::c_void, + token, + )) } None => { state.unread += 1; @@ -106,8 +221,9 @@ impl ExecCallback { } } }; - // ...guard is dropped, and only now do we call into the executor. - if let Some((callback_fn, user_data)) = armed { + // ...guard is dropped, and only now do we call into the executor. The + // token outlives the call and deregisters on the way out, panic or not. + if let Some((callback_fn, user_data, _token)) = armed { invoke_user_callback!(self.site, unsafe { callback_fn(user_data, 1) }); } } @@ -123,8 +239,20 @@ impl ExecCallback { /// messages have already arrived — the common startup race — see them. /// /// [`notify_one`]: Self::notify_one + /// # Blocking + /// + /// This returns only once every dispatch that captured the *outgoing* + /// `user_data` has finished, so that a caller clearing the slot may then + /// free what it pointed at. It therefore blocks for as long as an executor + /// callback already running on another thread takes to return — + /// `rmw_zenoh_cpp` has the same property, where the wait is on + /// `event_mutex_` instead. A callback re-entering on its own thread never + /// waits. pub fn set(&self, callback: Option, user_data: *mut crate::c_void) { - // Collect under the lock... + // Collect under the lock, registering the replay — if there is one — + // before releasing it, for the same reason `notify_one` does: a + // concurrent `set` must wait for this replay before freeing the pointer + // being handed to it here. let backlog = { let Ok(mut state) = self.state.lock() else { return; @@ -133,13 +261,21 @@ impl ExecCallback { state.user_data = user_data as usize; match callback { Some(callback_fn) if state.unread > 0 => { - Some((callback_fn, std::mem::take(&mut state.unread))) + let token = self.enter_dispatch(); + Some((callback_fn, std::mem::take(&mut state.unread), token)) } _ => None, } }; - // ...guard is dropped, and only now do we call into the executor. - if let Some((callback_fn, count)) = backlog { + + // ...guard is dropped. The slot now holds the incoming `user_data`, so + // any dispatch starting from here on uses the new pointer; the ones + // still registered captured the old one. Wait for exactly those, which + // is what lets the caller free it once we return. Our own replay + // registration is on this thread and so is not waited for. + self.wait_for_other_dispatches(); + + if let Some((callback_fn, count, _token)) = backlog { tracing::debug!( "[{}] replaying {} unread item(s) to a newly installed callback", self.site, @@ -169,9 +305,10 @@ impl std::fmt::Debug for ExecCallback { #[cfg(test)] mod tests { use super::*; - use std::sync::atomic::{AtomicUsize, Ordering}; + use std::panic::AssertUnwindSafe; + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::mpsc; - use std::time::Duration; + use std::time::{Duration, Instant}; /// How long a re-entrant call is given before we declare it deadlocked. /// The fixed code returns in microseconds; the pre-fix code never returns. @@ -189,8 +326,43 @@ mod tests { /// legal and unbounded re-entry is unbounded recursion, which is a /// property of this callback and not of the code under test. reentries_left: AtomicUsize, + /// Set once a callback is inside and about to block on `gate`. + entered: AtomicBool, + /// Holds a callback open so a concurrent `set` has something to wait for. + gate: Gate, + /// Ticket source for the ordering assertions. + seq: AtomicUsize, + /// Ticket taken as the blocking callback returns. + callback_exit: AtomicUsize, + /// Ticket taken as the concurrent `set` returns. + set_returned: AtomicUsize, + } + + /// A one-shot gate. Used to pin a callback in flight for as long as a test + /// needs, without a sleep deciding the outcome. + #[derive(Default)] + struct Gate { + open: Mutex, + condvar: Condvar, } + impl Gate { + fn wait(&self) { + let mut open = self.open.lock().unwrap(); + while !*open { + open = self.condvar.wait(open).unwrap(); + } + } + + fn open(&self) { + *self.open.lock().unwrap() = true; + self.condvar.notify_all(); + } + } + + /// No ticket taken yet. + const NO_TICKET: usize = usize::MAX; + impl Executor { fn new(site: &'static str, reentries: usize) -> Box { Box::new(Self { @@ -198,9 +370,18 @@ mod tests { calls: AtomicUsize::new(0), last_count: AtomicUsize::new(0), reentries_left: AtomicUsize::new(reentries), + entered: AtomicBool::new(false), + gate: Gate::default(), + seq: AtomicUsize::new(0), + callback_exit: AtomicUsize::new(NO_TICKET), + set_returned: AtomicUsize::new(NO_TICKET), }) } + fn ticket(&self) -> usize { + self.seq.fetch_add(1, Ordering::SeqCst) + } + fn user_data(&self) -> *mut crate::c_void { self as *const Self as *mut crate::c_void } @@ -234,6 +415,17 @@ mod tests { } } + /// A callback that parks inside the dispatch until the test releases it, + /// then takes a ticket on the way out. Lets a test observe whether a + /// concurrent `set` returned before or after the callback finished. + unsafe extern "C" fn blocks_until_released(user_data: *const std::ffi::c_void, _count: usize) { + let exec = unsafe { &*(user_data as *const Executor) }; + exec.calls.fetch_add(1, Ordering::SeqCst); + exec.entered.store(true, Ordering::SeqCst); + exec.gate.wait(); + exec.callback_exit.store(exec.ticket(), Ordering::SeqCst); + } + /// A callback that does not re-enter. unsafe extern "C" fn passive(user_data: *const std::ffi::c_void, count: usize) { let exec = unsafe { &*(user_data as *const Executor) }; @@ -345,6 +537,107 @@ mod tests { ); } + // --- the use-after-free detectors --- + + /// The lifetime guarantee `rmw_zenoh_cpp` gets from dispatching under + /// `event_mutex_`: once `set(None, ..)` returns, no callback is in flight, + /// so the caller may free what `user_data` pointed at. + /// + /// Collect-then-dispatch alone does not provide this. Without the in-flight + /// wait, `set` takes the lock, clears the slot and returns while a delivery + /// thread is still inside the callback holding the outgoing pointer — and + /// the next thing rclcpp does is destroy the entity. + #[test] + fn set_waits_for_a_dispatch_in_flight_on_another_thread() { + // Leaked on purpose. Dropping it here would assume exactly what is + // under test — that it is safe to free once `set` has returned. + let exec: &'static Executor = Box::leak(Executor::new("subscription", 0)); + exec.slot.set(Some(blocks_until_released), exec.user_data()); + + let delivery = std::thread::spawn(move || exec.slot.notify_one()); + + // Park until the callback is genuinely inside the dispatch. + let start = Instant::now(); + while !exec.entered.load(Ordering::SeqCst) { + assert!( + start.elapsed() < DEADLOCK_TIMEOUT, + "the callback never entered; the test cannot say anything" + ); + std::thread::yield_now(); + } + + let clearer = std::thread::spawn(move || { + exec.slot.set(None, std::ptr::null_mut()); + exec.set_returned.store(exec.ticket(), Ordering::SeqCst); + }); + + // The callback is still parked. A `set` that does not wait has already + // returned by now; one that waits cannot have. + std::thread::sleep(Duration::from_millis(250)); + assert_eq!( + exec.set_returned.load(Ordering::SeqCst), + NO_TICKET, + "set() returned while a callback was still in flight on another \ + thread — rclcpp is now free to destroy the entity and free the \ + user_data that callback is still holding" + ); + + exec.gate.open(); + delivery.join().expect("delivery thread"); + clearer.join().expect("clearing thread"); + + assert_eq!( + exec.callback_exit.load(Ordering::SeqCst), + 0, + "the callback should have taken the first ticket" + ); + assert_eq!( + exec.set_returned.load(Ordering::SeqCst), + 1, + "and set() the second — it must return strictly after the callback" + ); + } + + /// The other half of the same rule, and the reason the wait is scoped to + /// *other* threads: a callback re-entering `set` on its own thread must not + /// wait for itself. Waiting for all dispatches unconditionally would + /// reinstate the deadlock this type exists to remove, in a new place. + #[test] + fn set_does_not_wait_for_a_dispatch_on_its_own_thread() { + let calls = assert_completes("set() re-entered from its own dispatch", || { + let exec = Executor::new("service", 1); + exec.slot.set(Some(reenters_by_clearing), exec.user_data()); + // The callback runs on this thread and calls `set` from inside the + // dispatch. If the wait counted this thread, it would block here. + exec.slot.notify_one(); + exec.calls.load(Ordering::SeqCst) + }); + assert_eq!(calls, 1); + } + + /// A panicking callback must not leave its registration behind: the wait in + /// `set` would then never be satisfiable, turning the use-after-free into a + /// permanent hang. Covers the `DispatchToken` unwind path. + /// + /// The panic is caught, so the run prints one backtrace-ish line from the + /// default hook. That noise is expected. + #[test] + fn a_panicking_callback_does_not_wedge_later_sets() { + unsafe extern "C" fn panics(_user_data: *const std::ffi::c_void, _count: usize) { + panic!("executor callback exploded"); + } + + let exec = Executor::new("client", 0); + exec.slot.set(Some(panics), exec.user_data()); + + let unwound = std::panic::catch_unwind(AssertUnwindSafe(|| exec.slot.notify_one())); + assert!(unwound.is_err(), "the panic must still propagate"); + + assert_completes("set() after a callback panicked", move || { + exec.slot.set(None, std::ptr::null_mut()); + }); + } + // --- semantics preserved from the pre-fix code --- #[test] From 752a8a504a2acc0e7fcfa73bf324247c1b0a725c Mon Sep 17 00:00:00 2001 From: yuanyuyuan Date: Thu, 30 Jul 2026 09:50:26 +0800 Subject: [PATCH 3/4] test(rmw): exercise the reachable unwind path, not an FFI panic A panic out of an `extern "C"` callback aborts rather than unwinding, so the previous test killed the test binary (rc=101) instead of asserting anything -- and took the other tests' results with it. The reachable unwind is Rust-side, before the boundary: `invoke_user_callback!`'s debug assertion. Raise it directly and keep asserting what matters, that the dispatch registration is released so a later set() cannot block forever. --- crates/rmw-zenoh-rs/src/exec_callback.rs | 54 ++++++++++++++++++------ 1 file changed, 40 insertions(+), 14 deletions(-) diff --git a/crates/rmw-zenoh-rs/src/exec_callback.rs b/crates/rmw-zenoh-rs/src/exec_callback.rs index f1f2727ab..80f1ec4e6 100644 --- a/crates/rmw-zenoh-rs/src/exec_callback.rs +++ b/crates/rmw-zenoh-rs/src/exec_callback.rs @@ -62,6 +62,23 @@ //! thread must not wait — waiting for itself is the deadlock. A `set` on an //! unrelated thread must wait, because it is the one about to free the pointer. //! +//! ## What this still does not survive +//! +//! Two threads *both* dispatching this entity's callback *and* both re-entering +//! `set` from inside it will wait on each other. Each is in a live dispatch that +//! may still touch its `user_data` after `set` returns, so neither wait can +//! safely be skipped. +//! +//! Stated rather than hidden, because it is a real residual — but it is not a +//! regression against the reference. `rmw_zenoh_cpp` cannot survive re-entrant +//! `set` at all: `set_callback` takes `event_mutex_` (and replays its backlog +//! under it), so a callback dispatched from `trigger_callback` that calls +//! `set_callback` re-locks a non-recursive `std::mutex` on the thread that +//! already holds it. That deadlocks with one thread. This deadlocks only with +//! two threads mutually re-entering, and the single-threaded case — the one +//! rclcpp actually exercises when an executor detaches from inside a +//! notification — is the case this type makes work. +//! //! # Enforcement //! //! The mutex is a [`hiroz::reentrancy::TrackedMutex`], so its guards are counted @@ -238,7 +255,6 @@ impl ExecCallback { /// in a single call, which is what lets an executor that attaches after /// messages have already arrived — the common startup race — see them. /// - /// [`notify_one`]: Self::notify_one /// # Blocking /// /// This returns only once every dispatch that captured the *outgoing* @@ -248,6 +264,8 @@ impl ExecCallback { /// `rmw_zenoh_cpp` has the same property, where the wait is on /// `event_mutex_` instead. A callback re-entering on its own thread never /// waits. + /// + /// [`notify_one`]: Self::notify_one pub fn set(&self, callback: Option, user_data: *mut crate::c_void) { // Collect under the lock, registering the replay — if there is one — // before releasing it, for the same reason `notify_one` does: a @@ -615,25 +633,33 @@ mod tests { assert_eq!(calls, 1); } - /// A panicking callback must not leave its registration behind: the wait in - /// `set` would then never be satisfiable, turning the use-after-free into a - /// permanent hang. Covers the `DispatchToken` unwind path. + /// An unwind through a live dispatch must not leave its registration + /// behind: the wait in `set` would then never be satisfiable, turning the + /// use-after-free into a permanent hang. Covers `DispatchToken`'s `Drop`. + /// + /// The unwind is raised directly rather than from a callback, because a + /// panic *out of a callback* is not the reachable case: the callbacks are + /// `extern "C"`, and Rust aborts rather than unwinding across that + /// boundary — a panicking executor callback kills the process before any + /// `Drop` runs. What can unwind here is Rust code on this side of the + /// boundary, most notably `invoke_user_callback!`'s own debug assertion, + /// which fires *before* the call. /// - /// The panic is caught, so the run prints one backtrace-ish line from the - /// default hook. That noise is expected. + /// The panic is caught, so the run prints one line from the default hook. + /// That noise is expected. #[test] - fn a_panicking_callback_does_not_wedge_later_sets() { - unsafe extern "C" fn panics(_user_data: *const std::ffi::c_void, _count: usize) { - panic!("executor callback exploded"); - } - + fn an_unwind_through_a_dispatch_releases_its_registration() { let exec = Executor::new("client", 0); - exec.slot.set(Some(panics), exec.user_data()); - let unwound = std::panic::catch_unwind(AssertUnwindSafe(|| exec.slot.notify_one())); + let unwound = std::panic::catch_unwind(AssertUnwindSafe(|| { + let _token = exec.slot.enter_dispatch(); + panic!("something unwound mid-dispatch"); + })); assert!(unwound.is_err(), "the panic must still propagate"); - assert_completes("set() after a callback panicked", move || { + // If the registration leaked, this blocks forever: the wait is looking + // for a dispatch on another thread that has already gone. + assert_completes("set() after an unwind mid-dispatch", move || { exec.slot.set(None, std::ptr::null_mut()); }); } From 262184da045cc564d6d1c494685fec8643ab643f Mon Sep 17 00:00:00 2001 From: yuanyuyuan Date: Thu, 30 Jul 2026 10:02:11 +0800 Subject: [PATCH 4/4] test(rmw): stop the timeout message naming only one cause assert_completes is now shared by tests whose blocking has three different causes; asserting the original one for all of them misdescribes two. --- crates/rmw-zenoh-rs/src/exec_callback.rs | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/crates/rmw-zenoh-rs/src/exec_callback.rs b/crates/rmw-zenoh-rs/src/exec_callback.rs index 80f1ec4e6..710d62a8a 100644 --- a/crates/rmw-zenoh-rs/src/exec_callback.rs +++ b/crates/rmw-zenoh-rs/src/exec_callback.rs @@ -468,9 +468,11 @@ mod tests { match rx.recv_timeout(DEADLOCK_TIMEOUT) { Ok(v) => v, Err(_) => panic!( - "{what} did not return within {DEADLOCK_TIMEOUT:?} — the executor \ - callback was invoked while a guard on the same mutex was live, so \ - the re-entrant call is blocked on this thread's own lock" + "{what} did not return within {DEADLOCK_TIMEOUT:?}, so it is blocked. \ + The causes this suite distinguishes: a guard on the state mutex was \ + still live across the dispatch (the original defect); the in-flight \ + wait counted this thread and is waiting on itself; or a finished \ + dispatch failed to deregister and the wait can never be satisfied" ), } }