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,
}
}
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;
}
}
}
});
}
}
}
})
}
}
pub struct IpcClient;
impl IpcClient {
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()),
}
}
}