use crate::{Error, Result};
use async_trait::async_trait;
use futures::stream::Stream;
use serde::{Deserialize, Serialize};
use std::pin::Pin;
pub mod http;
pub mod stdio;
pub mod websocket;
#[async_trait]
pub trait Transport: Send + Sync {
async fn send(&mut self, message: &str) -> Result<()>;
async fn receive(&mut self) -> Result<Option<String>>;
async fn close(&mut self) -> Result<()>;
fn is_connected(&self) -> bool;
fn transport_type(&self) -> &'static str;
}
#[derive(Debug)]
pub enum TransportConfig {
Stdio,
WebSocket {
url: String,
},
Http {
url: String,
},
Custom(Box<dyn CustomTransportConfig>),
}
pub trait CustomTransportConfig: std::fmt::Debug + Send + Sync {
fn transport_type(&self) -> &'static str;
}
pub type MessageStream = Pin<Box<dyn Stream<Item = Result<String>> + Send>>;
pub struct TransportFactory;
impl TransportFactory {
pub async fn create(config: TransportConfig) -> Result<Box<dyn Transport>> {
match config {
TransportConfig::Stdio => Ok(Box::new(stdio::StdioTransport::new())),
TransportConfig::WebSocket { url } => {
Ok(Box::new(websocket::WebSocketTransport::new(&url).await?))
}
TransportConfig::Http { url } => Ok(Box::new(http::HttpTransport::new(&url).await?)),
TransportConfig::Custom(_) => {
Err(Error::internal("Custom transports not yet implemented"))
}
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TransportMessage {
pub data: String,
pub timestamp: chrono::DateTime<chrono::Utc>,
}
impl TransportMessage {
pub fn new(data: String) -> Self {
Self {
data,
timestamp: chrono::Utc::now(),
}
}
}
#[derive(Debug, Clone, Default)]
pub struct TransportStats {
pub messages_sent: u64,
pub messages_received: u64,
pub bytes_sent: u64,
pub bytes_received: u64,
pub connection_time: Option<chrono::DateTime<chrono::Utc>>,
pub last_activity: Option<chrono::DateTime<chrono::Utc>>,
}