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 contracts_map: HashMap<Bytes, HashSet<String>>,
60 proxy_token_addresses: HashMap<Address, Address>,
62 failed_components: HashSet<String>,
66 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
84pub 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 override_providers: HashMap<String, Arc<dyn StateOverrideProvider>>,
108 block_time_secs: u64,
111 chain: Chain,
114}
115
116fn 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 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 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 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 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 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 pub fn min_token_quality(&mut self, quality: u32) {
223 self.min_token_quality = quality;
224 }
225
226 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 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 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 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 pub async fn decode(&self, msg: &FeedMessage<H>) -> Result<Update, StreamDecodeError> {
354 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 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 {
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 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 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 let (impl_addr, proxy_state) = match state_guard
471 .proxy_token_addresses
472 .get(&original_address)
473 {
474 Some(impl_addr) => {
475 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 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 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 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 storage_by_address.insert(account.address, account);
527 }
528 }
529
530 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 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 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 'snapshot_loop: for (id, snapshot) in protocol_msg
592 .snapshots
593 .get_states()
594 .clone()
595 {
596 if !self.admits(protocol.as_str(), &snapshot) {
598 continue;
599 }
600
601 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 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 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 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 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 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 new_pairs.insert(id.clone(), component.clone());
690
691 components_to_store.insert(id.clone(), component);
693
694 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 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 updated_states.extend(new_components);
755
756 if let Some(deltas) = protocol_msg.deltas.clone() {
758 let mut state_guard = self.state.write().await;
760
761 let mut account_update_by_address: HashMap<Address, AccountUpdate> = HashMap::new();
762 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 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 let impl_addr = match state_guard
789 .proxy_token_addresses
790 .get(&original_address)
791 {
792 Some(impl_addr) => {
793 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 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 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 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 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 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 let mut pools_to_update = HashSet::new();
872 for (account, _update) in deltas.account_deltas {
873 pools_to_update.extend(
875 contracts_map
876 .get(&account)
877 .cloned()
878 .unwrap_or_default(),
879 );
880 pools_to_update.extend(
882 state_guard
883 .contracts_map
884 .get(&account)
885 .cloned()
886 .unwrap_or_default(),
887 );
888 }
889
890 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 for (id, update) in deltas.state_deltas {
925 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 updated_states.remove(&id);
945 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 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 msg_failed_components.insert(id.clone());
959 } else {
960 return Err(e);
961 }
962 }
963 }
964 }
965
966 for pool in pools_to_update {
968 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 updated_states.remove(&pool);
986 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 warn!(pool = pool, "Component not found in new_pairs or state, cannot add to removed_pairs");
995 }
996
997 msg_failed_components.insert(pool.clone());
999 } else {
1000 return Err(e);
1001 }
1002 }
1003 }
1004 }
1005 };
1006 }
1007
1008 let mut state_guard = self.state.write().await;
1010
1011 state_guard
1013 .failed_components
1014 .extend(msg_failed_components);
1015
1016 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 for (id, component) in new_pairs.iter() {
1050 state_guard
1051 .components
1052 .insert(id.clone(), component.clone());
1053 }
1054
1055 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 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 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 ¤t_block,
1143 &mut updated_states,
1144 &state_guard,
1145 &all_balances,
1146 )?;
1147 }
1148 }
1149
1150 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 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 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 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 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 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 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
1279fn 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
1300fn 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 "MockProtocolSim"
1412 }
1413
1414 fn typetag_deserialize(&self) {
1415 unreachable!()
1416 }
1417}
1418
1419#[cfg(test)]
1420mod tests {
1421 use std::str::FromStr;
1422
1423 use alloy::primitives::address;
1424 use mockall::predicate::*;
1425 use rstest::*;
1426 use tycho_client::feed::BlockHeader;
1427 use tycho_common::{models::Chain, Bytes};
1428
1429 use super::*;
1430
1431 fn header_at(number: u64, timestamp: u64, partial: Option<u32>) -> BlockHeader {
1432 BlockHeader {
1433 hash: Bytes::from([0u8; 32]),
1434 number,
1435 parent_hash: Bytes::from([0u8; 32]),
1436 revert: false,
1437 timestamp,
1438 partial_block_index: partial,
1439 }
1440 }
1441
1442 fn block_sensitive_state() -> Box<dyn ProtocolSim> {
1444 use crate::evm::protocol::{
1445 aerodrome_slipstreams::state::AerodromeSlipstreamsState,
1446 utils::{
1447 slipstreams::{dynamic_fee_module::DynamicFeeConfig, observations::Observation},
1448 uniswap::{tick_list::TickInfo, tick_math::get_sqrt_ratio_at_tick},
1449 },
1450 };
1451
1452 Box::new(
1453 AerodromeSlipstreamsState::new(
1454 "block-sensitive".to_string(),
1455 0,
1456 1_000_000_000_000_000_000,
1457 get_sqrt_ratio_at_tick(0).unwrap(),
1458 0,
1459 1,
1460 3000,
1461 1,
1462 0,
1463 vec![TickInfo::new(-120, 0).unwrap(), TickInfo::new(120, 0).unwrap()],
1464 vec![Observation { block_timestamp: 500, initialized: true, ..Default::default() }],
1465 DynamicFeeConfig::new(2700, 30_000, 0, true, 750),
1466 )
1467 .expect("state should build")
1468 .with_position_assumption(crate::protocol::models::BlockPositionAssumption::First),
1471 )
1472 }
1473
1474 #[test]
1475 fn confirmed_header_targets_the_next_block() {
1476 let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Base);
1477
1478 let execution_block = decoder.execution_block(&header_at(100, 1_000, None));
1479
1480 assert_eq!(execution_block.number(), 101);
1481 assert_eq!(execution_block.timestamp(), 1_000 + Chain::Base.block_time_secs());
1482 }
1483
1484 #[test]
1485 fn partial_header_targets_the_block_that_is_still_open() {
1486 let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
1487
1488 let execution_block = decoder.execution_block(&header_at(100, 1_000, Some(3)));
1489
1490 assert_eq!(execution_block.number(), 100);
1491 assert_eq!(execution_block.timestamp(), 1_000);
1492 }
1493
1494 #[test]
1495 fn refresh_re_emits_a_state_whose_fee_flipped_without_a_delta() {
1496 let mut stored = HashMap::from([("block-sensitive".to_string(), {
1500 let mut state = block_sensitive_state();
1501 state.apply_block(&BlockContext::new(100, 500));
1502 state
1503 })]);
1504 let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1505
1506 TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1507 &mut updated,
1508 &mut stored,
1509 &HashSet::new(),
1510 &HashMap::<String, ()>::new(),
1511 &BlockContext::new(101, 502),
1512 );
1513
1514 let emitted = updated
1515 .get("block-sensitive")
1516 .expect("a fee flip must be emitted even without a delta");
1517 assert_eq!(emitted.fee(), 750.0 / 1_000_000.0);
1518 assert_eq!(stored["block-sensitive"].fee(), 750.0 / 1_000_000.0);
1520 }
1521
1522 #[test]
1523 fn refresh_never_re_emits_failed_components() {
1524 let mut stored = HashMap::from([("zombie".to_string(), {
1527 let mut state = block_sensitive_state();
1528 state.apply_block(&BlockContext::new(100, 500));
1529 state
1530 })]);
1531 let failed = HashSet::from(["zombie".to_string()]);
1532 let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1533
1534 TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1535 &mut updated,
1536 &mut stored,
1537 &failed,
1538 &HashMap::<String, ()>::new(),
1539 &BlockContext::new(101, 502),
1540 );
1541
1542 assert!(updated.is_empty());
1543 }
1544
1545 #[test]
1546 fn refresh_never_re_emits_removed_components() {
1547 let mut stored = HashMap::from([("gone".to_string(), {
1550 let mut state = block_sensitive_state();
1551 state.apply_block(&BlockContext::new(100, 500));
1552 state
1553 })]);
1554 let removed = HashMap::from([("gone".to_string(), ())]);
1555 let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1556
1557 TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1558 &mut updated,
1559 &mut stored,
1560 &HashSet::new(),
1561 &removed,
1562 &BlockContext::new(101, 502),
1563 );
1564
1565 assert!(updated.is_empty());
1566 }
1567
1568 #[test]
1569 fn refresh_stays_quiet_when_no_fee_changed() {
1570 let mut stored: HashMap<String, Box<dyn ProtocolSim>> = HashMap::from([
1573 ("idle-sensitive".to_string(), {
1574 let mut state = block_sensitive_state();
1575 state.apply_block(&BlockContext::new(101, 502));
1576 state
1577 }),
1578 (
1579 "univ2".to_string(),
1580 Box::new(crate::evm::protocol::uniswap_v2::state::UniswapV2State::new(
1581 U256::from(1_000_000u64),
1582 U256::from(1_000_000u64),
1583 )) as Box<dyn ProtocolSim>,
1584 ),
1585 ]);
1586 let mut updated: HashMap<String, Box<dyn ProtocolSim>> = HashMap::new();
1587
1588 TychoStreamDecoder::<BlockHeader>::refresh_execution_block(
1589 &mut updated,
1590 &mut stored,
1591 &HashSet::new(),
1592 &HashMap::<String, ()>::new(),
1593 &BlockContext::new(102, 504),
1594 );
1595
1596 assert!(updated.is_empty());
1597 }
1598 use crate::evm::protocol::{curve::CurveState, uniswap_v2::state::UniswapV2State};
1599
1600 #[test]
1601 fn curve_vm_adapter_registration_flagged_deprecated() {
1602 assert!(!is_deprecated_curve_registration::<CurveState>("vm:curve"));
1604 assert!(is_deprecated_curve_registration::<UniswapV2State>("vm:curve"));
1606 assert!(!is_deprecated_curve_registration::<UniswapV2State>("uniswap_v2"));
1608 }
1609
1610 async fn setup_decoder(set_tokens: bool) -> TychoStreamDecoder<BlockHeader> {
1611 let mut decoder = TychoStreamDecoder::new(Chain::Ethereum);
1612 decoder.register_decoder::<UniswapV2State>("uniswap_v2");
1613 if set_tokens {
1614 let tokens = [
1615 Bytes::from("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").lpad(20, 0),
1616 Bytes::from("0xdac17f958d2ee523a2206206994597c13d831ec7").lpad(20, 0),
1617 ]
1618 .iter()
1619 .map(|addr| {
1620 let addr_str = format!("{addr:x}");
1621 (
1622 addr.clone(),
1623 Token::new(addr, &addr_str, 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1624 )
1625 })
1626 .collect();
1627 decoder.set_tokens(tokens).await;
1628 }
1629 decoder
1630 }
1631
1632 fn load_test_msg(name: &str) -> FeedMessage<BlockHeader> {
1633 use std::{fs, path::Path};
1634
1635 use tycho_client::feed::dto;
1636 let project_root = env!("CARGO_MANIFEST_DIR");
1637 let asset_path = Path::new(project_root).join(format!("tests/assets/decoder/{name}.json"));
1638 let json_data = fs::read_to_string(asset_path).expect("Failed to read test asset");
1639 let feed_msg: dto::FeedMessage<BlockHeader> =
1640 serde_json::from_str(&json_data).expect("Failed to deserialize FeedMsg json!");
1641 FeedMessage::from(feed_msg)
1642 }
1643
1644 #[tokio::test]
1645 async fn test_decode() {
1646 let decoder = setup_decoder(true).await;
1647
1648 let msg = load_test_msg("uniswap_v2_snapshot");
1649 let res1 = decoder
1650 .decode(&msg)
1651 .await
1652 .expect("decode failure");
1653 let msg = load_test_msg("uniswap_v2_delta");
1654 let res2 = decoder
1655 .decode(&msg)
1656 .await
1657 .expect("decode failure");
1658
1659 assert_eq!(res1.states.len(), 1);
1660 assert_eq!(res2.states.len(), 1);
1661 assert_eq!(res1.sync_states.len(), 1);
1662 assert_eq!(res2.sync_states.len(), 1);
1663 }
1664
1665 #[tokio::test]
1666 async fn test_decode_token_creation_delta_with_existing_proxy() {
1667 let decoder = setup_decoder(true).await;
1668 let msg = load_test_msg("uniswap_v2_delta_token_creation");
1669
1670 decoder
1673 .decode(&msg)
1674 .await
1675 .expect("first decode (proxy creation) failed");
1676
1677 decoder
1681 .decode(&msg)
1682 .await
1683 .expect("decode of a token Creation delta with an existing proxy failed");
1684 }
1685
1686 #[tokio::test]
1687 async fn test_decode_component_missing_token() {
1688 let decoder = setup_decoder(false).await;
1689 let tokens = [Bytes::from("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").lpad(20, 0)]
1690 .iter()
1691 .map(|addr| {
1692 let addr_str = format!("{addr:x}");
1693 (
1694 addr.clone(),
1695 Token::new(addr, &addr_str, 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
1696 )
1697 })
1698 .collect();
1699 decoder.set_tokens(tokens).await;
1700
1701 let msg = load_test_msg("uniswap_v2_snapshot");
1702 let res1 = decoder
1703 .decode(&msg)
1704 .await
1705 .expect("decode failure");
1706
1707 assert_eq!(res1.states.len(), 0);
1708 }
1709
1710 #[tokio::test]
1711 async fn test_decode_component_bad_id() {
1712 let decoder = setup_decoder(true).await;
1713 let msg = load_test_msg("uniswap_v2_snapshot_broken_id");
1714
1715 match decoder.decode(&msg).await {
1716 Err(StreamDecodeError::Fatal(msg)) => {
1717 assert_eq!(msg, "Component id mismatch");
1718 }
1719 Ok(_) => {
1720 panic!("Expected failures to be raised")
1721 }
1722 }
1723 }
1724
1725 #[rstest]
1726 #[case(true)]
1727 #[case(false)]
1728 #[tokio::test]
1729 async fn test_decode_component_bad_state(#[case] skip_failures: bool) {
1730 let mut decoder = setup_decoder(true).await;
1731 decoder.skip_state_decode_failures = skip_failures;
1732
1733 let msg = load_test_msg("uniswap_v2_snapshot_broken_state");
1734 match decoder.decode(&msg).await {
1735 Err(StreamDecodeError::Fatal(msg)) => {
1736 if !skip_failures {
1737 assert_eq!(msg, "Missing attributes reserve0");
1738 } else {
1739 panic!("Expected failures to be ignored. Err: {msg}")
1740 }
1741 }
1742 Ok(res) => {
1743 if !skip_failures {
1744 panic!("Expected failures to be raised")
1745 } else {
1746 assert_eq!(res.states.len(), 0);
1747 }
1748 }
1749 }
1750 }
1751
1752 #[tokio::test]
1753 async fn test_decode_updates_state_on_contract_change() {
1754 let decoder = setup_decoder(true).await;
1755
1756 let mut mock_state = MockProtocolSim::new();
1758
1759 mock_state
1760 .expect_clone_box()
1761 .times(1)
1762 .returning(|| {
1763 let mut cloned_mock_state = MockProtocolSim::new();
1764 cloned_mock_state
1766 .expect_delta_transition()
1767 .times(1)
1768 .returning(|_, _, _| Ok(()));
1769 cloned_mock_state
1770 .expect_clone_box()
1771 .times(1)
1772 .returning(|| Box::new(MockProtocolSim::new()));
1773 Box::new(cloned_mock_state)
1774 });
1775
1776 let pool_id =
1778 "0x93d199263632a4ef4bb438f1feb99e57b4b5f0bd0000000000000000000005c2".to_string();
1779 decoder
1780 .state
1781 .write()
1782 .await
1783 .states
1784 .insert(pool_id.clone(), Box::new(mock_state) as Box<dyn ProtocolSim>);
1785 decoder
1786 .state
1787 .write()
1788 .await
1789 .contracts_map
1790 .insert(
1791 Bytes::from("0xba12222222228d8ba445958a75a0704d566bf2c8").lpad(20, 0),
1792 HashSet::from([pool_id.clone()]),
1793 );
1794
1795 let msg = load_test_msg("balancer_v2_delta");
1797
1798 let _ = decoder
1800 .decode(&msg)
1801 .await
1802 .expect("decode failure");
1803
1804 }
1806
1807 #[tokio::test]
1810 async fn test_apply_deltas_ephemeral_sets_pending_overrides_on_the_clone_only() {
1811 use tycho_common::models::{
1812 blockchain::BlockAggregatedChanges, contract::AccountDelta,
1813 protocol::ProtocolComponentStateDelta, ChangeType as ModelChangeType,
1814 };
1815
1816 use crate::evm::protocol::uniswap_v4::state::{UniswapV4Fees, UniswapV4State};
1817
1818 let decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
1819 let pool_id = "0xhooked".to_string();
1820 let quiet_pool_id = "0xhooked-no-log".to_string();
1821 let pool = UniswapV4State::new(
1822 1000,
1823 U256::from(1u8) << 96,
1824 UniswapV4Fees::new(0, 0, 3000),
1825 0,
1826 60,
1827 vec![],
1828 )
1829 .unwrap();
1830 let hook = Bytes::from(vec![7u8; 20]);
1831 {
1832 let mut state = decoder.state.write().await;
1833 state
1834 .states
1835 .insert(pool_id.clone(), Box::new(pool.clone()));
1836 state
1837 .states
1838 .insert(quiet_pool_id.clone(), Box::new(pool));
1839 state
1840 .contracts_map
1841 .insert(hook.clone(), HashSet::from([pool_id.clone(), quiet_pool_id.clone()]));
1842 }
1843
1844 let deltas = BlockAggregatedChanges {
1845 extractor: "uniswap_v4_hooks".to_string(),
1846 state_deltas: HashMap::from([(
1847 pool_id.clone(),
1848 ProtocolComponentStateDelta {
1849 component_id: pool_id.clone(),
1850 updated_attributes: HashMap::from([(
1851 "liquidity".to_string(),
1852 Bytes::from(2000_u64.to_be_bytes().to_vec()),
1853 )]),
1854 deleted_attributes: HashSet::new(),
1855 created_attributes: HashSet::new(),
1856 },
1857 )]),
1858 account_deltas: HashMap::from([(
1859 hook.clone(),
1860 AccountDelta::new(
1861 Chain::Ethereum,
1862 hook.clone(),
1863 HashMap::from([(Bytes::from(vec![0u8]), Some(Bytes::from(vec![9u8])))]),
1864 None,
1865 None,
1866 ModelChangeType::Update,
1867 ),
1868 )]),
1869 ..Default::default()
1870 };
1871 let header = BlockHeader { number: 10, timestamp: 20, ..Default::default() };
1872
1873 let update = decoder
1874 .apply_deltas_ephemeral(
1875 &HashMap::from([("uniswap_v4_hooks".to_string(), deltas)]),
1876 header,
1877 )
1878 .await
1879 .unwrap();
1880
1881 let clone = update.states[&pool_id]
1882 .as_any()
1883 .downcast_ref::<UniswapV4State>()
1884 .unwrap();
1885 let pending = clone
1886 .pending_overrides()
1887 .expect("the clone carries the pending block's overrides");
1888 assert_eq!(
1889 pending.block,
1890 Some(BlockEnvOverrides { number: Some(10), timestamp: Some(20) })
1891 );
1892 assert_eq!(
1893 pending.storage.as_ref().unwrap()[&Address::from_slice(&hook)][&U256::ZERO],
1894 U256::from(9)
1895 );
1896 let quiet = update.states[&quiet_pool_id]
1897 .as_any()
1898 .downcast_ref::<UniswapV4State>()
1899 .expect("a pool whose hook was written is cloned without a delta of its own");
1900 assert!(quiet.pending_overrides().is_some());
1901
1902 let stored = decoder.state.read().await;
1903 for id in [&pool_id, &quiet_pool_id] {
1904 let stored = stored.states[id]
1905 .as_any()
1906 .downcast_ref::<UniswapV4State>()
1907 .unwrap();
1908 assert!(stored.pending_overrides().is_none(), "confirmed state is never overlaid");
1909 }
1910 }
1911
1912 #[test]
1913 fn test_generate_proxy_token_address() {
1914 let idx = 1;
1915 let generated_address =
1916 generate_proxy_token_address(idx).expect("proxy token address should be valid");
1917 assert_eq!(generated_address, address!("000000000000000000000000000000001badbabe"));
1918
1919 let idx = 123456;
1920 let generated_address =
1921 generate_proxy_token_address(idx).expect("proxy token address should be valid");
1922 assert_eq!(generated_address, address!("00000000000000000000000000001e240badbabe"));
1923 }
1924
1925 fn hooked_v4_msg() -> FeedMessage<BlockHeader> {
1931 use tycho_client::feed::synchronizer::{Snapshot, StateSyncMessage};
1932 use tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState};
1933
1934 let id = "0x00000000000000000000000000000000000000000000000000000000000000c4";
1935 let static_attributes = HashMap::from([
1936 ("key_lp_fee".to_string(), Bytes::from(500_i32.to_be_bytes().to_vec())),
1937 ("tick_spacing".to_string(), Bytes::from(60_i32.to_be_bytes().to_vec())),
1938 (
1939 "hooks".to_string(),
1940 Bytes::from_str("0x00000000000000000000000000000000000000c4").unwrap(),
1941 ),
1942 ]);
1943 let attributes = HashMap::from([
1944 ("liquidity".to_string(), Bytes::from(100_u64.to_be_bytes().to_vec())),
1945 ("tick".to_string(), Bytes::from(300_i32.to_be_bytes().to_vec())),
1946 (
1947 "sqrt_price_x96".to_string(),
1948 Bytes::from(
1949 79228162514264337593543950336_u128
1950 .to_be_bytes()
1951 .to_vec(),
1952 ),
1953 ),
1954 ("protocol_fees/zero2one".to_string(), Bytes::from(0_u32.to_be_bytes().to_vec())),
1955 ("protocol_fees/one2zero".to_string(), Bytes::from(0_u32.to_be_bytes().to_vec())),
1956 ("ticks/60/net_liquidity".to_string(), Bytes::from(400_i128.to_be_bytes().to_vec())),
1957 ]);
1958
1959 let snapshot = ComponentWithState {
1960 state: ProtocolComponentState::new(id, attributes, HashMap::new()),
1961 component: ProtocolComponent {
1962 id: id.to_string(),
1963 static_attributes,
1964 ..Default::default()
1965 },
1966 component_tvl: None,
1967 entrypoints: Vec::new(),
1968 };
1969
1970 FeedMessage {
1971 state_msgs: HashMap::from([(
1972 "uniswap_v4_hooks".to_string(),
1973 StateSyncMessage {
1974 header: header_at(1, 1_000, None),
1975 snapshots: Snapshot {
1976 states: HashMap::from([(id.to_string(), snapshot)]),
1977 vm_storage: HashMap::new(),
1978 },
1979 deltas: None,
1980 removed_components: HashMap::new(),
1981 },
1982 )]),
1983 sync_states: HashMap::new(),
1984 }
1985 }
1986
1987 fn hooked_v4_decoder(chain: Chain, context: DecoderContext) -> TychoStreamDecoder<BlockHeader> {
1989 let mut decoder = TychoStreamDecoder::new(chain);
1990 decoder.register_decoder_with_context::<crate::evm::protocol::uniswap_v4::state::UniswapV4State>(
1991 "uniswap_v4_hooks", context
1992 );
1993 decoder
1994 }
1995
1996 #[tokio::test]
1997 async fn test_unregistered_hook_is_rejected_off_the_generic_vm_chains() {
1998 let decoder = hooked_v4_decoder(Chain::Robinhood, DecoderContext::new());
1999
2000 let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
2001 panic!("a hooked pool on robinhood must not decode");
2002 };
2003
2004 assert!(message.contains("unsupported uniswap v4 hook"), "{message}");
2005 assert!(message.contains("robinhood"), "{message}");
2006 }
2007
2008 #[tokio::test]
2009 async fn test_decoder_chain_overrides_the_registered_context_chain() {
2010 let decoder =
2011 hooked_v4_decoder(Chain::Robinhood, DecoderContext::new().chain(Chain::Ethereum));
2012
2013 let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
2014 panic!("the decoder's chain must win over the context's");
2015 };
2016
2017 assert!(message.contains("robinhood"), "{message}");
2018 }
2019
2020 #[tokio::test]
2021 async fn test_unregistered_hook_takes_the_generic_path_on_ethereum() {
2022 let decoder = hooked_v4_decoder(Chain::Ethereum, DecoderContext::new());
2023
2024 let Err(StreamDecodeError::Fatal(message)) = decoder.decode(&hooked_v4_msg()).await else {
2025 panic!("the generic VM creator has no balance_owner attribute to work from");
2026 };
2027
2028 assert!(!message.contains("unsupported uniswap v4 hook"), "{message}");
2029 assert!(message.contains("balance_owner"), "{message}");
2030 }
2031
2032 #[tokio::test(flavor = "multi_thread")]
2033 async fn test_euler_hook_low_pool_manager_balance() {
2034 let mut decoder = TychoStreamDecoder::new(Chain::Ethereum);
2035
2036 decoder.register_decoder_with_context::<crate::evm::protocol::uniswap_v4::state::UniswapV4State>(
2037 "uniswap_v4_hooks", DecoderContext::new().vm_traces(true)
2038 );
2039
2040 let weth = Bytes::from_str("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").unwrap();
2041 let teth = Bytes::from_str("0xd11c452fc99cf405034ee446803b6f6c1f6d5ed8").unwrap();
2042 let tokens = HashMap::from([
2043 (
2044 weth.clone(),
2045 Token::new(&weth, "WETH", 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
2046 ),
2047 (
2048 teth.clone(),
2049 Token::new(&teth, "tETH", 18, 100, &[Some(100_000)], Chain::Ethereum, 100),
2050 ),
2051 ]);
2052
2053 decoder.set_tokens(tokens.clone()).await;
2054
2055 let msg = load_test_msg("euler_hook_snapshot");
2056 let res = decoder
2057 .decode(&msg)
2058 .await
2059 .expect("decode failure");
2060
2061 let pool_state = res
2062 .states
2063 .get("0xc70d7fbd7fcccdf726e02fed78548b40dc52502b097c7a1ee7d995f4d4396134")
2064 .expect("Couldn't find target pool");
2065 let amount_out = pool_state
2066 .get_amount_out(
2067 BigUint::from_str("1000000000000000000").unwrap(),
2068 tokens.get(&teth).unwrap(),
2069 tokens.get(&weth).unwrap(),
2070 )
2071 .expect("Get amount out failed");
2072
2073 assert_eq!(amount_out.amount, BigUint::from_str("1216190190361759119").unwrap());
2074 }
2075
2076 fn pons_feed_message(snapshot: ComponentWithState) -> FeedMessage<BlockHeader> {
2078 use tycho_client::feed::synchronizer::{Snapshot, StateSyncMessage};
2079
2080 use crate::evm::protocol::uniswap_v4::pons_fixture;
2081
2082 FeedMessage {
2083 state_msgs: HashMap::from([(
2084 "uniswap_v4_hooks".to_string(),
2085 StateSyncMessage {
2086 header: pons_fixture::header(),
2087 snapshots: Snapshot {
2088 states: HashMap::from([(pons_fixture::POOL_ID.to_string(), snapshot)]),
2089 vm_storage: HashMap::new(),
2090 },
2091 deltas: None,
2092 removed_components: HashMap::new(),
2093 },
2094 )]),
2095 sync_states: HashMap::new(),
2096 }
2097 }
2098
2099 async fn pons_stream_decoder() -> TychoStreamDecoder<BlockHeader> {
2101 use crate::evm::protocol::uniswap_v4::{
2102 hooks::hook_handler_creator::initialize_hook_handlers, pons_fixture,
2103 state::UniswapV4State,
2104 };
2105
2106 initialize_hook_handlers().expect("hook handler registration should succeed");
2107 let mut decoder = TychoStreamDecoder::new(Chain::Robinhood);
2108 decoder.register_decoder::<UniswapV4State>("uniswap_v4_hooks");
2109 decoder
2110 .set_tokens(pons_fixture::tokens())
2111 .await;
2112 decoder
2113 }
2114
2115 #[tokio::test]
2119 async fn test_pons_pool_decodes_through_the_stream_decoder_on_robinhood() {
2120 use crate::evm::protocol::uniswap_v4::{
2121 hooks::pons_v2::hook_handler::PonsV2HookHandler, pons_fixture, state::UniswapV4State,
2122 };
2123
2124 let decoder = pons_stream_decoder().await;
2125
2126 let result = decoder
2127 .decode(&pons_feed_message(pons_fixture::snapshot()))
2128 .await
2129 .expect("decode failure");
2130
2131 let emitted = result
2132 .states
2133 .get(pons_fixture::POOL_ID)
2134 .expect("the Pons pool must reach consumers");
2135 let pool = emitted
2136 .as_any()
2137 .downcast_ref::<UniswapV4State>()
2138 .expect("a uniswap_v4_hooks component decodes into a UniswapV4State");
2139 let handler = pool
2140 .hook
2141 .as_ref()
2142 .expect("a Pons pool must carry a hook handler");
2143 let pons = handler
2144 .as_any()
2145 .downcast_ref::<PonsV2HookHandler>()
2146 .expect("the registry must hand back the native Pons handler");
2147 assert_eq!(u32::from(pons.hook_fee_bps()), pons_fixture::HOOK_FEE_BPS);
2148 assert_eq!(u32::from(pons.creator_tax_bps()), pons_fixture::CREATOR_TAX_BPS);
2149
2150 let core = pons_fixture::decode_on(pons_fixture::hookless_snapshot(), Chain::Robinhood)
2151 .await
2152 .expect("the same pool with no hook must decode too");
2153 let (xlg, nvda) = (pons_fixture::xlg(), pons_fixture::nvda());
2154 let amount_in = BigUint::from(10u64).pow(18);
2155 for (token_in, token_out) in [(&xlg, &nvda), (&nvda, &xlg)] {
2156 let core_out = core
2157 .get_amount_out(amount_in.clone(), token_in, token_out)
2158 .expect("the reference pool quotes one whole token")
2159 .amount;
2160 let hooked_out = emitted
2161 .get_amount_out(amount_in.clone(), token_in, token_out)
2162 .expect("the emitted pool quotes one whole token")
2163 .amount;
2164
2165 let expected = pons_fixture::net_of_hook_take(&core_out);
2166 assert!(expected < core_out, "the hook must take something out of {core_out}");
2167 assert_eq!(hooked_out, expected, "{} -> {}", token_in.symbol, token_out.symbol);
2168 }
2169 }
2170
2171 #[tokio::test]
2174 async fn test_unknown_hook_component_is_not_emitted_on_robinhood() {
2175 use crate::evm::protocol::uniswap_v4::pons_fixture;
2176
2177 let mut decoder = pons_stream_decoder().await;
2178 decoder.skip_state_decode_failures(true);
2179
2180 let result = decoder
2181 .decode(&pons_feed_message(pons_fixture::unknown_hook_snapshot()))
2182 .await
2183 .expect("a skipped decode failure is not fatal");
2184
2185 assert!(
2186 !result
2187 .states
2188 .contains_key(pons_fixture::POOL_ID),
2189 "a pool whose hook cannot be modelled must not be quoted"
2190 );
2191 assert!(result.states.is_empty());
2192 }
2193
2194 fn component_with_id(id: &str) -> ComponentWithState {
2195 use tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState};
2196
2197 ComponentWithState {
2198 state: ProtocolComponentState::new(id, HashMap::new(), HashMap::new()),
2199 component: ProtocolComponent { id: id.to_string(), ..Default::default() },
2200 component_tvl: None,
2201 entrypoints: Vec::new(),
2202 }
2203 }
2204
2205 fn rejects_a(component: &ComponentWithState) -> bool {
2206 component.component.id != "a"
2207 }
2208
2209 fn rejects_b(component: &ComponentWithState) -> bool {
2210 component.component.id != "b"
2211 }
2212
2213 #[test]
2214 fn test_admits_requires_every_registered_filter() {
2215 let mut decoder = TychoStreamDecoder::<BlockHeader>::new(Chain::Ethereum);
2218 decoder.register_filter("x", rejects_a);
2219 decoder.register_filter("x", rejects_b);
2220
2221 assert!(!decoder.admits("x", &component_with_id("a")));
2222 assert!(!decoder.admits("x", &component_with_id("b")));
2223 assert!(decoder.admits("x", &component_with_id("c")));
2224 assert!(decoder.admits("y", &component_with_id("a")));
2225 }
2226}