diff --git a/crates/hiroz-py/src/node.rs b/crates/hiroz-py/src/node.rs index 6e43eb4dc..bfb9edb9f 100644 --- a/crates/hiroz-py/src/node.rs +++ b/crates/hiroz-py/src/node.rs @@ -158,6 +158,24 @@ pub struct PyZNode { next_sub_id: u64, } +impl Drop for PyZNode { + fn drop(&mut self) { + // Same hazard as `destroy_subscriber`, reached a different way: `del + // node` and interpreter shutdown run `tp_dealloc` with the GIL held. + // + // Dropping a `ZSub` no longer joins its delivery thread (it detaches -- + // see `CallbackDispatcher::drop`), so this wrapper is belt-and-braces + // rather than load-bearing. It is kept because it costs nothing and it + // documents the hazard for anyone who adds a barrier here later. + // Take the subscribers out and drop them with the GIL released. + if self.owned_subs.is_empty() { + return; + } + let subs = std::mem::take(&mut self.owned_subs); + Python::with_gil(move |py| py.allow_threads(move || drop(subs))); + } +} + #[allow(unsafe_op_in_unsafe_fn)] #[pymethods] impl PyZNode { @@ -410,14 +428,26 @@ impl PyZNode { /// /// Matches rclpy's `Node.destroy_subscription()`. Has no effect on queue-based /// subscribers (those are owned by the caller and dropped when they go out of scope). - fn destroy_subscriber(&mut self, sub: &PyZSubscriber) -> PyResult<()> { + fn destroy_subscriber(&mut self, py: Python<'_>, sub: &PyZSubscriber) -> PyResult<()> { let Some(id) = sub.owned_id else { return Err(pyo3::exceptions::PyValueError::new_err( "destroy_subscriber only applies to callback-based subscribers", )); }; if let Some(pos) = self.owned_subs.iter().position(|(sid, _)| *sid == id) { - self.owned_subs.swap_remove(pos); + let owned = self.owned_subs.swap_remove(pos); + // Drop with the GIL released. + // + // Dropping a `ZSub` no longer joins its delivery thread -- it + // detaches, see `CallbackDispatcher::drop` -- so this wrapper is + // belt-and-braces rather than load-bearing. It is kept because it + // costs nothing and it documents the hazard it used to prevent: + // that thread's callback body is `Python::with_gil`, a + // `#[pymethods]` fn runs with the GIL held, and any barrier added + // here would wait for a thread that is waiting for the GIL we hold. + // The interpreter would freeze with no exception and no traceback. + // Anyone reaching for `ZSub::close` here must keep this wrapper. + py.allow_threads(move || drop(owned)); } Ok(()) } diff --git a/crates/hiroz-py/src/pubsub.rs b/crates/hiroz-py/src/pubsub.rs index 16fada65e..9ff7e4c65 100644 --- a/crates/hiroz-py/src/pubsub.rs +++ b/crates/hiroz-py/src/pubsub.rs @@ -24,20 +24,30 @@ impl PyZPublisher { /// Publish a message /// /// Serializes the Python message (msgspec.Struct) to ZBuf and publishes (zero-copy path) - unsafe fn publish(&self, _py: Python, data: &Bound<'_, PyAny>) -> PyResult<()> { - // Serialize Python message directly to ZBuf (zero-copy) + unsafe fn publish(&self, py: Python, data: &Bound<'_, PyAny>) -> PyResult<()> { + // Serialize Python message directly to ZBuf (zero-copy). This touches + // Python objects, so it must run with the GIL held. let zbuf = hiroz_msgs::serialize_to_zbuf(&self.type_name, data)?; - // Publish the ZBuf directly - self.inner.publish(zbuf.into()).map_err(|e| e.into_pyerr()) + // Release the GIL for the publish itself. Zenoh delivers samples to + // local subscribers synchronously on the publishing thread, so a publish + // issued from inside a subscriber callback can block here; holding the + // GIL across it would freeze the whole interpreter — no exception, no + // traceback — instead of blocking just this thread. + py.allow_threads(|| self.inner.publish(zbuf.into())) + .map_err(|e| e.into_pyerr()) } /// Publish pre-serialized CDR bytes directly /// /// Use this for zero-copy forwarding of received messages (e.g., in a pong responder). /// The bytes should be in CDR format (as returned by recv_serialized/try_recv_serialized). - fn publish_raw(&self, data: &[u8]) -> PyResult<()> { - self.inner.publish(data.into()).map_err(|e| e.into_pyerr()) + fn publish_raw(&self, py: Python, data: &[u8]) -> PyResult<()> { + // Copy out of the Python buffer before dropping the GIL, then publish + // without it — see `publish` for why. + let payload: zenoh::bytes::ZBytes = data.into(); + py.allow_threads(|| self.inner.publish(payload)) + .map_err(|e| e.into_pyerr()) } /// Get the topic name (for debugging) diff --git a/crates/hiroz-py/tests/test_reentrant_publish.py b/crates/hiroz-py/tests/test_reentrant_publish.py new file mode 100644 index 000000000..fe3a59bc0 --- /dev/null +++ b/crates/hiroz-py/tests/test_reentrant_publish.py @@ -0,0 +1,435 @@ +#!/usr/bin/env python3 +"""Publishing from inside a subscriber callback must not freeze the interpreter. + +Two independent defects combined here: + +1. hiroz declared every subscriber as a zenoh-ext ``AdvancedSubscriber``, which + runs the user callback while holding a non-reentrant ``std::sync::Mutex``. + Zenoh delivers same-session samples synchronously on the publishing thread, so + a publish from inside a callback re-entered that mutex and deadlocked. +2. ``ZPublisher.publish`` held the GIL across the blocking publish, so the + deadlock froze the *whole* interpreter -- no exception, no traceback, only an + external kill. + +Defect 1 is fixed differently per durability, and both paths are covered here: + +* ``Volatile`` (the ROS 2 default) gets a plain zenoh subscriber. Samples that + arrived over a transport run inline on the zenoh RX worker; samples published + by the delivering thread itself are handed to a per-subscriber + ``hiroz-sub-drain`` thread. That second case is what a re-entrant publish + hits, so it *iterates* rather than recursing and needs no depth cap. +* ``TransientLocal`` keeps the advanced subscriber -- its reordering guarantees + need that mutex -- so *every* sample moves off the delivering thread onto the + drain thread. + +Either way the drain thread is a *Rust* thread and is therefore invisible to +``threading.enumerate()``; the leak test below reads ``/proc/self/task`` +instead, which is the only view Python has of it. + +Everything here is deadline-guarded so a regression fails the suite rather than +wedging it. That guard is itself only meaningful because of fix 2 -- with the GIL +held across the block, no watchdog and no deadline thread can run at all, which +is what ``test_interpreter_stays_alive_during_reentrant_publish`` pins down. +""" + +import gc +import itertools +import os +import threading +import time + +import hiroz_py +import pytest +from hiroz_py import std_msgs + +# Generous relative to the work done (a handful of intra-process publishes). +SCENARIO_TIMEOUT = 20.0 +HOPS = 4 + +# How far the self-feeding loop must run to prove it iterates rather than +# recurses. Two orders of magnitude past the point where the old inline path +# need, and well past the stack depth a recursive implementation survives. +ITERATION_TARGET = 2000 + +DRAIN_THREAD_NAME = "hiroz-sub-drain" + +DURABILITIES = ["volatile", "transient_local"] + +_topic_seq = itertools.count() + + +def _topic(stem): + """A fresh topic per scenario. + + TransientLocal replays history to late-joining subscribers, so reusing a + topic across scenarios would leak samples from one into the next and make a + hop chain appear to complete that never actually ran. + """ + return f"/reentrant_py_{stem}_{next(_topic_seq)}" + + +def _qos(durability): + return hiroz_py.QosProfile(durability=durability) + + +def _run_on_deadline(fn, timeout=SCENARIO_TIMEOUT): + """Run ``fn`` on a daemon thread and assert it returned within ``timeout``. + + The seeding publish is what blocks in the unfixed build, so it is what has to + be time-boxed. A daemon thread means a genuine deadlock leaves a stuck thread + behind but still lets the process exit with a real failure -- provided the + GIL is released across the block, which is fix 2. + """ + error = [] + + def target(): + try: + fn() + except BaseException as exc: # noqa: BLE001 - re-raised below + error.append(exc) + + driver = threading.Thread(target=target, daemon=True) + driver.start() + driver.join(timeout) + assert not driver.is_alive(), ( + f"publish() from inside a subscriber callback did not return within " + f"{timeout}s - re-entrant publish deadlocked" + ) + if error: + raise error[0] + + +# -------------------------------------------------------------------------- +# Preconditions -- fail loudly rather than skipping silently. +# +# The sibling Rust suite (crates/hiroz-tests/tests/reentrant_service.rs) is +# gated behind ``#![cfg(feature = "ros-msgs")]`` and reports a green "0 passed" +# when the feature is off. Nothing here may be able to do that: if the +# environment cannot express a scenario, that is a failure, not a skip. +# -------------------------------------------------------------------------- + + +def test_preconditions(context): + """The APIs every scenario below depends on must exist and behave.""" + assert hasattr(hiroz_py, "QosProfile"), "hiroz_py.QosProfile is missing" + + for durability in DURABILITIES: + profile = _qos(durability) + assert durability in repr(profile).lower(), ( + f"QosProfile(durability={durability!r}) did not round-trip: {profile!r}" + ) + + node = context.create_node("reentrant_precond").with_namespace("/test").build() + assert hasattr(node, "destroy_subscriber"), ( + "ZNode.destroy_subscriber is missing - the drain-thread leak test cannot " + "drop subscribers and must not be reported as passing" + ) + + # The leak test can only see the Rust drain thread through procfs. + assert os.path.isdir("/proc/self/task"), ( + "/proc/self/task is unavailable - the drain-thread leak test cannot run " + "on this platform and must not be reported as passing" + ) + + +# -------------------------------------------------------------------------- +# The 6-cell matrix: {volatile, transient_local} +# x {same-topic, two-topic cycle, unbounded feedback}. +# -------------------------------------------------------------------------- + + +@pytest.mark.parametrize("durability", DURABILITIES) +def test_same_topic_republish(context, durability): + """A callback republishing to its own topic must complete the hop chain.""" + node = ( + context.create_node(f"reentrant_same_{durability}") + .with_namespace("/test") + .build() + ) + topic = _topic("same") + qos = _qos(durability) + publisher = node.create_publisher(topic, std_msgs.String, qos=qos) + + seen = [] + done = threading.Event() + + def on_message(msg): + seen.append(msg.data) + hop = int(msg.data) + if hop < HOPS: + # The re-entrant publish. Before the fix this never returned. + publisher.publish(std_msgs.String(data=str(hop + 1))) + else: + done.set() + + node.create_subscriber(topic, std_msgs.String, qos=qos, callback=on_message) + time.sleep(0.5) + + _run_on_deadline(lambda: publisher.publish(std_msgs.String(data="0"))) + + assert done.wait(SCENARIO_TIMEOUT), f"only got {seen!r}, expected {HOPS + 1} hops" + assert seen == [str(i) for i in range(HOPS + 1)], seen + + +@pytest.mark.parametrize("durability", DURABILITIES) +def test_two_topic_cycle(context, durability): + """A -> B -> A across two subscribers must not deadlock either. + + Distinct from the same-topic case: the re-entrant publish lands on a + *different* subscriber, so this exercises the delivering thread rather than + one subscriber's own state. + """ + node = ( + context.create_node(f"reentrant_cycle_{durability}") + .with_namespace("/test") + .build() + ) + topic_a = _topic("cycle_a") + topic_b = _topic("cycle_b") + qos = _qos(durability) + pub_a = node.create_publisher(topic_a, std_msgs.String, qos=qos) + pub_b = node.create_publisher(topic_b, std_msgs.String, qos=qos) + + seen = [] + done = threading.Event() + + def hop_handler(which, forward_to): + def handler(msg): + seen.append((which, msg.data)) + hop = int(msg.data) + if hop < HOPS: + forward_to.publish(std_msgs.String(data=str(hop + 1))) + else: + done.set() + + return handler + + node.create_subscriber( + topic_a, std_msgs.String, qos=qos, callback=hop_handler("a", pub_b) + ) + node.create_subscriber( + topic_b, std_msgs.String, qos=qos, callback=hop_handler("b", pub_a) + ) + time.sleep(0.5) + + _run_on_deadline(lambda: pub_a.publish(std_msgs.String(data="0"))) + + assert done.wait(SCENARIO_TIMEOUT), f"cycle stalled after {seen!r}" + assert len(seen) == HOPS + 1, seen + assert [d for _, d in seen] == [str(i) for i in range(HOPS + 1)], seen + assert [t for t, _ in seen] == ["a", "b", "a", "b", "a"][: HOPS + 1], seen + + +@pytest.mark.parametrize("durability", DURABILITIES) +def test_self_feeding_loop_iterates(context, durability): + """A callback that republishes every time must loop, not recurse. + + This is the Python-level view of the change that made the ``afor`` + benchmark's ``intra`` cell expressible. hiroz used to deliver a same-session + sample inline on the publishing thread, so this shape was recursion: it grew + the stack; before this fix it deadlocked on the first hop rather than + it. Delivery now enqueues and returns, so each hop starts from a flat stack + and the loop runs for as long as it is fed. + + The assertion is the inverse of the one it replaces: the loop must run far + *past* the old cap. + """ + node = ( + context.create_node(f"reentrant_loop_{durability}") + .with_namespace("/test") + .build() + ) + topic = _topic("loop") + qos = _qos(durability) + publisher = node.create_publisher(topic, std_msgs.String, qos=qos) + + counter = itertools.count() + delivered = [] + reached = threading.Event() + + def on_message(msg): + n = next(counter) + delivered.append(n) + if n >= ITERATION_TARGET: + # Stop feeding so the test ends on its own; the loop's *ability* to + # keep going is what is under test. + reached.set() + return + publisher.publish(std_msgs.String(data=str(n))) + + node.create_subscriber(topic, std_msgs.String, qos=qos, callback=on_message) + time.sleep(0.5) + + _run_on_deadline(lambda: publisher.publish(std_msgs.String(data="0"))) + + assert reached.wait(SCENARIO_TIMEOUT), ( + f"self-feeding loop stalled at {len(delivered)} deliveries, target " + f"{ITERATION_TARGET}. Under the old inline dispatch this caps out at 16." + ) + assert len(delivered) >= ITERATION_TARGET, len(delivered) + + +# -------------------------------------------------------------------------- +# Fix 2: the interpreter must stay alive. +# -------------------------------------------------------------------------- + + +def test_interpreter_stays_alive_during_reentrant_publish(context): + """A watchdog thread must keep running across a re-entrant publish. + + This is the user-facing half of the fix and the only part a Rust test cannot + check. ``ZPublisher.publish`` used to hold the GIL across the blocking zenoh + publish, so a deadlock there did not stall one thread -- it stalled every + thread. No exception, no traceback, no way to diagnose it from inside the + process. ``py.allow_threads`` downgrades that to a single blocked thread. + + Failure mode on the unfixed build: this test does not fail, it *hangs the + whole process*, because the assertions below also need the GIL. That is + precisely the symptom, and it is why the unfixed baseline has to be measured + under an external wall-clock timeout rather than trusted to self-report. + """ + node = context.create_node("reentrant_watchdog").with_namespace("/test").build() + topic = _topic("watchdog") + publisher = node.create_publisher(topic, std_msgs.String) + + ticks = [] + stop = threading.Event() + + def watchdog(): + while not stop.is_set(): + ticks.append(time.monotonic()) + time.sleep(0.02) + + def on_message(msg): + hop = int(msg.data) + if hop < HOPS: + publisher.publish(std_msgs.String(data=str(hop + 1))) + + node.create_subscriber(topic, std_msgs.String, callback=on_message) + time.sleep(0.5) + + wd = threading.Thread(target=watchdog, daemon=True) + wd.start() + time.sleep(0.2) + before = len(ticks) + assert before > 0, "watchdog never started - the detector is broken, not the code" + + _run_on_deadline(lambda: publisher.publish(std_msgs.String(data="0"))) + + time.sleep(0.3) + after = len(ticks) + stop.set() + wd.join(5.0) + + assert after > before, ( + "the watchdog thread made no progress across the re-entrant publish - " + "the GIL was held across the block and the whole interpreter went dark" + ) + + # No single stall longer than a generous multiple of the tick interval. + gaps = [b - a for a, b in zip(ticks, ticks[1:])] + worst = max(gaps) if gaps else 0.0 + assert worst < 2.0, ( + f"watchdog stalled for {worst:.2f}s during the publish - the interpreter " + "was frozen even if it eventually recovered" + ) + + +# -------------------------------------------------------------------------- +# The TransientLocal dispatcher thread must not leak. +# -------------------------------------------------------------------------- + + +def _drain_threads(): + """Count live ``hiroz-sub-drain`` threads via procfs. + + ``threading.enumerate()`` cannot see these: they are Rust threads spawned by + ``CallbackDispatcher::spawn`` and were never registered with CPython. + ``/proc/self/task/*/comm`` is the only way Python can observe them, so it is + the only honest detector for a leak. + """ + total = 0 + for tid in os.listdir("/proc/self/task"): + try: + with open(f"/proc/self/task/{tid}/comm") as handle: + if handle.read().strip() == DRAIN_THREAD_NAME: + total += 1 + except OSError: + continue # thread exited between listdir and open + return total + + +def _settle(deadline, target=None): + """Wait up to ``deadline`` for the drain count to reach ``target``.""" + gc.collect() + end = time.monotonic() + deadline + current = _drain_threads() + while time.monotonic() < end: + if target is not None and current == target: + return current + time.sleep(0.2) + current = _drain_threads() + return current + + +def test_transient_local_dispatcher_threads_do_not_leak(context): + """Create and destroy TransientLocal subscribers; the count must return. + + Volatile *callback* subscribers also carry a dispatcher now (for the + session-local delivery path), so the leak concern is no longer + TransientLocal-only -- but TransientLocal spawns one unconditionally, which + makes it the tighter test of the same teardown path. + """ + node = context.create_node("reentrant_leak").with_namespace("/test").build() + qos = _qos("transient_local") + + baseline_py = len(threading.enumerate()) + baseline_drain = _settle(deadline=3.0) + + subs_per_cycle = 4 + cycles = 3 + peaks = [] + + for cycle in range(cycles): + subs = [ + node.create_subscriber( + _topic(f"leak_{cycle}_{i}"), + std_msgs.String, + qos=qos, + callback=lambda _msg: None, + ) + for i in range(subs_per_cycle) + ] + time.sleep(0.5) + + peak = _drain_threads() + peaks.append(peak) + # Precondition: if the dispatcher never spawned, this test is not + # exercising the TransientLocal path and must fail rather than pass. + assert peak >= baseline_drain + subs_per_cycle, ( + f"cycle {cycle}: expected at least {subs_per_cycle} new " + f"{DRAIN_THREAD_NAME!r} threads over a baseline of {baseline_drain}, " + f"saw {peak}. The TransientLocal dispatcher path was not exercised - " + "this test proves nothing in this state." + ) + + for sub in subs: + node.destroy_subscriber(sub) + subs.clear() + + settled = _settle(deadline=15.0, target=baseline_drain) + assert settled == baseline_drain, ( + f"cycle {cycle}: {settled - baseline_drain} {DRAIN_THREAD_NAME!r} " + f"thread(s) leaked after destroying {subs_per_cycle} subscribers " + f"(baseline {baseline_drain}, peak {peak})" + ) + + # Repeated cycles must not ratchet upward. + assert max(peaks) - min(peaks) <= 1, ( + f"drain-thread peak drifted across cycles: {peaks} - threads accumulate" + ) + + final_py = len(threading.enumerate()) + assert final_py <= baseline_py, ( + f"Python-level threads grew from {baseline_py} to {final_py}" + ) diff --git a/crates/hiroz-tests/tests/dispatch_backpressure.rs b/crates/hiroz-tests/tests/dispatch_backpressure.rs new file mode 100644 index 000000000..d49a6380e --- /dev/null +++ b/crates/hiroz-tests/tests/dispatch_backpressure.rs @@ -0,0 +1,331 @@ +//! The plain-path callback dispatcher must honour `KEEP_LAST(depth)`. +//! +//! A callback subscriber cannot run user code on the publishing thread — that is +//! what made re-entrant publishes recurse — so locally published samples are +//! handed to `CallbackDispatcher`'s queue and delivered on its own thread. That +//! queue is the subscriber's history buffer, and it must behave like the one the +//! queue-mode path uses (`BoundedQueue`): retain the last `depth` undelivered +//! samples and drop the *oldest* on overflow. +//! +//! Two properties are asserted, one per direction of the bound: +//! +//! * `KeepLast(n)` drops, and drops the oldest — an unbounded queue delivers the +//! *first* `n` after the burst instead of the last, so the assertion on +//! contents (not merely on count) fails if the bound is removed. +//! * `KeepAll` does not drop — the lossless branch is still reachable. +//! +//! Both are made deterministic by blocking the drain thread inside the first +//! callback until the whole burst has been published, so the race between +//! producer and drain thread is removed rather than slept on. + +mod common; + +use std::{ + num::NonZeroUsize, + sync::{ + Arc, Condvar, Mutex, + atomic::{AtomicBool, Ordering}, + mpsc, + }, + thread, + time::{Duration, Instant}, +}; + +use common::{TestRouter, create_hiroz_context_with_endpoint}; +use hiroz::{ + Builder, TypeHash, + qos::{QosDurability, QosHistory, QosProfile}, + ros_msg::MessageTypeInfo, +}; +use serde::{Deserialize, Serialize}; +use serial_test::serial; + +/// Budget for one scenario. Generous relative to the work done. +const SCENARIO_TIMEOUT: Duration = Duration::from_secs(30); + +/// History depth under test. Small enough that a 50-message burst overflows it +/// many times over, large enough that an off-by-one would be visible. +const DEPTH: usize = 4; + +/// Messages published after the drain thread is parked. Counter values are +/// `0..BURST`. +const BURST: u64 = 50; + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +struct Seq { + counter: u64, +} + +impl MessageTypeInfo for Seq { + fn type_name() -> &'static str { + "test_msgs::msg::dds_::Seq_" + } + + fn type_hash() -> TypeHash { + TypeHash::zero() + } +} + +impl hiroz::ros_msg::WithTypeInfo for Seq {} + +impl hiroz::msg::ZMessage for Seq { + type Serdes = hiroz::msg::SerdeCdrSerdes; +} + +/// Run `scenario` on its own thread and fail (rather than hang) if it does not +/// finish within [`SCENARIO_TIMEOUT`]. +fn run_with_deadline(name: &'static str, scenario: impl FnOnce() + Send + 'static) { + let (tx, rx) = mpsc::channel(); + thread::Builder::new() + .name(name.to_string()) + .spawn(move || { + scenario(); + let _ = tx.send(()); + }) + .expect("failed to spawn scenario thread"); + + if rx.recv_timeout(SCENARIO_TIMEOUT).is_err() { + panic!("`{name}` did not finish within {SCENARIO_TIMEOUT:?}"); + } +} + +/// A latch the first callback parks on until the publisher releases it. +#[derive(Default)] +struct Latch { + open: Mutex, + changed: Condvar, +} + +impl Latch { + fn wait(&self) { + let mut open = self.open.lock().expect("latch poisoned"); + while !*open { + open = self.changed.wait(open).expect("latch poisoned"); + } + } + + fn release(&self) { + *self.open.lock().expect("latch poisoned") = true; + self.changed.notify_all(); + } +} + +/// Publish `0..BURST` into a callback subscriber whose drain thread is parked on +/// the first message, then release it and return everything the callback saw. +/// +/// Returns once the delivered sequence has been quiet for 500 ms, so the caller +/// asserts on a settled result rather than on a snapshot mid-drain. +fn burst_through_dispatcher( + endpoint: &str, + node: &str, + topic: &str, + history: QosHistory, + durability: QosDurability, +) -> Vec { + let qos = QosProfile { + history, + durability, + ..Default::default() + }; + + let ctx = create_hiroz_context_with_endpoint(endpoint).expect("failed to create context"); + let node = ctx + .create_node(node) + .build() + .expect("failed to create node"); + + let publisher = node + .create_pub::(topic) + .with_qos(qos) + .build() + .expect("failed to create publisher"); + + let seen = Arc::new(Mutex::new(Vec::::new())); + let latch = Arc::new(Latch::default()); + let parked = Arc::new(AtomicBool::new(false)); + + let cb_seen = seen.clone(); + let cb_latch = latch.clone(); + let cb_parked = parked.clone(); + let _sub = node + .create_sub::(topic) + .with_qos(qos) + .build_with_callback(move |msg: Seq| { + cb_seen.lock().expect("seen poisoned").push(msg.counter); + if msg.counter == 0 { + // The drain thread is now provably inside the callback, so every + // subsequent publish lands in the queue rather than racing it. + cb_parked.store(true, Ordering::SeqCst); + cb_latch.wait(); + } + }) + .expect("failed to create subscriber"); + + thread::sleep(Duration::from_millis(300)); + + publisher + .publish(&Seq { counter: 0 }) + .expect("seed publish failed"); + + let deadline = Instant::now() + Duration::from_secs(5); + while !parked.load(Ordering::SeqCst) { + assert!( + Instant::now() < deadline, + "the drain thread never entered the callback" + ); + thread::sleep(Duration::from_millis(10)); + } + + for counter in 1..BURST { + publisher + .publish(&Seq { counter }) + .expect("burst publish failed"); + } + + latch.release(); + + // Settle: stop once the delivered sequence has not changed for 500 ms. + let mut last = 0usize; + let mut stable_since = Instant::now(); + let deadline = Instant::now() + Duration::from_secs(10); + loop { + let len = seen.lock().expect("seen poisoned").len(); + if len != last { + last = len; + stable_since = Instant::now(); + } else if stable_since.elapsed() >= Duration::from_millis(500) { + break; + } + assert!(Instant::now() < deadline, "delivery never settled"); + thread::sleep(Duration::from_millis(25)); + } + + seen.lock().expect("seen poisoned").clone() +} + +/// `KeepLast(DEPTH)` must retain the newest `DEPTH` undelivered samples. +/// +/// The first message is already out of the queue (the drain thread is parked +/// holding it), so the settled sequence is `[0]` followed by the last `DEPTH` of +/// the burst. Asserting the *values* rather than the count is what makes this a +/// drop-**oldest** detector: an unbounded queue yields `0,1,2,…` and a +/// drop-newest queue yields `0,1,2,3,4`. +#[test] +#[serial] +fn keep_last_drops_the_oldest_local_samples() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("keep_last_drops_oldest", move || { + let delivered = burst_through_dispatcher( + &endpoint, + "dispatch_keep_last", + "/dispatch_keep_last", + QosHistory::KeepLast(NonZeroUsize::new(DEPTH).unwrap()), + QosDurability::Volatile, + ); + + let mut expected = vec![0]; + expected.extend((BURST - DEPTH as u64)..BURST); + + assert_eq!( + delivered, expected, + "a KeepLast({DEPTH}) callback subscriber must deliver the seed plus the \ + last {DEPTH} of the burst, dropping the oldest in between" + ); + }); +} + +/// `KeepAll` must not drop: the queue is unbounded on that profile, matching +/// `BoundedQueue::new(usize::MAX)` on the queue-mode path. +#[test] +#[serial] +fn keep_all_delivers_every_local_sample() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("keep_all_lossless", move || { + let delivered = burst_through_dispatcher( + &endpoint, + "dispatch_keep_all", + "/dispatch_keep_all", + QosHistory::KeepAll, + QosDurability::Volatile, + ); + + let expected: Vec = (0..BURST).collect(); + assert_eq!( + delivered, expected, + "a KeepAll callback subscriber must not drop" + ); + }); +} + +/// The **advanced** path must honour `KeepLast(DEPTH)` too. +/// +/// `TransientLocal` routes through `AdvancedSubscriber`, a different construction +/// site with its own `CallbackDispatcher::spawn` call. That site passed +/// `DISPATCH_UNBOUNDED` unconditionally, so a `KeepLast` subscriber got an +/// unbounded queue — and because the advanced path's shim enqueues *remote* +/// samples too, that replaced zenoh's transport backpressure with unbounded +/// growth. Nothing detected it: every other test in this file takes the plain +/// path. +/// +/// Asserting values rather than a count makes this a drop-**oldest** detector in +/// both directions, exactly as on the plain path: unbounded yields `0,1,2,…`, +/// drop-newest yields the first `DEPTH + 1`. +#[test] +#[serial] +fn transient_local_keep_last_drops_the_oldest_local_samples() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("transient_local_keep_last_drops_oldest", move || { + let delivered = burst_through_dispatcher( + &endpoint, + "dispatch_tl_keep_last", + "/dispatch_tl_keep_last", + QosHistory::KeepLast(NonZeroUsize::new(DEPTH).unwrap()), + QosDurability::TransientLocal, + ); + + let mut expected = vec![0]; + expected.extend((BURST - DEPTH as u64)..BURST); + + assert_eq!( + delivered, expected, + "a TransientLocal KeepLast({DEPTH}) callback subscriber must deliver the \ + seed plus the last {DEPTH} of the burst; an unbounded advanced queue \ + delivers all {BURST}" + ); + }); +} + +/// `KeepAll` on the advanced path stays lossless. +/// +/// This is the half of the previous test's argument that survives: replaying +/// history and recovering missed samples is what `TransientLocal` is for, so a +/// profile that asks to keep everything must keep everything. Bounding by the +/// declared depth must not have collapsed this case too. +#[test] +#[serial] +fn transient_local_keep_all_delivers_every_local_sample() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("transient_local_keep_all_lossless", move || { + let delivered = burst_through_dispatcher( + &endpoint, + "dispatch_tl_keep_all", + "/dispatch_tl_keep_all", + QosHistory::KeepAll, + QosDurability::TransientLocal, + ); + + let expected: Vec = (0..BURST).collect(); + assert_eq!( + delivered, expected, + "a TransientLocal KeepAll callback subscriber must not drop" + ); + }); +} diff --git a/crates/hiroz-tests/tests/panic_guard_inline.rs b/crates/hiroz-tests/tests/panic_guard_inline.rs new file mode 100644 index 000000000..4a2985b1b --- /dev/null +++ b/crates/hiroz-tests/tests/panic_guard_inline.rs @@ -0,0 +1,266 @@ +//! The plain path's INLINE branch must survive a panicking callback. +//! +//! `CallbackDispatcher`'s drain loop wraps the user callback in `catch_unwind`. +//! The plain path does not always use that loop: `local_only_shim` enqueues only +//! the samples the delivering thread published itself, and calls the handler +//! inline for every other one — which is every sample that arrived over a +//! transport. +//! +//! That inline branch is the ROS 2 default. `qos_needs_advanced` is true only +//! for `TransientLocal` durability and the default is `Volatile`, so a default +//! subscriber takes the plain path and its inter-process traffic is exactly the +//! inline case. Before the guard was added there, a panicking callback unwound +//! out of hiroz and into a zenoh receive worker. +//! +//! # How this test avoids passing for the wrong reason +//! +//! Three properties, each of which would otherwise let it pass vacuously. +//! +//! 1. **The sample must be remote.** The shim's selector is +//! [`local_publish_active`], a THREAD-LOCAL depth counter — not session and +//! not process locality. A publish on the same thread takes the `if` arm and +//! exercises the already-guarded queue path, proving nothing about the +//! branch under test. Publisher and subscriber therefore use separate +//! contexts through a router, and the test asserts the delivering thread was +//! a zenoh receive worker rather than `hiroz-sub-drain`. +//! +//! 2. **Panics must unwind.** `catch_unwind` can only return `Err` where they +//! do, and this workspace's `[profile.opt]` sets `panic = "abort"`. The file +//! is gated so it SKIPS visibly on such a profile instead of passing. +//! +//! 3. **A positive control.** `delivery_continues_without_a_panic` runs the +//! identical shape with the panic removed. Without it, "the later samples +//! never arrived" would be indistinguishable from "delivery never worked +//! here at all". +//! +//! # Revert direction +//! +//! Remove the `catch_unwind` from `local_only_shim`'s `else` arm in +//! `crates/hiroz/src/pubsub.rs` — leaving a bare `handler(sample)` — and +//! `a_panicking_callback_on_a_remote_sample_does_not_stop_delivery` must fail. +//! +//! Part of #296 (tag G7). + +#![cfg(panic = "unwind")] + +mod common; + +use std::{ + collections::BTreeSet, + sync::{ + Arc, Mutex, + atomic::{AtomicU64, Ordering}, + mpsc, + }, + thread, + time::Duration, +}; + +use common::{TestRouter, create_hiroz_context_with_endpoint}; +use hiroz::{Builder, TypeHash, ros_msg::MessageTypeInfo}; +use serde::{Deserialize, Serialize}; +use serial_test::serial; + +/// Budget for one scenario. Generous relative to the work done. +const SCENARIO_TIMEOUT: Duration = Duration::from_secs(30); + +/// Samples published. The callback panics on exactly one of them. +const COUNT: u64 = 12; + +/// Which sample panics. Chosen so that several arrive before it and several +/// after, making "delivery stopped" distinguishable from "delivery never +/// started". +const PANIC_AT: u64 = 4; + +/// How long to wait for the delivered set to settle before asserting. +const SETTLE: Duration = Duration::from_secs(3); + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +struct Seq { + counter: u64, +} + +impl MessageTypeInfo for Seq { + fn type_name() -> &'static str { + "test_msgs::msg::dds_::Seq_" + } + + fn type_hash() -> TypeHash { + TypeHash::zero() + } +} + +impl hiroz::ros_msg::WithTypeInfo for Seq {} + +impl hiroz::msg::ZMessage for Seq { + type Serdes = hiroz::msg::SerdeCdrSerdes; +} + +/// Run `scenario` on its own thread and fail rather than hang. +fn run_with_deadline(name: &'static str, scenario: impl FnOnce() + Send + 'static) { + let (tx, rx) = mpsc::channel(); + thread::Builder::new() + .name(name.to_string()) + .spawn(move || { + scenario(); + let _ = tx.send(()); + }) + .expect("failed to spawn scenario thread"); + + if rx.recv_timeout(SCENARIO_TIMEOUT).is_err() { + panic!("`{name}` did not finish within {SCENARIO_TIMEOUT:?}"); + } +} + +/// What one run observed. +struct Observed { + /// Counter values the callback received. + seen: BTreeSet, + /// Names of the threads the callback ran on. + threads: BTreeSet, +} + +/// Publish `0..COUNT` from a SEPARATE context through `router`, into a +/// default-QoS callback subscriber. `panic_at` makes the callback panic on that +/// counter value; `None` is the positive control. +fn deliver_remote(endpoint: &str, topic: &str, panic_at: Option) -> Observed { + let sub_ctx = create_hiroz_context_with_endpoint(endpoint).expect("subscriber context"); + let pub_ctx = create_hiroz_context_with_endpoint(endpoint).expect("publisher context"); + + let sub_node = sub_ctx + .create_node("panic_guard_sub") + .build() + .expect("sub node"); + let pub_node = pub_ctx + .create_node("panic_guard_pub") + .build() + .expect("pub node"); + + let seen: Arc>> = Arc::new(Mutex::new(BTreeSet::new())); + let threads: Arc>> = Arc::new(Mutex::new(BTreeSet::new())); + let received = Arc::new(AtomicU64::new(0)); + + let c_seen = seen.clone(); + let c_threads = threads.clone(); + let c_recv = received.clone(); + + // Default QoS: Volatile, so this is the PLAIN path. + let _sub = sub_node + .create_sub::(topic) + .build_with_callback(move |msg: Seq| { + c_threads + .lock() + .expect("threads poisoned") + .insert(thread::current().name().unwrap_or("").to_string()); + c_seen.lock().expect("seen poisoned").insert(msg.counter); + c_recv.fetch_add(1, Ordering::Relaxed); + if panic_at == Some(msg.counter) { + panic!( + "deliberate panic in a subscriber callback, counter={}", + msg.counter + ); + } + }) + .expect("subscriber"); + + let zpub = pub_node + .create_pub::(topic) + .build() + .expect("publisher"); + + // Let discovery settle so the first samples are not lost to a race. A lost + // early sample would weaken the "delivery started" half of the assertion. + thread::sleep(Duration::from_millis(500)); + + for counter in 0..COUNT { + zpub.publish(&Seq { counter }).expect("publish"); + thread::sleep(Duration::from_millis(50)); + } + + // Wait for the delivered set to go quiet rather than sleeping a fixed time. + let mut last = u64::MAX; + for _ in 0..(SETTLE.as_millis() / 100) { + let now = received.load(Ordering::Relaxed); + if now == last { + break; + } + last = now; + thread::sleep(Duration::from_millis(100)); + } + + let seen = seen.lock().expect("seen poisoned").clone(); + let threads = threads.lock().expect("threads poisoned").clone(); + Observed { seen, threads } +} + +/// GUARD 1: the samples must have arrived over a transport, on a zenoh receive +/// worker. If they came in on `hiroz-sub-drain` the enqueue arm was taken and +/// the already-guarded path was tested instead. +fn assert_took_the_inline_branch(o: &Observed, what: &str) { + assert!( + !o.threads.is_empty(), + "{what}: the callback never ran, so nothing was tested" + ); + assert!( + !o.threads.contains("hiroz-sub-drain"), + "{what}: delivery ran on `hiroz-sub-drain`, so the sample took \ + local_only_shim's enqueue arm. That is the queue path, which the drain \ + loop already guards — this test measured the wrong branch. \ + Threads seen: {:?}", + o.threads + ); +} + +#[test] +#[serial] +fn a_panicking_callback_on_a_remote_sample_does_not_stop_delivery() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("panic_guard_inline", move || { + let o = deliver_remote(&endpoint, "/panic_guard_inline", Some(PANIC_AT)); + + assert_took_the_inline_branch(&o, "panic run"); + + assert!( + o.seen.contains(&PANIC_AT), + "the panicking sample itself never arrived, so the panic never \ + happened and nothing was tested. Seen: {:?}", + o.seen + ); + + let after: Vec = o.seen.iter().copied().filter(|c| *c > PANIC_AT).collect(); + assert!( + !after.is_empty(), + "no sample after the panic was delivered: the panicking callback \ + stopped the subscriber. The inline branch of `local_only_shim` \ + needs the same `catch_unwind` the drain loop has. Seen: {:?}", + o.seen + ); + }); +} + +/// GUARD 3: the positive control. Without it, a failure above could mean +/// "delivery never worked in this configuration" rather than "the panic +/// stopped it". +#[test] +#[serial] +fn delivery_continues_without_a_panic() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("panic_guard_control", move || { + let o = deliver_remote(&endpoint, "/panic_guard_control", None); + + assert_took_the_inline_branch(&o, "control run"); + + let after: Vec = o.seen.iter().copied().filter(|c| *c > PANIC_AT).collect(); + assert!( + !after.is_empty(), + "the control run delivered nothing after counter {PANIC_AT}, so this \ + configuration cannot deliver at all and the panic test above proves \ + nothing either way. Seen: {:?}", + o.seen + ); + }); +} diff --git a/crates/hiroz-tests/tests/reentrant_publish.rs b/crates/hiroz-tests/tests/reentrant_publish.rs new file mode 100644 index 000000000..4dc4407e5 --- /dev/null +++ b/crates/hiroz-tests/tests/reentrant_publish.rs @@ -0,0 +1,921 @@ +//! Re-entrant publish from inside a subscriber callback must not deadlock. +//! +//! Every hiroz subscriber used to be declared as a zenoh-ext `AdvancedSubscriber`, +//! whose sample callback runs the user closure while a non-reentrant +//! `std::sync::Mutex` guard is alive (`sub_callback` takes `zlock!(statesref)`, +//! then `handle_sample` calls the callback under it). Zenoh core dispatches +//! samples published on the same session synchronously on the publishing thread, +//! so publishing from inside a callback re-enters that mutex on the very thread +//! that already holds it — a deterministic self-deadlock. +//! +//! Each scenario runs on a dedicated thread and reports through a channel, so a +//! hang fails the test instead of wedging the harness. +//! +//! Two families of scenario run here, because hiroz uses two different subscriber +//! implementations depending on QoS: +//! +//! * **Volatile** (the ROS 2 default) declares a plain zenoh subscriber, whose +//! callback runs inline with no lock held. +//! * **TransientLocal** must keep the `AdvancedSubscriber` — history replay and +//! sample-miss recovery live there — so its user callback is moved onto a +//! dedicated delivery thread and zenoh-ext only ever gets an enqueue-only shim. +//! +//! Both must survive the same re-entrancy scenarios, and with the same +//! observable semantics, so each scenario is written twice. + +mod common; + +use std::{ + sync::{ + Arc, Mutex, + atomic::{AtomicBool, AtomicUsize, Ordering}, + mpsc, + }, + thread, + time::{Duration, Instant}, +}; + +use common::{TestRouter, create_hiroz_context_with_endpoint}; +use hiroz::{ + Builder, TypeHash, + pubsub::CloseOutcome, + qos::{QosDurability, QosProfile}, + ros_msg::MessageTypeInfo, +}; +use serde::{Deserialize, Serialize}; +use serial_test::serial; + +/// QoS that forces hiroz down the zenoh-ext `AdvancedSubscriber` path. +/// +/// This is the case the `qos_needs_advanced` gate deliberately does *not* cover: +/// the advanced entity is genuinely needed here, so the deadlock has to be solved +/// rather than side-stepped. +fn transient_local() -> QosProfile { + QosProfile { + durability: QosDurability::TransientLocal, + ..Default::default() + } +} + +/// Budget for one scenario. Generous relative to the work done (a handful of +/// intra-process publishes) — anything slower than this is a hang, not slowness. +const SCENARIO_TIMEOUT: Duration = Duration::from_secs(20); + +/// How long to wait for the expected number of deliveries once seeded. +const DELIVERY_TIMEOUT: Duration = Duration::from_secs(5); + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +struct Tick { + counter: u64, +} + +impl MessageTypeInfo for Tick { + fn type_name() -> &'static str { + "test_msgs::msg::dds_::Tick_" + } + + fn type_hash() -> TypeHash { + TypeHash::zero() + } +} + +impl hiroz::ros_msg::WithTypeInfo for Tick {} + +impl hiroz::msg::ZMessage for Tick { + type Serdes = hiroz::msg::SerdeCdrSerdes; +} + +/// Run `scenario` on its own thread and fail (rather than hang) if it does not +/// finish within [`SCENARIO_TIMEOUT`]. +/// +/// On timeout the worker thread is deliberately left running: it is blocked on a +/// mutex it can never acquire and cannot be unwound. Each scenario owns its own +/// context and router, and the test binary exits shortly after. +fn run_with_deadline(name: &'static str, scenario: impl FnOnce() + Send + 'static) { + let (tx, rx) = mpsc::channel(); + thread::Builder::new() + .name(name.to_string()) + .spawn(move || { + scenario(); + let _ = tx.send(()); + }) + .expect("failed to spawn scenario thread"); + + if rx.recv_timeout(SCENARIO_TIMEOUT).is_err() { + panic!( + "`{name}` did not finish within {SCENARIO_TIMEOUT:?} — \ + re-entrant publish from a subscriber callback deadlocked" + ); + } +} + +/// Block until `seen` reaches `expected`, or fail with what actually arrived. +fn await_deliveries(seen: &AtomicUsize, expected: usize) { + let deadline = Instant::now() + DELIVERY_TIMEOUT; + while seen.load(Ordering::SeqCst) < expected { + assert!( + Instant::now() < deadline, + "only {} of {expected} messages delivered", + seen.load(Ordering::SeqCst) + ); + thread::sleep(Duration::from_millis(20)); + } +} + +/// A callback that publishes back onto its own topic must make progress. +/// +/// The minimal shape of the bug: one subscriber, one publisher, one session. The +/// callback re-publishes for a bounded number of hops so the scenario terminates +/// by construction rather than relying on the fix to bound it. +#[test] +#[serial] +fn callback_republishing_on_same_topic_does_not_deadlock() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("same_topic", move || { + const HOPS: u64 = 4; + + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_same_topic") + .build() + .expect("failed to create node"); + + let publisher = Arc::new( + node.create_pub::("/reentrant_same") + .build() + .expect("failed to create publisher"), + ); + let seen = Arc::new(AtomicUsize::new(0)); + + let cb_pub = publisher.clone(); + let cb_seen = seen.clone(); + let _sub = node + .create_sub::("/reentrant_same") + .build_with_callback(move |msg: Tick| { + cb_seen.fetch_add(1, Ordering::SeqCst); + if msg.counter < HOPS { + // Re-entrant publish on the same session, from inside the + // subscriber callback. This is what used to deadlock. + cb_pub + .publish(&Tick { + counter: msg.counter + 1, + }) + .expect("re-entrant publish failed"); + } + }) + .expect("failed to create subscriber"); + + thread::sleep(Duration::from_millis(300)); + publisher + .publish(&Tick { counter: 0 }) + .expect("seed publish failed"); + + await_deliveries(&seen, HOPS as usize + 1); + }); +} + +/// Two callbacks publishing to each other's topic must make progress. +/// +/// Distinct from the same-topic case: two *different* subscribers, and so two +/// different zenoh-ext state mutexes, form a cycle — the shape a real ping-pong +/// node pair has. The deadlock still occurs, one frame later. +#[test] +#[serial] +fn callback_cycle_across_two_topics_does_not_deadlock() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("two_topics", move || { + const HOPS: u64 = 4; + + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_cycle") + .build() + .expect("failed to create node"); + + let pub_a = Arc::new( + node.create_pub::("/reentrant_a") + .build() + .expect("failed to create publisher a"), + ); + let pub_b = Arc::new( + node.create_pub::("/reentrant_b") + .build() + .expect("failed to create publisher b"), + ); + let seen = Arc::new(AtomicUsize::new(0)); + + let cb_pub_b = pub_b.clone(); + let cb_seen_a = seen.clone(); + let _sub_a = node + .create_sub::("/reentrant_a") + .build_with_callback(move |msg: Tick| { + cb_seen_a.fetch_add(1, Ordering::SeqCst); + if msg.counter < HOPS { + cb_pub_b + .publish(&Tick { + counter: msg.counter + 1, + }) + .expect("publish a->b failed"); + } + }) + .expect("failed to create subscriber a"); + + let cb_pub_a = pub_a.clone(); + let cb_seen_b = seen.clone(); + let _sub_b = node + .create_sub::("/reentrant_b") + .build_with_callback(move |msg: Tick| { + cb_seen_b.fetch_add(1, Ordering::SeqCst); + if msg.counter < HOPS { + cb_pub_a + .publish(&Tick { + counter: msg.counter + 1, + }) + .expect("publish b->a failed"); + } + }) + .expect("failed to create subscriber b"); + + thread::sleep(Duration::from_millis(300)); + pub_a + .publish(&Tick { counter: 0 }) + .expect("seed publish failed"); + + await_deliveries(&seen, HOPS as usize + 1); + }); +} + +/// A self-feeding callback loop must *iterate*, not recurse. +/// +/// This is the detector for the whole point of routing session-local delivery +/// through the dispatcher. A callback that republishes to its own topic used to +/// be reached inline from inside `publish()`, so the loop was recursion: it grew +/// the stack, so before this fix it could not be expressed as a loop at all and +/// drop samples past the cap to avoid a `SIGSEGV`. +/// +/// Now the sample is enqueued and the callback runs on the dispatcher thread, so +/// each iteration returns to a flat stack before the next begins. The loop runs +/// indefinitely at constant stack depth and constant queue depth. +/// +/// The assertion is deliberately the *inverse* of the old one: it requires the +/// loop to exceed the old cap by orders of magnitude. Against the pre-fix +/// sources this fails at 16. +#[test] +#[serial] +fn self_feeding_callback_loop_iterates_without_a_depth_cap() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("unbounded_loop", move || { + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_unbounded") + .build() + .expect("failed to create node"); + + let publisher = Arc::new( + node.create_pub::("/reentrant_unbounded") + .build() + .expect("failed to create publisher"), + ); + let seen = Arc::new(AtomicUsize::new(0)); + + let cb_pub = publisher.clone(); + let cb_seen = seen.clone(); + let _sub = node + .create_sub::("/reentrant_unbounded") + .build_with_callback(move |msg: Tick| { + // Stop feeding well before the deadline so the test ends on its + // own; the loop's *ability* to keep going is what is under test. + if cb_seen.fetch_add(1, Ordering::SeqCst) >= ITERATION_TARGET { + return; + } + let _ = cb_pub.publish(&Tick { + counter: msg.counter + 1, + }); + }) + .expect("failed to create subscriber"); + + thread::sleep(Duration::from_millis(300)); + publisher + .publish(&Tick { counter: 0 }) + .expect("seed publish failed"); + + await_deliveries(&seen, ITERATION_TARGET); + + let delivered = seen.load(Ordering::SeqCst); + assert!( + delivered >= ITERATION_TARGET, + "self-feeding loop stalled at {delivered} deliveries (target {ITERATION_TARGET}). \ + Under the old inline dispatch this deadlocks on the first hop." + ); + }); +} + +/// How far the self-feeding loops must run to prove they are iterative. +/// +/// Two orders of magnitude above the old depth cap of 16, and well past the +/// stack depth a recursive implementation survives: with the fix reverted these +/// same tests die with `fatal runtime error: stack overflow` rather than merely +/// falling short of the target. +/// +/// Deliberately not larger. At 20_000 these loops saturate a core for long +/// enough to perturb the *next* test binary in a sequential run — `parameter_tests` +/// began failing three service calls on timeout purely from the load, with no +/// code path in common (verified: the dispatcher branch is never taken in that +/// suite). The property under test is "does it iterate at all", which 2_000 +/// settles just as conclusively as 20_000. +const ITERATION_TARGET: usize = 2_000; + +// --------------------------------------------------------------------------- +// TransientLocal variants +// +// These exercise the `AdvancedSubscriber` path, which cannot avoid holding its +// state mutex across the user callback: `handle_sample` interleaves +// `callback.call(sample)` with mutation of `last_delivered`/`pending_samples`. +// Instead the callback zenoh-ext receives only enqueues, and hiroz runs the real +// callback on a dedicated delivery thread with no lock held. +// --------------------------------------------------------------------------- + +/// TransientLocal counterpart of +/// [`callback_republishing_on_same_topic_does_not_deadlock`]. +#[test] +#[serial] +fn transient_local_callback_republishing_on_same_topic_does_not_deadlock() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("tl_same_topic", move || { + const HOPS: u64 = 4; + + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_tl_same_topic") + .build() + .expect("failed to create node"); + + let publisher = Arc::new( + node.create_pub::("/reentrant_tl_same") + .with_qos(transient_local()) + .build() + .expect("failed to create publisher"), + ); + let seen = Arc::new(AtomicUsize::new(0)); + + let cb_pub = publisher.clone(); + let cb_seen = seen.clone(); + let _sub = node + .create_sub::("/reentrant_tl_same") + .with_qos(transient_local()) + .build_with_callback(move |msg: Tick| { + cb_seen.fetch_add(1, Ordering::SeqCst); + if msg.counter < HOPS { + cb_pub + .publish(&Tick { + counter: msg.counter + 1, + }) + .expect("re-entrant publish failed"); + } + }) + .expect("failed to create subscriber"); + + thread::sleep(Duration::from_millis(300)); + publisher + .publish(&Tick { counter: 0 }) + .expect("seed publish failed"); + + await_deliveries(&seen, HOPS as usize + 1); + }); +} + +/// TransientLocal counterpart of +/// [`callback_cycle_across_two_topics_does_not_deadlock`]. +/// +/// Each subscriber owns its own delivery thread, so this also covers a callback +/// on one delivery thread publishing into a *different* advanced subscriber. +#[test] +#[serial] +fn transient_local_callback_cycle_across_two_topics_does_not_deadlock() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("tl_two_topics", move || { + const HOPS: u64 = 4; + + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_tl_cycle") + .build() + .expect("failed to create node"); + + let pub_a = Arc::new( + node.create_pub::("/reentrant_tl_a") + .with_qos(transient_local()) + .build() + .expect("failed to create publisher a"), + ); + let pub_b = Arc::new( + node.create_pub::("/reentrant_tl_b") + .with_qos(transient_local()) + .build() + .expect("failed to create publisher b"), + ); + let seen = Arc::new(AtomicUsize::new(0)); + + let cb_pub_b = pub_b.clone(); + let cb_seen_a = seen.clone(); + let _sub_a = node + .create_sub::("/reentrant_tl_a") + .with_qos(transient_local()) + .build_with_callback(move |msg: Tick| { + cb_seen_a.fetch_add(1, Ordering::SeqCst); + if msg.counter < HOPS { + cb_pub_b + .publish(&Tick { + counter: msg.counter + 1, + }) + .expect("publish a->b failed"); + } + }) + .expect("failed to create subscriber a"); + + let cb_pub_a = pub_a.clone(); + let cb_seen_b = seen.clone(); + let _sub_b = node + .create_sub::("/reentrant_tl_b") + .with_qos(transient_local()) + .build_with_callback(move |msg: Tick| { + cb_seen_b.fetch_add(1, Ordering::SeqCst); + if msg.counter < HOPS { + cb_pub_a + .publish(&Tick { + counter: msg.counter + 1, + }) + .expect("publish b->a failed"); + } + }) + .expect("failed to create subscriber b"); + + thread::sleep(Duration::from_millis(300)); + pub_a + .publish(&Tick { counter: 0 }) + .expect("seed publish failed"); + + await_deliveries(&seen, HOPS as usize + 1); + }); +} + +/// TransientLocal counterpart of +/// [`self_feeding_callback_loop_iterates_without_a_depth_cap`]. +/// +/// This path always ran the callback on the dispatcher thread, so it was already +/// iterative — the depth cap was carried across the queue only to keep its +/// behaviour identical to the inline path's. Now that the inline path is gone +/// there is nothing to stay identical to, and both paths iterate. Asserting it +/// here keeps the two in step. +#[test] +#[serial] +fn transient_local_self_feeding_callback_loop_iterates() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("tl_unbounded_loop", move || { + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_tl_unbounded") + .build() + .expect("failed to create node"); + + let publisher = Arc::new( + node.create_pub::("/reentrant_tl_unbounded") + .with_qos(transient_local()) + .build() + .expect("failed to create publisher"), + ); + let seen = Arc::new(AtomicUsize::new(0)); + + let cb_pub = publisher.clone(); + let cb_seen = seen.clone(); + let _sub = node + .create_sub::("/reentrant_tl_unbounded") + .with_qos(transient_local()) + .build_with_callback(move |msg: Tick| { + if cb_seen.fetch_add(1, Ordering::SeqCst) >= ITERATION_TARGET { + return; + } + let _ = cb_pub.publish(&Tick { + counter: msg.counter + 1, + }); + }) + .expect("failed to create subscriber"); + + thread::sleep(Duration::from_millis(300)); + publisher + .publish(&Tick { counter: 0 }) + .expect("seed publish failed"); + + await_deliveries(&seen, ITERATION_TARGET); + + let delivered = seen.load(Ordering::SeqCst); + assert!( + delivered >= ITERATION_TARGET, + "self-feeding loop stalled at {delivered} deliveries (target {ITERATION_TARGET})" + ); + }); +} + +/// The `intra` closed loop: two topics, two callbacks, one session, each +/// callback feeding the other. This is the shape the `afor` benchmark's `intra` +/// cell drives, and the reason that cell could not produce a number. +/// +/// Under inline session-local dispatch this is not a loop at all — it is mutual +/// recursion on one thread, so it either overflows the stack or, with the depth +/// cap, stops dead at 16. Under dispatcher delivery each hop returns to a flat +/// stack, so the loop runs as long as it is fed. +#[test] +#[serial] +fn intra_closed_loop_runs_iteratively() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("intra_closed_loop", move || { + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("intra_closed_loop") + .build() + .expect("failed to create node"); + + let hello_pub = Arc::new( + node.create_pub::("/intra_hello") + .build() + .expect("failed to create hello publisher"), + ); + let world_pub = Arc::new( + node.create_pub::("/intra_world") + .build() + .expect("failed to create world publisher"), + ); + let seen = Arc::new(AtomicUsize::new(0)); + + // hello -> world + let to_world = world_pub.clone(); + let hello_seen = seen.clone(); + let _hello_sub = node + .create_sub::("/intra_hello") + .build_with_callback(move |msg: Tick| { + if hello_seen.fetch_add(1, Ordering::SeqCst) >= ITERATION_TARGET { + return; + } + let _ = to_world.publish(&Tick { + counter: msg.counter + 1, + }); + }) + .expect("failed to create hello subscriber"); + + // world -> hello + let to_hello = hello_pub.clone(); + let world_seen = seen.clone(); + let _world_sub = node + .create_sub::("/intra_world") + .build_with_callback(move |msg: Tick| { + if world_seen.fetch_add(1, Ordering::SeqCst) >= ITERATION_TARGET { + return; + } + let _ = to_hello.publish(&Tick { + counter: msg.counter + 1, + }); + }) + .expect("failed to create world subscriber"); + + thread::sleep(Duration::from_millis(300)); + hello_pub + .publish(&Tick { counter: 0 }) + .expect("seed publish failed"); + + await_deliveries(&seen, ITERATION_TARGET); + + let delivered = seen.load(Ordering::SeqCst); + assert!( + delivered >= ITERATION_TARGET, + "intra closed loop stalled at {delivered} hops (target {ITERATION_TARGET}) — \ + this is the `afor` intra cell's failure mode" + ); + }); +} + +/// The delivery thread must not reorder samples. +/// +/// `AdvancedSubscriber` exists to deliver samples in source order and to recover +/// missed ones; a handoff that reordered them would defeat its entire purpose. +/// The shim enqueues from inside `handle_sample` — i.e. under zenoh-ext's state +/// mutex, in exactly the order zenoh-ext chose to deliver — and a single thread +/// pops FIFO, so the observed order must be the publish order. +/// +/// **Ordering, not completeness.** The burst is far deeper than the declared +/// `KeepLast` depth, so the queue drops the oldest — exactly as +/// `rmw_zenoh_cpp`'s `add_new_message` does (`rmw_subscription_data.cpp`: +/// `size() >= adapted_qos_profile.depth` → `pop_front()`, with no +/// `TransientLocal` exemption). Asserting the delivered values were *all* +/// published would therefore assert a promise no RMW makes. +/// +/// Asserting a strictly increasing subsequence is the stronger test anyway: it +/// catches reordering whether or not anything was dropped, whereas an equality +/// check conflates the two failures. `keep_all_delivers_every_local_sample` in +/// `dispatch_backpressure.rs` covers losslessness on the profile that promises +/// it. +#[test] +#[serial] +fn transient_local_delivery_preserves_order() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("tl_ordering", move || { + const COUNT: u64 = 500; + + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_tl_order") + .build() + .expect("failed to create node"); + + let publisher = node + .create_pub::("/reentrant_tl_order") + .with_qos(transient_local()) + .build() + .expect("failed to create publisher"); + + let received: Arc>> = Arc::new(Mutex::new(Vec::new())); + let seen = Arc::new(AtomicUsize::new(0)); + + let cb_received = received.clone(); + let cb_seen = seen.clone(); + let _sub = node + .create_sub::("/reentrant_tl_order") + .with_qos(transient_local()) + .build_with_callback(move |msg: Tick| { + cb_received.lock().unwrap().push(msg.counter); + cb_seen.fetch_add(1, Ordering::SeqCst); + }) + .expect("failed to create subscriber"); + + thread::sleep(Duration::from_millis(300)); + for counter in 0..COUNT { + publisher + .publish(&Tick { counter }) + .expect("publish failed"); + } + + // Settle rather than wait for a fixed count: with a bounded queue the + // delivered total is a property of scheduling, not of the publish count. + let mut last = 0usize; + let mut stable_since = Instant::now(); + let deadline = Instant::now() + DELIVERY_TIMEOUT; + loop { + let len = seen.load(Ordering::SeqCst); + if len != last { + last = len; + stable_since = Instant::now(); + } else if stable_since.elapsed() >= Duration::from_millis(300) { + break; + } + assert!(Instant::now() < deadline, "delivery never settled"); + thread::sleep(Duration::from_millis(25)); + } + + let received = received.lock().unwrap(); + assert!( + !received.is_empty(), + "nothing was delivered — the scenario proved nothing" + ); + assert!( + received.windows(2).all(|w| w[0] < w[1]), + "delivery thread reordered samples: {received:?}" + ); + assert!( + received.iter().all(|&c| c < COUNT), + "delivered a counter that was never published: {received:?}" + ); + }); +} + +/// `close(deadline)` must shut the delivery thread down — no leak, no hang. +/// +/// `Drop` alone no longer does this. It guarantees that no *new* callback +/// starts, not that a running one has finished, so the thread winds down +/// asynchronously and this assertion would race. `close` is the barrier, and +/// this test is the migration BC5 asks callers who relied on the old behaviour +/// to make. +/// +/// The sentinel is owned by the user callback, which the delivery thread owns in +/// turn, so its `Drop` firing proves the thread actually exited and released the +/// closure rather than being detached and left running. +#[test] +#[serial] +fn transient_local_subscriber_drop_shuts_down_delivery_thread() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("tl_drop", move || { + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_tl_drop") + .build() + .expect("failed to create node"); + + let publisher = node + .create_pub::("/reentrant_tl_drop") + .with_qos(transient_local()) + .build() + .expect("failed to create publisher"); + + /// Fires when the delivery thread drops the user callback. + struct Sentinel(Arc); + impl Drop for Sentinel { + fn drop(&mut self) { + self.0.store(true, Ordering::SeqCst); + } + } + + let callback_dropped = Arc::new(AtomicBool::new(false)); + let sentinel = Sentinel(callback_dropped.clone()); + let seen = Arc::new(AtomicUsize::new(0)); + let cb_seen = seen.clone(); + + let sub = node + .create_sub::("/reentrant_tl_drop") + .with_qos(transient_local()) + .build_with_callback(move |_msg: Tick| { + // Keep the sentinel owned by the callback. + let _ = &sentinel; + cb_seen.fetch_add(1, Ordering::SeqCst); + }) + .expect("failed to create subscriber"); + + thread::sleep(Duration::from_millis(300)); + publisher + .publish(&Tick { counter: 0 }) + .expect("seed publish failed"); + await_deliveries(&seen, 1); + + assert!( + !callback_dropped.load(Ordering::SeqCst), + "callback was dropped while the subscriber was still alive" + ); + + // `close` is the barrier; `Drop` is not. If it hung, the deadline guard + // fails the test. + let outcome = sub.close(Duration::from_secs(5)); + assert_eq!( + outcome, + CloseOutcome::Joined, + "close() timed out: the drain thread did not exit within the deadline" + ); + + assert!( + callback_dropped.load(Ordering::SeqCst), + "close() reported Joined but the callback was never dropped — \ + the thread is leaked" + ); + }); +} + +/// Dropping the subscriber *from inside its own callback* must not deadlock. +/// +/// That drop runs on the delivery thread itself. `Drop` detaches rather than +/// joining, so there is no self-join to detect and no special case guarding it — +/// the thread observes `closed` and winds itself down once the callback returns. +/// +/// Calling `close` here instead would report `TimedOut` at once: a thread cannot +/// wait for itself. +#[test] +#[serial] +fn transient_local_subscriber_dropped_inside_its_own_callback_does_not_deadlock() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("tl_self_drop", move || { + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_tl_self_drop") + .build() + .expect("failed to create node"); + + let publisher = node + .create_pub::("/reentrant_tl_self_drop") + .with_qos(transient_local()) + .build() + .expect("failed to create publisher"); + + // The subscriber hands itself to its own callback, which drops it. + let slot: Arc>>> = Arc::new(Mutex::new(None)); + let dropped = Arc::new(AtomicBool::new(false)); + + let cb_slot = slot.clone(); + let cb_dropped = dropped.clone(); + let sub = node + .create_sub::("/reentrant_tl_self_drop") + .with_qos(transient_local()) + .build_with_callback(move |_msg: Tick| { + // Runs on the delivery thread. The drop detaches; the thread + // exits once this callback returns. + let taken = cb_slot.lock().unwrap().take(); + drop(taken); + cb_dropped.store(true, Ordering::SeqCst); + }) + .expect("failed to create subscriber"); + + *slot.lock().unwrap() = Some(Box::new(sub)); + + thread::sleep(Duration::from_millis(300)); + publisher + .publish(&Tick { counter: 0 }) + .expect("seed publish failed"); + + let deadline = Instant::now() + DELIVERY_TIMEOUT; + while !dropped.load(Ordering::SeqCst) { + assert!( + Instant::now() < deadline, + "subscriber was never dropped from its own callback" + ); + thread::sleep(Duration::from_millis(20)); + } + }); +} + +/// `async_publish` must also hand session-local delivery to the drain thread. +/// +/// Every other scenario in this file drives the *synchronous* `publish`, so the +/// async path's guard was unpinned. It is guarded differently, and the +/// difference is load-bearing: a thread-local held across an `.await` would be +/// observed on whatever thread resumed the task, so `async_publish` scopes +/// `LocalPublishGuard` to the `into_future()` call rather than the await. That +/// is only sufficient because zenoh resolves a put eagerly there -- +/// `IntoFuture for PublicationBuilder<_, PublicationBuilderPut>` is +/// `std::future::ready(self.wait())` (zenoh 1.9.0, `api/builders/publisher.rs`). +/// +/// That is an upstream implementation detail, not a documented contract. Should +/// zenoh ever make the put lazy, it would move outside the guard, session-local +/// samples would again be dispatched inline on the publishing thread, and every +/// deadlock this file exists to prevent would return on the async path with no +/// other test noticing. +/// +/// The detector is thread identity rather than a hang, so it fails immediately +/// and for a legible reason instead of timing out: with the guard covering the +/// put, the callback runs on `hiroz-sub-drain`; without it, on the publishing +/// thread. +#[test] +#[serial] +fn async_publish_delivers_off_the_publishing_thread() { + let router = TestRouter::new(); + let endpoint = router.endpoint().to_string(); + + run_with_deadline("async_publish_off_thread", move || { + let ctx = create_hiroz_context_with_endpoint(&endpoint).expect("failed to create context"); + let node = ctx + .create_node("reentrant_async_publish") + .build() + .expect("failed to create node"); + + let publisher = node + .create_pub::("/reentrant_async") + .build() + .expect("failed to create publisher"); + + let (tx, rx) = mpsc::channel(); + let _sub = node + .create_sub::("/reentrant_async") + .build_with_callback(move |_msg: Tick| { + let _ = tx.send(thread::current().id()); + }) + .expect("failed to create subscriber"); + + thread::sleep(Duration::from_millis(300)); + + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("failed to build runtime"); + let publishing_thread = runtime.block_on(async { + publisher + .async_publish(&Tick { counter: 0 }) + .await + .expect("async publish failed"); + thread::current().id() + }); + + let callback_thread = rx + .recv_timeout(DELIVERY_TIMEOUT) + .expect("async_publish produced no delivery"); + + assert_ne!( + callback_thread, publishing_thread, + "async_publish delivered the sample inline on the publishing thread: \ + the local-publish guard did not cover the put, so a callback that \ + publishes would recurse instead of iterate" + ); + }); +} diff --git a/crates/hiroz/src/common.rs b/crates/hiroz/src/common.rs index 5cf360b85..48203e5c2 100644 --- a/crates/hiroz/src/common.rs +++ b/crates/hiroz/src/common.rs @@ -19,6 +19,21 @@ pub(crate) enum DataHandler { } impl DataHandler { + /// Whether `handle` runs *user* code on the delivering thread. + /// + /// Only [`DataHandler::Callback`] does. The queue variants enqueue and + /// return. zenoh's own `FifoChannel` handler has the same structure. The + /// user's code therefore runs on whatever thread calls `recv()`. Nothing on + /// the delivery thread can re-enter hiroz. + /// + /// This function deliberately does not count `QueueWithNotifier`'s notifier. + /// The notifier is the rmw layer's wait-set wake. It must run promptly on + /// the delivery thread. It does not call back into hiroz. Issue #290 tracks + /// that nothing enforces this last claim. + pub(crate) fn runs_user_code(&self) -> bool { + matches!(self, DataHandler::Callback(_)) + } + pub(crate) fn handle(&self, data: T) { match self { DataHandler::Queue(queue) => { diff --git a/crates/hiroz/src/ffi/publisher.rs b/crates/hiroz/src/ffi/publisher.rs index 6b78a6389..ea0739816 100644 --- a/crates/hiroz/src/ffi/publisher.rs +++ b/crates/hiroz/src/ffi/publisher.rs @@ -36,6 +36,15 @@ impl RawPublisher { } pub fn publish_bytes(&self, data: &[u8]) -> Result<(), zenoh::Error> { + // The four `ZPub` publish paths take this same guard, for the same + // reason. `local_only_shim` hands a sample to the drain thread only + // while `LOCAL_PUBLISH_DEPTH` is set. Without the guard, hiroz does not + // mark a same-process raw publish as local. The subscriber's callback + // then runs inline on this thread. A callback that publishes back into + // its own topic recurses until the stack is gone (#249). + // `rmw-zenoh-rs` publishes through this path, so an unguarded publish + // here keeps the deadlock reachable from every rmw user. + let _local = crate::pubsub::LocalPublishGuard::enter(); self.inner .put(data) .attachment(self.new_attachment()) diff --git a/crates/hiroz/src/ffi/subscriber.rs b/crates/hiroz/src/ffi/subscriber.rs index 35294812b..3c573058c 100644 --- a/crates/hiroz/src/ffi/subscriber.rs +++ b/crates/hiroz/src/ffi/subscriber.rs @@ -2,14 +2,15 @@ use super::node::{CNode, get_node_ref}; use super::qos::CQosProfile; use super::{ErrorCode, cstr_to_str}; use std::ffi::c_char; -use zenoh_ext::AdvancedSubscriber; + +use crate::pubsub::SubscriberHandle; /// Callback type for receiving messages pub type MessageCallback = extern "C" fn(user_data: usize, data: *const u8, len: usize); /// Raw subscriber wrapper that keeps the zenoh subscriber alive pub struct RawSubscriber { - pub inner: AdvancedSubscriber<()>, + pub inner: SubscriberHandle, } /// Opaque subscriber handle for FFI diff --git a/crates/hiroz/src/node.rs b/crates/hiroz/src/node.rs index 90435aeed..950edfdd8 100644 --- a/crates/hiroz/src/node.rs +++ b/crates/hiroz/src/node.rs @@ -657,7 +657,10 @@ impl ZNode { { use crate::{ entity::{EndpointEntity, EndpointKind}, - pubsub::apply_transient_local_sub, + pubsub::{ + CallbackDispatcher, SubscriberHandle, apply_transient_local_sub, dispatch_capacity, + qos_needs_advanced, + }, topic_name, }; use zenoh_ext::AdvancedSubscriberBuilderExt; @@ -681,15 +684,49 @@ impl ZNode { }; let topic_ke = self.keyexpr_format.topic_key_expr(&entity)?; - let subscriber = self - .session - .declare_subscriber((*topic_ke).clone()) - .callback(move |sample| { - let payload = sample.payload().to_bytes(); - callback(&payload); - }) - .advanced(); - let subscriber = apply_transient_local_sub(subscriber, &entity.qos).wait()?; + let raw_callback = Arc::new(move |sample: zenoh::sample::Sample| { + let payload = sample.payload().to_bytes(); + callback(&payload); + }); + + // Same rule as the typed path: use zenoh-ext only when the QoS asks for + // advanced features. This callback is user (FFI) code on either arm. It + // therefore never runs on a thread that is inside a hiroz publish — see + // `pubsub::CallbackDispatcher`. + let subscriber = if qos_needs_advanced(&entity.qos) { + let dispatcher = CallbackDispatcher::spawn( + &qualified_topic, + raw_callback, + dispatch_capacity(&entity.qos), + )?; + let subscriber = self + .session + .declare_subscriber((*topic_ke).clone()) + .callback(dispatcher.always_shim()); + SubscriberHandle::Advanced { + subscriber: Box::new( + apply_transient_local_sub(subscriber.advanced(), &entity.qos).wait()?, + ), + // Always `Some` here: this path exists to deliver to an FFI + // callback, which is user code by definition. + dispatcher: Some(dispatcher), + } + } else { + let dispatcher = CallbackDispatcher::spawn( + &qualified_topic, + raw_callback.clone(), + dispatch_capacity(&entity.qos), + )?; + let subscriber = self + .session + .declare_subscriber((*topic_ke).clone()) + .callback(dispatcher.local_only_shim(raw_callback)) + .wait()?; + SubscriberHandle::Plain { + subscriber, + dispatcher: Some(dispatcher), + } + }; Ok(crate::ffi::subscriber::RawSubscriber { inner: subscriber }) } diff --git a/crates/hiroz/src/prelude.rs b/crates/hiroz/src/prelude.rs index c66e74783..84e37f0b1 100644 --- a/crates/hiroz/src/prelude.rs +++ b/crates/hiroz/src/prelude.rs @@ -23,7 +23,7 @@ pub use crate::Builder; /// Core runtime types. pub use crate::context::{ZContext, ZContextBuilder}; pub use crate::node::ZNode; -pub use crate::pubsub::{ZPub, ZSub}; +pub use crate::pubsub::{CloseOutcome, ZPub, ZSub}; pub use crate::service::{RequestId, ServiceReply, ServiceRequest, ZClient, ZServer}; pub use crate::action::server::{Accepted, Executing, Requested}; diff --git a/crates/hiroz/src/pubsub.rs b/crates/hiroz/src/pubsub.rs index b03bb30b6..600864713 100644 --- a/crates/hiroz/src/pubsub.rs +++ b/crates/hiroz/src/pubsub.rs @@ -30,6 +30,770 @@ use zenoh_ext::{ /// Matches rmw_zenoh_cpp's `SAMPLE_MISS_DETECTION_HEARTBEAT_PERIOD`. const SAMPLE_MISS_HEARTBEAT_PERIOD: Duration = Duration::from_millis(500); +thread_local! { + /// How many hiroz publish calls are currently on this thread's stack. + /// + /// A non-zero count means this thread produced any sample it is *about* to + /// deliver. It produced that sample synchronously, from inside `put`. See + /// [`local_publish_active`]. + static LOCAL_PUBLISH_DEPTH: std::cell::Cell = const { std::cell::Cell::new(0) }; +} + +/// RAII marker that hiroz holds for the duration of a publish. +/// +/// hiroz has no single choke point for publishing. Each of `ZPub`'s four +/// publish paths enters a guard itself and holds it across the zenoh `put`: +/// [`ZPub::publish`], [`ZPub::async_publish`], [`ZPub::publish_serialized`] and +/// [`ZPub::publish_sample`]. +/// +/// A fifth publish path added later must do the same. If it does not, +/// session-local delivery on that path runs inline on the publishing thread. +/// The deadlock this guard prevents then comes back (#249). +/// +/// The guard counts nesting rather than sets a flag. A callback can run on a +/// thread that is already inside a publish. A publish issued from that callback +/// then restores the correct depth when its own guard drops. +pub(crate) struct LocalPublishGuard; + +impl LocalPublishGuard { + pub(crate) fn enter() -> Self { + LOCAL_PUBLISH_DEPTH.with(|d| d.set(d.get() + 1)); + Self + } +} + +impl Drop for LocalPublishGuard { + fn drop(&mut self) { + LOCAL_PUBLISH_DEPTH.with(|d| d.set(d.get().saturating_sub(1))); + } +} + +/// Whether this thread is currently inside a hiroz publish. +/// +/// hiroz can reach a subscriber callback in two ways. This function +/// discriminates between them. That discrimination is what makes re-entrancy +/// structurally impossible without a cost on the inter-process path. +/// +/// * **true** — zenoh delivers the sample *synchronously on the publishing +/// thread*. It does this for same-session delivery: `Session::resolve_put` +/// drops the session lock and calls the local callbacks inline. It also does +/// this for two sessions that share one process with a direct in-process +/// route. That route is `send_push_consume` -> `route_data` -> the peer +/// session's callbacks, with no thread hop. A user callback that runs here +/// can publish into its own topic graph and *recurse* instead of iterate +/// (#249). So hiroz +/// enqueues and returns on this path, exactly as zenoh's own `FifoChannel` +/// handler does. The callback then runs on the dispatcher thread. +/// +/// * **false** — the sample arrived over a transport. A zenoh RX worker +/// delivers it (`ZRuntime::RX`, threads named `rx-N`). That worker is never +/// an application thread and never inside a hiroz publish. Nothing can +/// re-enter, so the callback runs inline. The inter-process path pays only +/// this thread-local read. +/// +/// This function deliberately keys on the *publishing thread*, not on zenoh's +/// `Locality`. Zenoh still delivers a `Locality::Remote`-tagged sample inline on +/// the publisher's thread when that sample crosses two sessions inside one +/// process. An `allowed_origin(SessionLocal)` split would therefore miss it. The +/// thread is the reliable signal. The origin is not. +fn local_publish_active() -> bool { + LOCAL_PUBLISH_DEPTH.with(|d| d.get()) != 0 +} + +/// Backlog size at which an *unbounded* [`CallbackDispatcher`] first warns. +/// +/// The threshold doubles after each warning. A persistently slow callback +/// therefore does not flood the log. +/// +/// The check excludes bounded dispatchers explicitly. It does not assume that +/// this value is out of their reach: a `KeepLast(1024)` subscriber has exactly +/// this capacity. Such a subscriber would otherwise report that its queue is +/// lossless immediately before it drops a sample. Bounded dispatchers warn on +/// drops instead. +const DISPATCH_BACKLOG_WARN_AT: usize = 1024; + +/// The capacity that makes a [`CallbackDispatcher`] unbounded, i.e. lossless. +pub(crate) const DISPATCH_UNBOUNDED: usize = usize::MAX; + +/// The dispatcher capacity implied by a subscriber's history QoS. +/// +/// This matches what [`ZSubBuilder::build`] gives the queue-mode +/// [`BoundedQueue`]. `KeepLast(depth)` keeps `depth` samples. `KeepAll` keeps +/// everything. A callback subscriber and a queue subscriber with the same QoS +/// therefore retain the same number of undelivered samples. Retention does not +/// depend on which hiroz API the caller chose. +/// +/// The two are *not* the same expression. The difference applies only to a zero +/// depth. Zero is the rmw spelling of "system default". [`QosProfile`] cannot +/// produce it, but it can arrive over the wire. This function floors it at 1; +/// the queue path passes it through. +/// +/// Retention still agrees. [`BoundedQueue::push`] evicts before it inserts +/// (`len >= capacity` → `pop_front`, then `push_back`), so a capacity of 0 also +/// retains exactly one sample. Only the bookkeeping differs: at capacity 0 +/// every push reports a drop, including the first push into an empty queue. +pub(crate) fn dispatch_capacity(qos: &hiroz_protocol::qos::QosProfile) -> usize { + match qos.history { + QosHistory::KeepLast(depth) => depth.max(1), + QosHistory::KeepAll => DISPATCH_UNBOUNDED, + } +} + +struct DispatchState { + /// Samples awaiting delivery, in the order zenoh decided to deliver them. + pending: std::collections::VecDeque, + /// [`CallbackDispatcher::drop`] sets this: stop delivering and exit. The + /// dispatcher discards queued but undelivered samples -- see + /// [`DispatchQueue::dequeue`] for why teardown does not drain them. + closed: bool, + /// Next backlog length that triggers a warning. Unbounded queues only. + warn_at: usize, + /// Samples discarded because the queue was at capacity. + dropped: u64, + /// Next `dropped` total that triggers a warning. + warn_dropped_at: u64, +} + +struct DispatchQueue { + state: Mutex, + ready: std::sync::Condvar, + /// Set by the drain thread as its last action, once the drain loop has + /// exited and the final user callback has returned. + /// + /// [`CallbackDispatcher::close`] waits on this rather than on + /// [`std::thread::JoinHandle::join`], because `join` has no timed form on + /// stable Rust and `close` must honour a deadline. It lives in its own + /// mutex so that a waiting `close` never contends with delivery. + finished: Mutex, + /// Signalled with `finished`. + done: std::sync::Condvar, + topic: String, + /// Maximum number of undelivered samples the queue retains. + /// [`DISPATCH_UNBOUNDED`] means lossless. Any smaller capacity drops the + /// *oldest* sample on overflow, exactly as [`BoundedQueue::push`] does. See + /// [`CallbackDispatcher`]'s "Backpressure" section for which path gets + /// which. + capacity: usize, +} + +impl DispatchQueue { + fn lock(&self) -> std::sync::MutexGuard<'_, DispatchState> { + // A panicking user callback must not wedge the subscriber: the queue + // holds no invariant that a partial mutation could break. + self.state.lock().unwrap_or_else(|e| e.into_inner()) + } + + /// The shim callback that hiroz hands to zenoh. + /// + /// This function may run with zenoh-ext's state mutex held (advanced path). + /// It may also run on the publishing thread inside `put` (local path). It + /// must therefore never publish and never block. + fn enqueue(&self, sample: Sample) { + let (backlog, dropped) = { + let mut state = self.lock(); + if state.closed { + return; + } + + // Drop the oldest sample. Never drop the newest. Never drop the + // incoming one. `BoundedQueue::push` makes the same choice. ROS + // `KEEP_LAST(depth)` describes the same choice. A bounded queue + // that *blocked* here would re-create the original deadlock — see + // the type's docs. + let dropped = if state.pending.len() >= self.capacity { + state.pending.pop_front(); + state.dropped = state.dropped.saturating_add(1); + if state.dropped >= state.warn_dropped_at { + state.warn_dropped_at = state.dropped.saturating_mul(2); + Some(state.dropped) + } else { + None + } + } else { + None + }; + + state.pending.push_back(sample); + let len = state.pending.len(); + // Unbounded queues only. A bounded queue *can* reach + // `DISPATCH_BACKLOG_WARN_AT`: nothing stops a subscriber from + // declaring `KeepLast(1024)` or deeper. It would then log that the + // queue is lossless and costs only memory. A bounded queue does the + // opposite. Bounded queues report drops instead. That warning is + // immediately below. It is the accurate one. + let backlog = if self.capacity == DISPATCH_UNBOUNDED && len >= state.warn_at { + state.warn_at = len.saturating_mul(2); + Some(len) + } else { + None + }; + (backlog, dropped) + }; + self.ready.notify_one(); + if let Some(len) = backlog { + warn!( + topic = %self.topic, + backlog = len, + "subscriber delivery backlog is growing; the callback is slower than the \ + publish rate. This queue is lossless, so the backlog costs memory." + ); + } + if let Some(total) = dropped { + warn!( + topic = %self.topic, + dropped = total, + capacity = self.capacity, + "subscriber delivery queue is full; dropping the oldest undelivered sample. \ + The callback is slower than the publish rate — raise the history depth or \ + make the callback cheaper." + ); + } + } + + /// Blocks until a sample arrives, or until the queue closes. A closed queue + /// returns `None`, which ends the drain loop. + /// + /// This function checks `closed` **before** `pending`. That order is the + /// difference between a bounded and an unbounded teardown. + /// + /// An earlier version drained the backlog first. `drop(subscriber)` then + /// ran a user callback for every queued sample before it returned. On the + /// unbounded (TransientLocal) path that cost is + /// `backlog × callback_duration` with no ceiling. A 1 kHz publisher against + /// a 5 ms callback leaves about 30 000 samples queued after 30 s, so the + /// drop blocks silently for minutes. It can also block *forever* if a + /// callback waits on anything the dropping thread must supply. + /// + /// Dropping a subscriber means "stop delivering to me". The dispatcher + /// therefore discards undelivered samples. It does not force them through a + /// callback the caller has already disposed of. Destroying an rclcpp + /// subscription does the same. Teardown costs at most one in-flight + /// callback. + fn dequeue(&self) -> Option { + let mut state = self.lock(); + loop { + if state.closed { + return None; + } + if let Some(entry) = state.pending.pop_front() { + return Some(entry); + } + state = self.ready.wait(state).unwrap_or_else(|e| e.into_inner()); + } + } +} + +/// Runs a subscriber's user callback on a dedicated thread. A FIFO queue feeds +/// that thread. +/// +/// This type is hiroz's equivalent of zenoh's `FifoChannel` handler. It is also +/// the equivalent of zenoh-python's `Callback(indirect=True)`. zenoh-python +/// installs that handler by default when you hand `declare_subscriber` a plain +/// callable. The delivery thread enqueues and returns. User code runs on this +/// type's thread. +/// +/// A sample takes this path for either of two independent reasons. +/// +/// 1. **This same thread published it** ([`local_publish_active`]) — the +/// session-local case. Inline delivery would let a callback that publishes +/// into its own topic graph recurse instead of iterate. The queue makes that +/// feedback loop *iterative*. hiroz therefore needs no re-entrancy depth +/// cap: nothing can reach a callback from inside `put`. +/// 2. **The subscriber is a zenoh-ext `AdvancedSubscriber`.** It invokes the +/// sample callback while it holds the `std::sync::Mutex` that guards its +/// reordering state. It *has* to. `handle_sample` interleaves +/// `callback.call(sample)` with mutation of `last_delivered` and +/// `pending_samples`. `deliver_and_flush` calls the callback, records the +/// delivered sequence number, then drains newly-contiguous pending samples +/// and calls the callback again. zenoh-ext cannot drop the guard before the +/// call the way `Session::resolve_put` does. The lock protects exactly the +/// state the delivery loop walks. +/// +/// A sample that matches neither reason arrived over a transport, on a zenoh RX +/// worker, for a plain subscriber. hiroz delivers it inline. It never touches +/// this queue. That is deliberate. The RX thread is not an application thread +/// and holds no hiroz lock, so nothing can re-enter. The inter-process path must +/// not pay for a hazard it does not have. +/// +/// # Ordering +/// +/// There is one producer path, one FIFO queue and one drain thread. The user +/// therefore observes exactly the order zenoh chose to deliver in. +/// +/// On the advanced path the shim enqueues from inside `handle_sample`, under +/// zenoh-ext's state mutex. Enqueue order therefore includes the several +/// back-to-back deliveries that one `deliver_and_flush` performs when it drains +/// pending samples. This changes only the thread the callback runs on. It does +/// not change the reordering and recovery guarantees that `AdvancedSubscriber` +/// exists to provide. +/// +/// One ordering property does *not* hold: order between the two paths. A plain +/// subscriber can receive both local and remote publications on one topic. It +/// now runs the local ones on this thread and the remote ones on an RX thread. +/// Their relative order is no longer guaranteed, and the two can overlap. +/// +/// This weakens no guarantee that hiroz was actually offering. Neither ROS 2 nor +/// zenoh guarantees ordering across distinct publishers. Several RX workers can +/// already invoke a plain zenoh subscriber concurrently. The change is real, so +/// this section states it rather than leaving a reader to find it later. +/// +/// # Backpressure +/// +/// The queue **never blocks its producer**. That is not a tuning choice. A +/// bounded queue that blocked would re-create the original deadlock in a new +/// form on both paths. +/// +/// On the advanced path the blocked thread sits inside `sub_callback` and holds +/// zenoh-ext's state mutex. On the local path it sits inside the user's own +/// `publish()`. In a closed feedback loop, the drain thread it waits on is the +/// very thread that must publish for the queue to drain. +/// +/// zenoh's own `FifoChannel` is bounded *and* blocking. It documents exactly +/// this cost: "a slow subscriber could block the underlying Zenoh thread" +/// (`fifo.rs`). hiroz does not adopt that failure mode. +/// +/// That leaves a choice between unbounded (lossless, grows without limit) and +/// bounded drop-oldest (lossy, constant memory). **Both paths take the same +/// bound from the same history depth. They differ only in which samples reach +/// it:** +/// +/// * **Plain path — bounded, drop-oldest, capacity from the subscriber's +/// history QoS** ([`dispatch_capacity`]). A plain subscriber is `Volatile` +/// with `KEEP_LAST(depth)`. It already promises only the last `depth` +/// undelivered samples. The queue-mode path enforces exactly that with +/// [`BoundedQueue`], from the same history depth. The two expressions differ +/// only for a zero depth, which [`dispatch_capacity`] documents. +/// +/// A callback subscriber that retained *every* undelivered sample would +/// honour a QoS stricter than its declared one. It would also let a tight +/// local publish loop with a slow callback grow the process until it dies. +/// The declared QoS permits it to discard the samples it would retain, so +/// that trade has no upside. Drop-oldest also preserves the relative order of +/// the samples that survive. +/// +/// * **Advanced path — the same bound, from the same history depth.** This matches +/// `rmw_zenoh_cpp`. Its `SubscriptionData::add_new_message` drops the oldest +/// sample once `message_queue_.size() >= adapted_qos_profile.depth`. It does +/// so for every arriving sample, with **no `TransientLocal` exemption**: the +/// check reads the history policy only. Upstream sizes its advanced-subscriber +/// cache the same way (`adv_sub_opts.history->max_samples = qos_.depth`). +/// +/// `KeepAll` maps to [`DISPATCH_UNBOUNDED`] because that profile asks for +/// losslessness. A declared `KeepLast(depth)` does not, so this path honours +/// the depth. [`Self::always_shim`] enqueues remote samples too, so an +/// unbounded queue here would also replace zenoh's transport backpressure with +/// unbounded in-process growth. +/// +/// Both implementations drop **silently**, as the ROS event API sees it. +/// Upstream raises `MESSAGE_LOST` from *sequence-number gaps* between +/// arriving messages. A depth-drop cannot produce such a gap. hiroz raises no +/// `MESSAGE_LOST` event either (#292). The escalating `warn!` below is the +/// only signal here. It is more visible than upstream's debug log. +/// +/// Two consequences follow. This section states them rather than leaving a +/// reader to find them later. +/// +/// **The bound applies to different samples on each path.** On the plain path +/// only *locally published* samples pass through this queue. zenoh delivers a +/// sample that arrived over a transport inline on an RX worker, and +/// backpressures it at the transport instead. On the advanced path +/// [`Self::always_shim`] enqueues everything, remote samples included. +/// +/// A slow callback therefore loses local samples and stalls remote ones on the +/// plain path. It loses either kind on the advanced path. That asymmetry follows +/// from delivering the two on different threads, which is what makes re-entrancy +/// impossible without a cost on the inter-process path. +/// +/// **On the advanced path a `KeepAll` subscriber has no backpressure at all.** +/// The queue is genuinely unbounded there, and remote samples enter it. A +/// publisher that outpaces the callback grows `pending` without limit. The +/// escalating backlog warning is the only signal. This honours the declared QoS +/// and is not a defect. Even so, `KeepAll` on a slow callback commits unbounded +/// memory. Choose it deliberately. +pub struct CallbackDispatcher { + queue: Arc, + thread: Option>, +} + +impl CallbackDispatcher { + /// Spawns the drain thread. + /// + /// Two callers share `handler`. The drain thread always calls it. The + /// *plain* path's shim also calls it inline, for samples that did not + /// originate on this thread. Use [`Self::always_shim`] or + /// [`Self::local_only_shim`] to obtain the callback to hand to zenoh. + /// + /// `capacity` is the number of undelivered samples the queue retains before + /// it drops the oldest. **All four construction sites pass + /// [`dispatch_capacity`]:** + /// + /// * the plain arm of the typed builder (this module), + /// * the advanced arm of the typed builder, + /// * the plain arm of the FFI raw subscriber (`node.rs`), + /// * the advanced arm of the FFI raw subscriber. + /// + /// A callback subscriber therefore retains what its history QoS declares on + /// every path. See the "Backpressure" section. + /// + /// Keep all four sites in step by hand. The PR gate compiles the FFI arms + /// but does not lint or test them, so a wrong constant there fails no check + /// (#291). + pub(crate) fn spawn(topic: &str, handler: Arc, capacity: usize) -> Result + where + F: Fn(Sample) + Send + Sync + 'static, + { + let queue = Arc::new(DispatchQueue { + state: Mutex::new(DispatchState { + pending: std::collections::VecDeque::new(), + closed: false, + warn_at: DISPATCH_BACKLOG_WARN_AT, + dropped: 0, + warn_dropped_at: 1, + }), + ready: std::sync::Condvar::new(), + finished: Mutex::new(false), + done: std::sync::Condvar::new(), + topic: topic.to_string(), + capacity, + }); + + let drain_queue = queue.clone(); + let drain_topic = topic.to_string(); + let thread = std::thread::Builder::new() + .name("hiroz-sub-drain".to_string()) + .spawn(move || { + while let Some(sample) = drain_queue.dequeue() { + // A panicking user callback must not kill the drain thread. + // If it did, the subscriber would stop delivering silently. + // + // This guard works only where panics unwind. This + // workspace's `[profile.opt]` sets `panic = "abort"`, so + // there the panic aborts the process before `catch_unwind` + // can return `Err`. Neither the recovery nor the log line + // runs on that profile. + // + // The guard is therefore effective for dev, test and + // `release` builds. CI runs all three. The guard is inert + // for `opt`. A build that opts into aborting on panic has + // opted out of surviving one. + if std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| (*handler)(sample))) + .is_err() + { + tracing::error!( + topic = %drain_topic, + "subscriber callback panicked; dropping the sample and continuing" + ); + } + } + // Last action of the thread. `close` waits for this. It is set + // even when the loop exits because the queue closed mid-flight, + // which is the case `close` exists to observe. + { + let mut finished = drain_queue + .finished + .lock() + .unwrap_or_else(|e| e.into_inner()); + *finished = true; + } + drain_queue.done.notify_all(); + }) + .map_err(|e| { + zenoh::Error::from(format!("failed to spawn subscriber delivery thread: {e}")) + })?; + + Ok(Self { + queue, + thread: Some(thread), + }) + } + + /// A shim that enqueues **every** sample. + /// + /// The advanced path uses this shim. zenoh-ext holds its state mutex across + /// the callback whatever the sample's origin. + pub(crate) fn always_shim(&self) -> impl Fn(Sample) + Send + Sync + 'static { + let queue = self.queue.clone(); + move |sample: Sample| queue.enqueue(sample) + } + + /// A shim that enqueues only the samples the delivering thread produced + /// itself. It calls `handler` inline for every other sample. The plain path + /// uses this shim. + /// + /// The inline branch is the inter-process hot path. A zenoh RX worker + /// delivers a sample that arrived over a transport. That worker is never + /// inside a hiroz publish, so [`local_publish_active`] returns false. This + /// one thread-local read is the only cost the path pays. + pub(crate) fn local_only_shim( + &self, + handler: Arc, + ) -> impl Fn(Sample) + Send + Sync + 'static + where + F: Fn(Sample) + Send + Sync + 'static, + { + let queue = self.queue.clone(); + let topic = self.queue.topic.clone(); + move |sample: Sample| { + if local_publish_active() { + queue.enqueue(sample); + } else { + // The inline branch runs on a zenoh RX worker, so it needs the + // same guard the drain loop has. Without one, a panicking + // callback unwinds out of hiroz and into zenoh's receive path. + // + // This is the DEFAULT profile, which is what makes it matter: + // `qos_needs_advanced` is true only for `TransientLocal`, and + // the ROS 2 default is `Volatile`. Guarding only the drain + // thread therefore left every remote sample on the common path + // unguarded -- the half that carries inter-process traffic. + // + // The `panic = "abort"` caveat on the drain loop's guard + // applies here identically: a build that opts into aborting on + // a panic has opted out of surviving one. + if std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| handler(sample))) + .is_err() + { + tracing::error!( + topic = %topic, + "subscriber callback panicked on the inline path; dropping the \ + sample and continuing" + ); + } + } + } + } +} + +/// What [`CallbackDispatcher::close`], [`SubscriberHandle::close`] and +/// [`ZSub::close`] observed. Branch on it; do not ignore it. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum CloseOutcome { + /// The drain thread exited before the deadline. No user callback from this + /// subscriber can still be running. + Joined, + /// The deadline expired first. The drain thread was **detached**, not + /// killed: a user callback from this subscriber may still be running, and + /// will run to completion. No further callback starts. + TimedOut, +} + +impl Drop for CallbackDispatcher { + /// Stops delivery and returns at once. **It never joins the drain thread.** + /// + /// Dropping a subscriber guarantees that no *new* callback starts. It does + /// not guarantee that a callback already running has finished. Call + /// [`CallbackDispatcher::close`] (or [`ZSub::close`]) to wait for that. + /// + /// An earlier version joined here. That made `drop(sub)` block for the + /// whole of an in-flight callback, with no timeout and no log, and it made + /// every Python drop site depend on remembering `py.allow_threads` — the + /// callback body is `Python::with_gil`, so joining it while holding the GIL + /// waits for a thread that is waiting for the GIL. See + /// #296 (tag G2) for the decision. + /// + /// Detaching is not a new state. This type already detached when the caller + /// dropped the subscriber from inside its own callback, because joining a + /// thread from itself deadlocks. That special case is gone: it is now the + /// general rule, so it needs no separate branch. + /// + /// Detaching cannot be memory-unsafe. The handler is `Arc` with + /// `F: Send + Sync + 'static` and the queue is an `Arc`, so the detached + /// thread borrows nothing from the dropping scope. + fn drop(&mut self) { + self.queue.lock().closed = true; + self.queue.ready.notify_all(); + // Dropping the `JoinHandle` detaches the thread. It does not stop the + // thread; it stops *waiting* for it. The thread observes `closed` on + // its next `dequeue` and exits. + drop(self.thread.take()); + } +} + +impl CallbackDispatcher { + /// Stops delivery and waits, for at most `deadline`, until the drain thread + /// has exited. + /// + /// This is the explicit barrier that [`Drop`] no longer provides. It + /// consumes the dispatcher, so it cannot be called twice. + /// + /// It **discards the undelivered backlog**, exactly as [`Drop`] does. This + /// is a teardown barrier, not a flush. Draining first would re-create the + /// unbounded teardown cost that [`DispatchQueue::dequeue`]'s + /// `closed`-before-`pending` order exists to remove. + /// + /// Returns [`CloseOutcome::Joined`] when the thread exited in time, and + /// [`CloseOutcome::TimedOut`] when it did not — in which case the thread is + /// detached and its in-flight callback runs to completion. + /// + /// Calling this from inside the subscriber's own callback returns + /// [`CloseOutcome::TimedOut`] immediately. That thread cannot finish while + /// it is waiting for itself, so there is nothing to wait for. + pub fn close(mut self, deadline: Duration) -> CloseOutcome { + self.queue.lock().closed = true; + self.queue.ready.notify_all(); + + let Some(thread) = self.thread.take() else { + return CloseOutcome::Joined; + }; + if thread.thread().id() == std::thread::current().id() { + return CloseOutcome::TimedOut; + } + + let start = std::time::Instant::now(); + let mut finished = self + .queue + .finished + .lock() + .unwrap_or_else(|e| e.into_inner()); + while !*finished { + let Some(remaining) = deadline.checked_sub(start.elapsed()) else { + break; + }; + if remaining.is_zero() { + break; + } + let (guard, timed_out) = self + .queue + .done + .wait_timeout(finished, remaining) + .unwrap_or_else(|e| e.into_inner()); + finished = guard; + if timed_out.timed_out() { + break; + } + } + let exited = *finished; + drop(finished); + + if !exited { + warn!( + topic = %self.queue.topic, + ?deadline, + "subscriber delivery thread did not exit before the close deadline; \ + detaching it. Its in-flight callback will run to completion." + ); + return CloseOutcome::TimedOut; + } + if thread.join().is_err() { + warn!( + topic = %self.queue.topic, + "subscriber delivery thread terminated abnormally" + ); + } + CloseOutcome::Joined + } +} + +/// Whether a QoS profile needs a zenoh-ext `AdvancedSubscriber`. +/// +/// Both callers are subscriber paths. `ZPubBuilder::build` calls `.advanced()` +/// unconditionally and does not consult this function. +/// +/// The advanced entities exist for history replay, sample-miss detection and +/// recovery, and entity detection. [`apply_transient_local_sub`] and +/// [`apply_transient_local_pub`] configure all of these for `TransientLocal` +/// durability only. +/// +/// For the ROS 2 default (`Volatile`) an unconfigured `AdvancedSubscriber` adds +/// no protocol behaviour. It *does* run the user callback while it holds a +/// non-reentrant `std::sync::Mutex`. In `advanced_subscriber.rs`, `sub_callback` +/// takes `zlock!(statesref)`. `handle_sample` then calls the callback under that +/// guard. zenoh delivers a session-local sample synchronously on the publishing +/// thread. A publish from inside such a callback therefore deadlocks that thread +/// against itself (#249). A `Volatile` subscriber pays the lock and gains +/// nothing, so hiroz declares a plain subscriber for it instead. +pub(crate) fn qos_needs_advanced(qos: &hiroz_protocol::qos::QosProfile) -> bool { + matches!(qos.durability, QosDurability::TransientLocal) +} + +/// The declared zenoh subscriber backing a [`ZSub`]. +/// +/// hiroz declares a plain subscriber unless the QoS profile actually configures +/// advanced features — see [`qos_needs_advanced`]. +pub enum SubscriberHandle { + /// A plain zenoh subscriber (the `Volatile` default). + /// + /// A sample that arrived over a transport runs inline on the zenoh RX + /// worker. hiroz hands a sample the delivering thread published itself to + /// the dispatcher — see [`CallbackDispatcher`]. `dispatcher` is `None` for + /// queue-mode subscribers. They run no user code on the delivery thread, so + /// they need no handoff. + Plain { + /// This field comes first in the declaration so that it drops first. + /// + /// Rust drops struct fields in declaration order, so this order is a + /// proof obligation rather than a style choice. See the same field on + /// [`SubscriberHandle::Advanced`] for why the undeclare must precede + /// the dispatcher's teardown. + subscriber: zenoh::pubsub::Subscriber<()>, + dispatcher: Option, + }, + /// A zenoh-ext advanced subscriber, used for `TransientLocal` durability. + /// + /// `dispatcher` is `Some` only when the handler runs user code. zenoh-ext + /// holds its state lock across the callback, so hiroz must move user code + /// off that thread. A queue-mode handler only enqueues into a + /// [`BoundedQueue`] and re-enters nothing. It can therefore run under that + /// lock safely and needs no thread of its own. + Advanced { + /// Boxed because it is several times larger than the plain variant. + /// + /// This field comes first in the declaration so that it drops first. + /// Undeclaring the + /// subscriber stops new samples from entering the queue. Only then does + /// the dispatcher discard its backlog and release its thread. + /// + /// Rust drops struct fields in declaration order, so this order is a + /// proof obligation rather than a style choice. The guarantee is weaker + /// than it looks. zenoh undeclares with `wait_callbacks: false`, so a + /// sample can still arrive afterwards. `enqueue` returns early once + /// `closed` is set, which makes such a sample harmless. + subscriber: Box>, + dispatcher: Option, + }, +} + +impl SubscriberHandle { + /// Undeclares the zenoh subscriber, then waits for at most `deadline` for + /// the dispatcher's drain thread to exit. See [`CallbackDispatcher::close`]. + /// + /// A queue-mode subscriber has no dispatcher and runs no user code on the + /// delivery thread, so it reports [`CloseOutcome::Joined`] at once. + pub fn close(self, deadline: Duration) -> CloseOutcome { + // Undeclare first, for the same reason the field order does it: no new + // sample should enter a queue that is about to be discarded. + let dispatcher = match self { + Self::Plain { + subscriber, + dispatcher, + } => { + drop(subscriber); + dispatcher + } + Self::Advanced { + subscriber, + dispatcher, + } => { + drop(subscriber); + dispatcher + } + }; + match dispatcher { + Some(dispatcher) => dispatcher.close(deadline), + None => CloseOutcome::Joined, + } + } +} + +impl std::fmt::Debug for SubscriberHandle { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Plain { .. } => f.write_str("SubscriberHandle::Plain"), + Self::Advanced { .. } => f.write_str("SubscriberHandle::Advanced"), + } + } +} + /// Query timeout for TransientLocal subscribers' initial history fetch. /// Matches rmw_zenoh_cpp's `query_timeout_ms = u64::max()` literally /// (`Duration::from_millis(u64::MAX)`, not `Duration::MAX`, to avoid any @@ -538,6 +1302,7 @@ where trace!("[PUB] Attached sn={}", sn); } + let _local = LocalPublishGuard::enter(); put_builder.wait() } @@ -572,7 +1337,20 @@ where if self.with_attachment { put_builder = put_builder.attachment(self.new_attachment()); } - put_builder.await + // The guard must cover the delivery. Delivery happens in `into_future`, + // not at the await. zenoh's `PublicationBuilder` future is + // `std::future::ready(self.wait())`, so the put completes before a + // future exists to poll. That put includes any inline local-subscriber + // dispatch. + // + // Scoping the guard here rather than across the `.await` also keeps + // this future `Send`. A thread-local guard held across an await point + // would not. + let fut = { + let _local = LocalPublishGuard::enter(); + std::future::IntoFuture::into_future(put_builder) + }; + fut.await } /// Publish pre-serialized data directly @@ -593,6 +1371,7 @@ where if self.with_attachment { put_builder = put_builder.attachment(self.new_attachment()); } + let _local = LocalPublishGuard::enter(); put_builder.wait() } @@ -609,6 +1388,7 @@ where if self.with_attachment { put_builder = put_builder.attachment(self.new_attachment()); } + let _local = LocalPublishGuard::enter(); put_builder.wait() } @@ -825,9 +1605,12 @@ where let events_mgr = Arc::new(Mutex::new(EventsManager::new(gid))); let loss = Arc::new(MessageLossTracker::new(events_mgr.clone())); - // Wrap handler with encoding validation if expected encoding is set + // Wrap the handler with encoding validation. This needs no re-entrancy + // accounting: nothing reaches a user callback from inside `put` — see + // `CallbackDispatcher`. let expected_encoding = self.expected_encoding.clone(); - let validated_handler = move |sample: Sample| { + let runs_user_code = handler.runs_user_code(); + let validated_handler = Arc::new(move |sample: Sample| { observe_loss(&loss, &sample); // Validate encoding if expected encoding is set if let Some(ref expected) = expected_encoding { @@ -847,23 +1630,96 @@ where } } handler.handle(sample) - }; - - // Build an AdvancedSubscriber and configure based on durability - let mut sub_builder = self - .session - .declare_subscriber(key_expr) - .callback(validated_handler) - .advanced(); - - // Apply locality restriction if specified - if let Some(locality) = self.locality { - sub_builder = sub_builder.allowed_origin(locality); - debug!("[SUB] Locality restriction: {:?}", locality); - } + }); - let sub_builder = apply_transient_local_sub(sub_builder, &self.entity.qos); - let inner = sub_builder.wait()?; + // Go through zenoh-ext only when the QoS profile actually configures + // advanced features. See `qos_needs_advanced`. + let inner = if qos_needs_advanced(&self.entity.qos) { + debug!("[SUB] Using AdvancedSubscriber (TransientLocal durability)"); + // `AdvancedSubscriber` holds its state lock across the callback and + // cannot avoid it. hiroz therefore enqueues *user* code and runs it + // on the dispatcher's thread. See `CallbackDispatcher`. + // + // Capacity comes from the history QoS, exactly as on the plain + // path. Do not pass `DISPATCH_UNBOUNDED` here. `always_shim` + // enqueues *remote* samples too, so an unbounded queue on every + // profile trades zenoh's transport backpressure for unbounded + // in-process growth. `dispatch_capacity` still maps `KeepAll` to + // `DISPATCH_UNBOUNDED`, which is where the user asked to be + // lossless. See `CallbackDispatcher`'s "Backpressure" section. + // + // A queue-mode handler is exempt. It only pushes into a + // `BoundedQueue` and re-enters nothing, so it runs safely under + // zenoh-ext's lock. A dispatcher would add a thread, a wake and a + // second queue in front of the bounded one, on every TransientLocal + // rmw subscription, for nothing. + let dispatcher = if runs_user_code { + Some(CallbackDispatcher::spawn( + &qualified_topic, + validated_handler.clone(), + dispatch_capacity(&self.entity.qos), + )?) + } else { + None + }; + // Boxed to a common type: `always_shim` returns an opaque `impl Fn`, + // so the two arms cannot share a `match` unerased. + let callback: Box = match dispatcher.as_ref() { + Some(d) => Box::new(d.always_shim()), + None => Box::new(move |sample: Sample| validated_handler(sample)), + }; + let mut sub_builder = self.session.declare_subscriber(key_expr).callback(callback); + if let Some(locality) = self.locality { + sub_builder = sub_builder.allowed_origin(locality); + debug!("[SUB] Locality restriction: {:?}", locality); + } + let sub_builder = apply_transient_local_sub(sub_builder.advanced(), &self.entity.qos); + SubscriberHandle::Advanced { + subscriber: Box::new(sub_builder.wait()?), + dispatcher, + } + } else if runs_user_code { + // A plain subscriber holds no lock across the callback. zenoh still + // delivers a same-thread publication *inline*, so a callback that + // publishes into its own topic graph would recurse (#249). Hand + // those samples to the dispatcher. Deliver everything else inline, + // which keeps the inter-process path at one thread-local read. The + // dispatcher bounds its queue at the history depth and drops the + // oldest sample. That is the same `KEEP_LAST(depth)` the queue-mode + // path enforces with `BoundedQueue`. + let dispatcher = CallbackDispatcher::spawn( + &qualified_topic, + validated_handler.clone(), + dispatch_capacity(&self.entity.qos), + )?; + let mut sub_builder = self + .session + .declare_subscriber(key_expr) + .callback(dispatcher.local_only_shim(validated_handler)); + if let Some(locality) = self.locality { + sub_builder = sub_builder.allowed_origin(locality); + debug!("[SUB] Locality restriction: {:?}", locality); + } + SubscriberHandle::Plain { + subscriber: sub_builder.wait()?, + dispatcher: Some(dispatcher), + } + } else { + // Queue mode: the delivery thread only enqueues, so there is no user + // code to move off it and no dispatcher to pay for. + let mut sub_builder = self + .session + .declare_subscriber(key_expr) + .callback(move |sample: Sample| validated_handler(sample)); + if let Some(locality) = self.locality { + sub_builder = sub_builder.allowed_origin(locality); + debug!("[SUB] Locality restriction: {:?}", locality); + } + SubscriberHandle::Plain { + subscriber: sub_builder.wait()?, + dispatcher: None, + } + }; let lv_ke = self .keyexpr_format @@ -907,6 +1763,23 @@ where /// pattern), so Python/Go callers do not need to assign the return value. /// Rust callers must store the `ZSub` in their node or context. /// + /// # Dropping is not a barrier + /// + /// **Dropping a subscriber guarantees that no *new* callback starts. It does + /// not guarantee that a callback already running has finished.** The drop + /// returns at once: it closes the queue, discards the undelivered backlog + /// and detaches the drain thread. + /// + /// A callback that was already running therefore runs to completion, and may + /// still be running after `drop(sub)` returns. Do not assume that dropping + /// the subscriber releases a resource your callback captured. + /// + /// Call [`ZSub::close`] when you need the barrier. It waits, with a deadline + /// you choose, and tells you whether the thread exited or was detached. + /// + /// Dropping from inside the callback itself is safe, and needs no special + /// case: the drop never waits for anything. + /// /// # Arguments /// /// * `callback` - A function that will be called with each deserialized message @@ -992,7 +1865,7 @@ where pub struct ZSub { pub entity: EndpointEntity, pub queue: Option>>, - _inner: AdvancedSubscriber<()>, + _inner: SubscriberHandle, _lv_token: LivelinessToken, events_mgr: Arc>, graph: Arc, @@ -1004,6 +1877,35 @@ pub struct ZSub { _phantom_data: PhantomData<(T, Q, S)>, } +impl ZSub { + /// Tears the subscriber down and waits, for at most `deadline`, until no + /// callback of this subscriber is still running. + /// + /// Use this when you need a barrier. Plain `drop` gives you none: it + /// guarantees that no *new* callback starts, and returns without waiting + /// for one that is already running. + /// + /// This **discards the undelivered backlog**, exactly as `drop` does. It is + /// a teardown barrier, not a flush. + /// + /// ```text + /// match sub.close(Duration::from_secs(5)) { + /// CloseOutcome::Joined => // nothing of mine is still running + /// CloseOutcome::TimedOut => // a callback is still running; it was detached + /// } + /// ``` + pub fn close(self, deadline: Duration) -> CloseOutcome { + let ZSub { + _inner, _lv_token, .. + } = self; + let outcome = _inner.close(deadline); + // Same order the field declarations impose on `drop`: the subscriber + // goes away, then the liveliness token leaves the ROS graph. + drop(_lv_token); + outcome + } +} + impl std::fmt::Debug for ZSub { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("ZSub") @@ -1263,6 +2165,250 @@ impl ZSub, + cv: std::sync::Condvar, + } + + impl Latch { + fn new() -> Arc { + Arc::new(Self { + set: Mutex::new(false), + cv: std::sync::Condvar::new(), + }) + } + + fn set(&self) { + *self.set.lock().unwrap() = true; + self.cv.notify_all(); + } + + fn wait_for(&self, deadline: Duration) -> bool { + let start = std::time::Instant::now(); + let mut set = self.set.lock().unwrap(); + while !*set { + let Some(remaining) = deadline.checked_sub(start.elapsed()) else { + break; + }; + if remaining.is_zero() { + break; + } + let (guard, res) = self.cv.wait_timeout(set, remaining).unwrap(); + set = guard; + if res.timed_out() { + break; + } + } + *set + } + } + + /// Long enough that a genuinely blocking teardown cannot pass by luck, + /// short enough that a red run reports quickly. + const D1_SHORT: Duration = Duration::from_secs(5); + /// The upper bound on a parked callback. Only a red run ever waits this + /// long, and only to clean up after the assertion has already been made. + const D1_LONG: Duration = Duration::from_secs(60); + + fn d1_sample() -> Sample { + zenoh::sample::SampleBuilder::put( + zenoh::key_expr::KeyExpr::try_from("hiroz/test/d1").unwrap(), + vec![0u8], + ) + .into() + } + + /// Spawns a dispatcher whose callback parks on `release` until the test + /// lets it go, and signals `entered` as it starts. Returns the dispatcher + /// with one sample already enqueued and its callback confirmed running. + fn d1_parked_dispatcher( + entered: &Arc, + release: &Arc, + calls: &Arc, + ) -> CallbackDispatcher { + let (e, r, c) = (entered.clone(), release.clone(), calls.clone()); + let handler = Arc::new(move |_sample: Sample| { + c.fetch_add(1, Ordering::SeqCst); + e.set(); + assert!( + r.wait_for(D1_LONG), + "the parked callback was never released; the test leaked a thread" + ); + }); + let dispatcher = CallbackDispatcher::spawn("/d1", handler, 8).unwrap(); + let shim = dispatcher.always_shim(); + shim(d1_sample()); + assert!( + entered.wait_for(D1_SHORT), + "the callback never started, so the teardown assertion would be vacuous" + ); + dispatcher + } + + #[test] + fn drop_returns_while_a_callback_is_still_running() { + let entered = Latch::new(); + let release = Latch::new(); + let dropped = Latch::new(); + let calls = Arc::new(AtomicUsize::new(0)); + + let dispatcher = d1_parked_dispatcher(&entered, &release, &calls); + + // Drop from another thread, so this thread can time the drop while the + // callback is still parked. + let dropped_signal = dropped.clone(); + let dropper = std::thread::spawn(move || { + drop(dispatcher); + dropped_signal.set(); + }); + + let returned = dropped.wait_for(D1_SHORT); + // Release before asserting: a failing run must still terminate. + release.set(); + dropper.join().unwrap(); + + assert!( + returned, + "drop(dispatcher) did not return within {D1_SHORT:?} while a callback was \ + still running. Dropping a subscriber must stop new callbacks, not wait \ + for the in-flight one -- see #296 (tag G2)." + ); + } + + #[test] + fn close_waits_for_an_in_flight_callback() { + let entered = Latch::new(); + let release = Latch::new(); + let calls = Arc::new(AtomicUsize::new(0)); + + let dispatcher = d1_parked_dispatcher(&entered, &release, &calls); + + let releaser = { + let release = release.clone(); + std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(200)); + release.set(); + }) + }; + + let outcome = dispatcher.close(D1_SHORT); + releaser.join().unwrap(); + + assert_eq!( + outcome, + CloseOutcome::Joined, + "close() must report that it joined when the callback finishes in time" + ); + } + + #[test] + fn close_reports_a_timeout_rather_than_blocking_forever() { + let entered = Latch::new(); + let release = Latch::new(); + let calls = Arc::new(AtomicUsize::new(0)); + + let dispatcher = d1_parked_dispatcher(&entered, &release, &calls); + + let start = std::time::Instant::now(); + let outcome = dispatcher.close(Duration::from_millis(300)); + let waited = start.elapsed(); + release.set(); + + assert_eq!( + outcome, + CloseOutcome::TimedOut, + "a callback that outlives the deadline must be reported, not waited on" + ); + assert!( + waited < D1_SHORT, + "close() waited {waited:?}, far past its 300ms deadline" + ); + } + + #[test] + fn close_discards_the_backlog_rather_than_flushing_it() { + let entered = Latch::new(); + let release = Latch::new(); + let calls = Arc::new(AtomicUsize::new(0)); + + let dispatcher = d1_parked_dispatcher(&entered, &release, &calls); + + // Three more samples queue up behind the parked callback. + let shim = dispatcher.always_shim(); + for _ in 0..3 { + shim(d1_sample()); + } + + let releaser = { + let release = release.clone(); + std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(200)); + release.set(); + }) + }; + let outcome = dispatcher.close(D1_SHORT); + releaser.join().unwrap(); + + assert_eq!(outcome, CloseOutcome::Joined); + assert_eq!( + calls.load(Ordering::SeqCst), + 1, + "close() is a teardown barrier, not a flush: the three queued samples \ + must be discarded, exactly as drop discards them" + ); + } + + #[test] + fn volatile_does_not_need_an_advanced_subscriber() { + let qos = hiroz_protocol::qos::QosProfile { + durability: QosDurability::Volatile, + ..Default::default() + }; + assert!( + !qos_needs_advanced(&qos), + "Volatile is the ROS 2 default. An AdvancedSubscriber adds no protocol \ + behaviour for it and runs the user callback under a non-reentrant mutex" + ); + } + + #[test] + fn transient_local_needs_an_advanced_subscriber() { + let qos = hiroz_protocol::qos::QosProfile { + durability: QosDurability::TransientLocal, + ..Default::default() + }; + assert!( + qos_needs_advanced(&qos), + "TransientLocal needs history replay and miss recovery" + ); + } + // ----------------------------------------------------------------------- // Topic name qualification (leading '/' is added when missing) // ----------------------------------------------------------------------- diff --git a/crates/hiroz/src/queue.rs b/crates/hiroz/src/queue.rs index a8cf6423a..296ab54f9 100644 --- a/crates/hiroz/src/queue.rs +++ b/crates/hiroz/src/queue.rs @@ -114,3 +114,50 @@ impl BoundedQueue { } } } + +#[cfg(test)] +mod tests { + use super::*; + + /// A zero capacity retains one sample, not zero. + /// + /// `pubsub::dispatch_capacity` floors a zero history depth at 1. The + /// queue-mode path passes the 0 straight through. The two sizing + /// expressions therefore differ. This test pins the reason that divergence + /// is harmless: `push` evicts *before* it inserts, so capacity 0 retains + /// one sample just as capacity 1 does. A reordering of `push` to + /// insert-then-evict makes a zero-depth queue discard every sample, and + /// this test then fails. + #[test] + fn zero_capacity_retains_one_sample() { + let q = BoundedQueue::new(0); + + assert!( + q.push(1), + "capacity 0 reports a drop even on the first push" + ); + assert_eq!(q.len(), 1, "capacity 0 must retain one sample, not zero"); + + q.push(2); + assert_eq!(q.len(), 1); + assert_eq!(q.try_recv(), Some(2), "the newest sample is the one kept"); + assert!(q.is_empty()); + } + + /// The capacity-1 case that the doc claims equivalence against. Retention + /// is the same. Capacity 1 reports no spurious drop on the first push. + #[test] + fn capacity_one_retains_one_sample_without_reporting_a_drop() { + let q = BoundedQueue::new(1); + + assert!( + !q.push(1), + "an empty capacity-1 queue must not report a drop" + ); + assert_eq!(q.len(), 1); + + assert!(q.push(2), "the second push evicts the first"); + assert_eq!(q.len(), 1); + assert_eq!(q.try_recv(), Some(2)); + } +}