1use std::{
6 collections::{HashMap, HashSet},
7 sync::LazyLock,
8 time::SystemTime,
9};
10
11use alloy::primitives::Address;
12use async_trait::async_trait;
13use futures::stream::BoxStream;
14use num_bigint::BigUint;
15use reqwest::Client;
16use tokio::time::{interval, timeout, Duration};
17use tracing::{error, info, warn};
18use tycho_common::{
19 models::{
20 protocol::{GetAmountOutParams, ProtocolComponent, ProtocolComponentState},
21 Chain,
22 },
23 simulation::indicatively_priced::SignedQuote,
24 Bytes,
25};
26
27use crate::{
28 rfq::{
29 client::RFQClient,
30 errors::RFQError,
31 models::TimestampHeader,
32 protocols::metric::models::{
33 MetricBidAskResponse, MetricMetadata, PaginatedMetadataResponse,
34 },
35 },
36 tycho_client::feed::synchronizer::{ComponentWithState, Snapshot, StateSyncMessage},
37};
38
39static METRIC_HTTP_CLIENT: LazyLock<Client> = LazyLock::new(Client::new);
40
41const METADATA_PAGE_SIZE: u32 = 500;
43
44#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
45pub struct MetricClient {
46 chain: Chain,
47 metadata_endpoint: String,
48 chain_endpoint: String,
50 tokens: HashSet<Bytes>,
51 tvl: f64,
52 #[serde(skip_serializing, default)]
53 api_key: Option<String>,
54 poll_time: Duration,
55 quote_timeout: Duration,
56 #[serde(default = "default_protocol_system")]
57 protocol_system: String,
58}
59
60fn default_protocol_system() -> String {
61 MetricClient::PROTOCOL_SYSTEM.to_string()
62}
63
64impl MetricClient {
65 pub const PROTOCOL_SYSTEM: &'static str = "rfq:metric";
66 pub const FALLBACK_PROTOCOL_SYSTEM: &'static str = "fallback:rfq:metric";
68
69 pub(super) fn via_fallback_router(mut self) -> Self {
70 self.protocol_system = Self::FALLBACK_PROTOCOL_SYSTEM.to_string();
71 self
72 }
73
74 pub fn new(
75 chain: Chain,
76 tokens: HashSet<Bytes>,
77 tvl: f64,
78 base_url: String,
79 api_key: Option<String>,
80 poll_time: Duration,
81 quote_timeout: Duration,
82 ) -> Result<Self, RFQError> {
83 let chain_id = chain_to_chain_id(chain)?;
84 let base_url = base_url.trim_end_matches('/');
85 let chain_endpoint = format!("{base_url}/public/v1/evm/{chain_id}");
86 Ok(Self {
87 chain,
88 metadata_endpoint: format!("{chain_endpoint}/metadata"),
89 chain_endpoint,
90 tokens,
91 tvl,
92 api_key,
93 poll_time,
94 quote_timeout,
95 protocol_system: Self::PROTOCOL_SYSTEM.to_string(),
96 })
97 }
98
99 fn http_client(&self) -> &Client {
100 &METRIC_HTTP_CLIENT
101 }
102
103 pub fn create_component_with_state(
104 &self,
105 component_id: String,
106 metadata: &MetricMetadata,
107 bid_ask: &MetricBidAskResponse,
108 tvl: f64,
109 ) -> ComponentWithState {
110 let protocol_component = ProtocolComponent {
111 id: component_id.clone(),
112 protocol_system: self.protocol_system.clone(),
113 protocol_type_name: "metric_pool".to_string(),
114 chain: self.chain,
115 tokens: vec![metadata.token0.clone(), metadata.token1.clone()],
116 contract_addresses: Vec::new(),
117 static_attributes: HashMap::new(),
118 ..Default::default()
119 };
120
121 let attributes = HashMap::from([
122 ("bid_adj".to_string(), Bytes::from(bid_ask.bid_adj.to_string().into_bytes())),
123 ("ask_adj".to_string(), Bytes::from(bid_ask.ask_adj.to_string().into_bytes())),
124 (
125 "total_token0_available".to_string(),
126 Bytes::from(
127 bid_ask
128 .total_token0_available
129 .as_ref()
130 .map(ToString::to_string)
131 .unwrap_or_default()
132 .into_bytes(),
133 ),
134 ),
135 (
136 "total_token1_available".to_string(),
137 Bytes::from(
138 bid_ask
139 .total_token1_available
140 .as_ref()
141 .map(ToString::to_string)
142 .unwrap_or_default()
143 .into_bytes(),
144 ),
145 ),
146 (
147 "server_ts".to_string(),
148 Bytes::from(
149 bid_ask
150 .server_ts
151 .to_string()
152 .into_bytes(),
153 ),
154 ),
155 (
156 "depth".to_string(),
157 Bytes::from(serde_json::to_vec(&bid_ask.depth).unwrap_or_default()),
158 ),
159 ]);
160
161 ComponentWithState {
162 state: ProtocolComponentState::new(&component_id, attributes, HashMap::new()),
163 component: protocol_component,
164 component_tvl: Some(tvl),
165 entrypoints: vec![],
166 }
167 }
168
169 async fn fetch_metadata(&self) -> Result<Vec<MetricMetadata>, RFQError> {
172 let mut pools = Vec::new();
173 let mut offset: u64 = 0;
174
175 loop {
176 let mut request = self
177 .http_client()
178 .get(&self.metadata_endpoint)
179 .header("accept", "application/json")
180 .query(&[
183 ("count", METADATA_PAGE_SIZE.to_string()),
184 ("offset", offset.to_string()),
185 ("include24h", "true".to_string()),
186 ]);
187
188 if let Some(api_key) = &self.api_key {
189 request = request.bearer_auth(api_key);
190 }
191
192 let response = request.send().await.map_err(|e| {
193 RFQError::ConnectionError(format!("Failed to fetch Metric metadata: {e}"))
194 })?;
195
196 if !response.status().is_success() {
197 return Err(RFQError::ConnectionError(format!(
198 "Metric metadata HTTP error {}: {}",
199 response.status(),
200 response
201 .text()
202 .await
203 .unwrap_or_default()
204 )));
205 }
206
207 let page: PaginatedMetadataResponse = response.json().await.map_err(|e| {
208 RFQError::ParsingError(format!("Failed to parse Metric metadata response: {e}"))
209 })?;
210
211 let page_len = page.data.len();
212 pools.extend(page.data);
213
214 match page.next_offset {
217 Some(next) if page_len > 0 && next > offset => offset = next,
218 _ => break,
219 }
220 }
221
222 Ok(pools)
223 }
224
225 async fn fetch_bid_ask(&self, pool: &Bytes) -> Result<MetricBidAskResponse, RFQError> {
226 let endpoint =
227 format!("{}/{}/bid_ask", self.chain_endpoint, bytes_to_address_string(pool)?);
228 let mut request = self
229 .http_client()
230 .get(endpoint)
231 .header("accept", "application/json");
232
233 if let Some(api_key) = &self.api_key {
234 request = request.bearer_auth(api_key);
235 }
236
237 let response = timeout(self.quote_timeout, request.send())
238 .await
239 .map_err(|_| {
240 RFQError::ConnectionError(format!(
241 "Metric bid/ask request timed out after {} seconds",
242 self.quote_timeout.as_secs()
243 ))
244 })?
245 .map_err(|e| {
246 RFQError::ConnectionError(format!("Failed to fetch Metric bid/ask: {e}"))
247 })?;
248
249 if !response.status().is_success() {
250 return Err(RFQError::ConnectionError(format!(
251 "Metric bid/ask HTTP error {}: {}",
252 response.status(),
253 response
254 .text()
255 .await
256 .unwrap_or_default()
257 )));
258 }
259
260 response.json().await.map_err(|e| {
261 RFQError::ParsingError(format!("Failed to parse Metric bid/ask response: {e}"))
262 })
263 }
264
265 fn find_pool<'a>(
266 &self,
267 metadata: &'a [MetricMetadata],
268 params: &GetAmountOutParams,
269 ) -> Result<&'a MetricMetadata, RFQError> {
270 metadata
271 .iter()
272 .find(|pool| {
273 (params.token_in == pool.token0 && params.token_out == pool.token1) ||
274 (params.token_in == pool.token1 && params.token_out == pool.token0)
275 })
276 .ok_or_else(|| {
277 RFQError::QuoteNotFound(format!(
278 "Metric pool not found for {} -> {}",
279 params.token_in, params.token_out
280 ))
281 })
282 }
283}
284
285#[async_trait]
286impl RFQClient for MetricClient {
287 fn stream(
288 &self,
289 ) -> BoxStream<'static, Result<(String, StateSyncMessage<TimestampHeader>), RFQError>> {
290 let client = self.clone();
291
292 Box::pin(async_stream::stream! {
293 let mut current_components: HashMap<String, ComponentWithState> = HashMap::new();
294 let mut ticker = interval(client.poll_time);
295
296 info!("Starting Metric polling every {} seconds", client.poll_time.as_secs());
297 loop {
298 ticker.tick().await;
299
300 let metadata = match client.fetch_metadata().await {
301 Ok(metadata) => metadata,
302 Err(e) => {
303 error!("Failed to fetch Metric metadata: {}", e);
304 continue;
305 }
306 };
307
308 let mut new_components = HashMap::new();
309 for pool in &metadata {
310 if !client.tokens.is_empty() &&
311 (!client.tokens.contains(&pool.token0) ||
312 !client.tokens.contains(&pool.token1))
313 {
314 continue;
315 }
316
317 let tvl = pool.tvl_fiat.unwrap_or(0.0);
320 if tvl < client.tvl {
321 continue;
322 }
323
324 let bid_ask = match client.fetch_bid_ask(&pool.pool_address).await {
325 Ok(bid_ask) => bid_ask,
326 Err(e) => {
327 warn!(
328 "Failed to fetch Metric bid/ask for pool {}: {}",
329 pool.pool_address, e
330 );
331 continue;
332 }
333 };
334 if !bid_ask.is_quotable() {
335 continue;
336 }
337
338 let component_id = pool.pool_address.to_string();
339 new_components.insert(
340 component_id.clone(),
341 client.create_component_with_state(component_id, pool, &bid_ask, tvl),
342 );
343 }
344
345 let removed_components: HashMap<String, ProtocolComponent> = current_components
346 .iter()
347 .filter(|(id, _)| !new_components.contains_key(*id))
348 .map(|(id, component)| (id.clone(), component.component.clone()))
349 .collect();
350
351 current_components = new_components.clone();
352 let timestamp = SystemTime::now()
353 .duration_since(SystemTime::UNIX_EPOCH)
354 .map_err(|_| RFQError::ParsingError("SystemTime before UNIX EPOCH".to_string()))?
355 .as_secs();
356
357 yield Ok(("metric".to_string(), StateSyncMessage {
358 header: TimestampHeader { timestamp },
359 snapshots: Snapshot { states: new_components, vm_storage: HashMap::new() },
360 deltas: None,
361 removed_components,
362 }));
363 }
364 })
365 }
366
367 async fn request_binding_quote(
368 &self,
369 params: &GetAmountOutParams,
370 ) -> Result<SignedQuote, RFQError> {
371 let metadata = self.fetch_metadata().await?;
372 self.find_pool(&metadata, params)?;
374
375 Ok(SignedQuote {
378 base_token: params.token_in.clone(),
379 quote_token: params.token_out.clone(),
380 amount_in: params.amount_in.clone(),
381 amount_out: BigUint::default(),
382 quote_attributes: HashMap::new(),
383 })
384 }
385}
386
387fn chain_to_chain_id(chain: Chain) -> Result<u64, RFQError> {
393 chain.try_id().map_err(|e| {
394 RFQError::FatalError(format!("Cannot resolve chain id for Metric on {chain}: {e}"))
395 })
396}
397
398fn bytes_to_address_string(address: &Bytes) -> Result<String, RFQError> {
399 if address.len() != 20 {
400 return Err(RFQError::InvalidInput(format!("Invalid EVM address length: {address}")));
401 }
402 Ok(Address::from_slice(address).to_checksum(None))
403}
404
405#[cfg(test)]
406mod tests {
407 use std::str::FromStr;
408
409 use rstest::rstest;
410
411 use super::*;
412 use crate::rfq::protocols::metric::{
413 client_builder::MetricClientBuilder,
414 models::{q64_to_f64, MetricDepth},
415 };
416
417 fn big(value: &str) -> BigUint {
418 value.parse().unwrap()
419 }
420
421 fn client() -> MetricClient {
422 MetricClient::new(
423 Chain::Ethereum,
424 HashSet::new(),
425 0.0,
426 "http://localhost:8080".to_string(),
427 None,
428 Duration::from_secs(1),
429 Duration::from_secs(1),
430 )
431 .unwrap()
432 }
433
434 fn live_client(chain: Chain) -> MetricClient {
435 let config = crate::rfq::constants::get_metric_config();
436 MetricClient::new(
437 chain,
438 HashSet::new(),
439 0.0,
440 config.base_url,
441 config.api_key,
442 Duration::from_secs(1),
443 Duration::from_secs(5),
444 )
445 .unwrap()
446 }
447
448 fn metadata() -> MetricMetadata {
449 MetricMetadata {
450 pool_address: Bytes::from_str("0xbF48bCf474d57fF82A3215319229e0DE1476A557").unwrap(),
451 token0: Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap(),
452 token1: Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap(),
453 tvl_fiat: Some(3000.0),
454 }
455 }
456
457 fn bid_ask() -> MetricBidAskResponse {
458 MetricBidAskResponse {
459 bid_adj: big("55340232221128654848000"),
460 ask_adj: big("55358678965202364400000"),
461 total_token0_available: Some(big("1000000000000000000")),
462 total_token1_available: Some(big("3000000000")),
463 server_ts: 1_770_053_095,
464 price_provider_status: Some("healthy".to_string()),
465 depth: MetricDepth::default(),
466 }
467 }
468
469 #[test]
470 fn test_chain_to_chain_id() {
471 assert_eq!(chain_to_chain_id(Chain::Ethereum).unwrap(), 1);
472 assert_eq!(chain_to_chain_id(Chain::Bsc).unwrap(), 56);
473 assert_eq!(chain_to_chain_id(Chain::Polygon).unwrap(), 137);
474 assert_eq!(chain_to_chain_id(Chain::Robinhood).unwrap(), 4663);
475 assert_eq!(chain_to_chain_id(Chain::Base).unwrap(), 8453);
476 assert_eq!(chain_to_chain_id(Chain::Arbitrum).unwrap(), 42161);
477 }
478
479 #[test]
480 fn test_endpoints_use_numeric_chain_id() {
481 let client = MetricClient::new(
482 Chain::Base,
483 HashSet::new(),
484 0.0,
485 "https://api.metric.xyz".to_string(),
486 None,
487 Duration::from_secs(1),
488 Duration::from_secs(1),
489 )
490 .unwrap();
491
492 assert_eq!(client.metadata_endpoint, "https://api.metric.xyz/public/v1/evm/8453/metadata");
493 assert_eq!(client.chain_endpoint, "https://api.metric.xyz/public/v1/evm/8453");
494 }
495
496 #[test]
497 fn test_builder_defaults_to_v1_base_url() {
498 let client = MetricClientBuilder::new(Chain::Ethereum)
499 .build()
500 .unwrap();
501
502 assert_eq!(client.metadata_endpoint, "https://api.metric.xyz/public/v1/evm/1/metadata");
503 }
504
505 #[test]
506 fn test_fallback_router_labels_components() {
507 let metadata = metadata();
508 let client = MetricClientBuilder::new(Chain::Ethereum)
509 .with_fallback_router()
510 .build()
511 .unwrap();
512
513 let component = client.create_component_with_state(
514 metadata.pool_address.to_string(),
515 &metadata,
516 &bid_ask(),
517 3000.0,
518 );
519
520 assert_eq!(component.component.protocol_system, MetricClient::FALLBACK_PROTOCOL_SYSTEM);
521 assert_eq!(
522 MetricClient::FALLBACK_PROTOCOL_SYSTEM,
523 tycho_execution::encoding::evm::METRIC_FALLBACK_PROTOCOL_SYSTEM
524 );
525 }
526
527 #[test]
528 fn test_component_attributes_round_trip_values() {
529 let metadata = metadata();
530 let component = client().create_component_with_state(
531 metadata.pool_address.to_string(),
532 &metadata,
533 &bid_ask(),
534 3000.0,
535 );
536
537 assert_eq!(component.component.protocol_system, MetricClient::PROTOCOL_SYSTEM);
538 assert_eq!(
539 component.component.tokens,
540 vec![metadata.token0.clone(), metadata.token1.clone()]
541 );
542 assert!(component
543 .component
544 .static_attributes
545 .is_empty());
546 assert_eq!(component.component.id, metadata.pool_address.to_string());
547 assert!(component
548 .component
549 .contract_addresses
550 .is_empty());
551 assert_eq!(
552 String::from_utf8(component.state.attributes["bid_adj"].to_vec()).unwrap(),
553 "55340232221128654848000"
554 );
555 assert_eq!(
556 String::from_utf8(component.state.attributes["server_ts"].to_vec()).unwrap(),
557 "1770053095"
558 );
559 }
560
561 #[rstest]
563 #[case::ethereum(Chain::Ethereum)]
564 #[case::bsc(Chain::Bsc)]
565 #[case::robinhood(Chain::Robinhood)]
566 #[case::base(Chain::Base)]
567 #[case::arbitrum(Chain::Arbitrum)]
568 #[tokio::test]
569 #[ignore = "hits Metric's public API"]
570 async fn test_live_metric_api_fetch_bid_ask_latest_fields(#[case] chain: Chain) {
571 let client = live_client(chain);
572 let metadata = client.fetch_metadata().await.unwrap();
573 assert!(!metadata.is_empty());
574
575 let mut last_error = None;
576 let mut selected = None;
577 for pool in &metadata {
578 match client
579 .fetch_bid_ask(&pool.pool_address)
580 .await
581 {
582 Ok(bid_ask) => {
583 if bid_ask.is_quotable() &&
584 !bid_ask.depth.asks.is_empty() &&
585 !bid_ask.depth.bids.is_empty()
586 {
587 selected = Some((pool, bid_ask));
588 break;
589 }
590 }
591 Err(error) => last_error = Some(error.to_string()),
592 }
593 }
594
595 let Some((_pool, bid_ask)) = selected else {
596 panic!(
597 "Metric live API on {chain} returned no quotable bid_ask response with ask and bid depth across {} pools; last error: {:?}",
598 metadata.len(),
599 last_error
600 );
601 };
602
603 let bid_price = bid_ask
604 .bid_price()
605 .unwrap()
606 .expect("the selected pool quotes a bid");
607 let ask_price = bid_ask
608 .ask_price()
609 .unwrap()
610 .expect("the selected pool quotes an ask");
611 assert!(bid_price.is_finite() && bid_price > 0.0);
612 assert!(ask_price.is_finite() && ask_price >= bid_price);
613 assert!(bid_ask.total_token0_available().is_ok());
614 assert!(bid_ask.total_token1_available().is_ok());
615 assert!(bid_ask.server_ts > 0);
616
617 for bin in bid_ask
618 .depth
619 .asks
620 .iter()
621 .chain(bid_ask.depth.bids.iter())
622 .take(6)
623 {
624 assert!(q64_to_f64(&bin.price)
625 .unwrap()
626 .is_finite());
627 assert!(bin.cumulative_input_volume > BigUint::ZERO);
630 }
631 }
632}