1use std::{str::FromStr, sync::Arc, time::Duration};
14
15use num_cpus;
16use rustc_hash::FxHashSet;
17use serde::{Deserialize, Serialize};
18use tokio::{
19 sync::broadcast,
20 task::{AbortHandle, JoinHandle},
21};
22use tycho_execution::encoding::evm::swap_encoder::swap_encoder_registry::SwapEncoderRegistry;
23#[cfg(feature = "experimental")]
24use tycho_simulation::evm::stream::BlockStepController;
25#[cfg(feature = "test-utils")]
26use tycho_simulation::tycho_ethereum::gas::{BlockGasPrice, GasPrice};
27use tycho_simulation::{
28 evm::pending::PendingBlockProcessor,
29 tycho_common::{models::Chain, traits::TxDeltaIndexer, Bytes},
30 tycho_core::models::Address,
31 tycho_ethereum::rpc::EthereumRpcClient,
32};
33
34use crate::{
35 algorithm::{AlgorithmConfig, AlgorithmError},
36 derived::{ComputationManager, ComputationManagerConfig, SharedDerivedDataRef},
37 encoding::{encoder::Encoder, fee_fetcher::RouterFeeFetcher, router_fees::SharedRouterFees},
38 feed::{
39 events::{MarketEvent, MarketEventHandler},
40 gas::GasPriceFetcher,
41 market_data::MarketData,
42 metrics_sampler::MetricsSampler,
43 tycho_feed::TychoFeed,
44 TychoFeedConfig,
45 },
46 graph::EdgeWeightUpdaterWithDerived,
47 price_guard::{
48 guard::PriceGuard, provider::PriceProvider, provider_registry::PriceProviderRegistry,
49 },
50 propamm_fallback::{
51 fee_tier_fetcher::FeeTierFetcher, SharedFeeTiers, PROPAMM_ROUTER_ADDRESS, PROPAMM_VENUES,
52 },
53 types::constants::native_token,
54 worker_pool::{
55 pool::{WorkerPool, WorkerPoolBuilder},
56 registry::UnknownAlgorithmError,
57 },
58 worker_pool_router::{
59 config::WorkerPoolRouterConfig, ExclusiveAccess, LiquidityScope, SolverPoolHandle,
60 WorkerPoolRouter,
61 },
62 Algorithm, Quote, QuoteRequest, SolveError,
63};
64
65pub mod defaults {
71 use std::time::Duration;
72
73 pub const MIN_TOKEN_QUALITY: i32 = 100;
75 pub const TRADED_N_DAYS_AGO: u64 = 3;
77 pub const TVL_BUFFER_RATIO: f64 = 1.1;
80 pub const GAS_REFRESH_INTERVAL: Duration = Duration::from_secs(30);
82 pub const METRICS_SAMPLE_INTERVAL: Duration = Duration::from_secs(10);
84 pub const ROUTER_FEE_REFRESH_INTERVAL: Duration = Duration::from_secs(300);
86 pub const FALLBACK_FEE_TIER_REFRESH_INTERVAL: Duration = Duration::from_secs(600);
89 pub const RECONNECT_DELAY: Duration = Duration::from_secs(5);
91 pub const ROUTER_MIN_RESPONSES: usize = 0;
94 pub const POOL_TASK_QUEUE_CAPACITY: usize = 1000;
96 pub const POOL_MIN_HOPS: usize = 1;
98 pub const POOL_MAX_HOPS: usize = 3;
100 pub const POOL_TIMEOUT_MS: u64 = 100;
102}
103
104const DEFAULT_TYCHO_USE_TLS: bool = true;
106const DEFAULT_DEPTH_SLIPPAGE_THRESHOLD: f64 = 0.01;
107const DEFAULT_ROUTER_TIMEOUT: Duration = Duration::from_secs(10);
110
111fn default_task_queue_capacity() -> usize {
114 defaults::POOL_TASK_QUEUE_CAPACITY
115}
116
117fn default_min_hops() -> usize {
118 defaults::POOL_MIN_HOPS
119}
120
121fn default_max_hops() -> usize {
122 defaults::POOL_MAX_HOPS
123}
124
125fn default_algo_timeout_ms() -> u64 {
126 defaults::POOL_TIMEOUT_MS
127}
128
129fn parse_connector_tokens(
130 raw: Option<&[String]>,
131) -> Result<Option<FxHashSet<Address>>, SolverBuildError> {
132 let Some(strings) = raw else {
133 return Ok(None);
134 };
135 let mut set = FxHashSet::with_capacity_and_hasher(strings.len(), Default::default());
136 for s in strings {
137 let addr = Address::from_str(s).map_err(|e| AlgorithmError::InvalidConfiguration {
138 reason: format!("connector_tokens: invalid address {s:?}: {e}"),
139 })?;
140 set.insert(addr);
141 }
142 Ok(Some(set))
143}
144
145#[must_use]
147#[derive(Debug, Clone, Serialize, Deserialize)]
148pub struct PoolConfig {
149 algorithm: String,
151 #[serde(default = "num_cpus::get")]
153 num_workers: usize,
154 #[serde(default = "default_task_queue_capacity")]
156 task_queue_capacity: usize,
157 #[serde(default = "default_min_hops")]
159 min_hops: usize,
160 #[serde(default = "default_max_hops")]
162 max_hops: usize,
163 #[serde(default = "default_algo_timeout_ms")]
165 timeout_ms: u64,
166 #[serde(default)]
168 max_routes: Option<usize>,
169 #[serde(default)]
172 connector_tokens: Option<Vec<String>>,
173 #[serde(default)]
175 liquidity_scope: Option<LiquidityScope>,
176}
177
178impl PoolConfig {
179 pub fn new(algorithm: impl Into<String>) -> Self {
182 Self {
183 algorithm: algorithm.into(),
184 num_workers: num_cpus::get(),
185 task_queue_capacity: defaults::POOL_TASK_QUEUE_CAPACITY,
186 min_hops: defaults::POOL_MIN_HOPS,
187 max_hops: defaults::POOL_MAX_HOPS,
188 timeout_ms: defaults::POOL_TIMEOUT_MS,
189 max_routes: None,
190 connector_tokens: None,
191 liquidity_scope: None,
192 }
193 }
194
195 pub fn algorithm(&self) -> &str {
197 &self.algorithm
198 }
199
200 pub fn liquidity_scope(&self) -> Option<LiquidityScope> {
202 self.liquidity_scope
203 }
204
205 pub fn with_liquidity_scope(mut self, scope: LiquidityScope) -> Self {
207 self.liquidity_scope = Some(scope);
208 self
209 }
210
211 pub fn num_workers(&self) -> usize {
213 self.num_workers
214 }
215
216 pub fn with_num_workers(mut self, num_workers: usize) -> Self {
218 self.num_workers = num_workers;
219 self
220 }
221
222 pub fn with_task_queue_capacity(mut self, task_queue_capacity: usize) -> Self {
224 self.task_queue_capacity = task_queue_capacity;
225 self
226 }
227
228 pub fn with_min_hops(mut self, min_hops: usize) -> Self {
230 self.min_hops = min_hops;
231 self
232 }
233
234 pub fn with_max_hops(mut self, max_hops: usize) -> Self {
236 self.max_hops = max_hops;
237 self
238 }
239
240 pub fn with_timeout_ms(mut self, timeout_ms: u64) -> Self {
242 self.timeout_ms = timeout_ms;
243 self
244 }
245
246 pub fn with_max_routes(mut self, max_routes: Option<usize>) -> Self {
248 self.max_routes = max_routes;
249 self
250 }
251
252 pub fn task_queue_capacity(&self) -> usize {
254 self.task_queue_capacity
255 }
256
257 pub fn min_hops(&self) -> usize {
259 self.min_hops
260 }
261
262 pub fn max_hops(&self) -> usize {
264 self.max_hops
265 }
266
267 pub fn timeout_ms(&self) -> u64 {
269 self.timeout_ms
270 }
271
272 pub fn max_routes(&self) -> Option<usize> {
274 self.max_routes
275 }
276
277 pub fn with_connector_tokens(mut self, tokens: Vec<String>) -> Self {
280 self.connector_tokens = Some(tokens);
281 self
282 }
283
284 pub fn connector_tokens(&self) -> Option<&[String]> {
286 self.connector_tokens.as_deref()
287 }
288}
289
290#[derive(Debug, thiserror::Error)]
292#[error("timed out after {timeout_ms}ms waiting for market data and derived computations")]
293pub struct WaitReadyError {
294 timeout_ms: u64,
295}
296
297#[non_exhaustive]
299#[derive(Debug, thiserror::Error)]
300pub enum SolverBuildError {
301 #[error("failed to create ethereum RPC client: {0}")]
303 RpcClient(String),
304 #[error(transparent)]
306 AlgorithmConfig(#[from] AlgorithmError),
307 #[error("failed to create computation manager: {0}")]
309 ComputationManager(String),
310 #[error("failed to create encoder: {0}")]
312 Encoder(String),
313 #[error("failed to create router fee fetcher: {0}")]
315 RouterFeeFetcher(String),
316 #[error("failed to create fallback fee tier fetcher: {0}")]
319 FeeTierFetcher(String),
320 #[error(transparent)]
322 UnknownAlgorithm(#[from] UnknownAlgorithmError),
323 #[error("gas token not configured for chain")]
325 GasToken,
326 #[error("no worker pools configured")]
328 NoPools,
329 #[error(
333 "every worker pool sets liquidity_scope = \"include_exclusive\"; requests without \
334 exclusive access would be served by no pool. Configure at least one public_only pool"
335 )]
336 NoPublicPool,
337 #[cfg(feature = "test-utils")]
339 #[error("replay failed: {0}")]
340 Replay(String),
341 #[error("feed setup failed before delivering pending processor: {0}")]
346 FeedSetup(String),
347 #[error("pending processor channel closed before processor was delivered")]
350 PendingChannelClosed,
351 #[cfg(feature = "experimental")]
354 #[error("step controller channel closed before controller was delivered")]
355 StepControllerChannelClosed,
356}
357
358enum PoolEntry {
360 BuiltIn {
361 name: String,
362 algorithm: String,
363 num_workers: usize,
364 task_queue_capacity: usize,
365 min_hops: usize,
366 max_hops: usize,
367 timeout_ms: u64,
368 max_routes: Option<usize>,
369 connector_tokens: Option<FxHashSet<Address>>,
370 liquidity_scope: Option<LiquidityScope>,
371 },
372 Custom(CustomPoolEntry),
373}
374
375impl PoolEntry {
376 fn liquidity_scope(&self) -> Option<LiquidityScope> {
378 match self {
379 PoolEntry::BuiltIn { liquidity_scope, .. } => *liquidity_scope,
380 PoolEntry::Custom(custom) => custom.liquidity_scope,
381 }
382 }
383}
384
385struct CustomPoolEntry {
387 name: String,
388 num_workers: usize,
389 task_queue_capacity: usize,
390 min_hops: usize,
391 max_hops: usize,
392 timeout_ms: u64,
393 max_routes: Option<usize>,
394 liquidity_scope: Option<LiquidityScope>,
395 configure: Box<dyn FnOnce(WorkerPoolBuilder) -> WorkerPoolBuilder + Send>,
397}
398
399struct BuiltComponents {
402 tycho_feed: TychoFeed,
403 gas_price_fetcher: GasPriceFetcher<EthereumRpcClient>,
404 router_fee_fetcher: Option<RouterFeeFetcher>,
405 fee_tier_fetcher: Option<FeeTierFetcher>,
406 computation_manager: ComputationManager,
407 computation_event_rx: broadcast::Receiver<MarketEvent>,
408 computation_shutdown_tx: broadcast::Sender<()>,
409 computation_shutdown_rx: broadcast::Receiver<()>,
410 router: WorkerPoolRouter,
411 worker_pools: Vec<WorkerPool>,
412 market_data: MarketData,
413 derived_data: SharedDerivedDataRef,
414 router_fees: SharedRouterFees,
415 chain: Chain,
416 router_address: Option<Bytes>,
417 pending_indexers: Vec<(String, Box<dyn TxDeltaIndexer>)>,
418 market_event_tx: broadcast::Sender<MarketEvent>,
419}
420
421#[must_use = "a builder does nothing until .build() is called"]
426pub struct FyndBuilder {
427 chain: Chain,
428 tycho_url: String,
429 rpc_url: String,
430 protocols: Vec<String>,
431 min_tvl: f64,
432 tycho_api_key: Option<String>,
433 tycho_use_tls: bool,
434 min_token_quality: i32,
435 traded_n_days_ago: u64,
436 tvl_buffer_ratio: f64,
437 gas_refresh_interval: Duration,
438 reconnect_delay: Duration,
439 blocklisted_components: FxHashSet<String>,
440 partial_blocks: bool,
441 router_timeout: Duration,
442 router_min_responses: usize,
443 encoder: Option<Encoder>,
444 calldata_watermark: Option<Vec<u8>>,
445 pools: Vec<PoolEntry>,
446 price_guard_enabled: bool,
447 price_providers: Vec<Box<dyn PriceProvider>>,
448 pending_indexers: Vec<(String, Box<dyn TxDeltaIndexer>)>,
449}
450
451impl FyndBuilder {
452 pub fn new(
454 chain: Chain,
455 tycho_url: impl Into<String>,
456 rpc_url: impl Into<String>,
457 protocols: Vec<String>,
458 min_tvl: f64,
459 ) -> Self {
460 Self {
461 chain,
462 tycho_url: tycho_url.into(),
463 rpc_url: rpc_url.into(),
464 protocols,
465 min_tvl,
466 tycho_api_key: None,
467 tycho_use_tls: DEFAULT_TYCHO_USE_TLS,
468 min_token_quality: defaults::MIN_TOKEN_QUALITY,
469 traded_n_days_ago: defaults::TRADED_N_DAYS_AGO,
470 tvl_buffer_ratio: defaults::TVL_BUFFER_RATIO,
471 gas_refresh_interval: defaults::GAS_REFRESH_INTERVAL,
472 reconnect_delay: defaults::RECONNECT_DELAY,
473 blocklisted_components: FxHashSet::default(),
474 partial_blocks: false,
475 router_timeout: DEFAULT_ROUTER_TIMEOUT,
476 router_min_responses: defaults::ROUTER_MIN_RESPONSES,
477 encoder: None,
478 calldata_watermark: None,
479 pools: Vec::new(),
480 price_guard_enabled: false,
481 price_providers: Vec::new(),
482 pending_indexers: Vec::new(),
483 }
484 }
485
486 pub fn chain(&self) -> Chain {
488 self.chain
489 }
490
491 pub fn tycho_api_key(mut self, key: impl Into<String>) -> Self {
493 self.tycho_api_key = Some(key.into());
494 self
495 }
496
497 pub fn min_tvl(mut self, min_tvl: f64) -> Self {
499 self.min_tvl = min_tvl;
500 self
501 }
502
503 pub fn tycho_use_tls(mut self, use_tls: bool) -> Self {
505 self.tycho_use_tls = use_tls;
506 self
507 }
508
509 pub fn min_token_quality(mut self, quality: i32) -> Self {
512 self.min_token_quality = quality;
513 self
514 }
515
516 pub fn traded_n_days_ago(mut self, days: u64) -> Self {
518 self.traded_n_days_ago = days;
519 self
520 }
521
522 pub fn tvl_buffer_ratio(mut self, ratio: f64) -> Self {
524 self.tvl_buffer_ratio = ratio;
525 self
526 }
527
528 pub fn gas_refresh_interval(mut self, interval: Duration) -> Self {
530 self.gas_refresh_interval = interval;
531 self
532 }
533
534 pub fn reconnect_delay(mut self, delay: Duration) -> Self {
536 self.reconnect_delay = delay;
537 self
538 }
539
540 pub fn blocklisted_components(mut self, components: impl IntoIterator<Item = String>) -> Self {
542 self.blocklisted_components = components.into_iter().collect();
543 self
544 }
545
546 pub fn partial_blocks(mut self, enabled: bool) -> Self {
552 self.partial_blocks = enabled;
553 self
554 }
555
556 pub fn worker_router_timeout(mut self, timeout: Duration) -> Self {
558 self.router_timeout = timeout;
559 self
560 }
561
562 pub fn worker_router_min_responses(mut self, min: usize) -> Self {
564 self.router_min_responses = min;
565 self
566 }
567
568 pub fn encoder(mut self, encoder: Encoder) -> Self {
570 self.encoder = Some(encoder);
571 self
572 }
573
574 pub fn calldata_watermark(mut self, watermark: impl Into<Vec<u8>>) -> Self {
578 self.calldata_watermark = Some(watermark.into());
579 self
580 }
581
582 pub fn algorithm(mut self, algorithm: impl Into<String>) -> Self {
584 self.pools.push(PoolEntry::BuiltIn {
585 name: "default".to_string(),
586 algorithm: algorithm.into(),
587 num_workers: num_cpus::get(),
588 task_queue_capacity: defaults::POOL_TASK_QUEUE_CAPACITY,
589 min_hops: defaults::POOL_MIN_HOPS,
590 max_hops: defaults::POOL_MAX_HOPS,
591 timeout_ms: defaults::POOL_TIMEOUT_MS,
592 max_routes: None,
593 connector_tokens: None,
594 liquidity_scope: None,
595 });
596 self
597 }
598
599 pub fn with_algorithm<A, F>(mut self, name: impl Into<String>, factory: F) -> Self
603 where
604 A: Algorithm + 'static,
605 A::GraphManager: MarketEventHandler + EdgeWeightUpdaterWithDerived + 'static,
606 F: Fn(AlgorithmConfig) -> A + Clone + Send + Sync + 'static,
607 {
608 let name = name.into();
609 let algo_name = name.clone();
610 let configure =
611 Box::new(move |builder: WorkerPoolBuilder| builder.with_algorithm(algo_name, factory));
612 self.pools
613 .push(PoolEntry::Custom(CustomPoolEntry {
614 name,
615 num_workers: num_cpus::get(),
616 task_queue_capacity: defaults::POOL_TASK_QUEUE_CAPACITY,
617 min_hops: defaults::POOL_MIN_HOPS,
618 max_hops: defaults::POOL_MAX_HOPS,
619 timeout_ms: defaults::POOL_TIMEOUT_MS,
620 max_routes: None,
621 liquidity_scope: None,
622 configure,
623 }));
624 self
625 }
626
627 pub fn add_default_price_providers(self) -> Self {
634 self.register_price_provider(Box::new(
635 crate::price_guard::hyperliquid::HyperliquidProvider::default(),
636 ))
637 .register_price_provider(Box::new(
638 crate::price_guard::binance_ws::BinanceWsProvider::default(),
639 ))
640 }
641
642 pub fn register_price_provider(mut self, provider: Box<dyn PriceProvider>) -> Self {
647 self.price_providers.push(provider);
648 self
649 }
650
651 pub fn with_pending_indexer(
657 mut self,
658 extractor: impl Into<String>,
659 indexer: Box<dyn TxDeltaIndexer>,
660 ) -> Self {
661 self.pending_indexers
662 .push((extractor.into(), indexer));
663 self
664 }
665
666 pub fn price_guard_enabled(mut self, enabled: bool) -> Self {
673 self.price_guard_enabled = enabled;
674 self
675 }
676
677 pub fn add_pool(
684 mut self,
685 name: impl Into<String>,
686 config: &PoolConfig,
687 ) -> Result<Self, SolverBuildError> {
688 let connector_tokens = parse_connector_tokens(config.connector_tokens())?;
689 self.pools.push(PoolEntry::BuiltIn {
690 name: name.into(),
691 algorithm: config.algorithm().to_string(),
692 num_workers: config.num_workers(),
693 task_queue_capacity: config.task_queue_capacity(),
694 min_hops: config.min_hops(),
695 max_hops: config.max_hops(),
696 timeout_ms: config.timeout_ms(),
697 max_routes: config.max_routes(),
698 connector_tokens,
699 liquidity_scope: config.liquidity_scope(),
700 });
701 Ok(self)
702 }
703
704 fn assemble_components(mut self) -> Result<BuiltComponents, SolverBuildError> {
707 if self.pools.is_empty() {
708 return Err(SolverBuildError::NoPools);
709 }
710
711 if self
715 .pools
716 .iter()
717 .all(|p| p.liquidity_scope() == Some(LiquidityScope::IncludeExclusive))
718 {
719 return Err(SolverBuildError::NoPublicPool);
720 }
721
722 if self.price_providers.is_empty() {
724 self = self.add_default_price_providers();
725 }
726
727 let market_data = MarketData::new_shared();
728
729 let tycho_feed_config = TychoFeedConfig::new(
730 self.tycho_url,
731 self.chain,
732 self.tycho_api_key,
733 self.tycho_use_tls,
734 self.protocols,
735 self.min_tvl,
736 )
737 .tvl_buffer_ratio(self.tvl_buffer_ratio)
738 .reconnect_delay(self.reconnect_delay)
739 .min_token_quality(self.min_token_quality)
740 .traded_n_days_ago(self.traded_n_days_ago)
741 .blocklisted_components(self.blocklisted_components)
742 .partial_blocks(self.partial_blocks);
743
744 let ethereum_client = EthereumRpcClient::new(self.rpc_url.as_str())
745 .map_err(|e| SolverBuildError::RpcClient(e.to_string()))?;
746
747 let gas_price_fetcher =
748 GasPriceFetcher::new(ethereum_client, market_data.clone(), self.gas_refresh_interval);
749
750 let tycho_feed = TychoFeed::new(tycho_feed_config, market_data.clone());
751 let market_event_tx = tycho_feed.event_sender();
752
753 let gas_token = native_token(&self.chain).map_err(|_| SolverBuildError::GasToken)?;
754 let computation_config = ComputationManagerConfig::new()
755 .with_gas_token(gas_token)
756 .with_depth_slippage_threshold(DEFAULT_DEPTH_SLIPPAGE_THRESHOLD);
757 let (computation_manager, _) =
760 ComputationManager::new(computation_config, market_data.clone())
761 .map_err(|e| SolverBuildError::ComputationManager(e.to_string()))?;
762
763 let derived_data: SharedDerivedDataRef = computation_manager.store();
764 let derived_event_tx = computation_manager.event_sender();
765
766 let computation_event_rx = tycho_feed.subscribe();
769 let (computation_shutdown_tx, computation_shutdown_rx) = broadcast::channel(1);
770
771 let mut solver_pool_handles: Vec<SolverPoolHandle> = Vec::new();
772 let mut worker_pools: Vec<WorkerPool> = Vec::new();
773 let fallback_fee_tiers = SharedFeeTiers::default();
775
776 let pools = std::mem::take(&mut self.pools);
777
778 for pool_entry in pools {
779 let pool_event_rx = tycho_feed.subscribe();
780 let derived_rx = derived_event_tx.subscribe();
781
782 let pool_scope = pool_entry
783 .liquidity_scope()
784 .unwrap_or_default();
785
786 let (worker_pool, task_handle) = match pool_entry {
787 PoolEntry::BuiltIn {
788 name,
789 algorithm,
790 num_workers,
791 task_queue_capacity,
792 min_hops,
793 max_hops,
794 timeout_ms,
795 max_routes,
796 connector_tokens,
797 liquidity_scope: _,
798 } => {
799 let mut algo_cfg = AlgorithmConfig::new(
800 min_hops,
801 max_hops,
802 Duration::from_millis(timeout_ms),
803 max_routes,
804 )?;
805 if let Some(tokens) = connector_tokens {
806 algo_cfg = algo_cfg.with_connector_tokens(tokens);
807 }
808 let builder = WorkerPoolBuilder::new()
809 .name(name)
810 .algorithm(algorithm)
811 .algorithm_config(algo_cfg)
812 .num_workers(num_workers)
813 .task_queue_capacity(task_queue_capacity)
814 .liquidity_scope(pool_scope)
815 .fallback_fee_tiers(fallback_fee_tiers.clone());
816 builder.build(
817 market_data.clone(),
818 Arc::clone(&derived_data),
819 pool_event_rx,
820 derived_rx,
821 )?
822 }
823 PoolEntry::Custom(custom) => {
824 let algo_cfg = AlgorithmConfig::new(
825 custom.min_hops,
826 custom.max_hops,
827 Duration::from_millis(custom.timeout_ms),
828 custom.max_routes,
829 )?;
830 let builder = WorkerPoolBuilder::new()
831 .name(custom.name)
832 .algorithm_config(algo_cfg)
833 .num_workers(custom.num_workers)
834 .task_queue_capacity(custom.task_queue_capacity)
835 .liquidity_scope(pool_scope)
836 .fallback_fee_tiers(fallback_fee_tiers.clone());
837 let builder = (custom.configure)(builder);
838 builder.build(
839 market_data.clone(),
840 Arc::clone(&derived_data),
841 pool_event_rx,
842 derived_rx,
843 )?
844 }
845 };
846
847 solver_pool_handles.push(
848 SolverPoolHandle::new(worker_pool.name(), task_handle)
849 .with_liquidity_scope(pool_scope),
850 );
851 worker_pools.push(worker_pool);
852 }
853
854 let encoder = match self.encoder {
855 Some(enc) => enc,
856 None => {
857 let registry = SwapEncoderRegistry::new(self.chain)
858 .add_default_encoders(None)
859 .map_err(|e| SolverBuildError::Encoder(e.to_string()))?;
860 Encoder::new(self.chain, registry)
861 .map_err(|e| SolverBuildError::Encoder(e.to_string()))?
862 }
863 };
864 let encoder = match self.calldata_watermark {
865 Some(watermark) => encoder.with_calldata_watermark(watermark),
866 None => encoder,
867 };
868
869 let chain = self.chain;
870 let router_address = encoder.router_address().cloned();
871 let router_fees = encoder.router_fees();
872
873 let router_fee_fetcher = match &router_address {
874 Some(addr) => Some(
875 RouterFeeFetcher::new(
876 self.rpc_url.as_str(),
877 addr,
878 router_fees.clone(),
879 defaults::ROUTER_FEE_REFRESH_INTERVAL,
880 )
881 .map_err(|e| SolverBuildError::RouterFeeFetcher(e.to_string()))?,
882 ),
883 None => {
884 tracing::warn!(
885 %chain,
886 "no Tycho router for this chain; running quote-only (encoding disabled)"
887 );
888 None
889 }
890 };
891
892 let fee_tier_fetcher = if chain == Chain::Ethereum {
898 let router = Bytes::from_str(PROPAMM_ROUTER_ADDRESS).map_err(|e| {
899 SolverBuildError::FeeTierFetcher(format!(
900 "PropAMMRouter address {PROPAMM_ROUTER_ADDRESS}: {e}"
901 ))
902 })?;
903 let venues = PROPAMM_VENUES
904 .iter()
905 .map(|venue| {
906 Bytes::from_str(venue).map_err(|e| {
907 SolverBuildError::FeeTierFetcher(format!("pAMM venue {venue}: {e}"))
908 })
909 })
910 .collect::<Result<Vec<Bytes>, _>>()?;
911 Some(
912 FeeTierFetcher::new(
913 self.rpc_url.as_str(),
914 &router,
915 &venues,
916 fallback_fee_tiers.clone(),
917 defaults::FALLBACK_FEE_TIER_REFRESH_INTERVAL,
918 )
919 .map_err(|e| SolverBuildError::FeeTierFetcher(e.to_string()))?,
920 )
921 } else {
922 None
923 };
924
925 let router_config = WorkerPoolRouterConfig::default()
928 .with_timeout(self.router_timeout)
929 .with_min_responses(self.router_min_responses);
930 let mut router = WorkerPoolRouter::new(solver_pool_handles, router_config, encoder);
931
932 if self.price_guard_enabled {
933 let mut registry = PriceProviderRegistry::new();
934 let mut worker_handles = Vec::new();
935 for mut provider in self.price_providers {
936 worker_handles.push(provider.start(market_data.clone()));
937 registry = registry.register(provider);
938 }
939 let price_guard = PriceGuard::new(registry, worker_handles);
940 router = router.with_price_guard(price_guard);
941 }
942
943 Ok(BuiltComponents {
944 tycho_feed,
945 gas_price_fetcher,
946 router_fee_fetcher,
947 fee_tier_fetcher,
948 computation_manager,
949 computation_event_rx,
950 computation_shutdown_tx,
951 computation_shutdown_rx,
952 router,
953 worker_pools,
954 market_data,
955 derived_data,
956 router_fees,
957 chain,
958 router_address,
959 pending_indexers: self.pending_indexers,
960 market_event_tx,
961 })
962 }
963
964 pub fn build(self) -> Result<Solver, SolverBuildError> {
970 let mut c = self.assemble_components()?;
971
972 let feed_handle = tokio::spawn(async move {
973 if let Err(e) = c.tycho_feed.run().await {
974 metrics::counter!("tycho_feed_failures_total").increment(1);
975 tracing::error!(error = %e, "tycho feed error");
976 }
977 });
978 let gas_price_handle = tokio::spawn(async move {
979 c.gas_price_fetcher.run().await;
980 });
981 let metrics_sampler =
982 MetricsSampler::new(c.market_data.clone(), defaults::METRICS_SAMPLE_INTERVAL);
983 let metrics_sampler_handle = tokio::spawn(async move { metrics_sampler.run().await });
984 let router_fee_handle = match c.router_fee_fetcher {
985 Some(fetcher) => tokio::spawn(async move { fetcher.run().await }),
986 None => tokio::spawn(async {}),
987 };
988 let fee_tier_handle = match c.fee_tier_fetcher {
989 Some(fetcher) => tokio::spawn(async move { fetcher.run().await }),
990 None => tokio::spawn(async {}),
991 };
992 let computation_handle = tokio::spawn(async move {
993 c.computation_manager
994 .run(c.computation_event_rx, c.computation_shutdown_rx)
995 .await;
996 });
997
998 Ok(Solver {
999 router: c.router,
1000 worker_pools: c.worker_pools,
1001 market_data: c.market_data,
1002 derived_data: c.derived_data,
1003 router_fees: c.router_fees,
1004 feed_handle,
1005 gas_price_handle,
1006 metrics_sampler_handle,
1007 router_fee_handle,
1008 fee_tier_handle,
1009 computation_handle,
1010 computation_shutdown_tx: c.computation_shutdown_tx,
1011 chain: c.chain,
1012 router_address: c.router_address,
1013 market_event_tx: c.market_event_tx,
1014 })
1015 }
1016
1017 pub async fn build_with_pending(
1029 self,
1030 ) -> Result<(Solver, PendingBlockProcessor), SolverBuildError> {
1031 let mut c = self.assemble_components()?;
1032
1033 let (pending_tx, pending_rx) =
1034 tokio::sync::oneshot::channel::<Result<PendingBlockProcessor, String>>();
1035
1036 let pending_indexers = c.pending_indexers;
1037 let feed_handle = tokio::spawn(async move {
1038 if let Err(e) = c
1039 .tycho_feed
1040 .run_with_pending(pending_tx, pending_indexers)
1041 .await
1042 {
1043 metrics::counter!("tycho_feed_failures_total").increment(1);
1044 tracing::error!(error = %e, "tycho feed error");
1045 }
1046 });
1047 let gas_price_handle = tokio::spawn(async move {
1048 c.gas_price_fetcher.run().await;
1049 });
1050 let metrics_sampler =
1051 MetricsSampler::new(c.market_data.clone(), defaults::METRICS_SAMPLE_INTERVAL);
1052 let metrics_sampler_handle = tokio::spawn(async move { metrics_sampler.run().await });
1053 let router_fee_handle = match c.router_fee_fetcher {
1054 Some(fetcher) => tokio::spawn(async move { fetcher.run().await }),
1055 None => tokio::spawn(async {}),
1056 };
1057 let fee_tier_handle = match c.fee_tier_fetcher {
1058 Some(fetcher) => tokio::spawn(async move { fetcher.run().await }),
1059 None => tokio::spawn(async {}),
1060 };
1061 let computation_handle = tokio::spawn(async move {
1062 c.computation_manager
1063 .run(c.computation_event_rx, c.computation_shutdown_rx)
1064 .await;
1065 });
1066
1067 let pending = pending_rx
1068 .await
1069 .map_err(|_| SolverBuildError::PendingChannelClosed)?
1070 .map_err(SolverBuildError::FeedSetup)?;
1071
1072 Ok((
1073 Solver {
1074 router: c.router,
1075 worker_pools: c.worker_pools,
1076 market_data: c.market_data,
1077 derived_data: c.derived_data,
1078 router_fees: c.router_fees,
1079 feed_handle,
1080 gas_price_handle,
1081 metrics_sampler_handle,
1082 router_fee_handle,
1083 fee_tier_handle,
1084 computation_handle,
1085 computation_shutdown_tx: c.computation_shutdown_tx,
1086 chain: c.chain,
1087 router_address: c.router_address,
1088 market_event_tx: c.market_event_tx,
1089 },
1090 pending,
1091 ))
1092 }
1093
1094 #[cfg(feature = "experimental")]
1109 pub async fn build_with_step_controller(
1110 self,
1111 ) -> Result<(Solver, BlockStepController), SolverBuildError> {
1112 let mut c = self.assemble_components()?;
1113
1114 let (controller_tx, controller_rx) =
1115 tokio::sync::oneshot::channel::<Result<BlockStepController, String>>();
1116
1117 let feed_handle = tokio::spawn(async move {
1118 if let Err(e) = c
1119 .tycho_feed
1120 .run_with_step_controller(controller_tx)
1121 .await
1122 {
1123 tracing::error!(error = %e, "tycho feed error");
1124 }
1125 });
1126 let gas_price_handle = tokio::spawn(async move {
1127 c.gas_price_fetcher.run().await;
1128 });
1129 let metrics_sampler =
1130 MetricsSampler::new(c.market_data.clone(), defaults::METRICS_SAMPLE_INTERVAL);
1131 let metrics_sampler_handle = tokio::spawn(async move { metrics_sampler.run().await });
1132 let router_fee_handle = match c.router_fee_fetcher {
1133 Some(fetcher) => tokio::spawn(async move { fetcher.run().await }),
1134 None => tokio::spawn(async {}),
1135 };
1136 let fee_tier_handle = match c.fee_tier_fetcher {
1137 Some(fetcher) => tokio::spawn(async move { fetcher.run().await }),
1138 None => tokio::spawn(async {}),
1139 };
1140 let computation_handle = tokio::spawn(async move {
1141 c.computation_manager
1142 .run(c.computation_event_rx, c.computation_shutdown_rx)
1143 .await;
1144 });
1145
1146 let controller = controller_rx
1147 .await
1148 .map_err(|_| SolverBuildError::StepControllerChannelClosed)?
1149 .map_err(SolverBuildError::FeedSetup)?;
1150
1151 Ok((
1152 Solver {
1153 router: c.router,
1154 worker_pools: c.worker_pools,
1155 market_data: c.market_data,
1156 derived_data: c.derived_data,
1157 router_fees: c.router_fees,
1158 feed_handle,
1159 gas_price_handle,
1160 metrics_sampler_handle,
1161 router_fee_handle,
1162 fee_tier_handle,
1163 computation_handle,
1164 computation_shutdown_tx: c.computation_shutdown_tx,
1165 chain: c.chain,
1166 router_address: c.router_address,
1167 market_event_tx: c.market_event_tx,
1168 },
1169 controller,
1170 ))
1171 }
1172} pub struct Solver {
1176 router: WorkerPoolRouter,
1177 worker_pools: Vec<WorkerPool>,
1178 market_data: MarketData,
1179 derived_data: SharedDerivedDataRef,
1180 router_fees: SharedRouterFees,
1181 feed_handle: JoinHandle<()>,
1182 gas_price_handle: JoinHandle<()>,
1183 metrics_sampler_handle: JoinHandle<()>,
1184 router_fee_handle: JoinHandle<()>,
1185 fee_tier_handle: JoinHandle<()>,
1186 computation_handle: JoinHandle<()>,
1187 computation_shutdown_tx: broadcast::Sender<()>,
1188 chain: Chain,
1189 router_address: Option<Bytes>,
1190 market_event_tx: broadcast::Sender<MarketEvent>,
1191}
1192
1193impl Solver {
1194 pub fn market_data(&self) -> MarketData {
1196 self.market_data.clone()
1197 }
1198
1199 pub fn router_address(&self) -> Option<&Bytes> {
1201 self.router_address.as_ref()
1202 }
1203
1204 pub fn derived_data(&self) -> SharedDerivedDataRef {
1206 Arc::clone(&self.derived_data)
1207 }
1208
1209 pub fn subscribe_market_events(&self) -> broadcast::Receiver<crate::feed::events::MarketEvent> {
1214 self.market_event_tx.subscribe()
1215 }
1216
1217 pub async fn quote(&self, request: QuoteRequest) -> Result<Quote, SolveError> {
1227 self.router
1228 .quote(request, ExclusiveAccess::Granted)
1229 .await
1230 }
1231
1232 pub async fn wait_until_ready(&self, timeout: Duration) -> Result<(), WaitReadyError> {
1249 const POLL_INTERVAL: Duration = Duration::from_millis(500);
1250
1251 let deadline = tokio::time::Instant::now() + timeout;
1252
1253 loop {
1254 let market_ready = self
1255 .market_data
1256 .read()
1257 .await
1258 .last_updated()
1259 .is_some();
1260 let derived_ready = self
1261 .derived_data
1262 .read()
1263 .await
1264 .derived_data_ready();
1265
1266 if market_ready && derived_ready {
1267 return Ok(());
1268 }
1269
1270 if tokio::time::Instant::now() >= deadline {
1271 return Err(WaitReadyError { timeout_ms: timeout.as_millis() as u64 });
1272 }
1273
1274 tokio::time::sleep(POLL_INTERVAL).await;
1275 }
1276 }
1277
1278 #[cfg(feature = "test-utils")]
1291 pub async fn from_recording(
1292 chain: Chain,
1293 updates: Vec<tycho_simulation::protocol::models::Update>,
1294 pools: std::collections::HashMap<String, PoolConfig>,
1295 gas_price_wei: Option<num_bigint::BigUint>,
1296 ) -> Result<Self, SolverBuildError> {
1297 if pools.is_empty() {
1298 return Err(SolverBuildError::NoPools);
1299 }
1300
1301 let market_data = MarketData::new_shared();
1302
1303 let feed_config =
1305 TychoFeedConfig::new("ws://replay".to_string(), chain, None, false, vec![], 0.0);
1306 let feed = TychoFeed::new(feed_config, market_data.clone());
1307 let market_event_tx = feed.event_sender();
1308 let _feed_rx = feed.subscribe();
1309
1310 for update in updates {
1311 feed.handle_tycho_message(update)
1312 .await
1313 .map_err(|e| SolverBuildError::Replay(e.to_string()))?;
1314 }
1315
1316 let gas_price = match gas_price_wei {
1318 Some(price) => price,
1319 None => {
1320 tracing::warn!("no recorded gas price, defaulting to 10 gwei");
1321 num_bigint::BigUint::from(10_000_000_000u64)
1322 }
1323 };
1324 let block_number = match market_data.read().await.last_updated() {
1325 Some(block) => block.number(),
1326 None => {
1327 tracing::warn!("no block number from replayed updates, defaulting to 0");
1328 0
1329 }
1330 };
1331 {
1332 let mut market = market_data.write().await;
1333 market.update_gas_price(BlockGasPrice {
1334 block_number,
1335 block_hash: Default::default(),
1336 block_timestamp: 0,
1337 pricing: GasPrice::Legacy { gas_price },
1338 });
1339 }
1340
1341 let gas_token = native_token(&chain).map_err(|_| SolverBuildError::GasToken)?;
1343 let computation_config = ComputationManagerConfig::new()
1344 .with_gas_token(gas_token)
1345 .with_depth_slippage_threshold(DEFAULT_DEPTH_SLIPPAGE_THRESHOLD);
1346 let (computation_manager, _) =
1347 ComputationManager::new(computation_config, market_data.clone())
1348 .map_err(|e| SolverBuildError::ComputationManager(e.to_string()))?;
1349
1350 let derived_data: SharedDerivedDataRef = computation_manager.store();
1351 let derived_event_tx = computation_manager.event_sender();
1352
1353 let computation_event_rx = feed.subscribe();
1354 let (computation_shutdown_tx, computation_shutdown_rx) = broadcast::channel(1);
1355
1356 let computation_handle = tokio::spawn(async move {
1357 computation_manager
1358 .run(computation_event_rx, computation_shutdown_rx)
1359 .await;
1360 });
1361
1362 let mut solver_pool_handles: Vec<SolverPoolHandle> = Vec::new();
1364 let mut worker_pools: Vec<WorkerPool> = Vec::new();
1365 let mut max_timeout_ms = 0u64;
1366
1367 for (name, pool_cfg) in &pools {
1368 let algo_cfg = AlgorithmConfig::new(
1369 pool_cfg.min_hops(),
1370 pool_cfg.max_hops(),
1371 Duration::from_millis(pool_cfg.timeout_ms()),
1372 pool_cfg.max_routes(),
1373 )?;
1374
1375 let pool_event_rx = feed.subscribe();
1376 let derived_rx = derived_event_tx.subscribe();
1377
1378 let (worker_pool, task_handle) = WorkerPoolBuilder::new()
1379 .name(name.clone())
1380 .algorithm(pool_cfg.algorithm().to_string())
1381 .algorithm_config(algo_cfg)
1382 .num_workers(pool_cfg.num_workers())
1383 .task_queue_capacity(pool_cfg.task_queue_capacity())
1384 .build(market_data.clone(), Arc::clone(&derived_data), pool_event_rx, derived_rx)?;
1385
1386 solver_pool_handles.push(SolverPoolHandle::new(worker_pool.name(), task_handle));
1387 max_timeout_ms = max_timeout_ms.max(pool_cfg.timeout_ms());
1388 worker_pools.push(worker_pool);
1389 }
1390
1391 let encoder = {
1393 let registry = SwapEncoderRegistry::new(chain)
1394 .add_default_encoders(None)
1395 .map_err(|e| SolverBuildError::Encoder(e.to_string()))?;
1396 Encoder::new(chain, registry).map_err(|e| SolverBuildError::Encoder(e.to_string()))?
1397 };
1398
1399 let router_address = encoder.router_address().cloned();
1400 let router_fees = encoder.router_fees();
1404 router_fees.set(crate::encoding::router_fees::RouterFees::new(
1405 100_000_000,
1406 0,
1407 0,
1408 rustc_hash::FxHashMap::default(),
1409 ));
1410 let router_config = WorkerPoolRouterConfig::default()
1411 .with_timeout(Duration::from_millis(max_timeout_ms.max(5000)))
1412 .with_min_responses(defaults::ROUTER_MIN_RESPONSES);
1413 let router = WorkerPoolRouter::new(solver_pool_handles, router_config, encoder);
1414
1415 let market_read = market_data.read().await;
1417 let added = market_read.component_topology();
1418 drop(market_read);
1419
1420 if market_event_tx
1421 .send(MarketEvent::MarketUpdated {
1422 added_components: added,
1423 removed_components: vec![],
1424 updated_components: vec![],
1425 })
1426 .is_err()
1427 {
1428 tracing::warn!("no receivers for initial MarketUpdated broadcast");
1429 }
1430
1431 let feed_handle = tokio::spawn(futures::future::pending::<()>());
1434 let gas_price_handle = tokio::spawn(async { });
1435 let metrics_sampler_handle = tokio::spawn(async { });
1436 let router_fee_handle = tokio::spawn(async { });
1437 let fee_tier_handle = tokio::spawn(async { });
1438
1439 Ok(Solver {
1440 router,
1441 worker_pools,
1442 market_data,
1443 derived_data,
1444 router_fees,
1445 feed_handle,
1446 gas_price_handle,
1447 metrics_sampler_handle,
1448 router_fee_handle,
1449 fee_tier_handle,
1450 computation_handle,
1451 computation_shutdown_tx,
1452 chain,
1453 router_address,
1454 market_event_tx,
1455 })
1456 }
1457
1458 pub fn shutdown(self) {
1460 let _ = self.computation_shutdown_tx.send(());
1461 for pool in self.worker_pools {
1462 pool.shutdown();
1463 }
1464 self.feed_handle.abort();
1465 self.gas_price_handle.abort();
1466 self.metrics_sampler_handle.abort();
1467 self.router_fee_handle.abort();
1468 self.fee_tier_handle.abort();
1469 }
1470
1471 pub fn into_parts(self) -> SolverParts {
1473 SolverParts {
1474 router: self.router,
1475 worker_pools: self.worker_pools,
1476 market_data: self.market_data,
1477 derived_data: self.derived_data,
1478 router_fees: self.router_fees,
1479 feed_handle: self.feed_handle,
1480 gas_price_handle: self.gas_price_handle,
1481 metrics_sampler_handle: self.metrics_sampler_handle,
1482 router_fee_handle: self.router_fee_handle,
1483 fee_tier_handle: self.fee_tier_handle,
1484 computation_handle: self.computation_handle,
1485 computation_shutdown_tx: self.computation_shutdown_tx,
1486 chain: self.chain,
1487 router_address: self.router_address,
1488 }
1489 }
1490}
1491
1492pub struct SolverParts {
1496 router: WorkerPoolRouter,
1498 worker_pools: Vec<WorkerPool>,
1500 market_data: MarketData,
1502 derived_data: SharedDerivedDataRef,
1504 router_fees: SharedRouterFees,
1506 feed_handle: JoinHandle<()>,
1508 gas_price_handle: JoinHandle<()>,
1510 metrics_sampler_handle: JoinHandle<()>,
1512 router_fee_handle: JoinHandle<()>,
1514 fee_tier_handle: JoinHandle<()>,
1519 computation_handle: JoinHandle<()>,
1521 computation_shutdown_tx: broadcast::Sender<()>,
1523 chain: Chain,
1525 router_address: Option<Bytes>,
1527}
1528
1529impl SolverParts {
1530 pub fn chain(&self) -> Chain {
1532 self.chain
1533 }
1534
1535 pub fn router_address(&self) -> Option<&Bytes> {
1537 self.router_address.as_ref()
1538 }
1539
1540 pub fn worker_pools(&self) -> &[WorkerPool] {
1542 &self.worker_pools
1543 }
1544
1545 pub fn market_data(&self) -> &MarketData {
1547 &self.market_data
1548 }
1549
1550 pub fn derived_data(&self) -> &SharedDerivedDataRef {
1552 &self.derived_data
1553 }
1554
1555 pub fn router_fees(&self) -> &SharedRouterFees {
1557 &self.router_fees
1558 }
1559
1560 pub fn fee_tier_abort_handle(&self) -> AbortHandle {
1566 self.fee_tier_handle.abort_handle()
1567 }
1568
1569 pub fn into_router(self) -> WorkerPoolRouter {
1571 self.router
1572 }
1573
1574 #[allow(clippy::type_complexity)]
1576 pub fn into_components(
1577 self,
1578 ) -> (
1579 WorkerPoolRouter,
1580 Vec<WorkerPool>,
1581 MarketData,
1582 SharedDerivedDataRef,
1583 JoinHandle<()>,
1584 JoinHandle<()>,
1585 JoinHandle<()>,
1586 JoinHandle<()>,
1587 JoinHandle<()>,
1588 broadcast::Sender<()>,
1589 ) {
1590 (
1591 self.router,
1592 self.worker_pools,
1593 self.market_data,
1594 self.derived_data,
1595 self.feed_handle,
1596 self.gas_price_handle,
1597 self.metrics_sampler_handle,
1598 self.router_fee_handle,
1599 self.computation_handle,
1600 self.computation_shutdown_tx,
1601 )
1602 }
1603}
1604
1605#[cfg(test)]
1606mod tests {
1607 use super::*;
1608
1609 #[test]
1612 fn test_unscoped_pool_resolves_to_public_only() {
1613 let config = PoolConfig::new("most_liquid");
1614 assert_eq!(config.liquidity_scope(), None);
1615 assert_eq!(
1616 config
1617 .liquidity_scope()
1618 .unwrap_or_default(),
1619 LiquidityScope::PublicOnly
1620 );
1621 }
1622
1623 #[test]
1626 fn test_build_all_exclusive_pools() {
1627 let config =
1628 PoolConfig::new("most_liquid").with_liquidity_scope(LiquidityScope::IncludeExclusive);
1629 let result = FyndBuilder::new(
1630 Chain::Ethereum,
1631 "wss://example.invalid",
1632 "https://example.invalid",
1633 vec!["uniswap_v2".to_string()],
1634 100.0,
1635 )
1636 .add_pool("exclusive", &config)
1637 .expect("add_pool should accept the config")
1638 .build();
1639
1640 assert!(matches!(result, Err(SolverBuildError::NoPublicPool)));
1641 }
1642}