use crate::{Result, QsshError, PqAlgorithm};
use crate::crypto::{PqKeyExchange, SymmetricCrypto};
use crate::transport::Transport;
use tokio::net::{TcpListener, TcpStream};
use std::sync::Arc;
use tokio::sync::RwLock;
use serde::{Serialize, Deserialize};
pub mod discovery;
pub mod nat;
#[derive(Debug, Clone, Copy, PartialEq)]
pub enum P2pMode {
Direct,
NatTraversal,
Relay,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerInfo {
pub id: String,
pub name: String,
pub addresses: Vec<String>,
pub public_key: Vec<u8>,
pub algorithms: Vec<PqAlgorithm>,
pub last_seen: i64,
}
pub struct P2pSession {
our_info: PeerInfo,
keypair: PqKeyExchange,
peers: Arc<RwLock<Vec<PeerInfo>>>,
connections: Arc<RwLock<Vec<P2pConnection>>>,
}
pub struct P2pConnection {
peer: PeerInfo,
transport: Transport,
mode: P2pMode,
}
impl P2pSession {
pub async fn new(name: String) -> Result<Self> {
let keypair = PqKeyExchange::new()?;
use sha3::{Sha3_256, Digest};
let mut hasher = Sha3_256::new();
hasher.update(&keypair.falcon_pk);
let id = hex::encode(&hasher.finalize()[..16]);
let our_info = PeerInfo {
id,
name,
addresses: Self::discover_addresses().await?,
public_key: keypair.falcon_pk.clone(),
algorithms: vec![PqAlgorithm::Falcon512, PqAlgorithm::SphincsPlus],
last_seen: chrono::Utc::now().timestamp(),
};
Ok(Self {
our_info,
keypair,
peers: Arc::new(RwLock::new(Vec::new())),
connections: Arc::new(RwLock::new(Vec::new())),
})
}
pub fn peer_info(&self) -> &PeerInfo {
&self.our_info
}
async fn discover_addresses() -> Result<Vec<String>> {
let mut addrs = Vec::new();
if let Ok(ifaces) = if_addrs::get_if_addrs() {
for iface in ifaces {
if !iface.is_loopback() {
addrs.push(format!("{}:22222", iface.ip()));
}
}
}
if addrs.is_empty() {
addrs.push("127.0.0.1:22222".to_string());
}
Ok(addrs)
}
pub async fn listen(&self, addr: &str) -> Result<()> {
let listener = TcpListener::bind(addr).await
.map_err(|e| QsshError::Connection(format!("Failed to bind: {}", e)))?;
log::info!("P2P listening on {}", addr);
loop {
let (stream, remote_addr) = listener.accept().await?;
log::info!("Incoming P2P connection from {}", remote_addr);
let session = self.clone_ref();
tokio::spawn(async move {
if let Err(e) = session.handle_incoming(stream).await {
log::error!("P2P connection error: {}", e);
}
});
}
}
pub async fn connect(&self, peer_addr: &str) -> Result<P2pConnection> {
log::info!("Connecting to peer at {}", peer_addr);
let stream = TcpStream::connect(peer_addr).await
.map_err(|e| QsshError::Connection(format!("Failed to connect: {}", e)))?;
let (transport, peer_info) = self.perform_outgoing_handshake(stream).await?;
let connection = P2pConnection {
peer: peer_info,
transport,
mode: P2pMode::Direct,
};
self.connections.write().await.push(connection);
Ok(self.connections.read().await.last()
.ok_or_else(|| QsshError::Connection("No connections available".into()))?
.clone_ref())
}
async fn handle_incoming(&self, stream: TcpStream) -> Result<()> {
let (transport, peer_info) = self.perform_incoming_handshake(stream).await?;
let connection = P2pConnection {
peer: peer_info.clone(),
transport,
mode: P2pMode::Direct,
};
self.connections.write().await.push(connection);
log::info!("P2P connection established with {}", peer_info.name);
Ok(())
}
async fn perform_outgoing_handshake(&self, mut stream: TcpStream) -> Result<(Transport, PeerInfo)> {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let our_data = bincode::serialize(&self.our_info)
.map_err(|e| QsshError::Protocol(format!("Serialization failed: {}", e)))?;
let len = (our_data.len() as u32).to_be_bytes();
stream.write_all(&len).await?;
stream.write_all(&our_data).await?;
let mut len_bytes = [0u8; 4];
stream.read_exact(&mut len_bytes).await?;
let len = u32::from_be_bytes(len_bytes) as usize;
let mut peer_data = vec![0u8; len];
stream.read_exact(&mut peer_data).await?;
let peer_info: PeerInfo = bincode::deserialize(&peer_data)
.map_err(|e| QsshError::Protocol(format!("Deserialization failed: {}", e)))?;
let shared_secret = self.perform_key_exchange(&mut stream, &peer_info, true).await?;
let crypto = SymmetricCrypto::from_shared_secret(&shared_secret)?;
let transport = Transport::new(stream, crypto)?;
Ok((transport, peer_info))
}
async fn perform_incoming_handshake(&self, mut stream: TcpStream) -> Result<(Transport, PeerInfo)> {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let mut len_bytes = [0u8; 4];
stream.read_exact(&mut len_bytes).await?;
let len = u32::from_be_bytes(len_bytes) as usize;
let mut peer_data = vec![0u8; len];
stream.read_exact(&mut peer_data).await?;
let peer_info: PeerInfo = bincode::deserialize(&peer_data)
.map_err(|e| QsshError::Protocol(format!("Deserialization failed: {}", e)))?;
let our_data = bincode::serialize(&self.our_info)
.map_err(|e| QsshError::Protocol(format!("Serialization failed: {}", e)))?;
let len = (our_data.len() as u32).to_be_bytes();
stream.write_all(&len).await?;
stream.write_all(&our_data).await?;
let shared_secret = self.perform_key_exchange(&mut stream, &peer_info, false).await?;
let crypto = SymmetricCrypto::from_shared_secret(&shared_secret)?;
let transport = Transport::new(stream, crypto)?;
Ok((transport, peer_info))
}
async fn perform_key_exchange(
&self,
stream: &mut TcpStream,
peer_info: &PeerInfo,
initiator: bool
) -> Result<Vec<u8>> {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let (our_share, our_signature) = self.keypair.create_key_share()?;
if initiator {
let len = (our_share.len() as u32).to_be_bytes();
stream.write_all(&len).await?;
stream.write_all(&our_share).await?;
let sig_len = (our_signature.len() as u32).to_be_bytes();
stream.write_all(&sig_len).await?;
stream.write_all(&our_signature).await?;
let mut len_bytes = [0u8; 4];
stream.read_exact(&mut len_bytes).await?;
let len = u32::from_be_bytes(len_bytes) as usize;
let mut peer_share = vec![0u8; len];
stream.read_exact(&mut peer_share).await?;
stream.read_exact(&mut len_bytes).await?;
let sig_len = u32::from_be_bytes(len_bytes) as usize;
let mut peer_signature = vec![0u8; sig_len];
stream.read_exact(&mut peer_signature).await?;
let verified_share = self.keypair.process_key_share(
&peer_info.public_key,
&peer_share,
&peer_signature
)?;
Ok(self.keypair.compute_shared_secret(
&our_share,
&verified_share,
self.our_info.id.as_bytes(),
peer_info.id.as_bytes()
))
} else {
let mut len_bytes = [0u8; 4];
stream.read_exact(&mut len_bytes).await?;
let len = u32::from_be_bytes(len_bytes) as usize;
let mut peer_share = vec![0u8; len];
stream.read_exact(&mut peer_share).await?;
stream.read_exact(&mut len_bytes).await?;
let sig_len = u32::from_be_bytes(len_bytes) as usize;
let mut peer_signature = vec![0u8; sig_len];
stream.read_exact(&mut peer_signature).await?;
let len = (our_share.len() as u32).to_be_bytes();
stream.write_all(&len).await?;
stream.write_all(&our_share).await?;
let sig_len = (our_signature.len() as u32).to_be_bytes();
stream.write_all(&sig_len).await?;
stream.write_all(&our_signature).await?;
let verified_share = self.keypair.process_key_share(
&peer_info.public_key,
&peer_share,
&peer_signature
)?;
Ok(self.keypair.compute_shared_secret(
&our_share,
&verified_share,
self.our_info.id.as_bytes(),
peer_info.id.as_bytes()
))
}
}
fn clone_ref(&self) -> P2pSession {
P2pSession {
our_info: self.our_info.clone(),
keypair: self.keypair.clone(),
peers: self.peers.clone(),
connections: self.connections.clone(),
}
}
}
impl P2pConnection {
pub async fn send(&self, data: &[u8]) -> Result<()> {
use crate::transport::{Message, ChannelMessage};
let msg = Message::Channel(ChannelMessage::Data {
channel_id: 0,
data: data.to_vec(),
});
self.transport.send_message(&msg).await
}
pub async fn receive(&self) -> Result<Vec<u8>> {
use crate::transport::Message;
match self.transport.receive_message::<Message>().await? {
Message::Channel(crate::transport::ChannelMessage::Data { data, .. }) => Ok(data),
_ => Err(QsshError::Protocol("Expected data message".into())),
}
}
pub fn peer_info(&self) -> &PeerInfo {
&self.peer
}
fn clone_ref(&self) -> P2pConnection {
P2pConnection {
peer: self.peer.clone(),
transport: self.transport.clone(),
mode: self.mode,
}
}
}
impl Clone for PqKeyExchange {
fn clone(&self) -> Self {
PqKeyExchange::new().expect("Failed to clone PqKeyExchange")
}
}