eyes-subscriber 0.1.4

Tracing subscriber for sending traces to Eyes (eyes.coreyja.com)
Documentation
use async_trait::async_trait;
use tokio::sync::{mpsc, oneshot};

use crate::EventData;

#[async_trait]
pub(crate) 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 => {
                // Process any remaining events in the channel
                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 => {
                // Channel closed and no shutdown signal
                break;
            }
        }
    }

    if let Err(e) = transport.close().await {
        eprintln!("Failed to close transport: {}", e);
    }

    let _ = completion_tx.send(());
}