1use std::{collections::BTreeSet, sync::Arc};
2
3use alloy_primitives::{Address, I256, U256, keccak256};
4use evm_fork_cache::{StateUpdate, StateView};
5use thiserror::Error;
6
7use crate::{FeedRegistration, OraclePriceUpdate};
8
9const OCR2_HOT_VARS_SLOT: u64 = 11;
10const OCR2_TRANSMISSIONS_SLOT: u64 = 12;
11const OCR2_LATEST_EPOCH_AND_ROUND_OFFSET_BITS: usize = 8;
12const OCR2_LATEST_AGGREGATOR_ROUND_ID_OFFSET_BITS: usize = 6 * 8;
13const OCR2_TRANSMISSION_OBSERVATIONS_TIMESTAMP_OFFSET_BITS: usize = 24 * 8;
14const OCR2_TRANSMISSION_TIMESTAMP_OFFSET_BITS: usize = 28 * 8;
15const OCR1_HOT_VARS_SLOT: u64 = 43;
16const OCR1_TRANSMISSIONS_SLOT: u64 = 44;
17const OCR1_LATEST_EPOCH_AND_ROUND_OFFSET_BITS: usize = 16 * 8;
18const OCR1_LATEST_AGGREGATOR_ROUND_ID_OFFSET_BITS: usize = 22 * 8;
19const OCR1_TRANSMISSION_TIMESTAMP_OFFSET_BITS: usize = 24 * 8;
20const UINT32_BITS: usize = 32;
21const UINT40_BITS: usize = 40;
22const INT192_BITS: usize = 192;
23
24#[derive(Debug, Error)]
26pub enum OracleStorageError {
27 #[error("oracle storage adapter `{adapter}` failed: {message}")]
29 Adapter {
30 adapter: &'static str,
32 message: String,
34 },
35}
36
37impl OracleStorageError {
38 fn adapter(adapter: &'static str, message: impl Into<String>) -> Self {
39 Self::Adapter {
40 adapter,
41 message: message.into(),
42 }
43 }
44}
45
46#[derive(Clone, Debug, PartialEq, Eq)]
48pub enum OracleStorageEffect {
49 StateUpdates {
51 adapter: &'static str,
53 updates: Vec<StateUpdate>,
55 },
56 FallbackPurge,
58}
59
60#[derive(Clone, Debug, PartialEq, Eq)]
62pub struct Ocr2TransmissionStorageUpdate {
63 pub aggregator: Address,
65 pub aggregator_round_id: U256,
67 pub answer: I256,
69 pub observations_timestamp: u64,
71 pub transmission_timestamp: u64,
73 pub epoch_and_round: u64,
75}
76
77#[derive(Clone, Debug, PartialEq, Eq)]
79pub struct Ocr1TransmissionStorageUpdate {
80 pub aggregator: Address,
82 pub aggregator_round_id: U256,
84 pub answer: I256,
86 pub transmission_timestamp: u64,
88 pub epoch_and_round: u64,
90}
91
92pub trait OracleStorageAdapter: Send + Sync {
98 fn name(&self) -> &'static str;
100
101 fn warm_slots_for_registration(
107 &self,
108 _registration: &FeedRegistration,
109 ) -> Vec<(Address, U256)> {
110 Vec::new()
111 }
112
113 fn state_updates_for_answer(
115 &self,
116 registration: &FeedRegistration,
117 update: &OraclePriceUpdate,
118 state: &dyn StateView,
119 ) -> Result<Option<Vec<StateUpdate>>, OracleStorageError>;
120
121 fn state_updates_for_ocr2_transmission(
123 &self,
124 _registration: &FeedRegistration,
125 _update: &Ocr2TransmissionStorageUpdate,
126 _state: &dyn StateView,
127 ) -> Result<Option<Vec<StateUpdate>>, OracleStorageError> {
128 Ok(None)
129 }
130
131 fn state_updates_for_ocr1_transmission(
133 &self,
134 _registration: &FeedRegistration,
135 _update: &Ocr1TransmissionStorageUpdate,
136 _state: &dyn StateView,
137 ) -> Result<Option<Vec<StateUpdate>>, OracleStorageError> {
138 Ok(None)
139 }
140
141 fn prefers_ocr2_new_transmission(&self, _registration: &FeedRegistration) -> bool {
143 false
144 }
145
146 fn prefers_ocr1_new_transmission(&self, _registration: &FeedRegistration) -> bool {
148 false
149 }
150}
151
152#[derive(Clone, Debug, Default)]
166pub struct ChainlinkOcr2StorageAdapter;
167
168impl ChainlinkOcr2StorageAdapter {
169 const NAME: &'static str = "chainlink-ocr2";
170
171 pub fn new() -> Self {
173 Self
174 }
175
176 pub fn hot_vars_slot() -> U256 {
178 U256::from(OCR2_HOT_VARS_SLOT)
179 }
180
181 pub fn latest_aggregator_round_id_mask() -> U256 {
183 uint_mask(UINT32_BITS) << OCR2_LATEST_AGGREGATOR_ROUND_ID_OFFSET_BITS
184 }
185
186 pub fn latest_epoch_and_round_mask() -> U256 {
188 uint_mask(UINT40_BITS) << OCR2_LATEST_EPOCH_AND_ROUND_OFFSET_BITS
189 }
190
191 pub fn latest_aggregator_round_id_value(round_id: u32) -> U256 {
193 U256::from(round_id) << OCR2_LATEST_AGGREGATOR_ROUND_ID_OFFSET_BITS
194 }
195
196 pub fn latest_epoch_and_round_value(epoch_and_round: u64) -> Result<U256, OracleStorageError> {
198 if epoch_and_round > uint_mask_u64(UINT40_BITS) {
199 return Err(OracleStorageError::adapter(
200 Self::NAME,
201 format!("OCR2 epochAndRound {epoch_and_round} does not fit uint40"),
202 ));
203 }
204
205 Ok(U256::from(epoch_and_round) << OCR2_LATEST_EPOCH_AND_ROUND_OFFSET_BITS)
206 }
207
208 pub fn transmission_slot(round_id: u32) -> U256 {
210 mapping_slot(U256::from(round_id), U256::from(OCR2_TRANSMISSIONS_SLOT))
211 }
212
213 pub fn pack_transmission_from_event(
219 answer: I256,
220 updated_at: u64,
221 ) -> Result<U256, OracleStorageError> {
222 let timestamp = u32::try_from(updated_at).map_err(|_| {
223 OracleStorageError::adapter(
224 Self::NAME,
225 format!("OCR2 timestamp {updated_at} does not fit uint32"),
226 )
227 })?;
228 Self::pack_transmission(answer, timestamp, timestamp)
229 }
230
231 pub fn pack_transmission(
233 answer: I256,
234 observations_timestamp: u32,
235 transmission_timestamp: u32,
236 ) -> Result<U256, OracleStorageError> {
237 let answer = encode_int192(answer)?;
238 Ok(answer
239 | (U256::from(observations_timestamp)
240 << OCR2_TRANSMISSION_OBSERVATIONS_TIMESTAMP_OFFSET_BITS)
241 | (U256::from(transmission_timestamp) << OCR2_TRANSMISSION_TIMESTAMP_OFFSET_BITS))
242 }
243}
244
245impl OracleStorageAdapter for ChainlinkOcr2StorageAdapter {
246 fn name(&self) -> &'static str {
247 Self::NAME
248 }
249
250 fn warm_slots_for_registration(&self, registration: &FeedRegistration) -> Vec<(Address, U256)> {
251 let Some(aggregator) = registration.current_aggregator else {
252 return Vec::new();
253 };
254 if registration_supports_ocr2(registration, aggregator) {
255 vec![(aggregator, Self::hot_vars_slot())]
256 } else {
257 Vec::new()
258 }
259 }
260
261 fn state_updates_for_answer(
262 &self,
263 registration: &FeedRegistration,
264 update: &OraclePriceUpdate,
265 state: &dyn StateView,
266 ) -> Result<Option<Vec<StateUpdate>>, OracleStorageError> {
267 let Some(aggregator) = registration.current_aggregator else {
268 return Ok(None);
269 };
270 if update.aggregator != aggregator || !registration_supports_ocr2(registration, aggregator)
271 {
272 return Ok(None);
273 }
274 if state.storage(aggregator, Self::hot_vars_slot()).is_none() {
275 return Ok(None);
276 }
277
278 let round_id = u32::try_from(update.event_round_id).map_err(|_| {
279 OracleStorageError::adapter(
280 Self::NAME,
281 format!(
282 "OCR2 round id {} does not fit aggregator-local uint32",
283 update.event_round_id
284 ),
285 )
286 })?;
287 let transmission =
288 Self::pack_transmission_from_event(update.raw_answer, update.updated_at)?;
289
290 Ok(Some(vec![
291 StateUpdate::slot_masked(
292 aggregator,
293 Self::hot_vars_slot(),
294 Self::latest_aggregator_round_id_mask(),
295 Self::latest_aggregator_round_id_value(round_id),
296 ),
297 StateUpdate::slot(aggregator, Self::transmission_slot(round_id), transmission),
298 ]))
299 }
300
301 fn state_updates_for_ocr2_transmission(
302 &self,
303 registration: &FeedRegistration,
304 update: &Ocr2TransmissionStorageUpdate,
305 state: &dyn StateView,
306 ) -> Result<Option<Vec<StateUpdate>>, OracleStorageError> {
307 let Some(aggregator) = registration.current_aggregator else {
308 return Ok(None);
309 };
310 if update.aggregator != aggregator || !registration_supports_ocr2(registration, aggregator)
311 {
312 return Ok(None);
313 }
314 if state.storage(aggregator, Self::hot_vars_slot()).is_none() {
315 return Ok(None);
316 }
317
318 let round_id = u32::try_from(update.aggregator_round_id).map_err(|_| {
319 OracleStorageError::adapter(
320 Self::NAME,
321 format!(
322 "OCR2 round id {} does not fit aggregator-local uint32",
323 update.aggregator_round_id
324 ),
325 )
326 })?;
327 let observations_timestamp =
328 u32::try_from(update.observations_timestamp).map_err(|_| {
329 OracleStorageError::adapter(
330 Self::NAME,
331 format!(
332 "OCR2 observations timestamp {} does not fit uint32",
333 update.observations_timestamp
334 ),
335 )
336 })?;
337 let transmission_timestamp =
338 u32::try_from(update.transmission_timestamp).map_err(|_| {
339 OracleStorageError::adapter(
340 Self::NAME,
341 format!(
342 "OCR2 transmission timestamp {} does not fit uint32",
343 update.transmission_timestamp
344 ),
345 )
346 })?;
347 let hot_vars_mask =
348 Self::latest_epoch_and_round_mask() | Self::latest_aggregator_round_id_mask();
349 let hot_vars_value = Self::latest_epoch_and_round_value(update.epoch_and_round)?
350 | Self::latest_aggregator_round_id_value(round_id);
351 let transmission = Self::pack_transmission(
352 update.answer,
353 observations_timestamp,
354 transmission_timestamp,
355 )?;
356
357 Ok(Some(vec![
358 StateUpdate::slot_masked(
359 aggregator,
360 Self::hot_vars_slot(),
361 hot_vars_mask,
362 hot_vars_value,
363 ),
364 StateUpdate::slot(aggregator, Self::transmission_slot(round_id), transmission),
365 ]))
366 }
367
368 fn prefers_ocr2_new_transmission(&self, registration: &FeedRegistration) -> bool {
369 registration
370 .current_aggregator
371 .is_some_and(|aggregator| registration_supports_ocr2(registration, aggregator))
372 }
373}
374
375fn registration_supports_ocr2(registration: &FeedRegistration, aggregator: Address) -> bool {
376 registration
377 .aggregator_layout
378 .as_ref()
379 .is_some_and(|layout| layout.aggregator == aggregator && layout.is_chainlink_ocr2_v1())
380}
381
382#[derive(Clone, Debug, Default)]
396pub struct ChainlinkOcr1StorageAdapter;
397
398impl ChainlinkOcr1StorageAdapter {
399 const NAME: &'static str = "chainlink-ocr1";
400
401 pub fn new() -> Self {
403 Self
404 }
405
406 pub fn hot_vars_slot() -> U256 {
408 U256::from(OCR1_HOT_VARS_SLOT)
409 }
410
411 pub fn latest_aggregator_round_id_mask() -> U256 {
413 uint_mask(UINT32_BITS) << OCR1_LATEST_AGGREGATOR_ROUND_ID_OFFSET_BITS
414 }
415
416 pub fn latest_epoch_and_round_mask() -> U256 {
418 uint_mask(UINT40_BITS) << OCR1_LATEST_EPOCH_AND_ROUND_OFFSET_BITS
419 }
420
421 pub fn latest_aggregator_round_id_value(round_id: u32) -> U256 {
423 U256::from(round_id) << OCR1_LATEST_AGGREGATOR_ROUND_ID_OFFSET_BITS
424 }
425
426 pub fn latest_epoch_and_round_value(epoch_and_round: u64) -> Result<U256, OracleStorageError> {
428 if epoch_and_round > uint_mask_u64(UINT40_BITS) {
429 return Err(OracleStorageError::adapter(
430 Self::NAME,
431 format!("OCR1 epochAndRound {epoch_and_round} does not fit uint40"),
432 ));
433 }
434
435 Ok(U256::from(epoch_and_round) << OCR1_LATEST_EPOCH_AND_ROUND_OFFSET_BITS)
436 }
437
438 pub fn transmission_slot(round_id: u32) -> U256 {
440 mapping_slot(U256::from(round_id), U256::from(OCR1_TRANSMISSIONS_SLOT))
441 }
442
443 pub fn pack_transmission(
445 answer: I256,
446 transmission_timestamp: u64,
447 ) -> Result<U256, OracleStorageError> {
448 let answer = encode_int192_with_adapter(answer, Self::NAME, "OCR1")?;
449 Ok(
450 answer
451 | (U256::from(transmission_timestamp) << OCR1_TRANSMISSION_TIMESTAMP_OFFSET_BITS),
452 )
453 }
454}
455
456impl OracleStorageAdapter for ChainlinkOcr1StorageAdapter {
457 fn name(&self) -> &'static str {
458 Self::NAME
459 }
460
461 fn warm_slots_for_registration(&self, registration: &FeedRegistration) -> Vec<(Address, U256)> {
462 let Some(aggregator) = registration.current_aggregator else {
463 return Vec::new();
464 };
465 if registration_supports_ocr1(registration, aggregator) {
466 vec![(aggregator, Self::hot_vars_slot())]
467 } else {
468 Vec::new()
469 }
470 }
471
472 fn state_updates_for_answer(
473 &self,
474 _registration: &FeedRegistration,
475 _update: &OraclePriceUpdate,
476 _state: &dyn StateView,
477 ) -> Result<Option<Vec<StateUpdate>>, OracleStorageError> {
478 Ok(None)
479 }
480
481 fn state_updates_for_ocr1_transmission(
482 &self,
483 registration: &FeedRegistration,
484 update: &Ocr1TransmissionStorageUpdate,
485 state: &dyn StateView,
486 ) -> Result<Option<Vec<StateUpdate>>, OracleStorageError> {
487 let Some(aggregator) = registration.current_aggregator else {
488 return Ok(None);
489 };
490 if update.aggregator != aggregator || !registration_supports_ocr1(registration, aggregator)
491 {
492 return Ok(None);
493 }
494 if state.storage(aggregator, Self::hot_vars_slot()).is_none() {
495 return Ok(None);
496 }
497
498 let round_id = u32::try_from(update.aggregator_round_id).map_err(|_| {
499 OracleStorageError::adapter(
500 Self::NAME,
501 format!(
502 "OCR1 round id {} does not fit aggregator-local uint32",
503 update.aggregator_round_id
504 ),
505 )
506 })?;
507 let hot_vars_mask =
508 Self::latest_epoch_and_round_mask() | Self::latest_aggregator_round_id_mask();
509 let hot_vars_value = Self::latest_epoch_and_round_value(update.epoch_and_round)?
510 | Self::latest_aggregator_round_id_value(round_id);
511 let transmission = Self::pack_transmission(update.answer, update.transmission_timestamp)?;
512
513 Ok(Some(vec![
514 StateUpdate::slot_masked(
515 aggregator,
516 Self::hot_vars_slot(),
517 hot_vars_mask,
518 hot_vars_value,
519 ),
520 StateUpdate::slot(aggregator, Self::transmission_slot(round_id), transmission),
521 ]))
522 }
523
524 fn prefers_ocr1_new_transmission(&self, registration: &FeedRegistration) -> bool {
525 registration
526 .current_aggregator
527 .is_some_and(|aggregator| registration_supports_ocr1(registration, aggregator))
528 }
529}
530
531fn registration_supports_ocr1(registration: &FeedRegistration, aggregator: Address) -> bool {
532 registration
533 .aggregator_layout
534 .as_ref()
535 .is_some_and(|layout| layout.aggregator == aggregator && layout.is_chainlink_ocr1())
536}
537
538#[derive(Clone, Default)]
540pub struct OracleStorageSync {
541 adapters: Vec<Arc<dyn OracleStorageAdapter>>,
542}
543
544impl std::fmt::Debug for OracleStorageSync {
545 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
546 f.debug_struct("OracleStorageSync")
547 .field("adapters", &self.adapters.len())
548 .finish()
549 }
550}
551
552impl OracleStorageSync {
553 pub fn new() -> Self {
555 Self::default()
556 }
557
558 pub fn chainlink_defaults() -> Self {
560 Self::new()
561 .with_adapter(Arc::new(ChainlinkOcr2StorageAdapter::new()))
562 .with_adapter(Arc::new(ChainlinkOcr1StorageAdapter::new()))
563 }
564
565 pub fn with_adapter(mut self, adapter: Arc<dyn OracleStorageAdapter>) -> Self {
567 self.adapters.push(adapter);
568 self
569 }
570
571 pub fn push_adapter(&mut self, adapter: Arc<dyn OracleStorageAdapter>) {
573 self.adapters.push(adapter);
574 }
575
576 pub fn is_empty(&self) -> bool {
578 self.adapters.is_empty()
579 }
580
581 pub fn warm_slots_for_registrations<'a>(
583 &self,
584 registrations: impl IntoIterator<Item = &'a FeedRegistration>,
585 ) -> Vec<(Address, U256)> {
586 let mut slots = Vec::new();
587 let mut seen = BTreeSet::new();
588 for registration in registrations {
589 for adapter in &self.adapters {
590 for slot in adapter.warm_slots_for_registration(registration) {
591 if seen.insert(slot) {
592 slots.push(slot);
593 }
594 }
595 }
596 }
597 slots
598 }
599
600 pub fn state_effect_for_answer(
602 &self,
603 registration: &FeedRegistration,
604 update: &OraclePriceUpdate,
605 state: &dyn StateView,
606 ) -> Result<OracleStorageEffect, OracleStorageError> {
607 for adapter in &self.adapters {
608 if let Some(updates) = adapter.state_updates_for_answer(registration, update, state)? {
609 return Ok(OracleStorageEffect::StateUpdates {
610 adapter: adapter.name(),
611 updates,
612 });
613 }
614 }
615
616 Ok(OracleStorageEffect::FallbackPurge)
617 }
618
619 pub fn state_effect_for_ocr2_transmission(
621 &self,
622 registration: &FeedRegistration,
623 update: &Ocr2TransmissionStorageUpdate,
624 state: &dyn StateView,
625 ) -> Result<OracleStorageEffect, OracleStorageError> {
626 for adapter in &self.adapters {
627 if let Some(updates) =
628 adapter.state_updates_for_ocr2_transmission(registration, update, state)?
629 {
630 return Ok(OracleStorageEffect::StateUpdates {
631 adapter: adapter.name(),
632 updates,
633 });
634 }
635 }
636
637 Ok(OracleStorageEffect::FallbackPurge)
638 }
639
640 pub fn state_effect_for_ocr1_transmission(
642 &self,
643 registration: &FeedRegistration,
644 update: &Ocr1TransmissionStorageUpdate,
645 state: &dyn StateView,
646 ) -> Result<OracleStorageEffect, OracleStorageError> {
647 for adapter in &self.adapters {
648 if let Some(updates) =
649 adapter.state_updates_for_ocr1_transmission(registration, update, state)?
650 {
651 return Ok(OracleStorageEffect::StateUpdates {
652 adapter: adapter.name(),
653 updates,
654 });
655 }
656 }
657
658 Ok(OracleStorageEffect::FallbackPurge)
659 }
660
661 pub fn prefers_ocr2_new_transmission(&self, registration: &FeedRegistration) -> bool {
663 self.adapters
664 .iter()
665 .any(|adapter| adapter.prefers_ocr2_new_transmission(registration))
666 }
667
668 pub fn prefers_ocr1_new_transmission(&self, registration: &FeedRegistration) -> bool {
670 self.adapters
671 .iter()
672 .any(|adapter| adapter.prefers_ocr1_new_transmission(registration))
673 }
674}
675
676fn mapping_slot(key: U256, base_slot: U256) -> U256 {
677 let mut preimage = [0_u8; 64];
678 preimage[..32].copy_from_slice(&key.to_be_bytes::<32>());
679 preimage[32..].copy_from_slice(&base_slot.to_be_bytes::<32>());
680 U256::from_be_slice(keccak256(preimage).as_slice())
681}
682
683fn encode_int192(value: I256) -> Result<U256, OracleStorageError> {
684 encode_int192_with_adapter(value, ChainlinkOcr2StorageAdapter::NAME, "OCR2")
685}
686
687fn encode_int192_with_adapter(
688 value: I256,
689 adapter: &'static str,
690 family: &'static str,
691) -> Result<U256, OracleStorageError> {
692 let raw = value.into_raw();
693 let encoded = raw & uint_mask(INT192_BITS);
694 let sign_bit_set = (encoded & (U256::from(1_u8) << (INT192_BITS - 1))) != U256::ZERO;
695 let high = raw >> INT192_BITS;
696 let expected_high = if sign_bit_set {
697 uint_mask(256 - INT192_BITS)
698 } else {
699 U256::ZERO
700 };
701
702 if high != expected_high {
703 return Err(OracleStorageError::adapter(
704 adapter,
705 format!("answer {value} does not fit {family} int192"),
706 ));
707 }
708
709 Ok(encoded)
710}
711
712fn uint_mask(bits: usize) -> U256 {
713 debug_assert!(bits <= 256);
714 match bits {
715 0 => U256::ZERO,
716 256 => U256::MAX,
717 bits => (U256::from(1_u8) << bits) - U256::from(1_u8),
718 }
719}
720
721fn uint_mask_u64(bits: usize) -> u64 {
722 debug_assert!(bits < 64);
723 (1_u64 << bits) - 1
724}