#[path = "../gen/gen_types.rs"] mod gen_types;
#[path = "../gen/gen_ffi.rs"] mod gen_ffi;
pub use gen_types::*; use gen_ffi::ffi; pub use gen_ffi::ByteHandler;
pub mod rpc;
pub fn init_dds(domain_id: i32, network_interface: &str, config_file: &str) {
ffi::boot_dds(domain_id, network_interface, config_file);
}
use std::sync::Arc;
use cxx::UniquePtr;
use cxx::memory::UniquePtrTarget;
use tokio::sync::watch;
use serde::de::DeserializeOwned;
use serde::Serialize;
pub trait DdsType {
fn dds_name() -> &'static str;
}
pub struct Subscriber {
id: i32,
}
impl Drop for Subscriber {
fn drop(&mut self) { if self.id > 0 { ffi::unsubscribe(self.id); } }
}
pub fn subscribe<D: DdsType + DeserializeOwned + Send + Sync + 'static>(topic: &str) -> (Subscriber, watch::Receiver<Option<D>>) {
let (tx, rx) = watch::channel(None);
let id = ffi::subscribe_any(topic, D::dds_name(), Box::new(ByteHandler::new(move |bytes: Vec<u8>| {
if let Ok(data) = cdr::deserialize::<D>(&bytes[..]) {
let _ = tx.send_replace(Some(data));
}
})));
(Subscriber { id }, rx)
}
pub trait TopicPublish {
type PubFfi: Send + Sync + UniquePtrTarget;
fn new_publisher(topic: &str) -> anyhow::Result<Arc<UniquePtr<Self::PubFfi>>>;
fn publish_ffi(pub_: &UniquePtr<Self::PubFfi>, bytes: Vec<u8>);
}
pub struct Publisher<T: TopicPublish>
where T::PubFfi: UniquePtrTarget
{
inner: Arc<UniquePtr<T::PubFfi>>,
}
impl<T: TopicPublish> Publisher<T>
where T::PubFfi: UniquePtrTarget
{
pub fn new(topic: &str) -> anyhow::Result<Self> {
Ok(Self { inner: T::new_publisher(topic)? })
}
pub fn publish(&self, d: &T) where T: Serialize {
if let Ok(bytes) = cdr::serialize::<_, _, cdr::CdrLe>(d, cdr::Infinite) {
T::publish_ffi(&self.inner, bytes);
}
}
}
macro_rules! topic_publish {
($rust:ty, $ffi_pub:ident, $new_fn:ident, $publish_fn:ident) => {
unsafe impl Send for ffi::$ffi_pub {}
unsafe impl Sync for ffi::$ffi_pub {}
impl TopicPublish for $rust {
type PubFfi = ffi::$ffi_pub;
fn new_publisher(topic: &str) -> anyhow::Result<Arc<UniquePtr<Self::PubFfi>>> {
let inner = ffi::$new_fn(topic);
if inner.is_null() { anyhow::bail!("创建 {} 失败", stringify!($rust)); }
Ok(Arc::new(inner))
}
fn publish_ffi(pub_: &UniquePtr<Self::PubFfi>, bytes: Vec<u8>) {
ffi::$publish_fn(pub_, bytes);
}
}
};
}
include!("../gen/gen_topic_impl.rs");
pub struct RpcRequestHandler {
callback: Box<dyn Fn(&str) -> String + Send + 'static>,
}
impl RpcRequestHandler {
pub fn new(callback: impl Fn(&str) -> String + Send + 'static) -> Self {
Self { callback: Box::new(callback) }
}
fn handle(&self, request: &str) -> String { (self.callback)(request) }
}
#[cxx::bridge(namespace = "unitree")]
pub mod rpc_ffi {
extern "Rust" {
type RpcRequestHandler;
fn handle(self: &RpcRequestHandler, request: &str) -> String;
}
unsafe extern "C++" {
include!("rpc_bridge.h");
type RpcClient;
fn new_rpc_client(service: &str) -> UniquePtr<RpcClient>;
fn register_api(self: &RpcClient, api_id: i32);
fn call(self: &RpcClient, api_id: i32, request: &str) -> String;
type RpcServer;
fn new_rpc_server(service: &str) -> UniquePtr<RpcServer>;
fn start(self: &RpcServer);
fn register_handler(self: &RpcServer, api_id: i32, handler: Box<RpcRequestHandler>);
}
}