ubiquity-mesh 0.1.1

Unix socket mesh for zero-port agent communication
Documentation
//! Mesh connection handling

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), // Connection closed
            Ok(n) => {
                let message = serde_json::from_slice(&buffer[..n])?;
                Ok(Some(message))
            }
            Err(e) => Err(UbiquityError::SocketError(e)),
        }
    }
}