diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index aad0b15a..70234f29 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -191,6 +191,14 @@ jobs: # Reset only the turtlesim instance as it is not destroyed at the end of the previous job source /opt/ros/humble/setup.bash && ros2 service call /reset std_srvs/srv/Empty & cargo run --example python-ros2-dataflow --features="ros2-examples" + - name: "c++-ros2-dataflow" + timeout-minutes: 30 + env: + QT_QPA_PLATFORM: offscreen + run: | + # Reset only the turtlesim instance as it is not destroyed at the end of the previous job + source /opt/ros/humble/setup.bash && ros2 service call /reset std_srvs/srv/Empty & + cargo run --example cxx-ros2-dataflow --features="ros2-examples" bench: name: "Bench" diff --git a/Cargo.lock b/Cargo.lock index 78ac0396..a2f623d0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -405,7 +405,7 @@ dependencies = [ "async-task", "concurrent-queue", "fastrand 2.0.1", - "futures-lite", + "futures-lite 1.13.0", "slab", ] @@ -420,7 +420,7 @@ dependencies = [ "async-io", "async-lock", "blocking", - "futures-lite", + "futures-lite 1.13.0", "once_cell", "tokio", ] @@ -435,7 +435,7 @@ dependencies = [ "autocfg", "cfg-if 1.0.0", "concurrent-queue", - "futures-lite", + "futures-lite 1.13.0", "log", "parking", "polling", @@ -466,7 +466,7 @@ dependencies = [ "blocking", "cfg-if 1.0.0", "event-listener 2.5.3", - "futures-lite", + "futures-lite 1.13.0", "rustix 0.37.25", "signal-hook", "windows-sys 0.48.0", @@ -499,7 +499,7 @@ dependencies = [ "futures-channel", "futures-core", "futures-io", - "futures-lite", + "futures-lite 1.13.0", "gloo-timers", "kv-log-macro", "log", @@ -748,7 +748,7 @@ dependencies = [ "async-task", "fastrand 2.0.1", "futures-io", - "futures-lite", + "futures-lite 1.13.0", "piper", "tracing", ] @@ -1582,7 +1582,14 @@ dependencies = [ "cxx", "cxx-build", "dora-node-api", + "dora-ros2-bridge", + "dora-ros2-bridge-msg-gen", "eyre", + "futures-lite 2.2.0", + "prettyplease", + "rust-format", + "serde", + "serde-big-array", ] [[package]] @@ -1675,11 +1682,12 @@ version = "0.1.0" dependencies = [ "array-init", "dora-daemon", - "dora-ros2-bridge-msg-gen-macro", + "dora-ros2-bridge-msg-gen", "eyre", "futures", "rand", "ros2-client", + "rust-format", "rustdds", "serde", "serde-big-array", @@ -1703,21 +1711,6 @@ dependencies = [ "thiserror", ] -[[package]] -name = "dora-ros2-bridge-msg-gen-macro" -version = "0.1.0" -dependencies = [ - "anyhow", - "dora-ros2-bridge-msg-gen", - "heck 0.3.3", - "nom", - "proc-macro2", - "quote", - "regex", - "syn 1.0.109", - "thiserror", -] - [[package]] name = "dora-ros2-bridge-python" version = "0.1.0" @@ -2119,6 +2112,19 @@ dependencies = [ "waker-fn", ] +[[package]] +name = "futures-lite" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "445ba825b27408685aaecefd65178908c36c6e96aaf6d8599419d46e624192ba" +dependencies = [ + "fastrand 2.0.1", + "futures-core", + "futures-io", + "parking", + "pin-project-lite", +] + [[package]] name = "futures-macro" version = "0.3.28" @@ -3782,9 +3788,9 @@ checksum = "b15813163c1d831bf4a13c3610c05c0d03b39feb07f7e09fa234dac9b15aaf39" [[package]] name = "parking" -version = "2.1.1" +version = "2.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e52c774a4c39359c1d1c52e43f73dd91a75a614652c825408eec30c95a9b2067" +checksum = "bb813b8af86854136c6922af0598d719255ecb2179515e6e7730d468f05c9cae" [[package]] name = "parking_lot" @@ -4629,6 +4635,17 @@ dependencies = [ "eyre", ] +[[package]] +name = "rust-format" +version = "0.3.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "60e7c00b6c3bf5e38a880eec01d7e829d12ca682079f8238a464def3c4b31627" +dependencies = [ + "prettyplease", + "proc-macro2", + "syn 1.0.109", +] + [[package]] name = "rust-ros2-dataflow-example-node" version = "0.3.2" @@ -4798,9 +4815,9 @@ checksum = "ef703b7cb59335eae2eb93ceb664c0eb7ea6bf567079d843e09420219668e072" [[package]] name = "safer-ffi" -version = "0.1.4" +version = "0.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "395ace5aff9629c7268ca8255aceb945525b2cb644015f3caec5131a6a537c11" +checksum = "4483c5ab47f222d2c297e73a520c9003e09e2fe1f1b04edcb572e6939f303003" dependencies = [ "inventory 0.1.11", "inventory 0.3.12", @@ -4816,9 +4833,9 @@ dependencies = [ [[package]] name = "safer_ffi-proc_macros" -version = "0.1.4" +version = "0.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9255504d5467bae9e07d58b8de446ba6739b29bf72e1fa35b2387e30d29dcbfe" +checksum = "bf04ebd3786110e64269a74eea58c5564dd92a1e790c0f6f9871d6fe1b8e34db" dependencies = [ "macro_rules_attribute", "prettyplease", diff --git a/Cargo.toml b/Cargo.toml index 6b27c18e..34fe4fbb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -27,7 +27,6 @@ members = [ "libraries/extensions/dora-record", "libraries/extensions/ros2-bridge", "libraries/extensions/ros2-bridge/msg-gen", - "libraries/extensions/ros2-bridge/msg-gen-macro", "libraries/extensions/ros2-bridge/python", ] @@ -59,6 +58,7 @@ dora-runtime = { version = "0.3.2", path = "binaries/runtime" } dora-daemon = { version = "0.3.2", path = "binaries/daemon" } dora-coordinator = { version = "0.3.2", path = "binaries/coordinator" } dora-ros2-bridge = { path = "libraries/extensions/ros2-bridge" } +dora-ros2-bridge-msg-gen = { path = "libraries/extensions/ros2-bridge/msg-gen" } dora-ros2-bridge-python = { path = "libraries/extensions/ros2-bridge/python" } arrow = "48.0.0" arrow-schema = "48.0.0" @@ -138,3 +138,8 @@ path = "examples/multiple-daemons/run.rs" [[example]] name = "cmake-dataflow" path = "examples/cmake-dataflow/run.rs" + +[[example]] +name = "cxx-ros2-dataflow" +path = "examples/c++-ros2-dataflow/run.rs" +required-features = ["ros2-examples"] diff --git a/apis/c++/node/Cargo.toml b/apis/c++/node/Cargo.toml index 77e58c25..ff27a2e8 100644 --- a/apis/c++/node/Cargo.toml +++ b/apis/c++/node/Cargo.toml @@ -14,11 +14,28 @@ crate-type = ["staticlib"] [features] default = ["tracing"] tracing = ["dora-node-api/tracing"] +ros2-bridge = [ + "dep:dora-ros2-bridge", + "dep:dora-ros2-bridge-msg-gen", + "dep:rust-format", + "dep:prettyplease", + "dep:serde", + "dep:serde-big-array", +] [dependencies] cxx = "1.0.73" dora-node-api = { workspace = true } eyre = "0.6.8" +dora-ros2-bridge = { workspace = true, optional = true } +futures-lite = { version = "2.2" } +serde = { version = "1.0.164", features = ["derive"], optional = true } +serde-big-array = { version = "0.5.1", optional = true } [build-dependencies] cxx-build = "1.0.73" +dora-ros2-bridge-msg-gen = { workspace = true, optional = true } +rust-format = { version = "0.3.4", features = [ + "pretty_please", +], optional = true } +prettyplease = { version = "0.1", features = ["verbatim"], optional = true } diff --git a/apis/c++/node/README.md b/apis/c++/node/README.md new file mode 100644 index 00000000..c18fb2e0 --- /dev/null +++ b/apis/c++/node/README.md @@ -0,0 +1,230 @@ +# Dora Node API for C++ + +Dora supports nodes written in C++ through this API crate. + +## Build + +- Clone the `dora` repository: + ```bash + > git clone https://github.com/dora-rs/dora.git + > cd dora + ``` +- Build the `dora-node-api-cxx` package: + ```bash + cargo build --package dora-node-api-cxx + ``` + - This will result in `dora-node-api.h` and `dora-node-api.cc` files in the `target/cxxbridge/dora-node-api-cxx` directory. +- Include the `dora-node-api.h` header file in your source file. +- Add the `dora-node-api.cc` file to your compile and link steps. + +### Build with ROS2 Bridge + +Dora features an experimental ROS2 Bridge that enables dora nodes to publish and subscribe to ROS2 topics. +To enable the bridge, use these steps: + +- Clone the `dora` repository (see above). +- Source the ROS2 setup files (see [ROS2 docs](https://docs.ros.org/en/rolling/Tutorials/Beginner-CLI-Tools/Configuring-ROS2-Environment.html#source-the-setup-files)) +- Optional: Source package-specific ROS2 setup files if you want to use custom package-specific ROS2 messages in the bridge (see [ROS2 docs](https://docs.ros.org/en/rolling/Tutorials/Beginner-Client-Libraries/Creating-Your-First-ROS2-Package.html#source-the-setup-file)) +- Build the `dora-node-api-cxx` package **with the `ros2-bridge` feature enabled**: + ```bash + cargo build --package dora-node-api-cxx --features ros2-bridge + ``` + - In addition to the `dora-node-api.h` and `dora-node-api.cc` files, this will place a `dora-ros2-bindings.h` and a `dora-ros2-bindings.cc` file in the `target/cxxbridge/dora-node-api-cxx` directory. +- Include both the `dora-node-api.h` and the `dora-ros2-bindings.h` header files in your source file. +- Add the `dora-node-api.cc` and `dora-ros2-bindings.cc` files to your compile and link steps. + +## Usage + +The `dora-node-api.h` header provides various functions to interact with Dora. + +### Init Dora Node + +All nodes need to register themselves with Dora at startup. +To do that, call the `init_dora_node()` function. +The function returns a `DoraNode` instance, which gives access to dora events and enables sending Dora outputs. + +```c++ +auto dora_node = init_dora_node(); +``` + +### Receiving Events + +The `dora_node.events` field is a stream of incoming events. +To wait for the next incoming event, call `dora_node.events->next()`: + +```c++ +auto event = dora_node.events->next(); +``` + +The `next` function returns an opaque `DoraEvent` type, which cannot be inspected from C++ directly. +Instead, use the following functions to read and destructure the event: + +- `event_type(event)` returns a `DoraEventType`, which describes the kind of event. For example, an event could be an input or a stop instruction. + - when receiving a `DoraEventType::AllInputsClosed`, the node should exit and not call `next` anymore +- Events of type `DoraEventType::Input` can be downcasted using `event_as_input`: + ```c++ + auto input = event_as_input(std::move(event)); + ``` + The function returns a `DoraInput` instance, which has an `id` and `data` field. + - The input `id` can be converted to a C++ string through `std::string(input.id)`. + - The `data` of inputs is currently of type [`rust::Vec`](https://cxx.rs/binding/vec.html). Use the provided methods for reading or converting the data. + - **Note:** In the future, we plan to change the data type to the [Apache Arrow](https://arrow.apache.org/) data format to support typed inputs. + +### Sending Outputs + +Nodes can send outputs using the `send_output` output function and the `dora_node.send_output` field. +Note that all outputs need to be listed in the dataflow YAML declaration file, otherwise an error will occur. + +**Example:** + +```c++ +// the data you want to send (NOTE: only byte vectors are supported right now) +std::vector out_vec{42}; +// create a Rust slice from the output vector +rust::Slice out_slice{out_vec.data(), out_vec.size()}; +// send the slice as output +auto result = send_output(dora_node.send_output, "output_id", out_slice); + +// check for errors +auto error = std::string(result.error); +if (!error.empty()) +{ + std::cerr << "Error: " << error << std::endl; + return -1; +} +``` + +## Using the ROS2 Bridge + +The `dora-ros2-bindings.h` contains function and struct definitions that allow interacting with ROS2 nodes. +Currently, the bridge supports publishing and subscribing to ROS2 topics. +In the future, we plan to support ROS2 services and ROS2 actions as well. + +### Initializing the ROS2 Context + +The first step is to initialize a ROS2 context: + +```c++ +auto ros2_context = init_ros2_context(); +``` + +### Creating Nodes + +After initializing a ROS2 context, you can use it to create ROS2 nodes: + +```c++ +auto node = ros2_context->new_node("/ros2_demo", "turtle_teleop"); +``` + +The first argument is the namespace of the node and the second argument is its name. + +### Creating Topics + +After creating a node, you can use one of the `create_topic_` functions to create a topic on it. +The `` describes the message type that will be sent on the topic. +The Dora ROS2 bridge automatically creates `create_topic_` functions for all messages types found in the sourced ROS2 environment. + +```c++ +auto vel_topic = node->create_topic_geometry_msgs_Twist("/turtle1", "cmd_vel", qos_default()); +``` + +The first argument is the namespace of the topic and the second argument is its name. +The third argument is the QoS (quality of service) setting for the topic. +It can be adjusted as desired, for example: + +```c++ +auto qos = qos_default(); +qos.durability = Ros2Durability::Volatile; +qos.liveliness = Ros2Liveliness::Automatic; +auto vel_topic = node->create_topic_geometry_msgs_Twist("/turtle1", "cmd_vel", qos); +``` + +### Publish + +After creating a topic, it is possible to publish messages on it. +First, create a publisher: + +```c++ +auto vel_publisher = node->create_publisher(vel_topic, qos); +``` + +The returned publisher is typed by the chosen topic. +It will only accept messages of the topic's type, otherwise a compile error will occur. + +After creating a publisher, you can use the `publish` function to publish one or more messages. +For example: + +```c++ +geometry_msgs::Twist twist = { + .linear = {.x = 1, .y = 0, .z = 0}, + .angular = {.x = 0, .y = 0, .z = 0.5} +}; +vel_publisher->publish(twist); +``` + +The `geometry_msgs::Twist` struct is automatically generated from the sourced ROS2 environment. +Since the publisher is typed, its `publish` method only accepts `geometry_msgs::Twist` messages. + + +### Subscriptions + +Subscribing to a topic is possible through the `create_subscription` function on nodes: + +```c++ +auto pose_topic = node->create_topic_turtlesim_Pose("/turtle1", "pose", qos_default()); +auto pose_subscription = node->create_subscription(pose_topic, qos_default(), event_stream); +``` + +The `topic` is the topic you want to subscribe to, created using a `create_topic_` function. +The second argument is the quality of service setting, which can be customized as described above. + +The third parameter is the event stream that the received messages should be merged into. +You can create such a event stream from Dora's event stream using the `dora_events_into_combined` function: + +```c++ +auto event_stream = dora_events_into_combined(std::move(dora_node.events)); +``` + +Multiple subscriptions can be merged into the same event stream. + +#### Receiving Messages from Merged Event Stream + +The merged event stream will receive all incoming events of the node, including Dora events and ROS2 messages. +To wait for the next incoming event, use its `next` method: + +```c++ +auto event = event_stream.next(); +``` + +This returns a `event` instance of type `CombinedEvent`, which can be downcasted to Dora events or ROS2 messages. +To handle an event, you should check it's type and then downcast it: + +- To check for a Dora event, you can use the `is_dora()` function. If it returns `true`, you can downcast the combined event to a Dora event using the `downcast_dora` function. +- ROS2 subscriptions support a `matches` function to check whether a combined event is an instance of the respective ROS2 subscription. If it returns true, you can downcast the event to the respective ROS2 message struct using the subscription's `downcast` function. + +**Example:** + +```c++ +if (event.is_dora()) +{ + auto dora_event = downcast_dora(std::move(event)); + // handle dora_event as described above + auto ty = event_type(dora_event); + if (ty == DoraEventType::Input) + { + auto input = event_as_input(std::move(dora_event)); + // etc + } + // .. else if +} +else if (pose_subscription->matches(event)) +{ + auto pose = pose_subscription->downcast(std::move(event)); + std::cout << "Received pose x:" << pose.x << ", y:" << pose.y << std::endl; +} +else +{ + std::cout << "received unexpected event" << std::endl; +} +``` + diff --git a/apis/c++/node/build.rs b/apis/c++/node/build.rs index dba00707..767bac2e 100644 --- a/apis/c++/node/build.rs +++ b/apis/c++/node/build.rs @@ -1,4 +1,146 @@ +use std::path::{Path, PathBuf}; + fn main() { - let _build = cxx_build::bridge("src/lib.rs"); + let mut bridge_files = vec![PathBuf::from("src/lib.rs")]; + #[cfg(feature = "ros2-bridge")] + bridge_files.push(ros2::generate()); + + let _build = cxx_build::bridges(&bridge_files); println!("cargo:rerun-if-changed=src/lib.rs"); + + // rename header files + let src_dir = target_dir() + .join("cxxbridge") + .join("dora-node-api-cxx") + .join("src"); + let target_dir = src_dir.parent().unwrap(); + std::fs::copy(src_dir.join("lib.rs.h"), target_dir.join("dora-node-api.h")).unwrap(); + std::fs::copy( + src_dir.join("lib.rs.cc"), + target_dir.join("dora-node-api.cc"), + ) + .unwrap(); + + #[cfg(feature = "ros2-bridge")] + ros2::generate_ros2_message_header(bridge_files.last().unwrap()); + + // to avoid unnecessary `mut` warning + bridge_files.clear(); +} + +fn target_dir() -> PathBuf { + std::env::var("CARGO_TARGET_DIR") + .map(PathBuf::from) + .unwrap_or_else(|_| { + let root = Path::new(env!("CARGO_MANIFEST_DIR")) + .ancestors() + .nth(3) + .unwrap(); + root.join("target") + }) +} + +#[cfg(feature = "ros2-bridge")] +mod ros2 { + use super::target_dir; + use std::path::{Component, Path, PathBuf}; + + pub fn generate() -> PathBuf { + use rust_format::Formatter; + let paths = ament_prefix_paths(); + let generated = dora_ros2_bridge_msg_gen::gen(paths.as_slice(), true); + let generated_string = rust_format::PrettyPlease::default() + .format_tokens(generated) + .unwrap(); + let out_dir = PathBuf::from(std::env::var("OUT_DIR").unwrap()); + let target_file = out_dir.join("ros2_bindings.rs"); + std::fs::write(&target_file, generated_string).unwrap(); + println!( + "cargo:rustc-env=ROS2_BINDINGS_PATH={}", + target_file.display() + ); + + target_file + } + + fn ament_prefix_paths() -> Vec { + let ament_prefix_path: String = match std::env::var("AMENT_PREFIX_PATH") { + Ok(path) => path, + Err(std::env::VarError::NotPresent) => { + println!("cargo:warning='AMENT_PREFIX_PATH not set'"); + String::new() + } + Err(std::env::VarError::NotUnicode(s)) => { + panic!( + "AMENT_PREFIX_PATH is not valid unicode: `{}`", + s.to_string_lossy() + ); + } + }; + println!("cargo:rerun-if-env-changed=AMENT_PREFIX_PATH"); + + let paths: Vec<_> = ament_prefix_path.split(':').map(PathBuf::from).collect(); + for path in &paths { + println!("cargo:rerun-if-changed={}", path.display()); + } + + paths + } + + pub fn generate_ros2_message_header(source_file: &Path) { + use std::io::Write as _; + + let out_dir = source_file.parent().unwrap(); + let relative_path = local_relative_path(&source_file) + .ancestors() + .nth(2) + .unwrap() + .join("out"); + let header_path = out_dir + .join("cxxbridge") + .join("include") + .join("dora-node-api-cxx") + .join(&relative_path) + .join("ros2_bindings.rs.h"); + let code_path = out_dir + .join("cxxbridge") + .join("sources") + .join("dora-node-api-cxx") + .join(&relative_path) + .join("ros2_bindings.rs.cc"); + + // copy message files to target directory + let target_path = target_dir() + .join("cxxbridge") + .join("dora-node-api-cxx") + .join("dora-ros2-bindings.h"); + + std::fs::copy(&header_path, &target_path).unwrap(); + println!("cargo:rerun-if-changed={}", header_path.display()); + + let mut node_header = + std::fs::File::open(target_path.with_file_name("dora-node-api.h")).unwrap(); + let mut code_file = std::fs::File::open(&code_path).unwrap(); + println!("cargo:rerun-if-changed={}", code_path.display()); + let mut code_target_file = + std::fs::File::create(target_path.with_file_name("dora-ros2-bindings.cc")).unwrap(); + + // copy both the node header and the code file to prevent import errors + std::io::copy(&mut node_header, &mut code_target_file).unwrap(); + std::io::copy(&mut code_file, &mut code_target_file).unwrap(); + code_target_file.flush().unwrap(); + } + + // copy from cxx-build source + fn local_relative_path(path: &Path) -> PathBuf { + let mut rel_path = PathBuf::new(); + for component in path.components() { + match component { + Component::Prefix(_) | Component::RootDir | Component::CurDir => {} + Component::ParentDir => drop(rel_path.pop()), // noop if empty + Component::Normal(name) => rel_path.push(name), + } + } + rel_path + } } diff --git a/apis/c++/node/src/lib.rs b/apis/c++/node/src/lib.rs index 61220ecf..d5c03528 100644 --- a/apis/c++/node/src/lib.rs +++ b/apis/c++/node/src/lib.rs @@ -1,10 +1,17 @@ +use std::any::Any; + use dora_node_api::{ self, arrow::array::{AsArray, BinaryArray}, + merged::{MergeExternal, MergedEvent}, Event, EventStream, }; use eyre::bail; +#[cfg(feature = "ros2-bridge")] +use dora_ros2_bridge::_core; +use futures_lite::{Stream, StreamExt}; + #[cxx::bridge] #[allow(clippy::needless_lifetimes)] mod ffi { @@ -31,14 +38,26 @@ mod ffi { error: String, } + pub struct CombinedEvents { + events: Box, + } + + pub struct CombinedEvent { + event: Box, + } + extern "Rust" { type Events; type OutputSender; type DoraEvent; + type MergedEvents; + type MergedDoraEvent; fn init_dora_node() -> Result; - fn next_event(inputs: &mut Box) -> Box; + fn dora_events_into_combined(events: Box) -> CombinedEvents; + fn next(self: &mut Events) -> Box; + fn next_event(events: &mut Box) -> Box; fn event_type(event: &Box) -> DoraEventType; fn event_as_input(event: Box) -> Result; fn send_output( @@ -46,9 +65,20 @@ mod ffi { id: String, data: &[u8], ) -> DoraResult; + + fn next(self: &mut CombinedEvents) -> CombinedEvent; + + fn is_dora(self: &CombinedEvent) -> bool; + fn downcast_dora(event: CombinedEvent) -> Result>; } } +#[cfg(feature = "ros2-bridge")] +pub mod ros2 { + pub use dora_ros2_bridge::*; + include!(env!("ROS2_BINDINGS_PATH")); +} + fn init_dora_node() -> eyre::Result { let (node, events) = dora_node_api::DoraNode::init_from_env()?; let events = Events(events); @@ -62,8 +92,24 @@ fn init_dora_node() -> eyre::Result { pub struct Events(EventStream); +impl Events { + fn next(&mut self) -> Box { + Box::new(DoraEvent(self.0.recv())) + } +} + fn next_event(events: &mut Box) -> Box { - Box::new(DoraEvent(events.0.recv())) + events.next() +} + +fn dora_events_into_combined(events: Box) -> ffi::CombinedEvents { + let events = events.0.map(MergedEvent::Dora); + ffi::CombinedEvents { + events: Box::new(MergedEvents { + events: Some(Box::new(events)), + next_id: 1, + }), + } } pub struct DoraEvent(Option); @@ -111,3 +157,59 @@ fn send_output(sender: &mut Box, id: String, data: &[u8]) -> ffi:: }; ffi::DoraResult { error } } + +pub struct MergedEvents { + events: Option> + Unpin>>, + next_id: u32, +} + +impl MergedEvents { + fn next(&mut self) -> MergedDoraEvent { + let event = futures_lite::future::block_on(self.events.as_mut().unwrap().next()); + MergedDoraEvent(event) + } + + pub fn merge(&mut self, events: impl Stream> + Unpin + 'static) -> u32 { + let id = self.next_id; + self.next_id += 1; + let events = Box::pin(events.map(move |event| ExternalEvent { event, id })); + + let inner = self.events.take().unwrap(); + let merged: Box + Unpin + 'static> = + Box::new(inner.merge_external(events).map(|event| match event { + MergedEvent::Dora(event) => MergedEvent::Dora(event), + MergedEvent::External(event) => MergedEvent::External(event.flatten()), + })); + self.events = Some(merged); + + id + } +} + +impl ffi::CombinedEvents { + fn next(&mut self) -> ffi::CombinedEvent { + ffi::CombinedEvent { + event: Box::new(self.events.next()), + } + } +} + +pub struct MergedDoraEvent(Option>); + +pub struct ExternalEvent { + pub event: Box, + pub id: u32, +} + +impl ffi::CombinedEvent { + fn is_dora(&self) -> bool { + matches!(&self.event.0, Some(MergedEvent::Dora(_))) + } +} + +fn downcast_dora(event: ffi::CombinedEvent) -> eyre::Result> { + match event.event.0 { + Some(MergedEvent::Dora(event)) => Ok(Box::new(DoraEvent(Some(event)))), + _ => eyre::bail!("not an external event"), + } +} diff --git a/apis/c/operator/operator_types.h b/apis/c/operator/operator_types.h index c6e3c6f8..9cf2f3d2 100644 --- a/apis/c/operator/operator_types.h +++ b/apis/c/operator/operator_types.h @@ -34,7 +34,7 @@ typedef struct Vec_uint8 { /** */ typedef struct DoraResult { /** */ - Vec_uint8_t error; + Vec_uint8_t * error; } DoraResult_t; /** */ diff --git a/apis/rust/operator/src/lib.rs b/apis/rust/operator/src/lib.rs index fd770c78..262fcdd3 100644 --- a/apis/rust/operator/src/lib.rs +++ b/apis/rust/operator/src/lib.rs @@ -64,9 +64,6 @@ impl DoraOutputSender<'_> { open_telemetry_context: String::new().into(), // TODO }, }); - match result.error { - None => Ok(()), - Some(error) => Err(error.into()), - } + result.into_result() } } diff --git a/apis/rust/operator/src/raw.rs b/apis/rust/operator/src/raw.rs index 46aa9638..2ff89c64 100644 --- a/apis/rust/operator/src/raw.rs +++ b/apis/rust/operator/src/raw.rs @@ -24,7 +24,7 @@ pub unsafe fn dora_init_operator() -> DoraInitResult { pub unsafe fn dora_drop_operator(operator_context: *mut c_void) -> DoraResult { let raw: *mut O = operator_context.cast(); - let _ = unsafe { Box::from_raw(raw) }; + drop(unsafe { Box::from_raw(raw) }); DoraResult { error: None } } @@ -40,9 +40,7 @@ pub unsafe fn dora_on_event( let event_variant = if let Some(input) = &mut event.input { let Some(data_array) = input.data_array.take() else { return OnEventResult { - result: DoraResult { - error: Some("data already taken".to_string().into()), - }, + result: DoraResult::from_error("data already taken".to_string()), status: DoraStatus::Continue, }; }; @@ -75,9 +73,7 @@ pub unsafe fn dora_on_event( status, }, Err(error) => OnEventResult { - result: DoraResult { - error: Some(error.into()), - }, + result: DoraResult::from_error(error), status: DoraStatus::Stop, }, } diff --git a/apis/rust/operator/types/src/lib.rs b/apis/rust/operator/types/src/lib.rs index e92d4538..7bff6234 100644 --- a/apis/rust/operator/types/src/lib.rs +++ b/apis/rust/operator/types/src/lib.rs @@ -15,7 +15,7 @@ use safer_ffi::{ closure::ArcDynFn1, derive_ReprC, ffi_export, }; -use std::path::Path; +use std::{ops::Deref, path::Path}; #[derive_ReprC] #[ffi_export] @@ -44,7 +44,31 @@ pub struct DoraDropOperator { #[repr(C)] #[derive(Debug)] pub struct DoraResult { - pub error: Option, + pub error: Option>, +} + +impl DoraResult { + pub const SUCCESS: Self = Self { error: None }; + + pub fn from_error(error: String) -> Self { + Self { + error: Some(Box::new(safer_ffi::String::from(error)).into()), + } + } + + pub fn error(&self) -> Option<&str> { + self.error.as_deref().map(|s| s.deref()) + } + + pub fn into_result(self) -> Result<(), String> { + match self.error { + None => Ok(()), + Some(error) => { + let converted = safer_ffi::boxed::Box_::into(error); + Err((*converted).into()) + } + } + } } #[derive_ReprC] @@ -175,7 +199,7 @@ pub unsafe fn dora_send_operator_output( match result() { Ok(output) => send_output.send_output.call(output), Err(error) => DoraResult { - error: Some(error.into()), + error: Some(Box::new(safer_ffi::String::from(error)).into()), }, } } diff --git a/binaries/runtime/src/operator/shared_lib.rs b/binaries/runtime/src/operator/shared_lib.rs index 6a91102c..495c1b6c 100644 --- a/binaries/runtime/src/operator/shared_lib.rs +++ b/binaries/runtime/src/operator/shared_lib.rs @@ -98,7 +98,7 @@ impl<'lib> SharedLibraryOperator<'lib> { let raw = match result.error { Some(error) => { let _ = init_done.send(Err(eyre!(error.to_string()))); - bail!("init_operator failed: {}", String::from(error)) + bail!("init_operator failed: {}", *error) } None => operator_context, }; @@ -126,11 +126,7 @@ impl<'lib> SharedLibraryOperator<'lib> { let arrow_array = match arrow::ffi::from_ffi(data_array, &schema) { Ok(a) => a, - Err(err) => { - return DoraResult { - error: Some(err.to_string().into()), - } - } + Err(err) => return DoraResult::from_error(err.to_string()), }; let total_len = required_data_size(&arrow_array); @@ -138,11 +134,7 @@ impl<'lib> SharedLibraryOperator<'lib> { let type_info = match copy_array_into_sample(&mut sample, &arrow_array) { Ok(t) => t, - Err(err) => { - return DoraResult { - error: Some(err.to_string().into()), - } - } + Err(err) => return DoraResult::from_error(err.to_string()), }; let event = OperatorEvent::Output { @@ -157,12 +149,10 @@ impl<'lib> SharedLibraryOperator<'lib> { .blocking_send(event) .map_err(|_| eyre!("failed to send output to runtime")); - let error = match result { - Ok(()) => None, - Err(_) => Some(String::from("runtime process closed unexpectedly").into()), - }; - - DoraResult { error } + match result { + Ok(()) => DoraResult::SUCCESS, + Err(_) => DoraResult::from_error("runtime process closed unexpectedly".into()), + } }); let reason = loop { @@ -261,7 +251,7 @@ impl<'lib> SharedLibraryOperator<'lib> { ) }; match error { - Some(error) => bail!("on_input failed: {}", String::from(error)), + Some(error) => bail!("on_input failed: {}", *error), None => match status { DoraStatus::Continue => {} DoraStatus::Stop => break StopReason::ExplicitStop, diff --git a/examples/c++-dataflow/node-rust-api/main.cc b/examples/c++-dataflow/node-rust-api/main.cc index 31415523..9407650a 100644 --- a/examples/c++-dataflow/node-rust-api/main.cc +++ b/examples/c++-dataflow/node-rust-api/main.cc @@ -13,7 +13,7 @@ int main() for (int i = 0; i < 20; i++) { - auto event = next_event(dora_node.events); + auto event = dora_node.events->next(); auto ty = event_type(event); if (ty == DoraEventType::AllInputsClosed) diff --git a/examples/c++-ros2-dataflow/.gitignore b/examples/c++-ros2-dataflow/.gitignore new file mode 100644 index 00000000..5761abcf --- /dev/null +++ b/examples/c++-ros2-dataflow/.gitignore @@ -0,0 +1 @@ +*.o diff --git a/examples/c++-ros2-dataflow/dataflow.yml b/examples/c++-ros2-dataflow/dataflow.yml new file mode 100644 index 00000000..7639c603 --- /dev/null +++ b/examples/c++-ros2-dataflow/dataflow.yml @@ -0,0 +1,8 @@ +nodes: + - id: cxx-node-rust-api + custom: + source: build/node_rust_api + inputs: + tick: dora/timer/millis/500 + outputs: + - pose diff --git a/examples/c++-ros2-dataflow/node-rust-api/main.cc b/examples/c++-ros2-dataflow/node-rust-api/main.cc new file mode 100644 index 00000000..4d8d0cb7 --- /dev/null +++ b/examples/c++-ros2-dataflow/node-rust-api/main.cc @@ -0,0 +1,84 @@ +#include "../build/dora-node-api.h" +#include "../build/dora-ros2-bindings.h" + +#include +#include +#include + +int main() +{ + std::cout << "HELLO FROM C++" << std::endl; + + auto dora_node = init_dora_node(); + auto merged_events = dora_events_into_combined(std::move(dora_node.events)); + + auto qos = qos_default(); + qos.durability = Ros2Durability::Volatile; + qos.liveliness = Ros2Liveliness::Automatic; + qos.reliable = true; + qos.max_blocking_time = 0.1; + + auto ros2_context = init_ros2_context(); + auto node = ros2_context->new_node("/ros2_demo", "turtle_teleop"); + auto vel_topic = node->create_topic_geometry_msgs_Twist("/turtle1", "cmd_vel", qos); + auto vel_publisher = node->create_publisher(vel_topic, qos); + auto pose_topic = node->create_topic_turtlesim_Pose("/turtle1", "pose", qos); + auto pose_subscription = node->create_subscription(pose_topic, qos, merged_events); + + std::random_device dev; + std::default_random_engine gen(dev()); + std::uniform_real_distribution<> dist(0., 1.); + + auto received_ticks = 0; + + for (int i = 0; i < 1000; i++) + { + auto event = merged_events.next(); + + if (event.is_dora()) + { + auto dora_event = downcast_dora(std::move(event)); + + auto ty = event_type(dora_event); + + if (ty == DoraEventType::AllInputsClosed) + { + break; + } + else if (ty == DoraEventType::Input) + { + auto input = event_as_input(std::move(dora_event)); + received_ticks += 1; + + std::cout << "Received input " << std::string(input.id) << std::endl; + + geometry_msgs::Twist twist = { + .linear = {.x = dist(gen) + 1, .y = 0, .z = 0}, + .angular = {.x = 0, .y = 0, .z = (dist(gen) - 0.5) * 5.0}}; + vel_publisher->publish(twist); + } + else + { + std::cerr << "Unknown event type " << static_cast(ty) << std::endl; + } + + if (received_ticks > 20) + { + break; + } + } + else if (pose_subscription->matches(event)) + { + auto pose = pose_subscription->downcast(std::move(event)); + std::cout << "Received pose x:" << pose.x << ", y:" << pose.y << std::endl; + } + else + { + std::cout << "received unexpected event" << std::endl; + } + } + + std::cout << "GOODBYE FROM C++ node (using Rust API)" << std::endl; + + return 0; +} diff --git a/examples/c++-ros2-dataflow/run.rs b/examples/c++-ros2-dataflow/run.rs new file mode 100644 index 00000000..4d5d242f --- /dev/null +++ b/examples/c++-ros2-dataflow/run.rs @@ -0,0 +1,165 @@ +use dora_tracing::set_up_tracing; +use eyre::{bail, Context}; +use std::{env::consts::EXE_SUFFIX, path::Path}; + +#[tokio::main] +async fn main() -> eyre::Result<()> { + set_up_tracing("c++-ros2-dataflow-exaple").wrap_err("failed to set up tracing")?; + + if cfg!(windows) { + tracing::error!( + "The c++ example does not work on Windows currently because of a linker error" + ); + return Ok(()); + } + + let root = Path::new(env!("CARGO_MANIFEST_DIR")); + let target = root.join("target"); + std::env::set_current_dir(root.join(file!()).parent().unwrap()) + .wrap_err("failed to set working dir")?; + + tokio::fs::create_dir_all("build").await?; + let build_dir = Path::new("build"); + + build_package("dora-node-api-cxx", &["ros2-bridge"]).await?; + let node_cxxbridge = target.join("cxxbridge").join("dora-node-api-cxx"); + tokio::fs::copy( + node_cxxbridge.join("dora-node-api.cc"), + build_dir.join("dora-node-api.cc"), + ) + .await?; + tokio::fs::copy( + node_cxxbridge.join("dora-node-api.h"), + build_dir.join("dora-node-api.h"), + ) + .await?; + tokio::fs::copy( + node_cxxbridge.join("dora-ros2-bindings.cc"), + build_dir.join("dora-ros2-bindings.cc"), + ) + .await?; + tokio::fs::copy( + node_cxxbridge.join("dora-ros2-bindings.h"), + build_dir.join("dora-ros2-bindings.h"), + ) + .await?; + + build_cxx_node( + root, + &[ + &dunce::canonicalize(Path::new("node-rust-api").join("main.cc"))?, + &dunce::canonicalize(build_dir.join("dora-ros2-bindings.cc"))?, + &dunce::canonicalize(build_dir.join("dora-node-api.cc"))?, + ], + "node_rust_api", + &["-l", "dora_node_api_cxx"], + ) + .await?; + + let dataflow = Path::new("dataflow.yml").to_owned(); + run_dataflow(&dataflow).await?; + + Ok(()) +} + +async fn build_package(package: &str, features: &[&str]) -> eyre::Result<()> { + let cargo = std::env::var("CARGO").unwrap(); + let mut cmd = tokio::process::Command::new(&cargo); + cmd.arg("build"); + cmd.arg("--package").arg(package); + if !features.is_empty() { + cmd.arg("--features").arg(features.join(",")); + } + if !cmd.status().await?.success() { + bail!("failed to compile {package}"); + }; + Ok(()) +} + +async fn build_cxx_node( + root: &Path, + paths: &[&Path], + out_name: &str, + args: &[&str], +) -> eyre::Result<()> { + let mut clang = tokio::process::Command::new("clang++"); + clang.args(paths); + clang.arg("-std=c++17"); + #[cfg(target_os = "linux")] + { + clang.arg("-l").arg("m"); + clang.arg("-l").arg("rt"); + clang.arg("-l").arg("dl"); + clang.arg("-pthread"); + } + #[cfg(target_os = "windows")] + { + clang.arg("-ladvapi32"); + clang.arg("-luserenv"); + clang.arg("-lkernel32"); + clang.arg("-lws2_32"); + clang.arg("-lbcrypt"); + clang.arg("-lncrypt"); + clang.arg("-lschannel"); + clang.arg("-lntdll"); + clang.arg("-liphlpapi"); + + clang.arg("-lcfgmgr32"); + clang.arg("-lcredui"); + clang.arg("-lcrypt32"); + clang.arg("-lcryptnet"); + clang.arg("-lfwpuclnt"); + clang.arg("-lgdi32"); + clang.arg("-lmsimg32"); + clang.arg("-lmswsock"); + clang.arg("-lole32"); + clang.arg("-lopengl32"); + clang.arg("-lsecur32"); + clang.arg("-lshell32"); + clang.arg("-lsynchronization"); + clang.arg("-luser32"); + clang.arg("-lwinspool"); + + clang.arg("-Wl,-nodefaultlib:libcmt"); + clang.arg("-D_DLL"); + clang.arg("-lmsvcrt"); + } + #[cfg(target_os = "macos")] + { + clang.arg("-framework").arg("CoreServices"); + clang.arg("-framework").arg("Security"); + clang.arg("-l").arg("System"); + clang.arg("-l").arg("resolv"); + clang.arg("-l").arg("pthread"); + clang.arg("-l").arg("c"); + clang.arg("-l").arg("m"); + } + clang.args(args); + clang.arg("-L").arg(root.join("target").join("debug")); + clang + .arg("--output") + .arg(Path::new("../build").join(format!("{out_name}{EXE_SUFFIX}"))); + if let Some(parent) = paths[0].parent() { + clang.current_dir(parent); + } + + if !clang.status().await?.success() { + bail!("failed to compile c++ node"); + }; + Ok(()) +} + +async fn run_dataflow(dataflow: &Path) -> eyre::Result<()> { + let cargo = std::env::var("CARGO").unwrap(); + let mut cmd = tokio::process::Command::new(&cargo); + cmd.arg("run"); + cmd.arg("--package").arg("dora-cli"); + cmd.arg("--") + .arg("daemon") + .arg("--run-dataflow") + .arg(dataflow); + if !cmd.status().await?.success() { + bail!("failed to run dataflow"); + }; + Ok(()) +} diff --git a/examples/rust-ros2-dataflow/node/src/main.rs b/examples/rust-ros2-dataflow/node/src/main.rs index 15c360b6..6d20472f 100644 --- a/examples/rust-ros2-dataflow/node/src/main.rs +++ b/examples/rust-ros2-dataflow/node/src/main.rs @@ -5,10 +5,12 @@ use dora_node_api::{ DoraNode, Event, }; use dora_ros2_bridge::{ - geometry_msgs::msg::{Twist, Vector3}, + messages::{ + geometry_msgs::msg::{Twist, Vector3}, + turtlesim::msg::Pose, + }, ros2_client::{self, ros2, NodeOptions}, rustdds::{self, policy}, - turtlesim::msg::Pose, }; use eyre::{eyre, Context}; diff --git a/libraries/extensions/ros2-bridge/Cargo.toml b/libraries/extensions/ros2-bridge/Cargo.toml index 7c55a1cb..20b2d03f 100644 --- a/libraries/extensions/ros2-bridge/Cargo.toml +++ b/libraries/extensions/ros2-bridge/Cargo.toml @@ -2,18 +2,18 @@ name = "dora-ros2-bridge" version = "0.1.0" edition = "2021" +links = "dora-ros2-bridge" # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html [features] default = ["generate-messages"] -generate-messages = ["dep:dora-ros2-bridge-msg-gen-macro"] +generate-messages = ["dep:dora-ros2-bridge-msg-gen", "dep:rust-format"] # enables examples that depend on a sourced ROS2 installation -ros2-examples = ["eyre", "tokio", "dora-daemon"] +ros2-examples = ["dep:eyre", "tokio", "dora-daemon"] [dependencies] array-init = "2.1.0" -dora-ros2-bridge-msg-gen-macro = { path = "msg-gen-macro", optional = true } serde = { version = "1.0.164", features = ["derive"] } serde-big-array = "0.5.1" widestring = "1.0.2" @@ -28,3 +28,9 @@ tracing-subscriber = "0.3.17" [dev-dependencies] rand = "0.8.5" futures = { version = "0.3.28", default-features = false } + +[build-dependencies] +dora-ros2-bridge-msg-gen = { workspace = true, optional = true } +rust-format = { version = "0.3.4", features = [ + "pretty_please", +], optional = true } diff --git a/libraries/extensions/ros2-bridge/build.rs b/libraries/extensions/ros2-bridge/build.rs new file mode 100644 index 00000000..dc3a701e --- /dev/null +++ b/libraries/extensions/ros2-bridge/build.rs @@ -0,0 +1,45 @@ +use std::path::PathBuf; + +#[cfg(not(feature = "generate-messages"))] +fn main() {} + +#[cfg(feature = "generate-messages")] +fn main() { + use rust_format::Formatter; + let paths = ament_prefix_paths(); + let generated = dora_ros2_bridge_msg_gen::gen(paths.as_slice(), false); + let generated_string = rust_format::PrettyPlease::default() + .format_tokens(generated) + .unwrap(); + let out_dir = PathBuf::from(std::env::var("OUT_DIR").unwrap()); + let target_file = out_dir.join("messages.rs"); + std::fs::write(&target_file, generated_string).unwrap(); + println!("cargo:rustc-env=MESSAGES_PATH={}", target_file.display()); + + #[cfg(feature = "cxx-bridge")] + let _build = cxx_build::bridge(&target_file); +} + +fn ament_prefix_paths() -> Vec { + let ament_prefix_path: String = match std::env::var("AMENT_PREFIX_PATH") { + Ok(path) => path, + Err(std::env::VarError::NotPresent) => { + println!("cargo:warning='AMENT_PREFIX_PATH not set'"); + String::new() + } + Err(std::env::VarError::NotUnicode(s)) => { + panic!( + "AMENT_PREFIX_PATH is not valid unicode: `{}`", + s.to_string_lossy() + ); + } + }; + println!("cargo:rerun-if-env-changed=AMENT_PREFIX_PATH"); + + let paths: Vec<_> = ament_prefix_path.split(':').map(PathBuf::from).collect(); + for path in &paths { + println!("cargo:rerun-if-changed={}", path.display()); + } + + paths +} diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/Cargo.toml b/libraries/extensions/ros2-bridge/msg-gen-macro/Cargo.toml deleted file mode 100644 index 434e42ab..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/Cargo.toml +++ /dev/null @@ -1,20 +0,0 @@ -[package] -name = "dora-ros2-bridge-msg-gen-macro" -version = "0.1.0" -edition = "2021" -authors = ["Yuma Hiramatsu "] -license = "Apache-2.0" - -[dependencies] -anyhow = "1.0" -heck = "0.3" -nom = "7" -proc-macro2 = "1.0" -quote = "1.0" -regex = "1" -syn = "1.0" -thiserror = "1.0" -dora-ros2-bridge-msg-gen = { path = "../msg-gen" } - -[lib] -proc-macro = true diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/build.rs b/libraries/extensions/ros2-bridge/msg-gen-macro/build.rs deleted file mode 100644 index 9d5bff75..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/build.rs +++ /dev/null @@ -1,25 +0,0 @@ -use std::path::Path; - -fn main() { - let ament_prefix_path = match std::env::var("AMENT_PREFIX_PATH") { - Ok(path) => path, - Err(std::env::VarError::NotPresent) => { - println!("cargo:warning='AMENT_PREFIX_PATH not set'"); - String::new() - } - Err(std::env::VarError::NotUnicode(s)) => { - panic!( - "AMENT_PREFIX_PATH is not valid unicode: `{}`", - s.to_string_lossy() - ); - } - }; - println!("cargo:rerun-if-env-changed=AMENT_PREFIX_PATH"); - - let paths = ament_prefix_path.split(':').map(Path::new); - for path in paths { - println!("cargo:rerun-if-changed={}", path.display()); - } - - println!("cargo:rustc-env=DETECTED_AMENT_PREFIX_PATH={ament_prefix_path}"); -} diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/src/lib.rs b/libraries/extensions/ros2-bridge/msg-gen-macro/src/lib.rs deleted file mode 100644 index ae6d02ab..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/src/lib.rs +++ /dev/null @@ -1,43 +0,0 @@ -// Based on https://github.com/rclrust/rclrust/tree/3a48dbb8f23a3d67d3031351da3ed236a354f039/rclrust-msg-gen - -#![warn( - rust_2018_idioms, - elided_lifetimes_in_paths, - clippy::all, - clippy::nursery -)] - -use std::path::Path; - -use dora_ros2_bridge_msg_gen::get_packages; -use proc_macro::TokenStream; -use quote::quote; - -#[proc_macro] -pub fn msg_include_all(_input: TokenStream) -> TokenStream { - let ament_prefix_path = std::env!("DETECTED_AMENT_PREFIX_PATH").trim(); - if ament_prefix_path.is_empty() { - quote! { - /// **No messages are available because the `AMENT_PREFIX_PATH` environment variable - /// was not set during build.** - pub const AMENT_PREFIX_PATH_NOT_SET: () = (); - } - .into() - } else { - let paths = ament_prefix_path - .split(':') - .map(Path::new) - .collect::>(); - - let packages = get_packages(&paths) - .unwrap() - .iter() - .map(|v| v.token_stream()) - .collect::>(); - - (quote! { - #(#packages)* - }) - .into() - } -} diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/action/Fibonacci.action b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/action/Fibonacci.action deleted file mode 100644 index b2591f28..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/action/Fibonacci.action +++ /dev/null @@ -1,8 +0,0 @@ -#goal definition -int32 order ---- -#result definition -int32[] sequence ---- -#feedback -int32[] sequence diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Arrays.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Arrays.msg deleted file mode 100644 index 4e2d280b..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Arrays.msg +++ /dev/null @@ -1,34 +0,0 @@ -# Arrays of different types -bool[3] bool_values -byte[3] byte_values -char[3] char_values -float32[3] float32_values -float64[3] float64_values -int8[3] int8_values -uint8[3] uint8_values -int16[3] int16_values -uint16[3] uint16_values -int32[3] int32_values -uint32[3] uint32_values -int64[3] int64_values -uint64[3] uint64_values -string[3] string_values -BasicTypes[3] basic_types_values -Constants[3] constants_values -Defaults[3] defaults_values -bool[3] bool_values_default [false, true, false] -byte[3] byte_values_default [0, 1, 255] -char[3] char_values_default [0, 1, 127] -float32[3] float32_values_default [1.125, 0.0, -1.125] -float64[3] float64_values_default [3.1415, 0.0, -3.1415] -int8[3] int8_values_default [0, 127, -128] -uint8[3] uint8_values_default [0, 1, 255] -int16[3] int16_values_default [0, 32767, -32768] -uint16[3] uint16_values_default [0, 1, 65535] -int32[3] int32_values_default [0, 2147483647, -2147483648] -uint32[3] uint32_values_default [0, 1, 4294967295] -int64[3] int64_values_default [0, 9223372036854775807, -9223372036854775808] -uint64[3] uint64_values_default [0, 1, 18446744073709551615] -string[3] string_values_default ["", "max value", "min value"] -# Regression test: check alignment of basic field after an array field is correct -int32 alignment_check diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/BasicTypes.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/BasicTypes.msg deleted file mode 100644 index eaf9b9ef..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/BasicTypes.msg +++ /dev/null @@ -1,13 +0,0 @@ -bool bool_value -byte byte_value -char char_value -float32 float32_value -float64 float64_value -int8 int8_value -uint8 uint8_value -int16 int16_value -uint16 uint16_value -int32 int32_value -uint32 uint32_value -int64 int64_value -uint64 uint64_value diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/BoundedSequences.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/BoundedSequences.msg deleted file mode 100644 index e369ffc3..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/BoundedSequences.msg +++ /dev/null @@ -1,34 +0,0 @@ -# Bounded sequences of different types -bool[<=3] bool_values -byte[<=3] byte_values -char[<=3] char_values -float32[<=3] float32_values -float64[<=3] float64_values -int8[<=3] int8_values -uint8[<=3] uint8_values -int16[<=3] int16_values -uint16[<=3] uint16_values -int32[<=3] int32_values -uint32[<=3] uint32_values -int64[<=3] int64_values -uint64[<=3] uint64_values -string[<=3] string_values -BasicTypes[<=3] basic_types_values -Constants[<=3] constants_values -Defaults[<=3] defaults_values -bool[<=3] bool_values_default [false, true, false] -byte[<=3] byte_values_default [0, 1, 255] -char[<=3] char_values_default [0, 1, 127] -float32[<=3] float32_values_default [1.125, 0.0, -1.125] -float64[<=3] float64_values_default [3.1415, 0.0, -3.1415] -int8[<=3] int8_values_default [0, 127, -128] -uint8[<=3] uint8_values_default [0, 1, 255] -int16[<=3] int16_values_default [0, 32767, -32768] -uint16[<=3] uint16_values_default [0, 1, 65535] -int32[<=3] int32_values_default [0, 2147483647, -2147483648] -uint32[<=3] uint32_values_default [0, 1, 4294967295] -int64[<=3] int64_values_default [0, 9223372036854775807, -9223372036854775808] -uint64[<=3] uint64_values_default [0, 1, 18446744073709551615] -string[<=3] string_values_default ["", "max value", "min value"] -# Regression test: check alignment of basic field after a sequence field is correct -int32 alignment_check diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Constants.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Constants.msg deleted file mode 100644 index c9968593..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Constants.msg +++ /dev/null @@ -1,13 +0,0 @@ -bool BOOL_CONST=true -byte BYTE_CONST=50 -char CHAR_CONST=100 -float32 FLOAT32_CONST=1.125 -float64 FLOAT64_CONST=1.125 -int8 INT8_CONST=-50 -uint8 UINT8_CONST=200 -int16 INT16_CONST=-1000 -uint16 UINT16_CONST=2000 -int32 INT32_CONST=-30000 -uint32 UINT32_CONST=60000 -int64 INT64_CONST=-40000000 -uint64 UINT64_CONST=50000000 diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Defaults.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Defaults.msg deleted file mode 100644 index eb4832f5..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Defaults.msg +++ /dev/null @@ -1,13 +0,0 @@ -bool bool_value true -byte byte_value 50 -char char_value 100 -float32 float32_value 1.125 -float64 float64_value 1.125 -int8 int8_value -50 -uint8 uint8_value 200 -int16 int16_value -1000 -uint16 uint16_value 2000 -int32 int32_value -30000 -uint32 uint32_value 60000 -int64 int64_value -40000000 -uint64 uint64_value 50000000 diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Empty.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Empty.msg deleted file mode 100644 index e69de29b..00000000 diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/MultiNested.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/MultiNested.msg deleted file mode 100644 index cf324dc8..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/MultiNested.msg +++ /dev/null @@ -1,10 +0,0 @@ -# Mulitple levels of nested messages -Arrays[3] array_of_arrays -BoundedSequences[3] array_of_bounded_sequences -UnboundedSequences[3] array_of_unbounded_sequences -Arrays[<=3] bounded_sequence_of_arrays -BoundedSequences[<=3] bounded_sequence_of_bounded_sequences -UnboundedSequences[<=3] bounded_sequence_of_unbounded_sequences -Arrays[] unbounded_sequence_of_arrays -BoundedSequences[] unbounded_sequence_of_bounded_sequences -UnboundedSequences[] unbounded_sequence_of_unbounded_sequences diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Nested.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Nested.msg deleted file mode 100644 index 9adb9a59..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Nested.msg +++ /dev/null @@ -1 +0,0 @@ -BasicTypes basic_types_value diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Strings.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Strings.msg deleted file mode 100644 index e7f5e45f..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/Strings.msg +++ /dev/null @@ -1,13 +0,0 @@ -string string_value -string string_value_default1 "Hello world!" -string string_value_default2 "Hello'world!" -string string_value_default3 'Hello"world!' -string string_value_default4 'Hello\'world!' -string string_value_default5 "Hello\"world!" -string STRING_CONST="Hello world!" -string<=22 bounded_string_value -string<=22 bounded_string_value_default1 "Hello world!" -string<=22 bounded_string_value_default2 "Hello'world!" -string<=22 bounded_string_value_default3 'Hello"world!' -string<=22 bounded_string_value_default4 'Hello\'world!' -string<=22 bounded_string_value_default5 "Hello\"world!" diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/UnboundedSequences.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/UnboundedSequences.msg deleted file mode 100644 index 36c02d96..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/UnboundedSequences.msg +++ /dev/null @@ -1,34 +0,0 @@ -# Unbounded sequences of different types -bool[] bool_values -byte[] byte_values -char[] char_values -float32[] float32_values -float64[] float64_values -int8[] int8_values -uint8[] uint8_values -int16[] int16_values -uint16[] uint16_values -int32[] int32_values -uint32[] uint32_values -int64[] int64_values -uint64[] uint64_values -string[] string_values -BasicTypes[] basic_types_values -Constants[] constants_values -Defaults[] defaults_values -bool[] bool_values_default [false, true, false] -byte[] byte_values_default [0, 1, 255] -char[] char_values_default [0, 1, 127] -float32[] float32_values_default [1.125, 0.0, -1.125] -float64[] float64_values_default [3.1415, 0.0, -3.1415] -int8[] int8_values_default [0, 127, -128] -uint8[] uint8_values_default [0, 1, 255] -int16[] int16_values_default [0, 32767, -32768] -uint16[] uint16_values_default [0, 1, 65535] -int32[] int32_values_default [0, 2147483647, -2147483648] -uint32[] uint32_values_default [0, 1, 4294967295] -int64[] int64_values_default [0, 9223372036854775807, -9223372036854775808] -uint64[] uint64_values_default [0, 1, 18446744073709551615] -string[] string_values_default ["", "max value", "min value"] -# Regression test: check alignment of basic field after a sequence field is correct -int32 alignment_check diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/WStrings.msg b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/WStrings.msg deleted file mode 100644 index 666d0ec6..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/msg/WStrings.msg +++ /dev/null @@ -1,10 +0,0 @@ -wstring wstring_value -wstring wstring_value_default1 "Hello world!" -wstring wstring_value_default2 "Hellö wörld!" -wstring wstring_value_default3 "ハローワールド" -#wstring WSTRING_CONST="Hello world!" -#wstring<=22 bounded_wstring_value -#wstring<=22 bounded_wstring_value_default1 "Hello world!" -wstring[3] array_of_wstrings -wstring[<=3] bounded_sequence_of_wstrings -wstring[] unbounded_sequence_of_wstrings diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/srv/Arrays.srv b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/srv/Arrays.srv deleted file mode 100644 index 4c96820a..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/srv/Arrays.srv +++ /dev/null @@ -1,63 +0,0 @@ -bool[3] bool_values -byte[3] byte_values -char[3] char_values -float32[3] float32_values -float64[3] float64_values -int8[3] int8_values -uint8[3] uint8_values -int16[3] int16_values -uint16[3] uint16_values -int32[3] int32_values -uint32[3] uint32_values -int64[3] int64_values -uint64[3] uint64_values -string[3] string_values -BasicTypes[3] basic_types_values -Constants[3] constants_values -Defaults[3] defaults_values -bool[3] bool_values_default [false, true, false] -byte[3] byte_values_default [0, 1, 255] -char[3] char_values_default [0, 1, 127] -float32[3] float32_values_default [1.125, 0.0, -1.125] -float64[3] float64_values_default [3.1415, 0.0, -3.1415] -int8[3] int8_values_default [0, 127, -128] -uint8[3] uint8_values_default [0, 1, 255] -int16[3] int16_values_default [0, 32767, -32768] -uint16[3] uint16_values_default [0, 1, 65535] -int32[3] int32_values_default [0, 2147483647, -2147483648] -uint32[3] uint32_values_default [0, 1, 4294967295] -int64[3] int64_values_default [0, 9223372036854775807, -9223372036854775808] -uint64[3] uint64_values_default [0, 1, 18446744073709551615] -string[3] string_values_default ["", "max value", "min value"] ---- -bool[3] bool_values -byte[3] byte_values -char[3] char_values -float32[3] float32_values -float64[3] float64_values -int8[3] int8_values -uint8[3] uint8_values -int16[3] int16_values -uint16[3] uint16_values -int32[3] int32_values -uint32[3] uint32_values -int64[3] int64_values -uint64[3] uint64_values -string[3] string_values -BasicTypes[3] basic_types_values -Constants[3] constants_values -Defaults[3] defaults_values -bool[3] bool_values_default [false, true, false] -byte[3] byte_values_default [0, 1, 255] -char[3] char_values_default [0, 1, 127] -float32[3] float32_values_default [1.125, 0.0, -1.125] -float64[3] float64_values_default [3.1415, 0.0, -3.1415] -int8[3] int8_values_default [0, 127, -128] -uint8[3] uint8_values_default [0, 1, 255] -int16[3] int16_values_default [0, 32767, -32768] -uint16[3] uint16_values_default [0, 1, 65535] -int32[3] int32_values_default [0, 2147483647, -2147483648] -uint32[3] uint32_values_default [0, 1, 4294967295] -int64[3] int64_values_default [0, 9223372036854775807, -9223372036854775808] -uint64[3] uint64_values_default [0, 1, 18446744073709551615] -string[3] string_values_default ["", "max value", "min value"] diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/srv/BasicTypes.srv b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/srv/BasicTypes.srv deleted file mode 100644 index d00cfd53..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/srv/BasicTypes.srv +++ /dev/null @@ -1,29 +0,0 @@ -bool bool_value -byte byte_value -char char_value -float32 float32_value -float64 float64_value -int8 int8_value -uint8 uint8_value -int16 int16_value -uint16 uint16_value -int32 int32_value -uint32 uint32_value -int64 int64_value -uint64 uint64_value -string string_value ---- -bool bool_value -byte byte_value -char char_value -float32 float32_value -float64 float64_value -int8 int8_value -uint8 uint8_value -int16 int16_value -uint16 uint16_value -int32 int32_value -uint32 uint32_value -int64 int64_value -uint64 uint64_value -string string_value diff --git a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/srv/Empty.srv b/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/srv/Empty.srv deleted file mode 100644 index ed97d539..00000000 --- a/libraries/extensions/ros2-bridge/msg-gen-macro/test_msgs/srv/Empty.srv +++ /dev/null @@ -1 +0,0 @@ ---- diff --git a/libraries/extensions/ros2-bridge/msg-gen/src/lib.rs b/libraries/extensions/ros2-bridge/msg-gen/src/lib.rs index 313db9b3..54e7155c 100644 --- a/libraries/extensions/ros2-bridge/msg-gen/src/lib.rs +++ b/libraries/extensions/ros2-bridge/msg-gen/src/lib.rs @@ -7,7 +7,226 @@ clippy::nursery )] +use std::path::Path; + +use quote::quote; + pub mod parser; pub mod types; pub use crate::parser::get_packages; + +pub fn gen

(paths: &[P], create_cxx_bridge: bool) -> proc_macro2::TokenStream +where + P: AsRef, +{ + let packages = get_packages(paths).unwrap(); + let mut shared_type_defs = Vec::new(); + let mut message_struct_impls = Vec::new(); + let mut message_topic_defs = Vec::new(); + let mut message_topic_impls = Vec::new(); + let mut aliases = Vec::new(); + for package in &packages { + for message in &package.messages { + let (def, imp) = message.struct_token_stream(&package.name, create_cxx_bridge); + shared_type_defs.push(def); + message_struct_impls.push(imp); + if create_cxx_bridge { + let (topic_def, topic_impl) = message.topic_def(&package.name); + message_topic_defs.push(topic_def); + message_topic_impls.push(topic_impl); + } + } + aliases.push(package.aliases_token_stream()); + } + + let (attributes, imports_and_functions, cxx_bridge_impls) = if create_cxx_bridge { + ( + quote! { #[cxx::bridge] }, + quote! { + #[allow(dead_code)] + extern "C++" { + type CombinedEvents = crate::ffi::CombinedEvents; + type CombinedEvent = crate::ffi::CombinedEvent; + } + + extern "Rust" { + type Ros2Context; + type Ros2Node; + fn init_ros2_context() -> Result>; + fn new_node(self: &Ros2Context, name_space: &str, base_name: &str) -> Result>; + fn qos_default() -> Ros2QosPolicies; + + #(#message_topic_defs)* + } + + #[derive(Debug, Clone)] + pub struct Ros2QosPolicies { + pub durability: Ros2Durability, + pub liveliness: Ros2Liveliness, + pub lease_duration: f64, + pub reliable: bool, + pub max_blocking_time: f64, + pub keep_all: bool, + pub keep_last: i32, + } + + /// DDS 2.2.3.4 DURABILITY + #[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] + pub enum Ros2Durability { + Volatile, + TransientLocal, + Transient, + Persistent, + } + + /// DDS 2.2.3.11 LIVELINESS + #[derive(Copy, Clone, Debug, PartialEq)] + pub enum Ros2Liveliness { + Automatic, + ManualByParticipant, + ManualByTopic, + } + }, + quote! { + struct Ros2Context(ros2_client::Context); + + fn init_ros2_context() -> eyre::Result> { + Ok(Box::new(Ros2Context(ros2_client::Context::new()?))) + } + + impl Ros2Context { + fn new_node(&self, name_space: &str, base_name: &str) -> eyre::Result> { + let name = ros2_client::NodeName::new(name_space, base_name).map_err(|e| eyre::eyre!(e))?; + let options = ros2_client::NodeOptions::new().enable_rosout(true); + let node = self.0.new_node(name, options)?; + Ok(Box::new(Ros2Node(node))) + } + } + + struct Ros2Node(ros2_client::Node); + + fn qos_default() -> ffi::Ros2QosPolicies { + ffi::Ros2QosPolicies::new(None, None, None, None, None, None, None) + } + + impl ffi::Ros2QosPolicies { + pub fn new( + durability: Option, + liveliness: Option, + reliable: Option, + keep_all: Option, + lease_duration: Option, + max_blocking_time: Option, + keep_last: Option, + ) -> Self { + Self { + durability: durability.unwrap_or(ffi::Ros2Durability::Volatile), + liveliness: liveliness.unwrap_or(ffi::Ros2Liveliness::Automatic), + lease_duration: lease_duration.unwrap_or(f64::INFINITY), + reliable: reliable.unwrap_or(false), + max_blocking_time: max_blocking_time.unwrap_or(0.0), + keep_all: keep_all.unwrap_or(false), + keep_last: keep_last.unwrap_or(1), + } + } + } + + impl From for rustdds::QosPolicies { + fn from(value: ffi::Ros2QosPolicies) -> Self { + rustdds::QosPolicyBuilder::new() + .durability(value.durability.into()) + .liveliness(value.liveliness.convert(value.lease_duration)) + .reliability(if value.reliable { + rustdds::policy::Reliability::Reliable { + max_blocking_time: rustdds::Duration::from_frac_seconds( + value.max_blocking_time, + ), + } + } else { + rustdds::policy::Reliability::BestEffort + }) + .history(if value.keep_all { + rustdds::policy::History::KeepAll + } else { + rustdds::policy::History::KeepLast { + depth: value.keep_last, + } + }) + .build() + } + } + + + + impl From for rustdds::policy::Durability { + fn from(value: ffi::Ros2Durability) -> Self { + match value { + ffi::Ros2Durability::Volatile => rustdds::policy::Durability::Volatile, + ffi::Ros2Durability::TransientLocal => rustdds::policy::Durability::TransientLocal, + ffi::Ros2Durability::Transient => rustdds::policy::Durability::Transient, + ffi::Ros2Durability::Persistent => rustdds::policy::Durability::Persistent, + _ => unreachable!(), // required because enums are represented as integers in bridge + } + } + } + + + impl ffi::Ros2Liveliness { + fn convert(self, lease_duration: f64) -> rustdds::policy::Liveliness { + let lease_duration = if lease_duration.is_infinite() { + rustdds::Duration::INFINITE + } else { + rustdds::Duration::from_frac_seconds(lease_duration) + }; + match self { + ffi::Ros2Liveliness::Automatic => rustdds::policy::Liveliness::Automatic { lease_duration }, + ffi::Ros2Liveliness::ManualByParticipant => { + rustdds::policy::Liveliness::ManualByParticipant { lease_duration } + } + ffi::Ros2Liveliness::ManualByTopic => rustdds::policy::Liveliness::ManualByTopic { lease_duration }, + _ => unreachable!(), // required because enums are represented as integers in bridge + } + } + } + }, + ) + } else { + ( + quote! {}, + quote! { + use serde::{Serialize, Deserialize}; + }, + quote! {}, + ) + }; + + quote! { + #attributes + mod ffi { + #imports_and_functions + + #[derive(Debug, Default, Clone, PartialEq, Eq, Serialize, Deserialize)] + pub struct U16String { + chars: Vec, + } + + #(#shared_type_defs)* + } + + + impl crate::_core::InternalDefault for ffi::U16String { + fn _default() -> Self { + Default::default() + } + } + + #(#message_struct_impls)* + + #cxx_bridge_impls + #(#message_topic_impls)* + + + #(#aliases)* + } +} diff --git a/libraries/extensions/ros2-bridge/msg-gen/src/parser/package.rs b/libraries/extensions/ros2-bridge/msg-gen/src/parser/package.rs index 3c3f7b9f..12e7da5c 100644 --- a/libraries/extensions/ros2-bridge/msg-gen/src/parser/package.rs +++ b/libraries/extensions/ros2-bridge/msg-gen/src/parser/package.rs @@ -92,10 +92,13 @@ fn get_ros_msgs_each_package>(root_dir: P) -> Result Ok(packages) } -pub fn get_packages(paths: &[&Path]) -> Result> { +pub fn get_packages

(paths: &[P]) -> Result> +where + P: AsRef, +{ let mut packages = paths .iter() - .map(|&path| get_ros_msgs_each_package(path)) + .map(get_ros_msgs_each_package) .collect::>>()? .into_iter() .flatten() diff --git a/libraries/extensions/ros2-bridge/msg-gen/src/types/message.rs b/libraries/extensions/ros2-bridge/msg-gen/src/types/message.rs index 72dc1337..b8d81173 100644 --- a/libraries/extensions/ros2-bridge/msg-gen/src/types/message.rs +++ b/libraries/extensions/ros2-bridge/msg-gen/src/types/message.rs @@ -1,5 +1,6 @@ use heck::SnakeCase; use quote::{format_ident, quote, ToTokens}; +use syn::Ident; use super::{primitives::*, sequences::Array, ConstantType, MemberType}; @@ -122,9 +123,215 @@ pub struct Message { } impl Message { - pub fn token_stream_with_mod(&self) -> impl ToTokens { + pub fn struct_token_stream( + &self, + package_name: &str, + gen_cxx_bridge: bool, + ) -> (impl ToTokens, impl ToTokens) { + let cxx_name = format_ident!("{}", self.name); + let struct_raw_name = format_ident!("{package_name}__{}", self.name); + + let rust_type_def_inner = self.members.iter().map(|m| m.rust_type_def(&self.package)); + let constants_def_inner = self.constants.iter().map(|c| c.token_stream()); + let rust_type_default_inner = self.members.iter().map(|m| m.default_value()); + + let attributes = if gen_cxx_bridge { + quote! { + #[namespace = #package_name] + #[cxx_name = #cxx_name] + } + } else { + quote! {} + }; + + if self.members.is_empty() { + (quote! {}, quote! {}) + } else { + let def = quote! { + #[allow(non_camel_case_types)] + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] + #attributes + pub struct #struct_raw_name { + #(#rust_type_def_inner)* + } + }; + let impls = quote! { + impl crate::_core::InternalDefault for ffi::#struct_raw_name { + fn _default() -> Self { + Self { + #(#rust_type_default_inner)* + } + } + } + + impl std::default::Default for ffi::#struct_raw_name { + #[inline] + fn default() -> Self { + crate::_core::InternalDefault::_default() + } + } + }; + + (def, impls) + } + } + + pub fn topic_def(&self, package_name: &str) -> (impl ToTokens, impl ToTokens) { + if self.members.is_empty() { + return (quote! {}, quote! {}); + }; + + let topic_name = format_ident!("Topic__{package_name}__{}", self.name); + let cxx_topic_name = format_ident!("Topic_{}", self.name); + let create_topic = format_ident!("new__Topic__{package_name}__{}", self.name); + let cxx_create_topic = format!("create_topic_{package_name}_{}", self.name); + + let publisher_name = format_ident!("Publisher__{package_name}__{}", self.name); + let cxx_publisher_name = format_ident!("Publisher_{}", self.name); + let create_publisher = format_ident!("new__Publisher__{package_name}__{}", self.name); + let cxx_create_publisher = format_ident!("create_publisher"); + + let struct_raw_name = format_ident!("{package_name}__{}", self.name); + let struct_raw_name_str = struct_raw_name.to_string(); + let self_name = &self.name; + + let publish = format_ident!("publish__{package_name}__{}", self.name); + let cxx_publish = format_ident!("publish"); + + let subscription_name = format_ident!("Subscription__{package_name}__{}", self.name); + let subscription_name_str = subscription_name.to_string(); + let cxx_subscription_name = format_ident!("Subscription_{}", self.name); + let create_subscription = format_ident!("new__Subscription__{package_name}__{}", self.name); + let cxx_create_subscription = format_ident!("create_subscription"); + + let matches = format_ident!("matches__{package_name}__{}", self.name); + let cxx_matches = format_ident!("matches"); + let downcast = format_ident!("downcast__{package_name}__{}", self.name); + let cxx_downcast = format_ident!("downcast"); + + let def = quote! { + #[namespace = #package_name] + #[cxx_name = #cxx_topic_name] + type #topic_name; + #[cxx_name = #cxx_create_topic] + fn #create_topic(self: &Ros2Node, name_space: &str, base_name: &str, qos: Ros2QosPolicies) -> Result>; + #[cxx_name = #cxx_create_publisher] + fn #create_publisher(self: &mut Ros2Node, topic: &Box<#topic_name>, qos: Ros2QosPolicies) -> Result>; + #[cxx_name = #cxx_create_subscription] + fn #create_subscription(self: &mut Ros2Node, topic: &Box<#topic_name>, qos: Ros2QosPolicies, events: &mut CombinedEvents) -> Result>; + + #[namespace = #package_name] + #[cxx_name = #cxx_publisher_name] + type #publisher_name; + #[namespace = #package_name] + #[cxx_name = #cxx_publish] + fn #publish(self: &mut #publisher_name, message: #struct_raw_name) -> Result<()>; + + #[namespace = #package_name] + #[cxx_name = #cxx_subscription_name] + type #subscription_name; + + #[namespace = #package_name] + #[cxx_name = #cxx_matches] + fn #matches(self: &#subscription_name, event: &CombinedEvent) -> bool; + #[namespace = #package_name] + #[cxx_name = #cxx_downcast] + fn #downcast(self: &#subscription_name, event: CombinedEvent) -> Result<#struct_raw_name>; + }; + let imp = quote! { + #[allow(non_camel_case_types)] + pub struct #topic_name(rustdds::Topic); + + impl Ros2Node { + #[allow(non_snake_case)] + pub fn #create_topic(&self, name_space: &str, base_name: &str, qos: ffi::Ros2QosPolicies) -> eyre::Result> { + let name = ros2_client::Name::new(name_space, base_name).map_err(|e| eyre::eyre!(e))?; + let type_name = ros2_client::MessageTypeName::new(#package_name, #self_name); + let topic = self.0.create_topic(&name, type_name, &qos.into())?; + Ok(Box::new(#topic_name(topic))) + } + + #[allow(non_snake_case)] + pub fn #create_publisher(&mut self, topic: &Box<#topic_name>, qos: ffi::Ros2QosPolicies) -> eyre::Result> { + let publisher = self.0.create_publisher(&topic.0, Some(qos.into()))?; + Ok(Box::new(#publisher_name(publisher))) + } + + #[allow(non_snake_case)] + pub fn #create_subscription(&mut self, topic: &Box<#topic_name>, qos: ffi::Ros2QosPolicies, events: &mut crate::ffi::CombinedEvents) -> eyre::Result> { + let subscription = self.0.create_subscription::(&topic.0, Some(qos.into()))?; + let stream = futures_lite::stream::unfold(subscription, |sub| async { + let item = sub.async_take().await; + let item_boxed: Box = Box::new(item); + Some((item_boxed, sub)) + }); + let id = events.events.merge(Box::pin(stream)); + + Ok(Box::new(#subscription_name { id })) + } + } + + #[allow(non_camel_case_types)] + pub struct #publisher_name(ros2_client::Publisher); + + impl #publisher_name { + #[allow(non_snake_case)] + fn #publish(&mut self, message: ffi::#struct_raw_name) -> eyre::Result<()> { + use eyre::Context; + self.0.publish(message).context("publish failed") + } + } + + #[allow(non_camel_case_types)] + pub struct #subscription_name { + id: u32, + } + + impl #subscription_name { + #[allow(non_snake_case)] + fn #matches(&self, event: &crate::ffi::CombinedEvent) -> bool { + match &event.event.as_ref().0 { + Some(crate::MergedEvent::External(event)) if event.id == self.id => true, + _ => false + } + } + #[allow(non_snake_case)] + fn #downcast(&self, event: crate::ffi::CombinedEvent) -> eyre::Result { + use eyre::WrapErr; + + match (*event.event).0 { + Some(crate::MergedEvent::External(event)) if event.id == self.id => { + let result = event.event.downcast::>() + .map_err(|_| eyre::eyre!("downcast to {} failed", #struct_raw_name_str))?; + + let (data, _info) = result.with_context(|| format!("failed to receive {} event", #subscription_name_str))?; + Ok(data) + }, + _ => eyre::bail!("not a {} event", #subscription_name_str), + } + } + } + }; + (def, imp) + } + + pub fn alias_token_stream(&self, package_name: &Ident) -> impl ToTokens { + let cxx_name = format_ident!("{}", self.name); + let struct_raw_name = format_ident!("{package_name}__{}", self.name); + + if self.members.is_empty() { + quote! {} + } else { + quote! { + pub use super::super::ffi::#struct_raw_name as #cxx_name; + } + } + } + + pub fn token_stream_with_mod(&self, gen_cxx_bridge: bool) -> impl ToTokens { let mod_name = format_ident!("_{}", self.name.to_snake_case()); - let inner = self.token_stream(); + let inner = self.token_stream_args(gen_cxx_bridge); + quote! { pub use #mod_name::*; mod #mod_name { @@ -134,6 +341,10 @@ impl Message { } pub fn token_stream(&self) -> impl ToTokens { + self.token_stream_args(false) + } + + pub fn token_stream_args(&self, gen_cxx_bridge: bool) -> impl ToTokens { let rust_type = format_ident!("{}", self.name); let raw_type = format_ident!("{}_Raw", self.name); let raw_ref_type = format_ident!("{}_RawRef", self.name); @@ -144,6 +355,13 @@ impl Message { self.members.clone() }; + let attributes = if gen_cxx_bridge { + let namespace = &self.name; + quote! { #[cxx::bridge(namespace = #namespace)] } + } else { + quote! {} + }; + let rust_type_def_inner = self.members.iter().map(|m| m.rust_type_def(&self.package)); let constants_def_inner = self.constants.iter().map(|c| c.token_stream()); let rust_type_default_inner = self.members.iter().map(|m| m.default_value()); @@ -178,10 +396,15 @@ impl Message { FFIToRust as _FFIToRust, }; - #[allow(non_camel_case_types)] - #[derive(std::fmt::Debug, std::clone::Clone, std::cmp::PartialEq, serde::Serialize, serde::Deserialize)] - pub struct #rust_type { - #(#rust_type_def_inner)* + pub use self::t::#rust_type; + + #attributes + mod t { + #[allow(non_camel_case_types)] + #[derive(std::fmt::Debug, std::clone::Clone, std::cmp::PartialEq, serde::Serialize, serde::Deserialize)] + pub struct #rust_type { + #(#rust_type_def_inner)* + } } impl #rust_type { diff --git a/libraries/extensions/ros2-bridge/msg-gen/src/types/package.rs b/libraries/extensions/ros2-bridge/msg-gen/src/types/package.rs index 3d5c4944..ab0ccb3c 100644 --- a/libraries/extensions/ros2-bridge/msg-gen/src/types/package.rs +++ b/libraries/extensions/ros2-bridge/msg-gen/src/types/package.rs @@ -26,13 +26,55 @@ impl Package { self.messages.is_empty() && self.services.is_empty() && self.actions.is_empty() } - fn messages_block(&self) -> impl ToTokens { + pub fn message_structs(&self, gen_cxx_bridge: bool) -> (impl ToTokens, impl ToTokens) { + if self.messages.is_empty() { + // empty msg + (quote! {}, quote! {}) + } else { + let items = self + .messages + .iter() + .map(|v| v.struct_token_stream(&self.name, gen_cxx_bridge)); + let defs = items.clone().map(|(def, _)| def); + let impls = items.clone().map(|(_, im)| im); + let def_tokens = quote! { + #(#defs)* + }; + let impl_tokens = quote! { + #(#impls)* + }; + (def_tokens, impl_tokens) + } + } + + fn message_aliases(&self, package_name: &Ident) -> impl ToTokens { if self.messages.is_empty() { quote! { // empty msg } } else { - let items = self.messages.iter().map(|v| v.token_stream_with_mod()); + let items = self + .messages + .iter() + .map(|v| v.alias_token_stream(package_name)); + quote! { + pub mod msg { + #(#items)* + } + } + } + } + + fn messages_block(&self, gen_cxx_bridge: bool) -> impl ToTokens { + if self.messages.is_empty() { + quote! { + // empty msg + } + } else { + let items = self + .messages + .iter() + .map(|v| v.token_stream_with_mod(gen_cxx_bridge)); quote! { pub mod msg { #(#items)* @@ -71,9 +113,20 @@ impl Package { } } - pub fn token_stream(&self) -> impl ToTokens { + pub fn aliases_token_stream(&self) -> impl ToTokens { + let package_name = Ident::new(&self.name, Span::call_site()); + let aliases = self.message_aliases(&package_name); + + quote! { + pub mod #package_name { + #aliases + } + } + } + + pub fn token_stream(&self, gen_cxx_bridge: bool) -> impl ToTokens { let name = Ident::new(&self.name, Span::call_site()); - let messages_block = self.messages_block(); + let messages_block = self.messages_block(gen_cxx_bridge); let services_block = self.services_block(); let actions_block = self.actions_block(); diff --git a/libraries/extensions/ros2-bridge/msg-gen/src/types/primitives.rs b/libraries/extensions/ros2-bridge/msg-gen/src/types/primitives.rs index abf2747b..56d00149 100644 --- a/libraries/extensions/ros2-bridge/msg-gen/src/types/primitives.rs +++ b/libraries/extensions/ros2-bridge/msg-gen/src/types/primitives.rs @@ -131,9 +131,9 @@ pub struct NamedType(pub String); impl NamedType { fn type_tokens(&self, package: &str) -> impl ToTokens { let package = Ident::new(package, Span::call_site()); - let namespace = Ident::new("msg", Span::call_site()); let name = Ident::new(&self.0, Span::call_site()); - quote! { crate::#package::#namespace::#name } + let ident = format_ident!("{package}__{name}"); + quote! { #ident } } fn raw_type_tokens(&self, package: &str) -> impl ToTokens { @@ -173,9 +173,9 @@ pub struct NamespacedType { impl NamespacedType { fn type_tokens(&self) -> impl ToTokens { let package = Ident::new(&self.package, Span::call_site()); - let namespace = Ident::new(&self.namespace, Span::call_site()); let name = Ident::new(&self.name, Span::call_site()); - quote! { crate::#package::#namespace::#name } + let ident = format_ident!("{package}__{name}"); + quote! { #ident } } fn raw_type_tokens(&self) -> impl ToTokens { @@ -215,9 +215,9 @@ impl GenericString { fn type_tokens(self) -> impl ToTokens { if self.is_wide() { - quote! { crate::_core::string::U16String } + quote! { U16String } } else { - quote! { ::std::string::String } + quote! { String } } } diff --git a/libraries/extensions/ros2-bridge/msg-gen/src/types/sequences.rs b/libraries/extensions/ros2-bridge/msg-gen/src/types/sequences.rs index 13a5c5f7..2c4c1872 100644 --- a/libraries/extensions/ros2-bridge/msg-gen/src/types/sequences.rs +++ b/libraries/extensions/ros2-bridge/msg-gen/src/types/sequences.rs @@ -41,7 +41,7 @@ pub struct Sequence { impl Sequence { pub fn type_tokens(&self, package: &str) -> impl ToTokens { let inner_type = self.value_type.type_tokens(package); - quote! { std::vec::Vec<#inner_type> } + quote! { Vec<#inner_type> } } pub fn raw_type_tokens(&self, package: &str) -> impl ToTokens { @@ -72,7 +72,7 @@ pub struct BoundedSequence { impl BoundedSequence { pub fn type_tokens(&self, package: &str) -> impl ToTokens { let inner_type = self.value_type.type_tokens(package); - quote! { std::vec::Vec<#inner_type> } + quote! { Vec<#inner_type> } } pub fn raw_type_tokens(&self, package: &str) -> impl ToTokens { diff --git a/libraries/extensions/ros2-bridge/src/lib.rs b/libraries/extensions/ros2-bridge/src/lib.rs index e2053e0b..da41c521 100644 --- a/libraries/extensions/ros2-bridge/src/lib.rs +++ b/libraries/extensions/ros2-bridge/src/lib.rs @@ -2,6 +2,8 @@ pub use ros2_client; pub use rustdds; #[cfg(feature = "generate-messages")] -dora_ros2_bridge_msg_gen_macro::msg_include_all!(); +pub mod messages { + include!(env!("MESSAGES_PATH")); +} pub mod _core; diff --git a/rust-toolchain.toml b/rust-toolchain.toml index e724aeb2..0a2102d4 100644 --- a/rust-toolchain.toml +++ b/rust-toolchain.toml @@ -1,3 +1,3 @@ [toolchain] -channel = "1.72" +channel = "1.76" components = ["rustfmt", "clippy"]