#![cfg_attr(docsrs, feature(doc_cfg))]
#![warn(missing_docs)]
#![warn(rustdoc::missing_crate_level_docs)]
#![warn(rustdoc::unescaped_backticks)]
use autd3_core::link::{AsyncLink, LinkError, RxMessage, TxMessage};
use autd3_protobuf::*;
use std::net::SocketAddr;
struct SimulatorInner {
client: simulator_client::SimulatorClient<tonic::transport::Channel>,
last_geometry_version: usize,
}
impl SimulatorInner {
async fn open(
addr: &SocketAddr,
geometry: &autd3_core::geometry::Geometry,
) -> Result<SimulatorInner, LinkError> {
tracing::info!("Connecting to simulator@{}", addr);
let conn = tonic::transport::Endpoint::new(format!("http://{}", addr))
.map_err(AUTDProtoBufError::from)?
.connect()
.await
.map_err(AUTDProtoBufError::from)?;
let mut client = simulator_client::SimulatorClient::new(conn);
client
.config_geomety(Geometry::from(geometry))
.await
.map_err(|e| {
tracing::error!("Failed to configure simulator geometry: {}", e);
AUTDProtoBufError::SendError("Failed to initialize simulator".to_string())
})?;
Ok(Self {
client,
last_geometry_version: geometry.version(),
})
}
async fn close(&mut self) -> Result<(), LinkError> {
self.client
.close(CloseRequest {})
.await
.map_err(AUTDProtoBufError::from)?;
Ok(())
}
async fn update(&mut self, geometry: &autd3_core::geometry::Geometry) -> Result<(), LinkError> {
if self.last_geometry_version == geometry.version() {
return Ok(());
}
self.last_geometry_version = geometry.version();
self.client
.update_geomety(Geometry::from(geometry))
.await
.map_err(|e| {
tracing::error!("Failed to update geometry: {}", e);
AUTDProtoBufError::SendError("Failed to update geometry".to_string())
})?;
Ok(())
}
async fn send(&mut self, tx: &[TxMessage]) -> Result<(), LinkError> {
self.client
.send_data(TxRawData::from(tx))
.await
.map_err(AUTDProtoBufError::from)?;
Ok(())
}
async fn receive(&mut self, rx: &mut [RxMessage]) -> Result<bool, LinkError> {
let rx_ = Vec::<RxMessage>::from_msg(
self.client
.read_data(ReadRequest {})
.await
.map_err(AUTDProtoBufError::from)?
.into_inner(),
)?;
if rx.len() == rx_.len() {
rx.copy_from_slice(&rx_);
Ok(true)
} else {
Ok(false)
}
}
}
pub struct Simulator {
addr: SocketAddr,
inner: Option<SimulatorInner>,
#[cfg(feature = "blocking")]
runtime: Option<tokio::runtime::Runtime>,
}
impl Simulator {
#[must_use]
pub const fn new(addr: SocketAddr) -> Simulator {
Simulator {
addr,
inner: None,
#[cfg(feature = "blocking")]
runtime: None,
}
}
}
#[cfg_attr(feature = "async-trait", autd3_core::async_trait)]
impl AsyncLink for Simulator {
async fn open(&mut self, geometry: &autd3_core::geometry::Geometry) -> Result<(), LinkError> {
self.inner = Some(SimulatorInner::open(&self.addr, geometry).await?);
Ok(())
}
async fn close(&mut self) -> Result<(), LinkError> {
if let Some(mut inner) = self.inner.take() {
inner.close().await?;
}
Ok(())
}
async fn update(&mut self, geometry: &autd3_core::geometry::Geometry) -> Result<(), LinkError> {
if let Some(inner) = self.inner.as_mut() {
inner.update(geometry).await?;
Ok(())
} else {
Err(LinkError::new("Link is closed"))
}
}
async fn send(&mut self, tx: &[TxMessage]) -> Result<(), LinkError> {
if let Some(inner) = self.inner.as_mut() {
inner.send(tx).await?;
Ok(())
} else {
Err(LinkError::new("Link is closed"))
}
}
async fn receive(&mut self, rx: &mut [RxMessage]) -> Result<(), LinkError> {
if let Some(inner) = self.inner.as_mut() {
inner.receive(rx).await?;
Ok(())
} else {
Err(LinkError::new("Link is closed"))
}
}
fn is_open(&self) -> bool {
self.inner.is_some()
}
}
#[cfg(feature = "blocking")]
use autd3_core::link::Link;
#[cfg_attr(docsrs, doc(cfg(feature = "blocking")))]
#[cfg(feature = "blocking")]
impl Link for Simulator {
fn open(&mut self, geometry: &autd3_core::derive::Geometry) -> Result<(), LinkError> {
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.expect("Failed to create runtime");
runtime.block_on(<Self as AsyncLink>::open(self, geometry))?;
self.runtime = Some(runtime);
Ok(())
}
fn close(&mut self) -> Result<(), LinkError> {
self.runtime
.as_ref()
.map_or(Err(LinkError::new("Link is closed")), |runtime| {
runtime.block_on(async {
if let Some(mut inner) = self.inner.take() {
inner.close().await?;
}
Ok(())
})
})
}
fn update(&mut self, geometry: &autd3_core::geometry::Geometry) -> Result<(), LinkError> {
self.runtime
.as_ref()
.map_or(Err(LinkError::new("Link is closed")), |runtime| {
runtime.block_on(async {
if let Some(inner) = self.inner.as_mut() {
inner.update(geometry).await?;
}
Ok(())
})
})
}
fn send(&mut self, tx: &[TxMessage]) -> Result<(), LinkError> {
self.runtime
.as_ref()
.map_or(Err(LinkError::new("Link is closed")), |runtime| {
runtime.block_on(async {
if let Some(inner) = self.inner.as_mut() {
inner.send(tx).await?;
Ok(())
} else {
Err(LinkError::new("Link is closed"))
}
})
})
}
fn receive(&mut self, rx: &mut [RxMessage]) -> Result<(), LinkError> {
self.runtime
.as_ref()
.map_or(Err(LinkError::new("Link is closed")), |runtime| {
runtime.block_on(async {
if let Some(inner) = self.inner.as_mut() {
inner.receive(rx).await?;
Ok(())
} else {
Err(LinkError::new("Link is closed"))
}
})
})
}
fn is_open(&self) -> bool {
self.runtime.is_some() && self.inner.is_some()
}
}