Skip to main content

tycho_simulation/rfq/protocols/hashflow/
client.rs

1use std::{
2    collections::{HashMap, HashSet},
3    str::FromStr,
4};
5
6use alloy::primitives::{utils::keccak256, Address, U256};
7use async_trait::async_trait;
8use futures::stream::BoxStream;
9use num_bigint::BigUint;
10use reqwest::Client;
11use serde::{Deserialize, Serialize};
12use tokio::time::{interval, timeout, Duration};
13use tracing::{error, info, warn};
14use tycho_common::{
15    models::{protocol::GetAmountOutParams, Chain},
16    simulation::indicatively_priced::SignedQuote,
17    Bytes,
18};
19
20use crate::{
21    evm::protocol::u256_num::biguint_to_u256,
22    rfq::{
23        client::RFQClient,
24        errors::RFQError,
25        models::{ComponentLayout, QuoteRule, TimestampHeader},
26        protocols::{
27            component,
28            hashflow::models::{
29                HashflowChain, HashflowMarketMakerLevels, HashflowMarketMakersResponse,
30                HashflowPriceLevelsResponse, HashflowQuoteRequest, HashflowQuoteResponse,
31                HashflowRFQ, HashflowRFQOptions,
32            },
33            maker_price_levels::{self, MakerLevelsWithTvl, MakerPriceLevels},
34        },
35    },
36    tycho_client::feed::synchronizer::{ComponentWithState, StateSyncMessage},
37    tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState},
38};
39
40#[derive(Clone, Debug, Serialize, Deserialize)]
41pub struct HashflowClient {
42    chain: Chain,
43    price_levels_endpoint: String,
44    market_makers_endpoint: String,
45    quote_endpoint: String,
46    // Tokens that we want prices for
47    tokens: HashSet<Bytes>,
48    // Min tvl value in the quote token.
49    tvl: f64,
50    #[serde(skip_serializing, default)]
51    auth_key: String,
52    #[serde(skip_serializing, default)]
53    auth_user: String,
54    // Quote tokens to normalize to for TVL purposes. Should have the same prices.
55    quote_tokens: HashSet<Bytes>,
56    poll_time: Duration,
57    quote_timeout: Duration,
58    #[serde(default = "default_protocol_system")]
59    protocol_system: String,
60    #[serde(default)]
61    component_layout: ComponentLayout,
62    /// How often one route may take quotes from Hashflow. Read only by the all-pairs layout.
63    #[serde(default = "HashflowClient::default_quote_rule")]
64    quote_rule: QuoteRule,
65}
66
67fn default_protocol_system() -> String {
68    HashflowClient::PROTOCOL_SYSTEM.to_string()
69}
70
71impl HashflowClient {
72    pub const PROTOCOL_SYSTEM: &'static str = "rfq:hashflow";
73    /// Components executed through Tycho's `HashflowFallbackRouter`.
74    pub const FALLBACK_PROTOCOL_SYSTEM: &'static str = "fallback:rfq:hashflow";
75    pub const DEFAULT_QUOTE_RULE: QuoteRule = QuoteRule::OncePerMaker;
76
77    pub(super) fn via_fallback_router(mut self) -> Self {
78        self.protocol_system = Self::FALLBACK_PROTOCOL_SYSTEM.to_string();
79        self
80    }
81
82    #[allow(clippy::too_many_arguments)]
83    pub fn new(
84        chain: Chain,
85        tokens: HashSet<Bytes>,
86        tvl: f64,
87        quote_tokens: HashSet<Bytes>,
88        auth_user: String,
89        auth_key: String,
90        poll_time: Duration,
91        quote_timeout: Duration,
92    ) -> Result<Self, RFQError> {
93        Ok(Self {
94            chain,
95            price_levels_endpoint: "https://api.hashflow.com/taker/v3/price-levels".to_string(),
96            market_makers_endpoint: "https://api.hashflow.com/taker/v3/market-makers".to_string(),
97            quote_endpoint: "https://api.hashflow.com/taker/v3/rfq".to_string(),
98            tokens,
99            tvl,
100            auth_key,
101            auth_user,
102            quote_tokens,
103            poll_time,
104            quote_timeout,
105            protocol_system: Self::PROTOCOL_SYSTEM.to_string(),
106            component_layout: ComponentLayout::PerPair,
107            quote_rule: Self::DEFAULT_QUOTE_RULE,
108        })
109    }
110
111    fn default_quote_rule() -> QuoteRule {
112        Self::DEFAULT_QUOTE_RULE
113    }
114
115    pub(crate) fn with_component_layout(mut self, component_layout: ComponentLayout) -> Self {
116        self.component_layout = component_layout;
117        self
118    }
119
120    pub(crate) fn with_quote_rule(mut self, quote_rule: QuoteRule) -> Self {
121        self.quote_rule = quote_rule;
122        self
123    }
124
125    /// The name the client's stream tags its messages with. The two layouts use different names,
126    /// so one stream builder can carry a client of each.
127    fn stream_name(&self) -> &'static str {
128        match self.component_layout {
129            ComponentLayout::PerPair => "hashflow",
130            ComponentLayout::AllPairs => "hashflow_all_pairs",
131        }
132    }
133
134    pub fn quote_rule(&self) -> QuoteRule {
135        self.quote_rule
136    }
137
138    #[cfg(test)]
139    pub(crate) fn with_quote_endpoint(mut self, quote_endpoint: String) -> Self {
140        self.quote_endpoint = quote_endpoint;
141        self
142    }
143
144    /// Normalize TVL to a common quote token for comparison
145    /// Returns the normalized TVL value, or 0.0 if normalization fails due to no liquidity
146    fn normalize_tvl(
147        &self,
148        raw_tvl: f64,
149        quote_token: Bytes,
150        levels_by_mm: &HashMap<String, Vec<HashflowMarketMakerLevels>>,
151    ) -> Result<f64, RFQError> {
152        // If the quote token is already in our approved quote token set, no conversion needed
153        if self.quote_tokens.contains(&quote_token) {
154            return Ok(raw_tvl);
155        }
156
157        // Try to find the price of the quote token in one of the approved quote tokens
158        // for normalization.
159        for approved_quote_token in &self.quote_tokens {
160            for mm_levels_inner in levels_by_mm.values() {
161                for quote_mm_level in mm_levels_inner {
162                    // Check for direct pair: quote_token/approved_quote_token
163                    if quote_mm_level.pair.base_token == quote_token &&
164                        quote_mm_level.pair.quote_token == *approved_quote_token
165                    {
166                        if let Some(price) = quote_mm_level.get_price(1.0) {
167                            return Ok(raw_tvl * price);
168                        }
169                    }
170                }
171            }
172        }
173
174        // If we can't normalize, return TVL 0 (pool will be filtered out)
175        Ok(0.0)
176    }
177
178    fn create_component_with_state(
179        &self,
180        component_id: String,
181        tokens: Vec<Bytes>,
182        mm_name: &str,
183        mm_level: &HashflowMarketMakerLevels,
184        tvl: f64,
185    ) -> ComponentWithState {
186        let protocol_component = ProtocolComponent {
187            id: component_id.clone(),
188            protocol_system: self.protocol_system.clone(),
189            protocol_type_name: "hashflow_pool".to_string(),
190            chain: self.chain,
191            tokens,
192            contract_addresses: vec![], // empty for RFQ
193            ..Default::default()
194        };
195
196        let mut attributes = HashMap::new();
197
198        // Store price levels as JSON string
199        if !mm_level.levels.is_empty() {
200            let levels_json = serde_json::to_string(&mm_level.levels).unwrap_or_default();
201            attributes.insert("levels".to_string(), levels_json.as_bytes().to_vec().into());
202        }
203        attributes.insert("mm".to_string(), mm_name.as_bytes().to_vec().into());
204
205        ComponentWithState {
206            state: ProtocolComponentState::new(&component_id, attributes, HashMap::new()),
207            component: protocol_component,
208            component_tvl: Some(tvl),
209            entrypoints: vec![],
210        }
211    }
212
213    /// Every market maker's levels on a pair whose tokens the client requested and whose TVL
214    /// clears the threshold, each with its normalized TVL.
215    fn levels_above_tvl_threshold<'a>(
216        &self,
217        levels_by_mm: &'a HashMap<String, Vec<HashflowMarketMakerLevels>>,
218    ) -> Result<Vec<MakerLevelsWithTvl<'a, HashflowMarketMakerLevels>>, RFQError> {
219        let mut kept = Vec::new();
220        for (mm_name, mm_levels) in levels_by_mm.iter() {
221            for mm_level in mm_levels {
222                let base_token = &mm_level.pair.base_token;
223                let quote_token = &mm_level.pair.quote_token;
224                if !(self.tokens.contains(base_token) && self.tokens.contains(quote_token)) {
225                    continue;
226                }
227                let normalized_tvl = self.normalize_tvl(
228                    mm_level.calculate_tvl(),
229                    mm_level.pair.quote_token.clone(),
230                    levels_by_mm,
231                )?;
232                if normalized_tvl < self.tvl {
233                    info!(
234                        "Filtering out component {} due to low TVL: {:.2} < {:.2}",
235                        pair_component_id(mm_level),
236                        normalized_tvl,
237                        self.tvl
238                    );
239                    continue;
240                }
241                kept.push(MakerLevelsWithTvl {
242                    market_maker: mm_name,
243                    levels: mm_level,
244                    tvl: normalized_tvl,
245                });
246            }
247        }
248        Ok(kept)
249    }
250
251    /// One component per pair in `levels_above_tvl_threshold`. When several market makers quote a
252    /// pair, the component holds one of them.
253    fn pair_components(
254        &self,
255        levels_by_mm: &HashMap<String, Vec<HashflowMarketMakerLevels>>,
256    ) -> Result<HashMap<String, ComponentWithState>, RFQError> {
257        let mut new_components = HashMap::new();
258        for MakerLevelsWithTvl { market_maker: mm_name, levels: mm_level, tvl } in
259            self.levels_above_tvl_threshold(levels_by_mm)?
260        {
261            let component_id = pair_component_id(mm_level);
262            let tokens = vec![mm_level.pair.base_token.clone(), mm_level.pair.quote_token.clone()];
263            let component_with_state = self.create_component_with_state(
264                component_id.clone(),
265                tokens,
266                mm_name,
267                mm_level,
268                tvl,
269            );
270            new_components.insert(component_id, component_with_state);
271        }
272        Ok(new_components)
273    }
274
275    /// The all-pairs component for one poll: every maker's price levels in
276    /// `levels_above_tvl_threshold`. `None` when there is none.
277    fn all_pairs_component(
278        &self,
279        levels_by_mm: &HashMap<String, Vec<HashflowMarketMakerLevels>>,
280    ) -> Result<Option<ComponentWithState>, RFQError> {
281        let mut price_levels = Vec::new();
282        for MakerLevelsWithTvl { market_maker: mm_name, levels: mm_level, tvl } in
283            self.levels_above_tvl_threshold(levels_by_mm)?
284        {
285            let maker_levels = MakerPriceLevels {
286                market_maker: mm_name.to_string(),
287                base_token: mm_level.pair.base_token.clone(),
288                quote_token: mm_level.pair.quote_token.clone(),
289                levels: mm_level.levels.clone(),
290            };
291            price_levels.push((maker_levels, tvl));
292        }
293        maker_price_levels::all_pairs_component(
294            Self::PROTOCOL_SYSTEM,
295            "hashflow_pool",
296            self.chain,
297            price_levels,
298            self.quote_rule,
299        )
300    }
301
302    /// Requests a firm quote. With `market_maker`, Hashflow is asked for that maker alone with
303    /// `doNotRetryWithOtherMakers`, so a declined request is an error and not a quote from a
304    /// maker the caller did not choose.
305    pub(crate) async fn request_quote(
306        &self,
307        params: &GetAmountOutParams,
308        market_maker: Option<&str>,
309    ) -> Result<SignedQuote, RFQError> {
310        let hashflow_chain = HashflowChain::from(self.chain);
311        // A fresh random address becomes the quote's effectiveTrader — the address Hashflow
312        // scopes its strictly increasing quote nonces to — so quotes never invalidate each
313        // other, at the cost of a cold nonce storage slot on Hashflow's router (~17k gas per
314        // swap). The receiver executes the trade on-chain, so it is Hashflow's trader.
315        let effective_trader = Bytes::from(Address::random().to_vec());
316        let quote_request = HashflowQuoteRequest {
317            source: self.auth_user.clone(),
318            base_chain: hashflow_chain.clone(),
319            quote_chain: hashflow_chain,
320            rfqs: vec![HashflowRFQ {
321                base_token: params.token_in.to_string(),
322                quote_token: params.token_out.to_string(),
323                base_token_amount: Some(params.amount_in.to_string()),
324                quote_token_amount: None,
325                trader: params.receiver.to_string(),
326                effective_trader: Some(effective_trader.to_string()),
327                market_makers: market_maker.map(|mm| vec![mm.to_string()]),
328                options: market_maker
329                    .map(|_| HashflowRFQOptions { do_not_retry_with_other_makers: true }),
330            }],
331            calldata: false,
332        };
333        self.send_quote_request(params, &quote_request, &effective_trader)
334            .await
335    }
336
337    async fn fetch_market_makers(&mut self) -> Result<Vec<String>, RFQError> {
338        let query_params = vec![
339            ("source", self.auth_user.clone()),
340            ("baseChainType", "evm".to_string()),
341            ("baseChainId", self.chain.id().to_string()),
342        ];
343
344        let http_client = Client::new();
345        let request = http_client
346            .get(&self.market_makers_endpoint)
347            .query(&query_params)
348            .header("accept", "application/json")
349            .header("Authorization", &self.auth_key);
350
351        let response = request.send().await.map_err(|e| {
352            RFQError::ConnectionError(format!("Failed to fetch market makers: {e}"))
353        })?;
354
355        if !response.status().is_success() {
356            return Err(RFQError::ConnectionError(format!(
357                "HTTP error {}: {}",
358                response.status(),
359                response
360                    .text()
361                    .await
362                    .unwrap_or_default()
363            )));
364        }
365
366        let mm_response: HashflowMarketMakersResponse = response.json().await.map_err(|e| {
367            RFQError::ParsingError(format!("Failed to parse market makers response: {e}"))
368        })?;
369
370        info!(
371            "Fetched {} market makers: {:?}",
372            mm_response.market_makers.len(),
373            mm_response.market_makers
374        );
375
376        Ok(mm_response.market_makers)
377    }
378
379    async fn fetch_price_levels(
380        &self,
381        market_makers: &Vec<String>,
382    ) -> Result<HashMap<String, Vec<HashflowMarketMakerLevels>>, RFQError> {
383        let mut query_params = vec![
384            ("source", self.auth_user.clone()),
385            ("baseChainType", "evm".to_string()),
386            ("baseChainId", self.chain.id().to_string()),
387        ];
388
389        // Add market makers as array parameters
390        for mm in market_makers {
391            query_params.push(("marketMakers[]", mm.clone()));
392        }
393
394        let http_client = Client::new();
395        let request = http_client
396            .get(&self.price_levels_endpoint)
397            .query(&query_params)
398            .header("accept", "application/json")
399            .header("Authorization", &self.auth_key);
400
401        let response = request
402            .send()
403            .await
404            .map_err(|e| RFQError::ConnectionError(format!("Failed to fetch price levels: {e}")))?;
405
406        if !response.status().is_success() {
407            return Err(RFQError::ConnectionError(format!(
408                "HTTP error {}: {}",
409                response.status(),
410                response
411                    .text()
412                    .await
413                    .unwrap_or_default()
414            )));
415        }
416
417        let price_response: HashflowPriceLevelsResponse = response.json().await.map_err(|e| {
418            RFQError::ParsingError(format!("Failed to parse price levels response: {e}"))
419        })?;
420
421        if price_response.status != "success" {
422            let error = match price_response.error {
423                Some(error) => error.to_string(),
424                None => "no error details".to_string(),
425            };
426            return Err(RFQError::InvalidInput(format!("API returned error status: {error}")));
427        }
428
429        price_response
430            .levels
431            .ok_or_else(|| RFQError::ParsingError("API response missing levels".to_string()))
432    }
433
434    async fn send_quote_request(
435        &self,
436        params: &GetAmountOutParams,
437        quote_request: &HashflowQuoteRequest,
438        effective_trader: &Bytes,
439    ) -> Result<SignedQuote, RFQError> {
440        let url = self.quote_endpoint.clone();
441
442        let start_time = std::time::Instant::now();
443        const MAX_RETRIES: u32 = 3;
444        let mut last_error = None;
445
446        for attempt in 0..MAX_RETRIES {
447            // Check if we have time remaining for this attempt
448            let elapsed = start_time.elapsed();
449            if elapsed >= self.quote_timeout {
450                return Err(last_error.unwrap_or_else(|| {
451                    RFQError::ConnectionError(format!(
452                        "Hashflow quote request timed out after {} seconds",
453                        self.quote_timeout.as_secs()
454                    ))
455                }));
456            }
457
458            let remaining_time = self.quote_timeout - elapsed;
459
460            let http_client = Client::new();
461            let request = http_client
462                .post(&url)
463                .json(&quote_request)
464                .header("accept", "application/json")
465                .header("Authorization", &self.auth_key);
466
467            let response = match timeout(remaining_time, request.send()).await {
468                Ok(Ok(resp)) => resp,
469                Ok(Err(e)) => {
470                    warn!(
471                        "Hashflow quote request failed (attempt {}/{}): {}",
472                        attempt + 1,
473                        MAX_RETRIES,
474                        e
475                    );
476                    last_error = Some(RFQError::ConnectionError(format!(
477                        "Failed to send Hashflow quote request: {e}"
478                    )));
479                    if attempt < MAX_RETRIES - 1 {
480                        tokio::time::sleep(Duration::from_millis(100)).await;
481                        continue;
482                    } else {
483                        return Err(last_error.unwrap());
484                    }
485                }
486                Err(_) => {
487                    return Err(RFQError::ConnectionError(format!(
488                        "Hashflow quote request timed out after {} seconds",
489                        self.quote_timeout.as_secs()
490                    )));
491                }
492            };
493
494            if response.status() != 200 {
495                let err_msg = match response.text().await {
496                    Ok(text) => text,
497                    Err(e) => {
498                        warn!(
499                            "Hashflow error response parsing failed (attempt {}/{}): {}",
500                            attempt + 1,
501                            MAX_RETRIES,
502                            e
503                        );
504                        last_error = Some(RFQError::ParsingError(format!(
505                            "Failed to read response text from Hashflow failed request: {e}"
506                        )));
507                        if attempt < MAX_RETRIES - 1 {
508                            tokio::time::sleep(Duration::from_millis(100)).await;
509                            continue;
510                        } else {
511                            return Err(last_error.unwrap());
512                        }
513                    }
514                };
515                last_error = Some(RFQError::FatalError(format!(
516                    "Failed to send Hashflow quote request: {err_msg}",
517                )));
518                if attempt < MAX_RETRIES - 1 {
519                    warn!(
520                        "Hashflow returned non-200 status (attempt {}/{}): {}",
521                        attempt + 1,
522                        MAX_RETRIES,
523                        err_msg
524                    );
525                    tokio::time::sleep(Duration::from_millis(100)).await;
526                    continue;
527                } else {
528                    return Err(last_error.unwrap());
529                }
530            }
531
532            let quote_response = match response
533                .json::<HashflowQuoteResponse>()
534                .await
535            {
536                Ok(resp) => resp,
537                Err(e) => {
538                    warn!(
539                        "Hashflow quote response parsing failed (attempt {}/{}): {}",
540                        attempt + 1,
541                        MAX_RETRIES,
542                        e
543                    );
544                    last_error = Some(RFQError::ParsingError(format!(
545                        "Failed to parse Hashflow quote response: {e}"
546                    )));
547                    if attempt < MAX_RETRIES - 1 {
548                        tokio::time::sleep(Duration::from_millis(100)).await;
549                        continue;
550                    } else {
551                        return Err(last_error.unwrap());
552                    }
553                }
554            };
555
556            match quote_response.status.as_str() {
557                "success" => {
558                    if let Some(quotes) = quote_response.quotes {
559                        if quotes.is_empty() {
560                            return Err(RFQError::QuoteNotFound(format!(
561                                "Hashflow quote not found for {} {} ->{}",
562                                params.amount_in, params.token_in, params.token_out,
563                            )));
564                        }
565                        // We assume there will be only one quote request at a time
566                        let quote = quotes[0].clone();
567                        quote.validate(params, effective_trader)?;
568
569                        let mut quote_attributes: HashMap<String, Bytes> = HashMap::new();
570                        quote_attributes.insert("pool".to_string(), quote.quote_data.pool);
571                        if let Some(external_account) = quote.quote_data.external_account {
572                            quote_attributes
573                                .insert("external_account".to_string(), external_account);
574                        } else {
575                            quote_attributes.insert(
576                                "external_account".to_string(),
577                                Bytes::from_str(&Address::ZERO.to_string()).map_err(|_| {
578                                    RFQError::ParsingError(
579                                        "Failed to parse zero address".to_string(),
580                                    )
581                                })?,
582                            );
583                        }
584                        quote_attributes.insert("trader".to_string(), quote.quote_data.trader);
585                        quote_attributes
586                            .insert("effective_trader".to_string(), effective_trader.clone());
587                        quote_attributes
588                            .insert("base_token".to_string(), quote.quote_data.base_token);
589                        quote_attributes
590                            .insert("quote_token".to_string(), quote.quote_data.quote_token);
591                        quote_attributes.insert(
592                            "base_token_amount".to_string(),
593                            Bytes::from(
594                                biguint_to_u256(
595                                    &BigUint::from_str(&quote.quote_data.base_token_amount)
596                                        .map_err(|_| {
597                                            RFQError::ParsingError(format!(
598                                                "Failed to parse base token amount: {}",
599                                                quote.quote_data.base_token_amount
600                                            ))
601                                        })?,
602                                )
603                                .to_be_bytes::<32>()
604                                .to_vec(),
605                            ),
606                        );
607                        quote_attributes.insert(
608                            "quote_token_amount".to_string(),
609                            Bytes::from(
610                                biguint_to_u256(
611                                    &BigUint::from_str(&quote.quote_data.quote_token_amount)
612                                        .map_err(|_| {
613                                            RFQError::ParsingError(format!(
614                                                "Failed to parse quote token amount: {}",
615                                                quote.quote_data.quote_token_amount
616                                            ))
617                                        })?,
618                                )
619                                .to_be_bytes::<32>()
620                                .to_vec(),
621                            ),
622                        );
623                        quote_attributes.insert(
624                            "quote_expiry".to_string(),
625                            Bytes::from(
626                                U256::from(quote.quote_data.quote_expiry)
627                                    .to_be_bytes::<32>()
628                                    .to_vec(),
629                            ),
630                        );
631                        quote_attributes.insert(
632                            "nonce".to_string(),
633                            Bytes::from(
634                                U256::from(quote.quote_data.nonce)
635                                    .to_be_bytes::<32>()
636                                    .to_vec(),
637                            ),
638                        );
639                        quote_attributes.insert("tx_id".to_string(), quote.quote_data.tx_id);
640                        quote_attributes.insert("signature".to_string(), quote.signature);
641
642                        let signed_quote = SignedQuote {
643                            base_token: params.token_in.clone(),
644                            quote_token: params.token_out.clone(),
645                            amount_in: BigUint::from_str(&quote.quote_data.base_token_amount)
646                                .map_err(|_| {
647                                    RFQError::ParsingError(format!(
648                                        "Failed to parse amount in string: {}",
649                                        quote.quote_data.base_token_amount
650                                    ))
651                                })?,
652                            amount_out: BigUint::from_str(&quote.quote_data.quote_token_amount)
653                                .map_err(|_| {
654                                    RFQError::ParsingError(format!(
655                                        "Failed to parse amount out string: {}",
656                                        quote.quote_data.quote_token_amount
657                                    ))
658                                })?,
659                            quote_attributes,
660                        };
661                        return Ok(signed_quote);
662                    } else {
663                        return Err(RFQError::QuoteNotFound(format!(
664                            "Hashflow quote not found for {} {} ->{}",
665                            params.amount_in, params.token_in, params.token_out,
666                        )));
667                    }
668                }
669                "fail" => {
670                    let Some(error) = quote_response.error else {
671                        return Err(RFQError::FatalError(
672                            "Hashflow API error: request failed without an error".to_string(),
673                        ));
674                    };
675                    return Err(RFQError::FatalError(format!("Hashflow API error: {error}")));
676                }
677                _ => {
678                    return Err(RFQError::FatalError(
679                        "Hashflow API error: Unknown status".to_string(),
680                    ));
681                }
682            }
683        }
684
685        Err(last_error.unwrap_or_else(|| {
686            RFQError::ConnectionError("Hashflow quote request failed after retries".to_string())
687        }))
688    }
689}
690
691/// The id of the per-pair component of `mm_level`'s pair.
692fn pair_component_id(mm_level: &HashflowMarketMakerLevels) -> String {
693    let pair_str = format!(
694        "hashflow_{}/{}",
695        hex::encode(&mm_level.pair.base_token),
696        hex::encode(&mm_level.pair.quote_token)
697    );
698    format!("{}", keccak256(pair_str.as_bytes()))
699}
700
701#[async_trait]
702impl RFQClient for HashflowClient {
703    fn stream(
704        &self,
705    ) -> BoxStream<'static, Result<(String, StateSyncMessage<TimestampHeader>), RFQError>> {
706        let mut client = self.clone();
707
708        Box::pin(async_stream::stream! {
709            let mut current_components: HashMap<String, ComponentWithState> = HashMap::new();
710            let mut ticker = interval(client.poll_time);
711
712            info!("Starting Hashflow price levels polling every {} seconds", client.poll_time.as_secs());
713            info!("TVL threshold: {:.2}", client.tvl);
714
715            loop {
716                ticker.tick().await;
717
718                let market_makers;
719                match client.fetch_market_makers().await {
720                    Ok(mms) => {
721                        market_makers = mms;
722                        info!("Successfully fetched market makers");
723                    }
724                    Err(e) => {
725                        info!("Failed to fetch market makers: {}", e);
726                        continue;
727                    }
728                }
729
730                let levels_by_mm = match client.fetch_price_levels(&market_makers).await {
731                    Ok(levels_by_mm) => levels_by_mm,
732                    Err(e) => {
733                        error!("Failed to fetch price levels from Hashflow API: {}", e);
734                        continue;
735                    }
736                };
737                info!("Fetched price levels from {} market makers", levels_by_mm.len());
738                let components = match client.component_layout {
739                    ComponentLayout::PerPair => client.pair_components(&levels_by_mm)?,
740                    ComponentLayout::AllPairs => match client.all_pairs_component(&levels_by_mm) {
741                        Ok(component) => component
742                            .into_iter()
743                            .map(|component| (component.component.id.clone(), component))
744                            .collect(),
745                        Err(e) => {
746                            error!("Failed to build the Hashflow component: {}", e);
747                            continue;
748                        }
749                    },
750                };
751                let timestamp = component::unix_timestamp()?;
752                let msg = component::poll_message(&mut current_components, components, timestamp);
753                yield Ok((client.stream_name().to_string(), msg));
754            }
755        })
756    }
757
758    async fn request_binding_quote(
759        &self,
760        params: &GetAmountOutParams,
761    ) -> Result<SignedQuote, RFQError> {
762        self.request_quote(params, None).await
763    }
764}
765
766#[cfg(test)]
767mod tests {
768    use std::{env, str::FromStr, time::Duration};
769
770    use dotenv::dotenv;
771    use futures::StreamExt;
772    use rstest::rstest;
773    use tokio::time::timeout;
774
775    use super::*;
776    use crate::rfq::{
777        constants::get_hashflow_auth,
778        models::PriceLevel,
779        protocols::{
780            component::{BOOKS_ATTRIBUTE, SWAP_DIRECTIONS_ATTRIBUTE},
781            hashflow::{client_builder::HashflowClientBuilder, models::HashflowPair},
782            test_utils::{
783                effective_trader_of, mock_quote_server, quote_params, read_request_body,
784                HASHFLOW_QUOTE_RESPONSE,
785            },
786        },
787    };
788
789    #[test]
790    fn test_normalize_tvl_same_quote_token() {
791        let client = create_test_client();
792        let levels = HashMap::new();
793
794        // USDC is in our quote tokens, so no normalization should happen
795        let result = client.normalize_tvl(
796            1000.0,
797            Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap(),
798            &levels,
799        );
800        assert!(result.is_ok());
801        assert_eq!(result.unwrap(), 1000.0);
802    }
803
804    #[test]
805    fn test_normalize_tvl_different_quote_token() {
806        let client = create_test_client();
807        let mut levels = HashMap::new();
808        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
809        let usdc = Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap();
810
811        // Create mock levels for ETH/USDC pair for normalization
812        let eth_usdc_level = HashflowMarketMakerLevels {
813            pair: HashflowPair { base_token: weth.clone(), quote_token: usdc },
814            levels: vec![PriceLevel { quantity: 1.0, price: 3000.0 } /* 1 ETH = 3000 USDC */],
815        };
816
817        levels.insert("test_mm".to_string(), vec![eth_usdc_level]);
818
819        // Test normalizing ETH TVL to USDC
820        let result = client.normalize_tvl(2.0, weth, &levels);
821        assert!(result.is_ok());
822        // 2 ETH * 3000 USDC/ETH = 6000 USDC
823        assert_eq!(result.unwrap(), 6000.0);
824    }
825
826    #[test]
827    fn test_fallback_router_labels_components() {
828        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
829        let usdc = Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap();
830        let level = HashflowMarketMakerLevels {
831            pair: HashflowPair { base_token: weth.clone(), quote_token: usdc.clone() },
832            levels: vec![PriceLevel { quantity: 1.0, price: 3000.0 }],
833        };
834        let builder = || HashflowClientBuilder::new(Chain::Ethereum, String::new(), String::new());
835        let label = |client: HashflowClient| {
836            client
837                .create_component_with_state(
838                    String::from("hashflow"),
839                    vec![weth.clone(), usdc.clone()],
840                    "test_mm",
841                    &level,
842                    0.0,
843                )
844                .component
845                .protocol_system
846        };
847
848        assert_eq!(label(builder().build().unwrap()), HashflowClient::PROTOCOL_SYSTEM);
849        assert_eq!(
850            label(
851                builder()
852                    .with_fallback_router()
853                    .build()
854                    .unwrap()
855            ),
856            HashflowClient::FALLBACK_PROTOCOL_SYSTEM
857        );
858        assert_eq!(
859            HashflowClient::FALLBACK_PROTOCOL_SYSTEM,
860            tycho_execution::encoding::evm::HASHFLOW_FALLBACK_PROTOCOL_SYSTEM
861        );
862    }
863
864    #[test]
865    fn test_normalize_tvl_no_conversion_available() {
866        let client = create_test_client();
867        let levels = HashMap::new();
868        let result = client.normalize_tvl(
869            1000.0,
870            Bytes::from_str("0x1234567890123456789012345678901234567890").unwrap(),
871            &levels,
872        );
873        assert!(result.is_ok());
874        assert_eq!(result.unwrap(), 0.0);
875    }
876
877    fn create_test_client() -> HashflowClient {
878        let quote_tokens = HashSet::from([
879            Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap(), // USDC
880            Bytes::from_str("0xdAC17F958D2ee523a2206206994597C13D831ec7").unwrap(), // USDT
881        ]);
882
883        HashflowClient::new(
884            Chain::Ethereum,
885            HashSet::new(),
886            1.0,
887            quote_tokens,
888            "test_user".to_string(),
889            "test_key".to_string(),
890            Duration::from_secs(5),
891            Duration::from_secs(5),
892        )
893        .unwrap()
894    }
895
896    #[rstest]
897    #[case::per_pair(ComponentLayout::PerPair)]
898    #[case::all_pairs(ComponentLayout::AllPairs)]
899    #[tokio::test]
900    #[ignore] // Requires network access and HASHFLOW_KEY environment variable
901    async fn test_hashflow_api_polling(#[case] layout: ComponentLayout) {
902        dotenv().expect("Missing .env file");
903        let auth = get_hashflow_auth().unwrap();
904
905        let wbtc = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
906        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
907
908        let tokens = HashSet::from([wbtc, weth.clone()]);
909
910        let quote_tokens = HashSet::from([
911            Bytes::from_str("0xa0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap(), // USDC
912            Bytes::from_str("0xdac17f958d2ee523a2206206994597c13d831ec7").unwrap(), // USDT
913        ]);
914
915        let client = HashflowClient::new(
916            Chain::Ethereum,
917            tokens,
918            1.0, // $1 minimum TVL - very low to capture most pairs
919            quote_tokens,
920            auth.user,
921            auth.key,
922            Duration::from_secs(1),
923            Duration::from_secs(5),
924        )
925        .unwrap()
926        .with_component_layout(layout);
927
928        let mut stream = client.stream();
929
930        let result = timeout(Duration::from_secs(10), async {
931            let mut message_count = 0;
932            let max_messages = 3;
933            let mut total_components_received = 0;
934
935            while let Some(result) = stream.next().await {
936                match result {
937                    Ok((component_id, msg)) => {
938                        println!("Received message with ID: {component_id}");
939
940                        assert!(!component_id.is_empty());
941                        assert_eq!(component_id, client.stream_name());
942                        assert!(msg.header.timestamp > 0);
943
944                        let snapshot = &msg.snapshots;
945                        total_components_received += snapshot.states.len();
946
947                        println!("Received {} components in this message (Total so far: {})",
948                                snapshot.states.len(), total_components_received);
949
950                        if layout == ComponentLayout::AllPairs {
951                            assert!(snapshot.states.len() <= 1, "one all-pairs component");
952                        }
953                        for (id, component_with_state) in &snapshot.states {
954                            let attributes = &component_with_state.state.attributes;
955                            match layout {
956                                ComponentLayout::PerPair => {
957                                    assert!(!attributes["levels"].is_empty());
958                                    assert!(!attributes["mm"].is_empty());
959                                }
960                                ComponentLayout::AllPairs => {
961                                    let expected_id = component::component_id(
962                                        HashflowClient::PROTOCOL_SYSTEM,
963                                        client.chain,
964                                    );
965                                    assert_eq!(id, &expected_id);
966                                    let price_levels: Vec<MakerPriceLevels> =
967                                        serde_json::from_slice(&attributes[BOOKS_ATTRIBUTE]).unwrap();
968                                    assert!(!price_levels.is_empty());
969                                    let directions = &component_with_state.component.static_attributes[SWAP_DIRECTIONS_ATTRIBUTE];
970                                    assert_eq!(directions.len() % 40, 0);
971                                }
972                            }
973                            if let Some(tvl) = component_with_state.component_tvl {
974                                assert!(tvl >= 1.0);
975                                println!("Component {id} TVL: ${tvl:.2}");
976                            }
977                        }
978
979                        message_count += 1;
980                        if message_count >= max_messages {
981                            break;
982                        }
983                    }
984                    Err(e) => {
985                        panic!("Stream error: {e}");
986                    }
987                }
988            }
989
990            assert!(message_count > 0, "Should have received at least one message");
991            assert!(total_components_received >= 1, "Should have received at least 1 component with $1 TVL threshold");
992            println!("Successfully received {message_count} messages with {total_components_received} total components");
993        })
994        .await;
995
996        match result {
997            Ok(_) => println!("Test completed successfully"),
998            Err(_) => panic!("Test timed out - no messages received within 5 seconds"),
999        }
1000    }
1001
1002    #[tokio::test]
1003    #[ignore] // Requires network access and setting proper env vars
1004    async fn test_request_binding_quote() {
1005        let wbtc = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
1006        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1007
1008        let auth_user = String::from("propellerheads");
1009        dotenv().expect("Missing .env file");
1010        let auth_key = env::var("HASHFLOW_KEY").unwrap();
1011
1012        let client = HashflowClient::new(
1013            Chain::Ethereum,
1014            HashSet::from_iter(vec![weth.clone(), wbtc.clone()]),
1015            10.0,
1016            HashSet::new(),
1017            auth_user,
1018            auth_key,
1019            Duration::from_secs(0),
1020            Duration::from_secs(5),
1021        )
1022        .unwrap();
1023
1024        let router = Bytes::from_str("0xfD0b31d2E955fA55e3fa641Fe90e08b677188d35").unwrap();
1025
1026        let params = GetAmountOutParams {
1027            amount_in: BigUint::from(1_000000000000000000u64),
1028            token_in: weth.clone(),
1029            token_out: wbtc.clone(),
1030            sender: router.clone(),
1031            receiver: router.clone(),
1032        };
1033        let quote = client
1034            .request_binding_quote(&params)
1035            .await
1036            .unwrap();
1037
1038        assert_eq!(quote.base_token, weth);
1039        assert_eq!(quote.quote_token, wbtc);
1040        assert_eq!(quote.amount_in, BigUint::from(1_000000000000000000u64));
1041
1042        // // Assuming the BTC - WETH price doesn't change too much at the time of running this
1043        assert!(quote.amount_out > BigUint::from(3000000u64));
1044
1045        assert_eq!(quote.quote_attributes.len(), 12);
1046        let expected_attributes = [
1047            "pool",
1048            "external_account",
1049            "trader",
1050            "effective_trader",
1051            "base_token",
1052            "quote_token",
1053            "base_token_amount",
1054            "quote_token_amount",
1055            "quote_expiry",
1056            "nonce",
1057            "tx_id",
1058            "signature",
1059        ];
1060        for attr in expected_attributes {
1061            assert!(
1062                quote
1063                    .quote_attributes
1064                    .contains_key(attr),
1065                "Missing attribute: {attr}"
1066            );
1067        }
1068        assert_eq!(
1069            quote
1070                .quote_attributes
1071                .get("trader")
1072                .unwrap(),
1073            &router
1074        );
1075    }
1076
1077    const QUOTE_RESPONSE_WITHOUT_EFFECTIVE_TRADER: &str = r#"{"status":"success","error":null,"rfqId":"test-rfq-id","internalRfqIds":null,"quotes":[{"quoteData":{"pool":"0x71D9750ECF0c5081FAE4E3EDC4253E52024b0B59","externalAccount":null,"trader":"0xfD0b31d2E955fA55e3fa641Fe90e08b677188d35","baseToken":"0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2","baseTokenAmount":"1000000000000000000","quoteToken":"0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599","quoteTokenAmount":"3329502","quoteExpiry":1707847360,"nonce":1707844960943648659,"txid":"0x0000000000000000000000000000000000000000000000000000000000000001"},"signature":"0x1234567890abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890abcdef12"}]}"#;
1078
1079    fn create_test_hashflow_client(
1080        quote_endpoint: String,
1081        quote_timeout: Duration,
1082    ) -> HashflowClient {
1083        let token_in = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1084        let token_out = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
1085
1086        HashflowClient {
1087            chain: Chain::Ethereum,
1088            price_levels_endpoint: "http://unused/price-levels".to_string(),
1089            market_makers_endpoint: "http://unused/market-makers".to_string(),
1090            quote_endpoint,
1091            tokens: HashSet::from([token_in, token_out]),
1092            tvl: 10.0,
1093            auth_key: "test_key".to_string(),
1094            auth_user: "test_user".to_string(),
1095            quote_tokens: HashSet::new(),
1096            poll_time: Duration::from_secs(0),
1097            quote_timeout,
1098            protocol_system: HashflowClient::PROTOCOL_SYSTEM.to_string(),
1099            component_layout: ComponentLayout::PerPair,
1100            quote_rule: QuoteRule::OncePerMaker,
1101        }
1102    }
1103
1104    #[tokio::test]
1105    async fn test_request_binding_quote_without_effective_trader() {
1106        // A response that drops the requested effectiveTrader would leave the quote in the
1107        // trader's shared nonce scope, so the client rejects it.
1108        let (addr, _) = mock_quote_server(0, QUOTE_RESPONSE_WITHOUT_EFFECTIVE_TRADER).await;
1109        let client = create_test_hashflow_client(
1110            format!("http://127.0.0.1:{}/rfq", addr.port()),
1111            Duration::from_secs(1),
1112        );
1113        let params = quote_params();
1114
1115        let err = client
1116            .request_binding_quote(&params)
1117            .await
1118            .unwrap_err();
1119
1120        assert!(format!("{err:?}").contains("Effective trader mismatch"));
1121    }
1122
1123    #[tokio::test]
1124    async fn test_request_binding_quote_field_mapping() {
1125        // The wire request carries the receiver as Hashflow's trader and a fresh random
1126        // address as the effectiveTrader — a new one per quote request.
1127        let (addr, request_log) = mock_quote_server(0, HASHFLOW_QUOTE_RESPONSE).await;
1128        let client = create_test_hashflow_client(
1129            format!("http://127.0.0.1:{}/rfq", addr.port()),
1130            Duration::from_secs(1),
1131        );
1132        let params = quote_params();
1133
1134        let first_quote = client
1135            .request_binding_quote(&params)
1136            .await
1137            .unwrap();
1138        client
1139            .request_binding_quote(&params)
1140            .await
1141            .unwrap();
1142
1143        let requests = request_log.lock().unwrap();
1144        assert_eq!(requests.len(), 2);
1145        for body in requests.iter() {
1146            assert!(
1147                body.contains(&format!("\"trader\":\"{}\"", params.receiver)),
1148                "trader is not the receiver: {body}"
1149            );
1150        }
1151        let first = effective_trader_of(&requests[0]);
1152        let second = effective_trader_of(&requests[1]);
1153        assert_eq!(first.len(), 42, "effective trader is not an address");
1154        assert_ne!(first, second, "effective traders are not unique per quote");
1155        assert_ne!(first, params.receiver.to_string(), "effective trader equals the trader");
1156        assert_eq!(
1157            first_quote
1158                .quote_attributes
1159                .get("effective_trader")
1160                .unwrap()
1161                .to_string(),
1162            first,
1163            "quote attributes do not carry the requested effective trader"
1164        );
1165    }
1166
1167    #[tokio::test]
1168    async fn test_request_quote_from_maker() {
1169        let (addr, request_log) = mock_quote_server(0, HASHFLOW_QUOTE_RESPONSE).await;
1170        let client = create_test_hashflow_client(
1171            format!("http://127.0.0.1:{}/rfq", addr.port()),
1172            Duration::from_secs(1),
1173        );
1174
1175        client
1176            .request_quote(&quote_params(), Some("mm1"))
1177            .await
1178            .unwrap();
1179
1180        let request = &request_log.lock().unwrap()[0];
1181        assert!(request.contains("\"marketMakers\":[\"mm1\"]"), "{request}");
1182        assert!(request.contains("\"doNotRetryWithOtherMakers\":true"), "{request}");
1183    }
1184
1185    #[tokio::test]
1186    async fn test_request_binding_quote_names_no_maker() {
1187        let (addr, request_log) = mock_quote_server(0, HASHFLOW_QUOTE_RESPONSE).await;
1188        let client = create_test_hashflow_client(
1189            format!("http://127.0.0.1:{}/rfq", addr.port()),
1190            Duration::from_secs(1),
1191        );
1192
1193        client
1194            .request_binding_quote(&quote_params())
1195            .await
1196            .unwrap();
1197
1198        let request = &request_log.lock().unwrap()[0];
1199        assert!(!request.contains("marketMakers"), "{request}");
1200        assert!(!request.contains("options"), "{request}");
1201    }
1202
1203    fn levels(pair: (&Bytes, &Bytes), levels: &[(f64, f64)]) -> HashflowMarketMakerLevels {
1204        HashflowMarketMakerLevels {
1205            pair: HashflowPair { base_token: pair.0.clone(), quote_token: pair.1.clone() },
1206            levels: levels
1207                .iter()
1208                .map(|&(quantity, price)| PriceLevel { quantity, price })
1209                .collect(),
1210        }
1211    }
1212
1213    /// WETH and WBTC quoted against USDC by two makers, with USDC as the TVL quote token.
1214    fn all_pairs_test_client() -> (HashflowClient, Bytes, Bytes, Bytes) {
1215        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1216        let wbtc = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
1217        let usdc = Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap();
1218        let mut client =
1219            create_test_hashflow_client("http://unused/rfq".to_string(), Duration::from_secs(1));
1220        client.tokens = HashSet::from([weth.clone(), wbtc.clone(), usdc.clone()]);
1221        client.quote_tokens = HashSet::from([usdc.clone()]);
1222        client.tvl = 100.0;
1223        (client, weth, wbtc, usdc)
1224    }
1225
1226    #[test]
1227    fn test_all_pairs_component_keeps_books_above_tvl_threshold() {
1228        let (client, weth, wbtc, usdc) = all_pairs_test_client();
1229        let levels_by_mm = HashMap::from([
1230            (
1231                "mm_b".to_string(),
1232                vec![
1233                    levels((&weth, &usdc), &[(1.0, 3000.0)]),
1234                    levels((&wbtc, &usdc), &[(0.001, 65.0)]),
1235                ],
1236            ),
1237            ("mm_a".to_string(), vec![levels((&weth, &usdc), &[(2.0, 3000.0)])]),
1238        ]);
1239
1240        let component = client
1241            .all_pairs_component(&levels_by_mm)
1242            .unwrap()
1243            .expect("two makers' price levels clear the threshold");
1244
1245        assert_eq!(
1246            component.component.id,
1247            component::component_id(HashflowClient::PROTOCOL_SYSTEM, client.chain)
1248        );
1249        assert_eq!(
1250            component.component_tvl,
1251            Some(9000.0),
1252            "the WBTC price levels are below the threshold"
1253        );
1254        let mut expected_tokens = vec![weth.clone(), usdc.clone()];
1255        expected_tokens.sort();
1256        assert_eq!(component.component.tokens, expected_tokens);
1257        let mut expected_directions = weth.to_vec();
1258        expected_directions.extend_from_slice(&usdc);
1259        assert_eq!(
1260            component.component.static_attributes[SWAP_DIRECTIONS_ATTRIBUTE].to_vec(),
1261            expected_directions
1262        );
1263        assert_eq!(
1264            component.component.static_attributes[QuoteRule::ATTRIBUTE].as_ref(),
1265            b"once_per_maker"
1266        );
1267
1268        let price_levels: Vec<MakerPriceLevels> =
1269            serde_json::from_slice(&component.state.attributes[BOOKS_ATTRIBUTE]).unwrap();
1270        assert_eq!(price_levels.len(), 2);
1271        assert_eq!(price_levels[0].market_maker, "mm_a", "price levels are sorted by maker");
1272        assert_eq!(price_levels[0].levels[0].quantity, 2.0);
1273        assert_eq!(price_levels[1].market_maker, "mm_b");
1274        assert_eq!(price_levels[1].base_token, weth);
1275    }
1276
1277    /// Each component carries the id and attributes the per-pair state decodes.
1278    #[test]
1279    fn test_pair_components_keep_pairs_above_tvl_threshold() {
1280        let (client, weth, wbtc, usdc) = all_pairs_test_client();
1281        let levels_by_mm = HashMap::from([(
1282            "mm_a".to_string(),
1283            vec![
1284                levels((&weth, &usdc), &[(1.0, 3000.0)]),
1285                levels((&wbtc, &usdc), &[(0.001, 65.0)]),
1286            ],
1287        )]);
1288
1289        let components = client
1290            .pair_components(&levels_by_mm)
1291            .unwrap();
1292
1293        let id = format!(
1294            "{}",
1295            keccak256(format!("hashflow_{}/{}", hex::encode(&weth), hex::encode(&usdc)).as_bytes())
1296        );
1297        assert_eq!(components.keys().collect::<Vec<_>>(), [&id], "the WBTC pair is below TVL");
1298        let component = &components[&id];
1299        assert_eq!(component.component.tokens, [weth, usdc]);
1300        assert_eq!(component.component_tvl, Some(3000.0));
1301        assert!(component
1302            .component
1303            .static_attributes
1304            .is_empty());
1305        assert_eq!(component.state.attributes["mm"].as_ref(), b"mm_a");
1306        let levels: Vec<PriceLevel> =
1307            serde_json::from_slice(&component.state.attributes["levels"]).unwrap();
1308        assert_eq!(levels, [PriceLevel { quantity: 1.0, price: 3000.0 }]);
1309    }
1310
1311    #[tokio::test]
1312    async fn test_hashflow_quote_timeout() {
1313        let (addr, _) = mock_quote_server(500, HASHFLOW_QUOTE_RESPONSE).await;
1314
1315        // Test 1: Client with short timeout (200ms) - should timeout
1316        let client_short_timeout = create_test_hashflow_client(
1317            format!("http://127.0.0.1:{}/rfq", addr.port()),
1318            Duration::from_millis(200),
1319        );
1320        let params = quote_params();
1321
1322        // This should timeout after 200ms
1323        let start = std::time::Instant::now();
1324        let result = client_short_timeout
1325            .request_binding_quote(&params)
1326            .await;
1327        let elapsed = start.elapsed();
1328
1329        // Verify that we got a timeout error
1330        assert!(result.is_err());
1331        let err = result.unwrap_err();
1332        match err {
1333            RFQError::ConnectionError(msg) => {
1334                assert!(msg.contains("timed out"), "Expected timeout error, got: {}", msg);
1335            }
1336            _ => panic!("Expected ConnectionError, got: {:?}", err),
1337        }
1338        // Should have timed out around 200ms, definitely less than 400ms
1339        assert!(
1340            elapsed.as_millis() >= 200 && elapsed.as_millis() < 400,
1341            "Expected timeout around 200ms, got: {:?}",
1342            elapsed
1343        );
1344
1345        // Test 2: Client with long timeout (1 second) - should wait and receive response
1346        // Note: With retry logic, we may need multiple attempts if the response is malformed,
1347        // so we need a longer timeout to account for retries
1348        let client_long_timeout = create_test_hashflow_client(
1349            format!("http://127.0.0.1:{}/rfq", addr.port()),
1350            Duration::from_secs(1),
1351        );
1352
1353        // This should wait for the response (500ms)
1354        let result = client_long_timeout
1355            .request_binding_quote(&params)
1356            .await;
1357
1358        // Should succeed - the server waits 500ms which is within the 1s timeout
1359        assert!(result.is_ok(), "Expected success, got: {:?}", result);
1360    }
1361
1362    /// Helper function to create a mock server that fails twice, then succeeds
1363    async fn create_retry_server() -> (std::net::SocketAddr, std::sync::Arc<std::sync::Mutex<u32>>)
1364    {
1365        use std::sync::{Arc, Mutex};
1366
1367        use tokio::{io::AsyncWriteExt, net::TcpListener};
1368
1369        let request_count = Arc::new(Mutex::new(0u32));
1370        let request_count_clone = request_count.clone();
1371
1372        let listener = TcpListener::bind("127.0.0.1:0")
1373            .await
1374            .unwrap();
1375        let addr = listener.local_addr().unwrap();
1376
1377        tokio::spawn(async move {
1378            while let Ok((mut stream, _)) = listener.accept().await {
1379                let count_clone = request_count_clone.clone();
1380                tokio::spawn(async move {
1381                    *count_clone.lock().unwrap() += 1;
1382                    let count = *count_clone.lock().unwrap();
1383                    println!("Mock server: Received request #{count}");
1384
1385                    let body = read_request_body(&mut stream).await;
1386                    if count <= 2 {
1387                        let response = "HTTP/1.1 500 Internal Server Error\r\nContent-Length: 21\r\n\r\nInternal Server Error";
1388                        let _ = stream
1389                            .write_all(response.as_bytes())
1390                            .await;
1391                    } else {
1392                        let json_response = HASHFLOW_QUOTE_RESPONSE
1393                            .replace("{{EFFECTIVE_TRADER}}", &effective_trader_of(&body));
1394                        let response = format!(
1395                            "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
1396                            json_response.len(),
1397                            json_response
1398                        );
1399                        let _ = stream
1400                            .write_all(response.as_bytes())
1401                            .await;
1402                    }
1403                    let _ = stream.flush().await;
1404                    let _ = stream.shutdown().await;
1405                });
1406            }
1407        });
1408
1409        tokio::time::sleep(Duration::from_millis(50)).await;
1410        (addr, request_count)
1411    }
1412
1413    #[tokio::test]
1414    async fn test_hashflow_quote_retry_on_bad_response() {
1415        let (addr, request_count) = create_retry_server().await;
1416
1417        let client = create_test_hashflow_client(
1418            format!("http://127.0.0.1:{}/rfq", addr.port()),
1419            Duration::from_secs(5),
1420        );
1421        let params = quote_params();
1422        let result = client
1423            .request_binding_quote(&params)
1424            .await;
1425
1426        assert!(result.is_ok(), "Expected success after retries, got: {:?}", result);
1427        let quote = result.unwrap();
1428
1429        // Verify the quote is parsed as expected
1430        assert_eq!(quote.amount_in, BigUint::from(1_000000000000000000u64));
1431        assert_eq!(quote.amount_out, BigUint::from(3329502u64));
1432
1433        // Verify exactly 3 requests were made (2 failures + 1 success)
1434        let final_count = *request_count.lock().unwrap();
1435        assert_eq!(final_count, 3, "Expected 3 requests, got {}", final_count);
1436    }
1437
1438    #[test]
1439    fn test_hashflow_client_serialize_deserialize_roundtrip() {
1440        let token_in = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1441        let token_out = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
1442        let quote_token = Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap();
1443
1444        let original = HashflowClient {
1445            chain: Chain::Ethereum,
1446            price_levels_endpoint: "https://api.hashflow.com/price_levels".to_string(),
1447            market_makers_endpoint: "https://api.hashflow.com/market_makers".to_string(),
1448            quote_endpoint: "https://api.hashflow.com/quote".to_string(),
1449            tokens: HashSet::from([token_in.clone(), token_out.clone()]),
1450            tvl: 50.5,
1451            auth_key: "secret_key".to_string(),
1452            auth_user: "secret_user".to_string(),
1453            quote_tokens: HashSet::from([quote_token.clone()]),
1454            poll_time: Duration::from_secs(10),
1455            quote_timeout: Duration::from_millis(5500),
1456            protocol_system: HashflowClient::PROTOCOL_SYSTEM.to_string(),
1457            component_layout: ComponentLayout::PerPair,
1458            quote_rule: QuoteRule::OncePerMaker,
1459        };
1460
1461        let serialized = serde_json::to_string(&original).unwrap();
1462        let deserialized: HashflowClient = serde_json::from_str(&serialized).unwrap();
1463
1464        // Fields that should round-trip correctly
1465        assert_eq!(deserialized.chain, original.chain);
1466        assert_eq!(deserialized.price_levels_endpoint, original.price_levels_endpoint);
1467        assert_eq!(deserialized.market_makers_endpoint, original.market_makers_endpoint);
1468        assert_eq!(deserialized.quote_endpoint, original.quote_endpoint);
1469        assert_eq!(deserialized.tokens, original.tokens);
1470        assert_eq!(deserialized.tvl, original.tvl);
1471        assert_eq!(deserialized.quote_tokens, original.quote_tokens);
1472        assert_eq!(deserialized.poll_time, original.poll_time);
1473        assert_eq!(deserialized.quote_timeout, original.quote_timeout);
1474
1475        // auth_key and auth_user should NOT round-trip (skip_serializing + default)
1476        assert_eq!(deserialized.auth_key, "");
1477        assert_eq!(deserialized.auth_user, "");
1478        assert_ne!(deserialized.auth_key, original.auth_key);
1479        assert_ne!(deserialized.auth_user, original.auth_user);
1480    }
1481
1482    #[test]
1483    fn test_hashflow_client_deserialize_with_credentials() {
1484        // When auth_key and auth_user are provided in JSON, they should be deserialized
1485        // (skip_serializing only affects serialization, not deserialization)
1486        let json = r#"{
1487            "chain": "ethereum",
1488            "price_levels_endpoint": "https://api.hashflow.com/price_levels",
1489            "market_makers_endpoint": "https://api.hashflow.com/market_makers",
1490            "quote_endpoint": "https://api.hashflow.com/quote",
1491            "tokens": [],
1492            "tvl": 10.0,
1493            "auth_key": "provided_key",
1494            "auth_user": "provided_user",
1495            "quote_tokens": [],
1496            "poll_time": {"secs": 10, "nanos": 0},
1497            "quote_timeout": {"secs": 30, "nanos": 0}
1498        }"#;
1499
1500        let client: HashflowClient = serde_json::from_str(json).unwrap();
1501
1502        // Credentials should be deserialized from JSON
1503        assert_eq!(client.auth_key, "provided_key");
1504        assert_eq!(client.auth_user, "provided_user");
1505    }
1506}