polyester/orderbook/
subscription.rs1use crate::errors::{Error, Result};
4use crate::models::{OrderBookDeltaUpdate, OrderbookData};
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<OrderbookData>,
16 stream: SnapshotThenStream<OrderbookData, OrderBookDeltaUpdate>,
17 closed: Arc<AtomicBool>,
18 bucket_ticks: Arc<Mutex<i64>>,
19 emit: Arc<dyn Fn() + Send + Sync>,
20 last_error: Arc<Mutex<Option<Error>>>,
21 tx_slot: Arc<Mutex<Option<mpsc::Sender<OrderbookData>>>>,
22}
23
24impl Subscription {
25 pub(crate) fn new(
26 rx: mpsc::Receiver<OrderbookData>,
27 stream: SnapshotThenStream<OrderbookData, OrderBookDeltaUpdate>,
28 closed: Arc<AtomicBool>,
29 bucket_ticks: Arc<Mutex<i64>>,
30 emit: Arc<dyn Fn() + Send + Sync>,
31 last_error: Arc<Mutex<Option<Error>>>,
32 tx_slot: Arc<Mutex<Option<mpsc::Sender<OrderbookData>>>>,
33 ) -> Self {
34 Self {
35 rx,
36 stream,
37 closed,
38 bucket_ticks,
39 emit,
40 last_error,
41 tx_slot,
42 }
43 }
44
45 pub fn updates(&mut self) -> &mut mpsc::Receiver<OrderbookData> {
47 &mut self.rx
48 }
49
50 pub fn err(&self) -> Option<Error> {
52 lock_unpoisoned(&self.last_error)
53 .clone()
54 .or_else(|| self.stream.err())
55 }
56
57 pub fn set_on_error<F>(&self, callback: F)
60 where
61 F: Fn(Error) + Send + Sync + 'static,
62 {
63 self.stream.set_on_error(callback);
64 }
65
66 pub fn set_bucket(&self, bucket: &str) -> Result<()> {
68 let ticks = crate::orderbook::parse_bucket_ticks(bucket)?;
69 *lock_unpoisoned(&self.bucket_ticks) = ticks;
70 if self.stream.is_ready() && !self.closed.load(Ordering::SeqCst) {
71 (self.emit)();
72 }
73 Ok(())
74 }
75
76 pub async fn refresh_snapshot(&self) -> Result<()> {
78 self.stream.refresh_snapshot().await
79 }
80
81 pub fn close(&self) {
83 self.closed.store(true, Ordering::SeqCst);
84 let _ = lock_unpoisoned(&self.tx_slot).take();
86 self.stream.close();
87 }
88}
89
90impl Drop for Subscription {
91 fn drop(&mut self) {
92 self.close();
93 }
94}