ubiquity_mesh/
connection.rs1use tokio::net::UnixStream;
4use tokio::io::{AsyncReadExt, AsyncWriteExt};
5use ubiquity_core::{MeshMessage, Result, UbiquityError};
6use tracing::{debug, error};
7
8pub struct MeshConnection {
9 stream: UnixStream,
10 agent_id: String,
11}
12
13impl MeshConnection {
14 pub fn new(stream: UnixStream, agent_id: String) -> Self {
15 Self { stream, agent_id }
16 }
17
18 pub async fn send_message(&mut self, message: &MeshMessage) -> Result<()> {
19 let data = serde_json::to_vec(message)?;
20 self.stream.write_all(&data).await
21 .map_err(|e| UbiquityError::SocketError(e))?;
22 Ok(())
23 }
24
25 pub async fn receive_message(&mut self) -> Result<Option<MeshMessage>> {
26 let mut buffer = vec![0u8; 4096];
27
28 match self.stream.read(&mut buffer).await {
29 Ok(0) => Ok(None), Ok(n) => {
31 let message = serde_json::from_slice(&buffer[..n])?;
32 Ok(Some(message))
33 }
34 Err(e) => Err(UbiquityError::SocketError(e)),
35 }
36 }
37}