Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 7 additions & 3 deletions .github/workflows/release-debs.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -49,9 +53,9 @@ jobs:
- name: Install IXWebSocket from source
run: |
cd /tmp
wget -q https://github.com/machinezone/IXWebSocket/archive/refs/tags/v11.4.6.tar.gz
tar xzf v11.4.6.tar.gz
cd IXWebSocket-11.4.6
wget -q https://github.com/machinezone/IXWebSocket/archive/refs/tags/v12.0.1.tar.gz
tar xzf v12.0.1.tar.gz
cd IXWebSocket-12.0.1
cmake -B build -DCMAKE_BUILD_TYPE=Release -DUSE_TLS=OFF -DUSE_ZLIB=ON
cmake --build build -j$(nproc)
cmake --install build
Expand Down
20 changes: 20 additions & 0 deletions .github/workflows/ros_lyrical.yaml
Original file line number Diff line number Diff line change
@@ -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
17 changes: 17 additions & 0 deletions CHANGELOG.rst
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,23 @@
Changelog for package pj_ros_bridge
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^

Forthcoming
-----------
* WebSocket permessage-deflate declined server-side (redundant with our own
ZSTD compression; was costing CPU for no bandwidth benefit).
* ROS2: ingest executor polled and drained (``ingest_poll_interval_ms``,
default 5 ms) instead of blocking ``spin()``, amortizing the per-wait-cycle
wait-set rebuild; publish/request/timeout timers moved to their own
executor thread so a long publish cycle never delays ingest.
* ROS2: ``min_qos_depth`` default raised ``1`` -> ``10`` so the ingest poll
interval can't overflow a shallow reader between polls.
* Dependencies: IXWebSocket 11.4.6 -> 12.0.1 (FetchContent and .deb builds),
Fast DDS 3.4.0 -> 3.4.3, CLI11 2.6.0 -> 2.6.2. Fixed the FastDDS/RTI
link failure when IXWebSocket is fetched with TLS.
* CI and Debian release for ROS 2 Lyrical.
* Removed zero-initializing copies in the ingest and serializer hot paths
(``resize()`` + ``memcpy`` -> direct-construct/``insert``).

0.9.0 (2026-07-11)
------------------
* Size-class frames: isolate heavy messages (``>= heavy_frame_threshold_bytes``)
Expand Down
9 changes: 5 additions & 4 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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)

Expand Down Expand Up @@ -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)

Expand Down Expand Up @@ -234,8 +234,9 @@ publish_rate: 50.0 # Hz
session_timeout: 10.0 # seconds
strip_large_messages: false # Opt-in: strip Image/PointCloud2/etc data fields
topic_whitelist: [".*"] # Full-match regex patterns restricting visible/subscribable topics
min_qos_depth: 1 # Minimum KEEP_LAST subscription depth after aggregating publisher depths
min_qos_depth: 10 # Minimum KEEP_LAST subscription depth after aggregating publisher depths
max_qos_depth: 100 # Maximum KEEP_LAST subscription depth after aggregating publisher depths
ingest_poll_interval_ms: 5.0 # Ingest executor poll interval (drains subscriptions each poll); 0 = blocking spin
topic_poll_interval: 1.0 # Seconds between topics_changed notification polls; 0 disables polling
client_backlog_size: 100 # Max frames queued per slow client before dropping the oldest (must be > 0)
heavy_frame_threshold_bytes: 262144 # Isolate messages >= this size into their own size-class frame; 0 disables
Expand Down
30 changes: 13 additions & 17 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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()
Expand Down Expand Up @@ -149,11 +154,12 @@ if(ament_cmake_FOUND)
pj_bridge_app
)

ament_target_dependencies(pj_bridge_ros2_lib PUBLIC
ament_index_cpp
rclcpp
sensor_msgs
nav_msgs
# target_link_libraries, not ament_target_dependencies: removed in newer distros (Lyrical)
target_link_libraries(pj_bridge_ros2_lib PUBLIC
ament_index_cpp::ament_index_cpp
rclcpp::rclcpp
${sensor_msgs_TARGETS}
${nav_msgs_TARGETS}
)

# ROS2 executable
Expand All @@ -165,10 +171,6 @@ if(ament_cmake_FOUND)
pj_bridge_ros2_lib
)

ament_target_dependencies(pj_bridge_ros2
rclcpp
)

# Install ROS2 executable
install(TARGETS pj_bridge_ros2
DESTINATION lib/${PROJECT_NAME}
Expand Down Expand Up @@ -276,12 +278,6 @@ if(BUILD_TESTING AND ament_cmake_FOUND)
data_path
)

ament_target_dependencies(${PROJECT_NAME}_tests
rclcpp
sensor_msgs
nav_msgs
)

# Sanitizer environment settings for CTest
if(ENABLE_TSAN)
set_tests_properties(${PROJECT_NAME}_tests PROPERTIES
Expand Down
13 changes: 7 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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

Expand All @@ -50,8 +50,9 @@ independently.
| `session_timeout` | double | 10.0 | Client timeout duration in seconds |
| `strip_large_messages` | bool | false | Opt-in: strip large arrays from Image, PointCloud2, LaserScan, OccupancyGrid messages |
| `topic_whitelist` | string array | `[".*"]` | Full-match regex patterns (ECMAScript) restricting visible/subscribable topics |
| `min_qos_depth` | int | 1 | ROS2 only: minimum KEEP_LAST subscription depth after aggregating publisher depths |
| `min_qos_depth` | int | 10 | ROS2 only: minimum KEEP_LAST subscription depth after aggregating publisher depths |
| `max_qos_depth` | int | 100 | ROS2 only: maximum KEEP_LAST subscription depth after aggregating publisher depths |
| `ingest_poll_interval_ms` | double | 5.0 | ROS2 only: interval between ingest executor polls; subscription callbacks drain their DDS reader each poll. `0` uses a blocking spin instead |
| `topic_poll_interval` | double | 1.0 | Seconds between `topics_changed` notification polls; `0` disables polling |
| `client_backlog_size` | int | 100 | Max binary frames queued per slow client before the oldest is dropped (must be `> 0`) |
| `heavy_frame_threshold_bytes` | int | 262144 | Isolate messages ≥ this size (bytes) into their own size-class frame so they don't starve small topics; `0` disables (must be `>= 0`) |
Expand Down
6 changes: 6 additions & 0 deletions app/src/message_buffer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,12 @@ void MessageBuffer::set_latched(const std::string& topic_name, bool latched) {
std::lock_guard<std::mutex> 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);
Expand Down
6 changes: 3 additions & 3 deletions app/src/message_serializer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<const uint8_t *>(data);
serialized_data_.insert(serialized_data_.end(), bytes, bytes + msg_size);
}

void AggregatedMessageSerializer::clear() {
Expand Down
2 changes: 2 additions & 0 deletions app/src/middleware/websocket_middleware.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,8 @@ tl::expected<void, std::string> WebSocketMiddleware::initialize(uint16_t port) {
}

server_ = std::make_shared<ix::WebSocketServer>(port, "0.0.0.0");
// Binary frames are already zstd-compressed; deflating them again is wasted work.
server_->disablePerMessageDeflate();

#ifdef IXWEBSOCKET_USE_TLS
if (tls_.has_value()) {
Expand Down
4 changes: 2 additions & 2 deletions conanfile.txt
Original file line number Diff line number Diff line change
@@ -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
Expand Down
22 changes: 21 additions & 1 deletion docs/API.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion docs/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
4 changes: 2 additions & 2 deletions fastdds/src/fastdds_subscription_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -80,8 +80,8 @@ class FastDdsSubscriptionManager::InternalReaderListener : public DataReaderList
continue;
}

auto cdr_data = std::make_shared<std::vector<std::byte>>(payload_.length);
std::memcpy(cdr_data->data(), payload_.data, payload_.length);
const auto* bytes = reinterpret_cast<const std::byte*>(payload_.data);
auto cdr_data = std::make_shared<std::vector<std::byte>>(bytes, bytes + payload_.length);

uint64_t timestamp_ns =
static_cast<uint64_t>(info.source_timestamp.seconds()) * 1'000'000'000ULL + info.source_timestamp.nanosec();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions ros2/include/pj_bridge_ros2/ros2_subscription_manager.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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_;
Expand Down
Loading
Loading