use crate::crypto::{
decrypt_payload, derive_session_key, encrypt_payload, pack_message, ratchet_forward, unpack_header, PacketType,
};
use crate::identity::NodeIdentity;
use crate::peer::{PathType, PeerStatus, PeerTable};
use base64::prelude::*;
use chrono::Utc;
use std::net::SocketAddr;
use std::sync::Arc;
use tokio::net::UdpSocket;
use tokio::sync::RwLock;
use x25519_dalek::{PublicKey, StaticSecret};
pub struct MeshTransport {
pub identity: Arc<NodeIdentity>,
pub peer_table: Arc<RwLock<PeerTable>>,
pub socket: Arc<UdpSocket>,
pub air_gap: Arc<RwLock<bool>>,
pub routing_table: Arc<RwLock<crate::routing::RoutingTable>>,
pub gossip_engine: Arc<RwLock<crate::gossip::GossipEngine>>,
}
impl MeshTransport {
pub async fn bind(
listen_port: u16,
identity: Arc<NodeIdentity>,
peer_table: Arc<RwLock<PeerTable>>,
air_gap: Arc<RwLock<bool>>,
) -> Result<Arc<Self>, String> {
let addr = format!("0.0.0.0:{}", listen_port);
let socket = UdpSocket::bind(&addr)
.await
.map_err(|e| format!("Failed to bind UDP transport socket on {}: {}", addr, e))?;
let transport = Arc::new(Self {
identity,
peer_table,
socket: Arc::new(socket),
air_gap,
routing_table: Arc::new(RwLock::new(crate::routing::RoutingTable::new())),
gossip_engine: Arc::new(RwLock::new(crate::gossip::GossipEngine::new())),
});
let recv_transport = transport.clone();
tokio::spawn(async move {
recv_transport.run_receiver().await;
});
Ok(transport)
}
async fn run_receiver(&self) {
let mut buf = [0u8; 65535];
loop {
match self.socket.recv_from(&mut buf).await {
Ok((len, sender_addr)) => {
if *self.air_gap.read().await {
continue; }
if len >= 7 && &buf[..4] == &crate::relay::RELAY_MAGIC {
if let Ok(relay_pkt) = crate::relay::RelayPacket::from_bytes(&buf[..len]) {
let p_table = self.peer_table.read().await;
let r_table = self.routing_table.read().await;
let action = crate::relay::RelayManager::process_packet_with_routing(
relay_pkt,
&self.identity.node_id,
&*p_table,
Some(&*r_table),
);
drop(p_table);
drop(r_table);
match action {
crate::relay::RelayAction::Forward { next_hop, packet } => {
let bytes = packet.to_bytes();
let _ = self.socket.send_to(&bytes, next_hop).await;
}
crate::relay::RelayAction::DeliverLocally(_pkt) => {}
crate::relay::RelayAction::Drop(_) => {}
}
continue;
}
}
if len < 45 {
continue;
}
let packet_data = &buf[..len];
if let Ok((ptype, sender_pubkey_bytes, counter, payload)) =
unpack_header(packet_data)
{
let sender_pubkey_b64 = BASE64_STANDARD.encode(sender_pubkey_bytes);
if self.gossip_engine.read().await.is_revoked(&sender_pubkey_b64) {
continue; }
let remote_public = PublicKey::from(sender_pubkey_bytes);
let mut table = self.peer_table.write().await;
if let Some(peer) = table.get_mut_by_pubkey(&sender_pubkey_b64) {
peer.bytes_recv += len as u64;
if peer.parsed_endpoint != Some(sender_addr) {
peer.parsed_endpoint = Some(sender_addr);
peer.roaming_events += 1;
}
let current_path = if sender_addr.is_ipv6() {
PathType::DirectIPv6(sender_addr)
} else {
PathType::DirectIPv4(sender_addr)
};
peer.add_or_update_path(current_path.clone());
let session_key = if let Some(ref k) = peer.session_key {
*k.as_bytes()
} else {
derive_session_key(&self.identity.secret, &remote_public)
};
match ptype {
PacketType::HandshakeInit => {
peer.status = PeerStatus::HandshakeInProgress;
peer.last_handshake = Some(Utc::now());
let resp_counter = peer.sequence_counter + 1;
peer.sequence_counter = resp_counter;
if let Ok(enc_payload) =
encrypt_payload(&session_key, resp_counter, b"SBM_HS_OK")
{
let reply = pack_message(
PacketType::HandshakeResp,
self.identity.public_key.as_bytes(),
resp_counter,
&enc_payload,
);
let _ = self.socket.send_to(&reply, sender_addr).await;
peer.bytes_sent += reply.len() as u64;
peer.status = PeerStatus::Connected;
}
}
PacketType::HandshakeResp => {
if decrypt_payload(&session_key, counter, payload).is_ok() {
peer.status = PeerStatus::Connected;
peer.last_handshake = Some(Utc::now());
}
}
PacketType::Ping => {
if decrypt_payload(&session_key, counter, payload).is_ok() {
let resp_counter = peer.sequence_counter + 1;
peer.sequence_counter = resp_counter;
if let Ok(enc) =
encrypt_payload(&session_key, resp_counter, b"PONG")
{
let pong = pack_message(
PacketType::Pong,
self.identity.public_key.as_bytes(),
resp_counter,
&enc,
);
let _ = self.socket.send_to(&pong, sender_addr).await;
peer.bytes_sent += pong.len() as u64;
}
}
}
PacketType::Pong => {
if decrypt_payload(&session_key, counter, payload).is_ok() {
if let Some(sent_time) = peer.last_ping_sent {
let rtt = (Utc::now() - sent_time)
.num_microseconds()
.unwrap_or(0)
as f64
/ 1000.0;
let rtt_val = rtt.max(0.1);
peer.rtt_ms = Some(rtt_val);
peer.record_path_success(¤t_path, rtt_val);
peer.status = PeerStatus::Connected;
}
}
}
PacketType::RekeyInit => {
if let Ok(dec) = decrypt_payload(&session_key, counter, payload) {
if dec.len() == 32 {
let mut initiator_eph_pub_bytes = [0u8; 32];
initiator_eph_pub_bytes.copy_from_slice(&dec);
let initiator_eph_pub = PublicKey::from(initiator_eph_pub_bytes);
let resp_eph_secret = StaticSecret::random_from_rng(rand::rngs::OsRng);
let resp_eph_pub = PublicKey::from(&resp_eph_secret);
let eph_shared = resp_eph_secret.diffie_hellman(&initiator_eph_pub);
let next_epoch = peer.session_epoch + 1;
let next_key = ratchet_forward(&session_key, eph_shared.as_bytes(), next_epoch);
let resp_counter = peer.sequence_counter + 1;
peer.sequence_counter = resp_counter;
if let Ok(enc_resp) = encrypt_payload(&session_key, resp_counter, resp_eph_pub.as_bytes()) {
let reply = pack_message(
PacketType::RekeyResp,
self.identity.public_key.as_bytes(),
resp_counter,
&enc_resp,
);
let _ = self.socket.send_to(&reply, sender_addr).await;
peer.bytes_sent += reply.len() as u64;
peer.session_key = Some(next_key);
peer.session_epoch = next_epoch;
peer.last_rekey = Some(Utc::now());
peer.rekey_count += 1;
peer.status = PeerStatus::Connected;
}
}
}
}
PacketType::RekeyResp => {
if let Ok(dec) = decrypt_payload(&session_key, counter, payload) {
if dec.len() == 32 {
if let Some(init_eph_secret_bytes) = peer.pending_ephemeral_secret.take() {
let init_eph_secret = StaticSecret::from(init_eph_secret_bytes);
let mut resp_eph_pub_bytes = [0u8; 32];
resp_eph_pub_bytes.copy_from_slice(&dec);
let resp_eph_pub = PublicKey::from(resp_eph_pub_bytes);
let eph_shared = init_eph_secret.diffie_hellman(&resp_eph_pub);
let next_epoch = peer.session_epoch + 1;
let next_key = ratchet_forward(&session_key, eph_shared.as_bytes(), next_epoch);
peer.session_key = Some(next_key);
peer.session_epoch = next_epoch;
peer.last_rekey = Some(Utc::now());
peer.rekey_count += 1;
peer.status = PeerStatus::Connected;
}
}
}
}
PacketType::Data => {
peer.status = PeerStatus::Connected;
peer.record_data_activity();
}
}
}
}
}
Err(_) => {
tokio::time::sleep(tokio::time::Duration::from_millis(50)).await;
}
}
}
}
pub async fn initiate_handshake(&self, pubkey_b64: &str) -> Result<(), String> {
if *self.air_gap.read().await {
return Err("Air-Gap is active. Network transmission blocked.".to_string());
}
let mut table = self.peer_table.write().await;
let peer = table
.get_mut_by_pubkey(pubkey_b64)
.ok_or_else(|| "Peer not found".to_string())?;
let endpoint = peer
.parsed_endpoint
.ok_or_else(|| "Peer has no endpoint IP:Port".to_string())?;
let pubkey_bytes = BASE64_STANDARD
.decode(pubkey_b64)
.map_err(|e| format!("Invalid pubkey base64: {}", e))?;
if pubkey_bytes.len() != 32 {
return Err("Invalid public key length".to_string());
}
let mut key_arr = [0u8; 32];
key_arr.copy_from_slice(&pubkey_bytes);
let remote_public = PublicKey::from(key_arr);
let session_key = derive_session_key(&self.identity.secret, &remote_public);
let counter = peer.sequence_counter + 1;
peer.sequence_counter = counter;
peer.status = PeerStatus::HandshakeInProgress;
peer.last_handshake = Some(Utc::now());
let encrypted = encrypt_payload(&session_key, counter, b"SBM_HANDSHAKE_INIT")?;
let packet = pack_message(
PacketType::HandshakeInit,
self.identity.public_key.as_bytes(),
counter,
&encrypted,
);
peer.bytes_sent += packet.len() as u64;
self.socket
.send_to(&packet, endpoint)
.await
.map_err(|e| format!("Failed to send handshake packet: {}", e))?;
Ok(())
}
pub async fn send_ping(&self, pubkey_b64: &str) -> Result<(), String> {
if *self.air_gap.read().await {
return Err("Air-Gap is active. Ping blocked.".to_string());
}
let mut table = self.peer_table.write().await;
let peer = table
.get_mut_by_pubkey(pubkey_b64)
.ok_or_else(|| "Peer not found".to_string())?;
let endpoint = peer
.parsed_endpoint
.ok_or_else(|| "Peer has no valid endpoint configured".to_string())?;
let pubkey_bytes = BASE64_STANDARD
.decode(pubkey_b64)
.map_err(|e| format!("Invalid pubkey base64: {}", e))?;
if pubkey_bytes.len() != 32 {
return Err("Invalid public key length".to_string());
}
let mut key_arr = [0u8; 32];
key_arr.copy_from_slice(&pubkey_bytes);
let remote_public = PublicKey::from(key_arr);
let session_key = derive_session_key(&self.identity.secret, &remote_public);
let counter = peer.sequence_counter + 1;
peer.sequence_counter = counter;
peer.last_ping_sent = Some(Utc::now());
let enc = encrypt_payload(&session_key, counter, b"PING")?;
let packet = pack_message(
PacketType::Ping,
self.identity.public_key.as_bytes(),
counter,
&enc,
);
peer.bytes_sent += packet.len() as u64;
self.socket
.send_to(&packet, endpoint)
.await
.map_err(|e| format!("Failed to send ping packet: {}", e))?;
Ok(())
}
pub async fn send_via_relay(
&self,
relay_endpoint: SocketAddr,
target_node_id: &str,
payload: &[u8],
) -> Result<(), String> {
if *self.air_gap.read().await {
return Err("Air-Gap is active. Relay transmission blocked.".to_string());
}
let relay_pkt = crate::relay::RelayPacket::new(
target_node_id,
&self.identity.node_id,
payload.to_vec(),
);
let wire_bytes = relay_pkt.to_bytes();
self.socket
.send_to(&wire_bytes, relay_endpoint)
.await
.map_err(|e| format!("Failed to send via relay: {}", e))?;
Ok(())
}
pub async fn send_obfuscated_packet(
&self,
target_endpoint: SocketAddr,
payload: &[u8],
target_bucket: Option<usize>,
) -> Result<(), String> {
if *self.air_gap.read().await {
return Err("Air-Gap is active. Transmission blocked.".to_string());
}
let padded = crate::obfuscation::apply_padding(payload, target_bucket);
self.socket
.send_to(&padded, target_endpoint)
.await
.map_err(|e| format!("Failed to send obfuscated packet: {}", e))?;
Ok(())
}
pub async fn initiate_rekey(&self, pubkey_b64: &str) -> Result<(), String> {
if *self.air_gap.read().await {
return Err("Air-Gap is active. Network transmission blocked.".to_string());
}
let mut table = self.peer_table.write().await;
let peer = table
.get_mut_by_pubkey(pubkey_b64)
.ok_or_else(|| "Peer not found".to_string())?;
let endpoint = peer
.parsed_endpoint
.ok_or_else(|| "Peer has no endpoint IP:Port".to_string())?;
let pubkey_bytes = BASE64_STANDARD
.decode(pubkey_b64)
.map_err(|e| format!("Invalid pubkey base64: {}", e))?;
if pubkey_bytes.len() != 32 {
return Err("Invalid public key length".to_string());
}
let mut key_arr = [0u8; 32];
key_arr.copy_from_slice(&pubkey_bytes);
let remote_public = PublicKey::from(key_arr);
let active_key = if let Some(ref k) = peer.session_key {
*k.as_bytes()
} else {
derive_session_key(&self.identity.secret, &remote_public)
};
let eph_secret = StaticSecret::random_from_rng(rand::rngs::OsRng);
let eph_pub = PublicKey::from(&eph_secret);
peer.pending_ephemeral_secret = Some(eph_secret.to_bytes());
let counter = peer.sequence_counter + 1;
peer.sequence_counter = counter;
let encrypted = encrypt_payload(&active_key, counter, eph_pub.as_bytes())?;
let packet = pack_message(
PacketType::RekeyInit,
self.identity.public_key.as_bytes(),
counter,
&encrypted,
);
peer.bytes_sent += packet.len() as u64;
self.socket
.send_to(&packet, endpoint)
.await
.map_err(|e| format!("Failed to send RekeyInit packet: {}", e))?;
Ok(())
}
pub async fn probe_candidate_paths(&self, pubkey_b64: &str) -> Result<usize, String> {
if *self.air_gap.read().await {
return Err("Air-Gap is active. Probe blocked.".to_string());
}
let mut table = self.peer_table.write().await;
let peer = table
.get_mut_by_pubkey(pubkey_b64)
.ok_or_else(|| "Peer not found".to_string())?;
let pubkey_bytes = BASE64_STANDARD
.decode(pubkey_b64)
.map_err(|e| format!("Invalid pubkey base64: {}", e))?;
if pubkey_bytes.len() != 32 {
return Err("Invalid public key length".to_string());
}
let mut key_arr = [0u8; 32];
key_arr.copy_from_slice(&pubkey_bytes);
let remote_public = PublicKey::from(key_arr);
let active_key = if let Some(ref k) = peer.session_key {
*k.as_bytes()
} else {
derive_session_key(&self.identity.secret, &remote_public)
};
let candidate_endpoints: Vec<SocketAddr> = peer
.candidate_paths
.iter()
.filter_map(|cp| match cp.path_type {
PathType::DirectIPv4(sa) | PathType::DirectIPv6(sa) => Some(sa),
PathType::Relay(_) => None,
})
.collect();
if candidate_endpoints.is_empty() {
return Ok(0);
}
let counter = peer.sequence_counter + 1;
peer.sequence_counter = counter;
peer.last_ping_sent = Some(Utc::now());
let enc = encrypt_payload(&active_key, counter, b"PING")?;
let packet = pack_message(
PacketType::Ping,
self.identity.public_key.as_bytes(),
counter,
&enc,
);
let mut sent_count = 0;
for ep in candidate_endpoints {
if self.socket.send_to(&packet, ep).await.is_ok() {
sent_count += 1;
}
}
peer.bytes_sent += (packet.len() * sent_count) as u64;
Ok(sent_count)
}
}