From 9eade2db69652026a783fcfd49f68db4c67932f1 Mon Sep 17 00:00:00 2001 From: Davide Faconti Date: Sun, 20 Sep 2026 12:43:49 +0200 Subject: [PATCH 1/9] perf(middleware): decline WebSocket permessage-deflate Binary frames are already zstd-compressed. When a client offered the extension, IXWebSocket deflated every frame again with zlib on the publish thread: 7-11% of process CPU in a perf profile, 19.2% -> 12.9% of a core for one client on an 83-topic workload. Co-Authored-By: Claude Fable 5.1 --- app/src/middleware/websocket_middleware.cpp | 3 ++ tests/unit/test_websocket_middleware.cpp | 31 +++++++++++++++++++++ 2 files changed, 34 insertions(+) diff --git a/app/src/middleware/websocket_middleware.cpp b/app/src/middleware/websocket_middleware.cpp index 901ff55..627f602 100644 --- a/app/src/middleware/websocket_middleware.cpp +++ b/app/src/middleware/websocket_middleware.cpp @@ -74,6 +74,9 @@ 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 costs more + // CPU than the zstd pass itself and gains nothing. + server_->disablePerMessageDeflate(); #ifdef IXWEBSOCKET_USE_TLS if (tls_.has_value()) { diff --git a/tests/unit/test_websocket_middleware.cpp b/tests/unit/test_websocket_middleware.cpp index 0dd9029..7b03401 100644 --- a/tests/unit/test_websocket_middleware.cpp +++ b/tests/unit/test_websocket_middleware.cpp @@ -200,6 +200,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()); From eb6d911a210fdf2caecadbd01b740eadd76bd3ec Mon Sep 17 00:00:00 2001 From: Davide Faconti Date: Sun, 20 Sep 2026 13:20:46 +0200 Subject: [PATCH 2/9] perf(ros2): poll + drain ingest, run timers on their own thread Profiling showed >50% of process CPU in rcl_wait: rmw_fastrtps re-attaches every subscription to the wait set on each wait cycle, and blocking spin() does one cycle per received message. - Subscription callbacks now drain their reader, so one cycle harvests a whole burst (the executor hands over only one message per cycle). - The ingest executor polls every ingest_poll_interval_ms (default 5, 0 = old blocking spin). - min_qos_depth default 1 -> 10 so a poll interval cannot overflow a shallow reader. - Publish/request/timeout timers moved to a dedicated callback group and executor thread: zstd no longer blocks ingest. 83 topics, ~700 msg/s, 1 client: 7.5% -> 3.4% of a core, identical message counts. Co-Authored-By: Claude Fable 5.1 --- ros2/src/generic_subscription_manager.cpp | 27 +++++- ros2/src/main.cpp | 42 ++++++++-- .../test_generic_subscription_manager.cpp | 84 +++++++++++++++++++ 3 files changed, 141 insertions(+), 12 deletions(-) diff --git a/ros2/src/generic_subscription_manager.cpp b/ros2/src/generic_subscription_manager.cpp index 22cfca3..44236ba 100644 --- a/ros2/src/generic_subscription_manager.cpp +++ b/ros2/src/generic_subscription_manager.cpp @@ -25,6 +25,10 @@ namespace pj_bridge { +namespace { +constexpr size_t kMaxDrainPerCallback = 1000; +} // 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,9 +107,25 @@ 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); + // The executor hands over one message per subscription per wait cycle, and + // each cycle rebuilds the whole wait set (O(subscriptions)). Drain whatever + // else the reader holds so a burst costs one cycle instead of one per message. + auto self = std::make_shared>(); + auto sub_callback = [topic_name, callback, self](std::shared_ptr msg) { + callback(topic_name, msg, get_current_time_ns()); + auto sub = self->lock(); + if (!sub) { + return; + } + // 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 info; + if (!sub->take_serialized(*extra, info)) { + break; + } + callback(topic_name, extra, get_current_time_ns()); + } }; rclcpp::QoS qos = adapt_qos(topic_name); @@ -113,6 +133,7 @@ bool GenericSubscriptionManager::subscribe( auto subscription = node_->create_generic_subscription(topic_name, topic_type, qos, sub_callback); + *self = subscription; subscriptions_[topic_name] = SubscriptionInfo{subscription, 1, transient_local}; return true; diff --git a/ros2/src/main.cpp b/ros2/src/main.cpp index 2f5b71b..8ab897e 100644 --- a/ros2/src/main.cpp +++ b/ros2/src/main.cpp @@ -22,6 +22,7 @@ #include #include #include +#include #include #include #include @@ -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(); @@ -159,14 +162,17 @@ 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 +183,34 @@ 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()); + std::thread timer_thread([&timer_executor]() { timer_executor.spin(); }); + + // Every executor wait cycle rebuilds the whole wait set, O(subscriptions); + // with blocking spin() that is one rebuild per received message (>50% of + // CPU with ~80 subscriptions). 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()) { + executor.spin_some(); + std::this_thread::sleep_for(poll_interval); + } + } else { + executor.spin(); + } - executor.spin(); + timer_executor.cancel(); + timer_thread.join(); // Graceful shutdown RCLCPP_INFO(node->get_logger(), "Shutting down bridge server..."); diff --git a/tests/unit/test_generic_subscription_manager.cpp b/tests/unit/test_generic_subscription_manager.cpp index 70d6c15..4b2c653 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,86 @@ 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"); +} From c590e3058a20bef63d93dbf811f0ba918a28d330 Mon Sep 17 00:00:00 2001 From: Davide Faconti Date: Sun, 20 Sep 2026 13:44:15 +0200 Subject: [PATCH 3/9] perf: heavy_frame_zstd_level knob, non-zeroing copies, idle-free ingest poll - heavy_frame_zstd_level (ROS2 param / --heavy-frame-zstd-level): zstd is ~60% of CPU on point-cloud workloads. Default stays 1 (no behaviour change); measured on 4 lidars, 52 MiB/s in: 1 = 19% of a core, 31 MB/s out; -5 = 15%, 39 MB/s; -100 = 10%, 51 MB/s. Wire-compatible. - Ingest and serializer copies no longer zero the buffer first. - ROS2 ingest loop blocks instead of polling while there are no subscriptions, so idle CPU matches the old blocking spin (0.45%). - Validation, docs and changelog for the new options. Co-Authored-By: Claude Fable 5.1 --- CHANGELOG.rst | 16 ++++++++ CLAUDE.md | 12 ++++-- README.md | 5 ++- app/include/pj_bridge/bridge_server.hpp | 7 ++++ app/include/pj_bridge/message_serializer.hpp | 5 ++- app/include/pj_bridge/protocol_constants.hpp | 5 +++ app/src/bridge_server.cpp | 3 +- app/src/message_serializer.cpp | 11 +++-- docs/API.md | 38 +++++++++++++++++- fastdds/src/fastdds_subscription_manager.cpp | 4 +- fastdds/src/main.cpp | 10 ++++- .../generic_subscription_manager.hpp | 1 + .../ros2_subscription_manager.hpp | 1 + ros2/src/generic_subscription_manager.cpp | 5 +++ ros2/src/main.cpp | 40 ++++++++++++++++--- ros2/src/ros2_subscription_manager.cpp | 8 +++- rti/src/main.cpp | 10 ++++- .../test_generic_subscription_manager.cpp | 8 ++-- tests/unit/test_message_serializer.cpp | 36 +++++++++++++++++ 19 files changed, 194 insertions(+), 31 deletions(-) diff --git a/CHANGELOG.rst b/CHANGELOG.rst index a555064..1d150d2 100644 --- a/CHANGELOG.rst +++ b/CHANGELOG.rst @@ -2,6 +2,22 @@ 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. +* New ``heavy_frame_zstd_level`` knob (ROS2 param, RTI/FastDDS + ``--heavy-frame-zstd-level``, default ``1``): CPU-vs-bandwidth zstd level + for heavy (size-class) frames only; wire-compatible (frame flags stay 0). +* 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..662f970 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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) @@ -234,11 +234,13 @@ 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 +heavy_frame_zstd_level: 1 # zstd level for heavy frames; any zstd level, negative = faster/larger tls: false # Enable TLS (wss://); requires certfile and keyfile certfile: "" # TLS server certificate file keyfile: "" # TLS private key file @@ -248,14 +250,16 @@ keyfile: "" # TLS private key file ```bash pj_bridge_rti --domains 0 1 --port 9090 --publish-rate 50 --session-timeout 10 \ --topic-whitelist ".*" --topic-poll-interval 1.0 --client-backlog-size 100 \ - --heavy-frame-threshold-bytes 262144 --certfile cert.pem --keyfile key.pem + --heavy-frame-threshold-bytes 262144 --heavy-frame-zstd-level 1 \ + --certfile cert.pem --keyfile key.pem ``` ### FastDDS (via CLI flags): ```bash pj_bridge_fastdds --domains 0 1 --port 9090 --publish-rate 50 --session-timeout 10 \ --topic-whitelist ".*" --topic-poll-interval 1.0 --client-backlog-size 100 \ - --heavy-frame-threshold-bytes 262144 --certfile cert.pem --keyfile key.pem + --heavy-frame-threshold-bytes 262144 --heavy-frame-zstd-level 1 \ + --certfile cert.pem --keyfile key.pem ``` See `docs/API.md` for full semantics of each option (topic whitelist matching rules, diff --git a/README.md b/README.md index 91aa6fc..75c2563 100644 --- a/README.md +++ b/README.md @@ -50,11 +50,13 @@ 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`) | +| `heavy_frame_zstd_level` | int | 1 | zstd compression level for heavy (size-class) frames; any zstd level, negative = faster/larger | | `tls` | bool | false | Enable TLS (`wss://`); requires `certfile` and `keyfile` | | `certfile` | string | `""` | TLS server certificate file | | `keyfile` | string | `""` | TLS private key file | @@ -73,6 +75,7 @@ independently. | `--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 (range `1`-`1000000`) | | `--heavy-frame-threshold-bytes` | int | 262144 | Isolate messages ≥ this size (bytes) into their own size-class frame; `0` disables (range `0`-`1000000000`) | +| `--heavy-frame-zstd-level` | int | 1 | zstd compression level for heavy (size-class) frames; negative = faster/larger (range: zstd min-max level) | | `--certfile` | string | (none) | TLS server certificate file; enables `wss://`, requires `--keyfile` | | `--keyfile` | string | (none) | TLS private key file; enables `wss://`, requires `--certfile` | | `--qos-profile` | string | (none) | RTI only: QoS profile XML file path | diff --git a/app/include/pj_bridge/bridge_server.hpp b/app/include/pj_bridge/bridge_server.hpp index 7ce7f42..95c01e3 100644 --- a/app/include/pj_bridge/bridge_server.hpp +++ b/app/include/pj_bridge/bridge_server.hpp @@ -61,6 +61,12 @@ struct BridgeServerConfig { /// its own size-class ("heavy") frame instead of being aggregated with light /// topics. 0 disables splitting (single aggregated frame). Default: 256 KiB. size_t heavy_frame_threshold_bytes = kDefaultHeavyFrameThresholdBytes; + /// zstd level for heavy frames; a CPU-vs-bandwidth knob (any zstd level, + /// negative = faster). Compression is ~60% of CPU on point-cloud workloads. + /// Measured on 4 lidars, 52 MiB/s in: level 1 = 19% of a core, 31 MB/s out; + /// -5 = 15%, 39 MB/s; -100 = 10%, 51 MB/s. Light frames always use + /// kDefaultZstdLevel. + int heavy_frame_zstd_level = kDefaultHeavyFrameZstdLevel; }; class BridgeServer { @@ -226,6 +232,7 @@ class BridgeServer { // Per-message byte size at or above which a topic is isolated into its own // size-class ("heavy") frame; 0 disables splitting. See publish_aggregated_messages(). size_t heavy_frame_threshold_bytes_; + int heavy_frame_zstd_level_; // State std::atomic initialized_; diff --git a/app/include/pj_bridge/message_serializer.hpp b/app/include/pj_bridge/message_serializer.hpp index 9df4185..018ffe8 100644 --- a/app/include/pj_bridge/message_serializer.hpp +++ b/app/include/pj_bridge/message_serializer.hpp @@ -26,6 +26,8 @@ #include #include +#include "pj_bridge/protocol_constants.hpp" + namespace pj_bridge { /** @@ -96,7 +98,8 @@ class AggregatedMessageSerializer { * kFrameFlagHeavy to mark an isolated large/size-class frame. * @return Vector containing header + compressed payload */ - std::vector finalize(uint32_t flags = 0); + /// @param compression_level any zstd level; negative levels trade ratio for speed. + std::vector finalize(uint32_t flags = 0, int compression_level = kDefaultZstdLevel); /** * @brief Compress data using ZSTD (compression level 1) diff --git a/app/include/pj_bridge/protocol_constants.hpp b/app/include/pj_bridge/protocol_constants.hpp index af57f3c..fb12a80 100644 --- a/app/include/pj_bridge/protocol_constants.hpp +++ b/app/include/pj_bridge/protocol_constants.hpp @@ -48,6 +48,11 @@ static constexpr uint32_t kFrameFlagHeavy = 0x1; /// threshold of 0 disables splitting (single aggregated frame, legacy behavior). static constexpr size_t kDefaultHeavyFrameThresholdBytes = 256 * 1024; // 256 KiB +/// zstd level for aggregated (light) frames. +static constexpr int kDefaultZstdLevel = 1; +/// zstd level for heavy frames; see BridgeServerConfig::heavy_frame_zstd_level. +static constexpr int kDefaultHeavyFrameZstdLevel = 1; + /// Schema encoding identifier for ROS2 message definitions inline constexpr const char* kSchemaEncodingRos2Msg = "ros2msg"; diff --git a/app/src/bridge_server.cpp b/app/src/bridge_server.cpp index fe3551e..a167540 100644 --- a/app/src/bridge_server.cpp +++ b/app/src/bridge_server.cpp @@ -69,6 +69,7 @@ BridgeServer::BridgeServer( publish_rate_(config.publish_rate), whitelist_(std::move(config.whitelist)), heavy_frame_threshold_bytes_(config.heavy_frame_threshold_bytes), + heavy_frame_zstd_level_(config.heavy_frame_zstd_level), initialized_(false), total_messages_published_(0), total_bytes_published_(0), @@ -1099,7 +1100,7 @@ void BridgeServer::publish_aggregated_messages() { AggregatedMessageSerializer heavy_serializer; heavy_serializer.serialize_message(topic, msg.timestamp_ns, msg.data->data(), msg.data->size()); GroupFrame heavy_frame; - heavy_frame.compressed_data = heavy_serializer.finalize(); + heavy_frame.compressed_data = heavy_serializer.finalize(0, heavy_frame_zstd_level_); heavy_frame.msg_count = 1; heavy_frame.client_ids = client_ids; heavy_frame.is_heavy = true; diff --git a/app/src/message_serializer.cpp b/app/src/message_serializer.cpp index 7a4d645..a3c3453 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() { @@ -67,7 +67,7 @@ size_t AggregatedMessageSerializer::get_message_count() const { return message_count_; } -std::vector AggregatedMessageSerializer::finalize(uint32_t flags) { +std::vector AggregatedMessageSerializer::finalize(uint32_t flags, int compression_level) { // Build 16-byte header (uncompressed) std::vector header(kBinaryHeaderSize); @@ -101,8 +101,7 @@ std::vector AggregatedMessageSerializer::finalize(uint32_t flags) { // Compress payload after header using persistent context size_t compressed_size = ZSTD_compressCCtx( cctx_, result.data() + kBinaryHeaderSize, max_compressed, serialized_data_.data(), serialized_data_.size(), - 1 // compression level - ); + compression_level); if (ZSTD_isError(compressed_size)) { throw std::runtime_error(std::string("ZSTD compression failed: ") + ZSTD_getErrorName(compressed_size)); diff --git a/docs/API.md b/docs/API.md index 2dd54eb..718f4ba 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. @@ -568,6 +588,22 @@ frame is configurable: Keep the threshold below the 1 MiB socket watermark so a single heavy message does not fill the socket buffer on its own. +Heavy frames use a separately configurable zstd compression level — a +CPU-vs-bandwidth knob independent of the threshold above. Light (aggregated) +frames always compress at level `1`. Any zstd level is valid (negative +levels trade ratio for speed); the frame `flags` stay `0` either way, so any +zstd decoder accepts the output regardless of level: + +- **ROS2**: int parameter `heavy_frame_zstd_level`, default `1`. Must be in + `[ZSTD_minCLevel(), ZSTD_maxCLevel()]`; the server refuses to start + otherwise. +- **FastDDS / RTI**: CLI flag `--heavy-frame-zstd-level`, default `1`, same + valid range. + +Measured on 4 lidars, 52 MiB/s in, 1 client (i7-13700H, P-core pinned, +performance governor): level `1` = 19.0% of a core / 31.0 MB/s out; `-5` = +14.8% / 39.0; `-20` = 14.4% / 45.6; `-100` = 10.3% / 50.6. + ## TLS / wss:// The bridge can optionally serve the WebSocket endpoint over TLS (`wss://`) 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/fastdds/src/main.cpp b/fastdds/src/main.cpp index bbd6b20..19861b0 100644 --- a/fastdds/src/main.cpp +++ b/fastdds/src/main.cpp @@ -18,6 +18,7 @@ */ #include +#include #include #include @@ -43,6 +44,7 @@ int main(int argc, char* argv[]) { double topic_poll_interval = 1.0; int client_backlog_size = 100; int heavy_frame_threshold_bytes = 262144; + int heavy_frame_zstd_level = pj_bridge::kDefaultHeavyFrameZstdLevel; std::string certfile; std::string keyfile; @@ -67,6 +69,11 @@ int main(int argc, char* argv[]) { "0 disables (keep below the 1 MiB socket watermark)") ->default_val(262144) ->check(CLI::Range(0, 1000000000)); + app.add_option( + "--heavy-frame-zstd-level", heavy_frame_zstd_level, + "zstd compression level for heavy (size-class) frames; negative = faster/larger") + ->default_val(pj_bridge::kDefaultHeavyFrameZstdLevel) + ->check(CLI::Range(ZSTD_minCLevel(), ZSTD_maxCLevel())); // Bound variable is already initialized to {".*"} (match everything); CLI11 // leaves it untouched if the flag is not passed, so no default_val() is // needed (and default_val() on a vector would round-trip through a @@ -94,6 +101,7 @@ int main(int argc, char* argv[]) { spdlog::info(" Topic poll interval: {:.1f} s", topic_poll_interval); spdlog::info(" Client backlog size: {}", client_backlog_size); spdlog::info(" Heavy frame threshold: {} bytes", heavy_frame_threshold_bytes); + spdlog::info(" Heavy frame zstd level: {}", heavy_frame_zstd_level); spdlog::info(" TLS: {}", tls_enabled ? "enabled" : "disabled"); auto whitelist_result = pj_bridge::WhitelistFilter::create(topic_whitelist); @@ -128,7 +136,7 @@ int main(int argc, char* argv[]) { pj_bridge::BridgeServer server( topic_source, sub_manager, middleware, {port, session_timeout, publish_rate, std::move(whitelist_result.value()), - static_cast(heavy_frame_threshold_bytes)}); + static_cast(heavy_frame_threshold_bytes), heavy_frame_zstd_level}); pj_bridge::run_standalone_event_loop( server, sub_manager, middleware, {port, publish_rate, session_timeout, stats_enabled, topic_poll_interval}); 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 44236ba..98f3cd1 100644 --- a/ros2/src/generic_subscription_manager.cpp +++ b/ros2/src/generic_subscription_manager.cpp @@ -171,6 +171,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 8ab897e..cf1a798 100644 --- a/ros2/src/main.cpp +++ b/ros2/src/main.cpp @@ -18,13 +18,14 @@ */ #include +#include #include #include #include -#include #include #include +#include #include #include "pj_bridge/bridge_server.hpp" @@ -52,6 +53,7 @@ int main(int argc, char** argv) { node->declare_parameter("topic_poll_interval", 1.0); node->declare_parameter("client_backlog_size", 100); node->declare_parameter("heavy_frame_threshold_bytes", 262144); + node->declare_parameter("heavy_frame_zstd_level", pj_bridge::kDefaultHeavyFrameZstdLevel); node->declare_parameter("tls", false); node->declare_parameter("certfile", ""); node->declare_parameter("keyfile", ""); @@ -67,6 +69,7 @@ int main(int argc, char** argv) { 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(); int64_t heavy_frame_threshold_bytes = node->get_parameter("heavy_frame_threshold_bytes").as_int(); + int64_t heavy_frame_zstd_level = node->get_parameter("heavy_frame_zstd_level").as_int(); bool tls_enabled = node->get_parameter("tls").as_bool(); std::string certfile = node->get_parameter("certfile").as_string(); std::string keyfile = node->get_parameter("keyfile").as_string(); @@ -74,9 +77,11 @@ 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, heavy_frame_zstd_level=%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, heavy_frame_zstd_level, + 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"); @@ -92,6 +97,22 @@ 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 (heavy_frame_zstd_level < ZSTD_minCLevel() || heavy_frame_zstd_level > ZSTD_maxCLevel()) { + RCLCPP_ERROR( + node->get_logger(), "Invalid heavy_frame_zstd_level: %ld (must be in [%d, %d])", heavy_frame_zstd_level, + ZSTD_minCLevel(), ZSTD_maxCLevel()); + 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(); @@ -148,7 +169,7 @@ int main(int argc, char** argv) { pj_bridge::BridgeServer server( topic_source, sub_manager, middleware, {port, session_timeout, publish_rate, std::move(whitelist_result.value()), - static_cast(heavy_frame_threshold_bytes)}); + static_cast(heavy_frame_threshold_bytes), static_cast(heavy_frame_zstd_level)}); if (!server.initialize()) { RCLCPP_ERROR(node->get_logger(), "Failed to initialize bridge server"); @@ -165,14 +186,16 @@ int main(int argc, char** argv) { // 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 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(); }, timer_group); - auto timeout_timer = node->create_wall_timer(1s, [&server]() { server.check_session_timeouts(); }, timer_group); + 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; @@ -202,6 +225,11 @@ int main(int argc, char** argv) { 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); } 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/rti/src/main.cpp b/rti/src/main.cpp index f419961..b363262 100644 --- a/rti/src/main.cpp +++ b/rti/src/main.cpp @@ -18,6 +18,7 @@ */ #include +#include #include #include @@ -44,6 +45,7 @@ int main(int argc, char* argv[]) { double topic_poll_interval = 1.0; int client_backlog_size = 100; int heavy_frame_threshold_bytes = 262144; + int heavy_frame_zstd_level = pj_bridge::kDefaultHeavyFrameZstdLevel; std::string certfile; std::string keyfile; @@ -69,6 +71,11 @@ int main(int argc, char* argv[]) { "0 disables (keep below the 1 MiB socket watermark)") ->default_val(262144) ->check(CLI::Range(0, 1000000000)); + app.add_option( + "--heavy-frame-zstd-level", heavy_frame_zstd_level, + "zstd compression level for heavy (size-class) frames; negative = faster/larger") + ->default_val(pj_bridge::kDefaultHeavyFrameZstdLevel) + ->check(CLI::Range(ZSTD_minCLevel(), ZSTD_maxCLevel())); // Bound variable is already initialized to {".*"} (match everything); CLI11 // leaves it untouched if the flag is not passed, so no default_val() is // needed (and default_val() on a vector would round-trip through a @@ -99,6 +106,7 @@ int main(int argc, char* argv[]) { spdlog::info(" Topic poll interval: {:.1f} s", topic_poll_interval); spdlog::info(" Client backlog size: {}", client_backlog_size); spdlog::info(" Heavy frame threshold: {} bytes", heavy_frame_threshold_bytes); + spdlog::info(" Heavy frame zstd level: {}", heavy_frame_zstd_level); spdlog::info(" TLS: {}", tls_enabled ? "enabled" : "disabled"); auto whitelist_result = pj_bridge::WhitelistFilter::create(topic_whitelist); @@ -133,7 +141,7 @@ int main(int argc, char* argv[]) { pj_bridge::BridgeServer server( topic_source, sub_manager, middleware, {port, session_timeout, publish_rate, std::move(whitelist_result.value()), - static_cast(heavy_frame_threshold_bytes)}); + static_cast(heavy_frame_threshold_bytes), heavy_frame_zstd_level}); pj_bridge::run_standalone_event_loop( server, sub_manager, middleware, {port, publish_rate, session_timeout, stats_enabled, topic_poll_interval}); diff --git a/tests/unit/test_generic_subscription_manager.cpp b/tests/unit/test_generic_subscription_manager.cpp index 4b2c653..1b02653 100644 --- a/tests/unit/test_generic_subscription_manager.cpp +++ b/tests/unit/test_generic_subscription_manager.cpp @@ -451,15 +451,13 @@ 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()); + 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) { + while (node_->count_publishers("/drain_burst_topic") == 0 && std::chrono::steady_clock::now() < deadline) { std::this_thread::sleep_for(std::chrono::milliseconds(10)); } } @@ -514,7 +512,7 @@ TEST_F(GenericSubscriptionManagerTest, DrainsAllPendingMessagesInOneExecutorPass 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"; + << 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_serializer.cpp b/tests/unit/test_message_serializer.cpp index 12a02a3..d18bd4c 100644 --- a/tests/unit/test_message_serializer.cpp +++ b/tests/unit/test_message_serializer.cpp @@ -530,3 +530,39 @@ TEST_F(MessageSerializerTest, FinalizeFlagsDoNotAlterPayload) { std::vector heavy_payload(heavy.begin() + 16, heavy.end()); EXPECT_EQ(plain_payload, heavy_payload); } + +// ============================================================================ +// Compression level — heavy_frame_zstd_level knob (finalize's 2nd argument) +// ============================================================================ + +TEST_F(MessageSerializerTest, FinalizeCompressionLevelRoundTrips) { + // A negative (fast) and a positive (default-ish) level should both decode + // to the identical uncompressed payload, with header fields unaffected by + // the compression level. + std::vector raw_data = {10, 20, 30, 40, 50}; + auto data = create_test_data(raw_data); + serializer_.serialize_message("/test", 12345, data.data(), data.size()); + const auto expected_payload = serializer_.get_serialized_data(); + + for (int level : {-5, 3}) { + auto result = serializer_.finalize(0, level); + ASSERT_GE(result.size(), 16u); + + uint32_t count; + std::memcpy(&count, result.data() + 4, sizeof(count)); + EXPECT_EQ(count, 1u); + + uint32_t uncompressed_size; + std::memcpy(&uncompressed_size, result.data() + 8, sizeof(uncompressed_size)); + EXPECT_EQ(uncompressed_size, expected_payload.size()); + + uint32_t flags; + std::memcpy(&flags, result.data() + 12, sizeof(flags)); + EXPECT_EQ(flags, 0u); + + std::vector compressed(result.begin() + 16, result.end()); + std::vector decompressed; + AggregatedMessageSerializer::decompress_zstd(compressed, decompressed); + EXPECT_EQ(decompressed, expected_payload); + } +} From c57a3130bd9c259f856c488dcaf37a7b3dfb7bbd Mon Sep 17 00:00:00 2001 From: Davide Faconti Date: Sun, 20 Sep 2026 13:52:38 +0200 Subject: [PATCH 4/9] fix(ros2): address review of the ingest rework - Drain moved into a GenericSubscription subclass (handle_serialized_message) instead of a callback holding a weak_ptr to its own subscription: the pointer was assigned after the subscription became executable, racing with the ingest thread. - Receive time now comes from the RMW received_timestamp when available; with batched delivery "now" would cluster timestamps and distort per-client rate limiting. - MessageBuffer::set_latched() seeds the retained sample from the buffer: with ingest on its own thread the sole latched sample can arrive before the subscriber's bookkeeping, and later clients would get no replay. - Timer thread is joined on every exit path and reports exceptions instead of terminating; take errors during drain are caught. Builds and passes subscription tests on Humble and Jazzy. Co-Authored-By: Claude Fable 5.1 --- app/src/message_buffer.cpp | 6 ++ ros2/src/generic_subscription_manager.cpp | 92 +++++++++++++++++------ ros2/src/main.cpp | 22 +++++- tests/unit/test_message_buffer.cpp | 13 ++++ 4 files changed, 108 insertions(+), 25 deletions(-) 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/ros2/src/generic_subscription_manager.cpp b/ros2/src/generic_subscription_manager.cpp index 98f3cd1..dfecb98 100644 --- a/ros2/src/generic_subscription_manager.cpp +++ b/ros2/src/generic_subscription_manager.cpp @@ -19,7 +19,10 @@ #include "pj_bridge_ros2/generic_subscription_manager.hpp" +#include + #include +#include #include "pj_bridge/time_utils.hpp" @@ -27,6 +30,65 @@ 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( @@ -107,33 +169,17 @@ bool GenericSubscriptionManager::subscribe( } try { - // The executor hands over one message per subscription per wait cycle, and - // each cycle rebuilds the whole wait set (O(subscriptions)). Drain whatever - // else the reader holds so a burst costs one cycle instead of one per message. - auto self = std::make_shared>(); - auto sub_callback = [topic_name, callback, self](std::shared_ptr msg) { - callback(topic_name, msg, get_current_time_ns()); - auto sub = self->lock(); - if (!sub) { - return; - } - // 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 info; - if (!sub->take_serialized(*extra, info)) { - break; - } - callback(topic_name, extra, get_current_time_ns()); - } - }; - 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); - *self = subscription; subscriptions_[topic_name] = SubscriptionInfo{subscription, 1, transient_local}; return true; diff --git a/ros2/src/main.cpp b/ros2/src/main.cpp index cf1a798..6290ee8 100644 --- a/ros2/src/main.cpp +++ b/ros2/src/main.cpp @@ -213,7 +213,25 @@ int main(int argc, char** argv) { // long publish cycle never delays ingest. rclcpp::executors::SingleThreadedExecutor timer_executor; timer_executor.add_callback_group(timer_group, node->get_node_base_interface()); - std::thread timer_thread([&timer_executor]() { timer_executor.spin(); }); + // 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 (>50% of @@ -238,7 +256,7 @@ int main(int argc, char** argv) { } timer_executor.cancel(); - timer_thread.join(); + timer_thread.thread.join(); // Graceful shutdown RCLCPP_INFO(node->get_logger(), "Shutting down bridge server..."); 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})); From 98c140b6ebabb4290c80e11b30da26aa67f6f583 Mon Sep 17 00:00:00 2001 From: Davide Faconti Date: Sun, 20 Sep 2026 14:30:35 +0200 Subject: [PATCH 5/9] build: link ROS2 targets directly instead of ament_target_dependencies ament_target_dependencies was removed in newer distros, so the package did not configure on Lyrical. Modern targets work on Humble, Jazzy and Lyrical. Co-Authored-By: Claude Fable 5.1 --- CMakeLists.txt | 21 ++++++--------------- 1 file changed, 6 insertions(+), 15 deletions(-) diff --git a/CMakeLists.txt b/CMakeLists.txt index f941358..a2d0ed5 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -149,11 +149,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 +166,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 +273,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 From 8367cabfddd379dc58d60dc25520b27a521a0c48 Mon Sep 17 00:00:00 2001 From: Davide Faconti Date: Sun, 20 Sep 2026 14:39:53 +0200 Subject: [PATCH 6/9] ci: build/test on Lyrical and release a Lyrical .deb No robostack-lyrical channel exists, so there is no pixi/conda/AppImage variant for Lyrical; colcon CI and the Debian release are covered. Co-Authored-By: Claude Fable 5.1 --- .github/workflows/release-debs.yaml | 4 ++++ .github/workflows/ros_lyrical.yaml | 20 ++++++++++++++++++++ README.md | 10 +++++----- 3 files changed, 29 insertions(+), 5 deletions(-) create mode 100644 .github/workflows/ros_lyrical.yaml diff --git a/.github/workflows/release-debs.yaml b/.github/workflows/release-debs.yaml index 3adaf15..94de8ba 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: 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/README.md b/README.md index 75c2563..8618324 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 From 8bacdcfdf201baecd0b1dfc5308ed44655587145 Mon Sep 17 00:00:00 2001 From: Davide Faconti Date: Sun, 20 Sep 2026 14:45:14 +0200 Subject: [PATCH 7/9] refactor: drop heavy_frame_zstd_level, no measurements in comments All frames keep zstd level 1 as before; the knob was a CPU-vs-bandwidth tradeoff not worth the configuration surface. FastDDS/RTI entry points are back to identical with main. Co-Authored-By: Claude Fable 5.1 --- CHANGELOG.rst | 3 -- CLAUDE.md | 7 ++-- README.md | 2 -- app/include/pj_bridge/bridge_server.hpp | 7 ---- app/include/pj_bridge/message_serializer.hpp | 5 +-- app/include/pj_bridge/protocol_constants.hpp | 5 --- app/src/bridge_server.cpp | 3 +- app/src/message_serializer.cpp | 5 +-- app/src/middleware/websocket_middleware.cpp | 3 +- docs/API.md | 16 --------- fastdds/src/main.cpp | 10 +----- ros2/src/main.cpp | 26 ++++---------- rti/src/main.cpp | 10 +----- tests/unit/test_message_serializer.cpp | 36 -------------------- 14 files changed, 17 insertions(+), 121 deletions(-) diff --git a/CHANGELOG.rst b/CHANGELOG.rst index 1d150d2..e0e39c0 100644 --- a/CHANGELOG.rst +++ b/CHANGELOG.rst @@ -12,9 +12,6 @@ Forthcoming 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. -* New ``heavy_frame_zstd_level`` knob (ROS2 param, RTI/FastDDS - ``--heavy-frame-zstd-level``, default ``1``): CPU-vs-bandwidth zstd level - for heavy (size-class) frames only; wire-compatible (frame flags stay 0). * Removed zero-initializing copies in the ingest and serializer hot paths (``resize()`` + ``memcpy`` -> direct-construct/``insert``). diff --git a/CLAUDE.md b/CLAUDE.md index 662f970..1a7e8f1 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -240,7 +240,6 @@ ingest_poll_interval_ms: 5.0 # Ingest executor poll interval (drains subscript 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 -heavy_frame_zstd_level: 1 # zstd level for heavy frames; any zstd level, negative = faster/larger tls: false # Enable TLS (wss://); requires certfile and keyfile certfile: "" # TLS server certificate file keyfile: "" # TLS private key file @@ -250,16 +249,14 @@ keyfile: "" # TLS private key file ```bash pj_bridge_rti --domains 0 1 --port 9090 --publish-rate 50 --session-timeout 10 \ --topic-whitelist ".*" --topic-poll-interval 1.0 --client-backlog-size 100 \ - --heavy-frame-threshold-bytes 262144 --heavy-frame-zstd-level 1 \ - --certfile cert.pem --keyfile key.pem + --heavy-frame-threshold-bytes 262144 --certfile cert.pem --keyfile key.pem ``` ### FastDDS (via CLI flags): ```bash pj_bridge_fastdds --domains 0 1 --port 9090 --publish-rate 50 --session-timeout 10 \ --topic-whitelist ".*" --topic-poll-interval 1.0 --client-backlog-size 100 \ - --heavy-frame-threshold-bytes 262144 --heavy-frame-zstd-level 1 \ - --certfile cert.pem --keyfile key.pem + --heavy-frame-threshold-bytes 262144 --certfile cert.pem --keyfile key.pem ``` See `docs/API.md` for full semantics of each option (topic whitelist matching rules, diff --git a/README.md b/README.md index 8618324..c066cdb 100644 --- a/README.md +++ b/README.md @@ -56,7 +56,6 @@ independently. | `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`) | -| `heavy_frame_zstd_level` | int | 1 | zstd compression level for heavy (size-class) frames; any zstd level, negative = faster/larger | | `tls` | bool | false | Enable TLS (`wss://`); requires `certfile` and `keyfile` | | `certfile` | string | `""` | TLS server certificate file | | `keyfile` | string | `""` | TLS private key file | @@ -75,7 +74,6 @@ independently. | `--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 (range `1`-`1000000`) | | `--heavy-frame-threshold-bytes` | int | 262144 | Isolate messages ≥ this size (bytes) into their own size-class frame; `0` disables (range `0`-`1000000000`) | -| `--heavy-frame-zstd-level` | int | 1 | zstd compression level for heavy (size-class) frames; negative = faster/larger (range: zstd min-max level) | | `--certfile` | string | (none) | TLS server certificate file; enables `wss://`, requires `--keyfile` | | `--keyfile` | string | (none) | TLS private key file; enables `wss://`, requires `--certfile` | | `--qos-profile` | string | (none) | RTI only: QoS profile XML file path | diff --git a/app/include/pj_bridge/bridge_server.hpp b/app/include/pj_bridge/bridge_server.hpp index 95c01e3..7ce7f42 100644 --- a/app/include/pj_bridge/bridge_server.hpp +++ b/app/include/pj_bridge/bridge_server.hpp @@ -61,12 +61,6 @@ struct BridgeServerConfig { /// its own size-class ("heavy") frame instead of being aggregated with light /// topics. 0 disables splitting (single aggregated frame). Default: 256 KiB. size_t heavy_frame_threshold_bytes = kDefaultHeavyFrameThresholdBytes; - /// zstd level for heavy frames; a CPU-vs-bandwidth knob (any zstd level, - /// negative = faster). Compression is ~60% of CPU on point-cloud workloads. - /// Measured on 4 lidars, 52 MiB/s in: level 1 = 19% of a core, 31 MB/s out; - /// -5 = 15%, 39 MB/s; -100 = 10%, 51 MB/s. Light frames always use - /// kDefaultZstdLevel. - int heavy_frame_zstd_level = kDefaultHeavyFrameZstdLevel; }; class BridgeServer { @@ -232,7 +226,6 @@ class BridgeServer { // Per-message byte size at or above which a topic is isolated into its own // size-class ("heavy") frame; 0 disables splitting. See publish_aggregated_messages(). size_t heavy_frame_threshold_bytes_; - int heavy_frame_zstd_level_; // State std::atomic initialized_; diff --git a/app/include/pj_bridge/message_serializer.hpp b/app/include/pj_bridge/message_serializer.hpp index 018ffe8..9df4185 100644 --- a/app/include/pj_bridge/message_serializer.hpp +++ b/app/include/pj_bridge/message_serializer.hpp @@ -26,8 +26,6 @@ #include #include -#include "pj_bridge/protocol_constants.hpp" - namespace pj_bridge { /** @@ -98,8 +96,7 @@ class AggregatedMessageSerializer { * kFrameFlagHeavy to mark an isolated large/size-class frame. * @return Vector containing header + compressed payload */ - /// @param compression_level any zstd level; negative levels trade ratio for speed. - std::vector finalize(uint32_t flags = 0, int compression_level = kDefaultZstdLevel); + std::vector finalize(uint32_t flags = 0); /** * @brief Compress data using ZSTD (compression level 1) diff --git a/app/include/pj_bridge/protocol_constants.hpp b/app/include/pj_bridge/protocol_constants.hpp index fb12a80..af57f3c 100644 --- a/app/include/pj_bridge/protocol_constants.hpp +++ b/app/include/pj_bridge/protocol_constants.hpp @@ -48,11 +48,6 @@ static constexpr uint32_t kFrameFlagHeavy = 0x1; /// threshold of 0 disables splitting (single aggregated frame, legacy behavior). static constexpr size_t kDefaultHeavyFrameThresholdBytes = 256 * 1024; // 256 KiB -/// zstd level for aggregated (light) frames. -static constexpr int kDefaultZstdLevel = 1; -/// zstd level for heavy frames; see BridgeServerConfig::heavy_frame_zstd_level. -static constexpr int kDefaultHeavyFrameZstdLevel = 1; - /// Schema encoding identifier for ROS2 message definitions inline constexpr const char* kSchemaEncodingRos2Msg = "ros2msg"; diff --git a/app/src/bridge_server.cpp b/app/src/bridge_server.cpp index a167540..fe3551e 100644 --- a/app/src/bridge_server.cpp +++ b/app/src/bridge_server.cpp @@ -69,7 +69,6 @@ BridgeServer::BridgeServer( publish_rate_(config.publish_rate), whitelist_(std::move(config.whitelist)), heavy_frame_threshold_bytes_(config.heavy_frame_threshold_bytes), - heavy_frame_zstd_level_(config.heavy_frame_zstd_level), initialized_(false), total_messages_published_(0), total_bytes_published_(0), @@ -1100,7 +1099,7 @@ void BridgeServer::publish_aggregated_messages() { AggregatedMessageSerializer heavy_serializer; heavy_serializer.serialize_message(topic, msg.timestamp_ns, msg.data->data(), msg.data->size()); GroupFrame heavy_frame; - heavy_frame.compressed_data = heavy_serializer.finalize(0, heavy_frame_zstd_level_); + heavy_frame.compressed_data = heavy_serializer.finalize(); heavy_frame.msg_count = 1; heavy_frame.client_ids = client_ids; heavy_frame.is_heavy = true; diff --git a/app/src/message_serializer.cpp b/app/src/message_serializer.cpp index a3c3453..72c70a2 100644 --- a/app/src/message_serializer.cpp +++ b/app/src/message_serializer.cpp @@ -67,7 +67,7 @@ size_t AggregatedMessageSerializer::get_message_count() const { return message_count_; } -std::vector AggregatedMessageSerializer::finalize(uint32_t flags, int compression_level) { +std::vector AggregatedMessageSerializer::finalize(uint32_t flags) { // Build 16-byte header (uncompressed) std::vector header(kBinaryHeaderSize); @@ -101,7 +101,8 @@ std::vector AggregatedMessageSerializer::finalize(uint32_t flags, int c // Compress payload after header using persistent context size_t compressed_size = ZSTD_compressCCtx( cctx_, result.data() + kBinaryHeaderSize, max_compressed, serialized_data_.data(), serialized_data_.size(), - compression_level); + 1 // compression level + ); if (ZSTD_isError(compressed_size)) { throw std::runtime_error(std::string("ZSTD compression failed: ") + ZSTD_getErrorName(compressed_size)); diff --git a/app/src/middleware/websocket_middleware.cpp b/app/src/middleware/websocket_middleware.cpp index 627f602..96c18f7 100644 --- a/app/src/middleware/websocket_middleware.cpp +++ b/app/src/middleware/websocket_middleware.cpp @@ -74,8 +74,7 @@ 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 costs more - // CPU than the zstd pass itself and gains nothing. + // Binary frames are already zstd-compressed; deflating them again is wasted work. server_->disablePerMessageDeflate(); #ifdef IXWEBSOCKET_USE_TLS diff --git a/docs/API.md b/docs/API.md index 718f4ba..52f9d8e 100644 --- a/docs/API.md +++ b/docs/API.md @@ -588,22 +588,6 @@ frame is configurable: Keep the threshold below the 1 MiB socket watermark so a single heavy message does not fill the socket buffer on its own. -Heavy frames use a separately configurable zstd compression level — a -CPU-vs-bandwidth knob independent of the threshold above. Light (aggregated) -frames always compress at level `1`. Any zstd level is valid (negative -levels trade ratio for speed); the frame `flags` stay `0` either way, so any -zstd decoder accepts the output regardless of level: - -- **ROS2**: int parameter `heavy_frame_zstd_level`, default `1`. Must be in - `[ZSTD_minCLevel(), ZSTD_maxCLevel()]`; the server refuses to start - otherwise. -- **FastDDS / RTI**: CLI flag `--heavy-frame-zstd-level`, default `1`, same - valid range. - -Measured on 4 lidars, 52 MiB/s in, 1 client (i7-13700H, P-core pinned, -performance governor): level `1` = 19.0% of a core / 31.0 MB/s out; `-5` = -14.8% / 39.0; `-20` = 14.4% / 45.6; `-100` = 10.3% / 50.6. - ## TLS / wss:// The bridge can optionally serve the WebSocket endpoint over TLS (`wss://`) diff --git a/fastdds/src/main.cpp b/fastdds/src/main.cpp index 19861b0..bbd6b20 100644 --- a/fastdds/src/main.cpp +++ b/fastdds/src/main.cpp @@ -18,7 +18,6 @@ */ #include -#include #include #include @@ -44,7 +43,6 @@ int main(int argc, char* argv[]) { double topic_poll_interval = 1.0; int client_backlog_size = 100; int heavy_frame_threshold_bytes = 262144; - int heavy_frame_zstd_level = pj_bridge::kDefaultHeavyFrameZstdLevel; std::string certfile; std::string keyfile; @@ -69,11 +67,6 @@ int main(int argc, char* argv[]) { "0 disables (keep below the 1 MiB socket watermark)") ->default_val(262144) ->check(CLI::Range(0, 1000000000)); - app.add_option( - "--heavy-frame-zstd-level", heavy_frame_zstd_level, - "zstd compression level for heavy (size-class) frames; negative = faster/larger") - ->default_val(pj_bridge::kDefaultHeavyFrameZstdLevel) - ->check(CLI::Range(ZSTD_minCLevel(), ZSTD_maxCLevel())); // Bound variable is already initialized to {".*"} (match everything); CLI11 // leaves it untouched if the flag is not passed, so no default_val() is // needed (and default_val() on a vector would round-trip through a @@ -101,7 +94,6 @@ int main(int argc, char* argv[]) { spdlog::info(" Topic poll interval: {:.1f} s", topic_poll_interval); spdlog::info(" Client backlog size: {}", client_backlog_size); spdlog::info(" Heavy frame threshold: {} bytes", heavy_frame_threshold_bytes); - spdlog::info(" Heavy frame zstd level: {}", heavy_frame_zstd_level); spdlog::info(" TLS: {}", tls_enabled ? "enabled" : "disabled"); auto whitelist_result = pj_bridge::WhitelistFilter::create(topic_whitelist); @@ -136,7 +128,7 @@ int main(int argc, char* argv[]) { pj_bridge::BridgeServer server( topic_source, sub_manager, middleware, {port, session_timeout, publish_rate, std::move(whitelist_result.value()), - static_cast(heavy_frame_threshold_bytes), heavy_frame_zstd_level}); + static_cast(heavy_frame_threshold_bytes)}); pj_bridge::run_standalone_event_loop( server, sub_manager, middleware, {port, publish_rate, session_timeout, stats_enabled, topic_poll_interval}); diff --git a/ros2/src/main.cpp b/ros2/src/main.cpp index 6290ee8..1de482e 100644 --- a/ros2/src/main.cpp +++ b/ros2/src/main.cpp @@ -18,7 +18,6 @@ */ #include -#include #include #include @@ -53,7 +52,6 @@ int main(int argc, char** argv) { node->declare_parameter("topic_poll_interval", 1.0); node->declare_parameter("client_backlog_size", 100); node->declare_parameter("heavy_frame_threshold_bytes", 262144); - node->declare_parameter("heavy_frame_zstd_level", pj_bridge::kDefaultHeavyFrameZstdLevel); node->declare_parameter("tls", false); node->declare_parameter("certfile", ""); node->declare_parameter("keyfile", ""); @@ -69,7 +67,6 @@ int main(int argc, char** argv) { 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(); int64_t heavy_frame_threshold_bytes = node->get_parameter("heavy_frame_threshold_bytes").as_int(); - int64_t heavy_frame_zstd_level = node->get_parameter("heavy_frame_zstd_level").as_int(); bool tls_enabled = node->get_parameter("tls").as_bool(); std::string certfile = node->get_parameter("certfile").as_string(); std::string keyfile = node->get_parameter("keyfile").as_string(); @@ -78,10 +75,9 @@ int main(int argc, char** argv) { 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, ingest_poll_interval_ms=%.1f, topic_poll_interval=%.1f s, " - "client_backlog_size=%ld, heavy_frame_zstd_level=%ld, tls=%s", + "client_backlog_size=%ld, tls=%s", port, publish_rate, session_timeout, strip_large_messages ? "true" : "false", min_qos_depth, max_qos_depth, - ingest_poll_interval_ms, topic_poll_interval, client_backlog_size, heavy_frame_zstd_level, - 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"); @@ -105,14 +101,6 @@ int main(int argc, char** argv) { return 1; } - if (heavy_frame_zstd_level < ZSTD_minCLevel() || heavy_frame_zstd_level > ZSTD_maxCLevel()) { - RCLCPP_ERROR( - node->get_logger(), "Invalid heavy_frame_zstd_level: %ld (must be in [%d, %d])", heavy_frame_zstd_level, - ZSTD_minCLevel(), ZSTD_maxCLevel()); - 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(); @@ -169,7 +157,7 @@ int main(int argc, char** argv) { pj_bridge::BridgeServer server( topic_source, sub_manager, middleware, {port, session_timeout, publish_rate, std::move(whitelist_result.value()), - static_cast(heavy_frame_threshold_bytes), static_cast(heavy_frame_zstd_level)}); + static_cast(heavy_frame_threshold_bytes)}); if (!server.initialize()) { RCLCPP_ERROR(node->get_logger(), "Failed to initialize bridge server"); @@ -234,10 +222,10 @@ int main(int argc, char** argv) { })}; // Every executor wait cycle rebuilds the whole wait set, O(subscriptions); - // with blocking spin() that is one rebuild per received message (>50% of - // CPU with ~80 subscriptions). 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. + // 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) { diff --git a/rti/src/main.cpp b/rti/src/main.cpp index b363262..f419961 100644 --- a/rti/src/main.cpp +++ b/rti/src/main.cpp @@ -18,7 +18,6 @@ */ #include -#include #include #include @@ -45,7 +44,6 @@ int main(int argc, char* argv[]) { double topic_poll_interval = 1.0; int client_backlog_size = 100; int heavy_frame_threshold_bytes = 262144; - int heavy_frame_zstd_level = pj_bridge::kDefaultHeavyFrameZstdLevel; std::string certfile; std::string keyfile; @@ -71,11 +69,6 @@ int main(int argc, char* argv[]) { "0 disables (keep below the 1 MiB socket watermark)") ->default_val(262144) ->check(CLI::Range(0, 1000000000)); - app.add_option( - "--heavy-frame-zstd-level", heavy_frame_zstd_level, - "zstd compression level for heavy (size-class) frames; negative = faster/larger") - ->default_val(pj_bridge::kDefaultHeavyFrameZstdLevel) - ->check(CLI::Range(ZSTD_minCLevel(), ZSTD_maxCLevel())); // Bound variable is already initialized to {".*"} (match everything); CLI11 // leaves it untouched if the flag is not passed, so no default_val() is // needed (and default_val() on a vector would round-trip through a @@ -106,7 +99,6 @@ int main(int argc, char* argv[]) { spdlog::info(" Topic poll interval: {:.1f} s", topic_poll_interval); spdlog::info(" Client backlog size: {}", client_backlog_size); spdlog::info(" Heavy frame threshold: {} bytes", heavy_frame_threshold_bytes); - spdlog::info(" Heavy frame zstd level: {}", heavy_frame_zstd_level); spdlog::info(" TLS: {}", tls_enabled ? "enabled" : "disabled"); auto whitelist_result = pj_bridge::WhitelistFilter::create(topic_whitelist); @@ -141,7 +133,7 @@ int main(int argc, char* argv[]) { pj_bridge::BridgeServer server( topic_source, sub_manager, middleware, {port, session_timeout, publish_rate, std::move(whitelist_result.value()), - static_cast(heavy_frame_threshold_bytes), heavy_frame_zstd_level}); + static_cast(heavy_frame_threshold_bytes)}); pj_bridge::run_standalone_event_loop( server, sub_manager, middleware, {port, publish_rate, session_timeout, stats_enabled, topic_poll_interval}); diff --git a/tests/unit/test_message_serializer.cpp b/tests/unit/test_message_serializer.cpp index d18bd4c..12a02a3 100644 --- a/tests/unit/test_message_serializer.cpp +++ b/tests/unit/test_message_serializer.cpp @@ -530,39 +530,3 @@ TEST_F(MessageSerializerTest, FinalizeFlagsDoNotAlterPayload) { std::vector heavy_payload(heavy.begin() + 16, heavy.end()); EXPECT_EQ(plain_payload, heavy_payload); } - -// ============================================================================ -// Compression level — heavy_frame_zstd_level knob (finalize's 2nd argument) -// ============================================================================ - -TEST_F(MessageSerializerTest, FinalizeCompressionLevelRoundTrips) { - // A negative (fast) and a positive (default-ish) level should both decode - // to the identical uncompressed payload, with header fields unaffected by - // the compression level. - std::vector raw_data = {10, 20, 30, 40, 50}; - auto data = create_test_data(raw_data); - serializer_.serialize_message("/test", 12345, data.data(), data.size()); - const auto expected_payload = serializer_.get_serialized_data(); - - for (int level : {-5, 3}) { - auto result = serializer_.finalize(0, level); - ASSERT_GE(result.size(), 16u); - - uint32_t count; - std::memcpy(&count, result.data() + 4, sizeof(count)); - EXPECT_EQ(count, 1u); - - uint32_t uncompressed_size; - std::memcpy(&uncompressed_size, result.data() + 8, sizeof(uncompressed_size)); - EXPECT_EQ(uncompressed_size, expected_payload.size()); - - uint32_t flags; - std::memcpy(&flags, result.data() + 12, sizeof(flags)); - EXPECT_EQ(flags, 0u); - - std::vector compressed(result.begin() + 16, result.end()); - std::vector decompressed; - AggregatedMessageSerializer::decompress_zstd(compressed, decompressed); - EXPECT_EQ(decompressed, expected_payload); - } -} From 8dd193656a0c6993be6f1734c4884e0fdad4171d Mon Sep 17 00:00:00 2001 From: Davide Faconti Date: Sun, 20 Sep 2026 14:54:44 +0200 Subject: [PATCH 8/9] build: IXWebSocket 12.0.1, Fast DDS 3.4.3, CLI11 2.6.2; fix fetched-TLS link IXWebSocket 12.x is fixes only (server header parsing, fd double close, DNS crash, move instead of copy in send buffering). pixi.lock is left alone: moving the Humble env to 12.x re-solves the whole RoboStack stack. The fetched static IXWebSocket does not export its OpenSSL dependency, so the standalone FastDDS/RTI link failed with TLS on (also on 11.4.6). Co-Authored-By: Claude Fable 5.1 --- .github/workflows/release-debs.yaml | 6 +++--- CHANGELOG.rst | 4 ++++ CLAUDE.md | 4 ++-- CMakeLists.txt | 9 +++++++-- conanfile.txt | 4 ++-- docs/ARCHITECTURE.md | 2 +- 6 files changed, 19 insertions(+), 10 deletions(-) diff --git a/.github/workflows/release-debs.yaml b/.github/workflows/release-debs.yaml index 94de8ba..3964e4f 100644 --- a/.github/workflows/release-debs.yaml +++ b/.github/workflows/release-debs.yaml @@ -53,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/CHANGELOG.rst b/CHANGELOG.rst index e0e39c0..ca84bcf 100644 --- a/CHANGELOG.rst +++ b/CHANGELOG.rst @@ -12,6 +12,10 @@ Forthcoming 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``). diff --git a/CLAUDE.md b/CLAUDE.md index 1a7e8f1..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) @@ -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) diff --git a/CMakeLists.txt b/CMakeLists.txt index a2d0ed5..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() 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/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 From ecad4d25d4b194c6b51683d07a9c6da3a77381e1 Mon Sep 17 00:00:00 2001 From: Davide Faconti Date: Sun, 20 Sep 2026 15:06:06 +0200 Subject: [PATCH 9/9] test: pin rcl_logging_spdlog so the test binary does not segfault at exit On Lyrical, rcl_logging_spdlog registers a periodic flush thread on the global spdlog registry and rcl unloads the library on shutdown. After the suites' repeated init/shutdown the thread resumed in unmapped code at exit (all tests passed, process returned -11). The bridge binary is unaffected. Co-Authored-By: Claude Fable 5.1 --- tests/unit/test_websocket_middleware.cpp | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/tests/unit/test_websocket_middleware.cpp b/tests/unit/test_websocket_middleware.cpp index 7b03401..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 @@ -737,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(); }