pub mod abstraction;
pub mod manager;
pub mod providers;
#[cfg(test)]
mod providers_test;
#[cfg(not(target_arch = "wasm32"))]
pub mod discovery;
#[cfg(not(target_arch = "wasm32"))]
pub mod udp_unified;
pub use abstraction::{TransportType, TransportCapabilities};
pub use discovery::DiscoveryTransport;
pub use auto_discovery::config::DiscoveryConfig;
#[cfg(not(target_arch = "wasm32"))]
pub mod tcp_simple;
#[cfg(not(target_arch = "wasm32"))]
pub mod tcp;
#[cfg(not(target_arch = "wasm32"))]
pub mod tcp_enhanced;
#[cfg(not(target_arch = "wasm32"))]
pub mod udp;
#[cfg(not(target_arch = "wasm32"))]
pub mod http_unified;
#[cfg(not(target_arch = "wasm32"))]
pub mod mdns_enhanced;
#[cfg(not(target_arch = "wasm32"))]
pub mod nat_traversal;
#[cfg(not(target_arch = "wasm32"))]
pub mod email_enhanced;
#[cfg(not(target_arch = "wasm32"))]
pub mod router;
#[cfg(not(target_arch = "wasm32"))]
pub mod llm_discovery;
use crate::{types::*, error::Result, circuit_breaker::{CircuitBreaker, RequestOutcome}};
use async_trait::async_trait;
use std::{time::{Duration, Instant}, sync::Arc};
use serde::{Serialize, Deserialize};
#[cfg(not(target_arch = "wasm32"))]
use tokio::time::timeout;
#[derive(Debug, Clone)]
pub enum TransportRoute {
#[cfg(not(target_arch = "wasm32"))]
DirectTcp {
address: String,
port: u16,
latency_ms: u32,
established_at: Instant,
},
#[cfg(not(target_arch = "wasm32"))]
DirectUdp {
address: String,
port: u16,
latency_ms: u32,
established_at: Instant,
},
#[cfg(not(target_arch = "wasm32"))]
Udp {
address: String,
latency: Duration,
reliability: f64,
},
#[cfg(not(target_arch = "wasm32"))]
WebSocket {
url: String,
latency: Duration,
reliability: f64,
},
#[cfg(not(target_arch = "wasm32"))]
Quic {
address: String,
latency: Duration,
reliability: f64,
multiplexed: bool,
},
#[cfg(not(target_arch = "wasm32"))]
LocalMdns {
service_name: String,
address: String,
port: u16,
latency_ms: u32,
discovered_at: Instant,
},
#[cfg(not(target_arch = "wasm32"))]
NatTraversal {
method: NatMethod,
external_address: String,
external_port: u16,
latency_ms: u32,
established_at: Instant,
},
#[cfg(not(target_arch = "wasm32"))]
FastEmailRelay {
relay_server: String,
estimated_latency_ms: u32,
},
#[cfg(not(target_arch = "wasm32"))]
StandardEmail {
estimated_latency_min: u32,
},
#[cfg(not(target_arch = "wasm32"))]
EmailDiscovery {
target_transport: Box<TransportRoute>,
},
#[cfg(target_arch = "wasm32")]
WebSocket {
url: String,
latency_ms: u32,
established_at: Instant,
},
#[cfg(target_arch = "wasm32")]
WebRtc {
peer_connection: String,
latency_ms: u32,
established_at: Instant,
},
#[cfg(target_arch = "wasm32")]
WebAssembly {
estimated_latency_ms: u32,
},
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum NatMethod {
Upnp,
Stun { server: String },
Turn { server: String, username: String },
IceCandidate,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ConnectionOffer {
pub entity_id: String,
#[cfg(not(target_arch = "wasm32"))]
pub tcp_endpoints: Vec<String>,
#[cfg(not(target_arch = "wasm32"))]
pub udp_endpoints: Vec<String>,
#[cfg(not(target_arch = "wasm32"))]
pub stun_servers: Vec<String>,
#[cfg(not(target_arch = "wasm32"))]
pub turn_servers: Vec<TurnServer>,
pub capabilities: Vec<String>,
pub public_key: String,
pub expires_at: chrono::DateTime<chrono::Utc>,
pub priority: u8,
#[cfg(target_arch = "wasm32")]
pub websocket_endpoints: Vec<String>,
#[cfg(target_arch = "wasm32")]
pub webrtc_endpoints: Vec<String>,
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TurnServer {
pub url: String,
pub username: String,
pub credential: String,
}
#[derive(Debug, Clone)]
pub struct TransportMetrics {
pub latency: Duration,
pub throughput_bps: u64,
pub packet_loss: f32,
pub jitter_ms: u32,
pub reliability_score: f32, pub last_updated: Instant,
}
impl Default for TransportMetrics {
fn default() -> Self {
Self {
latency: Duration::from_millis(50),
throughput_bps: 1_000_000, packet_loss: 0.0,
jitter_ms: 10,
reliability_score: 0.95,
last_updated: Instant::now(),
}
}
}
#[derive(Debug)]
pub struct HybridConnection {
pub primary: TransportRoute,
pub fallback: TransportRoute,
pub discovery_latency: Duration,
pub connection_latency: Duration,
pub total_setup_time: Duration,
pub metrics: TransportMetrics,
}
#[async_trait]
pub trait Transport: Send + Sync {
async fn send_message(&self, target: &str, message: &SecureMessage) -> Result<String>;
async fn send_message_with_breaker(
&self,
target: &str,
message: &SecureMessage,
circuit_breaker: Option<Arc<CircuitBreaker>>
) -> Result<String> {
if let Some(breaker) = circuit_breaker {
if !breaker.can_proceed().await {
return Err(crate::error::SynapseError::TransportError(
"Circuit breaker open - request rejected".to_string()
).into());
}
match self.send_message(target, message).await {
Ok(result) => {
breaker.record_outcome(RequestOutcome::Success).await;
Ok(result)
}
Err(e) => {
breaker.record_outcome(RequestOutcome::Failure(e.to_string())).await;
Err(e)
}
}
} else {
self.send_message(target, message).await
}
}
async fn receive_messages(&self) -> Result<Vec<SecureMessage>>;
async fn test_connectivity(&self, target: &str) -> Result<TransportMetrics>;
async fn test_connectivity_with_breaker(
&self,
target: &str,
circuit_breaker: Option<Arc<CircuitBreaker>>
) -> Result<TransportMetrics> {
if let Some(breaker) = circuit_breaker {
if !breaker.can_proceed().await {
return Err(crate::error::SynapseError::TransportError(
"Circuit breaker open - connectivity test rejected".to_string()
).into());
}
let start_time = Instant::now();
match self.test_connectivity(target).await {
Ok(metrics) => {
breaker.record_outcome(RequestOutcome::Success).await;
breaker.check_external_triggers(&metrics).await;
Ok(metrics)
}
Err(e) => {
let elapsed = start_time.elapsed();
if elapsed > Duration::from_secs(10) {
breaker.record_outcome(RequestOutcome::Timeout).await;
} else {
breaker.record_outcome(RequestOutcome::Failure(e.to_string())).await;
}
Err(e)
}
}
} else {
self.test_connectivity(target).await
}
}
async fn can_reach(&self, target: &str) -> bool;
fn get_capabilities(&self) -> Vec<String>;
fn estimated_latency(&self) -> Duration;
fn reliability_score(&self) -> f32;
}
pub struct TransportDiscovery {
discovery_timeout: Duration,
#[allow(dead_code)]
connectivity_cache: dashmap::DashMap<String, (TransportMetrics, Instant)>,
}
impl TransportDiscovery {
pub fn new() -> Self {
Self {
discovery_timeout: Duration::from_secs(10),
connectivity_cache: dashmap::DashMap::new(),
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn discover_transports(&mut self, target: &str) -> Result<Vec<TransportRoute>> {
let mut routes = Vec::new();
if let Ok(tcp_route) = self.test_direct_tcp(target).await {
routes.push(tcp_route);
}
if let Ok(udp_route) = self.test_direct_udp(target).await {
routes.push(udp_route);
}
if let Ok(mdns_route) = self.discover_mdns_peer(target).await {
routes.push(mdns_route);
}
if let Ok(nat_routes) = self.establish_nat_traversal(target).await {
routes.extend(nat_routes);
}
routes.push(TransportRoute::StandardEmail {
estimated_latency_min: 1 });
Ok(routes)
}
#[cfg(target_arch = "wasm32")]
pub async fn discover_transports(&mut self, _target: &str) -> Result<Vec<TransportRoute>> {
Ok(vec![TransportRoute::WebAssembly {
estimated_latency_ms: 50
}])
}
#[cfg(not(target_arch = "wasm32"))]
async fn test_direct_tcp(&self, target: &str) -> Result<TransportRoute> {
let ports = vec![8080, 8443, 9090, 7777];
for port in ports {
let address = format!("{}:{}", target, port);
let start = Instant::now();
match timeout(self.discovery_timeout, tokio::net::TcpStream::connect(&address)).await {
Ok(Ok(_stream)) => {
let latency = start.elapsed();
return Ok(TransportRoute::DirectTcp {
address: target.to_string(),
port,
latency_ms: latency.as_millis() as u32,
established_at: Instant::now(),
});
}
_ => continue,
}
}
Err(crate::error::SynapseError::TransportError("No direct TCP connection available".into()))
}
#[cfg(not(target_arch = "wasm32"))]
async fn test_direct_udp(&self, target: &str) -> Result<TransportRoute> {
let ports = vec![8080, 8443, 9090, 7777];
for port in ports {
let address = format!("{}:{}", target, port);
let start = Instant::now();
match timeout(
self.discovery_timeout,
tokio::net::UdpSocket::bind("0.0.0.0:0")
).await {
Ok(Ok(socket)) => {
if socket.connect(&address).await.is_ok() {
let latency = start.elapsed();
return Ok(TransportRoute::DirectUdp {
address: target.to_string(),
port,
latency_ms: latency.as_millis() as u32,
established_at: Instant::now(),
});
}
}
_ => continue,
}
}
Err(crate::error::SynapseError::TransportError("No direct UDP connection available".into()))
}
#[cfg(not(target_arch = "wasm32"))]
async fn discover_mdns_peer(&self, _target: &str) -> Result<TransportRoute> {
Err(crate::error::SynapseError::TransportError("mDNS discovery not yet implemented".into()))
}
#[cfg(not(target_arch = "wasm32"))]
async fn establish_nat_traversal(&self, _target: &str) -> Result<Vec<TransportRoute>> {
Ok(Vec::new())
}
}
pub struct TransportSelector {
discovery: TransportDiscovery,
performance_weights: TransportWeights,
}
#[derive(Debug, Clone)]
pub struct TransportWeights {
pub latency_weight: f32,
pub reliability_weight: f32,
pub throughput_weight: f32,
pub setup_time_weight: f32,
}
impl Default for TransportWeights {
fn default() -> Self {
Self {
latency_weight: 0.4,
reliability_weight: 0.3,
throughput_weight: 0.2,
setup_time_weight: 0.1,
}
}
}
impl TransportSelector {
pub fn new() -> Self {
Self {
discovery: TransportDiscovery::new(),
performance_weights: TransportWeights::default(),
}
}
#[cfg(not(target_arch = "wasm32"))]
pub async fn choose_optimal_transport(
&mut self,
target: &str,
urgency: abstraction::MessageUrgency
) -> Result<TransportRoute> {
let available_routes = self.discovery.discover_transports(target).await?;
match urgency {
abstraction::MessageUrgency::Critical | abstraction::MessageUrgency::RealTime => {
for route in available_routes {
match route {
TransportRoute::DirectTcp { latency_ms, .. }
| TransportRoute::DirectUdp { latency_ms, .. }
| TransportRoute::LocalMdns { latency_ms, .. }
| TransportRoute::NatTraversal { latency_ms, .. }
if latency_ms < 100 => return Ok(route),
_ => continue,
}
}
Err(crate::error::SynapseError::TransportError("No real-time transport available".into()))
}
abstraction::MessageUrgency::Interactive => {
self.select_best_route(available_routes, 1000).await
}
abstraction::MessageUrgency::Background | abstraction::MessageUrgency::Batch => {
self.select_most_reliable_route(available_routes).await
}
}
}
#[cfg(target_arch = "wasm32")]
pub async fn choose_optimal_transport(
&mut self,
_target: &str,
_urgency: abstraction::MessageUrgency
) -> Result<TransportRoute> {
Ok(TransportRoute::WebAssembly {
estimated_latency_ms: 50
})
}
#[cfg(not(target_arch = "wasm32"))]
async fn select_best_route(&self, routes: Vec<TransportRoute>, max_latency_ms: u32) -> Result<TransportRoute> {
let mut best_route = None;
let mut best_score = f32::MIN;
for route in routes {
let latency_ms = self.get_route_latency(&route);
if latency_ms <= max_latency_ms {
let score = self.calculate_route_score(&route);
if score > best_score {
best_score = score;
best_route = Some(route);
}
}
}
best_route.ok_or_else(|| {
crate::error::SynapseError::TransportError("No suitable transport found".into())
})
}
#[cfg(not(target_arch = "wasm32"))]
async fn select_most_reliable_route(&self, routes: Vec<TransportRoute>) -> Result<TransportRoute> {
for route in &routes {
match route {
TransportRoute::StandardEmail { .. } => return Ok(route.clone()),
TransportRoute::FastEmailRelay { .. } => return Ok(route.clone()),
_ => continue,
}
}
self.select_best_route(routes, u32::MAX).await
}
#[cfg(not(target_arch = "wasm32"))]
fn get_route_latency(&self, route: &TransportRoute) -> u32 {
match route {
TransportRoute::DirectTcp { latency_ms, .. }
| TransportRoute::DirectUdp { latency_ms, .. }
| TransportRoute::LocalMdns { latency_ms, .. }
| TransportRoute::NatTraversal { latency_ms, .. } => *latency_ms,
TransportRoute::FastEmailRelay { estimated_latency_ms, .. } => *estimated_latency_ms,
TransportRoute::StandardEmail { estimated_latency_min } => estimated_latency_min * 60 * 1000, TransportRoute::EmailDiscovery { .. } => 30_000, TransportRoute::Udp { .. } => 5, TransportRoute::WebSocket { .. } => 20, TransportRoute::Quic { .. } => 10, }
}
#[cfg(not(target_arch = "wasm32"))]
fn calculate_route_score(&self, route: &TransportRoute) -> f32 {
let latency_ms = self.get_route_latency(route) as f32;
let reliability = self.get_route_reliability(route);
let setup_time = self.get_route_setup_time(route);
let latency_score = 1.0 / (1.0 + latency_ms / 1000.0); let setup_score = 1.0 / (1.0 + setup_time);
latency_score * self.performance_weights.latency_weight +
reliability * self.performance_weights.reliability_weight +
setup_score * self.performance_weights.setup_time_weight
}
#[cfg(not(target_arch = "wasm32"))]
fn get_route_reliability(&self, route: &TransportRoute) -> f32 {
match route {
TransportRoute::StandardEmail { .. } => 0.95,
TransportRoute::FastEmailRelay { .. } => 0.90,
TransportRoute::DirectTcp { .. } => 0.85,
TransportRoute::DirectUdp { .. } => 0.80,
TransportRoute::LocalMdns { .. } => 0.95,
TransportRoute::NatTraversal { .. } => 0.70,
TransportRoute::EmailDiscovery { .. } => 0.95,
TransportRoute::Udp { .. } => 0.80,
TransportRoute::WebSocket { .. } => 0.85,
TransportRoute::Quic { .. } => 0.90,
}
}
#[cfg(not(target_arch = "wasm32"))]
fn get_route_setup_time(&self, route: &TransportRoute) -> f32 {
match route {
TransportRoute::DirectTcp { .. } => 0.1, TransportRoute::DirectUdp { .. } => 0.05, TransportRoute::LocalMdns { .. } => 0.05, TransportRoute::NatTraversal { .. } => 2.0, TransportRoute::FastEmailRelay { .. } => 1.0, TransportRoute::StandardEmail { .. } => 1.0, TransportRoute::EmailDiscovery { .. } => 30.0, TransportRoute::Udp { .. } => 0.05, TransportRoute::WebSocket { .. } => 0.2, TransportRoute::Quic { .. } => 0.1, }
}
}
pub use abstraction::*;
pub use manager::*;
#[cfg(not(target_arch = "wasm32"))]
pub use udp_unified::{UdpTransportImpl, UdpTransportFactory};
#[cfg(not(target_arch = "wasm32"))]
#[cfg(feature = "http")]
pub use http_unified::{HttpTransportImpl, HttpTransportFactory};
#[cfg(not(target_arch = "wasm32"))]
pub use tcp_simple::{SimpleTcpTransport, SimpleTcpTransportFactory};
#[cfg(not(target_arch = "wasm32"))]
pub use tcp::TcpTransport;
#[cfg(not(target_arch = "wasm32"))]
pub use tcp_enhanced::EnhancedTcpTransport;
#[cfg(not(target_arch = "wasm32"))]
pub use mdns_enhanced::{EnhancedMdnsTransport, MdnsConfig};
#[cfg(not(target_arch = "wasm32"))]
pub use nat_traversal::{NatTraversalTransport, IceCandidate};
#[cfg(not(target_arch = "wasm32"))]
#[cfg(feature = "email")]
pub use email_enhanced::{EmailEnhancedTransport, FastEmailRelay};
#[cfg(not(target_arch = "wasm32"))]
pub use router::MultiTransportRouter;
#[cfg(not(target_arch = "wasm32"))]
pub use llm_discovery::{
LlmDiscoveryManager, LlmDiscoveryConfig, DiscoveredLlm, LlmConnection,
LlmModelInfo, LlmConnectionInfo, LlmPerformanceMetrics, LlmStatus,
LlmRequest, LlmResponse, LlmResponseMetadata
};