Skip to main content

tycho_simulation/rfq/protocols/hashflow/
client.rs

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