diff --git a/.github/workflows/release-debs.yaml b/.github/workflows/release-debs.yaml index 3adaf15..3964e4f 100644 --- a/.github/workflows/release-debs.yaml +++ b/.github/workflows/release-debs.yaml @@ -27,6 +27,10 @@ jobs: os_codename: noble container: rostooling/setup-ros-docker:ubuntu-noble-ros-kilted-ros-base-latest runner: ubuntu-24.04 + - ros_distro: lyrical + os_codename: resolute + container: rostooling/setup-ros-docker:ubuntu-resolute-ros-lyrical-ros-base-latest + runner: ubuntu-24.04 runs-on: ${{ matrix.runner }} container: @@ -49,9 +53,9 @@ jobs: - name: Install IXWebSocket from source run: | cd /tmp - wget -q https://github.com/machinezone/IXWebSocket/archive/refs/tags/v11.4.6.tar.gz - tar xzf v11.4.6.tar.gz - cd IXWebSocket-11.4.6 + wget -q https://github.com/machinezone/IXWebSocket/archive/refs/tags/v12.0.1.tar.gz + tar xzf v12.0.1.tar.gz + cd IXWebSocket-12.0.1 cmake -B build -DCMAKE_BUILD_TYPE=Release -DUSE_TLS=OFF -DUSE_ZLIB=ON cmake --build build -j$(nproc) cmake --install build diff --git a/.github/workflows/ros_lyrical.yaml b/.github/workflows/ros_lyrical.yaml new file mode 100644 index 0000000..ee65ff8 --- /dev/null +++ b/.github/workflows/ros_lyrical.yaml @@ -0,0 +1,20 @@ +name: "ROS: Lyrical" + +on: + push: + branches: [development, main] + pull_request: + branches: [development, main] + workflow_dispatch: + +concurrency: + group: ${{ github.workflow }}-${{ github.ref }} + cancel-in-progress: true + +jobs: + build-and-test: + uses: ./.github/workflows/_ros.yaml + with: + ros_distro: lyrical + container: rostooling/setup-ros-docker:ubuntu-resolute-ros-lyrical-ros-base-latest + runner: ubuntu-24.04 diff --git a/CHANGELOG.rst b/CHANGELOG.rst index a555064..ca84bcf 100644 --- a/CHANGELOG.rst +++ b/CHANGELOG.rst @@ -2,6 +2,23 @@ Changelog for package pj_ros_bridge ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ +Forthcoming +----------- +* WebSocket permessage-deflate declined server-side (redundant with our own + ZSTD compression; was costing CPU for no bandwidth benefit). +* ROS2: ingest executor polled and drained (``ingest_poll_interval_ms``, + default 5 ms) instead of blocking ``spin()``, amortizing the per-wait-cycle + wait-set rebuild; publish/request/timeout timers moved to their own + executor thread so a long publish cycle never delays ingest. +* ROS2: ``min_qos_depth`` default raised ``1`` -> ``10`` so the ingest poll + interval can't overflow a shallow reader between polls. +* Dependencies: IXWebSocket 11.4.6 -> 12.0.1 (FetchContent and .deb builds), + Fast DDS 3.4.0 -> 3.4.3, CLI11 2.6.0 -> 2.6.2. Fixed the FastDDS/RTI + link failure when IXWebSocket is fetched with TLS. +* CI and Debian release for ROS 2 Lyrical. +* Removed zero-initializing copies in the ingest and serializer hot paths + (``resize()`` + ``memcpy`` -> direct-construct/``insert``). + 0.9.0 (2026-07-11) ------------------ * Size-class frames: isolate heavy messages (``>= heavy_frame_threshold_bytes``) diff --git a/CLAUDE.md b/CLAUDE.md index e6e6d26..5bf8581 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -51,7 +51,7 @@ make -j$(nproc) - **CLI11** — FetchContent (RTI backend only) ### FastDDS Backend (Conan-managed) -- **eProsima Fast DDS 3.4.0** — via Conan (`fast-dds/3.4.0`) +- **eProsima Fast DDS 3.4.3** — via Conan (`fast-dds/3.4.3`) - **eProsima Fast CDR 2.x** — transitive dependency via Conan - **CLI11** — FetchContent (for CLI parsing, shared with RTI) @@ -105,7 +105,7 @@ make -j$(nproc) ### Event Loop BridgeServer does NOT own timers. The entry point (`main.cpp`) drives the event loop: -- **ROS2**: `rclcpp` wall timers call `process_requests()`, `publish_aggregated_messages()`, `check_session_timeouts()` +- **ROS2**: `rclcpp` wall timers (`process_requests()`, `publish_aggregated_messages()`, `check_session_timeouts()`) run on their own dedicated executor thread, so a long publish/zstd cycle never delays ingest. The ingest executor itself is polled (`ingest_poll_interval_ms`, default 5 ms) via `spin_some()` instead of blocking in `spin()`; subscription callbacks drain their DDS reader on each poll. - **RTI**: `std::chrono` loop with `std::this_thread::sleep_for()` - **FastDDS**: `std::chrono` loop with `std::this_thread::sleep_for()` (same pattern as RTI) @@ -191,7 +191,7 @@ Then, for each message in the (compressed) payload: - CLI11 (FetchContent, for CLI parsing) ### FastDDS Backend -- eProsima Fast DDS 3.4.0 (Conan: `fast-dds/3.4.0`) +- eProsima Fast DDS 3.4.3 (Conan: `fast-dds/3.4.3`) - eProsima Fast CDR 2.x (transitive Conan dependency) - CLI11 (FetchContent, for CLI parsing) @@ -234,8 +234,9 @@ publish_rate: 50.0 # Hz session_timeout: 10.0 # seconds strip_large_messages: false # Opt-in: strip Image/PointCloud2/etc data fields topic_whitelist: [".*"] # Full-match regex patterns restricting visible/subscribable topics -min_qos_depth: 1 # Minimum KEEP_LAST subscription depth after aggregating publisher depths +min_qos_depth: 10 # Minimum KEEP_LAST subscription depth after aggregating publisher depths max_qos_depth: 100 # Maximum KEEP_LAST subscription depth after aggregating publisher depths +ingest_poll_interval_ms: 5.0 # Ingest executor poll interval (drains subscriptions each poll); 0 = blocking spin topic_poll_interval: 1.0 # Seconds between topics_changed notification polls; 0 disables polling client_backlog_size: 100 # Max frames queued per slow client before dropping the oldest (must be > 0) heavy_frame_threshold_bytes: 262144 # Isolate messages >= this size into their own size-class frame; 0 disables diff --git a/CMakeLists.txt b/CMakeLists.txt index f941358..5e75a07 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -69,8 +69,8 @@ if(NOT ixwebsocket_FOUND) message(STATUS "IXWebSocket not found — fetching via FetchContent") FetchContent_Declare( ixwebsocket - URL https://github.com/machinezone/IXWebSocket/archive/refs/tags/v11.4.6.tar.gz - URL_HASH SHA256=c024334f8e45980836c67008979a884d6dcc5ef067dd2eb1fa7241f4c17ddc32 + URL https://github.com/machinezone/IXWebSocket/archive/refs/tags/v12.0.1.tar.gz + URL_HASH SHA256=d23bdc91dbfe2b9ae13c322d539392d7a6b8b506560f41c90e227fa0f86a2405 ) set(USE_TLS ${PJ_BRIDGE_TLS} CACHE BOOL "" FORCE) if(PJ_BRIDGE_TLS) @@ -79,6 +79,11 @@ if(NOT ixwebsocket_FOUND) set(USE_ZLIB ON CACHE BOOL "" FORCE) set(IXWEBSOCKET_INSTALL OFF CACHE BOOL "" FORCE) FetchContent_MakeAvailable(ixwebsocket) + if(PJ_BRIDGE_TLS) + # The static library does not export its OpenSSL dependency to consumers. + find_package(OpenSSL REQUIRED) + target_link_libraries(ixwebsocket INTERFACE OpenSSL::SSL OpenSSL::Crypto) + endif() else() message(STATUS "Using system IXWebSocket") endif() @@ -149,11 +154,12 @@ if(ament_cmake_FOUND) pj_bridge_app ) - ament_target_dependencies(pj_bridge_ros2_lib PUBLIC - ament_index_cpp - rclcpp - sensor_msgs - nav_msgs + # target_link_libraries, not ament_target_dependencies: removed in newer distros (Lyrical) + target_link_libraries(pj_bridge_ros2_lib PUBLIC + ament_index_cpp::ament_index_cpp + rclcpp::rclcpp + ${sensor_msgs_TARGETS} + ${nav_msgs_TARGETS} ) # ROS2 executable @@ -165,10 +171,6 @@ if(ament_cmake_FOUND) pj_bridge_ros2_lib ) - ament_target_dependencies(pj_bridge_ros2 - rclcpp - ) - # Install ROS2 executable install(TARGETS pj_bridge_ros2 DESTINATION lib/${PROJECT_NAME} @@ -276,12 +278,6 @@ if(BUILD_TESTING AND ament_cmake_FOUND) data_path ) - ament_target_dependencies(${PROJECT_NAME}_tests - rclcpp - sensor_msgs - nav_msgs - ) - # Sanitizer environment settings for CTest if(ENABLE_TSAN) set_tests_properties(${PROJECT_NAME}_tests PROPERTIES diff --git a/README.md b/README.md index 91aa6fc..c066cdb 100644 --- a/README.md +++ b/README.md @@ -4,7 +4,7 @@ A high-performance bridge server that forwards middleware topic content over WebSocket to PlotJuggler clients. Three backends share a common core: -- **ROS2** (`pj_bridge_ros2`) — ROS2 Humble / Jazzy / Kilted via `rclcpp` +- **ROS2** (`pj_bridge_ros2`) — ROS2 Humble / Jazzy / Kilted / Lyrical via `rclcpp` - **FastDDS** (`pj_bridge_fastdds`) — eProsima Fast DDS 3.4 (standalone, no ROS2 required) - **RTI** (`pj_bridge_rti`) — RTI Connext DDS (build disabled, code preserved) @@ -34,10 +34,10 @@ independently. ## CI Status -| | Humble | Jazzy | Kilted | -|--|--------|-------|--------| -| **Pixi** | [![Pixi: Humble](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_humble.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_humble.yaml) | [![Pixi: Jazzy](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_jazzy.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_jazzy.yaml) | [![Pixi: Kilted](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_kilted.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_kilted.yaml) | -| **colcon** | [![ROS: Humble](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_humble.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_humble.yaml) | [![ROS: Jazzy](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_jazzy.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_jazzy.yaml) | [![ROS: Kilted](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_kilted.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_kilted.yaml) | +| | Humble | Jazzy | Kilted | Lyrical | +|--|--------|-------|--------|---------| +| **Pixi** | [![Pixi: Humble](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_humble.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_humble.yaml) | [![Pixi: Jazzy](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_jazzy.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_jazzy.yaml) | [![Pixi: Kilted](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_kilted.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/pixi_kilted.yaml) | — | +| **colcon** | [![ROS: Humble](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_humble.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_humble.yaml) | [![ROS: Jazzy](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_jazzy.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_jazzy.yaml) | [![ROS: Kilted](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_kilted.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_kilted.yaml) | [![ROS: Lyrical](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_lyrical.yaml/badge.svg?branch=main)](https://github.com/PlotJuggler/plotjuggler_bridge/actions/workflows/ros_lyrical.yaml) | ## Configuration Parameters @@ -50,8 +50,9 @@ independently. | `session_timeout` | double | 10.0 | Client timeout duration in seconds | | `strip_large_messages` | bool | false | Opt-in: strip large arrays from Image, PointCloud2, LaserScan, OccupancyGrid messages | | `topic_whitelist` | string array | `[".*"]` | Full-match regex patterns (ECMAScript) restricting visible/subscribable topics | -| `min_qos_depth` | int | 1 | ROS2 only: minimum KEEP_LAST subscription depth after aggregating publisher depths | +| `min_qos_depth` | int | 10 | ROS2 only: minimum KEEP_LAST subscription depth after aggregating publisher depths | | `max_qos_depth` | int | 100 | ROS2 only: maximum KEEP_LAST subscription depth after aggregating publisher depths | +| `ingest_poll_interval_ms` | double | 5.0 | ROS2 only: interval between ingest executor polls; subscription callbacks drain their DDS reader each poll. `0` uses a blocking spin instead | | `topic_poll_interval` | double | 1.0 | Seconds between `topics_changed` notification polls; `0` disables polling | | `client_backlog_size` | int | 100 | Max binary frames queued per slow client before the oldest is dropped (must be `> 0`) | | `heavy_frame_threshold_bytes` | int | 262144 | Isolate messages ≥ this size (bytes) into their own size-class frame so they don't starve small topics; `0` disables (must be `>= 0`) | diff --git a/app/src/message_buffer.cpp b/app/src/message_buffer.cpp index d0f4a27..855506e 100644 --- a/app/src/message_buffer.cpp +++ b/app/src/message_buffer.cpp @@ -55,6 +55,12 @@ void MessageBuffer::set_latched(const std::string& topic_name, bool latched) { std::lock_guard lock(mutex_); if (latched) { latched_topics_.insert(topic_name); + // Ingest may run on another thread: the retained sample can already be + // buffered by the time the subscriber's bookkeeping gets here. + auto it = topic_buffers_.find(topic_name); + if (latched_last_.count(topic_name) == 0 && it != topic_buffers_.end() && !it->second.empty()) { + latched_last_[topic_name] = it->second.back(); + } } else { latched_topics_.erase(topic_name); latched_last_.erase(topic_name); diff --git a/app/src/message_serializer.cpp b/app/src/message_serializer.cpp index 7a4d645..72c70a2 100644 --- a/app/src/message_serializer.cpp +++ b/app/src/message_serializer.cpp @@ -53,9 +53,9 @@ void AggregatedMessageSerializer::serialize_message( write_le(serialized_data_, msg_size); // Message data (CDR bytes) using memcpy - size_t old_size = serialized_data_.size(); - serialized_data_.resize(old_size + msg_size); - std::memcpy(serialized_data_.data() + old_size, data, msg_size); + // insert, not resize+memcpy: resize would zero the bytes first + const auto *bytes = reinterpret_cast(data); + serialized_data_.insert(serialized_data_.end(), bytes, bytes + msg_size); } void AggregatedMessageSerializer::clear() { diff --git a/app/src/middleware/websocket_middleware.cpp b/app/src/middleware/websocket_middleware.cpp index 901ff55..96c18f7 100644 --- a/app/src/middleware/websocket_middleware.cpp +++ b/app/src/middleware/websocket_middleware.cpp @@ -74,6 +74,8 @@ tl::expected WebSocketMiddleware::initialize(uint16_t port) { } server_ = std::make_shared(port, "0.0.0.0"); + // Binary frames are already zstd-compressed; deflating them again is wasted work. + server_->disablePerMessageDeflate(); #ifdef IXWEBSOCKET_USE_TLS if (tls_.has_value()) { diff --git a/conanfile.txt b/conanfile.txt index 92df9f1..5482dd0 100644 --- a/conanfile.txt +++ b/conanfile.txt @@ -1,6 +1,6 @@ [requires] -fast-dds/3.4.0 -cli11/2.6.0 +fast-dds/3.4.3 +cli11/2.6.2 [generators] CMakeDeps diff --git a/docs/API.md b/docs/API.md index 2dd54eb..52f9d8e 100644 --- a/docs/API.md +++ b/docs/API.md @@ -208,7 +208,7 @@ offers (so a burst from every publisher still fits in the subscription queue), then clamping the total to a configurable range — the same heuristic `foxglove_bridge`'s `determineQoS()` uses: -- **ROS2**: int parameters `min_qos_depth` (default `1`) and `max_qos_depth` +- **ROS2**: int parameters `min_qos_depth` (default `10`) and `max_qos_depth` (default `100`). A publisher that reports depth `0` (KEEP_ALL, or an RMW such as @@ -223,6 +223,26 @@ match what the discovered publishers offer (a RELIABLE subscription still switches to BEST_EFFORT if any publisher is BEST_EFFORT, and to TRANSIENT_LOCAL only if every publisher offers it). +`min_qos_depth` defaults to `10` (rather than `1`) so the ingest poll +interval below can't overflow a shallow reader between polls. + +## Ingest Poll Interval (ROS2 only) + +The ROS2 backend's ingest executor (the one running subscription callbacks) +polls via `spin_some()` on a timer instead of blocking in `spin()`: `rclcpp` +rebuilds its whole wait set — O(subscriptions) — on every wait cycle, which +with `spin()` means one rebuild per received message. Polling amortizes one +rebuild over every message that arrived in the poll interval; subscription +callbacks drain their DDS reader on each poll, so nothing is lost as long as +the reader depth (`min_qos_depth`) covers a burst. + +- **ROS2**: double parameter `ingest_poll_interval_ms` (default `5.0`, + milliseconds). `0` reverts to a blocking `spin()`. Must be `>= 0`; the + server refuses to start otherwise. + +Publish/request/session-timeout timers run on their own executor thread so a +long publish (zstd-compression) cycle never delays ingest. + ## Subscribe Subscribe to one or more topics. **Breaking change:** Subscribe now uses an additive model - it only adds topics without removing existing subscriptions. Use the `unsubscribe` command to remove topics. diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 6a86ba9..2935494 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -152,7 +152,7 @@ Unlike the RTI backend's 4-class two-level design (discovery + subscription mana - **FastDdsTopicSource**: Manages `DomainParticipant`s, discovers topics via `DomainParticipantListener::on_data_writer_discovery()`, resolves `DynamicType` from `TypeObjectRegistry`, generates IDL schema via `idl_serialize()`. Also provides `get_dynamic_type()` / `get_participant()` / `get_domain_id()` for use by the subscription manager. - **FastDdsSubscriptionManager**: Creates `DataReader`s with `DynamicPubSubType`, ref-counted subscriptions. Extracts CDR bytes by deserializing into `DynamicData` and re-serializing via `DynamicPubSubType::serialize()`. -FastDDS dependencies are managed via Conan (`fast-dds/3.4.0`). The backend is built standalone (not through colcon/ament). +FastDDS dependencies are managed via Conan (`fast-dds/3.4.3`). The backend is built standalone (not through colcon/ament). ## Design Decisions diff --git a/fastdds/src/fastdds_subscription_manager.cpp b/fastdds/src/fastdds_subscription_manager.cpp index 3750b0d..7b223ea 100644 --- a/fastdds/src/fastdds_subscription_manager.cpp +++ b/fastdds/src/fastdds_subscription_manager.cpp @@ -80,8 +80,8 @@ class FastDdsSubscriptionManager::InternalReaderListener : public DataReaderList continue; } - auto cdr_data = std::make_shared>(payload_.length); - std::memcpy(cdr_data->data(), payload_.data, payload_.length); + const auto* bytes = reinterpret_cast(payload_.data); + auto cdr_data = std::make_shared>(bytes, bytes + payload_.length); uint64_t timestamp_ns = static_cast(info.source_timestamp.seconds()) * 1'000'000'000ULL + info.source_timestamp.nanosec(); diff --git a/ros2/include/pj_bridge_ros2/generic_subscription_manager.hpp b/ros2/include/pj_bridge_ros2/generic_subscription_manager.hpp index cff1572..9a8abc6 100644 --- a/ros2/include/pj_bridge_ros2/generic_subscription_manager.hpp +++ b/ros2/include/pj_bridge_ros2/generic_subscription_manager.hpp @@ -54,6 +54,7 @@ class GenericSubscriptionManager { bool unsubscribe(const std::string& topic_name); bool is_subscribed(const std::string& topic_name) const; size_t get_reference_count(const std::string& topic_name) const; + size_t subscription_count() const; void unsubscribe_all(); /// True when the subscription for `topic_name` was created with diff --git a/ros2/include/pj_bridge_ros2/ros2_subscription_manager.hpp b/ros2/include/pj_bridge_ros2/ros2_subscription_manager.hpp index 2ab8b2b..6e1c8ef 100644 --- a/ros2/include/pj_bridge_ros2/ros2_subscription_manager.hpp +++ b/ros2/include/pj_bridge_ros2/ros2_subscription_manager.hpp @@ -50,6 +50,7 @@ class Ros2SubscriptionManager : public SubscriptionManagerInterface { void unsubscribe_all() override; bool is_transient_local(const std::string& topic_name) const override; bool is_subscribed(const std::string& topic_name) const override; + size_t subscription_count() const; private: GenericSubscriptionManager inner_manager_; diff --git a/ros2/src/generic_subscription_manager.cpp b/ros2/src/generic_subscription_manager.cpp index 22cfca3..dfecb98 100644 --- a/ros2/src/generic_subscription_manager.cpp +++ b/ros2/src/generic_subscription_manager.cpp @@ -19,12 +19,78 @@ #include "pj_bridge_ros2/generic_subscription_manager.hpp" +#include + #include +#include #include "pj_bridge/time_utils.hpp" namespace pj_bridge { +namespace { +constexpr size_t kMaxDrainPerCallback = 1000; + +// The executor hands over one message per subscription per wait cycle, and +// each cycle rebuilds the whole wait set (O(subscriptions)). This subscription +// drains whatever else its reader holds, so a burst costs one cycle instead of +// one per message. It overrides handle_serialized_message() rather than using +// a callback because only this hook sees the MessageInfo of the first message: +// with batched delivery, "now" is no longer a usable receive time. +class DrainingSubscription : public rclcpp::GenericSubscription { + public: + using Handler = std::function&, uint64_t)>; + + DrainingSubscription( + rclcpp::node_interfaces::NodeBaseInterface* node_base, std::shared_ptr ts_lib, + const std::string& topic_name, const std::string& topic_type, const rclcpp::QoS& qos, Handler handler) + : rclcpp::GenericSubscription( + node_base, std::move(ts_lib), topic_name, topic_type, qos, unused_callback(), + rclcpp::SubscriptionOptions()), + handler_(std::move(handler)) {} + + void handle_serialized_message( + const std::shared_ptr& message, const rclcpp::MessageInfo& info) override { + handler_(message, receive_time_ns(info)); + // Bounded so a publisher faster than we can drain cannot starve other topics. + for (size_t i = 0; i < kMaxDrainPerCallback; ++i) { + auto extra = std::make_shared(); + rclcpp::MessageInfo extra_info; + try { + if (!take_serialized(*extra, extra_info)) { + break; + } + } catch (const std::exception& e) { + RCLCPP_WARN(rclcpp::get_logger("pj_bridge"), "take failed on '%s': %s", get_topic_name(), e.what()); + break; + } + handler_(extra, receive_time_ns(extra_info)); + } + } + + private: + static uint64_t receive_time_ns(const rclcpp::MessageInfo& info) { + // Not every RMW fills received_timestamp. + const auto received = info.get_rmw_message_info().received_timestamp; + return received > 0 ? static_cast(received) : get_current_time_ns(); + } + +#if RCLCPP_VERSION_MAJOR >= 28 // Jazzy+: the constructor takes an AnySubscriptionCallback + static rclcpp::AnySubscriptionCallback> unused_callback() { + rclcpp::AnySubscriptionCallback> callback; + callback.set([](std::shared_ptr){}); + return callback; + } +#else + static std::function)> unused_callback() { + return [](std::shared_ptr) {}; + } +#endif + + Handler handler_; +}; +} // namespace + GenericSubscriptionManager::GenericSubscriptionManager( rclcpp::Node::SharedPtr node, size_t min_qos_depth, size_t max_qos_depth) : node_(node), min_qos_depth_(min_qos_depth), max_qos_depth_(max_qos_depth) {} @@ -103,15 +169,16 @@ bool GenericSubscriptionManager::subscribe( } try { - auto sub_callback = [topic_name, callback](std::shared_ptr msg) { - uint64_t receive_time = get_current_time_ns(); - callback(topic_name, msg, receive_time); - }; - rclcpp::QoS qos = adapt_qos(topic_name); bool transient_local = (qos.durability() == rclcpp::DurabilityPolicy::TransientLocal); - auto subscription = node_->create_generic_subscription(topic_name, topic_type, qos, sub_callback); + auto subscription = std::make_shared( + node_->get_node_base_interface().get(), rclcpp::get_typesupport_library(topic_type, "rosidl_typesupport_cpp"), + topic_name, topic_type, qos, + [topic_name, callback](const std::shared_ptr& msg, uint64_t receive_time_ns) { + callback(topic_name, msg, receive_time_ns); + }); + node_->get_node_topics_interface()->add_subscription(subscription, nullptr); subscriptions_[topic_name] = SubscriptionInfo{subscription, 1, transient_local}; @@ -150,6 +217,11 @@ bool GenericSubscriptionManager::is_subscribed(const std::string& topic_name) co return subscriptions_.find(topic_name) != subscriptions_.end(); } +size_t GenericSubscriptionManager::subscription_count() const { + std::lock_guard lock(mutex_); + return subscriptions_.size(); +} + size_t GenericSubscriptionManager::get_reference_count(const std::string& topic_name) const { std::lock_guard lock(mutex_); diff --git a/ros2/src/main.cpp b/ros2/src/main.cpp index 2f5b71b..1de482e 100644 --- a/ros2/src/main.cpp +++ b/ros2/src/main.cpp @@ -24,6 +24,7 @@ #include #include #include +#include #include #include "pj_bridge/bridge_server.hpp" @@ -45,7 +46,8 @@ int main(int argc, char** argv) { node->declare_parameter("session_timeout", 10.0); node->declare_parameter("strip_large_messages", false); node->declare_parameter>("topic_whitelist", {".*"}); - node->declare_parameter("min_qos_depth", 1); + node->declare_parameter("min_qos_depth", 10); + node->declare_parameter("ingest_poll_interval_ms", 5.0); node->declare_parameter("max_qos_depth", 100); node->declare_parameter("topic_poll_interval", 1.0); node->declare_parameter("client_backlog_size", 100); @@ -60,6 +62,7 @@ int main(int argc, char** argv) { bool strip_large_messages = node->get_parameter("strip_large_messages").as_bool(); std::vector topic_whitelist = node->get_parameter("topic_whitelist").as_string_array(); int64_t min_qos_depth = node->get_parameter("min_qos_depth").as_int(); + double ingest_poll_interval_ms = node->get_parameter("ingest_poll_interval_ms").as_double(); int64_t max_qos_depth = node->get_parameter("max_qos_depth").as_int(); double topic_poll_interval = node->get_parameter("topic_poll_interval").as_double(); int64_t client_backlog_size = node->get_parameter("client_backlog_size").as_int(); @@ -71,9 +74,10 @@ int main(int argc, char** argv) { RCLCPP_INFO( node->get_logger(), "Configuration: port=%d, publish_rate=%.1f Hz, session_timeout=%.1f s, strip_large_messages=%s, " - "min_qos_depth=%ld, max_qos_depth=%ld, topic_poll_interval=%.1f s, client_backlog_size=%ld, tls=%s", + "min_qos_depth=%ld, max_qos_depth=%ld, ingest_poll_interval_ms=%.1f, topic_poll_interval=%.1f s, " + "client_backlog_size=%ld, tls=%s", port, publish_rate, session_timeout, strip_large_messages ? "true" : "false", min_qos_depth, max_qos_depth, - topic_poll_interval, client_backlog_size, tls_enabled ? "true" : "false"); + ingest_poll_interval_ms, topic_poll_interval, client_backlog_size, tls_enabled ? "true" : "false"); if (tls_enabled && (certfile.empty() || keyfile.empty())) { RCLCPP_ERROR(node->get_logger(), "tls=true requires both 'certfile' and 'keyfile' parameters to be set"); @@ -89,6 +93,14 @@ int main(int argc, char** argv) { return 1; } + if (ingest_poll_interval_ms < 0.0) { + RCLCPP_ERROR( + node->get_logger(), "Invalid ingest_poll_interval_ms: %.1f (must be >= 0; 0 uses blocking spin)", + ingest_poll_interval_ms); + rclcpp::shutdown(); + return 1; + } + if (client_backlog_size <= 0) { RCLCPP_ERROR(node->get_logger(), "Invalid client_backlog_size: %ld (must be > 0)", client_backlog_size); rclcpp::shutdown(); @@ -159,14 +171,19 @@ int main(int argc, char** argv) { // ROS2 timers drive the event loop using namespace std::chrono_literals; - auto request_timer = node->create_wall_timer(10ms, [&server]() { server.process_requests(); }); + // Not added to the node's default executor (second argument false). + auto timer_group = node->create_callback_group(rclcpp::CallbackGroupType::MutuallyExclusive, false); + + auto request_timer = node->create_wall_timer( + 10ms, [&server]() { server.process_requests(); }, timer_group); auto publish_period = std::chrono::duration(1.0 / publish_rate); auto publish_timer = node->create_wall_timer( std::chrono::duration_cast(publish_period), - [&server]() { server.publish_aggregated_messages(); }); + [&server]() { server.publish_aggregated_messages(); }, timer_group); - auto timeout_timer = node->create_wall_timer(1s, [&server]() { server.check_session_timeouts(); }); + auto timeout_timer = node->create_wall_timer( + 1s, [&server]() { server.check_session_timeouts(); }, timer_group); // topic_poll_interval == 0 disables the pushed topic-advertisement poll. rclcpp::TimerBase::SharedPtr topic_poll_timer; @@ -177,16 +194,57 @@ int main(int argc, char** argv) { server.check_topic_changes(); topic_poll_timer = node->create_wall_timer( std::chrono::duration_cast(std::chrono::duration(topic_poll_interval)), - [&server]() { server.check_topic_changes(); }); + [&server]() { server.check_topic_changes(); }, timer_group); } - // Spin until shutdown. spin() blocks waiting for work and returns when - // rclcpp::shutdown() runs (e.g. on SIGINT) — unlike spin_some() in a - // loop, which returns immediately when idle and busy-spins a full core. + // Timers (publish/zstd, requests, timeouts) run on their own thread so a + // long publish cycle never delays ingest. + rclcpp::executors::SingleThreadedExecutor timer_executor; + timer_executor.add_callback_group(timer_group, node->get_node_base_interface()); + // Joins on every exit path: an exception unwinding past a joinable + // std::thread would call std::terminate(). + struct TimerThread { + rclcpp::executors::SingleThreadedExecutor& executor; + std::thread thread; + ~TimerThread() { + executor.cancel(); + if (thread.joinable()) { + thread.join(); + } + } + } timer_thread{timer_executor, std::thread([&timer_executor, &node]() { + try { + timer_executor.spin(); + } catch (const std::exception& e) { + RCLCPP_FATAL(node->get_logger(), "Timer thread failed: %s", e.what()); + rclcpp::shutdown(); + } + })}; + + // Every executor wait cycle rebuilds the whole wait set, O(subscriptions); + // with blocking spin() that is one rebuild per received message. Polling + // amortizes one rebuild over every message that arrived in the interval; + // subscription callbacks drain their reader, so nothing is lost as long + // as the reader depth covers a burst. rclcpp::executors::SingleThreadedExecutor executor; executor.add_node(node); + if (ingest_poll_interval_ms > 0.0) { + const auto poll_interval = std::chrono::duration(ingest_poll_interval_ms); + while (rclcpp::ok()) { + if (sub_manager->subscription_count() == 0) { + // Nothing to amortize: block instead of polling so idle costs nothing. + executor.spin_once(std::chrono::milliseconds(100)); + continue; + } + executor.spin_some(); + std::this_thread::sleep_for(poll_interval); + } + } else { + executor.spin(); + } - executor.spin(); + timer_executor.cancel(); + timer_thread.thread.join(); // Graceful shutdown RCLCPP_INFO(node->get_logger(), "Shutting down bridge server..."); diff --git a/ros2/src/ros2_subscription_manager.cpp b/ros2/src/ros2_subscription_manager.cpp index fd6de59..3f806b9 100644 --- a/ros2/src/ros2_subscription_manager.cpp +++ b/ros2/src/ros2_subscription_manager.cpp @@ -54,8 +54,8 @@ bool Ros2SubscriptionManager::subscribe(const std::string& topic_name, const std } const auto& rcl_msg = msg_to_use->get_rcl_serialized_message(); - auto data = std::make_shared>(rcl_msg.buffer_length); - std::memcpy(data->data(), rcl_msg.buffer, rcl_msg.buffer_length); + const auto* bytes = reinterpret_cast(rcl_msg.buffer); + auto data = std::make_shared>(bytes, bytes + rcl_msg.buffer_length); MessageCallback cb; { @@ -74,6 +74,10 @@ bool Ros2SubscriptionManager::unsubscribe(const std::string& topic_name) { return inner_manager_.unsubscribe(topic_name); } +size_t Ros2SubscriptionManager::subscription_count() const { + return inner_manager_.subscription_count(); +} + void Ros2SubscriptionManager::unsubscribe_all() { inner_manager_.unsubscribe_all(); } diff --git a/tests/unit/test_generic_subscription_manager.cpp b/tests/unit/test_generic_subscription_manager.cpp index 70d6c15..1b02653 100644 --- a/tests/unit/test_generic_subscription_manager.cpp +++ b/tests/unit/test_generic_subscription_manager.cpp @@ -24,6 +24,7 @@ #include #include #include +#include #include #include "pj_bridge_ros2/generic_subscription_manager.hpp" @@ -435,3 +436,84 @@ TEST_F(GenericSubscriptionManagerTest, IsTransientLocalFalseForVolatilePublisher TEST_F(GenericSubscriptionManagerTest, IsTransientLocalFalseForUnsubscribedTopic) { EXPECT_FALSE(manager_->is_transient_local("/never_subscribed_topic")); } + +// --------------------------------------------------------------------------- +// Draining a burst within one executor wait cycle +// +// An rclcpp executor hands over only ONE message per subscription per wait +// cycle. Without draining, a single spin_once() delivers 1 message even if +// 20 are queued in the reader's history. subscribe()'s callback must drain +// the rest via take_serialized() after delivering the one the executor gave +// it, so a single wait cycle empties a burst instead of taking one cycle per +// message. +// --------------------------------------------------------------------------- +TEST_F(GenericSubscriptionManagerTest, DrainsAllPendingMessagesInOneExecutorPass) { + // Publisher first, with enough depth (>=20) for the subscription's QoS + // (derived from the publisher via adapt_qos()) to hold a 20-message burst, + // and RELIABLE so nothing is dropped before the executor ever spins. + auto publisher = node_->create_publisher("/drain_burst_topic", rclcpp::QoS(50).reliable()); + + // Wait for the publisher to be discoverable before subscribing so + // adapt_qos() sees it and derives a >=20 depth. + { + auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(3); + while (node_->count_publishers("/drain_burst_topic") == 0 && std::chrono::steady_clock::now() < deadline) { + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + } + } + ASSERT_GT(node_->count_publishers("/drain_burst_topic"), 0u); + + std::atomic received_count{0}; + auto callback = [&received_count](const std::string&, const std::shared_ptr&, uint64_t) { + received_count++; + }; + + ASSERT_TRUE(manager_->subscribe("/drain_burst_topic", "std_msgs/msg/String", callback)); + + // Wait until the publisher sees the subscription matched. Nothing has been + // published yet, so spinning here would be harmless — but no spin is + // needed for local graph updates either. + { + auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(3); + while (publisher->get_subscription_count() == 0 && std::chrono::steady_clock::now() < deadline) { + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + } + } + ASSERT_GT(publisher->get_subscription_count(), 0u); + + // Publish a burst back-to-back with no executor running, so all 20 land in + // the reader's history before anything is ever delivered. + for (int i = 0; i < 20; ++i) { + std_msgs::msg::String msg; + msg.data = "msg"; + publisher->publish(msg); + } + std::this_thread::sleep_for(std::chrono::milliseconds(300)); + + // An executor's wait cycle hands over exactly one message per + // subscription; spin_once() (unlike spin_some(), which may loop through + // several wait cycles in one call) executes exactly one ready callback per + // call. A freshly-added node's very first wait cycle in this environment + // can spuriously find nothing ready yet even though data is already + // sitting in the reader (an executor/RMW wait-set warm-up quirk, not + // something the fix under test controls), so we retry spin_once() calls + // until something is delivered — but the discriminating check is that the + // *first* non-zero observation must already be all 20: without the drain + // loop in subscribe()'s callback, a single ready wait cycle only ever + // delivers 1 message, so received_count would stop at 1 instead of + // jumping straight to 20. + rclcpp::executors::SingleThreadedExecutor executor; + executor.add_node(node_); + int spin_calls = 0; + constexpr int kMaxSpinCalls = 30; + for (; spin_calls < kMaxSpinCalls && received_count.load() == 0; ++spin_calls) { + executor.spin_once(std::chrono::milliseconds(200)); + } + + ASSERT_GT(received_count.load(), 0) << "no message delivered after " << spin_calls << " spin_once() call(s)"; + EXPECT_EQ(received_count.load(), 20) << "the wait cycle that delivered the first message only delivered " + << received_count.load() << " of 20 — drain loop did not run"; + + executor.remove_node(node_); + manager_->unsubscribe("/drain_burst_topic"); +} diff --git a/tests/unit/test_message_buffer.cpp b/tests/unit/test_message_buffer.cpp index 6e7a8e1..370aa37 100644 --- a/tests/unit/test_message_buffer.cpp +++ b/tests/unit/test_message_buffer.cpp @@ -442,6 +442,19 @@ TEST_F(MessageBufferTest, GetLatchedReturnsNewestMessage) { EXPECT_EQ(latched->timestamp_ns, 2u); } +// Ingest runs on a different thread than subscribe bookkeeping, so the sole +// retained sample of a latched topic can land before set_latched(true). +TEST_F(MessageBufferTest, SetLatchedSeedsFromAlreadyBufferedSample) { + buffer_.add_message("/t", 7, create_test_data({7, 7})); + + buffer_.set_latched("/t", true); + + auto latched = buffer_.get_latched("/t"); + ASSERT_TRUE(latched.has_value()); + EXPECT_EQ(latched->timestamp_ns, 7u); + EXPECT_EQ(extract_data(latched->data), (std::vector{7, 7})); +} + TEST_F(MessageBufferTest, SetLatchedFalseClearsRetainedEntry) { buffer_.set_latched("/t", true); buffer_.add_message("/t", 1, create_test_data({1, 2, 3})); diff --git a/tests/unit/test_websocket_middleware.cpp b/tests/unit/test_websocket_middleware.cpp index 0dd9029..262931d 100644 --- a/tests/unit/test_websocket_middleware.cpp +++ b/tests/unit/test_websocket_middleware.cpp @@ -17,6 +17,7 @@ * along with pj_bridge. If not, see . */ +#include #include #include @@ -200,6 +201,37 @@ TEST_F(WebSocketMiddlewareTest, ClientConnectAndSendMessage) { std::this_thread::sleep_for(std::chrono::milliseconds(100)); } +// Binary frames are already zstd-compressed; negotiating permessage-deflate +// would re-compress them with zlib on the publish thread. +TEST_F(WebSocketMiddlewareTest, ServerDeclinesPerMessageDeflate) { + auto result = middleware_->initialize(18099); + ASSERT_TRUE(result.has_value()); + + std::mutex headers_mutex; + std::string extensions_header; + ix::WebSocket client; + client.setUrl("ws://127.0.0.1:18099"); + client.setPerMessageDeflateOptions(ix::WebSocketPerMessageDeflateOptions(true)); + client.setOnMessageCallback([&](const ix::WebSocketMessagePtr& msg) { + if (msg->type == ix::WebSocketMessageType::Open) { + std::lock_guard lock(headers_mutex); + auto it = msg->openInfo.headers.find("Sec-WebSocket-Extensions"); + if (it != msg->openInfo.headers.end()) { + extensions_header = it->second; + } + } + }); + client.start(); + ASSERT_TRUE(wait_for_client_open(client)) << "Client failed to connect"; + + { + std::lock_guard lock(headers_mutex); + EXPECT_EQ(extensions_header.find("permessage-deflate"), std::string::npos) + << "server negotiated: " << extensions_header; + } + client.stop(); +} + TEST_F(WebSocketMiddlewareTest, SendReplyToConnectedClient) { auto result = middleware_->initialize(18091); ASSERT_TRUE(result.has_value()); @@ -706,6 +738,12 @@ TEST(WebSocketMiddlewareTlsTest, TlsRoundTrip) { #endif // IXWEBSOCKET_USE_TLS int main(int argc, char** argv) { + // rcl_logging_spdlog registers a periodic flush thread on spdlog's global + // registry (which we share) and its code lives in that library. rcl unloads + // the library on shutdown, so after the suites' init/shutdown cycles the + // thread would resume in unmapped memory at exit. Pin the library. + dlopen("librcl_logging_spdlog.so", RTLD_NOW | RTLD_NODELETE); + testing::InitGoogleTest(&argc, argv); return RUN_ALL_TESTS(); }