Skip to main content

tycho_simulation/evm/
decoder.rs

1use std::{
2    collections::{hash_map::Entry, HashMap, HashSet},
3    future::Future,
4    pin::Pin,
5    sync::Arc,
6};
7
8use alloy::primitives::{Address, U256};
9use thiserror::Error;
10use tokio::sync::{watch, RwLock, RwLockReadGuard};
11use tracing::{debug, error, info, warn};
12use tycho_client::feed::{synchronizer::ComponentWithState, BlockHeader, FeedMessage, HeaderLike};
13use tycho_common::{
14    dto::{ChangeType, ProtocolStateDelta},
15    models::{blockchain::BlockAggregatedChanges, token::Token, Chain},
16    simulation::protocol_sim::{Balances, BlockContext, ProtocolSim},
17    Bytes,
18};
19#[cfg(test)]
20use {
21    mockall::mock,
22    num_bigint::BigUint,
23    std::any::Any,
24    tycho_common::simulation::{
25        errors::{SimulationError, TransitionError},
26        protocol_sim::GetAmountOutResult,
27    },
28};
29
30use crate::{
31    evm::{
32        engine_db::{update_engine, SHARED_TYCHO_DB},
33        override_stream::{OverrideSnapshot, StateOverrideProvider},
34        protocol::{
35            utils::bytes_to_address,
36            vm::{constants::ERC20_PROXY_BYTECODE, erc20_token::IMPLEMENTATION_SLOT},
37        },
38        tycho_models::{AccountUpdate, ResponseAccount},
39    },
40    protocol::{
41        errors::InvalidSnapshotError,
42        models::{DecoderContext, ProtocolComponent, TryFromWithBlock, Update},
43    },
44};
45
46#[derive(Error, Debug)]
47pub enum StreamDecodeError {
48    #[error("{0}")]
49    Fatal(String),
50}
51
52#[derive(Default)]
53struct DecoderState {
54    tokens: HashMap<Bytes, Token>,
55    states: HashMap<String, Box<dyn ProtocolSim>>,
56    components: HashMap<String, ProtocolComponent>,
57    // maps contract address to the pools they affect
58    contracts_map: HashMap<Bytes, HashSet<String>>,
59    // Maps original token address to their new proxy token address
60    proxy_token_addresses: HashMap<Address, Address>,
61    // Set of failed components, these are components that failed to decode and will not be emitted
62    // again TODO: handle more gracefully inside tycho-client. We could fetch the snapshot and
63    // try to decode it again.
64    failed_components: HashSet<String>,
65    // The block number of the last confirmed block decoded via `decode()`.
66    current_block_number: u64,
67}
68
69type DecodeFut =
70    Pin<Box<dyn Future<Output = Result<Box<dyn ProtocolSim>, InvalidSnapshotError>> + Send + Sync>>;
71type AccountBalances = HashMap<Bytes, HashMap<Bytes, Bytes>>;
72type RegistryFn<H> = dyn Fn(
73        ComponentWithState,
74        H,
75        AccountBalances,
76        Arc<RwLock<DecoderState>>,
77        Option<watch::Receiver<OverrideSnapshot>>,
78    ) -> DecodeFut
79    + Send
80    + Sync;
81type FilterFn = fn(&ComponentWithState) -> bool;
82
83/// A decoder to process raw messages.
84///
85/// This struct decodes incoming messages of type `FeedMessage` and converts it into the
86/// `BlockUpdate` struct.
87///
88/// # Important:
89/// - Supports registering exchanges and their associated filters for specific protocol components.
90/// - Allows the addition of client-side filters for custom conditions.
91///
92/// **Note:** Tokens provided via [`set_tokens`](Self::set_tokens) are used to decode startup
93/// snapshots and initialize protocol states. This is not an ongoing filter — components arriving
94/// after startup include their own token metadata.
95pub struct TychoStreamDecoder<H>
96where
97    H: HeaderLike,
98{
99    state: Arc<RwLock<DecoderState>>,
100    skip_state_decode_failures: bool,
101    min_token_quality: u32,
102    registry: HashMap<String, Box<RegistryFn<H>>>,
103    inclusion_filters: HashMap<String, Vec<FilterFn>>,
104    /// Live override providers keyed by `protocol_system`. A pool of that protocol subscribes to
105    /// its provider at creation time and reads fresh overrides on every simulation.
106    override_providers: HashMap<String, Arc<dyn StateOverrideProvider>>,
107    /// Seconds between blocks, used to project a confirmed block header onto the next block when
108    /// deriving the execution block for block-sensitive states.
109    block_time_secs: u64,
110    /// Chain every component this decoder handles lives on. Stamped into each registered
111    /// [`DecoderContext`] so chain-specific decoders cannot be pointed at the wrong chain.
112    chain: Chain,
113}
114
115/// Curve migrated from the generic VM adapter (`EVMPoolState`) to the native [`CurveState`]
116/// decoder. Returns true when `vm:curve` is registered with any other type — i.e. the deprecated
117/// VM-adapter path, still supported for a few releases before removal.
118fn is_deprecated_curve_registration<T: 'static>(exchange: &str) -> bool {
119    exchange == "vm:curve" &&
120        std::any::type_name::<T>() !=
121            std::any::type_name::<crate::evm::protocol::curve::CurveState>()
122}
123
124impl<H> TychoStreamDecoder<H>
125where
126    H: HeaderLike + Clone + Sync + Send + 'static + std::fmt::Debug,
127{
128    /// Creates a decoder for `chain`.
129    ///
130    /// # Panics
131    ///
132    /// Panics if `chain` is a custom chain with no registered config.
133    pub fn new(chain: Chain) -> Self {
134        Self {
135            state: Arc::new(RwLock::new(DecoderState::default())),
136            skip_state_decode_failures: false,
137            min_token_quality: 100,
138            registry: HashMap::new(),
139            inclusion_filters: HashMap::new(),
140            override_providers: HashMap::new(),
141            block_time_secs: chain.block_time_secs(),
142            chain,
143        }
144    }
145
146    /// The block a quote produced from `header` is expected to execute in.
147    ///
148    /// A partial (flashblock) header describes a block that is still open, so a quote can still
149    /// land in it. A confirmed header describes a closed block, so the quote targets the next one.
150    fn execution_block(&self, header: &BlockHeader) -> BlockContext {
151        if header.partial_block_index.is_some() {
152            BlockContext::new(header.number, header.timestamp)
153        } else {
154            BlockContext::new(header.number + 1, header.timestamp + self.block_time_secs)
155        }
156    }
157
158    /// Advances every state to `execution_block` via [`ProtocolSim::apply_block`].
159    ///
160    /// States already being emitted are advanced in place. Stored states absent from this
161    /// message are advanced in place under the write guard and cloned into `updated_states`
162    /// only when their quoting behavior changed — so an idle block-sensitive pool costs one
163    /// virtual call per message and zero clones, and consumers are only told about pools whose
164    /// quotes actually moved. Stored states of failed or removed components are left untouched,
165    /// so a component the consumer was told is gone is never re-emitted.
166    fn refresh_execution_block<C>(
167        updated_states: &mut HashMap<String, Box<dyn ProtocolSim>>,
168        stored_states: &mut HashMap<String, Box<dyn ProtocolSim>>,
169        failed_components: &HashSet<String>,
170        removed_components: &HashMap<String, C>,
171        execution_block: &BlockContext,
172    ) {
173        for state in updated_states.values_mut() {
174            state.apply_block(execution_block);
175        }
176        for (id, state) in stored_states.iter_mut() {
177            if failed_components.contains(id) ||
178                removed_components.contains_key(id) ||
179                updated_states.contains_key(id)
180            {
181                continue;
182            }
183            if state.apply_block(execution_block) {
184                updated_states.insert(id.clone(), state.clone_box());
185            }
186        }
187    }
188
189    /// Registers `provider` as the live override source for `protocol_system`.
190    ///
191    /// Pools of that protocol subscribe to it at creation time, so overrides apply from the first
192    /// simulation onward. A later call for the same `protocol_system` replaces the previous
193    /// provider.
194    pub fn set_override_provider(
195        &mut self,
196        protocol_system: String,
197        provider: Arc<dyn StateOverrideProvider>,
198    ) {
199        self.override_providers
200            .insert(protocol_system, provider);
201    }
202
203    /// Provides token metadata used to decode startup snapshots and initialize protocol states.
204    ///
205    /// This is not an ongoing stream filter. Components arriving after startup include their
206    /// own token metadata for decoding.
207    pub async fn set_tokens(&self, tokens: HashMap<Bytes, Token>) {
208        let mut guard = self.state.write().await;
209        guard.tokens = tokens;
210    }
211
212    pub fn skip_state_decode_failures(&mut self, skip: bool) {
213        self.skip_state_decode_failures = skip;
214    }
215
216    /// Sets the minimum token quality for decoding.
217    ///
218    /// Tokens arriving in stream deltas below this threshold are ignored. Defaults to 100.
219    /// Set this to the same value used in [`load_all_tokens()`](crate::utils::load_all_tokens) to
220    /// apply consistent filtering.
221    pub fn min_token_quality(&mut self, quality: u32) {
222        self.min_token_quality = quality;
223    }
224
225    /// Registers a decoder for a given exchange with a decoder context.
226    ///
227    /// This method maps an exchange identifier to a specific protocol simulation type.
228    /// The associated type must implement the `TryFromWithBlock` trait to enable decoding
229    /// of state updates from `ComponentWithState` objects. This allows the decoder to transform
230    /// the component data into the appropriate protocol simulation type based on the current
231    /// blockchain state and the provided block header.
232    /// For example, to register a decoder for the `uniswap_v2` exchange with an additional decoder
233    /// context, you must call this function with
234    /// `register_decoder_with_context::<UniswapV2State>("uniswap_v2", context)`.
235    /// This ensures that the exchange ID `uniswap_v2` is properly associated with the
236    /// `UniswapV2State` decoder for use in the protocol stream.
237    ///
238    /// The decoder's own chain is stamped into `context`, overriding any chain the caller set.
239    pub fn register_decoder_with_context<T>(&mut self, exchange: &str, mut context: DecoderContext)
240    where
241        T: ProtocolSim
242            + TryFromWithBlock<ComponentWithState, H, Error = InvalidSnapshotError>
243            + Send
244            + 'static,
245    {
246        if let Some(requested) = context
247            .chain
248            .filter(|requested| *requested != self.chain)
249        {
250            warn!(
251                exchange,
252                requested_chain = %requested,
253                decoder_chain = %self.chain,
254                "DecoderContext declares a different chain than the decoder; using the decoder's \
255                 chain"
256            );
257        }
258        context.chain = Some(self.chain);
259
260        if is_deprecated_curve_registration::<T>(exchange) {
261            warn!(
262                registered_type = std::any::type_name::<T>(),
263                "Registering \"vm:curve\" with the generic VM adapter is deprecated; register the \
264                 native `CurveState` decoder instead (`exchange::<CurveState>(\"vm:curve\", ...)`). \
265                 The VM-adapter path still works but will be removed in a future release."
266            );
267        }
268        let decoder = Box::new(
269            move |component: ComponentWithState,
270                  header: H,
271                  account_balances: AccountBalances,
272                  state: Arc<RwLock<DecoderState>>,
273                  live_override: Option<watch::Receiver<OverrideSnapshot>>| {
274                let mut context = context.clone();
275                context.live_override = live_override;
276                Box::pin(async move {
277                    let guard = state.read().await;
278                    T::try_from_with_header(
279                        component,
280                        header,
281                        &account_balances,
282                        &guard.tokens,
283                        &context,
284                    )
285                    .await
286                    .map(|c| Box::new(c) as Box<dyn ProtocolSim>)
287                }) as DecodeFut
288            },
289        );
290        self.registry
291            .insert(exchange.to_string(), decoder);
292    }
293
294    /// Registers a decoder for a given exchange.
295    ///
296    /// This method maps an exchange identifier to a specific protocol simulation type.
297    /// The associated type must implement the `TryFromWithBlock` trait to enable decoding
298    /// of state updates from `ComponentWithState` objects. This allows the decoder to transform
299    /// the component data into the appropriate protocol simulation type based on the current
300    /// blockchain state and the provided block header.
301    /// For example, to register a decoder for the `uniswap_v2` exchange, you must call
302    /// this function with `register_decoder::<UniswapV2State>("uniswap_v2", vm_attributes)`.
303    /// This ensures that the exchange ID `uniswap_v2` is properly associated with the
304    /// `UniswapV2State` decoder for use in the protocol stream.
305    pub fn register_decoder<T>(&mut self, exchange: &str)
306    where
307        T: ProtocolSim
308            + TryFromWithBlock<ComponentWithState, H, Error = InvalidSnapshotError>
309            + Send
310            + 'static,
311    {
312        let context = DecoderContext::new();
313        self.register_decoder_with_context::<T>(exchange, context);
314    }
315
316    /// Registers a client-side filter function for a given exchange.
317    ///
318    /// Associates a filter function with an exchange ID, enabling custom filtering of protocol
319    /// components. The filter function is applied client-side to refine the data received from the
320    /// stream. It can be used to exclude certain components based on attributes or conditions that
321    /// are not supported by the server-side filtering logic. This is particularly useful for
322    /// implementing custom behaviors, such as:
323    /// - Filtering out pools with specific attributes (e.g., unsupported features).
324    /// - Blacklisting pools based on custom criteria.
325    /// - Excluding pools that do not meet certain requirements (e.g., token pairs or liquidity
326    ///   constraints).
327    ///
328    /// For example, you might use a filter to exclude pools that are not fully supported in the
329    /// protocol, or to ignore pools with certain attributes that are irrelevant to your
330    /// application.
331    ///
332    /// Filters accumulate: registering a second predicate for the same exchange keeps the first,
333    /// and a component is admitted only when every registered predicate accepts it.
334    pub fn register_filter(&mut self, exchange: &str, predicate: FilterFn) {
335        self.inclusion_filters
336            .entry(exchange.to_string())
337            .or_default()
338            .push(predicate);
339    }
340
341    /// Whether every filter registered for `exchange` accepts `snapshot`. An exchange with no
342    /// registered filter admits every component.
343    fn admits(&self, exchange: &str, snapshot: &ComponentWithState) -> bool {
344        let Some(predicates) = self.inclusion_filters.get(exchange) else { return true };
345        predicates
346            .iter()
347            .all(|predicate| predicate(snapshot))
348    }
349
350    /// Decodes a `FeedMessage` into a `BlockUpdate` containing the updated states of protocol
351    /// components
352    pub async fn decode(&self, msg: &FeedMessage<H>) -> Result<Update, StreamDecodeError> {
353        // stores all states updated in this tick/msg
354        let mut updated_states = HashMap::new();
355        let mut new_pairs = HashMap::new();
356        let mut removed_pairs = HashMap::new();
357        let mut contracts_map = HashMap::new();
358        let mut msg_failed_components = HashSet::new();
359
360        let header = msg
361            .state_msgs
362            .values()
363            .next()
364            .ok_or_else(|| StreamDecodeError::Fatal("Missing block!".into()))?
365            .header
366            .clone();
367
368        let block_number_or_timestamp = header
369            .clone()
370            .block_number_or_timestamp();
371        let current_block = header.clone().block();
372        let is_partial = current_block
373            .as_ref()
374            .map(|h| h.partial_block_index.is_some())
375            .unwrap_or(false);
376
377        for (protocol, protocol_msg) in msg.state_msgs.iter() {
378            // Add any new tokens
379            if let Some(deltas) = protocol_msg.deltas.as_ref() {
380                let mut state_guard = self.state.write().await;
381
382                let new_tokens = deltas
383                    .new_tokens
384                    .iter()
385                    .filter(|(addr, t)| {
386                        t.quality >= self.min_token_quality &&
387                            !state_guard.tokens.contains_key(*addr)
388                    })
389                    .map(|(addr, t)| (addr.clone(), t.clone()))
390                    .collect::<HashMap<Bytes, Token>>();
391
392                if !new_tokens.is_empty() {
393                    debug!(n = new_tokens.len(), "NewTokens");
394                    state_guard.tokens.extend(new_tokens);
395                }
396            }
397
398            // Remove untracked components
399            {
400                let mut state_guard = self.state.write().await;
401                let removed_components: Vec<(String, ProtocolComponent)> = protocol_msg
402                    .removed_components
403                    .iter()
404                    .map(|(id, comp)| {
405                        if *id != comp.id {
406                            error!(
407                                "Component id mismatch in removed components {id} != {}",
408                                comp.id
409                            );
410                            return Err(StreamDecodeError::Fatal("Component id mismatch".into()));
411                        }
412
413                        let tokens = comp
414                            .tokens
415                            .iter()
416                            .flat_map(|addr| state_guard.tokens.get(addr).cloned())
417                            .collect::<Vec<_>>();
418
419                        if tokens.len() == comp.tokens.len() {
420                            Ok(Some((
421                                id.clone(),
422                                ProtocolComponent::from_with_tokens(comp.clone(), tokens),
423                            )))
424                        } else {
425                            Ok(None)
426                        }
427                    })
428                    .collect::<Result<Vec<Option<(String, ProtocolComponent)>>, StreamDecodeError>>(
429                    )?
430                    .into_iter()
431                    .flatten()
432                    .collect();
433
434                // Remove components from state and add to removed_pairs
435                for (id, component) in removed_components {
436                    state_guard.components.remove(&id);
437                    state_guard.states.remove(&id);
438                    removed_pairs.insert(id, component);
439                }
440
441                // UPDATE VM STORAGE
442                info!(
443                    "Processing {} contracts from snapshots",
444                    protocol_msg
445                        .snapshots
446                        .get_vm_storage()
447                        .len()
448                );
449
450                let mut proxy_token_accounts: HashMap<Address, AccountUpdate> = HashMap::new();
451                let mut storage_by_address: HashMap<Address, ResponseAccount> = HashMap::new();
452                for (key, value) in protocol_msg
453                    .snapshots
454                    .get_vm_storage()
455                    .iter()
456                {
457                    let account: ResponseAccount = value.clone().into();
458
459                    if state_guard.tokens.contains_key(key) {
460                        let original_address = account.address;
461                        // To work with Tycho's token overwrites system, if we get account
462                        // snapshots for a token we must handle them with a proxy/wrapper
463                        // contract.
464                        // Note: storage for the original contract must be set at the proxy
465                        // contract address. This is because the proxy contract uses
466                        // delegatecall to the original (implementation) contract.
467
468                        // Handle proxy token accounts
469                        let (impl_addr, proxy_state) = match state_guard
470                            .proxy_token_addresses
471                            .get(&original_address)
472                        {
473                            Some(impl_addr) => {
474                                // Token already has a proxy contract, simply update it.
475
476                                // Note: we apply the snapshot as an update. This is to cover the
477                                // case where a contract may be stale as it stopped being tracked
478                                // for some reason (e.g. due to a drop in tvl) and is now being
479                                // tracked again.
480                                let proxy_state = AccountUpdate::new(
481                                    original_address,
482                                    value.chain,
483                                    account.slots.clone(),
484                                    Some(account.native_balance),
485                                    None,
486                                    ChangeType::Update,
487                                );
488                                (*impl_addr, proxy_state)
489                            }
490                            None => {
491                                // Token does not have a proxy contract yet, create one
492
493                                // Assign original token contract to new address
494                                let impl_addr = generate_proxy_token_address(
495                                    state_guard.proxy_token_addresses.len() as u32,
496                                )?;
497                                state_guard
498                                    .proxy_token_addresses
499                                    .insert(original_address, impl_addr);
500
501                                // Add proxy token contract at original token address
502                                let proxy_state = create_proxy_token_account(
503                                    original_address,
504                                    Some(impl_addr),
505                                    &account.slots,
506                                    value.chain,
507                                    Some(account.native_balance),
508                                );
509
510                                (impl_addr, proxy_state)
511                            }
512                        };
513
514                        proxy_token_accounts.insert(original_address, proxy_state);
515
516                        // Assign original token contract to the implementation address
517                        let impl_update = ResponseAccount {
518                            address: impl_addr,
519                            slots: HashMap::new(),
520                            ..account.clone()
521                        };
522                        storage_by_address.insert(impl_addr, impl_update);
523                    } else {
524                        // Not a token, apply snapshot to the account at its original address
525                        storage_by_address.insert(account.address, account);
526                    }
527                }
528
529                // Split proxy accounts by change type:
530                // - Creation: new proxies that must overwrite any existing placeholder
531                // - Update: existing proxies whose storage is being refreshed (handled normally)
532                let mut proxy_creates: Vec<AccountUpdate> = Vec::new();
533                let mut proxy_updates: HashMap<Address, AccountUpdate> = HashMap::new();
534                for (addr, update) in proxy_token_accounts {
535                    if matches!(update.change, ChangeType::Creation) {
536                        proxy_creates.push(update);
537                    } else {
538                        proxy_updates.insert(addr, update);
539                    }
540                }
541
542                info!("Updating engine with {} contracts from snapshots", storage_by_address.len());
543                update_engine(
544                    SHARED_TYCHO_DB.clone(),
545                    header.clone().block(),
546                    Some(storage_by_address),
547                    proxy_updates,
548                )
549                .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
550
551                // Force-overwrite new proxy token accounts so that authoritative vm_storage data
552                // always wins over any empty placeholder previously inserted by engine setup
553                // (which uses init_account / init-if-not-exists).
554                if !proxy_creates.is_empty() {
555                    SHARED_TYCHO_DB
556                        .force_update_accounts(proxy_creates)
557                        .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
558                }
559                info!("Engine updated");
560                drop(state_guard);
561            }
562
563            // Construct a contract to token balances map: HashMap<ContractAddress,
564            // HashMap<TokenAddress, Balance>>
565            let account_balances = protocol_msg
566                .clone()
567                .snapshots
568                .get_vm_storage()
569                .iter()
570                .filter_map(|(addr, acc)| {
571                    if acc.token_balances.is_empty() {
572                        return None;
573                    }
574                    let balances = acc
575                        .token_balances
576                        .iter()
577                        .map(|(token_addr, ab)| (token_addr.clone(), ab.balance.clone()))
578                        .collect::<HashMap<Bytes, Bytes>>();
579                    Some((addr.clone(), balances))
580                })
581                .collect::<AccountBalances>();
582
583            let mut new_components = HashMap::new();
584            let mut count_token_skips = 0;
585            let mut components_to_store = HashMap::new();
586            {
587                let state_guard = self.state.read().await;
588
589                // PROCESS SNAPSHOTS
590                'snapshot_loop: for (id, snapshot) in protocol_msg
591                    .snapshots
592                    .get_states()
593                    .clone()
594                {
595                    // Skip any unsupported pools
596                    if !self.admits(protocol.as_str(), &snapshot) {
597                        continue;
598                    }
599
600                    // Construct component from snapshot
601                    let mut component_tokens = Vec::new();
602                    let mut new_tokens_accounts = HashMap::new();
603                    for token in snapshot.component.tokens.clone() {
604                        match state_guard.tokens.get(&token) {
605                            Some(token) => {
606                                component_tokens.push(token.clone());
607
608                                // If the token is not an existing proxy token, we need to add it to
609                                // the simulation engine
610                                let token_address = match bytes_to_address(&token.address) {
611                                    Ok(addr) => addr,
612                                    Err(_) => {
613                                        count_token_skips += 1;
614                                        msg_failed_components.insert(id.clone());
615                                        warn!(
616                                            "Token address could not be decoded {}, ignoring pool {:x?}",
617                                            token.address, id
618                                        );
619                                        continue 'snapshot_loop;
620                                    }
621                                };
622                                // Deploy a proxy account without an implementation set
623                                if !state_guard
624                                    .proxy_token_addresses
625                                    .contains_key(&token_address)
626                                {
627                                    new_tokens_accounts.insert(
628                                        token_address,
629                                        create_proxy_token_account(
630                                            token_address,
631                                            None,
632                                            &HashMap::new(),
633                                            snapshot.component.chain,
634                                            None,
635                                        ),
636                                    );
637                                }
638                            }
639                            None => {
640                                count_token_skips += 1;
641                                msg_failed_components.insert(id.clone());
642                                debug!("Token not found {}, ignoring pool {:x?}", token, id);
643                                continue 'snapshot_loop;
644                            }
645                        }
646                    }
647                    let component = ProtocolComponent::from_with_tokens(
648                        snapshot.component.clone(),
649                        component_tokens,
650                    );
651
652                    // Add new tokens to the simulation engine
653                    if !new_tokens_accounts.is_empty() {
654                        update_engine(
655                            SHARED_TYCHO_DB.clone(),
656                            header.clone().block(),
657                            None,
658                            new_tokens_accounts,
659                        )
660                        .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
661                    }
662
663                    // collect contracts:ids mapping for states that should update on contract
664                    // changes (non-manual updates)
665                    if !component
666                        .static_attributes
667                        .contains_key("manual_updates")
668                    {
669                        for contract in &component.contract_ids {
670                            contracts_map
671                                .entry(contract.clone())
672                                .or_insert_with(HashSet::new)
673                                .insert(id.clone());
674                        }
675                        // Add DCI contracts so changes to these contracts trigger
676                        // an update
677                        for (_, tracing) in snapshot.entrypoints.iter() {
678                            for contract in tracing.accessed_slots.keys().cloned() {
679                                contracts_map
680                                    .entry(contract)
681                                    .or_insert_with(HashSet::new)
682                                    .insert(id.clone());
683                            }
684                        }
685                    }
686
687                    // Collect new pairs (components)
688                    new_pairs.insert(id.clone(), component.clone());
689
690                    // Store component for later batch insertion
691                    components_to_store.insert(id.clone(), component);
692
693                    // Construct state from snapshot
694                    if let Some(state_decode_f) = self.registry.get(protocol.as_str()) {
695                        let live_override = self
696                            .override_providers
697                            .get(protocol.as_str())
698                            .and_then(|provider| provider.subscribe(protocol.as_str()));
699                        match state_decode_f(
700                            snapshot,
701                            header.clone(),
702                            account_balances.clone(),
703                            self.state.clone(),
704                            live_override,
705                        )
706                        .await
707                        {
708                            Ok(state) => {
709                                new_components.insert(id.clone(), state);
710                            }
711                            Err(e) => {
712                                if self.skip_state_decode_failures {
713                                    warn!(pool = id, error = %e, "StateDecodingFailure");
714                                    msg_failed_components.insert(id.clone());
715                                    continue 'snapshot_loop;
716                                } else {
717                                    error!(pool = id, error = %e, "StateDecodingFailure");
718                                    return Err(StreamDecodeError::Fatal(format!("{e}")));
719                                }
720                            }
721                        }
722                    } else if self.skip_state_decode_failures {
723                        warn!(pool = id, "MissingDecoderRegistration");
724                        msg_failed_components.insert(id.clone());
725                        continue 'snapshot_loop;
726                    } else {
727                        error!(pool = id, "MissingDecoderRegistration");
728                        return Err(StreamDecodeError::Fatal(format!(
729                            "Missing decoder registration for: {id}"
730                        )));
731                    }
732                }
733            }
734
735            // Batch insert components into state
736            if !components_to_store.is_empty() {
737                let mut state_guard = self.state.write().await;
738                for (id, component) in components_to_store {
739                    state_guard
740                        .components
741                        .insert(id, component);
742                }
743            }
744
745            if !protocol_msg.snapshots.states.is_empty() {
746                info!("Decoded {} snapshots for protocol {protocol}", new_components.len());
747            }
748            if count_token_skips > 0 {
749                info!("Skipped {count_token_skips} pools due to missing tokens");
750            }
751
752            //TODO: should we remove failed components for new_components?
753            updated_states.extend(new_components);
754
755            // PROCESS DELTAS
756            if let Some(deltas) = protocol_msg.deltas.clone() {
757                // Update engine with account changes
758                let mut state_guard = self.state.write().await;
759
760                let mut account_update_by_address: HashMap<Address, AccountUpdate> = HashMap::new();
761                // New proxy token accounts that must overwrite any existing placeholder.
762                let mut new_proxy_accounts: Vec<AccountUpdate> = Vec::new();
763                for (key, value) in deltas.account_deltas.iter() {
764                    let mut update: AccountUpdate = value.clone().into();
765
766                    // TEMP PATCH (ENG-4993)
767                    //
768                    // The indexer may emit Creation deltas with no code for EOA addresses.
769                    // Treat them as EOAs (empty code) rather than downgrading to Update, which
770                    // would skip init_account and cause "uninitialized account" warnings.
771                    if update.code.is_none() && matches!(update.change, ChangeType::Creation) {
772                        error!(
773                            update = ?update,
774                            "FaultyCreationDelta"
775                        );
776                        update.code = Some(vec![]);
777                    }
778
779                    if state_guard.tokens.contains_key(key) {
780                        let original_address = update.address;
781                        // If the account is a token, we need to handle it with a proxy contract.
782                        // Storage updates apply to the proxy contract (at original address).
783                        // Code updates (if any) apply to the token implementation contract (at
784                        // impl_addr).
785
786                        // Handle proxy contract updates
787                        let impl_addr = match state_guard
788                            .proxy_token_addresses
789                            .get(&original_address)
790                        {
791                            Some(impl_addr) => {
792                                // Token already has a proxy contract.
793
794                                // The proxy account already exists, so this is always a plain
795                                // storage update regardless of the incoming change type.
796                                let proxy_update = AccountUpdate {
797                                    code: None,
798                                    change: ChangeType::Update,
799                                    ..update.clone()
800                                };
801                                account_update_by_address.insert(original_address, proxy_update);
802
803                                *impl_addr
804                            }
805                            None => {
806                                // Token does not have a proxy contract yet, create one
807
808                                // Assign original token (implementation) contract to new proxy
809                                // address
810                                let impl_addr = generate_proxy_token_address(
811                                    state_guard.proxy_token_addresses.len() as u32,
812                                )?;
813                                state_guard
814                                    .proxy_token_addresses
815                                    .insert(original_address, impl_addr);
816
817                                // Create proxy token account with original account's storage (at
818                                // original address). Track it separately so it can be
819                                // force-overwritten and win over any placeholder that an engine
820                                // setup routine may have written earlier.
821                                let proxy_state = create_proxy_token_account(
822                                    original_address,
823                                    Some(impl_addr),
824                                    &update.slots,
825                                    update.chain,
826                                    update.balance,
827                                );
828                                new_proxy_accounts.push(proxy_state);
829
830                                impl_addr
831                            }
832                        };
833
834                        // Apply code update to token implementation contract
835                        if update.code.is_some() {
836                            let impl_update = AccountUpdate {
837                                address: impl_addr,
838                                slots: HashMap::new(),
839                                ..update.clone()
840                            };
841                            account_update_by_address.insert(impl_addr, impl_update);
842                        }
843                    } else {
844                        // Not a token, apply update to the account at its original address
845                        account_update_by_address.insert(update.address, update);
846                    }
847                }
848                drop(state_guard);
849
850                let state_guard = self.state.read().await;
851                info!("Updating engine with {} contract deltas", deltas.account_deltas.len());
852                update_engine(
853                    SHARED_TYCHO_DB.clone(),
854                    header.clone().block(),
855                    None,
856                    account_update_by_address,
857                )
858                .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
859
860                // Force-overwrite any newly-created proxy token accounts so they always win
861                // over placeholder entries inserted by engine setup.
862                if !new_proxy_accounts.is_empty() {
863                    SHARED_TYCHO_DB
864                        .force_update_accounts(new_proxy_accounts)
865                        .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
866                }
867                info!("Engine updated");
868
869                // Collect all pools related to the updated accounts
870                let mut pools_to_update = HashSet::new();
871                for (account, _update) in deltas.account_deltas {
872                    // get new pools related to the account updated
873                    pools_to_update.extend(
874                        contracts_map
875                            .get(&account)
876                            .cloned()
877                            .unwrap_or_default(),
878                    );
879                    // get existing pools related to the account updated
880                    pools_to_update.extend(
881                        state_guard
882                            .contracts_map
883                            .get(&account)
884                            .cloned()
885                            .unwrap_or_default(),
886                    );
887                }
888
889                // Collect all balance changes this block
890                let all_balances = Balances {
891                    component_balances: deltas
892                        .component_balances
893                        .iter()
894                        .map(|(pool_id, bals)| {
895                            let mut balances = HashMap::new();
896                            for (t, b) in bals {
897                                balances.insert(t.clone(), b.balance.clone());
898                            }
899                            pools_to_update.insert(pool_id.clone());
900                            (pool_id.clone(), balances)
901                        })
902                        .collect(),
903                    account_balances: deltas
904                        .account_balances
905                        .iter()
906                        .map(|(account, bals)| {
907                            let mut balances = HashMap::new();
908                            for (t, b) in bals {
909                                balances.insert(t.clone(), b.balance.clone());
910                            }
911                            pools_to_update.extend(
912                                contracts_map
913                                    .get(account)
914                                    .cloned()
915                                    .unwrap_or_default(),
916                            );
917                            (account.clone(), balances)
918                        })
919                        .collect(),
920                };
921
922                // update states with protocol state deltas (attribute changes etc.)
923                for (id, update) in deltas.state_deltas {
924                    // TODO: is this needed?
925                    let update_with_block = Self::add_block_info_to_delta(
926                        ProtocolStateDelta::from(update),
927                        current_block.clone(),
928                    );
929                    match Self::apply_update(
930                        &id,
931                        update_with_block,
932                        &mut updated_states,
933                        &state_guard,
934                        &all_balances,
935                    ) {
936                        Ok(_) => {
937                            pools_to_update.remove(&id);
938                        }
939                        Err(e) => {
940                            if self.skip_state_decode_failures {
941                                warn!(pool = id, error = %e, "Failed to apply state update, marking component as removed");
942                                // Remove from updated_states if it was there
943                                updated_states.remove(&id);
944                                // Try to get component from new_pairs first, then from state
945                                if let Some(component) = new_pairs.remove(&id) {
946                                    removed_pairs.insert(id.clone(), component);
947                                } else if let Some(component) = state_guard.components.get(&id) {
948                                    removed_pairs.insert(id.clone(), component.clone());
949                                } else {
950                                    // Component not found in new_pairs or state, this shouldn't
951                                    // happen
952                                    warn!(pool = id, "Component not found in new_pairs or state, cannot add to removed_pairs");
953                                }
954                                pools_to_update.remove(&id);
955
956                                // Add to failed components
957                                msg_failed_components.insert(id.clone());
958                            } else {
959                                return Err(e);
960                            }
961                        }
962                    }
963                }
964
965                // update remaining pools linked to updated contracts/updated balances
966                for pool in pools_to_update {
967                    // TODO: is this needed?
968                    let default_delta_with_block = Self::add_block_info_to_delta(
969                        ProtocolStateDelta::default(),
970                        current_block.clone(),
971                    );
972                    match Self::apply_update(
973                        &pool,
974                        default_delta_with_block,
975                        &mut updated_states,
976                        &state_guard,
977                        &all_balances,
978                    ) {
979                        Ok(_) => {}
980                        Err(e) => {
981                            if self.skip_state_decode_failures {
982                                warn!(pool = pool, error = %e, "Failed to apply contract/balance update, marking component as removed");
983                                // Remove from updated_states if it was there
984                                updated_states.remove(&pool);
985                                // Try to get component from new_pairs first, then from state
986                                if let Some(component) = new_pairs.remove(&pool) {
987                                    removed_pairs.insert(pool.clone(), component);
988                                } else if let Some(component) = state_guard.components.get(&pool) {
989                                    removed_pairs.insert(pool.clone(), component.clone());
990                                } else {
991                                    // Component not found in new_pairs or state, this shouldn't
992                                    // happen
993                                    warn!(pool = pool, "Component not found in new_pairs or state, cannot add to removed_pairs");
994                                }
995
996                                // Add to failed components
997                                msg_failed_components.insert(pool.clone());
998                            } else {
999                                return Err(e);
1000                            }
1001                        }
1002                    }
1003                }
1004            };
1005        }
1006
1007        // Persist the newly added/updated states
1008        let mut state_guard = self.state.write().await;
1009
1010        // Update failed components with any new ones
1011        state_guard
1012            .failed_components
1013            .extend(msg_failed_components);
1014
1015        // Remove any failed components from Updates
1016        // Perf: we could do it directly in the decoder logic to avoid some steps, but this logic is
1017        // complex and this is more robust.
1018        updated_states.retain(|id, _| {
1019            !state_guard
1020                .failed_components
1021                .contains(id)
1022        });
1023        new_pairs.retain(|id, _| {
1024            !state_guard
1025                .failed_components
1026                .contains(id)
1027        });
1028
1029        if let Some(header) = current_block.as_ref() {
1030            let execution_block = self.execution_block(header);
1031            let decoder_state = &mut *state_guard;
1032            Self::refresh_execution_block(
1033                &mut updated_states,
1034                &mut decoder_state.states,
1035                &decoder_state.failed_components,
1036                &removed_pairs,
1037                &execution_block,
1038            );
1039        }
1040
1041        state_guard
1042            .states
1043            .extend(updated_states.clone());
1044
1045        state_guard.current_block_number = block_number_or_timestamp;
1046
1047        // Add new components to persistent state
1048        for (id, component) in new_pairs.iter() {
1049            state_guard
1050                .components
1051                .insert(id.clone(), component.clone());
1052        }
1053
1054        // Remove components from persistent state
1055        for id in removed_pairs.keys() {
1056            state_guard.components.remove(id);
1057        }
1058
1059        for (key, values) in contracts_map {
1060            state_guard
1061                .contracts_map
1062                .entry(key)
1063                .or_insert_with(HashSet::new)
1064                .extend(values);
1065        }
1066
1067        // Send the tick with all updated states
1068        Ok(Update::new(block_number_or_timestamp, updated_states, new_pairs)
1069            .set_is_partial(is_partial)
1070            .set_removed_pairs(removed_pairs)
1071            .set_sync_states(msg.sync_states.clone()))
1072    }
1073
1074    /// Applies pending deltas from one or more `TxDeltaIndexer`s against the current confirmed
1075    /// state and returns an ephemeral `Update`.
1076    ///
1077    /// This is the read-only counterpart of `decode()`. It clones pool states, applies the
1078    /// supplied `pending_deltas`, and returns the result — **without writing back** to
1079    /// `DecoderState`. Calling this method twice with the same input produces identical results.
1080    ///
1081    /// Every state is rebuilt from `state_deltas` alone; nothing here writes to the VM database.
1082    /// A protocol decoding into the generic VM adapter therefore cannot take part: it re-reads
1083    /// pool state from that database, so its storage-derived values stay at the confirmed block —
1084    /// even though the delta's balance and block-environment attributes do get applied.
1085    ///
1086    /// # Parameters
1087    /// * `pending_deltas` — map from extractor name to the `BlockAggregatedChanges` produced by the
1088    ///   corresponding `TxDeltaIndexer::generate_deltas()` call.
1089    /// * `header` — the target block header. Its `block_number_or_timestamp()` is stamped on the
1090    ///   returned [`Update`]; its `block_number` and `block_timestamp` are injected into each state
1091    ///   delta so that protocols relying on block context (e.g. aerodrome slipstreams, etherfi)
1092    ///   receive correct values.
1093    pub async fn apply_deltas_ephemeral(
1094        &self,
1095        pending_deltas: &HashMap<String, BlockAggregatedChanges>,
1096        header: H,
1097    ) -> Result<Update, StreamDecodeError> {
1098        let block_number_or_timestamp = header
1099            .clone()
1100            .block_number_or_timestamp();
1101        let current_block = header.block();
1102        let state_guard = self.state.read().await;
1103
1104        let mut updated_states: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1105
1106        for deltas in pending_deltas.values() {
1107            let all_balances = Balances {
1108                component_balances: deltas
1109                    .component_balances
1110                    .iter()
1111                    .map(|(pool_id, bals)| {
1112                        let balances = bals
1113                            .iter()
1114                            .map(|(t, b)| (t.clone(), b.balance.clone()))
1115                            .collect();
1116                        (pool_id.clone(), balances)
1117                    })
1118                    .collect(),
1119                account_balances: HashMap::new(),
1120            };
1121
1122            for (id, state_delta) in &deltas.state_deltas {
1123                let dto_delta = Self::add_block_info_to_delta(
1124                    ProtocolStateDelta::from(state_delta.clone()),
1125                    current_block.clone(),
1126                );
1127                if let Err(e) = Self::apply_update(
1128                    id,
1129                    dto_delta,
1130                    &mut updated_states,
1131                    &state_guard,
1132                    &all_balances,
1133                ) {
1134                    warn!(pool = id, error = %e, "EphemeralDeltaTransitionError");
1135                }
1136            }
1137        }
1138
1139        // `header` is the block being built, so it already *is* the execution block — unlike
1140        // `decode`, there is nothing to project forward here. Only the delta-applied clones need
1141        // advancing: this path is read-only, and the stored states were already advanced to this
1142        // block by `decode()` on the confirmed stream.
1143        if let Some(header) = current_block.as_ref() {
1144            let execution_block = BlockContext::new(header.number, header.timestamp);
1145            for state in updated_states.values_mut() {
1146                state.apply_block(&execution_block);
1147            }
1148        }
1149
1150        Ok(Update::new(block_number_or_timestamp, updated_states, HashMap::new()))
1151    }
1152
1153    /// Add current block information (number and timestamp) to a ProtocolStateDelta.
1154    fn add_block_info_to_delta(
1155        mut delta: ProtocolStateDelta,
1156        block_header_opt: Option<BlockHeader>,
1157    ) -> ProtocolStateDelta {
1158        if let Some(header) = block_header_opt {
1159            // Add block_number and block_timestamp attributes to ensure pool states
1160            // receive current block information during delta_transition
1161            delta.updated_attributes.insert(
1162                "block_number".to_string(),
1163                Bytes::from(header.number.to_be_bytes().to_vec()),
1164            );
1165            delta.updated_attributes.insert(
1166                "block_timestamp".to_string(),
1167                Bytes::from(header.timestamp.to_be_bytes().to_vec()),
1168            );
1169        }
1170        delta
1171    }
1172
1173    fn apply_update(
1174        id: &String,
1175        update: ProtocolStateDelta,
1176        updated_states: &mut HashMap<String, Box<dyn ProtocolSim>>,
1177        state_guard: &RwLockReadGuard<'_, DecoderState>,
1178        all_balances: &Balances,
1179    ) -> Result<(), StreamDecodeError> {
1180        match updated_states.entry(id.clone()) {
1181            Entry::Occupied(mut entry) => {
1182                // If state exists in updated_states, apply the delta to it
1183                let state: &mut Box<dyn ProtocolSim> = entry.get_mut();
1184                state
1185                    .delta_transition(update, &state_guard.tokens, all_balances)
1186                    .map_err(|e| {
1187                        error!(pool = id, error = ?e, "DeltaTransitionError");
1188                        StreamDecodeError::Fatal(format!("TransitionFailure: {e:?}"))
1189                    })?;
1190            }
1191            Entry::Vacant(_) => {
1192                match state_guard.states.get(id) {
1193                    // If state does not exist in updated_states, apply the delta to the stored
1194                    // state
1195                    Some(stored_state) => {
1196                        let mut state = stored_state.clone();
1197                        state
1198                            .delta_transition(update, &state_guard.tokens, all_balances)
1199                            .map_err(|e| {
1200                                error!(pool = id, error = ?e, "DeltaTransitionError");
1201                                StreamDecodeError::Fatal(format!("TransitionFailure: {e:?}"))
1202                            })?;
1203                        updated_states.insert(id.clone(), state);
1204                    }
1205                    None => debug!(pool = id, reason = "MissingState", "DeltaTransitionError"),
1206                }
1207            }
1208        }
1209        Ok(())
1210    }
1211}
1212
1213/// Generate a proxy token address for a given token index
1214fn generate_proxy_token_address(idx: u32) -> Result<Address, StreamDecodeError> {
1215    let padded_idx = format!("{idx:x}");
1216    let padded_zeroes = "0".repeat(33 - padded_idx.len());
1217    let proxy_token_address = format!("{padded_zeroes}{padded_idx}BAdbaBe");
1218    let decoded = hex::decode(proxy_token_address).map_err(|e| {
1219        StreamDecodeError::Fatal(format!("Invalid proxy token address encoding: {e}"))
1220    })?;
1221
1222    const ADDRESS_LENGTH: usize = 20;
1223    if decoded.len() != ADDRESS_LENGTH {
1224        return Err(StreamDecodeError::Fatal(format!(
1225            "Invalid proxy token address length: expected {}, got {}",
1226            ADDRESS_LENGTH,
1227            decoded.len(),
1228        )));
1229    }
1230
1231    Ok(Address::from_slice(&decoded))
1232}
1233
1234/// Create a proxy token account for a token at a given address
1235///
1236/// The proxy token account is created at the original token address and points to the new token
1237/// address.
1238fn create_proxy_token_account(
1239    addr: Address,
1240    new_address: Option<Address>,
1241    storage: &HashMap<U256, U256>,
1242    chain: Chain,
1243    balance: Option<U256>,
1244) -> AccountUpdate {
1245    let mut slots = storage.clone();
1246    if let Some(new_address) = new_address {
1247        slots.insert(*IMPLEMENTATION_SLOT, U256::from_be_slice(new_address.as_slice()));
1248    }
1249
1250    AccountUpdate {
1251        address: addr,
1252        chain,
1253        slots,
1254        balance,
1255        code: Some(ERC20_PROXY_BYTECODE.to_vec()),
1256        change: ChangeType::Creation,
1257    }
1258}
1259
1260#[cfg(test)]
1261mock! {
1262    #[derive(Debug)]
1263    pub ProtocolSim {
1264        pub fn fee(&self) -> f64;
1265        pub fn spot_price(&self, base: &Token, quote: &Token) -> Result<f64, SimulationError>;
1266        pub fn get_amount_out(
1267            &self,
1268            amount_in: BigUint,
1269            token_in: &Token,
1270            token_out: &Token,
1271        ) -> Result<GetAmountOutResult, SimulationError>;
1272        pub fn get_limits(
1273            &self,
1274            sell_token: Bytes,
1275            buy_token: Bytes,
1276        ) -> Result<(BigUint, BigUint), SimulationError>;
1277        pub fn delta_transition(
1278            &mut self,
1279            delta: ProtocolStateDelta,
1280            tokens: &HashMap<Bytes, Token>,
1281            balances: &Balances,
1282        ) -> Result<(), TransitionError>;
1283        pub fn clone_box(&self) -> Box<dyn ProtocolSim>;
1284        pub fn eq(&self, other: &dyn ProtocolSim) -> bool;
1285    }
1286}
1287
1288#[cfg(test)]
1289crate::impl_non_serializable_protocol!(MockProtocolSim, "test protocol");
1290
1291#[cfg(test)]
1292impl ProtocolSim for MockProtocolSim {
1293    fn fee(&self) -> f64 {
1294        self.fee()
1295    }
1296
1297    fn spot_price(&self, base: &Token, quote: &Token) -> Result<f64, SimulationError> {
1298        self.spot_price(base, quote)
1299    }
1300
1301    fn get_amount_out(
1302        &self,
1303        amount_in: BigUint,
1304        token_in: &Token,
1305        token_out: &Token,
1306    ) -> Result<GetAmountOutResult, SimulationError> {
1307        self.get_amount_out(amount_in, token_in, token_out)
1308    }
1309
1310    fn get_limits(
1311        &self,
1312        sell_token: Bytes,
1313        buy_token: Bytes,
1314    ) -> Result<(BigUint, BigUint), SimulationError> {
1315        self.get_limits(sell_token, buy_token)
1316    }
1317
1318    fn delta_transition(
1319        &mut self,
1320        delta: ProtocolStateDelta,
1321        tokens: &HashMap<Bytes, Token>,
1322        balances: &Balances,
1323    ) -> Result<(), TransitionError> {
1324        self.delta_transition(delta, tokens, balances)
1325    }
1326
1327    fn clone_box(&self) -> Box<dyn ProtocolSim> {
1328        self.clone_box()
1329    }
1330
1331    fn as_any(&self) -> &dyn Any {
1332        panic!("MockProtocolSim does not support as_any")
1333    }
1334
1335    fn as_any_mut(&mut self) -> &mut dyn Any {
1336        panic!("MockProtocolSim does not support as_any_mut")
1337    }
1338
1339    fn eq(&self, other: &dyn ProtocolSim) -> bool {
1340        self.eq(other)
1341    }
1342
1343    fn typetag_name(&self) -> &'static str {
1344        unreachable!()
1345    }
1346
1347    fn typetag_deserialize(&self) {
1348        unreachable!()
1349    }
1350}
1351
1352#[cfg(test)]
1353mod tests {
1354    use std::str::FromStr;
1355
1356    use alloy::primitives::address;
1357    use mockall::predicate::*;
1358    use rstest::*;
1359    use tycho_client::feed::BlockHeader;
1360    use tycho_common::{models::Chain, Bytes};
1361
1362    use super::*;
1363
1364    fn header_at(number: u64, timestamp: u64, partial: Option<u32>) -> BlockHeader {
1365        BlockHeader {
1366            hash: Bytes::from([0u8; 32]),
1367            number,
1368            parent_hash: Bytes::from([0u8; 32]),
1369            revert: false,
1370            timestamp,
1371            partial_block_index: partial,
1372        }
1373    }
1374
1375    /// A block-sensitive state whose execution timestamp we can read back.
1376    fn block_sensitive_state() -> Box<dyn ProtocolSim> {
1377        use crate::evm::protocol::{
1378            aerodrome_slipstreams::state::AerodromeSlipstreamsState,
1379            utils::{
1380                slipstreams::{dynamic_fee_module::DynamicFeeConfig, observations::Observation},
1381                uniswap::{tick_list::TickInfo, tick_math::get_sqrt_ratio_at_tick},
1382            },
1383        };
1384
1385        Box::new(
1386            AerodromeSlipstreamsState::new(
1387                "block-sensitive".to_string(),
1388                0,
1389                1_000_000_000_000_000_000,
1390                get_sqrt_ratio_at_tick(0).unwrap(),
1391                0,
1392                1,
1393                3000,
1394                1,
1395                0,
1396                vec![TickInfo::new(-120, 0).unwrap(), TickInfo::new(120, 0).unwrap()],
1397                vec![Observation { block_timestamp: 500, initialized: true, ..Default::default() }],
1398                DynamicFeeConfig::new(2700, 30_000, 0, true, 750),
1399            )
1400            .expect("state should build")
1401            // These fixtures exercise the fee-flip path, which needs the optimistic mode:
1402            // under the worst-case default a flat-fee pool never flips.
1403            .with_position_assumption(crate::protocol::models::BlockPositionAssumption::First),
1404        )
1405    }
1406
1407    #[test]
1408    fn confirmed_header_targets_the_next_block() {
1409        let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Base);
1410
1411        let execution_block = decoder.execution_block(&header_at(100, 1_000, None));
1412
1413        assert_eq!(execution_block.number(), 101);
1414        assert_eq!(execution_block.timestamp(), 1_000 + Chain::Base.block_time_secs());
1415    }
1416
1417    #[test]
1418    fn partial_header_targets_the_block_that_is_still_open() {
1419        let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
1420
1421        let execution_block = decoder.execution_block(&header_at(100, 1_000, Some(3)));
1422
1423        assert_eq!(execution_block.number(), 100);
1424        assert_eq!(execution_block.timestamp(), 1_000);
1425    }
1426
1427    #[test]
1428    fn refresh_re_emits_a_state_whose_fee_flipped_without_a_delta() {
1429        // The pool traded in block 100 and has no delta afterwards. Crossing into block 101
1430        // flips its fee branch (dynamic -> initial), so the refresh must advance the stored
1431        // state in place and emit a copy to consumers.
1432        let mut stored = HashMap::from([("block-sensitive".to_string(), {
1433            let mut state = block_sensitive_state();
1434            state.apply_block(&BlockContext::new(100, 500));
1435            state
1436        })]);
1437        let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1438
1439        TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1440            &mut updated,
1441            &mut stored,
1442            &HashSet::new(),
1443            &HashMap::<String, ()>::new(),
1444            &BlockContext::new(101, 502),
1445        );
1446
1447        let emitted = updated
1448            .get("block-sensitive")
1449            .expect("a fee flip must be emitted even without a delta");
1450        assert_eq!(emitted.fee(), 750.0 / 1_000_000.0);
1451        // The stored copy was advanced in place as well.
1452        assert_eq!(stored["block-sensitive"].fee(), 750.0 / 1_000_000.0);
1453    }
1454
1455    #[test]
1456    fn refresh_never_re_emits_failed_components() {
1457        // A failed component's stored state may linger; the sweep must not resurrect it for
1458        // consumers that were told the component was removed — even when its fee flipped.
1459        let mut stored = HashMap::from([("zombie".to_string(), {
1460            let mut state = block_sensitive_state();
1461            state.apply_block(&BlockContext::new(100, 500));
1462            state
1463        })]);
1464        let failed = HashSet::from(["zombie".to_string()]);
1465        let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1466
1467        TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1468            &mut updated,
1469            &mut stored,
1470            &failed,
1471            &HashMap::<String, ()>::new(),
1472            &BlockContext::new(101, 502),
1473        );
1474
1475        assert!(updated.is_empty());
1476    }
1477
1478    #[test]
1479    fn refresh_never_re_emits_removed_components() {
1480        // The stored state of a component reported in `removed_pairs` stays in place; the sweep
1481        // must not surface it again, even when its fee flipped in the execution block.
1482        let mut stored = HashMap::from([("gone".to_string(), {
1483            let mut state = block_sensitive_state();
1484            state.apply_block(&BlockContext::new(100, 500));
1485            state
1486        })]);
1487        let removed = HashMap::from([("gone".to_string(), ())]);
1488        let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1489
1490        TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1491            &mut updated,
1492            &mut stored,
1493            &HashSet::new(),
1494            &removed,
1495            &BlockContext::new(101, 502),
1496        );
1497
1498        assert!(updated.is_empty());
1499    }
1500
1501    #[test]
1502    fn refresh_stays_quiet_when_no_fee_changed() {
1503        // An idle block-sensitive pool (fee branch unchanged) and a block-insensitive pool:
1504        // neither must be re-emitted.
1505        let mut stored: HashMap<String, Box<dyn ProtocolSim>> = HashMap::from([
1506            ("idle-sensitive".to_string(), {
1507                let mut state = block_sensitive_state();
1508                state.apply_block(&BlockContext::new(101, 502));
1509                state
1510            }),
1511            (
1512                "univ2".to_string(),
1513                Box::new(crate::evm::protocol::uniswap_v2::state::UniswapV2State::new(
1514                    U256::from(1_000_000u64),
1515                    U256::from(1_000_000u64),
1516                )) as Box<dyn ProtocolSim>,
1517            ),
1518        ]);
1519        let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1520
1521        TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1522            &mut updated,
1523            &mut stored,
1524            &HashSet::new(),
1525            &HashMap::<String, ()>::new(),
1526            &BlockContext::new(102, 504),
1527        );
1528
1529        assert!(updated.is_empty());
1530    }
1531    use crate::evm::protocol::{curve::CurveState, uniswap_v2::state::UniswapV2State};
1532
1533    #[test]
1534    fn curve_vm_adapter_registration_flagged_deprecated() {
1535        // The native decoder is the supported path — not flagged.
1536        assert!(!is_deprecated_curve_registration::<CurveState>("vm:curve"));
1537        // Any other type for vm:curve is the deprecated VM-adapter path.
1538        assert!(is_deprecated_curve_registration::<UniswapV2State>("vm:curve"));
1539        // Other exchanges are unaffected.
1540        assert!(!is_deprecated_curve_registration::<UniswapV2State>("uniswap_v2"));
1541    }
1542
1543    async fn setup_decoder(set_tokens: bool) -> TychoStreamDecoder<BlockHeader> {
1544        let mut decoder = TychoStreamDecoder::new(Chain::Ethereum);
1545        decoder.register_decoder::<UniswapV2State>("uniswap_v2");
1546        if set_tokens {
1547            let tokens = [
1548                Bytes::from("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").lpad(20, 0),
1549                Bytes::from("0xdac17f958d2ee523a2206206994597c13d831ec7").lpad(20, 0),
1550            ]
1551            .iter()
1552            .map(|addr| {
1553                let addr_str = format!("{addr:x}");
1554                (
1555                    addr.clone(),
1556                    Token::new(addr, &addr_str, 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1557                )
1558            })
1559            .collect();
1560            decoder.set_tokens(tokens).await;
1561        }
1562        decoder
1563    }
1564
1565    fn load_test_msg(name: &str) -> FeedMessage<BlockHeader> {
1566        use std::{fs, path::Path};
1567
1568        use tycho_client::feed::dto;
1569        let project_root = env!("CARGO_MANIFEST_DIR");
1570        let asset_path = Path::new(project_root).join(format!("tests/assets/decoder/{name}.json"));
1571        let json_data = fs::read_to_string(asset_path).expect("Failed to read test asset");
1572        let feed_msg: dto::FeedMessage<BlockHeader> =
1573            serde_json::from_str(&json_data).expect("Failed to deserialize FeedMsg json!");
1574        FeedMessage::from(feed_msg)
1575    }
1576
1577    #[tokio::test]
1578    async fn test_decode() {
1579        let decoder = setup_decoder(true).await;
1580
1581        let msg = load_test_msg("uniswap_v2_snapshot");
1582        let res1 = decoder
1583            .decode(&msg)
1584            .await
1585            .expect("decode failure");
1586        let msg = load_test_msg("uniswap_v2_delta");
1587        let res2 = decoder
1588            .decode(&msg)
1589            .await
1590            .expect("decode failure");
1591
1592        assert_eq!(res1.states.len(), 1);
1593        assert_eq!(res2.states.len(), 1);
1594        assert_eq!(res1.sync_states.len(), 1);
1595        assert_eq!(res2.sync_states.len(), 1);
1596    }
1597
1598    #[tokio::test]
1599    async fn test_decode_token_creation_delta_with_existing_proxy() {
1600        let decoder = setup_decoder(true).await;
1601        let msg = load_test_msg("uniswap_v2_delta_token_creation");
1602
1603        // First decode: the token has no proxy yet, so the Creation delta takes the
1604        // proxy-creating branch.
1605        decoder
1606            .decode(&msg)
1607            .await
1608            .expect("first decode (proxy creation) failed");
1609
1610        // Second decode: the proxy exists, so the same Creation delta must decode as a
1611        // storage update on the proxy account — not a code-less Creation, which the
1612        // engine rejects as "MissingCode".
1613        decoder
1614            .decode(&msg)
1615            .await
1616            .expect("decode of a token Creation delta with an existing proxy failed");
1617    }
1618
1619    #[tokio::test]
1620    async fn test_decode_component_missing_token() {
1621        let decoder = setup_decoder(false).await;
1622        let tokens = [Bytes::from("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").lpad(20, 0)]
1623            .iter()
1624            .map(|addr| {
1625                let addr_str = format!("{addr:x}");
1626                (
1627                    addr.clone(),
1628                    Token::new(addr, &addr_str, 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1629                )
1630            })
1631            .collect();
1632        decoder.set_tokens(tokens).await;
1633
1634        let msg = load_test_msg("uniswap_v2_snapshot");
1635        let res1 = decoder
1636            .decode(&msg)
1637            .await
1638            .expect("decode failure");
1639
1640        assert_eq!(res1.states.len(), 0);
1641    }
1642
1643    #[tokio::test]
1644    async fn test_decode_component_bad_id() {
1645        let decoder = setup_decoder(true).await;
1646        let msg = load_test_msg("uniswap_v2_snapshot_broken_id");
1647
1648        match decoder.decode(&msg).await {
1649            Err(StreamDecodeError::Fatal(msg)) => {
1650                assert_eq!(msg, "Component id mismatch");
1651            }
1652            Ok(_) => {
1653                panic!("Expected failures to be raised")
1654            }
1655        }
1656    }
1657
1658    #[rstest]
1659    #[case(true)]
1660    #[case(false)]
1661    #[tokio::test]
1662    async fn test_decode_component_bad_state(#[case] skip_failures: bool) {
1663        let mut decoder = setup_decoder(true).await;
1664        decoder.skip_state_decode_failures = skip_failures;
1665
1666        let msg = load_test_msg("uniswap_v2_snapshot_broken_state");
1667        match decoder.decode(&msg).await {
1668            Err(StreamDecodeError::Fatal(msg)) => {
1669                if !skip_failures {
1670                    assert_eq!(msg, "Missing attributes reserve0");
1671                } else {
1672                    panic!("Expected failures to be ignored. Err: {msg}")
1673                }
1674            }
1675            Ok(res) => {
1676                if !skip_failures {
1677                    panic!("Expected failures to be raised")
1678                } else {
1679                    assert_eq!(res.states.len(), 0);
1680                }
1681            }
1682        }
1683    }
1684
1685    #[tokio::test]
1686    async fn test_decode_updates_state_on_contract_change() {
1687        let decoder = setup_decoder(true).await;
1688
1689        // Create the mock instances
1690        let mut mock_state = MockProtocolSim::new();
1691
1692        mock_state
1693            .expect_clone_box()
1694            .times(1)
1695            .returning(|| {
1696                let mut cloned_mock_state = MockProtocolSim::new();
1697                // Expect `delta_transition` to be called once with any parameters
1698                cloned_mock_state
1699                    .expect_delta_transition()
1700                    .times(1)
1701                    .returning(|_, _, _| Ok(()));
1702                cloned_mock_state
1703                    .expect_clone_box()
1704                    .times(1)
1705                    .returning(|| Box::new(MockProtocolSim::new()));
1706                Box::new(cloned_mock_state)
1707            });
1708
1709        // Insert mock state into `updated_states`
1710        let pool_id =
1711            "0x93d199263632a4ef4bb438f1feb99e57b4b5f0bd0000000000000000000005c2".to_string();
1712        decoder
1713            .state
1714            .write()
1715            .await
1716            .states
1717            .insert(pool_id.clone(), Box::new(mock_state) as Box<dyn ProtocolSim>);
1718        decoder
1719            .state
1720            .write()
1721            .await
1722            .contracts_map
1723            .insert(
1724                Bytes::from("0xba12222222228d8ba445958a75a0704d566bf2c8").lpad(20, 0),
1725                HashSet::from([pool_id.clone()]),
1726            );
1727
1728        // Load a test message containing a contract update
1729        let msg = load_test_msg("balancer_v2_delta");
1730
1731        // Decode the message
1732        let _ = decoder
1733            .decode(&msg)
1734            .await
1735            .expect("decode failure");
1736
1737        // The mock framework will assert that `delta_transition` was called exactly once
1738    }
1739
1740    #[test]
1741    fn test_generate_proxy_token_address() {
1742        let idx = 1;
1743        let generated_address =
1744            generate_proxy_token_address(idx).expect("proxy token address should be valid");
1745        assert_eq!(generated_address, address!("000000000000000000000000000000001badbabe"));
1746
1747        let idx = 123456;
1748        let generated_address =
1749            generate_proxy_token_address(idx).expect("proxy token address should be valid");
1750        assert_eq!(generated_address, address!("00000000000000000000000000001e240badbabe"));
1751    }
1752
1753    /// A single-block feed message holding one tokenless Uniswap V4 pool whose `hooks` attribute
1754    /// points at an address no native handler is registered for.
1755    ///
1756    /// The component has no tokens and the message no VM storage, so decoding it neither reads
1757    /// nor writes the process-wide simulation engine and cannot disturb other tests.
1758    fn hooked_v4_msg() -> FeedMessage<BlockHeader> {
1759        use tycho_client::feed::synchronizer::{Snapshot, StateSyncMessage};
1760        use tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState};
1761
1762        let id = "0x00000000000000000000000000000000000000000000000000000000000000c4";
1763        let static_attributes = HashMap::from([
1764            ("key_lp_fee".to_string(), Bytes::from(500_i32.to_be_bytes().to_vec())),
1765            ("tick_spacing".to_string(), Bytes::from(60_i32.to_be_bytes().to_vec())),
1766            (
1767                "hooks".to_string(),
1768                Bytes::from_str("0x00000000000000000000000000000000000000c4").unwrap(),
1769            ),
1770        ]);
1771        let attributes = HashMap::from([
1772            ("liquidity".to_string(), Bytes::from(100_u64.to_be_bytes().to_vec())),
1773            ("tick".to_string(), Bytes::from(300_i32.to_be_bytes().to_vec())),
1774            (
1775                "sqrt_price_x96".to_string(),
1776                Bytes::from(
1777                    79228162514264337593543950336_u128
1778                        .to_be_bytes()
1779                        .to_vec(),
1780                ),
1781            ),
1782            ("protocol_fees/zero2one".to_string(), Bytes::from(0_u32.to_be_bytes().to_vec())),
1783            ("protocol_fees/one2zero".to_string(), Bytes::from(0_u32.to_be_bytes().to_vec())),
1784            ("ticks/60/net_liquidity".to_string(), Bytes::from(400_i128.to_be_bytes().to_vec())),
1785        ]);
1786
1787        let snapshot = ComponentWithState {
1788            state: ProtocolComponentState::new(id, attributes, HashMap::new()),
1789            component: ProtocolComponent {
1790                id: id.to_string(),
1791                static_attributes,
1792                ..Default::default()
1793            },
1794            component_tvl: None,
1795            entrypoints: Vec::new(),
1796        };
1797
1798        FeedMessage {
1799            state_msgs: HashMap::from([(
1800                "uniswap_v4_hooks".to_string(),
1801                StateSyncMessage {
1802                    header: header_at(1, 1_000, None),
1803                    snapshots: Snapshot {
1804                        states: HashMap::from([(id.to_string(), snapshot)]),
1805                        vm_storage: HashMap::new(),
1806                    },
1807                    deltas: None,
1808                    removed_components: HashMap::new(),
1809                },
1810            )]),
1811            sync_states: HashMap::new(),
1812        }
1813    }
1814
1815    /// A decoder for `chain` that decodes [`hooked_v4_msg`].
1816    fn hooked_v4_decoder(chain: Chain, context: DecoderContext) -> TychoStreamDecoder<BlockHeader> {
1817        let mut decoder = TychoStreamDecoder::new(chain);
1818        decoder.register_decoder_with_context::<crate::evm::protocol::uniswap_v4::state::UniswapV4State>(
1819            "uniswap_v4_hooks", context
1820        );
1821        decoder
1822    }
1823
1824    #[tokio::test]
1825    async fn test_unregistered_hook_is_rejected_off_the_generic_vm_chains() {
1826        let decoder = hooked_v4_decoder(Chain::Robinhood, DecoderContext::new());
1827
1828        let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
1829            panic!("a hooked pool on robinhood must not decode");
1830        };
1831
1832        assert!(message.contains("unsupported uniswap v4 hook"), "{message}");
1833        assert!(message.contains("robinhood"), "{message}");
1834    }
1835
1836    #[tokio::test]
1837    async fn test_decoder_chain_overrides_the_registered_context_chain() {
1838        let decoder =
1839            hooked_v4_decoder(Chain::Robinhood, DecoderContext::new().chain(Chain::Ethereum));
1840
1841        let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
1842            panic!("the decoder's chain must win over the context's");
1843        };
1844
1845        assert!(message.contains("robinhood"), "{message}");
1846    }
1847
1848    #[tokio::test]
1849    async fn test_unregistered_hook_takes_the_generic_path_on_ethereum() {
1850        let decoder = hooked_v4_decoder(Chain::Ethereum, DecoderContext::new());
1851
1852        let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
1853            panic!("the generic VM creator has no balance_owner attribute to work from");
1854        };
1855
1856        assert!(!message.contains("unsupported uniswap v4 hook"), "{message}");
1857        assert!(message.contains("balance_owner"), "{message}");
1858    }
1859
1860    #[tokio::test(flavor = "multi_thread")]
1861    async fn test_euler_hook_low_pool_manager_balance() {
1862        let mut decoder = TychoStreamDecoder::new(Chain::Ethereum);
1863
1864        decoder.register_decoder_with_context::<crate::evm::protocol::uniswap_v4::state::UniswapV4State>(
1865            "uniswap_v4_hooks", DecoderContext::new().vm_traces(true)
1866        );
1867
1868        let weth = Bytes::from_str("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").unwrap();
1869        let teth = Bytes::from_str("0xd11c452fc99cf405034ee446803b6f6c1f6d5ed8").unwrap();
1870        let tokens = HashMap::from([
1871            (
1872                weth.clone(),
1873                Token::new(&weth, "WETH", 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1874            ),
1875            (
1876                teth.clone(),
1877                Token::new(&teth, "tETH", 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1878            ),
1879        ]);
1880
1881        decoder.set_tokens(tokens.clone()).await;
1882
1883        let msg = load_test_msg("euler_hook_snapshot");
1884        let res = decoder
1885            .decode(&msg)
1886            .await
1887            .expect("decode failure");
1888
1889        let pool_state = res
1890            .states
1891            .get("0xc70d7fbd7fcccdf726e02fed78548b40dc52502b097c7a1ee7d995f4d4396134")
1892            .expect("Couldn't find target pool");
1893        let amount_out = pool_state
1894            .get_amount_out(
1895                BigUint::from_str("1000000000000000000").unwrap(),
1896                tokens.get(&teth).unwrap(),
1897                tokens.get(&weth).unwrap(),
1898            )
1899            .expect("Get amount out failed");
1900
1901        assert_eq!(amount_out.amount, BigUint::from_str("1216190190361759119").unwrap());
1902    }
1903
1904    /// A single-block feed message carrying `snapshot` as a `uniswap_v4_hooks` component.
1905    fn pons_feed_message(snapshot: ComponentWithState) -> FeedMessage<BlockHeader> {
1906        use tycho_client::feed::synchronizer::{Snapshot, StateSyncMessage};
1907
1908        use crate::evm::protocol::uniswap_v4::pons_fixture;
1909
1910        FeedMessage {
1911            state_msgs: HashMap::from([(
1912                "uniswap_v4_hooks".to_string(),
1913                StateSyncMessage {
1914                    header: pons_fixture::header(),
1915                    snapshots: Snapshot {
1916                        states: HashMap::from([(pons_fixture::POOL_ID.to_string(), snapshot)]),
1917                        vm_storage: HashMap::new(),
1918                    },
1919                    deltas: None,
1920                    removed_components: HashMap::new(),
1921                },
1922            )]),
1923            sync_states: HashMap::new(),
1924        }
1925    }
1926
1927    /// A decoder for Robinhood that knows the fixture pool's two tokens.
1928    async fn pons_stream_decoder() -> TychoStreamDecoder<BlockHeader> {
1929        use crate::evm::protocol::uniswap_v4::{
1930            hooks::hook_handler_creator::initialize_hook_handlers, pons_fixture,
1931            state::UniswapV4State,
1932        };
1933
1934        initialize_hook_handlers().expect("hook handler registration should succeed");
1935        let mut decoder = TychoStreamDecoder::new(Chain::Robinhood);
1936        decoder.register_decoder::<UniswapV4State>("uniswap_v4_hooks");
1937        decoder
1938            .set_tokens(pons_fixture::tokens())
1939            .await;
1940        decoder
1941    }
1942
1943    /// The whole stream path on Robinhood: a Pons pool arrives in a snapshot, is emitted as a
1944    /// `UniswapV4State` carrying the native handler, and quotes the hookless output less the fee
1945    /// and the tax `registerPool` froze for it.
1946    #[tokio::test]
1947    async fn test_pons_pool_decodes_through_the_stream_decoder_on_robinhood() {
1948        use crate::evm::protocol::uniswap_v4::{
1949            hooks::pons_v2::hook_handler::PonsV2HookHandler, pons_fixture, state::UniswapV4State,
1950        };
1951
1952        let decoder = pons_stream_decoder().await;
1953
1954        let result = decoder
1955            .decode(&pons_feed_message(pons_fixture::snapshot()))
1956            .await
1957            .expect("decode failure");
1958
1959        let emitted = result
1960            .states
1961            .get(pons_fixture::POOL_ID)
1962            .expect("the Pons pool must reach consumers");
1963        let pool = emitted
1964            .as_any()
1965            .downcast_ref::<UniswapV4State>()
1966            .expect("a uniswap_v4_hooks component decodes into a UniswapV4State");
1967        let handler = pool
1968            .hook
1969            .as_ref()
1970            .expect("a Pons pool must carry a hook handler");
1971        let pons = handler
1972            .as_any()
1973            .downcast_ref::<PonsV2HookHandler>()
1974            .expect("the registry must hand back the native Pons handler");
1975        assert_eq!(u32::from(pons.hook_fee_bps()), pons_fixture::HOOK_FEE_BPS);
1976        assert_eq!(u32::from(pons.creator_tax_bps()), pons_fixture::CREATOR_TAX_BPS);
1977
1978        let core = pons_fixture::decode_on(pons_fixture::hookless_snapshot(), Chain::Robinhood)
1979            .await
1980            .expect("the same pool with no hook must decode too");
1981        let (xlg, nvda) = (pons_fixture::xlg(), pons_fixture::nvda());
1982        let amount_in = BigUint::from(10u64).pow(18);
1983        for (token_in, token_out) in [(&xlg, &nvda), (&nvda, &xlg)] {
1984            let core_out = core
1985                .get_amount_out(amount_in.clone(), token_in, token_out)
1986                .expect("the reference pool quotes one whole token")
1987                .amount;
1988            let hooked_out = emitted
1989                .get_amount_out(amount_in.clone(), token_in, token_out)
1990                .expect("the emitted pool quotes one whole token")
1991                .amount;
1992
1993            let expected = pons_fixture::net_of_hook_take(&core_out);
1994            assert!(expected < core_out, "the hook must take something out of {core_out}");
1995            assert_eq!(hooked_out, expected, "{} -> {}", token_in.symbol, token_out.symbol);
1996        }
1997    }
1998
1999    /// The same pool behind a hook with no native handler never reaches consumers, even when the
2000    /// stream is configured to carry on past a decode failure.
2001    #[tokio::test]
2002    async fn test_unknown_hook_component_is_not_emitted_on_robinhood() {
2003        use crate::evm::protocol::uniswap_v4::pons_fixture;
2004
2005        let mut decoder = pons_stream_decoder().await;
2006        decoder.skip_state_decode_failures(true);
2007
2008        let result = decoder
2009            .decode(&pons_feed_message(pons_fixture::unknown_hook_snapshot()))
2010            .await
2011            .expect("a skipped decode failure is not fatal");
2012
2013        assert!(
2014            !result
2015                .states
2016                .contains_key(pons_fixture::POOL_ID),
2017            "a pool whose hook cannot be modelled must not be quoted"
2018        );
2019        assert!(result.states.is_empty());
2020    }
2021
2022    fn component_with_id(id: &str) -> ComponentWithState {
2023        use tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState};
2024
2025        ComponentWithState {
2026            state: ProtocolComponentState::new(id, HashMap::new(), HashMap::new()),
2027            component: ProtocolComponent { id: id.to_string(), ..Default::default() },
2028            component_tvl: None,
2029            entrypoints: Vec::new(),
2030        }
2031    }
2032
2033    fn rejects_a(component: &ComponentWithState) -> bool {
2034        component.component.id != "a"
2035    }
2036
2037    fn rejects_b(component: &ComponentWithState) -> bool {
2038        component.component.id != "b"
2039    }
2040
2041    #[test]
2042    fn test_admits_requires_every_registered_filter() {
2043        // Two filters on one exchange both apply: the second registration does not replace the
2044        // first, and a component must pass both. An exchange without filters admits everything.
2045        let mut decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
2046        decoder.register_filter("x", rejects_a);
2047        decoder.register_filter("x", rejects_b);
2048
2049        assert!(!decoder.admits("x", &component_with_id("a")));
2050        assert!(!decoder.admits("x", &component_with_id("b")));
2051        assert!(decoder.admits("x", &component_with_id("c")));
2052        assert!(decoder.admits("y", &component_with_id("a")));
2053    }
2054}