polyester/marketoverview/
subscription.rs1use crate::errors::{Error, Result};
4use crate::models::{MarketOverviewEntry, MarketOverviewList};
5use crate::realtime::{SnapshotThenStream, lock_unpoisoned};
6use std::sync::atomic::{AtomicBool, Ordering};
7use std::sync::{Arc, Mutex};
8use tokio::sync::mpsc;
9
10pub struct Subscription {
15 rx: mpsc::Receiver<Vec<MarketOverviewEntry>>,
16 stream: SnapshotThenStream<MarketOverviewList, MarketOverviewList>,
17 closed: Arc<AtomicBool>,
18 last_error: Arc<Mutex<Option<Error>>>,
19 tx_slot: Arc<Mutex<Option<mpsc::Sender<Vec<MarketOverviewEntry>>>>>,
20}
21
22impl Subscription {
23 pub(crate) fn new(
24 rx: mpsc::Receiver<Vec<MarketOverviewEntry>>,
25 stream: SnapshotThenStream<MarketOverviewList, MarketOverviewList>,
26 closed: Arc<AtomicBool>,
27 last_error: Arc<Mutex<Option<Error>>>,
28 tx_slot: Arc<Mutex<Option<mpsc::Sender<Vec<MarketOverviewEntry>>>>>,
29 ) -> Self {
30 Self {
31 rx,
32 stream,
33 closed,
34 last_error,
35 tx_slot,
36 }
37 }
38
39 pub fn updates(&mut self) -> &mut mpsc::Receiver<Vec<MarketOverviewEntry>> {
41 &mut self.rx
42 }
43
44 pub fn err(&self) -> Option<Error> {
46 lock_unpoisoned(&self.last_error)
47 .clone()
48 .or_else(|| self.stream.err())
49 }
50
51 pub fn set_on_error<F>(&self, callback: F)
54 where
55 F: Fn(Error) + Send + Sync + 'static,
56 {
57 self.stream.set_on_error(callback);
58 }
59
60 pub async fn refresh_snapshot(&self) -> Result<()> {
62 self.stream.refresh_snapshot().await
63 }
64
65 pub fn close(&self) {
67 self.closed.store(true, Ordering::SeqCst);
68 let _ = lock_unpoisoned(&self.tx_slot).take();
69 self.stream.close();
70 }
71}
72
73impl Drop for Subscription {
74 fn drop(&mut self) {
75 self.close();
76 }
77}