Skip to main content

tycho_simulation/rfq/protocols/metric/
client.rs

1//! Metric RFQ client for the `api.metric.xyz` API.
2//!
3//! Endpoint reference: <https://docs.metric.xyz/RSm94m71kqtGICv4iKRj/developers/api>
4
5use 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
41/// Page size for the paginated metadata endpoint. The API clamps `count` to `[1, 500]`.
42const 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    // Prefix ending at /public/v1/evm/{chain_id}; pool-specific endpoints are derived from it.
49    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    /// Fetches every configured pool by paging through the metadata endpoint until the API reports
156    /// no next page.
157    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                // `include24h=true` is required for the top-level `tvlFiat` field; without it the
167                // API omits TVL and every pool would fall below any non-zero threshold.
168                .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            // Stop when the API reports the last page, returns nothing, or fails to advance the
201            // offset (defensive guard against an infinite loop).
202            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                    // v1 metadata carries the fiat TVL directly, so no cross-pool price
304                    // normalization is needed.
305                    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        // Validates that a pool exists for the requested pair.
359        self.find_pool(&metadata, params)?;
360
361        // The v1 heartbeat updates the oracle on-chain every block, so no signed oracle-update args
362        // are relayed with the swap. The binding quote therefore carries no quote attributes.
363        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
373/// Resolves the EVM chain id Metric expects in its `/public/v1/evm/{chain_id}` paths.
374///
375/// Metric also serves Monad (143), HyperEVM (999), MegaETH (4326) and Avalanche (43114), but Tycho
376/// has no built-in `Chain` variant for them yet, so they are unsupported unless registered as a
377/// `Chain::Custom` with the matching chain id.
378fn 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    // Polygon is omitted: Metric lists it in `/public/v1/chains` but publishes no pools there yet.
526    #[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            // Deserialization already parsed the volumes; assert the input-driven depth walk's
592            // key field is populated in live responses.
593            assert!(bin.cumulative_input_volume > BigUint::ZERO);
594        }
595    }
596}