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: HashSet<Bytes>,
48 tvl: f64,
50 #[serde(skip_serializing, default)]
51 auth_key: String,
52 #[serde(skip_serializing, default)]
53 auth_user: String,
54 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 #[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 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 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 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 self.quote_tokens.contains("e_token) {
154 return Ok(raw_tvl);
155 }
156
157 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 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 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![], ..Default::default()
194 };
195
196 let mut attributes = HashMap::new();
197
198 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 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 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 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 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 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, "e_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 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 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("e_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 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("e.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("e.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("e.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("e.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
691fn 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 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 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 } ],
815 };
816
817 levels.insert("test_mm".to_string(), vec![eth_usdc_level]);
818
819 let result = client.normalize_tvl(2.0, weth, &levels);
821 assert!(result.is_ok());
822 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(), Bytes::from_str("0xdAC17F958D2ee523a2206206994597C13D831ec7").unwrap(), ]);
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] 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(), Bytes::from_str("0xdac17f958d2ee523a2206206994597c13d831ec7").unwrap(), ]);
914
915 let client = HashflowClient::new(
916 Chain::Ethereum,
917 tokens,
918 1.0, 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] 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(¶ms)
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 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 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(¶ms)
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 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(¶ms)
1136 .await
1137 .unwrap();
1138 client
1139 .request_binding_quote(¶ms)
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("e_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("e_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 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 #[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 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 let start = std::time::Instant::now();
1324 let result = client_short_timeout
1325 .request_binding_quote(¶ms)
1326 .await;
1327 let elapsed = start.elapsed();
1328
1329 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 assert!(
1340 elapsed.as_millis() >= 200 && elapsed.as_millis() < 400,
1341 "Expected timeout around 200ms, got: {:?}",
1342 elapsed
1343 );
1344
1345 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 let result = client_long_timeout
1355 .request_binding_quote(¶ms)
1356 .await;
1357
1358 assert!(result.is_ok(), "Expected success, got: {:?}", result);
1360 }
1361
1362 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(¶ms)
1424 .await;
1425
1426 assert!(result.is_ok(), "Expected success after retries, got: {:?}", result);
1427 let quote = result.unwrap();
1428
1429 assert_eq!(quote.amount_in, BigUint::from(1_000000000000000000u64));
1431 assert_eq!(quote.amount_out, BigUint::from(3329502u64));
1432
1433 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 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 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 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 assert_eq!(client.auth_key, "provided_key");
1504 assert_eq!(client.auth_user, "provided_user");
1505 }
1506}