Skip to main content

polyester/marketoverview/
subscription.rs

1//! Managed market overview subscription (Go `marketoverview.Subscription` parity).
2
3use 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
10/// Managed market-overview stream with snapshot prefetch and live merge.
11///
12/// Delivery contract: a full consumer queue fails the subscription with
13/// [`Error::QueueOverflow`] instead of silently dropping rows.
14pub 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    /// Receiver for merged overview rows.
40    pub fn updates(&mut self) -> &mut mpsc::Receiver<Vec<MarketOverviewEntry>> {
41        &mut self.rx
42    }
43
44    /// Terminal subscription error, if the stream failed (e.g. queue overflow).
45    pub fn err(&self) -> Option<Error> {
46        lock_unpoisoned(&self.last_error)
47            .clone()
48            .or_else(|| self.stream.err())
49    }
50
51    /// Register a callback for background transport, decode, snapshot, or
52    /// terminal buffering errors.
53    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    /// Refetch the REST snapshot.
61    pub async fn refresh_snapshot(&self) -> Result<()> {
62        self.stream.refresh_snapshot().await
63    }
64
65    /// Stop the subscription.
66    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}