templar_proxy_oracle_near_contract/
lib.rs

1#![allow(clippy::needless_pass_by_value)]
2
3use std::collections::{HashMap, HashSet};
4use std::ops::{Deref, DerefMut};
5
6use near_sdk::{
7    env, json_types::Base64VecU8, near, require, AccountId, Gas, PanicOnDefault, Promise,
8    PromiseOrValue,
9};
10use near_sdk_contract_tools::{owner::Owner, Owner};
11use templar_common::{
12    oracle::{
13        lazer::ext_pyth_lazer,
14        pyth::{ext_pyth, OracleResponse, PriceIdentifier},
15        redstone::{self, ext_redstone},
16    },
17    self_ext,
18    upgrade::{UpgradeSource, MIGRATE_METHOD},
19    versioned_state::{impl_versioned_state, StateVersion, VersionedState},
20    Decimal, Nanoseconds, UnwrapReject,
21};
22use templar_proxy_oracle_kernel::proxy::{
23    circuit_breaker::{
24        CircuitBreaker, CircuitBreakerOutcome, CircuitBreakerSet, CircuitBreakerSetConfig,
25    },
26    Proxy,
27};
28use templar_proxy_oracle_near_common::{
29    cache::{bounded_resolve_error_message, CachedProxyPrice, CachedProxyPriceStatus},
30    convert::{account_id_to_kernel, pyth_price_try_from_kernel, pyth_price_try_to_kernel},
31    event::{Event, MAX_MANUAL_TRIP_METADATA_LEN},
32    governance::ProxyOracleAdminInterface,
33    input::Source,
34    proxy::has_zero_weighted_source,
35    request::OracleRequest,
36    state,
37};
38
39mod callback_handler;
40use callback_handler::{callback_result, CallbackHandler, OracleType};
41
42type State = state::v1::State;
43
44pub(crate) fn emit_outcome<T>(price_id: PriceIdentifier, outcome: CircuitBreakerOutcome<T>) -> T {
45    for event in outcome.events {
46        Event::from_kernel(price_id, event).emit();
47    }
48    outcome.value
49}
50
51#[derive(Debug, Owner, PanicOnDefault)]
52#[near(contract_state)]
53pub struct Contract {
54    pub state: VersionedState<State>,
55}
56impl_versioned_state!(Contract, State, state::migration::Migration);
57
58impl Deref for Contract {
59    type Target = State;
60
61    fn deref(&self) -> &Self::Target {
62        &self.state
63    }
64}
65
66impl DerefMut for Contract {
67    fn deref_mut(&mut self) -> &mut Self::Target {
68        &mut self.state
69    }
70}
71
72#[near]
73impl Contract {
74    pub const GAS_FOR_PYTH_REQUEST: Gas = Gas::from_tgas(16).saturating_div(10);
75    pub const GAS_FOR_PYTH_LAZER_REQUEST: Gas = Gas::from_tgas(16).saturating_div(10);
76    pub const GAS_FOR_REDSONE_REQUEST: Gas = Gas::from_tgas(17).saturating_div(10);
77    pub const GAS_FOR_MIGRATE: Gas = Gas::from_tgas(250);
78
79    /// Initialize the oracle, owned by `owner_id` if given, otherwise by the
80    /// deployer (predecessor).
81    #[init]
82    pub fn new(owner_id: Option<AccountId>) -> Self {
83        let mut self_ = Self {
84            state: State::new(()),
85        };
86
87        let owner = owner_id.unwrap_or_else(env::predecessor_account_id);
88
89        Owner::init(&mut self_, &owner);
90
91        self_
92    }
93
94    // View methods
95
96    pub fn list_proxies(&self, offset: Option<u32>, count: Option<u32>) -> Vec<PriceIdentifier> {
97        self.state.list_proxies(offset, count)
98    }
99
100    pub fn get_proxy(&self, id: PriceIdentifier) -> Option<Proxy<Source>> {
101        self.state.get_proxy(id)
102    }
103
104    pub fn get_proxy_circuit_breaker_set(&self, id: PriceIdentifier) -> Option<CircuitBreakerSet> {
105        self.state.get_proxy_circuit_breaker_set(id)
106    }
107
108    pub fn get_cached_proxy_price(&self, id: PriceIdentifier) -> Option<CachedProxyPrice> {
109        self.state.get_cached_proxy_price(id)
110    }
111
112    pub fn list_cached_proxy_prices(
113        &self,
114        price_ids: Vec<PriceIdentifier>,
115    ) -> HashMap<PriceIdentifier, Option<CachedProxyPrice>> {
116        self.state.list_cached_proxy_prices(price_ids)
117    }
118
119    // Pyth interface
120
121    pub fn price_feed_exists(&self, price_identifier: PriceIdentifier) -> bool {
122        self.state.proxy_exists(&price_identifier)
123    }
124
125    pub const GAS_FOR_LIST_00_ENTRY: Gas = Gas::from_tgas(35).saturating_div(10);
126    pub fn list_ema_prices_no_older_than(
127        &self,
128        price_ids: Vec<PriceIdentifier>,
129        age: u64,
130    ) -> OracleResponse {
131        if price_ids.is_empty() {
132            return OracleResponse::new();
133        }
134
135        let max_age = Nanoseconds::from_secs(age);
136        let now = Nanoseconds::near_timestamp();
137        let mut results = OracleResponse::new();
138
139        for price_id in HashSet::<PriceIdentifier>::from_iter(price_ids) {
140            if !self.state.proxy_exists(&price_id) {
141                continue;
142            }
143
144            let price = self
145                .state
146                .get_cached_proxy_price(price_id)
147                .and_then(|cached| {
148                    cached
149                        .accepted_price_no_older_than(now, max_age)
150                        .and_then(pyth_price_try_from_kernel)
151                });
152            results.insert(price_id, price);
153        }
154
155        results
156    }
157
158    pub fn update_prices(
159        &mut self,
160        price_ids: Vec<PriceIdentifier>,
161    ) -> PromiseOrValue<HashMap<PriceIdentifier, CachedProxyPriceStatus>> {
162        if price_ids.is_empty() {
163            return PromiseOrValue::Value(HashMap::new());
164        }
165        let price_ids = HashSet::<PriceIdentifier>::from_iter(price_ids);
166
167        let mut invoked = Vec::with_capacity(price_ids.len());
168        let mut pyth_requests =
169            HashMap::<AccountId, HashSet<PriceIdentifier>>::with_capacity(price_ids.len());
170        let mut redstone_requests =
171            HashMap::<AccountId, HashSet<redstone::FeedId>>::with_capacity(price_ids.len());
172        let mut lazer_requests = HashMap::<AccountId, HashSet<u32>>::with_capacity(price_ids.len());
173        let mut transformer_promises = Vec::with_capacity(price_ids.len());
174        let now = Nanoseconds::near_timestamp();
175        let mut results = HashMap::new();
176
177        for price_id in &price_ids {
178            let proxy = match self.state.validated_proxy_entry(*price_id) {
179                Ok(Some(proxy)) => proxy,
180                Ok(None) => continue,
181                Err(error) => {
182                    let message = bounded_resolve_error_message(error.to_string());
183                    let status = self
184                        .state
185                        .cache_price_update_failure(*price_id, now, message);
186                    results.insert(*price_id, status);
187                    continue;
188                }
189            };
190            let pending = proxy.prepare_price_update();
191
192            invoked.push(pending.clone());
193
194            for source in pending.proxy.sources() {
195                let request = match source {
196                    Source::Request(request) => request,
197                    Source::Transformer(transformer) => {
198                        transformer_promises.push(transformer.call.promise());
199                        &transformer.request
200                    }
201                };
202
203                match request {
204                    OracleRequest::Pyth(p) => {
205                        pyth_requests
206                            .entry(p.oracle_id.clone())
207                            .or_default()
208                            .insert(p.price_id);
209                    }
210                    OracleRequest::RedStone(p) => {
211                        redstone_requests
212                            .entry(p.oracle_id.clone())
213                            .or_default()
214                            .insert(p.price_id.clone());
215                    }
216                    OracleRequest::Lazer(p) => {
217                        lazer_requests
218                            .entry(p.oracle_id.clone())
219                            .or_default()
220                            .insert(p.feed_id);
221                    }
222                }
223            }
224        }
225
226        let capacity = pyth_requests.len() + redstone_requests.len() + lazer_requests.len();
227        let mut oracle_order = Vec::with_capacity(capacity);
228        let mut oracle_promises = Vec::with_capacity(capacity);
229
230        for (oracle_id, price_ids) in pyth_requests {
231            oracle_order.push(OracleType::Pyth(oracle_id.clone()));
232            oracle_promises.push(
233                ext_pyth::ext(oracle_id)
234                    .with_static_gas(Self::GAS_FOR_PYTH_REQUEST)
235                    .list_ema_prices_unsafe(Vec::from_iter(price_ids)),
236            );
237        }
238
239        for (oracle_id, price_ids) in redstone_requests {
240            oracle_order.push(OracleType::RedStone(oracle_id.clone()));
241            oracle_promises.push(
242                ext_redstone::ext(oracle_id)
243                    .with_static_gas(Self::GAS_FOR_REDSONE_REQUEST)
244                    .read_price_data(Vec::from_iter(price_ids)),
245            );
246        }
247
248        for (oracle_id, feed_ids) in lazer_requests {
249            oracle_order.push(OracleType::Lazer(oracle_id.clone()));
250            oracle_promises.push(
251                ext_pyth_lazer::ext(oracle_id)
252                    .with_static_gas(Self::GAS_FOR_PYTH_LAZER_REQUEST)
253                    .get_feeds_data(Vec::from_iter(feed_ids)),
254            );
255        }
256
257        let Some(promise) = oracle_promises
258            .into_iter()
259            .chain(transformer_promises)
260            .reduce(near_sdk::Promise::and)
261        else {
262            return PromiseOrValue::Value(results);
263        };
264
265        PromiseOrValue::Promise(promise.then(
266            self_ext!(Self::GAS_FOR_UPDATE_01_CALLBACK).update_prices_01_consume_results(
267                oracle_order,
268                invoked,
269                results,
270            ),
271        ))
272    }
273
274    pub const GAS_FOR_UPDATE_01_CALLBACK: Gas = Gas::from_tgas(10);
275    #[private]
276    #[allow(
277        unused_mut,
278        reason = "near macro expansion checks the original binding"
279    )]
280    pub fn update_prices_01_consume_results(
281        &mut self,
282        oracle_order: Vec<OracleType>,
283        invoked: Vec<state::v1::PendingProxyPriceUpdate>,
284        mut results: HashMap<PriceIdentifier, CachedProxyPriceStatus>,
285    ) -> HashMap<PriceIdentifier, CachedProxyPriceStatus> {
286        let callback = CallbackHandler::new(&oracle_order);
287
288        let now = Nanoseconds::near_timestamp();
289
290        let mut i = oracle_order.len() as u64;
291        for pending in invoked {
292            let price_id = pending.price_id;
293            let mut prices = vec![];
294
295            for source in pending.proxy.sources() {
296                let source_result = match source {
297                    Source::Transformer(transformer) => {
298                        let price = callback.get(&transformer.request);
299                        let input = callback_result::<Decimal>(i);
300                        i += 1;
301
302                        price
303                            .zip(input)
304                            .and_then(|(price, input)| transformer.action.apply(price, input))
305                    }
306                    Source::Request(request) => callback.get(request),
307                };
308
309                prices.push(source_result.as_ref().and_then(pyth_price_try_to_kernel));
310            }
311
312            if let Some(status) =
313                self.state
314                    .finish_price_update_if_current(pending, now, |proxy, set| {
315                        match proxy.resolve(set, prices, now) {
316                            Ok(resolution) => match emit_outcome(price_id, resolution) {
317                                Ok(price) => CachedProxyPriceStatus::Accepted { price },
318                                Err(reason) => CachedProxyPriceStatus::Blocked { reason },
319                            },
320                            Err(error) => {
321                                let message = bounded_resolve_error_message(error.to_string());
322                                near_sdk::log!(
323                                    "Proxy resolve failed price_id={:?} error={}",
324                                    price_id,
325                                    message
326                                );
327                                CachedProxyPriceStatus::ResolveFailed { message }
328                            }
329                        }
330                    })
331            {
332                results.insert(price_id, status);
333            }
334        }
335
336        results
337    }
338}
339
340#[near]
341impl ProxyOracleAdminInterface for Contract {
342    fn admin_set_proxy(&mut self, id: PriceIdentifier, proxy: Option<Proxy<Source>>) {
343        self.assert_owner();
344        if proxy.as_ref().is_some_and(has_zero_weighted_source) {
345            env::panic_str("Proxy sources must have positive weights");
346        }
347        self.state.set_proxy(id, proxy);
348    }
349
350    fn admin_configure_circuit_breakers(
351        &mut self,
352        id: PriceIdentifier,
353        config: CircuitBreakerSetConfig,
354    ) {
355        self.assert_owner();
356        let result = self
357            .state
358            .proxy_entry_mut(id)
359            .unwrap_or_else(|| env::panic_str("Proxy not found"))
360            .configure_circuit_breakers(config)
361            .unwrap_or_else(|error| env::panic_str(&error.to_string()));
362        emit_outcome(id, result);
363    }
364
365    fn admin_add_circuit_breaker(
366        &mut self,
367        id: PriceIdentifier,
368        breaker_id: u32,
369        breaker: CircuitBreaker,
370    ) {
371        self.assert_owner();
372        let result = self
373            .state
374            .proxy_entry_mut(id)
375            .unwrap_or_else(|| env::panic_str("Proxy not found"))
376            .add_circuit_breaker(breaker_id, breaker)
377            .unwrap_or_reject();
378        emit_outcome(id, result);
379    }
380
381    fn admin_remove_circuit_breaker(&mut self, id: PriceIdentifier, breaker_id: u32) {
382        self.assert_owner();
383        let result = self
384            .state
385            .proxy_entry_mut(id)
386            .unwrap_or_else(|| env::panic_str("Proxy not found"))
387            .remove_circuit_breaker(breaker_id)
388            .unwrap_or_reject();
389        emit_outcome(id, result);
390    }
391
392    fn admin_set_manual_trip(
393        &mut self,
394        id: PriceIdentifier,
395        is_manually_tripped: bool,
396        metadata: Option<Base64VecU8>,
397    ) {
398        self.assert_owner();
399
400        require!(
401            metadata
402                .as_ref()
403                .is_none_or(|metadata| metadata.0.len() <= MAX_MANUAL_TRIP_METADATA_LEN),
404            "Manual trip metadata is too long"
405        );
406        let result = self
407            .state
408            .proxy_entry_mut(id)
409            .unwrap_or_else(|| env::panic_str("Proxy not found"))
410            .set_circuit_breaker_manual_trip(
411                is_manually_tripped,
412                account_id_to_kernel(env::predecessor_account_id().as_ref()),
413                metadata.map(|metadata| metadata.0),
414            );
415        if result.events.is_empty() {
416            return;
417        }
418
419        emit_outcome(id, result);
420    }
421
422    fn admin_rearm(&mut self, id: PriceIdentifier, breaker_id: u32, arming_delay_ns: Nanoseconds) {
423        self.assert_owner();
424
425        let armed_after_ns = Nanoseconds::from_ns(
426            Nanoseconds::near_timestamp()
427                .as_ns()
428                .checked_add(arming_delay_ns.as_ns())
429                .unwrap_or_else(|| env::panic_str("Rearm delay overflows ledger timestamp")),
430        );
431        let result = self
432            .state
433            .proxy_entry_mut(id)
434            .unwrap_or_else(|| env::panic_str("Proxy not found"))
435            .rearm(breaker_id, armed_after_ns)
436            .unwrap_or_reject();
437        emit_outcome(id, result);
438    }
439
440    fn admin_set_enforced(&mut self, id: PriceIdentifier, breaker_id: u32, is_enforced: bool) {
441        self.assert_owner();
442
443        let result = self
444            .state
445            .proxy_entry_mut(id)
446            .unwrap_or_else(|| env::panic_str("Proxy not found"))
447            .set_enforced(breaker_id, is_enforced)
448            .unwrap_or_reject();
449        emit_outcome(id, result);
450    }
451
452    fn admin_upgrade(&mut self, code: UpgradeSource, migrate_args: Base64VecU8) -> Promise {
453        self.assert_owner();
454        code.deploy_and_migrate(MIGRATE_METHOD, migrate_args, Self::GAS_FOR_MIGRATE)
455    }
456}
457
458#[cfg(target_arch = "wasm32")]
459mod custom_getrandom {
460    #![allow(clippy::no_mangle_with_rust_abi)]
461
462    use getrandom::{register_custom_getrandom, Error};
463    use near_sdk::env;
464
465    register_custom_getrandom!(custom_getrandom);
466
467    #[allow(clippy::unnecessary_wraps)]
468    pub fn custom_getrandom(buf: &mut [u8]) -> Result<(), Error> {
469        buf.copy_from_slice(&env::random_seed_array());
470        Ok(())
471    }
472}