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