1use std::{
2 borrow::Cow,
3 collections::{BTreeMap, BTreeSet},
4 sync::Arc,
5};
6
7use alloy_network::Ethereum;
8use alloy_primitives::{Address, B256, I256, U256};
9use alloy_rpc_types_eth::Filter;
10use evm_fork_cache::{
11 StateView,
12 reactive::{
13 ChainStatus, HandlerError, HandlerId, HandlerOutcome, HookSignal, InvalidationReason,
14 InvalidationRequest, LogInterest, ReactiveContext, ReactiveEffect, ReactiveHandler,
15 ReactiveInput, ReactiveInterest, ReportTag, RouteKeySpec, StateEffectQuality,
16 },
17 state_update::PurgeScope,
18};
19
20use crate::{
21 ANSWER_UPDATED_TOPIC, AnswerUpdated, FeedRegistration, OCR1_NEW_TRANSMISSION_TOPIC,
22 OCR2_NEW_TRANSMISSION_TOPIC, ORACLE_LEGACY_ANSWER_UPDATED_KIND, ORACLE_SIGNAL_NAMESPACE,
23 Ocr1NewTransmission, Ocr1TransmissionStorageUpdate, Ocr2NewTransmission,
24 Ocr2TransmissionStorageUpdate, OraclePriceUpdate, OracleSignalKind, OracleStorageEffect,
25 OracleStorageError, OracleStorageSync, OracleUpdate, OracleValueStatus, RoundData,
26 decode_answer_updated, decode_ocr1_new_transmission, decode_ocr2_new_transmission,
27 state::classify_round,
28};
29
30const HANDLER_ID: &str = "evm-oracle-state.chainlink";
31
32#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
33struct OracleDependencyKey {
34 aggregator: Address,
35 event: OracleDependencyEvent,
36}
37
38impl OracleDependencyKey {
39 fn new(aggregator: Address, event: OracleDependencyEvent) -> Self {
40 Self { aggregator, event }
41 }
42}
43
44#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
45enum OracleDependencyEvent {
46 AnswerUpdated,
47 Ocr1NewTransmission,
48 Ocr2NewTransmission,
49}
50
51impl OracleDependencyEvent {
52 fn topic(self) -> B256 {
53 match self {
54 Self::AnswerUpdated => ANSWER_UPDATED_TOPIC,
55 Self::Ocr1NewTransmission => OCR1_NEW_TRANSMISSION_TOPIC,
56 Self::Ocr2NewTransmission => OCR2_NEW_TRANSMISSION_TOPIC,
57 }
58 }
59}
60
61#[derive(Clone, Debug)]
63pub struct OracleReactiveHandler {
64 registrations_by_dependency: BTreeMap<OracleDependencyKey, Vec<FeedRegistration>>,
65 storage_sync: OracleStorageSync,
66}
67
68impl OracleReactiveHandler {
69 pub fn new(registrations: Vec<FeedRegistration>) -> Self {
71 Self::with_storage_sync(registrations, OracleStorageSync::chainlink_defaults())
72 }
73
74 pub fn with_storage_sync(
76 registrations: Vec<FeedRegistration>,
77 storage_sync: OracleStorageSync,
78 ) -> Self {
79 let mut registrations_by_dependency: BTreeMap<OracleDependencyKey, Vec<FeedRegistration>> =
80 BTreeMap::new();
81 for registration in registrations {
82 if !registration.source.uses_builtin_chainlink_handler() {
83 continue;
84 }
85 let mut seen_dependencies = BTreeSet::new();
86 for aggregator in registration
87 .source
88 .event_aggregators(registration.current_aggregator)
89 {
90 let wants_answer = registration
91 .source
92 .wants_answer_updated_from(registration.current_aggregator, aggregator);
93 let wants_ocr2 =
94 !wants_answer && storage_sync.prefers_ocr2_new_transmission(®istration);
95 let wants_ocr1 =
96 !wants_answer && storage_sync.prefers_ocr1_new_transmission(®istration);
97 let wants_answer = wants_answer || (!wants_ocr2 && !wants_ocr1);
98
99 if wants_answer {
100 insert_dependency_registration(
101 &mut registrations_by_dependency,
102 &mut seen_dependencies,
103 OracleDependencyKey::new(aggregator, OracleDependencyEvent::AnswerUpdated),
104 ®istration,
105 );
106 }
107 if wants_ocr1 {
108 insert_dependency_registration(
109 &mut registrations_by_dependency,
110 &mut seen_dependencies,
111 OracleDependencyKey::new(
112 aggregator,
113 OracleDependencyEvent::Ocr1NewTransmission,
114 ),
115 ®istration,
116 );
117 }
118 if wants_ocr2 {
119 insert_dependency_registration(
120 &mut registrations_by_dependency,
121 &mut seen_dependencies,
122 OracleDependencyKey::new(
123 aggregator,
124 OracleDependencyEvent::Ocr2NewTransmission,
125 ),
126 ®istration,
127 );
128 }
129 }
130 }
131 Self {
132 registrations_by_dependency,
133 storage_sync,
134 }
135 }
136
137 pub fn id(&self) -> HandlerId {
139 HandlerId::new(HANDLER_ID)
140 }
141
142 pub fn interests(&self) -> Vec<ReactiveInterest<Ethereum>> {
144 self.registrations_by_dependency
145 .keys()
146 .map(|key| log_interest(key.aggregator, key.event.topic()))
147 .collect()
148 }
149
150 pub fn handle_for_test(
152 &self,
153 ctx: &ReactiveContext,
154 input: &ReactiveInput<Ethereum>,
155 ) -> Result<HandlerOutcome, HandlerError> {
156 self.handle(ctx, input, &EmptyStateView)
157 }
158
159 fn handle_answer_updated(
160 &self,
161 ctx: &ReactiveContext,
162 log: &alloy_rpc_types_eth::Log,
163 state: &dyn StateView,
164 ) -> Result<HandlerOutcome, HandlerError> {
165 let aggregator = log.address();
166 let key = OracleDependencyKey::new(aggregator, OracleDependencyEvent::AnswerUpdated);
167 let Some(registrations) = self.registrations_by_dependency.get(&key) else {
168 return Ok(HandlerOutcome::empty(StateEffectQuality::NoStateEffect));
169 };
170
171 let event = decode_answer_updated(log)
172 .map_err(|error| HandlerError::new(format!("decode AnswerUpdated failed: {error}")))?;
173 let block_number = event
174 .block_number
175 .or_else(|| ctx.block.as_ref().map(|block| block.number));
176 let block_hash = log
177 .block_hash
178 .or_else(|| ctx.block.as_ref().map(|block| block.hash));
179 let log_index = event.log_index.or(ctx.log_index);
180
181 let mut effects = Vec::new();
182 let mut tags = Vec::new();
183 let mut needs_aggregator_purge = false;
184
185 let value_status = event_value_status(log.removed);
186 let applied_direct_effect = if log.removed {
187 apply_removed_log_fallback(&mut effects, &mut needs_aggregator_purge, registrations);
188 false
189 } else {
190 apply_answer_updated_dependency_storage_effect(
191 &self.storage_sync,
192 &mut effects,
193 &mut needs_aggregator_purge,
194 registrations,
195 &event,
196 block_number,
197 block_hash,
198 log_index,
199 value_status,
200 state,
201 )
202 .map_err(|error| HandlerError::new(format!("{error}")))?
203 };
204
205 for registration in registrations {
206 let normalized_answer = registration
207 .source
208 .normalize_answer_from_event(Some(aggregator), event.current);
209 let update = OracleUpdate {
210 id: registration.id.clone(),
211 proxy: registration.proxy,
212 aggregator,
213 round: RoundData {
214 round_id: event.round_id,
215 answer: normalized_answer,
216 started_at: event.updated_at,
217 updated_at: event.updated_at,
218 answered_in_round: event.round_id,
219 },
220 block_number,
221 log_index,
222 value_status,
223 };
224 let price_update = answer_updated_price_update(
225 registration,
226 aggregator,
227 &event,
228 normalized_answer,
229 block_number,
230 block_hash,
231 log_index,
232 value_status,
233 );
234 let hook_tags = hook_tags(registration, aggregator);
235 tags.extend(hook_tags.clone());
236 if !log.removed {
237 effects.push(ReactiveEffect::Hook(HookSignal {
238 namespace: Cow::Borrowed(ORACLE_SIGNAL_NAMESPACE),
239 kind: Cow::Borrowed(ORACLE_LEGACY_ANSWER_UPDATED_KIND),
240 labels: hook_tags.clone(),
241 payload: Some(Arc::new(update)),
242 }));
243 }
244 effects.push(ReactiveEffect::Hook(HookSignal {
245 namespace: Cow::Borrowed(ORACLE_SIGNAL_NAMESPACE),
246 kind: Cow::Borrowed(OracleSignalKind::PriceUpdate.as_str()),
247 labels: hook_tags,
248 payload: Some(Arc::new(price_update)),
249 }));
250 }
251
252 prepend_aggregator_purge(&mut effects, needs_aggregator_purge, aggregator);
253
254 Ok(HandlerOutcome {
255 effects,
256 quality: state_effect_quality(applied_direct_effect),
257 tags,
258 })
259 }
260
261 fn handle_ocr2_new_transmission(
262 &self,
263 ctx: &ReactiveContext,
264 log: &alloy_rpc_types_eth::Log,
265 state: &dyn StateView,
266 ) -> Result<HandlerOutcome, HandlerError> {
267 let aggregator = log.address();
268 let key = OracleDependencyKey::new(aggregator, OracleDependencyEvent::Ocr2NewTransmission);
269 let Some(registrations) = self.registrations_by_dependency.get(&key) else {
270 return Ok(HandlerOutcome::empty(StateEffectQuality::NoStateEffect));
271 };
272
273 let event = decode_ocr2_new_transmission(log).map_err(|error| {
274 HandlerError::new(format!("decode OCR2 NewTransmission failed: {error}"))
275 })?;
276 let transmission_timestamp = ctx
277 .block
278 .as_ref()
279 .and_then(|block| block.timestamp)
280 .or(log.block_timestamp)
281 .ok_or_else(|| {
282 HandlerError::new("OCR2 NewTransmission handling requires block timestamp")
283 })?;
284 let block_number = event
285 .block_number
286 .or_else(|| ctx.block.as_ref().map(|block| block.number));
287 let block_hash = log
288 .block_hash
289 .or_else(|| ctx.block.as_ref().map(|block| block.hash));
290 let log_index = event.log_index.or(ctx.log_index);
291
292 let mut effects = Vec::new();
293 let mut tags = Vec::new();
294 let mut needs_aggregator_purge = false;
295
296 let transmission_update = Ocr2TransmissionStorageUpdate {
297 aggregator,
298 aggregator_round_id: event.aggregator_round_id,
299 answer: event.answer,
300 observations_timestamp: event.observations_timestamp,
301 transmission_timestamp,
302 epoch_and_round: event.epoch_and_round,
303 };
304 let value_status = event_value_status(log.removed);
305 let applied_direct_effect = if log.removed {
306 apply_removed_log_fallback(&mut effects, &mut needs_aggregator_purge, registrations);
307 false
308 } else {
309 apply_ocr2_dependency_storage_effect(
310 &self.storage_sync,
311 &mut effects,
312 &mut needs_aggregator_purge,
313 registrations,
314 &transmission_update,
315 state,
316 )
317 .map_err(|error| HandlerError::new(format!("{error}")))?
318 };
319 for registration in registrations {
320 let price_update = ocr2_price_update(
321 registration,
322 aggregator,
323 &event,
324 transmission_timestamp,
325 block_number,
326 block_hash,
327 log_index,
328 value_status,
329 );
330 let hook_tags = hook_tags(registration, aggregator);
331 tags.extend(hook_tags.clone());
332 effects.push(ReactiveEffect::Hook(HookSignal {
333 namespace: Cow::Borrowed(ORACLE_SIGNAL_NAMESPACE),
334 kind: Cow::Borrowed(OracleSignalKind::PriceUpdate.as_str()),
335 labels: hook_tags,
336 payload: Some(Arc::new(price_update)),
337 }));
338 }
339
340 prepend_aggregator_purge(&mut effects, needs_aggregator_purge, aggregator);
341
342 Ok(HandlerOutcome {
343 effects,
344 quality: state_effect_quality(applied_direct_effect),
345 tags,
346 })
347 }
348
349 fn handle_ocr1_new_transmission(
350 &self,
351 ctx: &ReactiveContext,
352 log: &alloy_rpc_types_eth::Log,
353 state: &dyn StateView,
354 ) -> Result<HandlerOutcome, HandlerError> {
355 let aggregator = log.address();
356 let key = OracleDependencyKey::new(aggregator, OracleDependencyEvent::Ocr1NewTransmission);
357 let Some(registrations) = self.registrations_by_dependency.get(&key) else {
358 return Ok(HandlerOutcome::empty(StateEffectQuality::NoStateEffect));
359 };
360
361 let event = decode_ocr1_new_transmission(log).map_err(|error| {
362 HandlerError::new(format!("decode OCR1 NewTransmission failed: {error}"))
363 })?;
364 let transmission_timestamp = ctx
365 .block
366 .as_ref()
367 .and_then(|block| block.timestamp)
368 .or(log.block_timestamp)
369 .ok_or_else(|| {
370 HandlerError::new("OCR1 NewTransmission handling requires block timestamp")
371 })?;
372 let block_number = event
373 .block_number
374 .or_else(|| ctx.block.as_ref().map(|block| block.number));
375 let block_hash = log
376 .block_hash
377 .or_else(|| ctx.block.as_ref().map(|block| block.hash));
378 let log_index = event.log_index.or(ctx.log_index);
379
380 let mut effects = Vec::new();
381 let mut tags = Vec::new();
382 let mut needs_aggregator_purge = false;
383
384 let transmission_update = Ocr1TransmissionStorageUpdate {
385 aggregator,
386 aggregator_round_id: event.aggregator_round_id,
387 answer: event.answer,
388 transmission_timestamp,
389 epoch_and_round: event.epoch_and_round,
390 };
391 let value_status = event_value_status(log.removed);
392 let applied_direct_effect = if log.removed {
393 apply_removed_log_fallback(&mut effects, &mut needs_aggregator_purge, registrations);
394 false
395 } else {
396 apply_ocr1_dependency_storage_effect(
397 &self.storage_sync,
398 &mut effects,
399 &mut needs_aggregator_purge,
400 registrations,
401 &transmission_update,
402 state,
403 )
404 .map_err(|error| HandlerError::new(format!("{error}")))?
405 };
406 for registration in registrations {
407 let price_update = ocr1_price_update(
408 registration,
409 aggregator,
410 &event,
411 transmission_timestamp,
412 block_number,
413 block_hash,
414 log_index,
415 value_status,
416 );
417 let hook_tags = hook_tags(registration, aggregator);
418 tags.extend(hook_tags.clone());
419 effects.push(ReactiveEffect::Hook(HookSignal {
420 namespace: Cow::Borrowed(ORACLE_SIGNAL_NAMESPACE),
421 kind: Cow::Borrowed(OracleSignalKind::PriceUpdate.as_str()),
422 labels: hook_tags,
423 payload: Some(Arc::new(price_update)),
424 }));
425 }
426
427 prepend_aggregator_purge(&mut effects, needs_aggregator_purge, aggregator);
428
429 Ok(HandlerOutcome {
430 effects,
431 quality: state_effect_quality(applied_direct_effect),
432 tags,
433 })
434 }
435}
436
437fn insert_dependency_registration(
438 registrations_by_dependency: &mut BTreeMap<OracleDependencyKey, Vec<FeedRegistration>>,
439 seen_dependencies: &mut BTreeSet<OracleDependencyKey>,
440 key: OracleDependencyKey,
441 registration: &FeedRegistration,
442) {
443 if seen_dependencies.insert(key) {
444 registrations_by_dependency
445 .entry(key)
446 .or_default()
447 .push(registration.clone());
448 }
449}
450
451fn log_interest(aggregator: Address, topic: B256) -> ReactiveInterest<Ethereum> {
452 ReactiveInterest::Logs(LogInterest {
453 provider_filter: Filter::new().address(aggregator).event_signature(topic),
454 local_matcher: None,
455 route_key: Some(RouteKeySpec::EmitterAddress),
456 })
457}
458
459#[allow(clippy::too_many_arguments)]
460fn apply_answer_updated_dependency_storage_effect(
461 storage_sync: &OracleStorageSync,
462 effects: &mut Vec<ReactiveEffect>,
463 needs_aggregator_purge: &mut bool,
464 registrations: &[FeedRegistration],
465 event: &AnswerUpdated,
466 block_number: Option<u64>,
467 block_hash: Option<B256>,
468 log_index: Option<u64>,
469 value_status: OracleValueStatus,
470 state: &dyn StateView,
471) -> Result<bool, OracleStorageError> {
472 let mut direct_candidates = Vec::new();
473 let mut fallback_proxies = BTreeSet::new();
474
475 for registration in registrations {
476 if registration
477 .source
478 .supports_direct_chainlink_storage_effects()
479 {
480 direct_candidates.push(registration);
481 } else {
482 fallback_proxies.insert(registration.proxy);
483 }
484 }
485
486 let mut applied_direct_effect = false;
487 for registration in direct_candidates.iter().copied() {
488 let raw_update = answer_updated_price_update(
489 registration,
490 event.aggregator,
491 event,
492 event.current,
493 block_number,
494 block_hash,
495 log_index,
496 value_status,
497 );
498 if apply_state_updates(
499 effects,
500 storage_sync.state_effect_for_answer(registration, &raw_update, state)?,
501 ) {
502 applied_direct_effect = true;
503 break;
504 }
505 }
506
507 if !applied_direct_effect {
508 fallback_proxies.extend(
509 direct_candidates
510 .into_iter()
511 .map(|registration| registration.proxy),
512 );
513 }
514
515 append_fallback_invalidations(effects, needs_aggregator_purge, fallback_proxies);
516 Ok(applied_direct_effect)
517}
518
519fn apply_ocr2_dependency_storage_effect(
520 storage_sync: &OracleStorageSync,
521 effects: &mut Vec<ReactiveEffect>,
522 needs_aggregator_purge: &mut bool,
523 registrations: &[FeedRegistration],
524 update: &Ocr2TransmissionStorageUpdate,
525 state: &dyn StateView,
526) -> Result<bool, OracleStorageError> {
527 let mut fallback_proxies = BTreeSet::new();
528
529 for registration in registrations {
530 if apply_state_updates(
531 effects,
532 storage_sync.state_effect_for_ocr2_transmission(registration, update, state)?,
533 ) {
534 return Ok(true);
535 }
536 fallback_proxies.insert(registration.proxy);
537 }
538
539 append_fallback_invalidations(effects, needs_aggregator_purge, fallback_proxies);
540 Ok(false)
541}
542
543fn apply_ocr1_dependency_storage_effect(
544 storage_sync: &OracleStorageSync,
545 effects: &mut Vec<ReactiveEffect>,
546 needs_aggregator_purge: &mut bool,
547 registrations: &[FeedRegistration],
548 update: &Ocr1TransmissionStorageUpdate,
549 state: &dyn StateView,
550) -> Result<bool, OracleStorageError> {
551 let mut fallback_proxies = BTreeSet::new();
552
553 for registration in registrations {
554 if apply_state_updates(
555 effects,
556 storage_sync.state_effect_for_ocr1_transmission(registration, update, state)?,
557 ) {
558 return Ok(true);
559 }
560 fallback_proxies.insert(registration.proxy);
561 }
562
563 append_fallback_invalidations(effects, needs_aggregator_purge, fallback_proxies);
564 Ok(false)
565}
566
567fn apply_removed_log_fallback(
573 effects: &mut Vec<ReactiveEffect>,
574 needs_aggregator_purge: &mut bool,
575 registrations: &[FeedRegistration],
576) {
577 let proxies: BTreeSet<Address> = registrations
578 .iter()
579 .map(|registration| registration.proxy)
580 .collect();
581 append_fallback_invalidations(effects, needs_aggregator_purge, proxies);
582}
583
584fn apply_state_updates(effects: &mut Vec<ReactiveEffect>, effect: OracleStorageEffect) -> bool {
585 match effect {
586 OracleStorageEffect::StateUpdates { updates, .. } => {
587 effects.extend(updates.into_iter().map(ReactiveEffect::StateUpdate));
588 true
589 }
590 OracleStorageEffect::FallbackPurge => false,
591 }
592}
593
594fn append_fallback_invalidations(
595 effects: &mut Vec<ReactiveEffect>,
596 needs_aggregator_purge: &mut bool,
597 proxies: BTreeSet<Address>,
598) {
599 if proxies.is_empty() {
600 return;
601 }
602
603 *needs_aggregator_purge = true;
604 effects.extend(proxies.into_iter().map(|proxy| {
605 ReactiveEffect::Invalidate(InvalidationRequest {
606 scope: PurgeScope::AllStorage,
607 address: proxy,
608 reason: InvalidationReason::HandlerRequested,
609 })
610 }));
611}
612
613fn prepend_aggregator_purge(
614 effects: &mut Vec<ReactiveEffect>,
615 needs_aggregator_purge: bool,
616 aggregator: Address,
617) {
618 if needs_aggregator_purge {
619 effects.insert(
620 0,
621 ReactiveEffect::Invalidate(InvalidationRequest {
622 scope: PurgeScope::AllStorage,
623 address: aggregator,
624 reason: InvalidationReason::HandlerRequested,
625 }),
626 );
627 }
628}
629
630#[allow(clippy::too_many_arguments)]
631fn answer_updated_price_update(
632 registration: &FeedRegistration,
633 aggregator: Address,
634 event: &AnswerUpdated,
635 raw_answer: I256,
636 block_number: Option<u64>,
637 block_hash: Option<B256>,
638 log_index: Option<u64>,
639 value_status: OracleValueStatus,
640) -> OraclePriceUpdate {
641 let round = RoundData {
642 round_id: event.round_id,
643 answer: raw_answer,
644 started_at: event.updated_at,
645 updated_at: event.updated_at,
646 answered_in_round: event.round_id,
647 };
648 let round_status = classify_round(&round, event.updated_at, ®istration.staleness);
649 OraclePriceUpdate {
650 id: registration.id.clone(),
651 proxy: registration.proxy,
652 aggregator,
653 label: registration.label.clone(),
654 base: registration.base.clone(),
655 quote: registration.quote.clone(),
656 raw_answer,
657 decimals: registration.metadata.decimals,
658 event_round_id: event.round_id,
659 started_at: event.updated_at,
660 updated_at: event.updated_at,
661 block_number,
662 block_hash,
663 log_index,
664 round_status,
665 value_status,
666 source: registration.source.event_value_source(),
667 }
668}
669
670#[allow(clippy::too_many_arguments)]
671fn ocr2_price_update(
672 registration: &FeedRegistration,
673 aggregator: Address,
674 event: &Ocr2NewTransmission,
675 transmission_timestamp: u64,
676 block_number: Option<u64>,
677 block_hash: Option<B256>,
678 log_index: Option<u64>,
679 value_status: OracleValueStatus,
680) -> OraclePriceUpdate {
681 let raw_answer = registration
682 .source
683 .normalize_answer_from_event(Some(aggregator), event.answer);
684 let round = RoundData {
685 round_id: event.aggregator_round_id,
686 answer: raw_answer,
687 started_at: event.observations_timestamp,
688 updated_at: transmission_timestamp,
689 answered_in_round: event.aggregator_round_id,
690 };
691 let round_status = classify_round(&round, transmission_timestamp, ®istration.staleness);
692 OraclePriceUpdate {
693 id: registration.id.clone(),
694 proxy: registration.proxy,
695 aggregator,
696 label: registration.label.clone(),
697 base: registration.base.clone(),
698 quote: registration.quote.clone(),
699 raw_answer,
700 decimals: registration.metadata.decimals,
701 event_round_id: event.aggregator_round_id,
702 started_at: event.observations_timestamp,
703 updated_at: transmission_timestamp,
704 block_number,
705 block_hash,
706 log_index,
707 round_status,
708 value_status,
709 source: registration.source.event_value_source(),
710 }
711}
712
713#[allow(clippy::too_many_arguments)]
714fn ocr1_price_update(
715 registration: &FeedRegistration,
716 aggregator: Address,
717 event: &Ocr1NewTransmission,
718 transmission_timestamp: u64,
719 block_number: Option<u64>,
720 block_hash: Option<B256>,
721 log_index: Option<u64>,
722 value_status: OracleValueStatus,
723) -> OraclePriceUpdate {
724 let raw_answer = registration
725 .source
726 .normalize_answer_from_event(Some(aggregator), event.answer);
727 let round = RoundData {
728 round_id: event.aggregator_round_id,
729 answer: raw_answer,
730 started_at: transmission_timestamp,
731 updated_at: transmission_timestamp,
732 answered_in_round: event.aggregator_round_id,
733 };
734 let round_status = classify_round(&round, transmission_timestamp, ®istration.staleness);
735 OraclePriceUpdate {
736 id: registration.id.clone(),
737 proxy: registration.proxy,
738 aggregator,
739 label: registration.label.clone(),
740 base: registration.base.clone(),
741 quote: registration.quote.clone(),
742 raw_answer,
743 decimals: registration.metadata.decimals,
744 event_round_id: event.aggregator_round_id,
745 started_at: transmission_timestamp,
746 updated_at: transmission_timestamp,
747 block_number,
748 block_hash,
749 log_index,
750 round_status,
751 value_status,
752 source: registration.source.event_value_source(),
753 }
754}
755
756fn event_value_status(removed: bool) -> OracleValueStatus {
757 if removed {
758 OracleValueStatus::RequiresRepair
759 } else {
760 OracleValueStatus::EventPending
761 }
762}
763
764fn state_effect_quality(applied_direct_effect: bool) -> StateEffectQuality {
765 if applied_direct_effect {
766 StateEffectQuality::ExactFromInput
767 } else {
768 StateEffectQuality::RequiresRepair
769 }
770}
771
772fn hook_tags(registration: &FeedRegistration, aggregator: Address) -> Vec<ReportTag> {
773 vec![
774 ReportTag::new("feed_id", registration.id.to_string()),
775 ReportTag::new("proxy", format!("{:?}", registration.proxy)),
776 ReportTag::new("aggregator", format!("{aggregator:?}")),
777 ]
778}
779
780impl ReactiveHandler<Ethereum> for OracleReactiveHandler {
781 fn id(&self) -> HandlerId {
782 self.id()
783 }
784
785 fn interests(&self) -> Vec<ReactiveInterest<Ethereum>> {
786 self.interests()
787 }
788
789 fn handle(
790 &self,
791 ctx: &ReactiveContext,
792 input: &ReactiveInput<Ethereum>,
793 state: &dyn StateView,
794 ) -> Result<HandlerOutcome, HandlerError> {
795 let ReactiveInput::Log(log) = input else {
796 return Ok(HandlerOutcome::empty(StateEffectQuality::NoStateEffect));
797 };
798 if log.removed && !matches!(ctx.chain_status, ChainStatus::Reorged { .. }) {
799 return Ok(HandlerOutcome::empty(StateEffectQuality::NoStateEffect));
800 }
801
802 match log.topics().first().copied() {
803 Some(ANSWER_UPDATED_TOPIC) => self.handle_answer_updated(ctx, log, state),
804 Some(OCR1_NEW_TRANSMISSION_TOPIC) => self.handle_ocr1_new_transmission(ctx, log, state),
805 Some(OCR2_NEW_TRANSMISSION_TOPIC) => self.handle_ocr2_new_transmission(ctx, log, state),
806 _ => Ok(HandlerOutcome::empty(StateEffectQuality::NoStateEffect)),
807 }
808 }
809}
810
811struct EmptyStateView;
812
813impl StateView for EmptyStateView {
814 fn storage(&self, _address: Address, _slot: U256) -> Option<U256> {
815 None
816 }
817}