diff --git a/CMakeLists.txt b/CMakeLists.txt index ac0fa25..ce9192b 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -20,6 +20,7 @@ find_package(pluginlib REQUIRED) find_package(Boost REQUIRED COMPONENTS system) find_package(PkgConfig REQUIRED) pkg_check_modules(ZSTD REQUIRED libzstd) +find_package(zmqpp_vendor REQUIRED) if(BUILD_TESTING) find_package(launch_testing_ament_cmake REQUIRED) @@ -37,6 +38,10 @@ if(BUILD_TESTING) test/test_tcp.py TIMEOUT 2 # Sets a timeout for the test in seconds ) + add_launch_test( + test/test_zmq.py + TIMEOUT 10 # Sets a timeout for the test in seconds + ) endif() include_directories(include) @@ -55,6 +60,10 @@ add_library(tcp_interface SHARED src/network_interfaces/tcp_interface.cpp ) +add_library(zmq_interface SHARED + src/network_interfaces/zmq_interface.cpp +) + target_link_libraries(network_bridge PUBLIC ${std_msgs_TARGETS} ${tf2_msgs_TARGETS} @@ -76,6 +85,28 @@ target_link_libraries(tcp_interface PUBLIC ${Boost_LIBRARIES} ) +if(DEFINED zmqpp_vendor_INCLUDE_DIRS) + target_include_directories(zmq_interface PUBLIC + ${zmqpp_vendor_INCLUDE_DIRS} + ) +endif() + +if(DEFINED zmqpp_vendor_TARGETS) + target_link_libraries(zmq_interface PUBLIC + rclcpp::rclcpp + pluginlib::pluginlib + ${zmqpp_vendor_TARGETS} + ) +endif() + +if(DEFINED zmqpp_vendor_LIBRARIES) + target_link_libraries(zmq_interface PUBLIC + rclcpp::rclcpp + pluginlib::pluginlib + ${zmqpp_vendor_LIBRARIES} + ) +endif() + pluginlib_export_plugin_description_file(network_bridge network_interface_plugins.xml) install(TARGETS @@ -86,6 +117,7 @@ install(TARGETS install(TARGETS udp_interface tcp_interface + zmq_interface DESTINATION lib/ ) diff --git a/README.md b/README.md index 1598db8..d4cd007 100644 --- a/README.md +++ b/README.md @@ -1,22 +1,29 @@ # Network Bridge + [![CI](https://github.com/brow1633/network_bridge/actions/workflows/CI.yml/badge.svg)](https://github.com/brow1633/network_bridge/actions/workflows/CI.yml) -**Network Bridge** is a lightweight ROS2 node designed for robust communication between robotic systems over arbitrary network protocols. Supporting UDP and TCP protocols out of the box, this packages seamlessly bridges ROS2 topics across networks, facilitating effective remote communications between a base station and robotic systems, or between multiple robotic systems. +**Network Bridge** is a lightweight ROS2 node designed for robust communication between robotic systems over arbitrary network protocols. Supporting UDP, TCP, and ZMQ protocols out of the box, this packages seamlessly bridges ROS2 topics across networks, facilitating effective remote communications between a base station and robotic systems, or between multiple robotic systems. ## Installation + ### Installation via apt + Install with: + ``` sudo apt install ros--network-bridge ``` ### Building from Source + Simply clone the repository into your ROS2 workspace and build with `colcon build`. ## Usage ### Demo + #### TCP + ``` ros2 launch network_bridge tcp.launch.py @@ -26,6 +33,7 @@ ros2 topic echo /tcp2/MyDefaultTopic ``` #### UDP + ``` ros2 launch network_bridge udp.launch.py @@ -33,13 +41,29 @@ ros2 topic pub /udp1/MyDefaultTopic std_msgs/msg/String "data: 'Hello World'" ros2 topic echo /udp2/MyDefaultTopic ``` + +#### ZMQ + +``` +ros2 launch network_bridge zmq.launch.py + +ros2 topic pub /zmq1/MyDefaultTopic std_msgs/msg/String "data: 'Hello World'" + +ros2 topic echo /zmq2/MyDefaultTopic +``` + ### Configuration + Simply setup the network interface parameters and list your desired topics to get started. If you are using UDP over cellular data, it is recommended to setup a VPN to facilitate connection. Also, please note that **no encryption** occurs within this package. Currently, if you would like encryption, you must use a VPN. See `config/Udp1.yaml` for a description of all parameters, as well as the TCP example configuration files. + #### Minimal Example + The following configuration examples demonstrate a robot sending a message on `/gps/fix` over UDP to a basestation that will then re-publish the message. This works seamlessly on all message types, so long as they are built and sourced on both ends of the transmission. + #### Robot + ``` /udp_sender: ros__parameters: @@ -52,7 +76,9 @@ The following configuration examples demonstrate a robot sending a message on `/ topics: - "/gps/fix" ``` + #### Base Station + ``` /udp_receiver: ros__parameters: @@ -62,11 +88,14 @@ The following configuration examples demonstrate a robot sending a message on `/ remote_address: "192.168.1.2" send_port: 5001 ``` + #### Special case: TF + The TF topic `/tf` or `/tf_static` are handled as a special cases. The subscriber side will listen to all TF messages, accumulate them (similarly to a TF buffer) and send all of them at the specified rate. The behavior can be disabled or forced using the `is_tf` configuration. If some TFs need to be excluded or if the list of TFs to include is finite, one can use the include and exclude regex parameters. A transform is matched (hence excluded or included) if either the `frame_id` or `child_frame_id` are matching a pattern. + ``` /udp_sender: ros__parameters: @@ -92,21 +121,37 @@ If some TFs need to be excluded or if the list of TFs to include is finite, one ``` ### Choice of protocol + - **UDP**: Use UDP for low-latency, high-throughput communications, where occasional data loss is tolerable. Ideal for real-time telemetry data like sensor streams. - **TCP**: Opt for TCP when data integrity and reliability are critical. This ensures that control commands and state transitions are reliably delivered, though with potentially higher latency. +- **ZMQ**: Use ZMQ for reliable, high-performance messaging patterns. It provides a more robust and flexible communication architecture compared to raw TCP/UDP, abstracting away complex socket management. + +#### ZMQ Communication Patterns -Network protocols are implemented as pluginlib plugins, allowing the creation of arbitrary interfaces using the abstract class `include/network_interfaces/network_interface_base.hpp`. Any interface that can send and receive bytes could theoretically be implemented, including protocols that go beyond point-to-point communication, such as ZMQ. Please consider opening a pull request if you implement a new network interface. +ZMQ supports multiple messaging patterns. This package currently supports two, configurable via the `pattern` parameter: + +- **PUB/SUB** (`pattern: pub_sub`): One publisher broadcasts messages to multiple subscribers. Best for one-to-many data distribution (e.g., sensor streams to multiple consumers). Subscribers receive only the topics they subscribe to; late-joining subscribers miss messages sent before connection. + +- **PUSH/PULL** (`pattern: push_pull`): One pusher sends messages to a pool of pullers in a round-robin fashion. Best for load-balanced pipelines where each message must be processed by exactly one consumer. Provides back-pressure and queuing, unlike PUB/SUB. + +### Network Protocol Implementation + +Network protocols are implemented as pluginlib plugins, allowing the creation of arbitrary interfaces using the abstract class `include/network_interfaces/network_interface_base.hpp`. Any interface that can send and receive bytes could theoretically be implemented, including protocols that go beyond point-to-point communication. Please consider opening a pull request if you implement a new network interface. ### Tuning + This node can be launched with logger level DEBUG, which provides useful information for tuning the compression, rate and stale message parameters. For each message that is sent, the receiving side will output the number of bytes received, the decompressed size in bytes and the transmission delay. ### Contributing + Thank you for considering contributing! #### Code Formatting + Python code is formatted with `black`, and C++ is formatted with `uncrustify`. #### Pre-commit hooks + To ease the friction of linting, there are pre-commit hooks that you can install: ```bash @@ -120,4 +165,5 @@ pre-commit install # Run on commit automatically which will reformat code automatically when you commit changes. ## Acknowledgements -This package was developed for use in the Indy Autonomous Challenge by the Purdue AI Racing team. Inspiration was taken from mqtt_client (https://github.com/ika-rwth-aachen/mqtt_client/). + +This package was developed for use in the Indy Autonomous Challenge by the Purdue AI Racing team. Inspiration was taken from mqtt_client (). diff --git a/config/Zmq1.yaml b/config/Zmq1.yaml new file mode 100644 index 0000000..8480df9 --- /dev/null +++ b/config/Zmq1.yaml @@ -0,0 +1,20 @@ +# Configuration for the ZMQ Server (PUSH/PULL or PUB/SUB) +/**/zmq_bridge_server: + ros__parameters: + network_interface: "network_bridge::ZmqInterface" + + ZmqInterface: + role: "server" + # pattern: "pub_sub" # Options: pub_sub (default) or push_pull + port: 5555 + + default_rate: 10.0 + default_zstd_level: 3 + publish_stale_data: False + + # The server subscribes to this topic and pushes it to ZMQ + topics: + - "/test_topic" + + subscribe_namespace: "/zmq1" + publish_namespace: "/zmq1" diff --git a/config/Zmq1PushPull.yaml b/config/Zmq1PushPull.yaml new file mode 100644 index 0000000..baad6f5 --- /dev/null +++ b/config/Zmq1PushPull.yaml @@ -0,0 +1,20 @@ +# Configuration for the ZMQ Server (PUSH/PULL or PUB/SUB) +/**/zmq_bridge_server_push: + ros__parameters: + network_interface: "network_bridge::ZmqInterface" + + ZmqInterface: + role: "server" + pattern: "push_pull" # Options: pub_sub (default) or push_pull + port: 5556 + + default_rate: 10.0 + default_zstd_level: 3 + publish_stale_data: False + + # The server subscribes to this topic and pushes it to ZMQ + topics: + - "/test_topic" + + subscribe_namespace: "/zmq1/push_pull" + publish_namespace: "/zmq1/push_pull" diff --git a/config/Zmq2.yaml b/config/Zmq2.yaml new file mode 100644 index 0000000..6ab23cd --- /dev/null +++ b/config/Zmq2.yaml @@ -0,0 +1,17 @@ +# Configuration for the ZMQ Client (PUSH/PULL or PUB/SUB) +/**/zmq_bridge_client: + ros__parameters: + network_interface: "network_bridge::ZmqInterface" + + ZmqInterface: + role: "client" + # pattern: "pub_sub" # Options: pub_sub (default) or push_pull + remote_address: "127.0.0.1" + port: 5555 + + default_rate: 10.0 + default_zstd_level: 3 + publish_stale_data: False + + subscribe_namespace: "/zmq2" + publish_namespace: "/zmq2" diff --git a/config/Zmq2PushPull.yaml b/config/Zmq2PushPull.yaml new file mode 100644 index 0000000..613f5f2 --- /dev/null +++ b/config/Zmq2PushPull.yaml @@ -0,0 +1,17 @@ +# Configuration for the ZMQ Client (PUSH/PULL or PUB/SUB) +/**/zmq_bridge_client_pull: + ros__parameters: + network_interface: "network_bridge::ZmqInterface" + + ZmqInterface: + role: "client" + pattern: "push_pull" # Options: pub_sub (default) or push_pull + remote_address: "127.0.0.1" + port: 5556 + + default_rate: 10.0 + default_zstd_level: 3 + publish_stale_data: False + + subscribe_namespace: "/zmq2/push_pull" + publish_namespace: "/zmq2/push_pull" diff --git a/include/network_interfaces/zmq_interface.hpp b/include/network_interfaces/zmq_interface.hpp new file mode 100644 index 0000000..aa312a1 --- /dev/null +++ b/include/network_interfaces/zmq_interface.hpp @@ -0,0 +1,99 @@ +/* +============================================================================== +MIT License + +Copyright (c) 2024 Ethan M Brown +Copyright (c) 2026 PAL Robotics + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. +============================================================================== +*/ + +#pragma once + +#include +#include +#include +#include + +#include + +#include "network_interfaces/network_interface_base.hpp" + +namespace network_bridge +{ + +/** + * @class ZmqInterface + * @brief Represents a ZMQ network interface. + * + * The ZmqInterface class is a concrete implementation of the NetworkInterface + * abstract class. It provides functionality for opening, closing, receiving and + * writing data to a ZMQ interface. It also handles receiving data + * asynchronously and provides error handling capabilities. + */ +class ZmqInterface : public NetworkInterface +{ +public: + ZmqInterface() + : NetworkInterface() + { + ready_ = false; + failed_ = false; + } + + virtual ~ZmqInterface() {close();} + +protected: + /** + * @brief Initializes interface by loading parameters. + * + * Called from NetworkInterface::initialize() + */ + void initialize_() override; + +public: + bool has_failed() const override; + bool is_ready() const override; + void open() override; + void close() override; + void write(const std::vector & data) override; + +protected: + void load_parameters(); + void setup_server(); + void setup_client(); + void receive_thread(); + +private: + zmqpp::context context_; + std::shared_ptr socket_; + + std::string role_; + std::string pattern_; + std::string remote_address_; + int port_; + std::atomic ready_; + std::atomic failed_; + std::atomic shutting_down_; + + std::thread packet_thread_; +}; + +} // namespace network_bridge diff --git a/launch/zmq.launch.py b/launch/zmq.launch.py new file mode 100644 index 0000000..5cf4be9 --- /dev/null +++ b/launch/zmq.launch.py @@ -0,0 +1,28 @@ +from ament_index_python.packages import get_package_share_directory +from launch import LaunchDescription +from launch_ros.actions import Node + + +def generate_launch_description(): + config_dir = get_package_share_directory("network_bridge") + zmq1_config = config_dir + "/config/Zmq1.yaml" + zmq2_config = config_dir + "/config/Zmq2.yaml" + + return LaunchDescription( + [ + Node( + package="network_bridge", + executable="network_bridge", + name="zmq_bridge_server", + output="screen", + parameters=[zmq1_config], + ), + Node( + package="network_bridge", + executable="network_bridge", + name="zmq_bridge_client", + output="screen", + parameters=[zmq2_config], + ), + ] + ) diff --git a/network_interface_plugins.xml b/network_interface_plugins.xml index 3ed802c..df92cd9 100644 --- a/network_interface_plugins.xml +++ b/network_interface_plugins.xml @@ -14,4 +14,12 @@ + + + + + A ZeroMQ network interface for the network bridge. + + + diff --git a/package.xml b/package.xml index c68d187..927833a 100644 --- a/package.xml +++ b/package.xml @@ -24,6 +24,7 @@ tf2_msgs tf2_ros pluginlib + zmqpp_vendor diff --git a/src/network_interfaces/zmq_interface.cpp b/src/network_interfaces/zmq_interface.cpp new file mode 100644 index 0000000..e340d37 --- /dev/null +++ b/src/network_interfaces/zmq_interface.cpp @@ -0,0 +1,170 @@ +/* +============================================================================== +MIT License + +Copyright (c) 2024 Ethan M Brown +Copyright (c) 2026 PAL Robotics + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. +============================================================================== +*/ + +#include + +#include "network_interfaces/zmq_interface.hpp" + +namespace network_bridge +{ + +void ZmqInterface::initialize_() {load_parameters();} + +void ZmqInterface::load_parameters() +{ + std::string prefix = "ZmqInterface."; + node_->declare_parameter(prefix + "role", std::string("")); + node_->declare_parameter(prefix + "pattern", std::string("pub_sub")); + node_->declare_parameter(prefix + "remote_address", std::string("")); + node_->declare_parameter(prefix + "port", 0); + + node_->get_parameter(prefix + "role", role_); + node_->get_parameter(prefix + "pattern", pattern_); + node_->get_parameter(prefix + "remote_address", remote_address_); + node_->get_parameter(prefix + "port", port_); + + RCLCPP_INFO(node_->get_logger(), "role_: %s", role_.c_str()); + RCLCPP_INFO(node_->get_logger(), "pattern_: %s", pattern_.c_str()); + RCLCPP_INFO( + node_->get_logger(), "Remote Address: %s", + remote_address_.c_str()); + RCLCPP_INFO(node_->get_logger(), "Remote Port: %d", port_); +} + +void ZmqInterface::open() +{ + shutting_down_ = false; + failed_ = false; + ready_ = false; + if (role_ == "server") { + setup_server(); + ready_ = true; + } else if (role_ == "client") { + setup_client(); + ready_ = true; + packet_thread_ = + std::thread(std::bind(&ZmqInterface::receive_thread, this)); + } else { + RCLCPP_ERROR( + node_->get_logger(), "Invalid role specified: %s", + role_.c_str()); + failed_ = true; + return; + } +} + +bool ZmqInterface::is_ready() const {return ready_ && !failed_;} + +bool ZmqInterface::has_failed() const {return failed_;} + +void ZmqInterface::close() +{ + if (shutting_down_.exchange(true)) { + return; + } + ready_ = false; + + if (socket_) { + try { + socket_->close(); + } catch (const zmqpp::exception & e) { + RCLCPP_ERROR(node_->get_logger(), "ZMQ exception: %s", e.what()); + } + } + + if (packet_thread_.joinable()) { + packet_thread_.join(); + } + + try { + context_.terminate(); + } catch (const zmqpp::exception & e) { + RCLCPP_ERROR(node_->get_logger(), "ZMQ exception: %s", e.what()); + } +} + +void ZmqInterface::receive_thread() +{ + zmqpp::poller poller; + poller.add(*socket_, zmqpp::poller::poll_in); + while (!shutting_down_ && rclcpp::ok()) { + if (poller.poll(100)) { + zmqpp::message msg; + socket_->receive(msg); + const void * data = msg.raw_data(0); + size_t size = msg.size(0); + recv_cb_( + std::span(static_cast(data), size)); + } + } +} + +void ZmqInterface::setup_server() +{ + zmqpp::socket_type type = + (pattern_ == "pub_sub") ? zmqpp::socket_type::pub : zmqpp::socket_type::push; + socket_ = std::make_shared(context_, type); + try { + socket_->bind("tcp://*:" + std::to_string(port_)); + } catch (const zmqpp::exception & e) { + RCLCPP_ERROR(node_->get_logger(), "Bind failed: %s", e.what()); + failed_ = true; + return; + } + RCLCPP_INFO(node_->get_logger(), "Server bound to port %d", port_); +} + +void ZmqInterface::setup_client() +{ + zmqpp::socket_type type = + (pattern_ == "pub_sub") ? zmqpp::socket_type::sub : zmqpp::socket_type::pull; + socket_ = std::make_shared(context_, type); + if (pattern_ == "pub_sub") { + socket_->subscribe(""); + } + try { + socket_->connect("tcp://" + remote_address_ + ":" + std::to_string(port_)); + } catch (const zmqpp::exception & e) { + RCLCPP_ERROR(node_->get_logger(), "Connect failed: %s", e.what()); + failed_ = true; + return; + } + RCLCPP_INFO(node_->get_logger(), "Client connected to port %d", port_); +} + +void ZmqInterface::write(const std::vector & data) +{ + zmqpp::message msg; + msg.add_raw(data.data(), data.size()); + socket_->send(msg); +} + +} // namespace network_bridge + +PLUGINLIB_EXPORT_CLASS( + network_bridge::ZmqInterface, + network_bridge::NetworkInterface) diff --git a/test/test_zmq.py b/test/test_zmq.py new file mode 100644 index 0000000..f5d0d58 --- /dev/null +++ b/test/test_zmq.py @@ -0,0 +1,149 @@ +import time +import unittest + +from ament_index_python.packages import get_package_share_directory +import launch +import launch.actions +import launch_ros.actions +import launch_testing +import rclpy +from rclpy.node import Node +from rclpy.task import Future + +from std_msgs.msg import String + + +def generate_test_description(): + config = get_package_share_directory("network_bridge") + "/config/" + zmq1 = launch_ros.actions.Node( + package="network_bridge", + executable="network_bridge", + name="zmq_bridge_server", + output="screen", + parameters=[config + "Zmq1.yaml"], + arguments=["--ros-args", "--log-level", "debug", "--log-level", "rcl:=info"], + ) + + zmq2 = launch_ros.actions.Node( + package="network_bridge", + executable="network_bridge", + name="zmq_bridge_client", + output="screen", + parameters=[config + "Zmq2.yaml"], + arguments=["--ros-args", "--log-level", "debug", "--log-level", "rcl:=info"], + ) + + zmq1_push_pull = launch_ros.actions.Node( + package="network_bridge", + executable="network_bridge", + name="zmq_bridge_server_push", + output="screen", + parameters=[config + "Zmq1PushPull.yaml"], + arguments=["--ros-args", "--log-level", "debug", "--log-level", "rcl:=info"] + ) + + zmq2_push_pull = launch_ros.actions.Node( + package="network_bridge", + executable="network_bridge", + name="zmq_bridge_client_pull", + output="screen", + parameters=[config + "Zmq2PushPull.yaml"], + arguments=["--ros-args", "--log-level", "debug", "--log-level", "rcl:=info"] + ) + + return launch.LaunchDescription( + [ + zmq1, + launch.actions.TimerAction(period=0.1, actions=[zmq2]), + launch.actions.TimerAction(period=0.2, actions=[zmq1_push_pull]), + launch.actions.TimerAction(period=0.3, actions=[zmq2_push_pull]), + launch_testing.actions.ReadyToTest(), + ] + ) + + +class ZmqTestNode(Node): + + def __init__(self, name, pub_topic, sub_topic): + super().__init__(name) + self.test_message_received = Future() + self.received_msg = None + self.publisher = self.create_publisher(String, pub_topic, 10) + self.subscriber = self.create_subscription( + String, sub_topic, self.listener_callback, 10 + ) + + def publish(self, msg): + self.publisher.publish(msg) + + def listener_callback(self, msg): + self.received_msg = msg + self.test_message_received.set_result(True) + + +class TestZmq(unittest.TestCase): + + @classmethod + def setUpClass(cls): + rclpy.init() + + @classmethod + def tearDownClass(cls): + rclpy.shutdown() + + def test_node_output(self, proc_output): + proc_output.assertWaitFor("Server bound", timeout=0.5) + proc_output.assertWaitFor("Client connected", timeout=0.5) + + node = ZmqTestNode( + "zmq_test_node", + "/zmq1/test_topic", + "/zmq2/test_topic") + time.sleep(1.5) + + test_msg = String() + test_msg.data = "Testing123" + node.publish(test_msg) + + try: + rclpy.spin_until_future_complete( + node, node.test_message_received, timeout_sec=10.0 + ) + self.assertTrue( + node.test_message_received.done(), "Timeout on message receival." + ) + self.assertEqual( + node.received_msg.data, + "Testing123", + "The received message did not match the expected output.", + ) + finally: + node.destroy_node() + + node_push_pull = ZmqTestNode( + "zmq_push_pull_test_node", + "/zmq1/push_pull/test_topic", + "/zmq2/push_pull/test_topic") + time.sleep(1.5) + + node_push_pull.publish(test_msg) + + try: + rclpy.spin_until_future_complete( + node_push_pull, + node_push_pull.test_message_received, + timeout_sec=10.0 + ) + self.assertTrue( + node_push_pull.test_message_received.done(), "Timeout on message receival." + ) + self.assertEqual( + node_push_pull.received_msg.data, + "Testing123", + "The received message did not match the expected output.", + ) + finally: + node_push_pull.destroy_node() + +if __name__ == "__main__": + launch_testing.main()