Skip to main content

sb_mesh/
ipc.rs

1use crate::config::MeshConfig;
2use crate::identity::NodeIdentity;
3use crate::peer::{PeerState, PeerTable};
4use crate::transport::MeshTransport;
5use serde::{Deserialize, Serialize};
6use std::sync::Arc;
7use tokio::io::{AsyncReadExt, AsyncWriteExt};
8use tokio::sync::RwLock;
9
10pub const IPC_PIPE_NAME: &str = r"\\.\pipe\sb-mesh";
11pub const IPC_TCP_PORT: u16 = 58890;
12
13#[derive(Debug, Clone, Serialize, Deserialize)]
14#[serde(tag = "action", content = "params")]
15pub enum IpcRequest {
16    GetStatus,
17    ListPeers,
18    SendPing { target: String },
19    SetAirGap { enabled: bool },
20}
21
22#[derive(Debug, Clone, Serialize, Deserialize)]
23#[serde(tag = "status", content = "data")]
24pub enum IpcResponse {
25    Ok(IpcPayload),
26    Error(String),
27}
28
29#[derive(Debug, Clone, Serialize, Deserialize)]
30#[serde(untagged)]
31pub enum IpcPayload {
32    Status(NodeStatusPayload),
33    Peers(Vec<PeerSummaryPayload>),
34    Ack(String),
35}
36
37#[derive(Debug, Clone, Serialize, Deserialize)]
38pub struct NodeStatusPayload {
39    pub node_id: String,
40    pub wireguard_pubkey: String,
41    pub listen_port: u16,
42    pub air_gap_killswitch: bool,
43    pub lan_discovery: bool,
44    pub zk_attestation: String,
45    pub peer_count: usize,
46}
47
48#[derive(Debug, Clone, Serialize, Deserialize)]
49pub struct PeerSummaryPayload {
50    pub callsign: String,
51    pub node_id: String,
52    pub public_key_base64: String,
53    pub endpoint: Option<String>,
54    pub overlay_ip: Option<String>,
55    pub status: String,
56    pub rtt_ms: Option<f64>,
57    pub bytes_sent: u64,
58    pub bytes_recv: u64,
59    pub zk_attested: bool,
60}
61
62impl From<&PeerState> for PeerSummaryPayload {
63    fn from(p: &PeerState) -> Self {
64        Self {
65            callsign: p.config.callsign.clone(),
66            node_id: p.config.node_id.clone(),
67            public_key_base64: p.config.public_key_base64.clone(),
68            endpoint: p.config.endpoint.clone(),
69            overlay_ip: p.config.overlay_ip.clone(),
70            status: p.status.badge_text().to_string(),
71            rtt_ms: p.rtt_ms,
72            bytes_sent: p.bytes_sent,
73            bytes_recv: p.bytes_recv,
74            zk_attested: p.zk_attested,
75        }
76    }
77}
78
79pub struct IpcServer {
80    pub identity: Arc<NodeIdentity>,
81    pub transport: Arc<MeshTransport>,
82    pub peer_table: Arc<RwLock<PeerTable>>,
83    pub config: Arc<RwLock<MeshConfig>>,
84}
85
86impl IpcServer {
87    pub fn new(
88        identity: Arc<NodeIdentity>,
89        transport: Arc<MeshTransport>,
90        peer_table: Arc<RwLock<PeerTable>>,
91        config: Arc<RwLock<MeshConfig>>,
92    ) -> Self {
93        Self {
94            identity,
95            transport,
96            peer_table,
97            config,
98        }
99    }
100
101    /// Spawns the IPC listener in the background (using local TCP loopback 127.0.0.1:58890 for universal cross-platform reliability)
102    pub fn start(&self) -> tokio::task::JoinHandle<()> {
103        let identity = self.identity.clone();
104        let transport = self.transport.clone();
105        let peer_table = self.peer_table.clone();
106        let config = self.config.clone();
107
108        tokio::spawn(async move {
109            let addr = format!("127.0.0.1:{}", IPC_TCP_PORT);
110            if let Ok(listener) = tokio::net::TcpListener::bind(&addr).await {
111                loop {
112                    if let Ok((mut socket, _)) = listener.accept().await {
113                        let id_c = identity.clone();
114                        let tr_c = transport.clone();
115                        let pt_c = peer_table.clone();
116                        let cfg_c = config.clone();
117
118                        tokio::spawn(async move {
119                            let mut buf = vec![0u8; 4096];
120                            if let Ok(n) = socket.read(&mut buf).await {
121                                if n > 0 {
122                                    let req_str = String::from_utf8_lossy(&buf[..n]);
123                                    let resp = match serde_json::from_str::<IpcRequest>(&req_str) {
124                                        Ok(IpcRequest::GetStatus) => {
125                                            let cfg = cfg_c.read().await;
126                                            let pt = pt_c.read().await;
127                                            let payload = NodeStatusPayload {
128                                                node_id: id_c.node_id.clone(),
129                                                wireguard_pubkey: id_c.wireguard_pubkey.clone(),
130                                                listen_port: cfg.listen_port,
131                                                air_gap_killswitch: cfg.air_gap_killswitch,
132                                                lan_discovery: cfg.lan_discovery,
133                                                zk_attestation: "ACTIVE (Phase 2 / 0-Leak)".to_string(),
134                                                peer_count: pt.peers.len(),
135                                            };
136                                            IpcResponse::Ok(IpcPayload::Status(payload))
137                                        }
138                                        Ok(IpcRequest::ListPeers) => {
139                                            let pt = pt_c.read().await;
140                                            let list: Vec<PeerSummaryPayload> = pt.list().into_iter().map(Into::into).collect();
141                                            IpcResponse::Ok(IpcPayload::Peers(list))
142                                        }
143                                        Ok(IpcRequest::SendPing { target }) => {
144                                            let pt = pt_c.read().await;
145                                            let found = pt.peers.values().find(|p| {
146                                                p.config.callsign == target
147                                                    || p.config.node_id == target
148                                                    || p.config.public_key_base64 == target
149                                            }).map(|p| p.config.public_key_base64.clone());
150                                            drop(pt);
151
152                                            if let Some(pk) = found {
153                                                match tr_c.send_ping(&pk).await {
154                                                    Ok(_) => IpcResponse::Ok(IpcPayload::Ack(format!("Ping sent to {}", target))),
155                                                    Err(e) => IpcResponse::Error(format!("Ping failed: {}", e)),
156                                                }
157                                            } else {
158                                                IpcResponse::Error(format!("Peer not found: {}", target))
159                                            }
160                                        }
161                                        Ok(IpcRequest::SetAirGap { enabled }) => {
162                                            let mut cfg = cfg_c.write().await;
163                                            cfg.air_gap_killswitch = enabled;
164                                            let mut air = tr_c.air_gap.write().await;
165                                            *air = enabled;
166                                            IpcResponse::Ok(IpcPayload::Ack(format!("Air-gap killswitch set to {}", enabled)))
167                                        }
168                                        Err(e) => IpcResponse::Error(format!("Invalid IPC request: {}", e)),
169                                    };
170
171                                    if let Ok(resp_bytes) = serde_json::to_vec(&resp) {
172                                        let _ = socket.write_all(&resp_bytes).await;
173                                    }
174                                }
175                            }
176                        });
177                    }
178                }
179            }
180        })
181    }
182}
183
184/// Client helper for external processes (such as sb-browser and CLI)
185pub struct IpcClient;
186
187impl IpcClient {
188    /// Send an IPC request to the running sb-mesh daemon
189    pub async fn query(req: IpcRequest) -> Result<IpcResponse, String> {
190        let addr = format!("127.0.0.1:{}", IPC_TCP_PORT);
191        let mut stream = tokio::net::TcpStream::connect(&addr)
192            .await
193            .map_err(|e| format!("Could not connect to sb-mesh daemon at {}: {}", addr, e))?;
194
195        let req_bytes = serde_json::to_vec(&req)
196            .map_err(|e| format!("Serialization error: {}", e))?;
197        stream.write_all(&req_bytes)
198            .await
199            .map_err(|e| format!("IPC write error: {}", e))?;
200
201        let mut buf = vec![0u8; 65535];
202        let n = stream.read(&mut buf)
203            .await
204            .map_err(|e| format!("IPC read error: {}", e))?;
205
206        serde_json::from_slice(&buf[..n])
207            .map_err(|e| format!("Failed to parse IPC response: {}", e))
208    }
209
210    pub async fn get_status() -> Result<NodeStatusPayload, String> {
211        match Self::query(IpcRequest::GetStatus).await? {
212            IpcResponse::Ok(IpcPayload::Status(s)) => Ok(s),
213            IpcResponse::Error(e) => Err(e),
214            _ => Err("Unexpected response payload".to_string()),
215        }
216    }
217
218    pub async fn list_peers() -> Result<Vec<PeerSummaryPayload>, String> {
219        match Self::query(IpcRequest::ListPeers).await? {
220            IpcResponse::Ok(IpcPayload::Peers(p)) => Ok(p),
221            IpcResponse::Error(e) => Err(e),
222            _ => Err("Unexpected response payload".to_string()),
223        }
224    }
225}
226