1use std::{
2 collections::{hash_map::Entry, HashMap, HashSet},
3 future::Future,
4 pin::Pin,
5 sync::Arc,
6};
7
8use alloy::primitives::{Address, U256};
9use thiserror::Error;
10use tokio::sync::{watch, RwLock, RwLockReadGuard};
11use tracing::{debug, error, info, warn};
12use tycho_client::feed::{synchronizer::ComponentWithState, BlockHeader, FeedMessage, HeaderLike};
13use tycho_common::{
14 dto::{ChangeType, ProtocolStateDelta},
15 models::{blockchain::BlockAggregatedChanges, token::Token, Chain},
16 simulation::protocol_sim::{Balances, BlockContext, ProtocolSim},
17 Bytes,
18};
19#[cfg(test)]
20use {
21 mockall::mock,
22 num_bigint::BigUint,
23 std::any::Any,
24 tycho_common::simulation::{
25 errors::{SimulationError, TransitionError},
26 protocol_sim::GetAmountOutResult,
27 },
28};
29
30use crate::{
31 evm::{
32 engine_db::{update_engine, SHARED_TYCHO_DB},
33 override_stream::{OverrideSnapshot, StateOverrideProvider},
34 protocol::{
35 utils::bytes_to_address,
36 vm::{constants::ERC20_PROXY_BYTECODE, erc20_token::IMPLEMENTATION_SLOT},
37 },
38 tycho_models::{AccountUpdate, ResponseAccount},
39 },
40 protocol::{
41 errors::InvalidSnapshotError,
42 models::{DecoderContext, ProtocolComponent, TryFromWithBlock, Update},
43 },
44};
45
46#[derive(Error, Debug)]
47pub enum StreamDecodeError {
48 #[error("{0}")]
49 Fatal(String),
50}
51
52#[derive(Default)]
53struct DecoderState {
54 tokens: HashMap<Bytes, Token>,
55 states: HashMap<String, Box<dyn ProtocolSim>>,
56 components: HashMap<String, ProtocolComponent>,
57 contracts_map: HashMap<Bytes, HashSet<String>>,
59 proxy_token_addresses: HashMap<Address, Address>,
61 failed_components: HashSet<String>,
65 current_block_number: u64,
67}
68
69type DecodeFut =
70 Pin<Box<dyn Future<Output = Result<Box<dyn ProtocolSim>, InvalidSnapshotError>> + Send + Sync>>;
71type AccountBalances = HashMap<Bytes, HashMap<Bytes, Bytes>>;
72type RegistryFn<H> = dyn Fn(
73 ComponentWithState,
74 H,
75 AccountBalances,
76 Arc<RwLock<DecoderState>>,
77 Option<watch::Receiver<OverrideSnapshot>>,
78 ) -> DecodeFut
79 + Send
80 + Sync;
81type FilterFn = fn(&ComponentWithState) -> bool;
82
83pub struct TychoStreamDecoder<H>
96where
97 H: HeaderLike,
98{
99 state: Arc<RwLock<DecoderState>>,
100 skip_state_decode_failures: bool,
101 min_token_quality: u32,
102 registry: HashMap<String, Box<RegistryFn<H>>>,
103 inclusion_filters: HashMap<String, Vec<FilterFn>>,
104 override_providers: HashMap<String, Arc<dyn StateOverrideProvider>>,
107 block_time_secs: u64,
110 chain: Chain,
113}
114
115fn is_deprecated_curve_registration<T: 'static>(exchange: &str) -> bool {
119 exchange == "vm:curve" &&
120 std::any::type_name::<T>() !=
121 std::any::type_name::<crate::evm::protocol::curve::CurveState>()
122}
123
124impl<H> TychoStreamDecoder<H>
125where
126 H: HeaderLike + Clone + Sync + Send + 'static + std::fmt::Debug,
127{
128 pub fn new(chain: Chain) -> Self {
134 Self {
135 state: Arc::new(RwLock::new(DecoderState::default())),
136 skip_state_decode_failures: false,
137 min_token_quality: 100,
138 registry: HashMap::new(),
139 inclusion_filters: HashMap::new(),
140 override_providers: HashMap::new(),
141 block_time_secs: chain.block_time_secs(),
142 chain,
143 }
144 }
145
146 fn execution_block(&self, header: &BlockHeader) -> BlockContext {
151 if header.partial_block_index.is_some() {
152 BlockContext::new(header.number, header.timestamp)
153 } else {
154 BlockContext::new(header.number + 1, header.timestamp + self.block_time_secs)
155 }
156 }
157
158 fn refresh_execution_block<C>(
167 updated_states: &mut HashMap<String, Box<dyn ProtocolSim>>,
168 stored_states: &mut HashMap<String, Box<dyn ProtocolSim>>,
169 failed_components: &HashSet<String>,
170 removed_components: &HashMap<String, C>,
171 execution_block: &BlockContext,
172 ) {
173 for state in updated_states.values_mut() {
174 state.apply_block(execution_block);
175 }
176 for (id, state) in stored_states.iter_mut() {
177 if failed_components.contains(id) ||
178 removed_components.contains_key(id) ||
179 updated_states.contains_key(id)
180 {
181 continue;
182 }
183 if state.apply_block(execution_block) {
184 updated_states.insert(id.clone(), state.clone_box());
185 }
186 }
187 }
188
189 pub fn set_override_provider(
195 &mut self,
196 protocol_system: String,
197 provider: Arc<dyn StateOverrideProvider>,
198 ) {
199 self.override_providers
200 .insert(protocol_system, provider);
201 }
202
203 pub async fn set_tokens(&self, tokens: HashMap<Bytes, Token>) {
208 let mut guard = self.state.write().await;
209 guard.tokens = tokens;
210 }
211
212 pub fn skip_state_decode_failures(&mut self, skip: bool) {
213 self.skip_state_decode_failures = skip;
214 }
215
216 pub fn min_token_quality(&mut self, quality: u32) {
222 self.min_token_quality = quality;
223 }
224
225 pub fn register_decoder_with_context<T>(&mut self, exchange: &str, mut context: DecoderContext)
240 where
241 T: ProtocolSim
242 + TryFromWithBlock<ComponentWithState, H, Error = InvalidSnapshotError>
243 + Send
244 + 'static,
245 {
246 if let Some(requested) = context
247 .chain
248 .filter(|requested| *requested != self.chain)
249 {
250 warn!(
251 exchange,
252 requested_chain = %requested,
253 decoder_chain = %self.chain,
254 "DecoderContext declares a different chain than the decoder; using the decoder's \
255 chain"
256 );
257 }
258 context.chain = Some(self.chain);
259
260 if is_deprecated_curve_registration::<T>(exchange) {
261 warn!(
262 registered_type = std::any::type_name::<T>(),
263 "Registering \"vm:curve\" with the generic VM adapter is deprecated; register the \
264 native `CurveState` decoder instead (`exchange::<CurveState>(\"vm:curve\", ...)`). \
265 The VM-adapter path still works but will be removed in a future release."
266 );
267 }
268 let decoder = Box::new(
269 move |component: ComponentWithState,
270 header: H,
271 account_balances: AccountBalances,
272 state: Arc<RwLock<DecoderState>>,
273 live_override: Option<watch::Receiver<OverrideSnapshot>>| {
274 let mut context = context.clone();
275 context.live_override = live_override;
276 Box::pin(async move {
277 let guard = state.read().await;
278 T::try_from_with_header(
279 component,
280 header,
281 &account_balances,
282 &guard.tokens,
283 &context,
284 )
285 .await
286 .map(|c| Box::new(c) as Box<dyn ProtocolSim>)
287 }) as DecodeFut
288 },
289 );
290 self.registry
291 .insert(exchange.to_string(), decoder);
292 }
293
294 pub fn register_decoder<T>(&mut self, exchange: &str)
306 where
307 T: ProtocolSim
308 + TryFromWithBlock<ComponentWithState, H, Error = InvalidSnapshotError>
309 + Send
310 + 'static,
311 {
312 let context = DecoderContext::new();
313 self.register_decoder_with_context::<T>(exchange, context);
314 }
315
316 pub fn register_filter(&mut self, exchange: &str, predicate: FilterFn) {
335 self.inclusion_filters
336 .entry(exchange.to_string())
337 .or_default()
338 .push(predicate);
339 }
340
341 fn admits(&self, exchange: &str, snapshot: &ComponentWithState) -> bool {
344 let Some(predicates) = self.inclusion_filters.get(exchange) else { return true };
345 predicates
346 .iter()
347 .all(|predicate| predicate(snapshot))
348 }
349
350 pub async fn decode(&self, msg: &FeedMessage<H>) -> Result<Update, StreamDecodeError> {
353 let mut updated_states = HashMap::new();
355 let mut new_pairs = HashMap::new();
356 let mut removed_pairs = HashMap::new();
357 let mut contracts_map = HashMap::new();
358 let mut msg_failed_components = HashSet::new();
359
360 let header = msg
361 .state_msgs
362 .values()
363 .next()
364 .ok_or_else(|| StreamDecodeError::Fatal("Missing block!".into()))?
365 .header
366 .clone();
367
368 let block_number_or_timestamp = header
369 .clone()
370 .block_number_or_timestamp();
371 let current_block = header.clone().block();
372 let is_partial = current_block
373 .as_ref()
374 .map(|h| h.partial_block_index.is_some())
375 .unwrap_or(false);
376
377 for (protocol, protocol_msg) in msg.state_msgs.iter() {
378 if let Some(deltas) = protocol_msg.deltas.as_ref() {
380 let mut state_guard = self.state.write().await;
381
382 let new_tokens = deltas
383 .new_tokens
384 .iter()
385 .filter(|(addr, t)| {
386 t.quality >= self.min_token_quality &&
387 !state_guard.tokens.contains_key(*addr)
388 })
389 .map(|(addr, t)| (addr.clone(), t.clone()))
390 .collect::<HashMap<Bytes, Token>>();
391
392 if !new_tokens.is_empty() {
393 debug!(n = new_tokens.len(), "NewTokens");
394 state_guard.tokens.extend(new_tokens);
395 }
396 }
397
398 {
400 let mut state_guard = self.state.write().await;
401 let removed_components: Vec<(String, ProtocolComponent)> = protocol_msg
402 .removed_components
403 .iter()
404 .map(|(id, comp)| {
405 if *id != comp.id {
406 error!(
407 "Component id mismatch in removed components {id} != {}",
408 comp.id
409 );
410 return Err(StreamDecodeError::Fatal("Component id mismatch".into()));
411 }
412
413 let tokens = comp
414 .tokens
415 .iter()
416 .flat_map(|addr| state_guard.tokens.get(addr).cloned())
417 .collect::<Vec<_>>();
418
419 if tokens.len() == comp.tokens.len() {
420 Ok(Some((
421 id.clone(),
422 ProtocolComponent::from_with_tokens(comp.clone(), tokens),
423 )))
424 } else {
425 Ok(None)
426 }
427 })
428 .collect::<Result<Vec<Option<(String, ProtocolComponent)>>, StreamDecodeError>>(
429 )?
430 .into_iter()
431 .flatten()
432 .collect();
433
434 for (id, component) in removed_components {
436 state_guard.components.remove(&id);
437 state_guard.states.remove(&id);
438 removed_pairs.insert(id, component);
439 }
440
441 info!(
443 "Processing {} contracts from snapshots",
444 protocol_msg
445 .snapshots
446 .get_vm_storage()
447 .len()
448 );
449
450 let mut proxy_token_accounts: HashMap<Address, AccountUpdate> = HashMap::new();
451 let mut storage_by_address: HashMap<Address, ResponseAccount> = HashMap::new();
452 for (key, value) in protocol_msg
453 .snapshots
454 .get_vm_storage()
455 .iter()
456 {
457 let account: ResponseAccount = value.clone().into();
458
459 if state_guard.tokens.contains_key(key) {
460 let original_address = account.address;
461 let (impl_addr, proxy_state) = match state_guard
470 .proxy_token_addresses
471 .get(&original_address)
472 {
473 Some(impl_addr) => {
474 let proxy_state = AccountUpdate::new(
481 original_address,
482 value.chain,
483 account.slots.clone(),
484 Some(account.native_balance),
485 None,
486 ChangeType::Update,
487 );
488 (*impl_addr, proxy_state)
489 }
490 None => {
491 let impl_addr = generate_proxy_token_address(
495 state_guard.proxy_token_addresses.len() as u32,
496 )?;
497 state_guard
498 .proxy_token_addresses
499 .insert(original_address, impl_addr);
500
501 let proxy_state = create_proxy_token_account(
503 original_address,
504 Some(impl_addr),
505 &account.slots,
506 value.chain,
507 Some(account.native_balance),
508 );
509
510 (impl_addr, proxy_state)
511 }
512 };
513
514 proxy_token_accounts.insert(original_address, proxy_state);
515
516 let impl_update = ResponseAccount {
518 address: impl_addr,
519 slots: HashMap::new(),
520 ..account.clone()
521 };
522 storage_by_address.insert(impl_addr, impl_update);
523 } else {
524 storage_by_address.insert(account.address, account);
526 }
527 }
528
529 let mut proxy_creates: Vec<AccountUpdate> = Vec::new();
533 let mut proxy_updates: HashMap<Address, AccountUpdate> = HashMap::new();
534 for (addr, update) in proxy_token_accounts {
535 if matches!(update.change, ChangeType::Creation) {
536 proxy_creates.push(update);
537 } else {
538 proxy_updates.insert(addr, update);
539 }
540 }
541
542 info!("Updating engine with {} contracts from snapshots", storage_by_address.len());
543 update_engine(
544 SHARED_TYCHO_DB.clone(),
545 header.clone().block(),
546 Some(storage_by_address),
547 proxy_updates,
548 )
549 .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
550
551 if !proxy_creates.is_empty() {
555 SHARED_TYCHO_DB
556 .force_update_accounts(proxy_creates)
557 .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
558 }
559 info!("Engine updated");
560 drop(state_guard);
561 }
562
563 let account_balances = protocol_msg
566 .clone()
567 .snapshots
568 .get_vm_storage()
569 .iter()
570 .filter_map(|(addr, acc)| {
571 if acc.token_balances.is_empty() {
572 return None;
573 }
574 let balances = acc
575 .token_balances
576 .iter()
577 .map(|(token_addr, ab)| (token_addr.clone(), ab.balance.clone()))
578 .collect::<HashMap<Bytes, Bytes>>();
579 Some((addr.clone(), balances))
580 })
581 .collect::<AccountBalances>();
582
583 let mut new_components = HashMap::new();
584 let mut count_token_skips = 0;
585 let mut components_to_store = HashMap::new();
586 {
587 let state_guard = self.state.read().await;
588
589 'snapshot_loop: for (id, snapshot) in protocol_msg
591 .snapshots
592 .get_states()
593 .clone()
594 {
595 if !self.admits(protocol.as_str(), &snapshot) {
597 continue;
598 }
599
600 let mut component_tokens = Vec::new();
602 let mut new_tokens_accounts = HashMap::new();
603 for token in snapshot.component.tokens.clone() {
604 match state_guard.tokens.get(&token) {
605 Some(token) => {
606 component_tokens.push(token.clone());
607
608 let token_address = match bytes_to_address(&token.address) {
611 Ok(addr) => addr,
612 Err(_) => {
613 count_token_skips += 1;
614 msg_failed_components.insert(id.clone());
615 warn!(
616 "Token address could not be decoded {}, ignoring pool {:x?}",
617 token.address, id
618 );
619 continue 'snapshot_loop;
620 }
621 };
622 if !state_guard
624 .proxy_token_addresses
625 .contains_key(&token_address)
626 {
627 new_tokens_accounts.insert(
628 token_address,
629 create_proxy_token_account(
630 token_address,
631 None,
632 &HashMap::new(),
633 snapshot.component.chain,
634 None,
635 ),
636 );
637 }
638 }
639 None => {
640 count_token_skips += 1;
641 msg_failed_components.insert(id.clone());
642 debug!("Token not found {}, ignoring pool {:x?}", token, id);
643 continue 'snapshot_loop;
644 }
645 }
646 }
647 let component = ProtocolComponent::from_with_tokens(
648 snapshot.component.clone(),
649 component_tokens,
650 );
651
652 if !new_tokens_accounts.is_empty() {
654 update_engine(
655 SHARED_TYCHO_DB.clone(),
656 header.clone().block(),
657 None,
658 new_tokens_accounts,
659 )
660 .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
661 }
662
663 if !component
666 .static_attributes
667 .contains_key("manual_updates")
668 {
669 for contract in &component.contract_ids {
670 contracts_map
671 .entry(contract.clone())
672 .or_insert_with(HashSet::new)
673 .insert(id.clone());
674 }
675 for (_, tracing) in snapshot.entrypoints.iter() {
678 for contract in tracing.accessed_slots.keys().cloned() {
679 contracts_map
680 .entry(contract)
681 .or_insert_with(HashSet::new)
682 .insert(id.clone());
683 }
684 }
685 }
686
687 new_pairs.insert(id.clone(), component.clone());
689
690 components_to_store.insert(id.clone(), component);
692
693 if let Some(state_decode_f) = self.registry.get(protocol.as_str()) {
695 let live_override = self
696 .override_providers
697 .get(protocol.as_str())
698 .and_then(|provider| provider.subscribe(protocol.as_str()));
699 match state_decode_f(
700 snapshot,
701 header.clone(),
702 account_balances.clone(),
703 self.state.clone(),
704 live_override,
705 )
706 .await
707 {
708 Ok(state) => {
709 new_components.insert(id.clone(), state);
710 }
711 Err(e) => {
712 if self.skip_state_decode_failures {
713 warn!(pool = id, error = %e, "StateDecodingFailure");
714 msg_failed_components.insert(id.clone());
715 continue 'snapshot_loop;
716 } else {
717 error!(pool = id, error = %e, "StateDecodingFailure");
718 return Err(StreamDecodeError::Fatal(format!("{e}")));
719 }
720 }
721 }
722 } else if self.skip_state_decode_failures {
723 warn!(pool = id, "MissingDecoderRegistration");
724 msg_failed_components.insert(id.clone());
725 continue 'snapshot_loop;
726 } else {
727 error!(pool = id, "MissingDecoderRegistration");
728 return Err(StreamDecodeError::Fatal(format!(
729 "Missing decoder registration for: {id}"
730 )));
731 }
732 }
733 }
734
735 if !components_to_store.is_empty() {
737 let mut state_guard = self.state.write().await;
738 for (id, component) in components_to_store {
739 state_guard
740 .components
741 .insert(id, component);
742 }
743 }
744
745 if !protocol_msg.snapshots.states.is_empty() {
746 info!("Decoded {} snapshots for protocol {protocol}", new_components.len());
747 }
748 if count_token_skips > 0 {
749 info!("Skipped {count_token_skips} pools due to missing tokens");
750 }
751
752 updated_states.extend(new_components);
754
755 if let Some(deltas) = protocol_msg.deltas.clone() {
757 let mut state_guard = self.state.write().await;
759
760 let mut account_update_by_address: HashMap<Address, AccountUpdate> = HashMap::new();
761 let mut new_proxy_accounts: Vec<AccountUpdate> = Vec::new();
763 for (key, value) in deltas.account_deltas.iter() {
764 let mut update: AccountUpdate = value.clone().into();
765
766 if update.code.is_none() && matches!(update.change, ChangeType::Creation) {
772 error!(
773 update = ?update,
774 "FaultyCreationDelta"
775 );
776 update.code = Some(vec![]);
777 }
778
779 if state_guard.tokens.contains_key(key) {
780 let original_address = update.address;
781 let impl_addr = match state_guard
788 .proxy_token_addresses
789 .get(&original_address)
790 {
791 Some(impl_addr) => {
792 let proxy_update = AccountUpdate {
797 code: None,
798 change: ChangeType::Update,
799 ..update.clone()
800 };
801 account_update_by_address.insert(original_address, proxy_update);
802
803 *impl_addr
804 }
805 None => {
806 let impl_addr = generate_proxy_token_address(
811 state_guard.proxy_token_addresses.len() as u32,
812 )?;
813 state_guard
814 .proxy_token_addresses
815 .insert(original_address, impl_addr);
816
817 let proxy_state = create_proxy_token_account(
822 original_address,
823 Some(impl_addr),
824 &update.slots,
825 update.chain,
826 update.balance,
827 );
828 new_proxy_accounts.push(proxy_state);
829
830 impl_addr
831 }
832 };
833
834 if update.code.is_some() {
836 let impl_update = AccountUpdate {
837 address: impl_addr,
838 slots: HashMap::new(),
839 ..update.clone()
840 };
841 account_update_by_address.insert(impl_addr, impl_update);
842 }
843 } else {
844 account_update_by_address.insert(update.address, update);
846 }
847 }
848 drop(state_guard);
849
850 let state_guard = self.state.read().await;
851 info!("Updating engine with {} contract deltas", deltas.account_deltas.len());
852 update_engine(
853 SHARED_TYCHO_DB.clone(),
854 header.clone().block(),
855 None,
856 account_update_by_address,
857 )
858 .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
859
860 if !new_proxy_accounts.is_empty() {
863 SHARED_TYCHO_DB
864 .force_update_accounts(new_proxy_accounts)
865 .map_err(|e| StreamDecodeError::Fatal(e.to_string()))?;
866 }
867 info!("Engine updated");
868
869 let mut pools_to_update = HashSet::new();
871 for (account, _update) in deltas.account_deltas {
872 pools_to_update.extend(
874 contracts_map
875 .get(&account)
876 .cloned()
877 .unwrap_or_default(),
878 );
879 pools_to_update.extend(
881 state_guard
882 .contracts_map
883 .get(&account)
884 .cloned()
885 .unwrap_or_default(),
886 );
887 }
888
889 let all_balances = Balances {
891 component_balances: deltas
892 .component_balances
893 .iter()
894 .map(|(pool_id, bals)| {
895 let mut balances = HashMap::new();
896 for (t, b) in bals {
897 balances.insert(t.clone(), b.balance.clone());
898 }
899 pools_to_update.insert(pool_id.clone());
900 (pool_id.clone(), balances)
901 })
902 .collect(),
903 account_balances: deltas
904 .account_balances
905 .iter()
906 .map(|(account, bals)| {
907 let mut balances = HashMap::new();
908 for (t, b) in bals {
909 balances.insert(t.clone(), b.balance.clone());
910 }
911 pools_to_update.extend(
912 contracts_map
913 .get(account)
914 .cloned()
915 .unwrap_or_default(),
916 );
917 (account.clone(), balances)
918 })
919 .collect(),
920 };
921
922 for (id, update) in deltas.state_deltas {
924 let update_with_block = Self::add_block_info_to_delta(
926 ProtocolStateDelta::from(update),
927 current_block.clone(),
928 );
929 match Self::apply_update(
930 &id,
931 update_with_block,
932 &mut updated_states,
933 &state_guard,
934 &all_balances,
935 ) {
936 Ok(_) => {
937 pools_to_update.remove(&id);
938 }
939 Err(e) => {
940 if self.skip_state_decode_failures {
941 warn!(pool = id, error = %e, "Failed to apply state update, marking component as removed");
942 updated_states.remove(&id);
944 if let Some(component) = new_pairs.remove(&id) {
946 removed_pairs.insert(id.clone(), component);
947 } else if let Some(component) = state_guard.components.get(&id) {
948 removed_pairs.insert(id.clone(), component.clone());
949 } else {
950 warn!(pool = id, "Component not found in new_pairs or state, cannot add to removed_pairs");
953 }
954 pools_to_update.remove(&id);
955
956 msg_failed_components.insert(id.clone());
958 } else {
959 return Err(e);
960 }
961 }
962 }
963 }
964
965 for pool in pools_to_update {
967 let default_delta_with_block = Self::add_block_info_to_delta(
969 ProtocolStateDelta::default(),
970 current_block.clone(),
971 );
972 match Self::apply_update(
973 &pool,
974 default_delta_with_block,
975 &mut updated_states,
976 &state_guard,
977 &all_balances,
978 ) {
979 Ok(_) => {}
980 Err(e) => {
981 if self.skip_state_decode_failures {
982 warn!(pool = pool, error = %e, "Failed to apply contract/balance update, marking component as removed");
983 updated_states.remove(&pool);
985 if let Some(component) = new_pairs.remove(&pool) {
987 removed_pairs.insert(pool.clone(), component);
988 } else if let Some(component) = state_guard.components.get(&pool) {
989 removed_pairs.insert(pool.clone(), component.clone());
990 } else {
991 warn!(pool = pool, "Component not found in new_pairs or state, cannot add to removed_pairs");
994 }
995
996 msg_failed_components.insert(pool.clone());
998 } else {
999 return Err(e);
1000 }
1001 }
1002 }
1003 }
1004 };
1005 }
1006
1007 let mut state_guard = self.state.write().await;
1009
1010 state_guard
1012 .failed_components
1013 .extend(msg_failed_components);
1014
1015 updated_states.retain(|id, _| {
1019 !state_guard
1020 .failed_components
1021 .contains(id)
1022 });
1023 new_pairs.retain(|id, _| {
1024 !state_guard
1025 .failed_components
1026 .contains(id)
1027 });
1028
1029 if let Some(header) = current_block.as_ref() {
1030 let execution_block = self.execution_block(header);
1031 let decoder_state = &mut *state_guard;
1032 Self::refresh_execution_block(
1033 &mut updated_states,
1034 &mut decoder_state.states,
1035 &decoder_state.failed_components,
1036 &removed_pairs,
1037 &execution_block,
1038 );
1039 }
1040
1041 state_guard
1042 .states
1043 .extend(updated_states.clone());
1044
1045 state_guard.current_block_number = block_number_or_timestamp;
1046
1047 for (id, component) in new_pairs.iter() {
1049 state_guard
1050 .components
1051 .insert(id.clone(), component.clone());
1052 }
1053
1054 for id in removed_pairs.keys() {
1056 state_guard.components.remove(id);
1057 }
1058
1059 for (key, values) in contracts_map {
1060 state_guard
1061 .contracts_map
1062 .entry(key)
1063 .or_insert_with(HashSet::new)
1064 .extend(values);
1065 }
1066
1067 Ok(Update::new(block_number_or_timestamp, updated_states, new_pairs)
1069 .set_is_partial(is_partial)
1070 .set_removed_pairs(removed_pairs)
1071 .set_sync_states(msg.sync_states.clone()))
1072 }
1073
1074 pub async fn apply_deltas_ephemeral(
1094 &self,
1095 pending_deltas: &HashMap<String, BlockAggregatedChanges>,
1096 header: H,
1097 ) -> Result<Update, StreamDecodeError> {
1098 let block_number_or_timestamp = header
1099 .clone()
1100 .block_number_or_timestamp();
1101 let current_block = header.block();
1102 let state_guard = self.state.read().await;
1103
1104 let mut updated_states: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1105
1106 for deltas in pending_deltas.values() {
1107 let all_balances = Balances {
1108 component_balances: deltas
1109 .component_balances
1110 .iter()
1111 .map(|(pool_id, bals)| {
1112 let balances = bals
1113 .iter()
1114 .map(|(t, b)| (t.clone(), b.balance.clone()))
1115 .collect();
1116 (pool_id.clone(), balances)
1117 })
1118 .collect(),
1119 account_balances: HashMap::new(),
1120 };
1121
1122 for (id, state_delta) in &deltas.state_deltas {
1123 let dto_delta = Self::add_block_info_to_delta(
1124 ProtocolStateDelta::from(state_delta.clone()),
1125 current_block.clone(),
1126 );
1127 if let Err(e) = Self::apply_update(
1128 id,
1129 dto_delta,
1130 &mut updated_states,
1131 &state_guard,
1132 &all_balances,
1133 ) {
1134 warn!(pool = id, error = %e, "EphemeralDeltaTransitionError");
1135 }
1136 }
1137 }
1138
1139 if let Some(header) = current_block.as_ref() {
1144 let execution_block = BlockContext::new(header.number, header.timestamp);
1145 for state in updated_states.values_mut() {
1146 state.apply_block(&execution_block);
1147 }
1148 }
1149
1150 Ok(Update::new(block_number_or_timestamp, updated_states, HashMap::new()))
1151 }
1152
1153 fn add_block_info_to_delta(
1155 mut delta: ProtocolStateDelta,
1156 block_header_opt: Option<BlockHeader>,
1157 ) -> ProtocolStateDelta {
1158 if let Some(header) = block_header_opt {
1159 delta.updated_attributes.insert(
1162 "block_number".to_string(),
1163 Bytes::from(header.number.to_be_bytes().to_vec()),
1164 );
1165 delta.updated_attributes.insert(
1166 "block_timestamp".to_string(),
1167 Bytes::from(header.timestamp.to_be_bytes().to_vec()),
1168 );
1169 }
1170 delta
1171 }
1172
1173 fn apply_update(
1174 id: &String,
1175 update: ProtocolStateDelta,
1176 updated_states: &mut HashMap<String, Box<dyn ProtocolSim>>,
1177 state_guard: &RwLockReadGuard<'_, DecoderState>,
1178 all_balances: &Balances,
1179 ) -> Result<(), StreamDecodeError> {
1180 match updated_states.entry(id.clone()) {
1181 Entry::Occupied(mut entry) => {
1182 let state: &mut Box<dyn ProtocolSim> = entry.get_mut();
1184 state
1185 .delta_transition(update, &state_guard.tokens, all_balances)
1186 .map_err(|e| {
1187 error!(pool = id, error = ?e, "DeltaTransitionError");
1188 StreamDecodeError::Fatal(format!("TransitionFailure: {e:?}"))
1189 })?;
1190 }
1191 Entry::Vacant(_) => {
1192 match state_guard.states.get(id) {
1193 Some(stored_state) => {
1196 let mut state = stored_state.clone();
1197 state
1198 .delta_transition(update, &state_guard.tokens, all_balances)
1199 .map_err(|e| {
1200 error!(pool = id, error = ?e, "DeltaTransitionError");
1201 StreamDecodeError::Fatal(format!("TransitionFailure: {e:?}"))
1202 })?;
1203 updated_states.insert(id.clone(), state);
1204 }
1205 None => debug!(pool = id, reason = "MissingState", "DeltaTransitionError"),
1206 }
1207 }
1208 }
1209 Ok(())
1210 }
1211}
1212
1213fn generate_proxy_token_address(idx: u32) -> Result<Address, StreamDecodeError> {
1215 let padded_idx = format!("{idx:x}");
1216 let padded_zeroes = "0".repeat(33 - padded_idx.len());
1217 let proxy_token_address = format!("{padded_zeroes}{padded_idx}BAdbaBe");
1218 let decoded = hex::decode(proxy_token_address).map_err(|e| {
1219 StreamDecodeError::Fatal(format!("Invalid proxy token address encoding: {e}"))
1220 })?;
1221
1222 const ADDRESS_LENGTH: usize = 20;
1223 if decoded.len() != ADDRESS_LENGTH {
1224 return Err(StreamDecodeError::Fatal(format!(
1225 "Invalid proxy token address length: expected {}, got {}",
1226 ADDRESS_LENGTH,
1227 decoded.len(),
1228 )));
1229 }
1230
1231 Ok(Address::from_slice(&decoded))
1232}
1233
1234fn create_proxy_token_account(
1239 addr: Address,
1240 new_address: Option<Address>,
1241 storage: &HashMap<U256, U256>,
1242 chain: Chain,
1243 balance: Option<U256>,
1244) -> AccountUpdate {
1245 let mut slots = storage.clone();
1246 if let Some(new_address) = new_address {
1247 slots.insert(*IMPLEMENTATION_SLOT, U256::from_be_slice(new_address.as_slice()));
1248 }
1249
1250 AccountUpdate {
1251 address: addr,
1252 chain,
1253 slots,
1254 balance,
1255 code: Some(ERC20_PROXY_BYTECODE.to_vec()),
1256 change: ChangeType::Creation,
1257 }
1258}
1259
1260#[cfg(test)]
1261mock! {
1262 #[derive(Debug)]
1263 pub ProtocolSim {
1264 pub fn fee(&self) -> f64;
1265 pub fn spot_price(&self, base: &Token, quote: &Token) -> Result<f64, SimulationError>;
1266 pub fn get_amount_out(
1267 &self,
1268 amount_in: BigUint,
1269 token_in: &Token,
1270 token_out: &Token,
1271 ) -> Result<GetAmountOutResult, SimulationError>;
1272 pub fn get_limits(
1273 &self,
1274 sell_token: Bytes,
1275 buy_token: Bytes,
1276 ) -> Result<(BigUint, BigUint), SimulationError>;
1277 pub fn delta_transition(
1278 &mut self,
1279 delta: ProtocolStateDelta,
1280 tokens: &HashMap<Bytes, Token>,
1281 balances: &Balances,
1282 ) -> Result<(), TransitionError>;
1283 pub fn clone_box(&self) -> Box<dyn ProtocolSim>;
1284 pub fn eq(&self, other: &dyn ProtocolSim) -> bool;
1285 }
1286}
1287
1288#[cfg(test)]
1289crate::impl_non_serializable_protocol!(MockProtocolSim, "test protocol");
1290
1291#[cfg(test)]
1292impl ProtocolSim for MockProtocolSim {
1293 fn fee(&self) -> f64 {
1294 self.fee()
1295 }
1296
1297 fn spot_price(&self, base: &Token, quote: &Token) -> Result<f64, SimulationError> {
1298 self.spot_price(base, quote)
1299 }
1300
1301 fn get_amount_out(
1302 &self,
1303 amount_in: BigUint,
1304 token_in: &Token,
1305 token_out: &Token,
1306 ) -> Result<GetAmountOutResult, SimulationError> {
1307 self.get_amount_out(amount_in, token_in, token_out)
1308 }
1309
1310 fn get_limits(
1311 &self,
1312 sell_token: Bytes,
1313 buy_token: Bytes,
1314 ) -> Result<(BigUint, BigUint), SimulationError> {
1315 self.get_limits(sell_token, buy_token)
1316 }
1317
1318 fn delta_transition(
1319 &mut self,
1320 delta: ProtocolStateDelta,
1321 tokens: &HashMap<Bytes, Token>,
1322 balances: &Balances,
1323 ) -> Result<(), TransitionError> {
1324 self.delta_transition(delta, tokens, balances)
1325 }
1326
1327 fn clone_box(&self) -> Box<dyn ProtocolSim> {
1328 self.clone_box()
1329 }
1330
1331 fn as_any(&self) -> &dyn Any {
1332 panic!("MockProtocolSim does not support as_any")
1333 }
1334
1335 fn as_any_mut(&mut self) -> &mut dyn Any {
1336 panic!("MockProtocolSim does not support as_any_mut")
1337 }
1338
1339 fn eq(&self, other: &dyn ProtocolSim) -> bool {
1340 self.eq(other)
1341 }
1342
1343 fn typetag_name(&self) -> &'static str {
1344 unreachable!()
1345 }
1346
1347 fn typetag_deserialize(&self) {
1348 unreachable!()
1349 }
1350}
1351
1352#[cfg(test)]
1353mod tests {
1354 use std::str::FromStr;
1355
1356 use alloy::primitives::address;
1357 use mockall::predicate::*;
1358 use rstest::*;
1359 use tycho_client::feed::BlockHeader;
1360 use tycho_common::{models::Chain, Bytes};
1361
1362 use super::*;
1363
1364 fn header_at(number: u64, timestamp: u64, partial: Option<u32>) -> BlockHeader {
1365 BlockHeader {
1366 hash: Bytes::from([0u8; 32]),
1367 number,
1368 parent_hash: Bytes::from([0u8; 32]),
1369 revert: false,
1370 timestamp,
1371 partial_block_index: partial,
1372 }
1373 }
1374
1375 fn block_sensitive_state() -> Box<dyn ProtocolSim> {
1377 use crate::evm::protocol::{
1378 aerodrome_slipstreams::state::AerodromeSlipstreamsState,
1379 utils::{
1380 slipstreams::{dynamic_fee_module::DynamicFeeConfig, observations::Observation},
1381 uniswap::{tick_list::TickInfo, tick_math::get_sqrt_ratio_at_tick},
1382 },
1383 };
1384
1385 Box::new(
1386 AerodromeSlipstreamsState::new(
1387 "block-sensitive".to_string(),
1388 0,
1389 1_000_000_000_000_000_000,
1390 get_sqrt_ratio_at_tick(0).unwrap(),
1391 0,
1392 1,
1393 3000,
1394 1,
1395 0,
1396 vec![TickInfo::new(-120, 0).unwrap(), TickInfo::new(120, 0).unwrap()],
1397 vec![Observation { block_timestamp: 500, initialized: true, ..Default::default() }],
1398 DynamicFeeConfig::new(2700, 30_000, 0, true, 750),
1399 )
1400 .expect("state should build")
1401 .with_position_assumption(crate::protocol::models::BlockPositionAssumption::First),
1404 )
1405 }
1406
1407 #[test]
1408 fn confirmed_header_targets_the_next_block() {
1409 let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Base);
1410
1411 let execution_block = decoder.execution_block(&header_at(100, 1_000, None));
1412
1413 assert_eq!(execution_block.number(), 101);
1414 assert_eq!(execution_block.timestamp(), 1_000 + Chain::Base.block_time_secs());
1415 }
1416
1417 #[test]
1418 fn partial_header_targets_the_block_that_is_still_open() {
1419 let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
1420
1421 let execution_block = decoder.execution_block(&header_at(100, 1_000, Some(3)));
1422
1423 assert_eq!(execution_block.number(), 100);
1424 assert_eq!(execution_block.timestamp(), 1_000);
1425 }
1426
1427 #[test]
1428 fn refresh_re_emits_a_state_whose_fee_flipped_without_a_delta() {
1429 let mut stored = HashMap::from([("block-sensitive".to_string(), {
1433 let mut state = block_sensitive_state();
1434 state.apply_block(&BlockContext::new(100, 500));
1435 state
1436 })]);
1437 let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1438
1439 TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1440 &mut updated,
1441 &mut stored,
1442 &HashSet::new(),
1443 &HashMap::<String, ()>::new(),
1444 &BlockContext::new(101, 502),
1445 );
1446
1447 let emitted = updated
1448 .get("block-sensitive")
1449 .expect("a fee flip must be emitted even without a delta");
1450 assert_eq!(emitted.fee(), 750.0 / 1_000_000.0);
1451 assert_eq!(stored["block-sensitive"].fee(), 750.0 / 1_000_000.0);
1453 }
1454
1455 #[test]
1456 fn refresh_never_re_emits_failed_components() {
1457 let mut stored = HashMap::from([("zombie".to_string(), {
1460 let mut state = block_sensitive_state();
1461 state.apply_block(&BlockContext::new(100, 500));
1462 state
1463 })]);
1464 let failed = HashSet::from(["zombie".to_string()]);
1465 let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1466
1467 TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1468 &mut updated,
1469 &mut stored,
1470 &failed,
1471 &HashMap::<String, ()>::new(),
1472 &BlockContext::new(101, 502),
1473 );
1474
1475 assert!(updated.is_empty());
1476 }
1477
1478 #[test]
1479 fn refresh_never_re_emits_removed_components() {
1480 let mut stored = HashMap::from([("gone".to_string(), {
1483 let mut state = block_sensitive_state();
1484 state.apply_block(&BlockContext::new(100, 500));
1485 state
1486 })]);
1487 let removed = HashMap::from([("gone".to_string(), ())]);
1488 let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1489
1490 TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1491 &mut updated,
1492 &mut stored,
1493 &HashSet::new(),
1494 &removed,
1495 &BlockContext::new(101, 502),
1496 );
1497
1498 assert!(updated.is_empty());
1499 }
1500
1501 #[test]
1502 fn refresh_stays_quiet_when_no_fee_changed() {
1503 let mut stored: HashMap<String, Box<dyn ProtocolSim>> = HashMap::from([
1506 ("idle-sensitive".to_string(), {
1507 let mut state = block_sensitive_state();
1508 state.apply_block(&BlockContext::new(101, 502));
1509 state
1510 }),
1511 (
1512 "univ2".to_string(),
1513 Box::new(crate::evm::protocol::uniswap_v2::state::UniswapV2State::new(
1514 U256::from(1_000_000u64),
1515 U256::from(1_000_000u64),
1516 )) as Box<dyn ProtocolSim>,
1517 ),
1518 ]);
1519 let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1520
1521 TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1522 &mut updated,
1523 &mut stored,
1524 &HashSet::new(),
1525 &HashMap::<String, ()>::new(),
1526 &BlockContext::new(102, 504),
1527 );
1528
1529 assert!(updated.is_empty());
1530 }
1531 use crate::evm::protocol::{curve::CurveState, uniswap_v2::state::UniswapV2State};
1532
1533 #[test]
1534 fn curve_vm_adapter_registration_flagged_deprecated() {
1535 assert!(!is_deprecated_curve_registration::<CurveState>("vm:curve"));
1537 assert!(is_deprecated_curve_registration::<UniswapV2State>("vm:curve"));
1539 assert!(!is_deprecated_curve_registration::<UniswapV2State>("uniswap_v2"));
1541 }
1542
1543 async fn setup_decoder(set_tokens: bool) -> TychoStreamDecoder<BlockHeader> {
1544 let mut decoder = TychoStreamDecoder::new(Chain::Ethereum);
1545 decoder.register_decoder::<UniswapV2State>("uniswap_v2");
1546 if set_tokens {
1547 let tokens = [
1548 Bytes::from("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").lpad(20, 0),
1549 Bytes::from("0xdac17f958d2ee523a2206206994597c13d831ec7").lpad(20, 0),
1550 ]
1551 .iter()
1552 .map(|addr| {
1553 let addr_str = format!("{addr:x}");
1554 (
1555 addr.clone(),
1556 Token::new(addr, &addr_str, 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1557 )
1558 })
1559 .collect();
1560 decoder.set_tokens(tokens).await;
1561 }
1562 decoder
1563 }
1564
1565 fn load_test_msg(name: &str) -> FeedMessage<BlockHeader> {
1566 use std::{fs, path::Path};
1567
1568 use tycho_client::feed::dto;
1569 let project_root = env!("CARGO_MANIFEST_DIR");
1570 let asset_path = Path::new(project_root).join(format!("tests/assets/decoder/{name}.json"));
1571 let json_data = fs::read_to_string(asset_path).expect("Failed to read test asset");
1572 let feed_msg: dto::FeedMessage<BlockHeader> =
1573 serde_json::from_str(&json_data).expect("Failed to deserialize FeedMsg json!");
1574 FeedMessage::from(feed_msg)
1575 }
1576
1577 #[tokio::test]
1578 async fn test_decode() {
1579 let decoder = setup_decoder(true).await;
1580
1581 let msg = load_test_msg("uniswap_v2_snapshot");
1582 let res1 = decoder
1583 .decode(&msg)
1584 .await
1585 .expect("decode failure");
1586 let msg = load_test_msg("uniswap_v2_delta");
1587 let res2 = decoder
1588 .decode(&msg)
1589 .await
1590 .expect("decode failure");
1591
1592 assert_eq!(res1.states.len(), 1);
1593 assert_eq!(res2.states.len(), 1);
1594 assert_eq!(res1.sync_states.len(), 1);
1595 assert_eq!(res2.sync_states.len(), 1);
1596 }
1597
1598 #[tokio::test]
1599 async fn test_decode_token_creation_delta_with_existing_proxy() {
1600 let decoder = setup_decoder(true).await;
1601 let msg = load_test_msg("uniswap_v2_delta_token_creation");
1602
1603 decoder
1606 .decode(&msg)
1607 .await
1608 .expect("first decode (proxy creation) failed");
1609
1610 decoder
1614 .decode(&msg)
1615 .await
1616 .expect("decode of a token Creation delta with an existing proxy failed");
1617 }
1618
1619 #[tokio::test]
1620 async fn test_decode_component_missing_token() {
1621 let decoder = setup_decoder(false).await;
1622 let tokens = [Bytes::from("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").lpad(20, 0)]
1623 .iter()
1624 .map(|addr| {
1625 let addr_str = format!("{addr:x}");
1626 (
1627 addr.clone(),
1628 Token::new(addr, &addr_str, 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1629 )
1630 })
1631 .collect();
1632 decoder.set_tokens(tokens).await;
1633
1634 let msg = load_test_msg("uniswap_v2_snapshot");
1635 let res1 = decoder
1636 .decode(&msg)
1637 .await
1638 .expect("decode failure");
1639
1640 assert_eq!(res1.states.len(), 0);
1641 }
1642
1643 #[tokio::test]
1644 async fn test_decode_component_bad_id() {
1645 let decoder = setup_decoder(true).await;
1646 let msg = load_test_msg("uniswap_v2_snapshot_broken_id");
1647
1648 match decoder.decode(&msg).await {
1649 Err(StreamDecodeError::Fatal(msg)) => {
1650 assert_eq!(msg, "Component id mismatch");
1651 }
1652 Ok(_) => {
1653 panic!("Expected failures to be raised")
1654 }
1655 }
1656 }
1657
1658 #[rstest]
1659 #[case(true)]
1660 #[case(false)]
1661 #[tokio::test]
1662 async fn test_decode_component_bad_state(#[case] skip_failures: bool) {
1663 let mut decoder = setup_decoder(true).await;
1664 decoder.skip_state_decode_failures = skip_failures;
1665
1666 let msg = load_test_msg("uniswap_v2_snapshot_broken_state");
1667 match decoder.decode(&msg).await {
1668 Err(StreamDecodeError::Fatal(msg)) => {
1669 if !skip_failures {
1670 assert_eq!(msg, "Missing attributes reserve0");
1671 } else {
1672 panic!("Expected failures to be ignored. Err: {msg}")
1673 }
1674 }
1675 Ok(res) => {
1676 if !skip_failures {
1677 panic!("Expected failures to be raised")
1678 } else {
1679 assert_eq!(res.states.len(), 0);
1680 }
1681 }
1682 }
1683 }
1684
1685 #[tokio::test]
1686 async fn test_decode_updates_state_on_contract_change() {
1687 let decoder = setup_decoder(true).await;
1688
1689 let mut mock_state = MockProtocolSim::new();
1691
1692 mock_state
1693 .expect_clone_box()
1694 .times(1)
1695 .returning(|| {
1696 let mut cloned_mock_state = MockProtocolSim::new();
1697 cloned_mock_state
1699 .expect_delta_transition()
1700 .times(1)
1701 .returning(|_, _, _| Ok(()));
1702 cloned_mock_state
1703 .expect_clone_box()
1704 .times(1)
1705 .returning(|| Box::new(MockProtocolSim::new()));
1706 Box::new(cloned_mock_state)
1707 });
1708
1709 let pool_id =
1711 "0x93d199263632a4ef4bb438f1feb99e57b4b5f0bd0000000000000000000005c2".to_string();
1712 decoder
1713 .state
1714 .write()
1715 .await
1716 .states
1717 .insert(pool_id.clone(), Box::new(mock_state) as Box<dyn ProtocolSim>);
1718 decoder
1719 .state
1720 .write()
1721 .await
1722 .contracts_map
1723 .insert(
1724 Bytes::from("0xba12222222228d8ba445958a75a0704d566bf2c8").lpad(20, 0),
1725 HashSet::from([pool_id.clone()]),
1726 );
1727
1728 let msg = load_test_msg("balancer_v2_delta");
1730
1731 let _ = decoder
1733 .decode(&msg)
1734 .await
1735 .expect("decode failure");
1736
1737 }
1739
1740 #[test]
1741 fn test_generate_proxy_token_address() {
1742 let idx = 1;
1743 let generated_address =
1744 generate_proxy_token_address(idx).expect("proxy token address should be valid");
1745 assert_eq!(generated_address, address!("000000000000000000000000000000001badbabe"));
1746
1747 let idx = 123456;
1748 let generated_address =
1749 generate_proxy_token_address(idx).expect("proxy token address should be valid");
1750 assert_eq!(generated_address, address!("00000000000000000000000000001e240badbabe"));
1751 }
1752
1753 fn hooked_v4_msg() -> FeedMessage<BlockHeader> {
1759 use tycho_client::feed::synchronizer::{Snapshot, StateSyncMessage};
1760 use tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState};
1761
1762 let id = "0x00000000000000000000000000000000000000000000000000000000000000c4";
1763 let static_attributes = HashMap::from([
1764 ("key_lp_fee".to_string(), Bytes::from(500_i32.to_be_bytes().to_vec())),
1765 ("tick_spacing".to_string(), Bytes::from(60_i32.to_be_bytes().to_vec())),
1766 (
1767 "hooks".to_string(),
1768 Bytes::from_str("0x00000000000000000000000000000000000000c4").unwrap(),
1769 ),
1770 ]);
1771 let attributes = HashMap::from([
1772 ("liquidity".to_string(), Bytes::from(100_u64.to_be_bytes().to_vec())),
1773 ("tick".to_string(), Bytes::from(300_i32.to_be_bytes().to_vec())),
1774 (
1775 "sqrt_price_x96".to_string(),
1776 Bytes::from(
1777 79228162514264337593543950336_u128
1778 .to_be_bytes()
1779 .to_vec(),
1780 ),
1781 ),
1782 ("protocol_fees/zero2one".to_string(), Bytes::from(0_u32.to_be_bytes().to_vec())),
1783 ("protocol_fees/one2zero".to_string(), Bytes::from(0_u32.to_be_bytes().to_vec())),
1784 ("ticks/60/net_liquidity".to_string(), Bytes::from(400_i128.to_be_bytes().to_vec())),
1785 ]);
1786
1787 let snapshot = ComponentWithState {
1788 state: ProtocolComponentState::new(id, attributes, HashMap::new()),
1789 component: ProtocolComponent {
1790 id: id.to_string(),
1791 static_attributes,
1792 ..Default::default()
1793 },
1794 component_tvl: None,
1795 entrypoints: Vec::new(),
1796 };
1797
1798 FeedMessage {
1799 state_msgs: HashMap::from([(
1800 "uniswap_v4_hooks".to_string(),
1801 StateSyncMessage {
1802 header: header_at(1, 1_000, None),
1803 snapshots: Snapshot {
1804 states: HashMap::from([(id.to_string(), snapshot)]),
1805 vm_storage: HashMap::new(),
1806 },
1807 deltas: None,
1808 removed_components: HashMap::new(),
1809 },
1810 )]),
1811 sync_states: HashMap::new(),
1812 }
1813 }
1814
1815 fn hooked_v4_decoder(chain: Chain, context: DecoderContext) -> TychoStreamDecoder<BlockHeader> {
1817 let mut decoder = TychoStreamDecoder::new(chain);
1818 decoder.register_decoder_with_context::<crate::evm::protocol::uniswap_v4::state::UniswapV4State>(
1819 "uniswap_v4_hooks", context
1820 );
1821 decoder
1822 }
1823
1824 #[tokio::test]
1825 async fn test_unregistered_hook_is_rejected_off_the_generic_vm_chains() {
1826 let decoder = hooked_v4_decoder(Chain::Robinhood, DecoderContext::new());
1827
1828 let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
1829 panic!("a hooked pool on robinhood must not decode");
1830 };
1831
1832 assert!(message.contains("unsupported uniswap v4 hook"), "{message}");
1833 assert!(message.contains("robinhood"), "{message}");
1834 }
1835
1836 #[tokio::test]
1837 async fn test_decoder_chain_overrides_the_registered_context_chain() {
1838 let decoder =
1839 hooked_v4_decoder(Chain::Robinhood, DecoderContext::new().chain(Chain::Ethereum));
1840
1841 let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
1842 panic!("the decoder's chain must win over the context's");
1843 };
1844
1845 assert!(message.contains("robinhood"), "{message}");
1846 }
1847
1848 #[tokio::test]
1849 async fn test_unregistered_hook_takes_the_generic_path_on_ethereum() {
1850 let decoder = hooked_v4_decoder(Chain::Ethereum, DecoderContext::new());
1851
1852 let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
1853 panic!("the generic VM creator has no balance_owner attribute to work from");
1854 };
1855
1856 assert!(!message.contains("unsupported uniswap v4 hook"), "{message}");
1857 assert!(message.contains("balance_owner"), "{message}");
1858 }
1859
1860 #[tokio::test(flavor = "multi_thread")]
1861 async fn test_euler_hook_low_pool_manager_balance() {
1862 let mut decoder = TychoStreamDecoder::new(Chain::Ethereum);
1863
1864 decoder.register_decoder_with_context::<crate::evm::protocol::uniswap_v4::state::UniswapV4State>(
1865 "uniswap_v4_hooks", DecoderContext::new().vm_traces(true)
1866 );
1867
1868 let weth = Bytes::from_str("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").unwrap();
1869 let teth = Bytes::from_str("0xd11c452fc99cf405034ee446803b6f6c1f6d5ed8").unwrap();
1870 let tokens = HashMap::from([
1871 (
1872 weth.clone(),
1873 Token::new(&weth, "WETH", 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1874 ),
1875 (
1876 teth.clone(),
1877 Token::new(&teth, "tETH", 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1878 ),
1879 ]);
1880
1881 decoder.set_tokens(tokens.clone()).await;
1882
1883 let msg = load_test_msg("euler_hook_snapshot");
1884 let res = decoder
1885 .decode(&msg)
1886 .await
1887 .expect("decode failure");
1888
1889 let pool_state = res
1890 .states
1891 .get("0xc70d7fbd7fcccdf726e02fed78548b40dc52502b097c7a1ee7d995f4d4396134")
1892 .expect("Couldn't find target pool");
1893 let amount_out = pool_state
1894 .get_amount_out(
1895 BigUint::from_str("1000000000000000000").unwrap(),
1896 tokens.get(&teth).unwrap(),
1897 tokens.get(&weth).unwrap(),
1898 )
1899 .expect("Get amount out failed");
1900
1901 assert_eq!(amount_out.amount, BigUint::from_str("1216190190361759119").unwrap());
1902 }
1903
1904 fn pons_feed_message(snapshot: ComponentWithState) -> FeedMessage<BlockHeader> {
1906 use tycho_client::feed::synchronizer::{Snapshot, StateSyncMessage};
1907
1908 use crate::evm::protocol::uniswap_v4::pons_fixture;
1909
1910 FeedMessage {
1911 state_msgs: HashMap::from([(
1912 "uniswap_v4_hooks".to_string(),
1913 StateSyncMessage {
1914 header: pons_fixture::header(),
1915 snapshots: Snapshot {
1916 states: HashMap::from([(pons_fixture::POOL_ID.to_string(), snapshot)]),
1917 vm_storage: HashMap::new(),
1918 },
1919 deltas: None,
1920 removed_components: HashMap::new(),
1921 },
1922 )]),
1923 sync_states: HashMap::new(),
1924 }
1925 }
1926
1927 async fn pons_stream_decoder() -> TychoStreamDecoder<BlockHeader> {
1929 use crate::evm::protocol::uniswap_v4::{
1930 hooks::hook_handler_creator::initialize_hook_handlers, pons_fixture,
1931 state::UniswapV4State,
1932 };
1933
1934 initialize_hook_handlers().expect("hook handler registration should succeed");
1935 let mut decoder = TychoStreamDecoder::new(Chain::Robinhood);
1936 decoder.register_decoder::<UniswapV4State>("uniswap_v4_hooks");
1937 decoder
1938 .set_tokens(pons_fixture::tokens())
1939 .await;
1940 decoder
1941 }
1942
1943 #[tokio::test]
1947 async fn test_pons_pool_decodes_through_the_stream_decoder_on_robinhood() {
1948 use crate::evm::protocol::uniswap_v4::{
1949 hooks::pons_v2::hook_handler::PonsV2HookHandler, pons_fixture, state::UniswapV4State,
1950 };
1951
1952 let decoder = pons_stream_decoder().await;
1953
1954 let result = decoder
1955 .decode(&pons_feed_message(pons_fixture::snapshot()))
1956 .await
1957 .expect("decode failure");
1958
1959 let emitted = result
1960 .states
1961 .get(pons_fixture::POOL_ID)
1962 .expect("the Pons pool must reach consumers");
1963 let pool = emitted
1964 .as_any()
1965 .downcast_ref::<UniswapV4State>()
1966 .expect("a uniswap_v4_hooks component decodes into a UniswapV4State");
1967 let handler = pool
1968 .hook
1969 .as_ref()
1970 .expect("a Pons pool must carry a hook handler");
1971 let pons = handler
1972 .as_any()
1973 .downcast_ref::<PonsV2HookHandler>()
1974 .expect("the registry must hand back the native Pons handler");
1975 assert_eq!(u32::from(pons.hook_fee_bps()), pons_fixture::HOOK_FEE_BPS);
1976 assert_eq!(u32::from(pons.creator_tax_bps()), pons_fixture::CREATOR_TAX_BPS);
1977
1978 let core = pons_fixture::decode_on(pons_fixture::hookless_snapshot(), Chain::Robinhood)
1979 .await
1980 .expect("the same pool with no hook must decode too");
1981 let (xlg, nvda) = (pons_fixture::xlg(), pons_fixture::nvda());
1982 let amount_in = BigUint::from(10u64).pow(18);
1983 for (token_in, token_out) in [(&xlg, &nvda), (&nvda, &xlg)] {
1984 let core_out = core
1985 .get_amount_out(amount_in.clone(), token_in, token_out)
1986 .expect("the reference pool quotes one whole token")
1987 .amount;
1988 let hooked_out = emitted
1989 .get_amount_out(amount_in.clone(), token_in, token_out)
1990 .expect("the emitted pool quotes one whole token")
1991 .amount;
1992
1993 let expected = pons_fixture::net_of_hook_take(&core_out);
1994 assert!(expected < core_out, "the hook must take something out of {core_out}");
1995 assert_eq!(hooked_out, expected, "{} -> {}", token_in.symbol, token_out.symbol);
1996 }
1997 }
1998
1999 #[tokio::test]
2002 async fn test_unknown_hook_component_is_not_emitted_on_robinhood() {
2003 use crate::evm::protocol::uniswap_v4::pons_fixture;
2004
2005 let mut decoder = pons_stream_decoder().await;
2006 decoder.skip_state_decode_failures(true);
2007
2008 let result = decoder
2009 .decode(&pons_feed_message(pons_fixture::unknown_hook_snapshot()))
2010 .await
2011 .expect("a skipped decode failure is not fatal");
2012
2013 assert!(
2014 !result
2015 .states
2016 .contains_key(pons_fixture::POOL_ID),
2017 "a pool whose hook cannot be modelled must not be quoted"
2018 );
2019 assert!(result.states.is_empty());
2020 }
2021
2022 fn component_with_id(id: &str) -> ComponentWithState {
2023 use tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState};
2024
2025 ComponentWithState {
2026 state: ProtocolComponentState::new(id, HashMap::new(), HashMap::new()),
2027 component: ProtocolComponent { id: id.to_string(), ..Default::default() },
2028 component_tvl: None,
2029 entrypoints: Vec::new(),
2030 }
2031 }
2032
2033 fn rejects_a(component: &ComponentWithState) -> bool {
2034 component.component.id != "a"
2035 }
2036
2037 fn rejects_b(component: &ComponentWithState) -> bool {
2038 component.component.id != "b"
2039 }
2040
2041 #[test]
2042 fn test_admits_requires_every_registered_filter() {
2043 let mut decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
2046 decoder.register_filter("x", rejects_a);
2047 decoder.register_filter("x", rejects_b);
2048
2049 assert!(!decoder.admits("x", &component_with_id("a")));
2050 assert!(!decoder.admits("x", &component_with_id("b")));
2051 assert!(decoder.admits("x", &component_with_id("c")));
2052 assert!(decoder.admits("y", &component_with_id("a")));
2053 }
2054}