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        unreachable!()
1411    }
1412
1413    fn typetag_deserialize(&self) {
1414        unreachable!()
1415    }
1416}
1417
1418#[cfg(test)]
1419mod tests {
1420    use std::str::FromStr;
1421
1422    use alloy::primitives::address;
1423    use mockall::predicate::*;
1424    use rstest::*;
1425    use tycho_client::feed::BlockHeader;
1426    use tycho_common::{models::Chain, Bytes};
1427
1428    use super::*;
1429
1430    fn header_at(number: u64, timestamp: u64, partial: Option<u32>) -> BlockHeader {
1431        BlockHeader {
1432            hash: Bytes::from([0u8; 32]),
1433            number,
1434            parent_hash: Bytes::from([0u8; 32]),
1435            revert: false,
1436            timestamp,
1437            partial_block_index: partial,
1438        }
1439    }
1440
1441    /// A block-sensitive state whose execution timestamp we can read back.
1442    fn block_sensitive_state() -> Box<dyn ProtocolSim> {
1443        use crate::evm::protocol::{
1444            aerodrome_slipstreams::state::AerodromeSlipstreamsState,
1445            utils::{
1446                slipstreams::{dynamic_fee_module::DynamicFeeConfig, observations::Observation},
1447                uniswap::{tick_list::TickInfo, tick_math::get_sqrt_ratio_at_tick},
1448            },
1449        };
1450
1451        Box::new(
1452            AerodromeSlipstreamsState::new(
1453                "block-sensitive".to_string(),
1454                0,
1455                1_000_000_000_000_000_000,
1456                get_sqrt_ratio_at_tick(0).unwrap(),
1457                0,
1458                1,
1459                3000,
1460                1,
1461                0,
1462                vec![TickInfo::new(-120, 0).unwrap(), TickInfo::new(120, 0).unwrap()],
1463                vec![Observation { block_timestamp: 500, initialized: true, ..Default::default() }],
1464                DynamicFeeConfig::new(2700, 30_000, 0, true, 750),
1465            )
1466            .expect("state should build")
1467            // These fixtures exercise the fee-flip path, which needs the optimistic mode:
1468            // under the worst-case default a flat-fee pool never flips.
1469            .with_position_assumption(crate::protocol::models::BlockPositionAssumption::First),
1470        )
1471    }
1472
1473    #[test]
1474    fn confirmed_header_targets_the_next_block() {
1475        let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Base);
1476
1477        let execution_block = decoder.execution_block(&header_at(100, 1_000, None));
1478
1479        assert_eq!(execution_block.number(), 101);
1480        assert_eq!(execution_block.timestamp(), 1_000 + Chain::Base.block_time_secs());
1481    }
1482
1483    #[test]
1484    fn partial_header_targets_the_block_that_is_still_open() {
1485        let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
1486
1487        let execution_block = decoder.execution_block(&header_at(100, 1_000, Some(3)));
1488
1489        assert_eq!(execution_block.number(), 100);
1490        assert_eq!(execution_block.timestamp(), 1_000);
1491    }
1492
1493    #[test]
1494    fn refresh_re_emits_a_state_whose_fee_flipped_without_a_delta() {
1495        // The pool traded in block 100 and has no delta afterwards. Crossing into block 101
1496        // flips its fee branch (dynamic -> initial), so the refresh must advance the stored
1497        // state in place and emit a copy to consumers.
1498        let mut stored = HashMap::from([("block-sensitive".to_string(), {
1499            let mut state = block_sensitive_state();
1500            state.apply_block(&BlockContext::new(100, 500));
1501            state
1502        })]);
1503        let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1504
1505        TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1506            &mut updated,
1507            &mut stored,
1508            &HashSet::new(),
1509            &HashMap::<String, ()>::new(),
1510            &BlockContext::new(101, 502),
1511        );
1512
1513        let emitted = updated
1514            .get("block-sensitive")
1515            .expect("a fee flip must be emitted even without a delta");
1516        assert_eq!(emitted.fee(), 750.0 / 1_000_000.0);
1517        // The stored copy was advanced in place as well.
1518        assert_eq!(stored["block-sensitive"].fee(), 750.0 / 1_000_000.0);
1519    }
1520
1521    #[test]
1522    fn refresh_never_re_emits_failed_components() {
1523        // A failed component's stored state may linger; the sweep must not resurrect it for
1524        // consumers that were told the component was removed — even when its fee flipped.
1525        let mut stored = HashMap::from([("zombie".to_string(), {
1526            let mut state = block_sensitive_state();
1527            state.apply_block(&BlockContext::new(100, 500));
1528            state
1529        })]);
1530        let failed = HashSet::from(["zombie".to_string()]);
1531        let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1532
1533        TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1534            &mut updated,
1535            &mut stored,
1536            &failed,
1537            &HashMap::<String, ()>::new(),
1538            &BlockContext::new(101, 502),
1539        );
1540
1541        assert!(updated.is_empty());
1542    }
1543
1544    #[test]
1545    fn refresh_never_re_emits_removed_components() {
1546        // The stored state of a component reported in `removed_pairs` stays in place; the sweep
1547        // must not surface it again, even when its fee flipped in the execution block.
1548        let mut stored = HashMap::from([("gone".to_string(), {
1549            let mut state = block_sensitive_state();
1550            state.apply_block(&BlockContext::new(100, 500));
1551            state
1552        })]);
1553        let removed = HashMap::from([("gone".to_string(), ())]);
1554        let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1555
1556        TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1557            &mut updated,
1558            &mut stored,
1559            &HashSet::new(),
1560            &removed,
1561            &BlockContext::new(101, 502),
1562        );
1563
1564        assert!(updated.is_empty());
1565    }
1566
1567    #[test]
1568    fn refresh_stays_quiet_when_no_fee_changed() {
1569        // An idle block-sensitive pool (fee branch unchanged) and a block-insensitive pool:
1570        // neither must be re-emitted.
1571        let mut stored: HashMap<String, Box<dyn ProtocolSim>> = HashMap::from([
1572            ("idle-sensitive".to_string(), {
1573                let mut state = block_sensitive_state();
1574                state.apply_block(&BlockContext::new(101, 502));
1575                state
1576            }),
1577            (
1578                "univ2".to_string(),
1579                Box::new(crate::evm::protocol::uniswap_v2::state::UniswapV2State::new(
1580                    U256::from(1_000_000u64),
1581                    U256::from(1_000_000u64),
1582                )) as Box<dyn ProtocolSim>,
1583            ),
1584        ]);
1585        let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1586
1587        TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1588            &mut updated,
1589            &mut stored,
1590            &HashSet::new(),
1591            &HashMap::<String, ()>::new(),
1592            &BlockContext::new(102, 504),
1593        );
1594
1595        assert!(updated.is_empty());
1596    }
1597    use crate::evm::protocol::{curve::CurveState, uniswap_v2::state::UniswapV2State};
1598
1599    #[test]
1600    fn curve_vm_adapter_registration_flagged_deprecated() {
1601        // The native decoder is the supported path — not flagged.
1602        assert!(!is_deprecated_curve_registration::<CurveState>("vm:curve"));
1603        // Any other type for vm:curve is the deprecated VM-adapter path.
1604        assert!(is_deprecated_curve_registration::<UniswapV2State>("vm:curve"));
1605        // Other exchanges are unaffected.
1606        assert!(!is_deprecated_curve_registration::<UniswapV2State>("uniswap_v2"));
1607    }
1608
1609    async fn setup_decoder(set_tokens: bool) -> TychoStreamDecoder<BlockHeader> {
1610        let mut decoder = TychoStreamDecoder::new(Chain::Ethereum);
1611        decoder.register_decoder::<UniswapV2State>("uniswap_v2");
1612        if set_tokens {
1613            let tokens = [
1614                Bytes::from("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").lpad(20, 0),
1615                Bytes::from("0xdac17f958d2ee523a2206206994597c13d831ec7").lpad(20, 0),
1616            ]
1617            .iter()
1618            .map(|addr| {
1619                let addr_str = format!("{addr:x}");
1620                (
1621                    addr.clone(),
1622                    Token::new(addr, &addr_str, 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1623                )
1624            })
1625            .collect();
1626            decoder.set_tokens(tokens).await;
1627        }
1628        decoder
1629    }
1630
1631    fn load_test_msg(name: &str) -> FeedMessage<BlockHeader> {
1632        use std::{fs, path::Path};
1633
1634        use tycho_client::feed::dto;
1635        let project_root = env!("CARGO_MANIFEST_DIR");
1636        let asset_path = Path::new(project_root).join(format!("tests/assets/decoder/{name}.json"));
1637        let json_data = fs::read_to_string(asset_path).expect("Failed to read test asset");
1638        let feed_msg: dto::FeedMessage<BlockHeader> =
1639            serde_json::from_str(&json_data).expect("Failed to deserialize FeedMsg json!");
1640        FeedMessage::from(feed_msg)
1641    }
1642
1643    #[tokio::test]
1644    async fn test_decode() {
1645        let decoder = setup_decoder(true).await;
1646
1647        let msg = load_test_msg("uniswap_v2_snapshot");
1648        let res1 = decoder
1649            .decode(&msg)
1650            .await
1651            .expect("decode failure");
1652        let msg = load_test_msg("uniswap_v2_delta");
1653        let res2 = decoder
1654            .decode(&msg)
1655            .await
1656            .expect("decode failure");
1657
1658        assert_eq!(res1.states.len(), 1);
1659        assert_eq!(res2.states.len(), 1);
1660        assert_eq!(res1.sync_states.len(), 1);
1661        assert_eq!(res2.sync_states.len(), 1);
1662    }
1663
1664    #[tokio::test]
1665    async fn test_decode_token_creation_delta_with_existing_proxy() {
1666        let decoder = setup_decoder(true).await;
1667        let msg = load_test_msg("uniswap_v2_delta_token_creation");
1668
1669        // First decode: the token has no proxy yet, so the Creation delta takes the
1670        // proxy-creating branch.
1671        decoder
1672            .decode(&msg)
1673            .await
1674            .expect("first decode (proxy creation) failed");
1675
1676        // Second decode: the proxy exists, so the same Creation delta must decode as a
1677        // storage update on the proxy account — not a code-less Creation, which the
1678        // engine rejects as "MissingCode".
1679        decoder
1680            .decode(&msg)
1681            .await
1682            .expect("decode of a token Creation delta with an existing proxy failed");
1683    }
1684
1685    #[tokio::test]
1686    async fn test_decode_component_missing_token() {
1687        let decoder = setup_decoder(false).await;
1688        let tokens = [Bytes::from("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").lpad(20, 0)]
1689            .iter()
1690            .map(|addr| {
1691                let addr_str = format!("{addr:x}");
1692                (
1693                    addr.clone(),
1694                    Token::new(addr, &addr_str, 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1695                )
1696            })
1697            .collect();
1698        decoder.set_tokens(tokens).await;
1699
1700        let msg = load_test_msg("uniswap_v2_snapshot");
1701        let res1 = decoder
1702            .decode(&msg)
1703            .await
1704            .expect("decode failure");
1705
1706        assert_eq!(res1.states.len(), 0);
1707    }
1708
1709    #[tokio::test]
1710    async fn test_decode_component_bad_id() {
1711        let decoder = setup_decoder(true).await;
1712        let msg = load_test_msg("uniswap_v2_snapshot_broken_id");
1713
1714        match decoder.decode(&msg).await {
1715            Err(StreamDecodeError::Fatal(msg)) => {
1716                assert_eq!(msg, "Component id mismatch");
1717            }
1718            Ok(_) => {
1719                panic!("Expected failures to be raised")
1720            }
1721        }
1722    }
1723
1724    #[rstest]
1725    #[case(true)]
1726    #[case(false)]
1727    #[tokio::test]
1728    async fn test_decode_component_bad_state(#[case] skip_failures: bool) {
1729        let mut decoder = setup_decoder(true).await;
1730        decoder.skip_state_decode_failures = skip_failures;
1731
1732        let msg = load_test_msg("uniswap_v2_snapshot_broken_state");
1733        match decoder.decode(&msg).await {
1734            Err(StreamDecodeError::Fatal(msg)) => {
1735                if !skip_failures {
1736                    assert_eq!(msg, "Missing attributes reserve0");
1737                } else {
1738                    panic!("Expected failures to be ignored. Err: {msg}")
1739                }
1740            }
1741            Ok(res) => {
1742                if !skip_failures {
1743                    panic!("Expected failures to be raised")
1744                } else {
1745                    assert_eq!(res.states.len(), 0);
1746                }
1747            }
1748        }
1749    }
1750
1751    #[tokio::test]
1752    async fn test_decode_updates_state_on_contract_change() {
1753        let decoder = setup_decoder(true).await;
1754
1755        // Create the mock instances
1756        let mut mock_state = MockProtocolSim::new();
1757
1758        mock_state
1759            .expect_clone_box()
1760            .times(1)
1761            .returning(|| {
1762                let mut cloned_mock_state = MockProtocolSim::new();
1763                // Expect `delta_transition` to be called once with any parameters
1764                cloned_mock_state
1765                    .expect_delta_transition()
1766                    .times(1)
1767                    .returning(|_, _, _| Ok(()));
1768                cloned_mock_state
1769                    .expect_clone_box()
1770                    .times(1)
1771                    .returning(|| Box::new(MockProtocolSim::new()));
1772                Box::new(cloned_mock_state)
1773            });
1774
1775        // Insert mock state into `updated_states`
1776        let pool_id =
1777            "0x93d199263632a4ef4bb438f1feb99e57b4b5f0bd0000000000000000000005c2".to_string();
1778        decoder
1779            .state
1780            .write()
1781            .await
1782            .states
1783            .insert(pool_id.clone(), Box::new(mock_state) as Box<dyn ProtocolSim>);
1784        decoder
1785            .state
1786            .write()
1787            .await
1788            .contracts_map
1789            .insert(
1790                Bytes::from("0xba12222222228d8ba445958a75a0704d566bf2c8").lpad(20, 0),
1791                HashSet::from([pool_id.clone()]),
1792            );
1793
1794        // Load a test message containing a contract update
1795        let msg = load_test_msg("balancer_v2_delta");
1796
1797        // Decode the message
1798        let _ = decoder
1799            .decode(&msg)
1800            .await
1801            .expect("decode failure");
1802
1803        // The mock framework will assert that `delta_transition` was called exactly once
1804    }
1805
1806    /// The pending block's account deltas must reach the clone of every pool linked to a written
1807    /// account, never a pool the write cannot affect, and never the stored confirmed state.
1808    #[tokio::test]
1809    async fn test_apply_deltas_ephemeral_sets_pending_overrides_on_the_clone_only() {
1810        use tycho_common::models::{
1811            blockchain::BlockAggregatedChanges, contract::AccountDelta,
1812            protocol::ProtocolComponentStateDelta, ChangeType as ModelChangeType,
1813        };
1814
1815        use crate::evm::protocol::uniswap_v4::state::{UniswapV4Fees, UniswapV4State};
1816
1817        let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
1818        let pool_id = "0xhooked".to_string();
1819        let quiet_pool_id = "0xhooked-no-log".to_string();
1820        let unlinked_pool_id = "0xunlinked".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.clone()));
1839            state
1840                .states
1841                .insert(unlinked_pool_id.clone(), Box::new(pool));
1842            state
1843                .contracts_map
1844                .insert(hook.clone(), HashSet::from([pool_id.clone(), quiet_pool_id.clone()]));
1845        }
1846
1847        let deltas = BlockAggregatedChanges {
1848            extractor: "uniswap_v4_hooks".to_string(),
1849            state_deltas: [&pool_id, &unlinked_pool_id]
1850                .into_iter()
1851                .map(|id| {
1852                    let delta = ProtocolComponentStateDelta {
1853                        component_id: id.clone(),
1854                        updated_attributes: HashMap::from([(
1855                            "liquidity".to_string(),
1856                            Bytes::from(2000_u64.to_be_bytes().to_vec()),
1857                        )]),
1858                        deleted_attributes: HashSet::new(),
1859                        created_attributes: HashSet::new(),
1860                    };
1861                    (id.clone(), delta)
1862                })
1863                .collect(),
1864            account_deltas: HashMap::from([(
1865                hook.clone(),
1866                AccountDelta::new(
1867                    Chain::Ethereum,
1868                    hook.clone(),
1869                    HashMap::from([(Bytes::from(vec![0u8]), Some(Bytes::from(vec![9u8])))]),
1870                    None,
1871                    None,
1872                    ModelChangeType::Update,
1873                ),
1874            )]),
1875            ..Default::default()
1876        };
1877        let header = BlockHeader { number: 10, timestamp: 20, ..Default::default() };
1878
1879        let update = decoder
1880            .apply_deltas_ephemeral(
1881                &HashMap::from([("uniswap_v4_hooks".to_string(), deltas)]),
1882                header,
1883            )
1884            .await
1885            .unwrap();
1886
1887        let clone = update.states[&pool_id]
1888            .as_any()
1889            .downcast_ref::<UniswapV4State>()
1890            .unwrap();
1891        let pending = clone
1892            .pending_overrides()
1893            .expect("the clone carries the pending block's overrides");
1894        assert_eq!(
1895            pending.block,
1896            Some(BlockEnvOverrides { number: Some(10), timestamp: Some(20) })
1897        );
1898        assert_eq!(
1899            pending.storage.as_ref().unwrap()[&Address::from_slice(&hook)][&U256::ZERO],
1900            U256::from(9)
1901        );
1902        let quiet = update.states[&quiet_pool_id]
1903            .as_any()
1904            .downcast_ref::<UniswapV4State>()
1905            .expect("a pool whose hook was written is cloned without a delta of its own");
1906        assert!(quiet.pending_overrides().is_some());
1907        let unlinked = update.states[&unlinked_pool_id]
1908            .as_any()
1909            .downcast_ref::<UniswapV4State>()
1910            .unwrap();
1911        assert!(
1912            unlinked.pending_overrides().is_none(),
1913            "a pool linked to no written account gets no overrides"
1914        );
1915
1916        let stored = decoder.state.read().await;
1917        for id in [&pool_id, &quiet_pool_id, &unlinked_pool_id] {
1918            let stored = stored.states[id]
1919                .as_any()
1920                .downcast_ref::<UniswapV4State>()
1921                .unwrap();
1922            assert!(stored.pending_overrides().is_none(), "confirmed state is never overlaid");
1923        }
1924    }
1925
1926    #[test]
1927    fn test_generate_proxy_token_address() {
1928        let idx = 1;
1929        let generated_address =
1930            generate_proxy_token_address(idx).expect("proxy token address should be valid");
1931        assert_eq!(generated_address, address!("000000000000000000000000000000001badbabe"));
1932
1933        let idx = 123456;
1934        let generated_address =
1935            generate_proxy_token_address(idx).expect("proxy token address should be valid");
1936        assert_eq!(generated_address, address!("00000000000000000000000000001e240badbabe"));
1937    }
1938
1939    /// A single-block feed message holding one tokenless Uniswap V4 pool whose `hooks` attribute
1940    /// points at an address no native handler is registered for.
1941    ///
1942    /// The component has no tokens and the message no VM storage, so decoding it neither reads
1943    /// nor writes the process-wide simulation engine and cannot disturb other tests.
1944    fn hooked_v4_msg() -> FeedMessage<BlockHeader> {
1945        use tycho_client::feed::synchronizer::{Snapshot, StateSyncMessage};
1946        use tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState};
1947
1948        let id = "0x00000000000000000000000000000000000000000000000000000000000000c4";
1949        let static_attributes = HashMap::from([
1950            ("key_lp_fee".to_string(), Bytes::from(500_i32.to_be_bytes().to_vec())),
1951            ("tick_spacing".to_string(), Bytes::from(60_i32.to_be_bytes().to_vec())),
1952            (
1953                "hooks".to_string(),
1954                Bytes::from_str("0x00000000000000000000000000000000000000c4").unwrap(),
1955            ),
1956        ]);
1957        let attributes = HashMap::from([
1958            ("liquidity".to_string(), Bytes::from(100_u64.to_be_bytes().to_vec())),
1959            ("tick".to_string(), Bytes::from(300_i32.to_be_bytes().to_vec())),
1960            (
1961                "sqrt_price_x96".to_string(),
1962                Bytes::from(
1963                    79228162514264337593543950336_u128
1964                        .to_be_bytes()
1965                        .to_vec(),
1966                ),
1967            ),
1968            ("protocol_fees/zero2one".to_string(), Bytes::from(0_u32.to_be_bytes().to_vec())),
1969            ("protocol_fees/one2zero".to_string(), Bytes::from(0_u32.to_be_bytes().to_vec())),
1970            ("ticks/60/net_liquidity".to_string(), Bytes::from(400_i128.to_be_bytes().to_vec())),
1971        ]);
1972
1973        let snapshot = ComponentWithState {
1974            state: ProtocolComponentState::new(id, attributes, HashMap::new()),
1975            component: ProtocolComponent {
1976                id: id.to_string(),
1977                static_attributes,
1978                ..Default::default()
1979            },
1980            component_tvl: None,
1981            entrypoints: Vec::new(),
1982        };
1983
1984        FeedMessage {
1985            state_msgs: HashMap::from([(
1986                "uniswap_v4_hooks".to_string(),
1987                StateSyncMessage {
1988                    header: header_at(1, 1_000, None),
1989                    snapshots: Snapshot {
1990                        states: HashMap::from([(id.to_string(), snapshot)]),
1991                        vm_storage: HashMap::new(),
1992                    },
1993                    deltas: None,
1994                    removed_components: HashMap::new(),
1995                },
1996            )]),
1997            sync_states: HashMap::new(),
1998        }
1999    }
2000
2001    /// A decoder for `chain` that decodes [`hooked_v4_msg`].
2002    fn hooked_v4_decoder(chain: Chain, context: DecoderContext) -> TychoStreamDecoder<BlockHeader> {
2003        let mut decoder = TychoStreamDecoder::new(chain);
2004        decoder.register_decoder_with_context::<crate::evm::protocol::uniswap_v4::state::UniswapV4State>(
2005            "uniswap_v4_hooks", context
2006        );
2007        decoder
2008    }
2009
2010    #[tokio::test]
2011    async fn test_unregistered_hook_is_rejected_off_the_generic_vm_chains() {
2012        let decoder = hooked_v4_decoder(Chain::Robinhood, DecoderContext::new());
2013
2014        let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
2015            panic!("a hooked pool on robinhood must not decode");
2016        };
2017
2018        assert!(message.contains("unsupported uniswap v4 hook"), "{message}");
2019        assert!(message.contains("robinhood"), "{message}");
2020    }
2021
2022    #[tokio::test]
2023    async fn test_decoder_chain_overrides_the_registered_context_chain() {
2024        let decoder =
2025            hooked_v4_decoder(Chain::Robinhood, DecoderContext::new().chain(Chain::Ethereum));
2026
2027        let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
2028            panic!("the decoder's chain must win over the context's");
2029        };
2030
2031        assert!(message.contains("robinhood"), "{message}");
2032    }
2033
2034    #[tokio::test]
2035    async fn test_unregistered_hook_takes_the_generic_path_on_ethereum() {
2036        let decoder = hooked_v4_decoder(Chain::Ethereum, DecoderContext::new());
2037
2038        let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
2039            panic!("the generic VM creator has no balance_owner attribute to work from");
2040        };
2041
2042        assert!(!message.contains("unsupported uniswap v4 hook"), "{message}");
2043        assert!(message.contains("balance_owner"), "{message}");
2044    }
2045
2046    #[tokio::test(flavor = "multi_thread")]
2047    async fn test_euler_hook_low_pool_manager_balance() {
2048        let mut decoder = TychoStreamDecoder::new(Chain::Ethereum);
2049
2050        decoder.register_decoder_with_context::<crate::evm::protocol::uniswap_v4::state::UniswapV4State>(
2051            "uniswap_v4_hooks", DecoderContext::new().vm_traces(true)
2052        );
2053
2054        let weth = Bytes::from_str("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").unwrap();
2055        let teth = Bytes::from_str("0xd11c452fc99cf405034ee446803b6f6c1f6d5ed8").unwrap();
2056        let tokens = HashMap::from([
2057            (
2058                weth.clone(),
2059                Token::new(&weth, "WETH", 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
2060            ),
2061            (
2062                teth.clone(),
2063                Token::new(&teth, "tETH", 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
2064            ),
2065        ]);
2066
2067        decoder.set_tokens(tokens.clone()).await;
2068
2069        let msg = load_test_msg("euler_hook_snapshot");
2070        let res = decoder
2071            .decode(&msg)
2072            .await
2073            .expect("decode failure");
2074
2075        let pool_state = res
2076            .states
2077            .get("0xc70d7fbd7fcccdf726e02fed78548b40dc52502b097c7a1ee7d995f4d4396134")
2078            .expect("Couldn't find target pool");
2079        let amount_out = pool_state
2080            .get_amount_out(
2081                BigUint::from_str("1000000000000000000").unwrap(),
2082                tokens.get(&teth).unwrap(),
2083                tokens.get(&weth).unwrap(),
2084            )
2085            .expect("Get amount out failed");
2086
2087        assert_eq!(amount_out.amount, BigUint::from_str("1216190190361759119").unwrap());
2088    }
2089
2090    /// A single-block feed message carrying `snapshot` as a `uniswap_v4_hooks` component.
2091    fn pons_feed_message(snapshot: ComponentWithState) -> FeedMessage<BlockHeader> {
2092        use tycho_client::feed::synchronizer::{Snapshot, StateSyncMessage};
2093
2094        use crate::evm::protocol::uniswap_v4::pons_fixture;
2095
2096        FeedMessage {
2097            state_msgs: HashMap::from([(
2098                "uniswap_v4_hooks".to_string(),
2099                StateSyncMessage {
2100                    header: pons_fixture::header(),
2101                    snapshots: Snapshot {
2102                        states: HashMap::from([(pons_fixture::POOL_ID.to_string(), snapshot)]),
2103                        vm_storage: HashMap::new(),
2104                    },
2105                    deltas: None,
2106                    removed_components: HashMap::new(),
2107                },
2108            )]),
2109            sync_states: HashMap::new(),
2110        }
2111    }
2112
2113    /// A decoder for Robinhood that knows the fixture pool's two tokens.
2114    async fn pons_stream_decoder() -> TychoStreamDecoder<BlockHeader> {
2115        use crate::evm::protocol::uniswap_v4::{
2116            hooks::hook_handler_creator::initialize_hook_handlers, pons_fixture,
2117            state::UniswapV4State,
2118        };
2119
2120        initialize_hook_handlers().expect("hook handler registration should succeed");
2121        let mut decoder = TychoStreamDecoder::new(Chain::Robinhood);
2122        decoder.register_decoder::<UniswapV4State>("uniswap_v4_hooks");
2123        decoder
2124            .set_tokens(pons_fixture::tokens())
2125            .await;
2126        decoder
2127    }
2128
2129    /// The whole stream path on Robinhood: a Pons pool arrives in a snapshot, is emitted as a
2130    /// `UniswapV4State` carrying the native handler, and quotes the hookless output less the fee
2131    /// and the tax `registerPool` froze for it.
2132    #[tokio::test]
2133    async fn test_pons_pool_decodes_through_the_stream_decoder_on_robinhood() {
2134        use crate::evm::protocol::uniswap_v4::{
2135            hooks::pons_v2::hook_handler::PonsV2HookHandler, pons_fixture, state::UniswapV4State,
2136        };
2137
2138        let decoder = pons_stream_decoder().await;
2139
2140        let result = decoder
2141            .decode(&pons_feed_message(pons_fixture::snapshot()))
2142            .await
2143            .expect("decode failure");
2144
2145        let emitted = result
2146            .states
2147            .get(pons_fixture::POOL_ID)
2148            .expect("the Pons pool must reach consumers");
2149        let pool = emitted
2150            .as_any()
2151            .downcast_ref::<UniswapV4State>()
2152            .expect("a uniswap_v4_hooks component decodes into a UniswapV4State");
2153        let handler = pool
2154            .hook
2155            .as_ref()
2156            .expect("a Pons pool must carry a hook handler");
2157        let pons = handler
2158            .as_any()
2159            .downcast_ref::<PonsV2HookHandler>()
2160            .expect("the registry must hand back the native Pons handler");
2161        assert_eq!(u32::from(pons.hook_fee_bps()), pons_fixture::HOOK_FEE_BPS);
2162        assert_eq!(u32::from(pons.creator_tax_bps()), pons_fixture::CREATOR_TAX_BPS);
2163
2164        let core = pons_fixture::decode_on(pons_fixture::hookless_snapshot(), Chain::Robinhood)
2165            .await
2166            .expect("the same pool with no hook must decode too");
2167        let (xlg, nvda) = (pons_fixture::xlg(), pons_fixture::nvda());
2168        let amount_in = BigUint::from(10u64).pow(18);
2169        for (token_in, token_out) in [(&xlg, &nvda), (&nvda, &xlg)] {
2170            let core_out = core
2171                .get_amount_out(amount_in.clone(), token_in, token_out)
2172                .expect("the reference pool quotes one whole token")
2173                .amount;
2174            let hooked_out = emitted
2175                .get_amount_out(amount_in.clone(), token_in, token_out)
2176                .expect("the emitted pool quotes one whole token")
2177                .amount;
2178
2179            let expected = pons_fixture::net_of_hook_take(&core_out);
2180            assert!(expected < core_out, "the hook must take something out of {core_out}");
2181            assert_eq!(hooked_out, expected, "{} -> {}", token_in.symbol, token_out.symbol);
2182        }
2183    }
2184
2185    /// The same pool behind a hook with no native handler never reaches consumers, even when the
2186    /// stream is configured to carry on past a decode failure.
2187    #[tokio::test]
2188    async fn test_unknown_hook_component_is_not_emitted_on_robinhood() {
2189        use crate::evm::protocol::uniswap_v4::pons_fixture;
2190
2191        let mut decoder = pons_stream_decoder().await;
2192        decoder.skip_state_decode_failures(true);
2193
2194        let result = decoder
2195            .decode(&pons_feed_message(pons_fixture::unknown_hook_snapshot()))
2196            .await
2197            .expect("a skipped decode failure is not fatal");
2198
2199        assert!(
2200            !result
2201                .states
2202                .contains_key(pons_fixture::POOL_ID),
2203            "a pool whose hook cannot be modelled must not be quoted"
2204        );
2205        assert!(result.states.is_empty());
2206    }
2207
2208    fn component_with_id(id: &str) -> ComponentWithState {
2209        use tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState};
2210
2211        ComponentWithState {
2212            state: ProtocolComponentState::new(id, HashMap::new(), HashMap::new()),
2213            component: ProtocolComponent { id: id.to_string(), ..Default::default() },
2214            component_tvl: None,
2215            entrypoints: Vec::new(),
2216        }
2217    }
2218
2219    fn rejects_a(component: &ComponentWithState) -> bool {
2220        component.component.id != "a"
2221    }
2222
2223    fn rejects_b(component: &ComponentWithState) -> bool {
2224        component.component.id != "b"
2225    }
2226
2227    #[test]
2228    fn test_admits_requires_every_registered_filter() {
2229        // Two filters on one exchange both apply: the second registration does not replace the
2230        // first, and a component must pass both. An exchange without filters admits everything.
2231        let mut decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
2232        decoder.register_filter("x", rejects_a);
2233        decoder.register_filter("x", rejects_b);
2234
2235        assert!(!decoder.admits("x", &component_with_id("a")));
2236        assert!(!decoder.admits("x", &component_with_id("b")));
2237        assert!(decoder.admits("x", &component_with_id("c")));
2238        assert!(decoder.admits("y", &component_with_id("a")));
2239    }
2240}