|
| 1 | +//! Oracle |
| 2 | +//! |
| 3 | +//! The Oracle service is respoinsible for reacting to all remote/on-chain events. |
| 4 | +
|
| 5 | +use { |
| 6 | + crate::agent::{ |
| 7 | + solana::{ |
| 8 | + key_store::KeyStore, |
| 9 | + network::{ |
| 10 | + Config, |
| 11 | + Network, |
| 12 | + }, |
| 13 | + }, |
| 14 | + state::oracle::{ |
| 15 | + Oracle, |
| 16 | + PricePublishingMetadata, |
| 17 | + }, |
| 18 | + }, |
| 19 | + anyhow::Result, |
| 20 | + solana_account_decoder::UiAccountEncoding, |
| 21 | + solana_client::{ |
| 22 | + nonblocking::{ |
| 23 | + pubsub_client::PubsubClient, |
| 24 | + rpc_client::RpcClient, |
| 25 | + }, |
| 26 | + rpc_config::{ |
| 27 | + RpcAccountInfoConfig, |
| 28 | + RpcProgramAccountsConfig, |
| 29 | + }, |
| 30 | + rpc_response::{ |
| 31 | + Response, |
| 32 | + RpcKeyedAccount, |
| 33 | + }, |
| 34 | + }, |
| 35 | + solana_sdk::{ |
| 36 | + account::Account, |
| 37 | + commitment_config::CommitmentConfig, |
| 38 | + pubkey::Pubkey, |
| 39 | + }, |
| 40 | + std::{ |
| 41 | + collections::HashMap, |
| 42 | + sync::Arc, |
| 43 | + }, |
| 44 | + tokio::{ |
| 45 | + select, |
| 46 | + sync::watch::Sender, |
| 47 | + }, |
| 48 | + tokio_stream::StreamExt, |
| 49 | +}; |
| 50 | + |
| 51 | +pub async fn oracle<S>( |
| 52 | + config: Config, |
| 53 | + network: Network, |
| 54 | + state: Arc<S>, |
| 55 | + publisher_permissions_tx: Sender<HashMap<Pubkey, HashMap<Pubkey, PricePublishingMetadata>>>, |
| 56 | +) where |
| 57 | + S: Oracle, |
| 58 | + S: Send + Sync + 'static, |
| 59 | +{ |
| 60 | + loop { |
| 61 | + let _ = run( |
| 62 | + config.clone(), |
| 63 | + network, |
| 64 | + state.clone(), |
| 65 | + publisher_permissions_tx.clone(), |
| 66 | + ) |
| 67 | + .await; |
| 68 | + } |
| 69 | +} |
| 70 | + |
| 71 | +async fn run<S>( |
| 72 | + config: Config, |
| 73 | + network: Network, |
| 74 | + state: Arc<S>, |
| 75 | + publisher_permissions_tx: Sender<HashMap<Pubkey, HashMap<Pubkey, PricePublishingMetadata>>>, |
| 76 | +) -> Result<()> |
| 77 | +where |
| 78 | + S: Oracle, |
| 79 | + S: Send + Sync + 'static, |
| 80 | +{ |
| 81 | + let key_store = KeyStore::new(config.key_store.clone())?; |
| 82 | + let mut poll_interval = tokio::time::interval(config.oracle.poll_interval_duration); |
| 83 | + |
| 84 | + // Setup PubsubClient to listen for account changes on the Oracle program. |
| 85 | + let client = PubsubClient::new(config.wss_url.as_str()).await?; |
| 86 | + let (mut notifier, _unsub) = { |
| 87 | + let program_key = key_store.program_key; |
| 88 | + let commitment = config.oracle.commitment; |
| 89 | + let config = RpcProgramAccountsConfig { |
| 90 | + account_config: RpcAccountInfoConfig { |
| 91 | + commitment: Some(CommitmentConfig { commitment }), |
| 92 | + encoding: Some(UiAccountEncoding::Base64Zstd), |
| 93 | + ..Default::default() |
| 94 | + }, |
| 95 | + filters: None, |
| 96 | + with_context: Some(true), |
| 97 | + }; |
| 98 | + client.program_subscribe(&program_key, Some(config)).await |
| 99 | + }?; |
| 100 | + |
| 101 | + // Setup an RpcClient for manual polling. |
| 102 | + let client = { |
| 103 | + Arc::new(RpcClient::new_with_timeout_and_commitment( |
| 104 | + config.rpc_url, |
| 105 | + config.rpc_timeout, |
| 106 | + CommitmentConfig { |
| 107 | + commitment: config.oracle.commitment, |
| 108 | + }, |
| 109 | + )) |
| 110 | + }; |
| 111 | + |
| 112 | + loop { |
| 113 | + select! { |
| 114 | + update = notifier.next() => handle_account_update(state.clone(), network, update).await, |
| 115 | + _ = poll_interval.tick() => handle_poll_updates( |
| 116 | + state.clone(), |
| 117 | + network, |
| 118 | + key_store.mapping_key, |
| 119 | + client.clone(), |
| 120 | + config.oracle.max_lookup_batch_size, |
| 121 | + config.oracle.subscriber_enabled, |
| 122 | + publisher_permissions_tx.clone(), |
| 123 | + ).await, |
| 124 | + } |
| 125 | + } |
| 126 | +} |
| 127 | + |
| 128 | +/// When an account RPC Subscription update is receiveed. |
| 129 | +/// |
| 130 | +/// We check if the account is one we're aware of and tracking, and if so, spawn |
| 131 | +/// a small background task that handles that update. We only do this for price |
| 132 | +/// accounts, all other accounts are handled below in the poller. |
| 133 | +async fn handle_account_update<S>( |
| 134 | + state: Arc<S>, |
| 135 | + network: Network, |
| 136 | + update: Option<Response<RpcKeyedAccount>>, |
| 137 | +) where |
| 138 | + S: Oracle, |
| 139 | + S: Send + Sync + 'static, |
| 140 | +{ |
| 141 | + match update { |
| 142 | + Some(update) => { |
| 143 | + match update.value.account.decode::<Account>() { |
| 144 | + Some(account) => { |
| 145 | + tokio::spawn(async move { |
| 146 | + if let Err(err) = async move { |
| 147 | + let key: Pubkey = update.value.pubkey.as_str().try_into()?; |
| 148 | + Oracle::handle_price_account_update(&*state, network, &key, &account) |
| 149 | + .await |
| 150 | + } |
| 151 | + .await |
| 152 | + { |
| 153 | + tracing::error!(err = ?err, "Failed to handle account update."); |
| 154 | + } |
| 155 | + }); |
| 156 | + } |
| 157 | + |
| 158 | + None => { |
| 159 | + tracing::error!( |
| 160 | + update = ?update, |
| 161 | + "Failed to decode account from update.", |
| 162 | + ); |
| 163 | + } |
| 164 | + }; |
| 165 | + } |
| 166 | + None => { |
| 167 | + tracing::debug!("subscriber closed connection"); |
| 168 | + } |
| 169 | + } |
| 170 | +} |
| 171 | + |
| 172 | +/// On poll lookup all Pyth Mapping/Product/Price accounts and sync. |
| 173 | +async fn handle_poll_updates<S>( |
| 174 | + state: Arc<S>, |
| 175 | + network: Network, |
| 176 | + mapping_key: Pubkey, |
| 177 | + rpc: Arc<RpcClient>, |
| 178 | + max_lookup_batch_size: usize, |
| 179 | + enabled: bool, |
| 180 | + publisher_permissions_tx: Sender<HashMap<Pubkey, HashMap<Pubkey, PricePublishingMetadata>>>, |
| 181 | +) where |
| 182 | + S: Oracle, |
| 183 | + S: Send + Sync + 'static, |
| 184 | +{ |
| 185 | + if !enabled { |
| 186 | + return; |
| 187 | + } |
| 188 | + |
| 189 | + tokio::spawn(async move { |
| 190 | + if let Err(err) = async { |
| 191 | + Oracle::poll_updates( |
| 192 | + &*state, |
| 193 | + mapping_key, |
| 194 | + &rpc, |
| 195 | + max_lookup_batch_size, |
| 196 | + publisher_permissions_tx, |
| 197 | + ) |
| 198 | + .await?; |
| 199 | + Oracle::sync_global_store(&*state, network).await |
| 200 | + } |
| 201 | + .await |
| 202 | + { |
| 203 | + tracing::error!(err = ?err, "Failed to handle poll updates."); |
| 204 | + } |
| 205 | + }); |
| 206 | +} |
0 commit comments