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