Skip to main content

tycho_simulation/evm/
decoder.rs

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