Skip to main content

evm_oracle_state/
reactive.rs

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/// Reactive-runtime bridge for Chainlink-compatible oracle feeds.
62#[derive(Clone, Debug)]
63pub struct OracleReactiveHandler {
64    registrations_by_dependency: BTreeMap<OracleDependencyKey, Vec<FeedRegistration>>,
65    storage_sync: OracleStorageSync,
66}
67
68impl OracleReactiveHandler {
69    /// Create a handler from current feed registrations.
70    pub fn new(registrations: Vec<FeedRegistration>) -> Self {
71        Self::with_storage_sync(registrations, OracleStorageSync::chainlink_defaults())
72    }
73
74    /// Create a handler from current feed registrations and storage adapters.
75    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(&registration);
95                let wants_ocr1 =
96                    !wants_answer && storage_sync.prefers_ocr1_new_transmission(&registration);
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                        &registration,
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                        &registration,
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                        &registration,
127                    );
128                }
129            }
130        }
131        Self {
132            registrations_by_dependency,
133            storage_sync,
134        }
135    }
136
137    /// Stable handler id.
138    pub fn id(&self) -> HandlerId {
139        HandlerId::new(HANDLER_ID)
140    }
141
142    /// Interests for current aggregators.
143    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    /// Test helper that avoids requiring callers to construct a cache view.
151    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
567/// Removed (reorged-out) logs describe values the chain rolled back. Never
568/// re-apply them as exact storage writes — purge the affected feed proxies and
569/// the emitting aggregator so raw reads refetch canonical state. The paired
570/// hooks still fire with [`OracleValueStatus::RequiresRepair`] so trackers
571/// mark the affected snapshots for authoritative reconciliation.
572fn 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, &registration.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, &registration.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, &registration.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}