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 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
184pub struct IpcClient;
186
187impl IpcClient {
188 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