use crate::{
types::SecureMessage,
error::Result,
};
use async_trait::async_trait;
use std::{
time::{Duration, Instant},
collections::HashMap,
};
use dashmap::DashMap;
use serde::{Serialize, Deserialize};
use super::router::ConnectionOffer;
#[async_trait]
pub trait Transport: Send + Sync {
fn transport_type(&self) -> TransportType;
fn capabilities(&self) -> TransportCapabilities;
async fn can_reach(&self, target: &TransportTarget) -> bool;
async fn estimate_metrics(&self, target: &TransportTarget) -> Result<TransportEstimate>;
async fn send_message(&self, target: &TransportTarget, message: &SecureMessage) -> Result<DeliveryReceipt>;
async fn receive_messages(&self) -> Result<Vec<IncomingMessage>>;
async fn test_connectivity(&self, target: &TransportTarget) -> Result<ConnectivityResult>;
async fn start(&self) -> Result<()>;
async fn stop(&self) -> Result<()>;
async fn status(&self) -> TransportStatus;
async fn metrics(&self) -> TransportMetrics;
async fn send_connection_offer(&self, target: &str, offer: ConnectionOffer) -> Result<String> {
let _ = (target, offer);
Err(crate::error::SynapseError::TransportError("Connection offers not supported by this transport".to_string()))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum TransportType {
Tcp,
Udp,
WebSocket,
Http,
Email,
AutoDiscovery, Quic,
Custom(u32), }
impl std::fmt::Display for TransportType {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
TransportType::Tcp => write!(f, "TCP"),
TransportType::Udp => write!(f, "UDP"),
TransportType::WebSocket => write!(f, "WebSocket"),
TransportType::Http => write!(f, "HTTP"),
TransportType::Email => write!(f, "Email"),
TransportType::AutoDiscovery => write!(f, "Auto-Discovery"),
TransportType::Quic => write!(f, "QUIC"),
TransportType::Custom(id) => write!(f, "Custom({})", id),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TransportCapabilities {
pub max_message_size: usize,
pub reliable: bool,
pub real_time: bool,
pub broadcast: bool,
pub bidirectional: bool,
pub encrypted: bool,
pub network_spanning: bool,
pub supported_urgencies: Vec<MessageUrgency>,
pub features: Vec<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum MessageUrgency {
Critical,
RealTime,
Interactive,
Background,
Batch,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TransportTarget {
pub identifier: String,
pub address: Option<String>,
pub preferred_transports: Vec<TransportType>,
pub required_capabilities: Vec<String>,
pub urgency: MessageUrgency,
}
impl TransportTarget {
pub fn new(identifier: String) -> Self {
Self {
identifier,
address: None,
preferred_transports: Vec::new(),
required_capabilities: Vec::new(),
urgency: MessageUrgency::Interactive,
}
}
pub fn with_address(mut self, address: String) -> Self {
self.address = Some(address);
self
}
pub fn with_urgency(mut self, urgency: MessageUrgency) -> Self {
self.urgency = urgency;
self
}
pub fn prefer_transport(mut self, transport: TransportType) -> Self {
self.preferred_transports.push(transport);
self
}
pub fn require_capability(mut self, capability: String) -> Self {
self.required_capabilities.push(capability);
self
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TransportEstimate {
pub latency: Duration,
pub reliability: f64,
pub bandwidth: u64,
pub cost: f64,
pub available: bool,
pub confidence: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DeliveryReceipt {
pub message_id: String,
pub transport_used: TransportType,
pub delivery_time: Duration,
pub target_reached: String,
pub confirmation: DeliveryConfirmation,
pub metadata: HashMap<String, String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum DeliveryConfirmation {
Sent,
Delivered,
Received,
Acknowledged,
}
#[derive(Debug, Clone)]
pub struct IncomingMessage {
pub message: SecureMessage,
pub transport_type: TransportType,
pub source: String,
pub received_timestamp: u64,
pub metadata: HashMap<String, String>,
}
impl IncomingMessage {
pub fn new(message: SecureMessage, transport_type: TransportType, source: String) -> Self {
Self {
message,
transport_type,
source,
received_timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
metadata: HashMap::new(),
}
}
pub fn received_at(&self) -> Duration {
Duration::from_secs(self.received_timestamp)
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ConnectivityResult {
pub connected: bool,
pub rtt: Option<Duration>,
pub error: Option<String>,
pub quality: f64,
pub details: HashMap<String, String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum TransportStatus {
Stopped,
Starting,
Running,
Degraded,
Failed,
Stopping,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TransportMetrics {
pub transport_type: TransportType,
pub messages_sent: u64,
pub messages_received: u64,
pub send_failures: u64,
pub receive_failures: u64,
pub bytes_sent: u64,
pub bytes_received: u64,
pub average_latency_ms: u64,
pub reliability_score: f64,
pub active_connections: u32,
pub last_updated_timestamp: u64,
pub custom_metrics: HashMap<String, f64>,
}
impl Default for TransportMetrics {
fn default() -> Self {
Self {
transport_type: TransportType::Tcp,
messages_sent: 0,
messages_received: 0,
send_failures: 0,
receive_failures: 0,
bytes_sent: 0,
bytes_received: 0,
average_latency_ms: 0,
reliability_score: 1.0,
active_connections: 0,
last_updated_timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs(),
custom_metrics: HashMap::new(),
}
}
}
impl TransportMetrics {
pub fn average_latency(&self) -> Duration {
Duration::from_millis(self.average_latency_ms)
}
pub fn set_average_latency(&mut self, latency: Duration) {
self.average_latency_ms = latency.as_millis() as u64;
}
pub fn last_updated(&self) -> Instant {
Instant::now() }
pub fn touch(&mut self) {
self.last_updated_timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DeliveryEstimate {
pub latency: Duration,
pub reliability: f64,
pub throughput_estimate: u64,
pub cost_score: f64,
}
impl TransportCapabilities {
pub fn tcp() -> Self {
Self {
max_message_size: 64 * 1024 * 1024, reliable: true,
real_time: false,
broadcast: false,
bidirectional: true,
encrypted: false, network_spanning: true,
supported_urgencies: vec![
MessageUrgency::Interactive,
MessageUrgency::Background,
MessageUrgency::Batch,
],
features: vec![
"connection_oriented".to_string(),
"stream_based".to_string(),
"flow_control".to_string(),
],
}
}
pub fn udp() -> Self {
Self {
max_message_size: 65507, reliable: false,
real_time: true,
broadcast: true,
bidirectional: true,
encrypted: false,
network_spanning: true,
supported_urgencies: vec![
MessageUrgency::Critical,
MessageUrgency::RealTime,
MessageUrgency::Interactive,
],
features: vec![
"connectionless".to_string(),
"datagram_based".to_string(),
"low_overhead".to_string(),
"multicast".to_string(),
],
}
}
pub fn email() -> Self {
Self {
max_message_size: 25 * 1024 * 1024, reliable: true,
real_time: false,
broadcast: true,
bidirectional: true,
encrypted: true, network_spanning: true,
supported_urgencies: vec![
MessageUrgency::Background,
MessageUrgency::Batch,
],
features: vec![
"store_and_forward".to_string(),
"federation".to_string(),
"authentication".to_string(),
"persistent".to_string(),
],
}
}
pub fn auto_discovery() -> Self {
Self {
max_message_size: 1024, reliable: false,
real_time: true,
broadcast: true,
bidirectional: true,
encrypted: false,
network_spanning: false, supported_urgencies: vec![
MessageUrgency::Critical,
MessageUrgency::RealTime,
],
features: vec![
"service_discovery".to_string(),
"zero_configuration".to_string(),
"local_network".to_string(),
"multicast".to_string(),
],
}
}
pub fn websocket() -> Self {
Self {
max_message_size: 16 * 1024 * 1024, reliable: true,
real_time: true,
broadcast: false,
bidirectional: true,
encrypted: true, network_spanning: true,
supported_urgencies: vec![
MessageUrgency::Critical,
MessageUrgency::RealTime,
MessageUrgency::Interactive,
],
features: vec![
"web_compatible".to_string(),
"full_duplex".to_string(),
"frame_based".to_string(),
"http_upgrade".to_string(),
],
}
}
pub fn quic() -> Self {
Self {
max_message_size: 1024 * 1024 * 1024, reliable: true,
real_time: true,
broadcast: false,
bidirectional: true,
encrypted: true, network_spanning: true,
supported_urgencies: vec![
MessageUrgency::Critical,
MessageUrgency::RealTime,
MessageUrgency::Interactive,
MessageUrgency::Background,
],
features: vec![
"multiplexed_streams".to_string(),
"zero_rtt".to_string(),
"connection_migration".to_string(),
"modern_crypto".to_string(),
"congestion_control".to_string(),
],
}
}
pub fn http() -> Self {
Self {
max_message_size: 10 * 1024 * 1024, reliable: true,
real_time: false, broadcast: false,
bidirectional: true,
encrypted: false, network_spanning: true,
supported_urgencies: vec![
MessageUrgency::Interactive,
MessageUrgency::Background,
MessageUrgency::Batch,
],
features: vec![
"firewall_friendly".to_string(),
"web_compatible".to_string(),
"request_response".to_string(),
"standard_protocol".to_string(),
],
}
}
pub fn https() -> Self {
let mut caps = Self::http();
caps.encrypted = true;
caps.features.push("encrypted".to_string());
caps.features.push("authenticated".to_string());
caps
}
}
#[async_trait]
pub trait TransportFactory: Send + Sync {
async fn create_transport(&self, config: &HashMap<String, String>) -> Result<Box<dyn Transport>>;
fn transport_type(&self) -> TransportType;
fn default_config(&self) -> HashMap<String, String>;
fn validate_config(&self, config: &HashMap<String, String>) -> Result<()>;
}
pub struct TcpTransportFactory;
#[async_trait]
impl TransportFactory for TcpTransportFactory {
async fn create_transport(&self, _config: &HashMap<String, String>) -> Result<Box<dyn Transport>> {
let transport = crate::transport::tcp_simple::SimpleTcpTransport::new();
Ok(Box::new(transport))
}
fn transport_type(&self) -> TransportType {
TransportType::Tcp
}
fn default_config(&self) -> HashMap<String, String> {
let mut config = HashMap::new();
config.insert("listen_port".to_string(), "0".to_string());
config.insert("connection_timeout_ms".to_string(), "30000".to_string());
config.insert("max_message_size".to_string(), "1048576".to_string()); config
}
fn validate_config(&self, config: &HashMap<String, String>) -> Result<()> {
if let Some(port_str) = config.get("listen_port") {
if port_str.parse::<u16>().is_err() {
return Err(crate::error::SynapseError::Config(
"Invalid port number".to_string()
));
}
}
Ok(())
}
}
pub struct EmailTransportFactory;
#[async_trait]
impl TransportFactory for EmailTransportFactory {
async fn create_transport(&self, _config: &HashMap<String, String>) -> Result<Box<dyn Transport>> {
Err(crate::error::SynapseError::TransportError("Email transport not yet implemented".to_string()))
}
fn transport_type(&self) -> TransportType {
TransportType::Email
}
fn default_config(&self) -> HashMap<String, String> {
HashMap::new()
}
fn validate_config(&self, _config: &HashMap<String, String>) -> Result<()> {
Ok(())
}
}
pub struct MdnsTransportFactory;
#[async_trait]
impl TransportFactory for MdnsTransportFactory {
async fn create_transport(&self, _config: &HashMap<String, String>) -> Result<Box<dyn Transport>> {
Err(crate::error::SynapseError::TransportError("mDNS transport temporarily disabled".to_string()))
}
fn transport_type(&self) -> TransportType {
TransportType::AutoDiscovery
}
fn default_config(&self) -> HashMap<String, String> {
let mut config = HashMap::new();
config.insert("service_name".to_string(), "_synapse._tcp.local".to_string());
config.insert("local_port".to_string(), "0".to_string());
config.insert("discovery_timeout_ms".to_string(), "5000".to_string());
config.insert("max_message_size".to_string(), "65507".to_string()); config
}
fn validate_config(&self, config: &HashMap<String, String>) -> Result<()> {
if let Some(service_name) = config.get("service_name") {
if !service_name.contains("._tcp.") && !service_name.contains("._udp.") {
return Err(crate::error::SynapseError::Config(
"Service name must include protocol (_tcp. or _udp.)".to_string()
));
}
}
Ok(())
}
}
#[cfg(feature = "http")]
pub struct HttpTransportFactory;
#[cfg(feature = "http")]
#[async_trait]
impl TransportFactory for HttpTransportFactory {
async fn create_transport(&self, config: &HashMap<String, String>) -> Result<Box<dyn Transport>> {
let transport = crate::transport::http_unified::HttpTransportImpl::new(config).await?;
Ok(Box::new(transport))
}
fn transport_type(&self) -> TransportType {
TransportType::Http
}
fn default_config(&self) -> HashMap<String, String> {
let mut config = HashMap::new();
config.insert("use_https".to_string(), "true".to_string());
config.insert("server_port".to_string(), "0".to_string()); config.insert("server_address".to_string(), "127.0.0.1".to_string());
config.insert("timeout_ms".to_string(), "30000".to_string());
config.insert("max_message_size".to_string(), "10485760".to_string()); config.insert("user_agent".to_string(), "Synapse-HTTP-Transport/1.0".to_string());
config
}
fn validate_config(&self, config: &HashMap<String, String>) -> Result<()> {
if let Some(port_str) = config.get("server_port") {
if port_str.parse::<u16>().is_err() {
return Err(crate::error::SynapseError::Config(
"Invalid server port number".to_string()
));
}
}
if let Some(timeout_str) = config.get("timeout_ms") {
if timeout_str.parse::<u64>().is_err() {
return Err(crate::error::SynapseError::Config(
"Invalid timeout value".to_string()
));
}
}
if let Some(size_str) = config.get("max_message_size") {
if size_str.parse::<usize>().is_err() {
return Err(crate::error::SynapseError::Config(
"Invalid max message size".to_string()
));
}
}
Ok(())
}
}
pub struct UnifiedTransportManager {
transports: HashMap<TransportType, Box<dyn Transport>>,
target_preferences: HashMap<String, TransportType>, metrics_cache: DashMap<String, HashMap<TransportType, TransportEstimate>>, failover_policies: HashMap<TransportType, Vec<TransportType>>, config: UnifiedTransportConfig,
}
#[derive(Debug, Clone)]
pub struct UnifiedTransportConfig {
pub default_transport: TransportType,
pub enable_automatic_failover: bool,
pub prefer_real_time_for_sync: bool,
pub metrics_cache_seconds: u64,
pub optimize_for_bandwidth: bool,
}
impl Default for UnifiedTransportConfig {
fn default() -> Self {
Self {
default_transport: TransportType::WebSocket,
enable_automatic_failover: true,
prefer_real_time_for_sync: true,
metrics_cache_seconds: 60,
optimize_for_bandwidth: false,
}
}
}
impl UnifiedTransportManager {
pub async fn new(config: UnifiedTransportConfig) -> Result<Self> {
let mut manager = Self {
transports: HashMap::new(),
target_preferences: HashMap::new(),
metrics_cache: DashMap::new(),
failover_policies: HashMap::new(),
config,
};
manager.setup_default_failover_policies();
Ok(manager)
}
pub fn register_transport(&mut self, transport: Box<dyn Transport>) -> Result<()> {
let transport_type = transport.transport_type();
if self.transports.contains_key(&transport_type) {
return Err(crate::error::SynapseError::TransportError(format!("Transport already registered: {}", transport_type)));
}
self.transports.insert(transport_type, transport);
Ok(())
}
pub async fn send_message(&self, target: &TransportTarget, message: &SecureMessage) -> Result<DeliveryReceipt> {
let transport_type = if let Some(preferred) = self.target_preferences.get(&target.identifier) {
*preferred
} else {
self.determine_best_transport(target, message).await?
};
if let Some(transport) = self.transports.get(&transport_type) {
match transport.send_message(target, message).await {
Ok(receipt) => {
return Ok(receipt);
}
Err(err) => {
if self.config.enable_automatic_failover {
return self.try_failover_transports(target, message, transport_type).await;
}
return Err(err);
}
}
}
Err(crate::error::SynapseError::TransportError(format!("No suitable transport found for target: {}", target.identifier)))
}
async fn try_failover_transports(
&self,
target: &TransportTarget,
message: &SecureMessage,
failed_transport: TransportType,
) -> Result<DeliveryReceipt> {
if let Some(failover_list) = self.failover_policies.get(&failed_transport) {
for transport_type in failover_list {
if let Some(transport) = self.transports.get(transport_type) {
if transport.can_reach(target).await {
match transport.send_message(target, message).await {
Ok(receipt) => return Ok(receipt),
Err(_) => continue, }
}
}
}
}
Err(crate::error::SynapseError::TransportError(format!("All transports failed for target: {}", target.identifier)))
}
async fn determine_best_transport(&self, target: &TransportTarget, _message: &SecureMessage) -> Result<TransportType> {
if !target.preferred_transports.is_empty() {
for preferred in &target.preferred_transports {
if self.transports.contains_key(preferred) {
return Ok(*preferred);
}
}
}
let mut available_transports = Vec::new();
for (transport_type, transport) in &self.transports {
if transport.can_reach(target).await {
let metrics = transport.estimate_metrics(target).await?;
available_transports.push((*transport_type, metrics));
}
}
if available_transports.is_empty() {
return Err(crate::error::SynapseError::TransportError(format!("No transports can reach target: {}", target.identifier)));
}
available_transports.sort_by(|(_, a), (_, b)| {
match (a.available, b.available) {
(true, false) => return std::cmp::Ordering::Less,
(false, true) => return std::cmp::Ordering::Greater,
_ => {} }
let reliability_cmp = b.reliability.partial_cmp(&a.reliability).unwrap_or(std::cmp::Ordering::Equal);
if reliability_cmp != std::cmp::Ordering::Equal {
return reliability_cmp;
}
a.latency.partial_cmp(&b.latency).unwrap_or(std::cmp::Ordering::Equal)
});
Ok(available_transports[0].0)
}
fn setup_default_failover_policies(&mut self) {
self.failover_policies.insert(
TransportType::WebSocket,
vec![TransportType::Http, TransportType::Email]
);
self.failover_policies.insert(
TransportType::Http,
vec![TransportType::WebSocket, TransportType::Email]
);
self.failover_policies.insert(
TransportType::Email,
vec![TransportType::Http, TransportType::WebSocket]
);
self.failover_policies.insert(
TransportType::Quic,
vec![TransportType::WebSocket, TransportType::Http, TransportType::Email]
);
}
pub async fn receive_messages(&self) -> Result<Vec<IncomingMessage>> {
let mut all_messages = Vec::new();
for transport in self.transports.values() {
match transport.receive_messages().await {
Ok(mut messages) => all_messages.append(&mut messages),
Err(_) => continue, }
}
Ok(all_messages)
}
pub async fn receive_from_transport(&self, transport_type: TransportType) -> Result<Vec<IncomingMessage>> {
if let Some(transport) = self.transports.get(&transport_type) {
return transport.receive_messages().await;
}
Err(crate::error::SynapseError::TransportError(format!("Transport not found: {}", transport_type)))
}
pub async fn start_all_transports(&self) -> Result<()> {
for transport in self.transports.values() {
if let Err(e) = transport.start().await {
tracing::error!("Failed to start transport {}: {}", transport.transport_type(), e);
}
}
Ok(())
}
pub async fn stop_all_transports(&self) -> Result<()> {
for transport in self.transports.values() {
if let Err(e) = transport.stop().await {
tracing::error!("Failed to stop transport {}: {}", transport.transport_type(), e);
}
}
Ok(())
}
}
pub fn create_standard_factories() -> Vec<Box<dyn TransportFactory>> {
let factories: Vec<Box<dyn TransportFactory>> = vec![
Box::new(TcpTransportFactory),
Box::new(super::udp_unified::UdpTransportFactory),
Box::new(EmailTransportFactory),
Box::new(MdnsTransportFactory),
];
#[cfg(feature = "http")]
{
let mut factories = factories;
factories.push(Box::new(HttpTransportFactory));
factories
}
#[cfg(not(feature = "http"))]
{
factories
}
}