msg: add message translation node for ROS

This commit is contained in:
Beat Küng
2025-02-11 13:19:25 +01:00
committed by Silvan Fuhrer
parent 975ec30c9c
commit f6bfa9812e
32 changed files with 3200 additions and 0 deletions
@@ -8,6 +8,7 @@ if [ $# -gt 0 ]; then
fi
exec find boards msg src platforms test \
-path msg/translation_node -prune -o \
-path platforms/nuttx/NuttX -prune -o \
-path platforms/qurt/dspal -prune -o \
-path src/drivers/ins/vectornav/libvnc -prune -o \
+34
View File
@@ -0,0 +1,34 @@
#! /bin/bash
# Copy msgs and the message translation node into a ROS workspace directory
DIR=$( cd "$( dirname "${BASH_SOURCE[0]}" )" && pwd )
PX4_SRC_DIR="$DIR/.."
WS_DIR="$1"
if [ ! -e "${WS_DIR}" ]
then
echo "Usage: $0 <ros_ws_dir>"
exit 1
fi
WS_DIR="$WS_DIR"/src
if [ ! -e "${WS_DIR}" ]
then
echo "'src' directory not found inside ROS workspace (${WS_DIR})"
exit 1
fi
cp -ar "${PX4_SRC_DIR}"/msg/translation_node "${WS_DIR}"
cp -ar "${PX4_SRC_DIR}"/msg/px4_msgs_old "${WS_DIR}"
PX4_MSGS_DIR="${WS_DIR}"/px4_msgs
if [ ! -e "${PX4_MSGS_DIR}" ]
then
git clone https://github.com/PX4/px4_msgs.git "${PX4_MSGS_DIR}"
rm -rf "${PX4_MSGS_DIR}"/msg/*.msg
rm -rf "${PX4_MSGS_DIR}"/msg/versioned/*.msg
rm -rf "${PX4_MSGS_DIR}"/srv/*.srv
fi
cp -ar "${PX4_SRC_DIR}"/msg/*.msg "${PX4_MSGS_DIR}"/msg
mkdir -p "${PX4_MSGS_DIR}"/msg/versioned
cp -ar "${PX4_SRC_DIR}"/msg/versioned/*.msg "${PX4_MSGS_DIR}"/msg/versioned
cp -ar "${PX4_SRC_DIR}"/srv/*.srv "${PX4_MSGS_DIR}"/srv
+76
View File
@@ -0,0 +1,76 @@
cmake_minimum_required(VERSION 3.5)
project(px4_msgs_old)
list(INSERT CMAKE_MODULE_PATH 0 "${CMAKE_CURRENT_SOURCE_DIR}/cmake")
if(CMAKE_COMPILER_IS_GNUCXX OR CMAKE_CXX_COMPILER_ID MATCHES "Clang")
add_compile_options(-Wall -Wextra)
endif()
find_package(ament_cmake REQUIRED)
find_package(builtin_interfaces REQUIRED)
find_package(rosidl_default_generators REQUIRED)
# ##############################################################################
# Generate ROS messages, ROS2 interfaces and IDL files #
# ##############################################################################
# get all msg files
set(MSGS_DIR "${CMAKE_CURRENT_SOURCE_DIR}/msg")
file(GLOB PX4_MSGS RELATIVE "${CMAKE_CURRENT_SOURCE_DIR}" "${MSGS_DIR}/*.msg")
# get all srv files
set(SRVS_DIR "${CMAKE_CURRENT_SOURCE_DIR}/srv")
file(GLOB PX4_SRVS RELATIVE "${CMAKE_CURRENT_SOURCE_DIR}" "${SRVS_DIR}/*.srv")
# For the versioned topics, replace the namespace (px4_msgs_old -> px4_msgs) and message type name (<msg>Vx -> <msg>),
# so that DDS does not reject the subscription/publication due to mismatching type
# rosidl_typesupport_fastrtps_cpp
set(rosidl_typesupport_fastrtps_cpp_BIN ${CMAKE_CURRENT_BINARY_DIR}/rosidl_typesupport_fastrtps_cpp_wrapper.py)
file(TOUCH ${rosidl_typesupport_fastrtps_cpp_BIN})
# rosidl_typesupport_fastrtps_c
set(rosidl_typesupport_fastrtps_c_BIN ${CMAKE_CURRENT_BINARY_DIR}/rosidl_typesupport_fastrtps_c_wrapper.py)
file(TOUCH ${rosidl_typesupport_fastrtps_c_BIN})
# rosidl_typesupport_introspection_cpp (for cyclonedds)
set(rosidl_typesupport_introspection_cpp_BIN ${CMAKE_CURRENT_BINARY_DIR}/rosidl_typesupport_introspection_cpp_wrapper.py)
file(TOUCH ${rosidl_typesupport_introspection_cpp_BIN})
# Generate introspection typesupport for C and C++ and IDL files
if(PX4_MSGS)
rosidl_generate_interfaces(${PROJECT_NAME}
${PX4_MSGS}
${PX4_SRVS}
DEPENDENCIES builtin_interfaces
ADD_LINTER_TESTS
)
endif()
# rosidl_typesupport_fastrtps_cpp
set(rosidl_typesupport_fastrtps_cpp_orig ${rosidl_typesupport_fastrtps_cpp_DIR})
string(REPLACE "share/rosidl_typesupport_fastrtps_cpp/cmake" "lib/rosidl_typesupport_fastrtps_cpp/rosidl_typesupport_fastrtps_cpp"
rosidl_typesupport_fastrtps_cpp_orig ${rosidl_typesupport_fastrtps_cpp_DIR})
set(original_script_path ${rosidl_typesupport_fastrtps_cpp_orig})
configure_file(rename_msg_type.py.in ${rosidl_typesupport_fastrtps_cpp_BIN} @ONLY)
# rosidl_typesupport_fastrtps_c
set(rosidl_typesupport_fastrtps_c_orig ${rosidl_typesupport_fastrtps_c_DIR})
string(REPLACE "share/rosidl_typesupport_fastrtps_c/cmake" "lib/rosidl_typesupport_fastrtps_c/rosidl_typesupport_fastrtps_c"
rosidl_typesupport_fastrtps_c_orig ${rosidl_typesupport_fastrtps_c_DIR})
set(original_script_path ${rosidl_typesupport_fastrtps_c_orig})
configure_file(rename_msg_type.py.in ${rosidl_typesupport_fastrtps_c_BIN} @ONLY)
# rosidl_typesupport_introspection_cpp
set(rosidl_typesupport_introspection_cpp_orig ${rosidl_typesupport_introspection_cpp_DIR})
string(REPLACE "share/rosidl_typesupport_introspection_cpp/cmake" "lib/rosidl_typesupport_introspection_cpp/rosidl_typesupport_introspection_cpp"
rosidl_typesupport_introspection_cpp_orig ${rosidl_typesupport_introspection_cpp_DIR})
set(original_script_path ${rosidl_typesupport_introspection_cpp_orig})
configure_file(rename_msg_type.py.in ${rosidl_typesupport_introspection_cpp_BIN} @ONLY)
ament_export_dependencies(rosidl_default_runtime)
ament_package()
+25
View File
@@ -0,0 +1,25 @@
<?xml version="1.0"?>
<?xml-model href="http://download.ros.org/schema/package_format3.xsd" schematypens="http://www.w3.org/2001/XMLSchema"?>
<package format="3">
<name>px4_msgs_old</name>
<version>2.0.1</version>
<description>Package with the ROS-equivalent of PX4 uORB msgs (old message definitions)</description>
<maintainer email="info@px4.io">PX4</maintainer>
<license>BSD 3-Clause</license>
<buildtool_depend>ament_cmake</buildtool_depend>
<buildtool_depend>rosidl_default_generators</buildtool_depend>
<depend>builtin_interfaces</depend>
<depend>ros_environment</depend>
<exec_depend>rosidl_default_runtime</exec_depend>
<test_depend>ament_lint_common</test_depend>
<member_of_group>rosidl_interface_packages</member_of_group>
<export>
<build_type>ament_cmake</build_type>
</export>
</package>
+39
View File
@@ -0,0 +1,39 @@
#! /bin/python
import sys
import subprocess
import json
import os
import re
original_script = "@original_script_path@"
args = sys.argv[1:]
json_file = [arg for arg in args if arg.endswith('.json')][0]
proc = subprocess.run(['python3', original_script] + args)
proc.check_returncode()
def replace_namespace_and_type(content: str):
# Replace namespace type
content = content.replace('"px4_msgs_old"', '"px4_msgs"')
content = content.replace('"px4_msgs_old::msg"', '"px4_msgs::msg"')
# Replace versioned type with non-versioned one
content = re.sub(r'("[a-zA-Z0-9]+)V[0-9]+"', '\\1"', content)
# Services
content = content.replace('"px4_msgs_old::srv"', '"px4_msgs::srv"')
content = re.sub(r'("[a-zA-Z0-9]+)V[0-9]+_Request"', '\\1_Request"', content)
content = re.sub(r'("[a-zA-Z0-9]+)V[0-9]+_Response"', '\\1_Response"', content)
return content
with open(json_file, 'r') as f:
data = json.load(f)
output_dir = data['output_dir']
# Iterate files recursively
for root, dirs, files in os.walk(output_dir):
for file in files:
with open(os.path.join(root, file), 'r+') as f:
content = f.read()
f.seek(0)
f.write(replace_namespace_and_type(content))
f.truncate()
+82
View File
@@ -0,0 +1,82 @@
cmake_minimum_required(VERSION 3.8)
project(translation_node)
if(CMAKE_COMPILER_IS_GNUCXX OR CMAKE_CXX_COMPILER_ID MATCHES "Clang")
add_compile_options(-Wall -Wextra -Wpedantic -Wno-unused-parameter -Werror)
endif()
# find dependencies
find_package(ament_cmake REQUIRED)
find_package(rclcpp REQUIRED)
find_package(px4_msgs REQUIRED)
find_package(px4_msgs_old REQUIRED)
if(DEFINED ENV{ROS_DISTRO})
set(ROS_DISTRO $ENV{ROS_DISTRO})
else()
set(ROS_DISTRO "rolling")
endif()
add_library(${PROJECT_NAME}_lib
src/monitor.cpp
src/pub_sub_graph.cpp
src/service_graph.cpp
src/translations.cpp
)
ament_target_dependencies(${PROJECT_NAME}_lib rclcpp px4_msgs px4_msgs_old)
add_executable(${PROJECT_NAME}_bin
src/main.cpp
)
target_link_libraries(${PROJECT_NAME}_bin ${PROJECT_NAME}_lib)
target_include_directories(${PROJECT_NAME}_bin PUBLIC src)
ament_target_dependencies(${PROJECT_NAME}_bin rclcpp px4_msgs px4_msgs_old)
install(TARGETS
${PROJECT_NAME}_bin
DESTINATION lib/${PROJECT_NAME})
option(DISABLE_SERVICES "Disable services" OFF)
if(${ROS_DISTRO} STREQUAL "humble")
message(WARNING "Disabling services for ROS humble (API is not supported)")
target_compile_definitions(${PROJECT_NAME}_lib PRIVATE DISABLE_SERVICES)
set(DISABLE_SERVICES ON)
endif()
if(BUILD_TESTING)
find_package(std_msgs REQUIRED)
find_package(ament_lint_auto REQUIRED)
find_package(ament_cmake_gtest REQUIRED)
find_package(rosidl_default_generators REQUIRED)
ament_lint_auto_find_test_dependencies()
set(SRV_FILES
test/srv/TestV0.srv
test/srv/TestV1.srv
test/srv/TestV2.srv
)
rosidl_generate_interfaces(${PROJECT_NAME} ${SRV_FILES})
# Unit tests
set(TEST_SRC
test/graph.cpp
test/main.cpp
test/pub_sub.cpp
)
if (NOT DISABLE_SERVICES)
list(APPEND TEST_SRC test/services.cpp)
endif()
ament_add_gtest(${PROJECT_NAME}_unit_tests
${TEST_SRC}
)
target_include_directories(${PROJECT_NAME}_unit_tests PRIVATE ${CMAKE_CURRENT_LIST_DIR})
target_compile_options(${PROJECT_NAME}_unit_tests PRIVATE -Wno-error=sign-compare) # There is a warning from gtest internal
target_link_libraries(${PROJECT_NAME}_unit_tests ${PROJECT_NAME}_lib)
rosidl_get_typesupport_target(cpp_typesupport_target ${PROJECT_NAME} "rosidl_typesupport_cpp")
target_link_libraries(${PROJECT_NAME}_unit_tests "${cpp_typesupport_target}")
ament_target_dependencies(${PROJECT_NAME}_unit_tests
std_msgs
rclcpp
)
endif()
ament_package()
+6
View File
@@ -0,0 +1,6 @@
# Message Translations
This package contains a message translation node and a set of old message conversion methods.
It allows to run applications that are compiled with one set of message versions against a PX4 with another set of message versions, without having to change either the application or the PX4 side.
For details, see https://docs.px4.io/main/en/ros2/px4_ros2_msg_translation_node.html.
+27
View File
@@ -0,0 +1,27 @@
<?xml version="1.0"?>
<?xml-model href="http://download.ros.org/schema/package_format3.xsd" schematypens="http://www.w3.org/2001/XMLSchema"?>
<package format="3">
<name>translation_node</name>
<version>0.0.0</version>
<description>Message version translation node</description>
<maintainer email="info@px4.io">PX4</maintainer>
<license>BSD 3-Clause</license>
<buildtool_depend>ament_cmake</buildtool_depend>
<buildtool_depend>rosidl_default_generators</buildtool_depend>
<member_of_group>rosidl_interface_packages</member_of_group>
<test_depend>ament_lint_auto</test_depend>
<test_depend>ament_lint_common</test_depend>
<test_depend>ament_cmake_gtest</test_depend>
<test_depend>std_msgs</test_depend>
<depend>rclcpp</depend>
<depend>px4_msgs</depend>
<depend>px4_msgs_old</depend>
<export>
<build_type>ament_cmake</build_type>
</export>
</package>
File diff suppressed because it is too large Load Diff
+39
View File
@@ -0,0 +1,39 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#include <memory>
#include <rclcpp/rclcpp.hpp>
#include "../translations/all_translations.h"
#include "pub_sub_graph.h"
#include "service_graph.h"
#include "monitor.h"
using namespace std::chrono_literals;
class RosTranslationNode : public rclcpp::Node
{
public:
RosTranslationNode() : Node("translation_node")
{
_pub_sub_graph = std::make_unique<PubSubGraph>(*this, RegisteredTranslations::instance().topicTranslations());
_service_graph = std::make_unique<ServiceGraph>(*this, RegisteredTranslations::instance().serviceTranslations());
_monitor = std::make_unique<Monitor>(*this, _pub_sub_graph.get(), _service_graph.get());
}
private:
std::unique_ptr<PubSubGraph> _pub_sub_graph;
std::unique_ptr<ServiceGraph> _service_graph;
rclcpp::TimerBase::SharedPtr _node_update_timer;
std::unique_ptr<Monitor> _monitor;
};
int main(int argc, char * argv[])
{
rclcpp::init(argc, argv);
rclcpp::spin(std::make_shared<RosTranslationNode>());
rclcpp::shutdown();
return 0;
}
+60
View File
@@ -0,0 +1,60 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#include "monitor.h"
using namespace std::chrono_literals;
Monitor::Monitor(rclcpp::Node &node, PubSubGraph* pub_sub_graph, ServiceGraph* service_graph)
: _node(node), _pub_sub_graph(pub_sub_graph), _service_graph(service_graph) {
// Monitor subscriptions & publishers
// TODO: event-based
_node_update_timer = _node.create_wall_timer(1s, [this]() {
updateNow();
});
}
void Monitor::updateNow() {
// Topics
if (_pub_sub_graph != nullptr) {
std::vector<PubSubGraph::TopicInfo> topic_info;
const auto topics = _node.get_topic_names_and_types();
for (const auto &[topic_name, topic_types]: topics) {
auto publishers = _node.get_publishers_info_by_topic(topic_name);
auto subscribers = _node.get_subscriptions_info_by_topic(topic_name);
// Filter out self
int num_publishers = 0;
for (const auto &publisher: publishers) {
num_publishers += publisher.node_name() != _node.get_name();
}
int num_subscribers = 0;
for (const auto &subscriber: subscribers) {
num_subscribers += subscriber.node_name() != _node.get_name();
}
if (num_subscribers > 0 || num_publishers > 0) {
topic_info.emplace_back(PubSubGraph::TopicInfo{topic_name, num_subscribers, num_publishers});
}
}
_pub_sub_graph->updateCurrentTopics(topic_info);
}
// Services
#ifndef DISABLE_SERVICES // ROS Humble does not support the count_services() call
if (_service_graph != nullptr) {
std::vector<ServiceGraph::ServiceInfo> service_info;
const auto services = _node.get_service_names_and_types();
for (const auto& [service_name, service_types] : services) {
const int num_services = _node.get_node_graph_interface()->count_services(service_name);
const int num_clients = _node.get_node_graph_interface()->count_clients(service_name);
// We cannot filter out our own node, as we don't have that info.
// We could use `get_service_names_and_types_by_node`, but then we would not get
// services by non-ros nodes (e.g. microxrce dds bridge)
service_info.emplace_back(ServiceGraph::ServiceInfo{service_name, num_services, num_clients});
}
_service_graph->updateCurrentServices(service_info);
}
#endif
}
+23
View File
@@ -0,0 +1,23 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
#include <rclcpp/rclcpp.hpp>
#include "pub_sub_graph.h"
#include "service_graph.h"
#include <functional>
class Monitor {
public:
explicit Monitor(rclcpp::Node &node, PubSubGraph* pub_sub_graph, ServiceGraph* service_graph);
void updateNow();
private:
rclcpp::Node &_node;
PubSubGraph* _pub_sub_graph{nullptr};
ServiceGraph* _service_graph{nullptr};
rclcpp::TimerBase::SharedPtr _node_update_timer;
};
+195
View File
@@ -0,0 +1,195 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#include "pub_sub_graph.h"
#include "util.h"
PubSubGraph::PubSubGraph(rclcpp::Node &node, const TopicTranslations &translations) : _node(node) {
std::unordered_map<std::string, std::set<MessageVersionType>> known_versions;
for (const auto& topic : translations.topics()) {
const std::string full_topic_name = getFullTopicName(_node.get_effective_namespace(), topic.id.topic_name);
_known_topics_warned.insert({full_topic_name, false});
const MessageIdentifier id{full_topic_name, topic.id.version};
NodeDataPubSub node_data{topic.subscription_factory, topic.publication_factory, id, topic.max_serialized_message_size};
_pub_sub_graph.addNodeIfNotExists(id, std::move(node_data), topic.message_buffer);
known_versions[full_topic_name].insert(id.version);
}
auto get_full_topic_names = [this](std::vector<MessageIdentifier> ids) {
for (auto& id : ids) {
id.topic_name = getFullTopicName(_node.get_effective_namespace(), id.topic_name);
}
return ids;
};
for (const auto& translation : translations.translations()) {
const std::vector<MessageIdentifier> inputs = get_full_topic_names(translation.inputs);
const std::vector<MessageIdentifier> outputs = get_full_topic_names(translation.outputs);
_pub_sub_graph.addTranslation(translation.cb, inputs, outputs);
}
printTopicInfo(known_versions);
handleLargestTopic(known_versions);
}
void PubSubGraph::updateCurrentTopics(const std::vector<TopicInfo> &topics) {
_pub_sub_graph.iterateNodes([](const MessageIdentifier& type, const Graph<NodeDataPubSub>::MessageNodePtr& node) {
node->data().has_external_publisher = false;
node->data().has_external_subscriber = false;
node->data().visited = false;
});
for (const auto& info : topics) {
const auto [non_versioned_topic_name, version] = getNonVersionedTopicName(info.topic_name);
auto maybe_node = _pub_sub_graph.findNode({non_versioned_topic_name, version});
if (!maybe_node) {
auto known_topic_iter = _known_topics_warned.find(non_versioned_topic_name);
if (known_topic_iter != _known_topics_warned.end() && !known_topic_iter->second) {
RCLCPP_WARN(_node.get_logger(), "No translation available for version %i of topic %s", version, non_versioned_topic_name.c_str());
known_topic_iter->second = true;
}
continue;
}
const auto& node = maybe_node.value();
if (info.num_publishers > 0) {
node->data().has_external_publisher = true;
}
if (info.num_subscribers > 0) {
node->data().has_external_subscriber = true;
}
}
// Iterate connected graph segments
_pub_sub_graph.iterateNodes([this](const MessageIdentifier& type, const Graph<NodeDataPubSub>::MessageNodePtr& node) {
if (node->data().visited) {
return;
}
node->data().visited = true;
// Count the number of external subscribers and publishers for each connected graph
int num_publishers = 0;
int num_subscribers = 0;
int num_subscribers_without_publisher = 0;
_pub_sub_graph.iterateBFS(node, [&](const Graph<NodeDataPubSub>::MessageNodePtr& node) {
if (node->data().has_external_publisher) {
++num_publishers;
}
if (node->data().has_external_subscriber) {
++num_subscribers;
if (!node->data().has_external_publisher) {
++num_subscribers_without_publisher;
}
}
});
// We need to instantiate publishers and subscribers if:
// - there are multiple publishers and at least 1 subscriber
// - there is 1 publisher and at least 1 subscriber on another node
// Note that in case of splitting or merging topics, this might create more entities than actually needed
const bool require_translation = (num_publishers >= 2 && num_subscribers >= 1)
|| (num_publishers == 1 && num_subscribers_without_publisher >= 1);
if (require_translation) {
_pub_sub_graph.iterateBFS(node, [&](const Graph<NodeDataPubSub>::MessageNodePtr& node) {
node->data().visited = true;
// Has subscriber(s)?
if (node->data().has_external_subscriber && !node->data().publication) {
RCLCPP_INFO(_node.get_logger(), "Found subscriber for topic '%s', version: %i, adding publisher", node->data().topic_name.c_str(), node->data().version);
node->data().publication = node->data().publication_factory(_node);
} else if (!node->data().has_external_subscriber && node->data().publication) {
RCLCPP_INFO(_node.get_logger(), "No subscribers for topic '%s', version: %i, removing publisher", node->data().topic_name.c_str(), node->data().version);
node->data().publication.reset();
}
// Has publisher(s)?
if (node->data().has_external_publisher && !node->data().subscription) {
RCLCPP_INFO(_node.get_logger(), "Found publisher for topic '%s', version: %i, adding subscriber", node->data().topic_name.c_str(), node->data().version);
node->data().subscription = node->data().subscription_factory(_node, [this, node_cpy=node]() {
onSubscriptionUpdate(node_cpy);
});
} else if (!node->data().has_external_publisher && node->data().subscription) {
RCLCPP_INFO(_node.get_logger(), "No publishers for topic '%s', version: %i, removing subscriber", node->data().topic_name.c_str(), node->data().version);
node->data().subscription.reset();
}
});
} else {
// Reset any publishers or subscribers
_pub_sub_graph.iterateBFS(node, [&](const Graph<NodeDataPubSub>::MessageNodePtr& node) {
node->data().visited = true;
if (node->data().publication) {
RCLCPP_INFO(_node.get_logger(), "Removing publisher for topic '%s', version: %i",
node->data().topic_name.c_str(), node->data().version);
node->data().publication.reset();
}
if (node->data().subscription) {
RCLCPP_INFO(_node.get_logger(), "Removing subscriber for topic '%s', version: %i",
node->data().topic_name.c_str(), node->data().version);
node->data().subscription.reset();
}
});
}
});
}
void PubSubGraph::onSubscriptionUpdate(const Graph<NodeDataPubSub>::MessageNodePtr& node) {
_pub_sub_graph.translate(
node,
[this](const Graph<NodeDataPubSub>::MessageNodePtr& node) {
if (node->data().publication != nullptr) {
const auto ret = rcl_publish(node->data().publication->get_publisher_handle().get(),
node->buffer().get(), nullptr);
if (ret != RCL_RET_OK) {
RCLCPP_WARN_ONCE(_node.get_logger(), "Failed to publish on topic '%s', version: %i",
node->data().topic_name.c_str(), node->data().version);
}
}
});
}
void PubSubGraph::printTopicInfo(const std::unordered_map<std::string, std::set<MessageVersionType>>& known_versions) const {
// Print info about known versions
RCLCPP_INFO(_node.get_logger(), "Registered pub/sub topics and versions:");
for (const auto& [topic_name, version_set] : known_versions) {
if (version_set.empty()) {
continue;
}
const std::string versions = std::accumulate(std::next(version_set.begin()), version_set.end(),
std::to_string(*version_set.begin()), // start with first element
[](std::string a, auto&& b) {
return std::move(a) + ", " + std::to_string(b);
});
RCLCPP_INFO(_node.get_logger(), "- %s: %s", topic_name.c_str(), versions.c_str());
}
}
void PubSubGraph::handleLargestTopic(const std::unordered_map<std::string, std::set<MessageVersionType>> &known_versions) {
// FastDDS caches some type information per DDS participant when first creating a publisher or subscriber for a given
// type. The information that is relevant for us is the maximum serialized message size.
// Since different versions can have different sizes, we need to ensure the first publication or subscription
// happens with the version of the largest size. Otherwise, an out-of-memory exception can be triggered.
// And the type must continue to be in use (so we cannot delete it)
for (const auto& [topic_name, versions] : known_versions) {
size_t max_serialized_message_size = 0;
const PublicationFactoryCB* publication_factory_for_max = nullptr;
for (auto version : versions) {
const auto& node = _pub_sub_graph.findNode(MessageIdentifier{topic_name, version});
assert(node);
const auto& node_data = node.value()->data();
if (node_data.max_serialized_message_size > max_serialized_message_size) {
max_serialized_message_size = node_data.max_serialized_message_size;
publication_factory_for_max = &node_data.publication_factory;
}
}
if (publication_factory_for_max) {
_largest_topic_publications.emplace_back((*publication_factory_for_max)(_node));
}
}
}
+58
View File
@@ -0,0 +1,58 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
#include <rclcpp/rclcpp.hpp>
#include <utility>
#include "translations.h"
#include "translation_util.h"
#include "graph.h"
class PubSubGraph {
public:
struct TopicInfo {
std::string topic_name; ///< fully qualified topic name (with namespace)
int num_subscribers; ///< does not include this node's subscribers
int num_publishers; ///< does not include this node's publishers
};
PubSubGraph(rclcpp::Node& node, const TopicTranslations& translations);
void updateCurrentTopics(const std::vector<TopicInfo>& topics);
private:
struct NodeDataPubSub {
explicit NodeDataPubSub(SubscriptionFactoryCB subscription_factory, PublicationFactoryCB publication_factory,
const MessageIdentifier& id, size_t max_serialized_message_size)
: subscription_factory(std::move(subscription_factory)), publication_factory(std::move(publication_factory)),
topic_name(id.topic_name), version(id.version), max_serialized_message_size(max_serialized_message_size)
{ }
const SubscriptionFactoryCB subscription_factory;
const PublicationFactoryCB publication_factory;
const std::string topic_name;
const MessageVersionType version;
const size_t max_serialized_message_size;
// Keep track if there's currently a publisher/subscriber
bool has_external_publisher{false};
bool has_external_subscriber{false};
rclcpp::SubscriptionBase::SharedPtr subscription;
rclcpp::PublisherBase::SharedPtr publication;
bool visited{false};
};
void onSubscriptionUpdate(const Graph<NodeDataPubSub>::MessageNodePtr& node);
void printTopicInfo(const std::unordered_map<std::string, std::set<MessageVersionType>>& known_versions) const;
void handleLargestTopic(const std::unordered_map<std::string, std::set<MessageVersionType>>& known_versions);
rclcpp::Node& _node;
Graph<NodeDataPubSub> _pub_sub_graph;
std::unordered_map<std::string, bool> _known_topics_warned;
std::vector<rclcpp::PublisherBase::SharedPtr> _largest_topic_publications;
};
File diff suppressed because it is too large Load Diff
+76
View File
@@ -0,0 +1,76 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
#include <rclcpp/rclcpp.hpp>
#include <utility>
#include "translations.h"
#include "translation_util.h"
#include "graph.h"
class ServiceGraph {
public:
struct ServiceInfo {
std::string service_name; ///< fully qualified service name (with namespace)
int num_services; ///< This can include a service created by the translation node
int num_clients; ///< This can include a client created by the translation node
};
ServiceGraph(rclcpp::Node &node, const ServiceTranslations& translations);
void updateCurrentServices(const std::vector<ServiceInfo>& services);
private:
struct NodeDataService;
using GraphForService = Graph<std::shared_ptr<NodeDataService>>;
void printServiceInfo(const std::unordered_map<std::string, std::set<MessageVersionType>> &known_versions) const;
void handleLargestTopic(const std::unordered_map<std::string, std::set<MessageVersionType>>& known_versions);
void onNewRequest(std::shared_ptr<rmw_request_id_t> req_id, GraphForService::MessageNodePtr node);
void onResponse(rmw_request_id_t& req_id, GraphForService::MessageNodePtr node);
void cleanupStaleRequests();
struct Request {
std::shared_ptr<rmw_request_id_t> original_request_id;
std::shared_ptr<NodeDataService> original_node_data{nullptr};
rclcpp::Time timestamp_received;
};
struct NodeDataService {
explicit NodeDataService(const Service& service, const MessageIdentifier& id)
: service_factory(service.service_factory), client_factory(service.client_factory),
service_name(id.topic_name), version(id.version),
publication_factory{service.publication_factory_request, service.publication_factory_response},
max_serialized_message_size{service.max_serialized_message_size_request, service.max_serialized_message_size_response}
{ }
const ServiceFactoryCB service_factory;
const ClientFactoryCB client_factory;
const std::string service_name;
const MessageVersionType version;
const std::array<NamedPublicationFactoryCB, 2> publication_factory; // Request/Response
const std::array<size_t, 2> max_serialized_message_size;
// Keep track if there's currently a client/service
bool has_service{false};
bool has_client{false};
rclcpp::ClientBase::SharedPtr client;
ClientSendCB client_send_cb;
rclcpp::ServiceBase::SharedPtr service;
std::unordered_map<int64_t, Request> ongoing_requests; ///< Ongoing service calls for this node
bool visited{false};
};
rclcpp::Node& _node;
GraphForService _request_graph;
GraphForService _response_graph;
std::unordered_map<std::string, bool> _known_services_warned;
rclcpp::TimerBase::SharedPtr _cleanup_timer;
std::vector<rclcpp::PublisherBase::SharedPtr> _largest_topic_publications;
};
+64
View File
@@ -0,0 +1,64 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
#include <memory>
#include <utility>
#include <type_traits>
#include <vector>
/**
* Helper struct to store template parameter packs
*/
template <typename... Args>
struct Pack {
};
/**
* Struct for a template parameter pack with access to the individual types
*/
template<typename ...Types>
struct TypesArray {
template<typename T, typename...OtherTypes>
struct TypeHelper {
using Type = T;
using Next = TypeHelper<OtherTypes..., void>;
};
using Type1 = typename TypeHelper<Types...>::Type;
using Type2 = typename TypeHelper<Types...>::Next::Type;
using Type3 = typename TypeHelper<Types...>::Next::Next::Type;
using Type4 = typename TypeHelper<Types...>::Next::Next::Next::Type;
using Type5 = typename TypeHelper<Types...>::Next::Next::Next::Next::Type;
using Type6 = typename TypeHelper<Types...>::Next::Next::Next::Next::Next::Type;
using args = Pack<Types...>;
};
/**
* Helper for call_translation_function()
*/
template<typename F, typename MessageType, typename... ArgsIn, typename... ArgsOut, size_t... Is, size_t... Os>
inline void call_translation_function_impl(F f, Pack<ArgsIn...>, Pack<ArgsOut...>,
const std::vector<std::shared_ptr<MessageType>>& messages_in,
std::vector<std::shared_ptr<MessageType>>& messages_out,
std::integer_sequence<size_t, Is...>, std::integer_sequence<size_t, Os...>)
{
f(*static_cast<const ArgsIn*>(messages_in[Is].get())..., *static_cast<ArgsOut*>(messages_out[Os].get())...);
}
/**
* Call a translation function F which takes the arguments (const ArgsIn&..., ArgsOut&...),
* by passing messages_in and messages_out as arguments.
* Note that sizeof(ArgsIn) == messages_in.length() && sizeof(ArgsOut) == messages_out.length() must hold.
*/
template<typename F, typename MessageType, typename... ArgsIn, typename... ArgsOut>
inline void call_translation_function(F f, Pack<ArgsIn...> pack_in, Pack<ArgsOut...> pack_out,
const std::vector<std::shared_ptr<MessageType>>& messages_in,
std::vector<std::shared_ptr<MessageType>>& messages_out) {
call_translation_function_impl(f, pack_in, pack_out, messages_in, messages_out,
std::index_sequence_for<ArgsIn...>{}, std::index_sequence_for<ArgsOut...>{});
}
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,5 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#include "translations.h"
+91
View File
@@ -0,0 +1,91 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
#include <string>
#include <cstdint>
#include <unordered_map>
#include <utility>
#include <functional>
#include <vector>
#include <memory>
#include <tuple>
#include "util.h"
#include "graph.h"
#include <rclcpp/rclcpp.hpp>
using TranslationCB = std::function<void(const std::vector<MessageBuffer>&, std::vector<MessageBuffer>&)>;
using SubscriptionFactoryCB = std::function<rclcpp::SubscriptionBase::SharedPtr(rclcpp::Node&, const std::function<void()>& on_topic_cb)>;
using PublicationFactoryCB = std::function<rclcpp::PublisherBase::SharedPtr(rclcpp::Node&)>;
using NamedPublicationFactoryCB = std::function<rclcpp::PublisherBase::SharedPtr(rclcpp::Node&, const std::string&)>;
using ServiceFactoryCB = std::function<rclcpp::ServiceBase::SharedPtr(rclcpp::Node&, const std::function<void(std::shared_ptr<rmw_request_id_t> req_id)>& on_request_cb)>;
using ClientSendCB = std::function<int64_t(MessageBuffer)>;
using ClientFactoryCB = std::function<std::tuple<rclcpp::ClientBase::SharedPtr, ClientSendCB>(rclcpp::Node&, const std::function<void(rmw_request_id_t&)>& on_response_cb)>;
struct Topic {
MessageIdentifier id;
SubscriptionFactoryCB subscription_factory;
PublicationFactoryCB publication_factory;
std::shared_ptr<void> message_buffer;
size_t max_serialized_message_size{};
};
struct Service {
MessageIdentifier id;
ServiceFactoryCB service_factory;
ClientFactoryCB client_factory;
NamedPublicationFactoryCB publication_factory_request;
NamedPublicationFactoryCB publication_factory_response;
std::shared_ptr<void> message_buffer_request;
size_t max_serialized_message_size_request{};
std::shared_ptr<void> message_buffer_response;
size_t max_serialized_message_size_response{};
};
struct Translation {
TranslationCB cb;
std::vector<MessageIdentifier> inputs;
std::vector<MessageIdentifier> outputs;
};
class TopicTranslations {
public:
TopicTranslations() = default;
void addTopic(Topic topic) { _topics.push_back(std::move(topic)); }
void addTranslation(Translation translation) { _translations.push_back(std::move(translation)); }
const std::vector<Topic>& topics() const { return _topics; }
const std::vector<Translation>& translations() const { return _translations; }
private:
std::vector<Topic> _topics;
std::vector<Translation> _translations;
};
class ServiceTranslations {
public:
ServiceTranslations() = default;
void addNode(Service node) { _nodes.push_back(std::move(node)); }
void addRequestTranslation(Translation translation) { _request_translations.push_back(std::move(translation)); }
void addResponseTranslation(Translation translation) { _response_translations.push_back(std::move(translation)); }
const std::vector<Service>& nodes() const { return _nodes; }
const std::vector<Translation>& requestTranslations() const { return _request_translations; }
const std::vector<Translation>& responseTranslations() const { return _response_translations; }
private:
std::vector<Service> _nodes;
std::vector<Translation> _request_translations;
std::vector<Translation> _response_translations;
};
+51
View File
@@ -0,0 +1,51 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
#include <string>
#include <cstdint>
using MessageVersionType = uint32_t;
static inline std::string getVersionedTopicName(const std::string& topic_name, MessageVersionType version) {
// version == 0 can be used to transition from non-versioned topics to versioned ones
if (version == 0) {
return topic_name;
}
return topic_name + "_v" + std::to_string(version);
}
static inline std::pair<std::string, MessageVersionType> getNonVersionedTopicName(const std::string& topic_name) {
// topic name has the form <name>_v<version>, or just <name> (with version=0)
auto pos = topic_name.find_last_of("_v");
// Ensure there's at least one more char after the found string
if (pos == std::string::npos || pos + 2 > topic_name.length()) {
return std::make_pair(topic_name, 0);
}
std::string non_versioned_topic_name = topic_name.substr(0, pos - 1);
std::string version = topic_name.substr(pos + 1);
// Ensure only digits are in the version string
for (char c : version) {
if (!std::isdigit(c)) {
return std::make_pair(topic_name, 0);
}
}
return std::make_pair(non_versioned_topic_name, std::stol(version));
}
/**
* Get the full topic name, including namespace from a topic name.
* namespace_name should be set to Node::get_effective_namespace()
*/
static inline std::string getFullTopicName(const std::string& namespace_name, const std::string& topic_name) {
std::string full_topic_name = topic_name;
if (!full_topic_name.empty() && full_topic_name[0] != '/') {
if (namespace_name.empty() || namespace_name.back() != '/') {
full_topic_name = '/' + full_topic_name;
}
full_topic_name = namespace_name + full_topic_name;
}
return full_topic_name;
}
File diff suppressed because it is too large Load Diff
+16
View File
@@ -0,0 +1,16 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#include <gtest/gtest.h>
#include <rclcpp/rclcpp.hpp>
int main(int argc, char ** argv)
{
rclcpp::init(argc, argv);
testing::InitGoogleTest(&argc, argv);
const int ret = RUN_ALL_TESTS();
rclcpp::shutdown();
return ret;
}
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
+4
View File
@@ -0,0 +1,4 @@
uint32 MESSAGE_VERSION = 0
uint8 request_a
---
uint64 response_a
+4
View File
@@ -0,0 +1,4 @@
uint32 MESSAGE_VERSION = 1
uint64 request_a
---
uint8 response_a
+6
View File
@@ -0,0 +1,6 @@
uint32 MESSAGE_VERSION = 2
uint8 request_a
uint64 request_b
---
uint16 response_a
uint64 response_b
@@ -0,0 +1,11 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
#include <translation_util.h>
//#include "example_translation_direct_v1.h"
//#include "example_translation_multi_v2.h"
//#include "example_translation_service_v1.h"
@@ -0,0 +1,30 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
// Translate ExampleTopic v0 <--> v1
#include <px4_msgs_old/msg/example_topic_v0.hpp>
#include <px4_msgs/msg/example_topic.hpp>
class ExampleTopicV1Translation {
public:
using MessageOlder = px4_msgs_old::msg::ExampleTopicV0;
static_assert(MessageOlder::MESSAGE_VERSION == 0);
using MessageNewer = px4_msgs::msg::ExampleTopic;
static_assert(MessageNewer::MESSAGE_VERSION == 1);
static constexpr const char* kTopic = "fmu/out/example_topic";
static void fromOlder(const MessageOlder &msg_older, MessageNewer &msg_newer) {
// Set msg_newer from msg_older
}
static void toOlder(const MessageNewer &msg_newer, MessageOlder &msg_older) {
// Set msg_older from msg_newer
}
};
REGISTER_TOPIC_TRANSLATION_DIRECT(ExampleTopicV1Translation);
@@ -0,0 +1,42 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
// Translate ExampleTopic and OtherTopic v1 <--> v2
#include <px4_msgs_old/msg/example_topic_v1.hpp>
#include <px4_msgs_old/msg/other_topic_v1.hpp>
#include <px4_msgs/msg/example_topic.hpp>
#include <px4_msgs/msg/other_topic.hpp>
class ExampleTopicOtherTopicV2Translation {
public:
using MessagesOlder = TypesArray<px4_msgs_old::msg::ExampleTopicV1, px4_msgs_old::msg::OtherTopicV1>;
static constexpr const char* kTopicsOlder[] = {
"fmu/out/example_topic",
"fmu/out/other_topic",
};
static_assert(px4_msgs_old::msg::ExampleTopicV1::MESSAGE_VERSION == 1);
static_assert(px4_msgs_old::msg::OtherTopicV1::MESSAGE_VERSION == 1);
using MessagesNewer = TypesArray<px4_msgs::msg::ExampleTopic, px4_msgs::msg::OtherTopic>;
static constexpr const char* kTopicsNewer[] = {
"fmu/out/example_topic",
"fmu/out/other_topic",
};
static_assert(px4_msgs::msg::ExampleTopic::MESSAGE_VERSION == 2);
static_assert(px4_msgs::msg::OtherTopic::MESSAGE_VERSION == 2);
static void fromOlder(const MessagesOlder::Type1 &msg_older1, const MessagesOlder::Type2 &msg_older2,
MessagesNewer::Type1 &msg_newer1, MessagesNewer::Type2 &msg_newer2) {
// Set msg_newer1, msg_newer2 from msg_older1, msg_older2
}
static void toOlder(const MessagesNewer::Type1 &msg_newer1, const MessagesNewer::Type2 &msg_newer2,
MessagesOlder::Type1 &msg_older1, MessagesOlder::Type2 &msg_older2) {
// Set msg_older1, msg_older2 from msg_newer1, msg_newer2
}
};
REGISTER_TOPIC_TRANSLATION(ExampleTopicOtherTopicV2Translation);
@@ -0,0 +1,38 @@
/****************************************************************************
* Copyright (c) 2024 PX4 Development Team.
* SPDX-License-Identifier: BSD-3-Clause
****************************************************************************/
#pragma once
// Translate ExampleService v0 <--> v1
#include <px4_msgs_old/srv/example_service_v0.hpp>
#include <px4_msgs/srv/example_service.hpp>
class ExampleServiceV1Translation {
public:
using MessageOlder = px4_msgs_old::srv::ExampleServiceV0;
static_assert(MessageOlder::Request::MESSAGE_VERSION == 0);
using MessageNewer = px4_msgs::srv::ExampleService;
static_assert(MessageNewer::Request::MESSAGE_VERSION == 1);
static constexpr const char* kTopic = "fmu/example_service";
static void fromOlder(const MessageOlder::Request &msg_older, MessageNewer::Request &msg_newer) {
// Request: set msg_newer from msg_older
}
static void toOlder(const MessageNewer::Request &msg_newer, MessageOlder::Request &msg_older) {
// Request: set msg_older from msg_newer
}
static void fromOlder(const MessageOlder::Response &msg_older, MessageNewer::Response &msg_newer) {
// Response: set msg_newer from msg_older
}
static void toOlder(const MessageNewer::Response &msg_newer, MessageOlder::Response &msg_older) {
// Response: set msg_older from msg_newer
}
};
REGISTER_SERVICE_TRANSLATION_DIRECT(ExampleServiceV1Translation);