Skip to main content

polyester/orderbook/
subscription.rs

1//! Managed orderbook subscription (Go `orderbook.Subscription` parity).
2
3use 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
10/// Managed orderbook stream with snapshot prefetch and sequence-checked deltas.
11///
12/// Delivery contract: a full consumer queue fails the subscription with
13/// [`Error::QueueOverflow`] instead of silently dropping books.
14pub 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    /// Receiver for merged orderbook snapshots.
46    pub fn updates(&mut self) -> &mut mpsc::Receiver<OrderbookData> {
47        &mut self.rx
48    }
49
50    /// Terminal subscription error, if the stream failed (e.g. queue overflow).
51    pub fn err(&self) -> Option<Error> {
52        lock_unpoisoned(&self.last_error)
53            .clone()
54            .or_else(|| self.stream.err())
55    }
56
57    /// Register a callback for background transport, decode, snapshot, or
58    /// terminal buffering errors.
59    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    /// Change the active price bucket and re-emit the current book.
67    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    /// Refetch the REST snapshot.
77    pub async fn refresh_snapshot(&self) -> Result<()> {
78        self.stream.refresh_snapshot().await
79    }
80
81    /// Stop the subscription.
82    pub fn close(&self) {
83        self.closed.store(true, Ordering::SeqCst);
84        // Drop the sender so `updates().recv()` unblocks with None.
85        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}