templar_proxy_oracle_near_contract/
lib.rs1#![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 #[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 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 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}