use crate::Result;
use crate::common::EchoClient;
use crate::unix::datagram_protocol::{UnixDatagramExt, UnixDatagramProtocol};
use crate::unix::stream_protocol::{ManagedUnixStream, Protocol, StreamExt};
use async_trait::async_trait;
use std::path::PathBuf;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::UnixDatagram;
pub struct UnixStreamEchoClient {
stream: ManagedUnixStream,
}
impl UnixStreamEchoClient {
pub async fn connect(socket_path: PathBuf) -> Result<Self> {
let stream = Protocol::connect_unix(&socket_path).await?;
Ok(Self { stream })
}
}
#[async_trait]
impl EchoClient for UnixStreamEchoClient {
async fn echo(&mut self, data: &[u8]) -> Result<Vec<u8>> {
self.stream
.write_all(data)
.await
.map_err(crate::EchoError::Unix)?;
self.stream.flush().await.map_err(crate::EchoError::Unix)?;
let mut buffer = vec![0u8; data.len()];
let mut response = Vec::new();
let mut total_read = 0;
while total_read < data.len() {
let n = self
.stream
.read(&mut buffer)
.await
.map_err(crate::EchoError::Unix)?;
if n == 0 {
break; }
response.extend_from_slice(&buffer[..n]);
total_read += n;
}
Ok(response)
}
}
pub struct UnixDatagramEchoClient {
socket: UnixDatagram,
server_path: PathBuf,
}
impl UnixDatagramEchoClient {
pub async fn connect(server_path: PathBuf) -> Result<Self> {
let socket = UnixDatagramProtocol::connect_unix(&server_path).await?;
Ok(Self {
socket,
server_path,
})
}
}
#[async_trait]
impl EchoClient for UnixDatagramEchoClient {
async fn echo(&mut self, data: &[u8]) -> Result<Vec<u8>> {
self.socket
.send_to(data, &self.server_path)
.await
.map_err(crate::EchoError::Unix)?;
let mut buffer = vec![0u8; data.len()];
let (len, _) = self
.socket
.recv_from(&mut buffer)
.await
.map_err(crate::EchoError::Unix)?;
Ok(buffer[..len].to_vec())
}
}