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