use std::sync::Arc;
use tokio::sync::Mutex;
pub mod error;
pub mod models;
mod config;
mod lnm;
mod repositories;
mod state;
pub use config::StreamClientConfig;
use error::Result;
use lnm::LnmStreamRepo;
pub use repositories::StreamRepository;
pub use state::StreamConnectionStatus;
pub type StreamConnection = Arc<dyn StreamRepository>;
pub struct StreamClient {
config: StreamClientConfig,
conn: Mutex<Option<StreamConnection>>,
}
impl StreamClient {
pub fn new(config: impl Into<StreamClientConfig>) -> Arc<Self> {
Arc::new(Self {
config: config.into(),
conn: Mutex::new(None),
})
}
pub async fn connect(&self) -> Result<StreamConnection> {
let mut conn_guard = self.conn.lock().await;
if let Some(conn) = conn_guard.as_ref() {
match conn.connection_status().await {
StreamConnectionStatus::Connected | StreamConnectionStatus::Reconnecting => {
return Ok(conn.clone());
}
StreamConnectionStatus::DisconnectInitiated
| StreamConnectionStatus::Disconnected
| StreamConnectionStatus::Failed(_) => {}
}
}
let new_conn = Arc::new(LnmStreamRepo::new(self.config.clone()).await?);
*conn_guard = Some(new_conn.clone());
Ok(new_conn)
}
pub async fn reset(&self) {
let mut conn_guard = self.conn.lock().await;
*conn_guard = None;
}
}