sb-mesh 0.1.1

S&B Sovereign Mesh (sb-mesh) — User-Space P2P Overlay Network, WireGuard-compatible Crypto & TUI
Documentation
use crate::config::MeshConfig;
use crate::identity::NodeIdentity;
use crate::peer::{PeerState, PeerTable};
use crate::transport::MeshTransport;
use serde::{Deserialize, Serialize};
use std::sync::Arc;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::sync::RwLock;

pub const IPC_PIPE_NAME: &str = r"\\.\pipe\sb-mesh";
pub const IPC_TCP_PORT: u16 = 58890;

#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "action", content = "params")]
pub enum IpcRequest {
    GetStatus,
    ListPeers,
    SendPing { target: String },
    SetAirGap { enabled: bool },
}

#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum IpcResponse {
    Ok(IpcPayload),
    Error(String),
}

#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(untagged)]
pub enum IpcPayload {
    Status(NodeStatusPayload),
    Peers(Vec<PeerSummaryPayload>),
    Ack(String),
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NodeStatusPayload {
    pub node_id: String,
    pub wireguard_pubkey: String,
    pub listen_port: u16,
    pub air_gap_killswitch: bool,
    pub lan_discovery: bool,
    pub zk_attestation: String,
    pub peer_count: usize,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerSummaryPayload {
    pub callsign: String,
    pub node_id: String,
    pub public_key_base64: String,
    pub endpoint: Option<String>,
    pub overlay_ip: Option<String>,
    pub status: String,
    pub rtt_ms: Option<f64>,
    pub bytes_sent: u64,
    pub bytes_recv: u64,
    pub zk_attested: bool,
}

impl From<&PeerState> for PeerSummaryPayload {
    fn from(p: &PeerState) -> Self {
        Self {
            callsign: p.config.callsign.clone(),
            node_id: p.config.node_id.clone(),
            public_key_base64: p.config.public_key_base64.clone(),
            endpoint: p.config.endpoint.clone(),
            overlay_ip: p.config.overlay_ip.clone(),
            status: p.status.badge_text().to_string(),
            rtt_ms: p.rtt_ms,
            bytes_sent: p.bytes_sent,
            bytes_recv: p.bytes_recv,
            zk_attested: p.zk_attested,
        }
    }
}

pub struct IpcServer {
    pub identity: Arc<NodeIdentity>,
    pub transport: Arc<MeshTransport>,
    pub peer_table: Arc<RwLock<PeerTable>>,
    pub config: Arc<RwLock<MeshConfig>>,
}

impl IpcServer {
    pub fn new(
        identity: Arc<NodeIdentity>,
        transport: Arc<MeshTransport>,
        peer_table: Arc<RwLock<PeerTable>>,
        config: Arc<RwLock<MeshConfig>>,
    ) -> Self {
        Self {
            identity,
            transport,
            peer_table,
            config,
        }
    }

    /// Spawns the IPC listener in the background (using local TCP loopback 127.0.0.1:58890 for universal cross-platform reliability)
    pub fn start(&self) -> tokio::task::JoinHandle<()> {
        let identity = self.identity.clone();
        let transport = self.transport.clone();
        let peer_table = self.peer_table.clone();
        let config = self.config.clone();

        tokio::spawn(async move {
            let addr = format!("127.0.0.1:{}", IPC_TCP_PORT);
            if let Ok(listener) = tokio::net::TcpListener::bind(&addr).await {
                loop {
                    if let Ok((mut socket, _)) = listener.accept().await {
                        let id_c = identity.clone();
                        let tr_c = transport.clone();
                        let pt_c = peer_table.clone();
                        let cfg_c = config.clone();

                        tokio::spawn(async move {
                            let mut buf = vec![0u8; 4096];
                            if let Ok(n) = socket.read(&mut buf).await {
                                if n > 0 {
                                    let req_str = String::from_utf8_lossy(&buf[..n]);
                                    let resp = match serde_json::from_str::<IpcRequest>(&req_str) {
                                        Ok(IpcRequest::GetStatus) => {
                                            let cfg = cfg_c.read().await;
                                            let pt = pt_c.read().await;
                                            let payload = NodeStatusPayload {
                                                node_id: id_c.node_id.clone(),
                                                wireguard_pubkey: id_c.wireguard_pubkey.clone(),
                                                listen_port: cfg.listen_port,
                                                air_gap_killswitch: cfg.air_gap_killswitch,
                                                lan_discovery: cfg.lan_discovery,
                                                zk_attestation: "ACTIVE (Phase 2 / 0-Leak)".to_string(),
                                                peer_count: pt.peers.len(),
                                            };
                                            IpcResponse::Ok(IpcPayload::Status(payload))
                                        }
                                        Ok(IpcRequest::ListPeers) => {
                                            let pt = pt_c.read().await;
                                            let list: Vec<PeerSummaryPayload> = pt.list().into_iter().map(Into::into).collect();
                                            IpcResponse::Ok(IpcPayload::Peers(list))
                                        }
                                        Ok(IpcRequest::SendPing { target }) => {
                                            let pt = pt_c.read().await;
                                            let found = pt.peers.values().find(|p| {
                                                p.config.callsign == target
                                                    || p.config.node_id == target
                                                    || p.config.public_key_base64 == target
                                            }).map(|p| p.config.public_key_base64.clone());
                                            drop(pt);

                                            if let Some(pk) = found {
                                                match tr_c.send_ping(&pk).await {
                                                    Ok(_) => IpcResponse::Ok(IpcPayload::Ack(format!("Ping sent to {}", target))),
                                                    Err(e) => IpcResponse::Error(format!("Ping failed: {}", e)),
                                                }
                                            } else {
                                                IpcResponse::Error(format!("Peer not found: {}", target))
                                            }
                                        }
                                        Ok(IpcRequest::SetAirGap { enabled }) => {
                                            let mut cfg = cfg_c.write().await;
                                            cfg.air_gap_killswitch = enabled;
                                            let mut air = tr_c.air_gap.write().await;
                                            *air = enabled;
                                            IpcResponse::Ok(IpcPayload::Ack(format!("Air-gap killswitch set to {}", enabled)))
                                        }
                                        Err(e) => IpcResponse::Error(format!("Invalid IPC request: {}", e)),
                                    };

                                    if let Ok(resp_bytes) = serde_json::to_vec(&resp) {
                                        let _ = socket.write_all(&resp_bytes).await;
                                    }
                                }
                            }
                        });
                    }
                }
            }
        })
    }
}

/// Client helper for external processes (such as sb-browser and CLI)
pub struct IpcClient;

impl IpcClient {
    /// Send an IPC request to the running sb-mesh daemon
    pub async fn query(req: IpcRequest) -> Result<IpcResponse, String> {
        let addr = format!("127.0.0.1:{}", IPC_TCP_PORT);
        let mut stream = tokio::net::TcpStream::connect(&addr)
            .await
            .map_err(|e| format!("Could not connect to sb-mesh daemon at {}: {}", addr, e))?;

        let req_bytes = serde_json::to_vec(&req)
            .map_err(|e| format!("Serialization error: {}", e))?;
        stream.write_all(&req_bytes)
            .await
            .map_err(|e| format!("IPC write error: {}", e))?;

        let mut buf = vec![0u8; 65535];
        let n = stream.read(&mut buf)
            .await
            .map_err(|e| format!("IPC read error: {}", e))?;

        serde_json::from_slice(&buf[..n])
            .map_err(|e| format!("Failed to parse IPC response: {}", e))
    }

    pub async fn get_status() -> Result<NodeStatusPayload, String> {
        match Self::query(IpcRequest::GetStatus).await? {
            IpcResponse::Ok(IpcPayload::Status(s)) => Ok(s),
            IpcResponse::Error(e) => Err(e),
            _ => Err("Unexpected response payload".to_string()),
        }
    }

    pub async fn list_peers() -> Result<Vec<PeerSummaryPayload>, String> {
        match Self::query(IpcRequest::ListPeers).await? {
            IpcResponse::Ok(IpcPayload::Peers(p)) => Ok(p),
            IpcResponse::Error(e) => Err(e),
            _ => Err("Unexpected response payload".to_string()),
        }
    }
}