use tokio::net::UnixStream;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use ubiquity_core::{MeshMessage, Result, UbiquityError};
use tracing::{debug, error};
pub struct MeshConnection {
stream: UnixStream,
agent_id: String,
}
impl MeshConnection {
pub fn new(stream: UnixStream, agent_id: String) -> Self {
Self { stream, agent_id }
}
pub async fn send_message(&mut self, message: &MeshMessage) -> Result<()> {
let data = serde_json::to_vec(message)?;
self.stream.write_all(&data).await
.map_err(|e| UbiquityError::SocketError(e))?;
Ok(())
}
pub async fn receive_message(&mut self) -> Result<Option<MeshMessage>> {
let mut buffer = vec![0u8; 4096];
match self.stream.read(&mut buffer).await {
Ok(0) => Ok(None), Ok(n) => {
let message = serde_json::from_slice(&buffer[..n])?;
Ok(Some(message))
}
Err(e) => Err(UbiquityError::SocketError(e)),
}
}
}