Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
fb85084
docs: message transforms design (v1: Cloudini)
facontidavide Sep 20, 2026
1d4e607
docs: message transforms implementation plan; align spec
facontidavide Sep 20, 2026
94eb944
build: require C++20
facontidavide Sep 20, 2026
44f5d98
feat(transform): MessageTransform interface
facontidavide Sep 20, 2026
2e5cc43
feat(transform): TransformSet profile parsing, matching and counters
facontidavide Sep 20, 2026
6936698
feat(transform): advertise transformed type and schema
facontidavide Sep 20, 2026
fcf2588
fix(transform): contain transform exceptions; skip non-accepting rules
facontidavide Sep 20, 2026
ab92dca
feat(transform): ROS2 ingest hook; strip becomes a transform
facontidavide Sep 20, 2026
65aabaa
feat(transform): transform_profile parameter; strip_large_messages as…
facontidavide Sep 20, 2026
c46f8a4
test(transform): failed transform drops the sample on the ROS2 path
facontidavide Sep 20, 2026
4b91f93
feat(transform): Cloudini point cloud compression
facontidavide Sep 20, 2026
ff10ea2
docs: message transforms
facontidavide Sep 20, 2026
1719604
docs: message transform measurements
facontidavide Sep 20, 2026
26ec88b
fix(transform): validate clouds before encoding; harden the Cloudini …
facontidavide Sep 20, 2026
30aded7
test(transform): wait for the ingest thread in the concurrency test
facontidavide Sep 20, 2026
352cafb
refactor(transform): simplify after /simplify review; fix log-once ke…
facontidavide Sep 20, 2026
a61d665
build: Cloudini 1.3.1; drop the embedding workarounds and duplicated …
facontidavide Sep 20, 2026
872b48c
feat(transform)!: match_type is mandatory; a transform/type mismatch …
facontidavide Sep 20, 2026
c02385e
fix(transform): startup errors say why and where; harden tests
facontidavide Sep 20, 2026
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
32 changes: 32 additions & 0 deletions CHANGELOG.rst
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,38 @@
Changelog for package pj_ros_bridge
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^

Forthcoming
-----------
* Message transforms (ROS2 only): operator-configured, off by default, a
topic is transformed for all subscribers or none — no per-client
negotiation. New ``transform_profile`` parameter points to a JSON file of
ordered per-topic rules (``match_type`` exact and required, optionally
narrowed by a ``match_topic`` full-match regex; first match wins; a transform
paired with a type it cannot handle is a startup error). A transformed topic is advertised
with its transform's output type/schema, plus an optional ``source_type``
on ``get_topics`` entries (not yet on ``topics_changed`` entries); new
``message_transforms`` server capability.
* New ``cloudini`` transform: ``sensor_msgs/msg/PointCloud2`` ->
``point_cloud_interfaces/msg/CompressedPointCloud2`` using `Cloudini
<https://github.com/facontidavide/cloudini>`_ point cloud compression
(lossy at a configurable resolution; its own second compression stage is
disabled since the bridge already ZSTD-compresses every frame). Only
built when ``cloudini_lib`` is available (``find_package``, else fetched
at configure time unless ``-DPJ_BRIDGE_FETCH_CLOUDINI=OFF``); a profile
naming ``cloudini`` on a build without it fails at startup.
* ``strip_large_messages`` is now sugar for a ``strip`` transform rule
appended after the profile's own rules (an explicit profile rule for the
same type wins).
* **Behavior change:** a transform (including ``strip``) that fails on a
message now drops the sample instead of forwarding it — previously
``strip_large_messages`` forwarded the original, unstripped message on
failure.
* Per-topic transform statistics (samples, drops, compression ratio,
microseconds/sample) logged with the final statistics at shutdown.
* The project now requires **C++20** (``std::span``, and ``cloudini_lib``'s
PUBLIC ``cxx_std_20`` requirement).
* 304 unit tests passing.

0.10.0 (2026-09-20)
-------------------
* WebSocket permessage-deflate declined server-side (redundant with our own
Expand Down
15 changes: 12 additions & 3 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ make -j$(nproc)

### Code Quality

- **C++ Standard**: C++20 (required for `std::span`, and because `cloudini_lib` exports `cxx_std_20` as a PUBLIC requirement)
- **Thread Safety**: Use mutexes, document thread safety in class comments
- **Testing**: Unit tests required for all core components (gtest)
- **Code Formatting**: `pre-commit run -a` before committing (clang-format)
Expand Down Expand Up @@ -123,7 +124,7 @@ BridgeServer does NOT own timers. The entry point (`main.cpp`) drives the event

6. **Ros2TopicSource** (`ros2/`) — Wraps `TopicDiscovery` + `SchemaExtractor`. Schema encoding: `"ros2msg"`.

7. **Ros2SubscriptionManager** (`ros2/`) — Wraps `GenericSubscriptionManager` + optional `MessageStripper`. Converts `rclcpp::SerializedMessage` → `shared_ptr<vector<byte>>` via memcpy.
7. **Ros2SubscriptionManager** (`ros2/`) — Wraps `GenericSubscriptionManager`. On `subscribe()`, resolves the topic's `TransformSet` binding (if any) to find the real source type and subscribes with it; in the message callback, a bound transform's `apply()` runs on the rcl buffer before forwarding (replacing the old inline `MessageStripper` call — `strip` is now a transform, see item 13). Converts `rclcpp::SerializedMessage` → `shared_ptr<vector<byte>>` via memcpy (or via the transform's output when one is bound).

8. **RtiTopicSource** (`rti/`) — Wraps `DdsTopicDiscovery`. Schema encoding: `"omgidl"`. (Build disabled)

Expand All @@ -133,6 +134,10 @@ BridgeServer does NOT own timers. The entry point (`main.cpp`) drives the event

11. **FastDdsSubscriptionManager** (`fastdds/`) — Directly implements `SubscriptionManagerInterface`. Creates `DataReader`s with `DynamicPubSubType`, deserializes into `DynamicData` and re-serializes to extract CDR bytes.

12. **TransformSet** (`app/`) — Owns the ordered `transform_profile` rules, the transform name → `TransformFactory` map, and the per-topic `BoundTransform` instances (created lazily on first match, `first match wins`). Thread-safe: written from the request thread (`get_topics`/`subscribe`), read from the ingest thread. Also tracks per-topic transform statistics (samples, drops, compression ratio, µs/sample), logged at shutdown.

13. **TransformingTopicSource** (`app/`) — Decorator over a `TopicSourceInterface`. Rewrites `get_topics()`/`get_schema()` to the transform's output type/schema for matched topics (recording `source_type`) and passes everything else through unchanged. `strip` and `cloudini` are the two built-in transforms (`ros2/src/strip_transform.cpp`, `app/src/cloudini_transform.cpp`); `strip` replaces the old inline `MessageStripper` call in `Ros2SubscriptionManager` (see item 7).

### Communication Pattern

**WebSocket** (single port, default 9090):
Expand Down Expand Up @@ -181,6 +186,9 @@ Then, for each message in the (compressed) payload:
- **nlohmann/json** — JSON (`find_package(REQUIRED)`)
- **tl::expected** — error handling (header-only, vendored in 3rdparty/)

### Core (optional)
- **Cloudini** (`cloudini_lib`) — point cloud compression for the `cloudini` transform. `find_package(cloudini_lib)` first, else fetched at configure time (pinned tag `1.3.1`, static) unless `-DPJ_BRIDGE_FETCH_CLOUDINI=OFF`. A build with neither simply lacks the `cloudini` transform (`PJ_BRIDGE_HAS_CLOUDINI` undefined); `strip` and the rest of the transform machinery are unaffected.

### ROS2 Backend
- `rclcpp`, `ament_index_cpp`, `ament_cmake`
- `sensor_msgs`, `nav_msgs` (for message stripper)
Expand All @@ -197,7 +205,7 @@ Then, for each message in the (compressed) payload:

## Testing

### Test Count: 231 unit tests across 11 test suites
### Test Count: 304 unit tests

### Commands
```bash
Expand Down Expand Up @@ -233,6 +241,7 @@ port: 9090 # WebSocket port
publish_rate: 50.0 # Hz
session_timeout: 10.0 # seconds
strip_large_messages: false # Opt-in: strip Image/PointCloud2/etc data fields
transform_profile: "" # Path to a per-topic message-transform profile JSON (e.g. Cloudini); empty disables
topic_whitelist: [".*"] # Full-match regex patterns restricting visible/subscribable topics
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
Expand Down Expand Up @@ -282,5 +291,5 @@ identity/capabilities object, and TLS setup).

**Last Updated**: 2026-07-06
**Project Phase**: Unified multi-backend architecture
**Test Status**: 231 unit tests passing (all sanitizers clean)
**Test Status**: 304 unit tests passing (all sanitizers clean)
**Executables**: `pj_bridge_ros2` (ROS2), `pj_bridge_rti` (RTI DDS, disabled), `pj_bridge_fastdds` (FastDDS)
42 changes: 41 additions & 1 deletion CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ if(CMAKE_COMPILER_IS_GNUCXX OR CMAKE_CXX_COMPILER_ID MATCHES "Clang")
add_compile_options(-Wall -Wextra -Wpedantic)
endif()

set(CMAKE_CXX_STANDARD 17)
set(CMAKE_CXX_STANDARD 20)
set(CMAKE_CXX_STANDARD_REQUIRED ON)

# Sanitizer options
Expand Down Expand Up @@ -88,6 +88,25 @@ else()
message(STATUS "Using system IXWebSocket")
endif()

# Cloudini (optional) — point cloud compression transform
option(PJ_BRIDGE_FETCH_CLOUDINI "Fetch Cloudini when cloudini_lib is not installed" ON)
# 1.3.1 validates malformed clouds itself and can be embedded without workarounds.
find_package(cloudini_lib 1.3.1 QUIET)
if(NOT cloudini_lib_FOUND AND PJ_BRIDGE_FETCH_CLOUDINI)
message(STATUS "cloudini_lib not found — fetching via FetchContent")
FetchContent_Declare(
cloudini
URL https://github.com/facontidavide/cloudini/archive/refs/tags/1.3.1.tar.gz
URL_HASH SHA256=d8c0e265bb841f59002e996bc957b4f0845b5399deb62b3247dc11115b48b4d9
SOURCE_SUBDIR cloudini_lib
)
# Reuse the zstd found above instead of a second, vendored copy in the same binary.
set(CLOUDINI_FORCE_VENDORED_DEPS OFF CACHE BOOL "" FORCE)
FetchContent_MakeAvailable(cloudini)
set_target_properties(cloudini_lib PROPERTIES POSITION_INDEPENDENT_CODE ON)
set(cloudini_lib_FOUND TRUE)
endif()

# spdlog — available via rosdep (libspdlog-dev) or pixi (conda-forge)
find_package(spdlog REQUIRED)

Expand All @@ -113,6 +132,8 @@ add_library(pj_bridge_app STATIC
app/src/bridge_server.cpp
app/src/standalone_event_loop.cpp
app/src/whitelist_filter.cpp
app/src/transform_set.cpp
app/src/transforming_topic_source.cpp
)

target_include_directories(pj_bridge_app PUBLIC
Expand All @@ -127,6 +148,21 @@ target_link_libraries(pj_bridge_app PUBLIC
nlohmann_json::nlohmann_json
)

if(cloudini_lib_FOUND)
# A standalone install exports cloudini::cloudini_lib, an ament install cloudini_lib::cloudini_lib.
if(TARGET cloudini::cloudini_lib)
set(_pj_cloudini_target cloudini::cloudini_lib)
else()
set(_pj_cloudini_target cloudini_lib::cloudini_lib)
endif()
target_sources(pj_bridge_app PRIVATE app/src/cloudini_transform.cpp)
target_link_libraries(pj_bridge_app PRIVATE ${_pj_cloudini_target})
target_compile_definitions(pj_bridge_app PUBLIC PJ_BRIDGE_HAS_CLOUDINI)
message(STATUS "Cloudini transform: enabled")
else()
message(STATUS "Cloudini transform: disabled (cloudini_lib not found)")
endif()

# ============================================================================
# pj_bridge_ros2 — ROS2 backend (only if ament_cmake found)
# ============================================================================
Expand All @@ -142,6 +178,7 @@ if(ament_cmake_FOUND)
ros2/src/schema_extractor.cpp
ros2/src/generic_subscription_manager.cpp
ros2/src/message_stripper.cpp
ros2/src/strip_transform.cpp
ros2/src/ros2_topic_source.cpp
ros2/src/ros2_subscription_manager.cpp
)
Expand Down Expand Up @@ -266,11 +303,14 @@ if(BUILD_TESTING AND ament_cmake_FOUND)
tests/unit/test_bridge_server.cpp
tests/unit/test_protocol_constants.cpp
tests/unit/test_whitelist_filter.cpp
tests/unit/test_transform_set.cpp
tests/unit/test_transforming_topic_source.cpp
tests/unit/test_topic_discovery.cpp
tests/unit/test_schema_extractor.cpp
tests/unit/test_generic_subscription_manager.cpp
tests/unit/test_ros2_subscription_manager.cpp
tests/unit/test_message_stripper.cpp
tests/unit/test_cloudini_transform.cpp
)

target_link_libraries(${PROJECT_NAME}_tests
Expand Down
4 changes: 4 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ independently.
- **Multi-Client Support**: Multiple clients can connect simultaneously with shared subscriptions
- **Runtime Schema Discovery**: Automatic extraction of message schemas from installed ROS2 packages on the server side.
- **Large Message Stripping** (opt-in): Optional stripping of large array fields (Image, PointCloud2, LaserScan, OccupancyGrid) to reduce bandwidth while preserving metadata. Disabled by default — full message data is forwarded; enable with `strip_large_messages:=true` for low-bandwidth links
- **Message Transforms** (opt-in, ROS2 only): Operator-configured per-topic transforms (e.g. `cloudini` point cloud compression) applied before buffering, replacing a topic's payload and advertised type/schema for all subscribers. Off by default (`transform_profile:=""`); see [docs/API.md](docs/API.md#message-transforms-ros2-only)
- **Topic Whitelist**: Restrict which topics are visible/subscribable via full-match regex patterns (`topic_whitelist` / `--topic-whitelist`), mirroring foxglove_bridge's option of the same name
- **QoS Depth Heuristics** (ROS2 only): KEEP_LAST subscription depth is derived from the discovered publishers' depths and clamped to a configurable `[min_qos_depth, max_qos_depth]` range
- **Pushed Topic Advertisement** (opt-in): Clients can subscribe to a `topics_changed` notification instead of polling `get_topics`, at a configurable `topic_poll_interval`
Expand All @@ -49,6 +50,7 @@ independently.
| `publish_rate` | double | 50.0 | Aggregation publish rate in Hz |
| `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 |
| `transform_profile` | string | `""` | ROS2 only: path to a JSON profile of per-topic message transforms (e.g. Cloudini point cloud compression); empty disables the feature |
| `topic_whitelist` | string array | `[".*"]` | Full-match regex patterns (ECMAScript) restricting visible/subscribable topics |
| `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 |
Expand Down Expand Up @@ -113,6 +115,8 @@ All dependencies (spdlog, nlohmann_json, ZSTD) are provided by the dependency ma

TLS (`wss://`) support depends on IXWebSocket being built with OpenSSL. The CMake option `PJ_BRIDGE_TLS` (default `ON`) controls this for the FetchContent path (`-DPJ_BRIDGE_TLS=OFF` to disable); a system/conda-provided IXWebSocket must likewise have been built with TLS. See [docs/API.md](docs/API.md#tls--wss) for details.

Cloudini (the `cloudini` message transform) is optional: `find_package(cloudini_lib)` is tried first (version >= 1.3.1), otherwise it is fetched at configure time and linked statically. Builds without network access must either have `cloudini_lib` installed or pass `-DPJ_BRIDGE_FETCH_CLOUDINI=OFF`, which builds the bridge without that transform.

### ROS2 — Pixi

[Pixi](https://pixi.sh) manages the full toolchain including ROS2 via [RoboStack](https://robostack.github.io/).
Expand Down
34 changes: 34 additions & 0 deletions app/include/pj_bridge/cloudini_transform.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
/*
* Copyright (C) 2026 Davide Faconti
*
* This file is part of pj_bridge.
*
* pj_bridge is free software: you can redistribute it and/or modify
* it under the terms of the GNU Affero General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* pj_bridge is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU Affero General Public License for more details.
*
* You should have received a copy of the GNU Affero General Public License
* along with pj_bridge. If not, see <https://www.gnu.org/licenses/>.
*/

#pragma once

#include "pj_bridge/message_transform.hpp"

namespace pj_bridge {

/// `cloudini` transform: sensor_msgs PointCloud2 -> point_cloud_interfaces
/// CompressedPointCloud2. Works on CDR bytes, so it serves any backend whose
/// PointCloud2 is CDR-encoded. Only declared when Cloudini is available.
///
/// Params: resolution (float, default 0.001), fields ({name: resolution}, 0
/// removes the field), viz_preprocessing (bool, default false).
TransformFactory make_cloudini_transform_factory();

} // namespace pj_bridge
65 changes: 65 additions & 0 deletions app/include/pj_bridge/message_transform.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
/*
* Copyright (C) 2026 Davide Faconti
*
* This file is part of pj_bridge.
*
* pj_bridge is free software: you can redistribute it and/or modify
* it under the terms of the GNU Affero General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* pj_bridge is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU Affero General Public License for more details.
*
* You should have received a copy of the GNU Affero General Public License
* along with pj_bridge. If not, see <https://www.gnu.org/licenses/>.
*/

#pragma once

#include <cstddef>
#include <functional>
#include <memory>
#include <nlohmann/json.hpp>
#include <span>
#include <string>
#include <vector>

#include "tl/expected.hpp"

namespace pj_bridge {

/// One instance per topic, called from a single thread (the ingest executor).
/// Implementations SHOULD keep codec contexts and scratch buffers as members.
class MessageTransform {
public:
virtual ~MessageTransform() = default;

/// `in` views the middleware's buffer and is valid only during the call.
/// `out` is the final storage later held by MessageBuffer: clear and fill it.
/// On error the caller DROPS the sample. It must never forward `in` instead:
/// the client was told the output type at subscribe time.
virtual tl::expected<void, std::string> apply(std::span<const std::byte> in, std::vector<std::byte>& out) = 0;

/// Type name advertised to clients in place of `source_type`. Default: unchanged.
virtual std::string output_type(const std::string& source_type) const {
return source_type;
}

/// Schema advertised to clients, given the untransformed one. Default: unchanged.
virtual std::string output_schema(const std::string& source_schema) const {
return source_schema;
}
};

struct TransformFactory {
/// Can this transform handle messages of `source_type`?
std::function<bool(const std::string& source_type)> accepts;
/// Reject unknown or malformed params. Called once per rule at startup.
std::function<tl::expected<void, std::string>(const nlohmann::json& params)> check_params;
std::function<std::unique_ptr<MessageTransform>(const std::string& source_type, const nlohmann::json& params)> create;
};

} // namespace pj_bridge
1 change: 1 addition & 0 deletions app/include/pj_bridge/protocol_constants.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ inline constexpr const char* kServerCapabilities[] = {
"topics_changed", // pushed topic advertisement (subscribe_topic_updates)
"per_topic_rate_limit", // subscribe entries accept {name, max_rate_hz}
"size_class_frames", // large topics isolated into own frames (header flag bit0 = heavy)
"message_transforms", // topics may be advertised with a transformed type (+ source_type)
};

} // namespace pj_bridge
5 changes: 3 additions & 2 deletions app/include/pj_bridge/topic_source_interface.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,9 @@ namespace pj_bridge {

/// A discovered topic with its fully-qualified name and message type.
struct TopicInfo {
std::string name; ///< e.g. "/sensor/imu"
std::string type; ///< e.g. "sensor_msgs/msg/Imu"
std::string name; ///< e.g. "/sensor/imu"
std::string type; ///< e.g. "sensor_msgs/msg/Imu"
std::string source_type{}; ///< set only for transformed topics: the type before the transform
};

/// Abstract interface for discovering topics and retrieving their schemas.
Expand Down
Loading
Loading