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    #[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    /// Components executed through Tycho's `MetricFallbackRouter`.
67    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    /// Fetches every configured pool by paging through the metadata endpoint until the API reports
170    /// no next page.
171    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                // `include24h=true` is required for the top-level `tvlFiat` field; without it the
181                // API omits TVL and every pool would fall below any non-zero threshold.
182                .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            // Stop when the API reports the last page, returns nothing, or fails to advance the
215            // offset (defensive guard against an infinite loop).
216            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                    // v1 metadata carries the fiat TVL directly, so no cross-pool price
318                    // normalization is needed.
319                    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        // Validates that a pool exists for the requested pair.
373        self.find_pool(&metadata, params)?;
374
375        // The v1 heartbeat updates the oracle on-chain every block, so no signed oracle-update args
376        // are relayed with the swap. The binding quote therefore carries no quote attributes.
377        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
387/// Resolves the EVM chain id Metric expects in its `/public/v1/evm/{chain_id}` paths.
388///
389/// Metric also serves Monad (143), HyperEVM (999), MegaETH (4326) and Avalanche (43114), but Tycho
390/// has no built-in `Chain` variant for them yet, so they are unsupported unless registered as a
391/// `Chain::Custom` with the matching chain id.
392fn 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    // Polygon is omitted: Metric lists it in `/public/v1/chains` but publishes no pools there yet.
562    #[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            // Deserialization already parsed the volumes; assert the input-driven depth walk's
628            // key field is populated in live responses.
629            assert!(bin.cumulative_input_volume > BigUint::ZERO);
630        }
631    }
632}