Skip to main content

drift_rs/
auction_subscriber.rs

1use std::sync::Mutex;
2
3use solana_account_decoder_client_types::UiAccountEncoding;
4use solana_sdk::commitment_config::CommitmentConfig;
5
6use crate::{
7    drift_idl::accounts::User,
8    memcmp::{get_user_filter, get_user_with_auction_filter},
9    types::SdkResult,
10    websocket_program_account_subscriber::{
11        ProgramAccountUpdate, WebsocketProgramAccountOptions, WebsocketProgramAccountSubscriber,
12    },
13    SdkError, UnsubHandle,
14};
15
16pub struct AuctionSubscriberConfig {
17    pub commitment: CommitmentConfig,
18    pub resub_timeout_ms: Option<u64>,
19    pub url: String,
20}
21
22/// Subscribes to all user auction events across all markets
23///
24/// DEV: take care it is not dropped or the Auction stream will unsubscribe
25pub struct AuctionSubscriber {
26    subscriber: WebsocketProgramAccountSubscriber,
27    unsub: Mutex<Option<UnsubHandle>>,
28}
29
30impl AuctionSubscriber {
31    pub const SUBSCRIPTION_ID: &'static str = "auction";
32
33    pub fn new(config: AuctionSubscriberConfig) -> Self {
34        let filters = vec![get_user_filter(), get_user_with_auction_filter()];
35        let websocket_options = WebsocketProgramAccountOptions {
36            filters,
37            commitment: config.commitment,
38            encoding: UiAccountEncoding::Base64Zstd,
39        };
40
41        Self {
42            subscriber: WebsocketProgramAccountSubscriber::new(config.url, websocket_options),
43            unsub: Mutex::new(None),
44        }
45    }
46
47    /// Start the auction subscription task
48    ///
49    /// * `handler_fn` - fn to invoke on each update
50    ///
51    /// this class sends the entire User account, the callback is required to
52    /// interpret the diff e.g orders added/removed
53    ///
54    pub fn subscribe<F>(&self, handler_fn: F)
55    where
56        F: 'static + Send + Fn(&ProgramAccountUpdate<User>),
57    {
58        let mut guard = self.unsub.try_lock().expect("uncontested");
59        let unsub = self.subscriber.subscribe(Self::SUBSCRIPTION_ID, handler_fn);
60        guard.replace(unsub);
61    }
62
63    /// Unsubscribe stopping the auction subscription task
64    pub fn unsubscribe(self) -> SdkResult<()> {
65        let mut guard = self.unsub.lock().expect("acquired");
66        if let Some(unsub) = guard.take() {
67            if unsub.send(()).is_err() {
68                log::error!("unsub failed");
69                return Err(SdkError::CouldntUnsubscribe);
70            }
71        }
72
73        Ok(())
74    }
75}
76
77#[cfg(feature = "rpc_tests")]
78mod tests {
79    use super::*;
80    use crate::utils::test_envs::mainnet_endpoint;
81
82    #[tokio::test]
83    async fn test_auction_subscriber() {
84        env_logger::init();
85
86        let config = AuctionSubscriberConfig {
87            commitment: CommitmentConfig::confirmed(),
88            resub_timeout_ms: None,
89            url: mainnet_endpoint(),
90        };
91
92        let mut auction_subscriber = AuctionSubscriber::new(config);
93
94        let emitter = auction_subscriber.event_emitter.clone();
95
96        emitter.subscribe(move |event| {
97            log::info!("{:?}", event.now.elapsed());
98        });
99
100        let _ = auction_subscriber.subscribe().await;
101
102        tokio::time::sleep(tokio::time::Duration::from_secs(60)).await;
103
104        let _ = auction_subscriber.unsubscribe().await;
105
106        tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
107    }
108}