use async_trait::async_trait;
use tokio::sync::{mpsc, oneshot};
use crate::EventData;
#[async_trait]
pub trait Transport: Send + Sync + 'static {
async fn connect(&mut self) -> Result<(), TransportError>;
async fn send(&mut self, event: EventData) -> Result<(), TransportError>;
async fn close(&mut self) -> Result<(), TransportError>;
}
#[derive(Debug, thiserror::Error)]
pub enum TransportError {
#[error("Connection error: {0}")]
Connection(String),
#[error("Send error: {0}")]
Send(String),
#[error("Configuration error: {0}")]
Configuration(String),
}
pub async fn run_transport_loop(
mut transport: Box<dyn Transport>,
mut receiver: mpsc::UnboundedReceiver<EventData>,
mut shutdown_rx: oneshot::Receiver<()>,
completion_tx: oneshot::Sender<()>,
) {
if let Err(e) = transport.connect().await {
eprintln!("Failed to connect transport: {}", e);
let _ = completion_tx.send(());
return;
}
loop {
tokio::select! {
Some(event) = receiver.recv() => {
if let Err(e) = transport.send(event).await {
eprintln!("Failed to send event: {}", e);
}
}
_ = &mut shutdown_rx => {
while let Ok(event) = receiver.try_recv() {
if let Err(e) = transport.send(event).await {
eprintln!("Failed to send event during shutdown: {}", e);
}
}
break;
}
else => {
break;
}
}
}
if let Err(e) = transport.close().await {
eprintln!("Failed to close transport: {}", e);
}
let _ = completion_tx.send(());
}