use crate::errors::{Error, Result};
use crate::models::{MarketOverviewEntry, MarketOverviewList};
use crate::realtime::{SnapshotThenStream, lock_unpoisoned};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use tokio::sync::mpsc;
pub struct Subscription {
rx: mpsc::Receiver<Vec<MarketOverviewEntry>>,
stream: SnapshotThenStream<MarketOverviewList, MarketOverviewList>,
closed: Arc<AtomicBool>,
last_error: Arc<Mutex<Option<Error>>>,
tx_slot: Arc<Mutex<Option<mpsc::Sender<Vec<MarketOverviewEntry>>>>>,
}
impl Subscription {
pub(crate) fn new(
rx: mpsc::Receiver<Vec<MarketOverviewEntry>>,
stream: SnapshotThenStream<MarketOverviewList, MarketOverviewList>,
closed: Arc<AtomicBool>,
last_error: Arc<Mutex<Option<Error>>>,
tx_slot: Arc<Mutex<Option<mpsc::Sender<Vec<MarketOverviewEntry>>>>>,
) -> Self {
Self {
rx,
stream,
closed,
last_error,
tx_slot,
}
}
pub fn updates(&mut self) -> &mut mpsc::Receiver<Vec<MarketOverviewEntry>> {
&mut self.rx
}
pub fn err(&self) -> Option<Error> {
lock_unpoisoned(&self.last_error)
.clone()
.or_else(|| self.stream.err())
}
pub fn set_on_error<F>(&self, callback: F)
where
F: Fn(Error) + Send + Sync + 'static,
{
self.stream.set_on_error(callback);
}
pub async fn refresh_snapshot(&self) -> Result<()> {
self.stream.refresh_snapshot().await
}
pub fn close(&self) {
self.closed.store(true, Ordering::SeqCst);
let _ = lock_unpoisoned(&self.tx_slot).take();
self.stream.close();
}
}
impl Drop for Subscription {
fn drop(&mut self) {
self.close();
}
}