use alloc::vec;
use alloc::vec::Vec;
#[allow(unused_imports)]
use log::{debug, error, info, trace, warn};
use crate::channel::{ChannelActor, ReaderWriterChannel, ReaderWriterChannelIo};
use crate::client::{ChannelConfig, RpcClientConfig};
use crate::io::{Reader, Writer};
pub trait AsyncDelay {
fn delay() -> impl Future<Output = ()>;
}
pub struct AsyncRpcClient<'a, R: Reader, W: Writer, D: AsyncDelay> {
io: ReaderWriterChannelIo<'a, R, W>,
cmd_ch_config: ChannelConfig,
rsp_ch_config: ChannelConfig,
_delay: core::marker::PhantomData<D>,
}
impl<'a, R: Reader, W: Writer, D: AsyncDelay> AsyncRpcClient<'a, R, W, D> {
pub fn new(reader: &'a mut R, writer: &'a mut W, config: RpcClientConfig) -> Self {
let (cmd_ch_config, rsp_ch_config) = Self::get_channel_configs(config);
Self {
io: ReaderWriterChannelIo::new(reader, writer),
cmd_ch_config,
rsp_ch_config,
_delay: core::marker::PhantomData,
}
}
pub async fn request(&mut self, command: &[u8]) -> Result<Vec<u8>, crate::Error> {
debug!("Starting RPC request ({} bytes)", command.len());
let mut cmd_ch = self.cmd_channel().await?;
cmd_ch.publish_bytes(command).await?;
debug!("Command sent to target");
let mut rsp_ch = self.rsp_channel().await?;
let response_size = loop {
if let Some(size) = rsp_ch.data_available().await? {
debug!("Response available ({} bytes)", size);
break size;
}
D::delay().await;
};
let mut response_buf = vec![0u8; response_size];
let received_size = rsp_ch.consume_bytes(&mut response_buf).await?;
if received_size != response_size {
warn!(
"Expected {} bytes, received {} bytes",
response_size, received_size
);
response_buf.truncate(received_size);
}
debug!("RPC request completed ({} bytes received)", received_size);
Ok(response_buf)
}
fn get_channel_configs(config: RpcClientConfig) -> (ChannelConfig, ChannelConfig) {
match config {
RpcClientConfig::Direct {
cmd_ch_ptr,
cmd_ch_size,
rsp_ch_ptr,
rsp_ch_size,
} => (
ChannelConfig::Direct {
ptr: cmd_ch_ptr,
size: cmd_ch_size,
},
ChannelConfig::Direct {
ptr: rsp_ch_ptr,
size: rsp_ch_size,
},
),
RpcClientConfig::FromTarget {
cmd_ch_ptr,
rsp_ch_ptr,
} => (
ChannelConfig::FromTarget { ptr: cmd_ch_ptr },
ChannelConfig::FromTarget { ptr: rsp_ch_ptr },
),
}
}
async fn cmd_channel<'method>(
&'method mut self,
) -> Result<ReaderWriterChannel<'method, 'a, R, W>, crate::Error> {
match self.cmd_ch_config {
ChannelConfig::Direct { ptr, size } => {
ReaderWriterChannel::new(&mut self.io, ChannelActor::Producer, ptr, size).await
}
ChannelConfig::FromTarget { ptr } => {
ReaderWriterChannel::from_target(&mut self.io, ChannelActor::Producer, ptr).await
}
}
}
async fn rsp_channel<'method>(
&'method mut self,
) -> Result<ReaderWriterChannel<'method, 'a, R, W>, crate::Error> {
match self.rsp_ch_config {
ChannelConfig::Direct { ptr, size } => {
ReaderWriterChannel::new(&mut self.io, ChannelActor::Consumer, ptr, size).await
}
ChannelConfig::FromTarget { ptr } => {
ReaderWriterChannel::from_target(&mut self.io, ChannelActor::Consumer, ptr).await
}
}
}
}