1use std::{
2 collections::{HashMap, HashSet},
3 str::FromStr,
4 sync::LazyLock,
5 time::SystemTime,
6};
7
8use alloy::primitives::{utils::keccak256, Address};
9use async_trait::async_trait;
10use futures::stream::BoxStream;
11use num_bigint::BigUint;
12use reqwest::Client;
13use serde::{Deserialize, Serialize};
14use tokio::time::{interval, timeout, Duration};
15use tracing::{error, info, warn};
16use tycho_common::{
17 models::{protocol::GetAmountOutParams, Chain},
18 simulation::indicatively_priced::SignedQuote,
19 Bytes,
20};
21
22use super::models::{NativeOrderbookEntry, NativeOrderbookSide, NativePriceData, NativePriceLevel};
23use crate::{
24 rfq::{
25 client::RFQClient,
26 errors::RFQError,
27 models::TimestampHeader,
28 protocols::{
29 native::models::{
30 FirmQuoteRequest, FirmQuoteResponse, NativeApiErrorResponse, NativeSupportedChain,
31 },
32 utils::bytes_to_address,
33 },
34 },
35 tycho_client::feed::synchronizer::{ComponentWithState, Snapshot, StateSyncMessage},
36 tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState},
37};
38
39static NATIVE_HTTP_CLIENT: LazyLock<Client> = LazyLock::new(Client::new);
40const MAX_QUOTE_ATTEMPTS: u32 = 3;
41const TRANSIENT_RETRY_DELAY: Duration = Duration::from_millis(100);
42const NATIVE_API_RETRY_DELAY: Duration = Duration::from_secs(1);
43const TRADE_RFQT_SELECTOR: [u8; 4] = [0x70, 0x83, 0x52, 0x7c];
45const MIN_TRADE_RFQT_CALLDATA_LEN: usize = 4 + 3 * 32;
46const ACTUAL_SELLER_AMOUNT_OFFSET: usize = 4 + 32;
47const ACTUAL_MIN_OUTPUT_AMOUNT_OFFSET: usize = 4 + 2 * 32;
48
49enum QuoteAttemptError {
50 Retry { error: RFQError, delay: Duration },
51 Fatal(RFQError),
52}
53
54impl QuoteAttemptError {
55 fn into_error(self) -> RFQError {
56 match self {
57 Self::Retry { error, .. } | Self::Fatal(error) => error,
58 }
59 }
60}
61
62#[derive(Default)]
63struct AggregatedLevels {
64 levels: Vec<NativePriceLevel>,
65 minimum: f64,
67}
68
69impl AggregatedLevels {
70 fn extend(&mut self, levels: Vec<NativePriceLevel>, minimum: f64) {
71 self.levels.extend(levels);
72 self.minimum = self.minimum.max(minimum);
73 }
74}
75
76#[derive(Clone, Debug, Serialize, Deserialize)]
77pub struct NativeClient {
78 chain: Chain,
79 endpoint: String,
80 #[serde(skip_serializing, default)]
81 api_key: String,
82 tokens: HashSet<Bytes>,
83 tvl: f64,
84 quote_tokens: HashSet<Bytes>,
85 poll_time: Duration,
86 quote_timeout: Duration,
87}
88
89impl NativeClient {
90 pub const PROTOCOL_SYSTEM: &'static str = "rfq:native";
91 pub const DEFAULT_ENDPOINT: &'static str = "https://v2.api.native.org/swap-api-v2/v1";
92
93 fn classify_api_error(error: &NativeApiErrorResponse) -> QuoteAttemptError {
96 let message = format!("Native API error {}: {}", error.code, error.message);
97 match error.code {
98 301016 | 405030 => QuoteAttemptError::Retry {
100 error: RFQError::QuoteNotFound(message),
101 delay: NATIVE_API_RETRY_DELAY,
102 },
103 201005 => QuoteAttemptError::Retry {
104 error: RFQError::ConnectionError(message),
105 delay: NATIVE_API_RETRY_DELAY,
106 },
107 101010 | 171037 | 171011 | 171015 | 171055 | 101007 => {
111 QuoteAttemptError::Fatal(RFQError::QuoteNotFound(message))
112 }
113 131003 | 131004 | 131011 | 171018 | 171053 | 131005 => {
115 QuoteAttemptError::Fatal(RFQError::InvalidInput(message))
116 }
117 201001 => QuoteAttemptError::Fatal(RFQError::FatalError(message)),
118 _ => QuoteAttemptError::Fatal(RFQError::FatalError(format!(
119 "Unknown Native API error {}: {}",
120 error.code, error.message
121 ))),
122 }
123 }
124
125 pub fn new(
126 chain: Chain,
127 api_key: String,
128 tokens: HashSet<Bytes>,
129 tvl: f64,
130 quote_tokens: HashSet<Bytes>,
131 poll_time: Duration,
132 quote_timeout: Duration,
133 ) -> Result<Self, RFQError> {
134 NativeSupportedChain::try_from(chain).map_err(RFQError::InvalidInput)?;
135 if poll_time.is_zero() {
136 return Err(RFQError::InvalidInput(
137 "Native polling interval must be greater than zero".to_string(),
138 ))
139 }
140 Ok(Self {
141 chain,
142 endpoint: Self::DEFAULT_ENDPOINT.to_string(),
143 api_key,
144 tokens,
145 tvl,
146 quote_tokens,
147 poll_time,
148 quote_timeout,
149 })
150 }
151
152 fn select_tvl_conversion_book<'a>(
153 &self,
154 quote_address: &Bytes,
155 books: &'a HashMap<String, NativePriceData>,
156 ) -> Option<&'a NativePriceData> {
157 books
158 .values()
159 .filter(|candidate| {
160 candidate.base_address == *quote_address &&
163 self.quote_tokens
164 .contains(&candidate.quote_address)
165 })
166 .filter_map(|candidate| {
167 candidate
168 .calculate_tvl(None)
169 .map(|liquidity| (candidate, liquidity))
170 })
171 .max_by(|(candidate_a, liquidity_a), (candidate_b, liquidity_b)| {
172 liquidity_a
173 .total_cmp(liquidity_b)
174 .then_with(|| {
177 candidate_b
178 .quote_address
179 .as_ref()
180 .cmp(candidate_a.quote_address.as_ref())
181 })
182 })
183 .map(|(candidate, _)| candidate)
184 }
185
186 fn create_component_with_state(
187 &self,
188 component_id: String,
189 tokens: Vec<Bytes>,
190 book: NativePriceData,
191 tvl: f64,
192 ) -> ComponentWithState {
193 let protocol_component = ProtocolComponent {
194 id: component_id.clone(),
195 protocol_system: Self::PROTOCOL_SYSTEM.to_string(),
196 protocol_type_name: "native_relay_pool".to_string(),
197 chain: self.chain,
198 tokens,
199 contract_addresses: vec![],
200 static_attributes: Default::default(),
201 change: Default::default(),
202 creation_tx: Default::default(),
203 created_at: Default::default(),
204 };
205
206 let mut attributes = HashMap::new();
207
208 let book_json = serde_json::to_string(&book).unwrap_or_default();
209 attributes.insert("book".to_string(), book_json.as_bytes().to_vec().into());
210
211 ComponentWithState {
212 state: ProtocolComponentState::new(&component_id, attributes, HashMap::new()),
213 component: protocol_component,
214 component_tvl: Some(tvl),
215 entrypoints: vec![],
216 }
217 }
218
219 async fn fetch_orderbook(&self) -> Result<Vec<NativeOrderbookEntry>, RFQError> {
220 let chain = NativeSupportedChain::try_from(self.chain).map_err(RFQError::InvalidInput)?;
221 let response = NATIVE_HTTP_CLIENT
222 .get(format!("{}/orderbook", self.endpoint))
223 .query(&[("chain", chain.as_str()), ("showNative", "0x0")])
226 .header("accept", "application/json")
227 .header("apikey", &self.api_key)
228 .send()
229 .await
230 .map_err(|e| RFQError::ConnectionError(e.to_string()))?;
231
232 let status = response.status();
233 let response_body = response
234 .bytes()
235 .await
236 .map_err(|e| RFQError::ConnectionError(e.to_string()))?;
237
238 if let Ok(api_error) = serde_json::from_slice::<NativeApiErrorResponse>(&response_body) {
240 return Err(Self::classify_api_error(&api_error).into_error());
241 }
242
243 if !status.is_success() {
244 return Err(RFQError::ConnectionError(format!(
245 "Native Relay orderbook HTTP error {}: {}",
246 status,
247 String::from_utf8_lossy(&response_body)
248 )));
249 }
250
251 serde_json::from_slice(&response_body).map_err(|e| {
252 RFQError::ParsingError(format!("Failed to parse Native Relay orderbook: {e}"))
253 })
254 }
255
256 fn group_orderbook(
257 &self,
258 entries: Vec<NativeOrderbookEntry>,
259 ) -> HashMap<String, NativePriceData> {
260 let mut entries_by_pair: HashMap<(Bytes, Bytes), Vec<NativeOrderbookEntry>> =
261 HashMap::new();
262
263 for entry in entries {
264 let pair = if entry.base_address.as_ref() <= entry.quote_address.as_ref() {
265 (entry.base_address.clone(), entry.quote_address.clone())
266 } else {
267 (entry.quote_address.clone(), entry.base_address.clone())
268 };
269 entries_by_pair
270 .entry(pair)
271 .or_default()
272 .push(entry);
273 }
274
275 let mut books = HashMap::new();
276 for ((token0, token1), entries) in entries_by_pair {
277 let token0_is_quote = self.quote_tokens.contains(&token0);
282 let token1_is_quote = self.quote_tokens.contains(&token1);
283 let (base_address, quote_address) = match (token0_is_quote, token1_is_quote) {
284 (true, false) => (token1.clone(), token0.clone()),
285 (false, true) => (token0.clone(), token1.clone()),
286 (true, true) | (false, false) => (token0.clone(), token1.clone()),
287 };
288 let mut direct_bids = AggregatedLevels::default();
289 let mut direct_asks = AggregatedLevels::default();
290 let mut mirrored_bids = AggregatedLevels::default();
291 let mut mirrored_asks = AggregatedLevels::default();
292
293 for entry in entries {
298 let levels: Vec<_> = entry
302 .levels
303 .into_iter()
304 .filter(|level| level.quantity != 0.0)
305 .collect();
306 if levels.is_empty() {
307 continue
308 }
309
310 let is_direct =
311 entry.base_address == base_address && entry.quote_address == quote_address;
312 if is_direct {
313 match entry.side {
314 NativeOrderbookSide::Bid => {
315 direct_bids.extend(levels, entry.minimum_in_base)
316 }
317 NativeOrderbookSide::Ask => {
318 direct_asks.extend(levels, entry.minimum_in_base)
319 }
320 }
321 } else {
322 let levels = NativePriceData::invert_price_levels(&levels);
323 match entry.side {
324 NativeOrderbookSide::Bid => {
325 mirrored_asks.extend(levels, entry.minimum_in_base)
326 }
327 NativeOrderbookSide::Ask => {
328 mirrored_bids.extend(levels, entry.minimum_in_base)
329 }
330 }
331 }
332 }
333
334 let (bids, minimum_in_base, minimum_out_quote) = if direct_bids.levels.is_empty() {
338 (mirrored_bids.levels, 0.0, mirrored_bids.minimum)
339 } else {
340 (direct_bids.levels, direct_bids.minimum, 0.0)
341 };
342 let (asks, minimum_in_quote, minimum_out_base) = if direct_asks.levels.is_empty() {
343 (mirrored_asks.levels, mirrored_asks.minimum, 0.0)
344 } else {
345 (direct_asks.levels, 0.0, direct_asks.minimum)
346 };
347 let pair = format!("native_{}/{}", hex::encode(&token0), hex::encode(&token1));
350 let component_id = keccak256(pair.as_bytes()).to_string();
351 books.insert(
352 component_id,
353 NativePriceData {
354 base_address,
355 quote_address,
356 minimum_in_base,
357 minimum_in_quote,
358 minimum_out_base,
359 minimum_out_quote,
360 bids,
361 asks,
362 },
363 );
364 }
365
366 books
367 }
368
369 fn process_quote_response(
370 quote_response: FirmQuoteResponse,
371 params: &GetAmountOutParams,
372 ) -> Result<SignedQuote, RFQError> {
373 if !quote_response.success {
375 return Err(RFQError::QuoteNotFound(format!(
376 "Native Relay quote request failed: {}",
377 quote_response.error_message
378 )));
379 }
380
381 if quote_response.router_version != "6" {
382 return Err(RFQError::ParsingError(format!(
383 "Unexpected Native router version: expected 6, got {}",
384 quote_response.router_version
385 )));
386 }
387
388 let order = quote_response
390 .orders
391 .first()
392 .ok_or_else(|| {
393 RFQError::QuoteNotFound(format!(
394 "No Native Relay orders for {} {} -> {}",
395 params.amount_in, params.token_in, params.token_out,
396 ))
397 })?;
398
399 let seller_token = bytes_to_address(¶ms.token_in)?;
401 let buyer_token = bytes_to_address(¶ms.token_out)?;
402 let order_seller_token = Address::from_str(&order.seller_token).map_err(|e| {
403 RFQError::ParsingError(format!(
404 "Invalid Native seller token {}: {e}",
405 order.seller_token
406 ))
407 })?;
408 let order_buyer_token = Address::from_str(&order.buyer_token).map_err(|e| {
409 RFQError::ParsingError(format!("Invalid Native buyer token {}: {e}", order.buyer_token))
410 })?;
411 if order_seller_token != seller_token || order_buyer_token != buyer_token {
412 return Err(RFQError::ParsingError(format!(
413 "Native Relay quote token mismatch: expected {}/{}, got {}/{}",
414 seller_token, buyer_token, order_seller_token, order_buyer_token
415 )));
416 }
417
418 let receiver = bytes_to_address(¶ms.receiver)?;
419 let order_recipient = Address::from_str(&order.recipient).map_err(|e| {
420 RFQError::ParsingError(format!(
421 "Invalid Native order recipient {}: {e}",
422 order.recipient
423 ))
424 })?;
425 if order_recipient != receiver {
426 return Err(RFQError::ParsingError(format!(
427 "Native Relay quote recipient mismatch: expected {receiver}, got {order_recipient}"
428 )));
429 }
430
431 let now = SystemTime::now()
433 .duration_since(SystemTime::UNIX_EPOCH)
434 .map_err(|_| RFQError::ParsingError("SystemTime before UNIX EPOCH!".to_string()))?
435 .as_secs();
436
437 if order.deadline_timestamp <= now {
438 return Err(RFQError::QuoteNotFound(format!(
439 "Native Relay quote already expired: deadline {} <= now {}",
440 order.deadline_timestamp, now
441 )));
442 }
443
444 let quoted_amount_in = BigUint::from_str("e_response.amount_in).map_err(|_| {
448 RFQError::ParsingError(format!(
449 "Failed to parse amount_in: {}",
450 quote_response.amount_in
451 ))
452 })?;
453 if quoted_amount_in != params.amount_in {
454 return Err(RFQError::ParsingError(format!(
455 "Native Relay quote input amount mismatch: expected {}, got {}",
456 params.amount_in, quoted_amount_in
457 )));
458 }
459
460 let signed_amount_in = BigUint::from_str(&order.seller_token_amount).map_err(|_| {
461 RFQError::ParsingError(format!(
462 "Failed to parse signed seller token amount: {}",
463 order.seller_token_amount
464 ))
465 })?;
466 if signed_amount_in != params.amount_in {
467 return Err(RFQError::ParsingError(format!(
468 "Native Relay signed input amount mismatch: expected {}, got {}",
469 params.amount_in, signed_amount_in
470 )));
471 }
472
473 let quoted_value = BigUint::from_str("e_response.tx_request.value).map_err(|_| {
480 RFQError::ParsingError(format!(
481 "Failed to parse Native txRequest.value: {}",
482 quote_response.tx_request.value
483 ))
484 })?;
485 let expected_value =
486 if seller_token == Address::ZERO { quoted_amount_in.clone() } else { BigUint::ZERO };
487 if quoted_value != expected_value {
488 return Err(RFQError::ParsingError(format!(
489 "Native Relay payable value mismatch: expected {}, got {}",
490 expected_value, quoted_value
491 )));
492 }
493
494 if quote_response
495 .tx_request
496 .calldata
497 .is_empty()
498 {
499 return Err(RFQError::QuoteNotFound(
500 "Native Relay quote did not include calldata".to_string(),
501 ));
502 }
503 let calldata = hex::decode(
505 quote_response
506 .tx_request
507 .calldata
508 .trim_start_matches("0x"),
509 )
510 .map_err(|e| RFQError::ParsingError(format!("Failed to decode calldata: {e}")))?;
511
512 if calldata.len() < MIN_TRADE_RFQT_CALLDATA_LEN {
513 return Err(RFQError::ParsingError(format!(
514 "Native tradeRFQT calldata too short: expected at least {} bytes, got {}",
515 MIN_TRADE_RFQT_CALLDATA_LEN,
516 calldata.len()
517 )));
518 }
519
520 if calldata[..TRADE_RFQT_SELECTOR.len()] != TRADE_RFQT_SELECTOR {
521 return Err(RFQError::ParsingError(format!(
522 "Unexpected Native V6 selector: expected 0x{}, got 0x{}",
523 hex::encode(TRADE_RFQT_SELECTOR),
524 hex::encode(&calldata[..TRADE_RFQT_SELECTOR.len()]),
525 )));
526 }
527
528 if quote_response.amount_in_offset as usize != ACTUAL_SELLER_AMOUNT_OFFSET ||
533 quote_response.amount_out_minimum_offset as usize != ACTUAL_MIN_OUTPUT_AMOUNT_OFFSET
534 {
535 return Err(RFQError::ParsingError(format!(
536 "Unexpected Native V6 override offsets: expected {}/{} but got {}/{}",
537 ACTUAL_SELLER_AMOUNT_OFFSET,
538 ACTUAL_MIN_OUTPUT_AMOUNT_OFFSET,
539 quote_response.amount_in_offset,
540 quote_response.amount_out_minimum_offset,
541 )));
542 }
543
544 if calldata[ACTUAL_SELLER_AMOUNT_OFFSET..ACTUAL_MIN_OUTPUT_AMOUNT_OFFSET]
545 .iter()
546 .any(|byte| *byte != 0)
547 {
548 return Err(RFQError::ParsingError(
549 "Native actualSellerAmount override must be zero".to_string(),
550 ));
551 }
552 if calldata[ACTUAL_MIN_OUTPUT_AMOUNT_OFFSET..MIN_TRADE_RFQT_CALLDATA_LEN]
553 .iter()
554 .any(|byte| *byte != 0)
555 {
556 return Err(RFQError::ParsingError(
557 "Native actualMinOutputAmount override must be zero".to_string(),
558 ));
559 }
560
561 let target = Bytes::from_str("e_response.tx_request.target).map_err(|_| {
562 RFQError::ParsingError(format!(
563 "Failed to parse router target address: {}",
564 quote_response.tx_request.target
565 ))
566 })?;
567
568 let mut quote_attributes: HashMap<String, Bytes> = HashMap::new();
569 quote_attributes.insert("target".to_string(), target);
570 quote_attributes.insert("calldata".to_string(), Bytes::from(calldata));
571 quote_attributes.insert(
572 "deadline_timestamp".to_string(),
573 Bytes::from(
574 order
575 .deadline_timestamp
576 .to_be_bytes()
577 .to_vec(),
578 ),
579 );
580
581 Ok(SignedQuote {
582 base_token: params.token_in.clone(),
583 quote_token: params.token_out.clone(),
584 amount_in: quoted_amount_in,
585 amount_out: BigUint::from_str("e_response.amount_out).map_err(|_| {
586 RFQError::ParsingError(format!(
587 "Failed to parse amount_out: {}",
588 quote_response.amount_out
589 ))
590 })?,
591 quote_attributes,
592 })
593 }
594
595 async fn try_quote(
596 &self,
597 request_data: &FirmQuoteRequest,
598 params: &GetAmountOutParams,
599 ) -> Result<SignedQuote, QuoteAttemptError> {
600 let response = NATIVE_HTTP_CLIENT
601 .get(format!("{}/firm-quote", self.endpoint))
602 .query(request_data)
603 .header("apikey", &self.api_key)
604 .send()
605 .await
606 .map_err(|e| QuoteAttemptError::Retry {
607 error: RFQError::ConnectionError(format!(
608 "Failed to make Native quote request: {e}"
609 )),
610 delay: TRANSIENT_RETRY_DELAY,
611 })?;
612
613 let status = response.status();
614 let response_body = response
615 .bytes()
616 .await
617 .map_err(|e| QuoteAttemptError::Retry {
618 error: RFQError::ConnectionError(format!(
619 "Failed to read Native quote response: {e}"
620 )),
621 delay: TRANSIENT_RETRY_DELAY,
622 })?;
623
624 if let Ok(api_error) = serde_json::from_slice::<NativeApiErrorResponse>(&response_body) {
626 return Err(Self::classify_api_error(&api_error));
627 }
628
629 if !status.is_success() {
630 let response_text = String::from_utf8_lossy(&response_body);
631 if status.is_server_error() {
632 return Err(QuoteAttemptError::Retry {
633 error: RFQError::ConnectionError(format!(
634 "Native quote server error ({status}): {response_text}"
635 )),
636 delay: TRANSIENT_RETRY_DELAY,
637 });
638 }
639
640 return Err(QuoteAttemptError::Fatal(RFQError::ConnectionError(format!(
641 "Unexpected Native quote HTTP response ({status}): {response_text}"
642 ))));
643 }
644
645 let quote_response =
646 serde_json::from_slice::<FirmQuoteResponse>(&response_body).map_err(|e| {
647 QuoteAttemptError::Retry {
648 error: RFQError::ParsingError(format!(
649 "Failed to parse Native quote response: {e}"
650 )),
651 delay: TRANSIENT_RETRY_DELAY,
652 }
653 })?;
654
655 Self::process_quote_response(quote_response, params).map_err(QuoteAttemptError::Fatal)
656 }
657}
658
659#[async_trait]
660impl RFQClient for NativeClient {
661 fn stream(
662 &self,
663 ) -> BoxStream<'static, Result<(String, StateSyncMessage<TimestampHeader>), RFQError>> {
664 let client = self.clone();
665
666 Box::pin(async_stream::stream! {
667 let mut current_components: HashMap<String, ComponentWithState> = HashMap::new();
668 let mut ticker = interval(client.poll_time);
669
670 loop {
671 ticker.tick().await;
672
673 let books = match client.fetch_orderbook().await {
677 Ok(entries) => client.group_orderbook(entries),
678 Err(e) => {
679 error!("Failed to fetch Native Relay orderbook: {}", e);
680 continue;
681 }
682 };
683
684 let mut new_components = HashMap::new();
685
686 for (component_id, book) in &books {
687 if !client.tokens.contains(&book.base_address) ||
690 !client.tokens.contains(&book.quote_address)
691 {
692 continue;
693 }
694
695 let quote_price_data = if client.quote_tokens.contains(&book.quote_address) {
696 None
697 } else {
698 client.select_tvl_conversion_book(&book.quote_address, &books)
702 };
703
704 if !client.quote_tokens.contains(&book.quote_address) &&
705 quote_price_data.is_none()
706 {
707 continue;
708 }
709
710 let Some(incoming_tvl) = book.calculate_tvl(quote_price_data) else {
711 warn!("Skipping Native Relay market {component_id} because its TVL is unavailable or non-finite");
712 continue;
713 };
714
715 if incoming_tvl < client.tvl {
716 info!("Filtering out Native Relay market {} due to low TVL: {:.2} < {:.2}", component_id, incoming_tvl, client.tvl);
717 continue;
718 }
719
720 let tokens = vec![book.base_address.clone(), book.quote_address.clone()];
721 let component_with_state = client.create_component_with_state(
722 component_id.clone(),
723 tokens,
724 book.clone(),
725 incoming_tvl,
726 );
727 new_components.insert(component_id.clone(), component_with_state);
728 }
729
730 let removed_components: HashMap<String, ProtocolComponent> = current_components
733 .iter()
734 .filter(|&(id, _)| !new_components.contains_key(id))
735 .map(|(k, v)| (k.clone(), v.component.clone()))
736 .collect();
737
738 current_components = new_components.clone();
739
740 let snapshot = Snapshot {
741 states: new_components,
742 vm_storage: HashMap::new(),
743 };
744
745 let timestamp = SystemTime::now()
748 .duration_since(SystemTime::UNIX_EPOCH)
749 .unwrap_or_default()
750 .as_secs();
751
752 let msg = StateSyncMessage::<TimestampHeader> {
753 header: TimestampHeader { timestamp },
754 snapshots: snapshot,
755 deltas: None,
756 removed_components,
757 };
758
759 yield Ok(("native".to_string(), msg));
760 }
761 })
762 }
763
764 async fn request_binding_quote(
765 &self,
766 params: &GetAmountOutParams,
767 ) -> Result<SignedQuote, RFQError> {
768 let receiver = bytes_to_address(¶ms.receiver)?;
769 let token_in = bytes_to_address(¶ms.token_in)?;
770 let token_out = bytes_to_address(¶ms.token_out)?;
771
772 let chain = NativeSupportedChain::try_from(self.chain).map_err(RFQError::FatalError)?;
773
774 let request_data = FirmQuoteRequest {
775 src_chain: chain,
776 dst_chain: chain,
777 from_address: receiver.to_string(),
778 amount_wei: params.amount_in.to_string(),
779 token_in: token_in.to_string(),
780 token_out: token_out.to_string(),
781 version: 6,
782 allow_multihop: false,
783 };
784 let mut last_error = None;
785
786 let attempts = async {
787 for attempt in 1..=MAX_QUOTE_ATTEMPTS {
788 match self
789 .try_quote(&request_data, params)
790 .await
791 {
792 Ok(quote) => return Ok(quote),
793 Err(QuoteAttemptError::Fatal(error)) => return Err(error),
794 Err(QuoteAttemptError::Retry { error, delay }) => {
795 warn!(
796 "Native quote attempt {}/{} failed: {}",
797 attempt, MAX_QUOTE_ATTEMPTS, error
798 );
799 last_error = Some(error);
800
801 if attempt < MAX_QUOTE_ATTEMPTS {
802 tokio::time::sleep(delay).await;
803 }
804 }
805 }
806 }
807
808 Err(last_error.take().unwrap_or_else(|| {
809 RFQError::ConnectionError(
810 "Native quote request failed after all attempts".to_string(),
811 )
812 }))
813 };
814
815 let result = timeout(self.quote_timeout, attempts).await;
818 match result {
819 Ok(result) => result,
820 Err(_) => Err(last_error.unwrap_or_else(|| {
821 RFQError::ConnectionError(format!(
822 "Native quote request timed out after {:?}",
823 self.quote_timeout
824 ))
825 })),
826 }
827 }
828}
829
830#[cfg(test)]
831mod tests {
832 use std::{
833 collections::{HashMap, HashSet},
834 str::FromStr,
835 sync::{
836 atomic::{AtomicUsize, Ordering},
837 Arc,
838 },
839 };
840
841 use futures::StreamExt;
842 use rstest::rstest;
843 use tokio::{
844 io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
845 net::TcpListener,
846 };
847 use tycho_common::models::Chain;
848
849 use super::*;
850 use crate::rfq::protocols::native::client_builder::NativeClientBuilder;
851
852 fn successful_quote_json(amount_in: &str) -> serde_json::Value {
853 let calldata = format!("0x7083527c{:064x}{:064x}{:064x}", 0x60u8, 0u8, 0u8);
854 serde_json::json!({
855 "success": true,
856 "orders": [{
857 "pool": "0x1111111111111111111111111111111111111111",
858 "signer": "0x2222222222222222222222222222222222222222",
859 "recipient": "0x4444444444444444444444444444444444444444",
860 "sellerToken": "0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2",
861 "buyerToken": "0xdAC17F958D2ee523a2206206994597C13D831ec7",
862 "effectiveSellerTokenAmount": amount_in,
863 "sellerTokenAmount": amount_in,
864 "buyerTokenAmount": "2",
865 "deadlineTimestamp": u64::MAX,
866 "nonce": 1,
867 "quoteId": "test-quote",
868 "multiHop": false,
869 "signature": "",
870 "externalSwapCalldata": "",
871 "amountOutMinimum": "2",
872 "widgetFee": {
873 "signer": "0x0000000000000000000000000000000000000000",
874 "feeRecipient": "0x0000000000000000000000000000000000000000",
875 "feeRate": 0.0
876 },
877 "widgetFeeSignature": ""
878 }],
879 "widgetFee": {
880 "signer": "0x0000000000000000000000000000000000000000",
881 "feeRecipient": "0x0000000000000000000000000000000000000000",
882 "feeRate": 0.0
883 },
884 "widgetFeeSignature": "",
885 "recipient": "0x4444444444444444444444444444444444444444",
886 "amountIn": amount_in,
887 "amountOut": "2",
888 "amountOutBeforeFee": "2",
889 "fallbackSwapDataArray": null,
890 "tokenTransferFeeOnPercent": 0.0,
891 "txRequest": {
892 "target": "0x4777A6B3A9A889ABfd4C7666Bdd2a7AB633293be",
893 "calldata": calldata,
894 "value": "0"
895 },
896 "source": [6],
897 "errorMessage": "",
898 "router_version": "6",
899 "toWrap": false,
900 "toUnwrap": false,
901 "amountInOffset": 36,
902 "amountOutMinimumOffset": 68
903 })
904 }
905
906 fn successful_quote_response(amount_in: &str) -> FirmQuoteResponse {
907 serde_json::from_value(successful_quote_json(amount_in)).unwrap()
908 }
909
910 fn conversion_book(
911 base_address: Bytes,
912 quote_address: Bytes,
913 quantity: f64,
914 price: f64,
915 ) -> NativePriceData {
916 NativePriceData {
917 base_address,
918 quote_address,
919 minimum_in_base: 0.0,
920 minimum_in_quote: 0.0,
921 minimum_out_base: 0.0,
922 minimum_out_quote: 0.0,
923 bids: vec![NativePriceLevel { quantity, price }],
924 asks: vec![],
925 }
926 }
927
928 #[test]
929 fn test_native_client_serialization() {
930 let mut tokens = HashSet::new();
931 tokens.insert(Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap());
932 tokens.insert(Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap());
933
934 let client = NativeClient::new(
935 Chain::Ethereum,
936 "test-api-key".to_string(),
937 tokens,
938 10.0,
939 HashSet::new(),
940 Duration::from_secs(1),
941 Duration::from_secs(5),
942 )
943 .unwrap();
944
945 let serialized = serde_json::to_string(&client).unwrap();
946 let deserialized: NativeClient = serde_json::from_str(&serialized).unwrap();
947
948 assert_eq!(deserialized.chain, client.chain);
949 assert_eq!(deserialized.endpoint, client.endpoint);
950 assert_eq!(deserialized.tokens, client.tokens);
951 assert_eq!(deserialized.tvl, client.tvl);
952 assert!(deserialized.api_key.is_empty());
953 }
954
955 #[rstest]
956 #[case::ethereum(Chain::Ethereum, "ethereum")]
957 #[case::bsc(Chain::Bsc, "bsc")]
958 #[case::arbitrum(Chain::Arbitrum, "arbitrum")]
959 #[case::base(Chain::Base, "base")]
960 #[case::robinhood(Chain::Robinhood, "robinhood")]
961 fn accepts_every_chain_native_serves(#[case] chain: Chain, #[case] native_name: &str) {
962 let client = NativeClient::new(
963 chain,
964 "test-api-key".to_string(),
965 HashSet::new(),
966 0.0,
967 HashSet::new(),
968 Duration::from_secs(1),
969 Duration::from_secs(5),
970 );
971
972 assert!(client.is_ok());
973 assert_eq!(
974 NativeSupportedChain::try_from(chain)
975 .unwrap()
976 .as_str(),
977 native_name
978 );
979 }
980
981 #[test]
982 fn rejects_unsupported_chain_at_construction() {
983 let result = NativeClient::new(
984 Chain::Polygon,
985 "test-api-key".to_string(),
986 HashSet::new(),
987 0.0,
988 HashSet::new(),
989 Duration::from_secs(1),
990 Duration::from_secs(5),
991 );
992
993 assert!(matches!(result, Err(RFQError::InvalidInput(_))));
994 }
995
996 #[test]
997 fn builder_rejects_zero_poll_time() {
998 let result = NativeClientBuilder::new(Chain::Ethereum, "test-api-key".to_string())
999 .poll_time(Duration::ZERO)
1000 .build();
1001
1002 assert!(matches!(
1003 result,
1004 Err(RFQError::InvalidInput(message))
1005 if message == "Native polling interval must be greater than zero"
1006 ));
1007 }
1008
1009 #[test]
1010 fn selects_most_liquid_tvl_conversion_book() {
1011 let weth = Bytes::from_str("0x3333333333333333333333333333333333333333").unwrap();
1012 let usdc = Bytes::from_str("0x1111111111111111111111111111111111111111").unwrap();
1013 let usdt = Bytes::from_str("0x2222222222222222222222222222222222222222").unwrap();
1014 let wbtc = Bytes::from_str("0x4444444444444444444444444444444444444444").unwrap();
1015 let unapproved = Bytes::from_str("0x5555555555555555555555555555555555555555").unwrap();
1016 let client = NativeClient::new(
1017 Chain::Ethereum,
1018 "test-api-key".to_string(),
1019 HashSet::from([
1020 weth.clone(),
1021 usdc.clone(),
1022 usdt.clone(),
1023 wbtc.clone(),
1024 unapproved.clone(),
1025 ]),
1026 0.0,
1027 HashSet::from([usdc.clone(), usdt.clone()]),
1028 Duration::from_secs(1),
1029 Duration::from_secs(5),
1030 )
1031 .unwrap();
1032 let books = HashMap::from([
1033 ("usdc".to_string(), conversion_book(weth.clone(), usdc.clone(), 1.0, 100.0)),
1034 ("usdt".to_string(), conversion_book(weth.clone(), usdt.clone(), 2.0, 100.0)),
1035 ("unrelated".to_string(), conversion_book(wbtc, usdc, 1_000.0, 100.0)),
1036 ("unapproved".to_string(), conversion_book(weth.clone(), unapproved, 2_000.0, 100.0)),
1037 ]);
1038
1039 let selected = client
1040 .select_tvl_conversion_book(&weth, &books)
1041 .expect("one conversion book");
1042
1043 assert_eq!(selected.quote_address, usdt);
1044 }
1045
1046 #[test]
1047 fn selects_lower_quote_address_for_equal_tvl_conversion_liquidity() {
1048 let weth = Bytes::from_str("0x3333333333333333333333333333333333333333").unwrap();
1049 let lower_quote = Bytes::from_str("0x1111111111111111111111111111111111111111").unwrap();
1050 let higher_quote = Bytes::from_str("0x2222222222222222222222222222222222222222").unwrap();
1051 let client = NativeClient::new(
1052 Chain::Ethereum,
1053 "test-api-key".to_string(),
1054 HashSet::from([weth.clone(), lower_quote.clone(), higher_quote.clone()]),
1055 0.0,
1056 HashSet::from([lower_quote.clone(), higher_quote.clone()]),
1057 Duration::from_secs(1),
1058 Duration::from_secs(5),
1059 )
1060 .unwrap();
1061 let books = HashMap::from([
1062 ("higher".to_string(), conversion_book(weth.clone(), higher_quote, 2.0, 100.0)),
1063 ("lower".to_string(), conversion_book(weth.clone(), lower_quote.clone(), 1.0, 200.0)),
1064 ]);
1065
1066 let selected = client
1067 .select_tvl_conversion_book(&weth, &books)
1068 .expect("one conversion book");
1069
1070 assert_eq!(selected.quote_address, lower_quote);
1071 }
1072
1073 #[test]
1074 fn creates_indexer_compatible_component_from_relay_orderbook() {
1075 let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1076 let usdt = Bytes::from_str("0xdac17f958d2ee523a2206206994597c13d831ec7").unwrap();
1077 let client = NativeClient::new(
1078 Chain::Ethereum,
1079 "test-api-key".to_string(),
1080 HashSet::from([weth.clone(), usdt.clone()]),
1081 0.0,
1082 HashSet::from([usdt.clone()]),
1083 Duration::from_secs(1),
1084 Duration::from_secs(5),
1085 )
1086 .unwrap();
1087
1088 let books = client.group_orderbook(vec![
1089 NativeOrderbookEntry {
1090 base_address: weth.clone(),
1091 quote_address: usdt.clone(),
1092 minimum_in_base: 0.0,
1093 side: NativeOrderbookSide::Bid,
1094 levels: vec![NativePriceLevel { quantity: 0.0001, price: 3213.12345 }],
1095 },
1096 NativeOrderbookEntry {
1097 base_address: weth.clone(),
1098 quote_address: usdt.clone(),
1099 minimum_in_base: 0.0,
1100 side: NativeOrderbookSide::Ask,
1101 levels: vec![NativePriceLevel { quantity: 2.0, price: 3214.0 }],
1102 },
1103 NativeOrderbookEntry {
1104 base_address: usdt.clone(),
1105 quote_address: weth.clone(),
1106 minimum_in_base: 100.0,
1107 side: NativeOrderbookSide::Bid,
1108 levels: vec![NativePriceLevel { quantity: 6428.0, price: 1.0 / 3214.0 }],
1109 },
1110 NativeOrderbookEntry {
1111 base_address: usdt.clone(),
1112 quote_address: weth.clone(),
1113 minimum_in_base: 100.0,
1114 side: NativeOrderbookSide::Ask,
1115 levels: vec![NativePriceLevel { quantity: 0.321312345, price: 1.0 / 3213.12345 }],
1116 },
1117 ]);
1118
1119 let (component_id, book) = books
1120 .into_iter()
1121 .next()
1122 .expect("one grouped book");
1123 let component = client.create_component_with_state(
1124 component_id.clone(),
1125 vec![book.base_address.clone(), book.quote_address.clone()],
1126 book.clone(),
1127 book.calculate_tvl(None)
1128 .expect("TVL should be finite"),
1129 );
1130
1131 assert_eq!(component.component.id, component_id);
1132 assert_eq!(component.component.protocol_system, NativeClient::PROTOCOL_SYSTEM);
1133 assert_eq!(component.component.protocol_type_name, "native_relay_pool");
1134 assert_eq!(component.component.tokens, vec![weth, usdt]);
1135 assert_eq!(component.state.component_id, component_id);
1136
1137 let encoded_book = component
1138 .state
1139 .attributes
1140 .get("book")
1141 .expect("book attribute");
1142 let decoded_book: NativePriceData = serde_json::from_slice(encoded_book).unwrap();
1143 assert_eq!(decoded_book.bids.len(), 1);
1144 assert_eq!(decoded_book.asks.len(), 1);
1145 assert_eq!(decoded_book.bids[0].quantity, 0.0001);
1146 assert_eq!(decoded_book.bids[0].price, 3213.12345);
1147 }
1148
1149 #[test]
1150 fn uses_stable_component_id_and_direction_when_merging_mirrored_books() {
1151 let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1152 let usdc = Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap();
1153 let client = NativeClient::new(
1154 Chain::Ethereum,
1155 "test-api-key".to_string(),
1156 HashSet::from([weth.clone(), usdc.clone()]),
1157 0.0,
1158 HashSet::from([usdc.clone()]),
1159 Duration::from_secs(1),
1160 Duration::from_secs(5),
1161 )
1162 .unwrap();
1163 let entries = vec![
1164 NativeOrderbookEntry {
1165 base_address: weth.clone(),
1166 quote_address: usdc.clone(),
1167 minimum_in_base: 100_000_000_000.0,
1168 side: NativeOrderbookSide::Bid,
1169 levels: vec![NativePriceLevel { quantity: 1.0, price: 2_000.0 }],
1170 },
1171 NativeOrderbookEntry {
1172 base_address: weth.clone(),
1173 quote_address: usdc.clone(),
1174 minimum_in_base: 300_000_000_000.0,
1175 side: NativeOrderbookSide::Ask,
1176 levels: vec![NativePriceLevel { quantity: 1.0, price: 2_100.0 }],
1177 },
1178 NativeOrderbookEntry {
1179 base_address: usdc.clone(),
1180 quote_address: weth.clone(),
1181 minimum_in_base: 100.0,
1182 side: NativeOrderbookSide::Bid,
1183 levels: vec![NativePriceLevel { quantity: 2_000.0, price: 0.0005 }],
1184 },
1185 NativeOrderbookEntry {
1186 base_address: usdc.clone(),
1187 quote_address: weth.clone(),
1188 minimum_in_base: 250.0,
1189 side: NativeOrderbookSide::Bid,
1190 levels: vec![NativePriceLevel { quantity: 2_000.0, price: 0.0005 }],
1191 },
1192 NativeOrderbookEntry {
1193 base_address: usdc.clone(),
1194 quote_address: weth.clone(),
1195 minimum_in_base: 400.0,
1196 side: NativeOrderbookSide::Ask,
1197 levels: vec![NativePriceLevel { quantity: 2.0, price: 0.5 }],
1198 },
1199 ];
1200
1201 let forward_only = client.group_orderbook(entries[..2].to_vec());
1202 let reverse_entries = entries[2..].to_vec();
1203 let reverse_only = client.group_orderbook(reverse_entries.clone());
1204 let reversed_reverse_only = client.group_orderbook(
1205 reverse_entries
1206 .into_iter()
1207 .rev()
1208 .collect(),
1209 );
1210 assert_eq!(reverse_only, reversed_reverse_only);
1211 let pair = format!("native_{}/{}", hex::encode(&usdc), hex::encode(&weth));
1212 let component_id = keccak256(pair.as_bytes()).to_string();
1213
1214 let forward_book = forward_only
1215 .get(&component_id)
1216 .expect("forward-only book uses the stable component ID");
1217 assert_eq!(forward_book.base_address, weth);
1218 assert_eq!(forward_book.quote_address, usdc);
1219 assert_eq!(forward_book.minimum_in_base, 100_000_000_000.0);
1220 assert_eq!(forward_book.minimum_in_quote, 0.0);
1221 assert_eq!(forward_book.minimum_out_base, 300_000_000_000.0);
1222 assert_eq!(forward_book.minimum_out_quote, 0.0);
1223
1224 let reverse_book = reverse_only
1225 .get(&component_id)
1226 .expect("reverse-only book uses the stable component ID");
1227 assert_eq!(reverse_book.base_address, weth);
1228 assert_eq!(reverse_book.quote_address, usdc);
1229 assert_eq!(reverse_book.minimum_in_base, 0.0);
1230 assert_eq!(reverse_book.minimum_in_quote, 250.0);
1231 assert_eq!(reverse_book.minimum_out_base, 0.0);
1232 assert_eq!(reverse_book.minimum_out_quote, 400.0);
1233 assert_eq!(reverse_book.bids, vec![NativePriceLevel { quantity: 1.0, price: 2.0 }]);
1234 assert_eq!(
1235 reverse_book.asks,
1236 vec![
1237 NativePriceLevel { quantity: 1.0, price: 2_000.0 },
1238 NativePriceLevel { quantity: 1.0, price: 2_000.0 },
1239 ]
1240 );
1241
1242 let mixed = client.group_orderbook(vec![entries[0].clone(), entries[2].clone()]);
1243 let mixed_book = mixed.get(&component_id).unwrap();
1244 assert_eq!(mixed_book.minimum_in_base, 100_000_000_000.0);
1245 assert_eq!(mixed_book.minimum_in_quote, 100.0);
1246 assert_eq!(mixed_book.minimum_out_base, 0.0);
1247 assert_eq!(mixed_book.minimum_out_quote, 0.0);
1248 assert_eq!(mixed_book.bids, vec![NativePriceLevel { quantity: 1.0, price: 2_000.0 }]);
1249 assert_eq!(mixed_book.asks, vec![NativePriceLevel { quantity: 1.0, price: 2_000.0 }]);
1250
1251 let books = client.group_orderbook(entries.clone());
1252 let reversed_books = client.group_orderbook(entries.into_iter().rev().collect());
1253
1254 assert_eq!(books, reversed_books);
1255 assert_eq!(books.len(), 1);
1256 let book = books.get(&component_id).unwrap();
1257 assert_eq!(book.base_address, weth);
1258 assert_eq!(book.quote_address, usdc);
1259 assert_eq!(book.minimum_in_base, 100_000_000_000.0);
1260 assert_eq!(book.minimum_in_quote, 0.0);
1261 assert_eq!(book.minimum_out_base, 300_000_000_000.0);
1262 assert_eq!(book.minimum_out_quote, 0.0);
1263 assert_eq!(book.bids, vec![NativePriceLevel { quantity: 1.0, price: 2_000.0 }]);
1264 assert_eq!(book.asks, vec![NativePriceLevel { quantity: 1.0, price: 2_100.0 }]);
1265 }
1266
1267 #[test]
1268 fn zero_only_direct_side_does_not_suppress_mirrored_liquidity() {
1269 let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1270 let usdc = Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap();
1271 let client = NativeClient::new(
1272 Chain::Ethereum,
1273 "test-api-key".to_string(),
1274 HashSet::from([weth.clone(), usdc.clone()]),
1275 0.0,
1276 HashSet::from([usdc.clone()]),
1277 Duration::from_secs(1),
1278 Duration::from_secs(5),
1279 )
1280 .unwrap();
1281
1282 let books = client.group_orderbook(vec![
1283 NativeOrderbookEntry {
1284 base_address: weth.clone(),
1285 quote_address: usdc.clone(),
1286 minimum_in_base: 999.0,
1287 side: NativeOrderbookSide::Bid,
1288 levels: vec![NativePriceLevel { quantity: 0.0, price: 2_000.0 }],
1289 },
1290 NativeOrderbookEntry {
1291 base_address: usdc,
1292 quote_address: weth,
1293 minimum_in_base: 250.0,
1294 side: NativeOrderbookSide::Ask,
1295 levels: vec![NativePriceLevel { quantity: 2.0, price: 0.5 }],
1296 },
1297 ]);
1298
1299 let book = books
1300 .values()
1301 .next()
1302 .expect("one grouped book");
1303 assert_eq!(book.bids, vec![NativePriceLevel { quantity: 1.0, price: 2.0 }]);
1304 assert_eq!(book.minimum_in_base, 0.0);
1305 assert_eq!(book.minimum_out_quote, 250.0);
1306 }
1307
1308 #[rstest]
1309 #[case::token0_only(true, false, true)]
1310 #[case::token1_only(false, true, false)]
1311 #[case::both_tokens(true, true, false)]
1312 #[case::neither_token(false, false, false)]
1313 fn selects_stable_direction_for_quote_token_preferences(
1314 #[case] token0_is_quote: bool,
1315 #[case] token1_is_quote: bool,
1316 #[case] expected_quote_is_token0: bool,
1317 ) {
1318 let token0 = Bytes::from_str("0x1111111111111111111111111111111111111111").unwrap();
1319 let token1 = Bytes::from_str("0x2222222222222222222222222222222222222222").unwrap();
1320 assert!(token0.as_ref() < token1.as_ref());
1321
1322 let mut quote_tokens = HashSet::new();
1323 if token0_is_quote {
1324 quote_tokens.insert(token0.clone());
1325 }
1326 if token1_is_quote {
1327 quote_tokens.insert(token1.clone());
1328 }
1329 let client = NativeClient::new(
1330 Chain::Ethereum,
1331 "test-api-key".to_string(),
1332 HashSet::from([token0.clone(), token1.clone()]),
1333 0.0,
1334 quote_tokens,
1335 Duration::from_secs(1),
1336 Duration::from_secs(5),
1337 )
1338 .unwrap();
1339
1340 let books = client.group_orderbook(vec![NativeOrderbookEntry {
1342 base_address: token1.clone(),
1343 quote_address: token0.clone(),
1344 minimum_in_base: 0.0,
1345 side: NativeOrderbookSide::Bid,
1346 levels: vec![NativePriceLevel { quantity: 1.0, price: 1.0 }],
1347 }]);
1348
1349 let book = books
1350 .values()
1351 .next()
1352 .expect("one grouped book");
1353 let (expected_base, expected_quote) =
1354 if expected_quote_is_token0 { (token1, token0) } else { (token0, token1) };
1355 assert_eq!(book.base_address, expected_base);
1356 assert_eq!(book.quote_address, expected_quote);
1357 }
1358
1359 fn create_test_quote_params() -> GetAmountOutParams {
1360 GetAmountOutParams {
1361 amount_in: BigUint::from(1u64),
1362 token_in: Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap(),
1363 token_out: Bytes::from_str("0xdac17f958d2ee523a2206206994597c13d831ec7").unwrap(),
1364 sender: Bytes::from_str("0x3333333333333333333333333333333333333333").unwrap(),
1365 receiver: Bytes::from_str("0x4444444444444444444444444444444444444444").unwrap(),
1366 }
1367 }
1368
1369 #[test]
1370 fn accepts_quote_with_requested_input_amount() {
1371 let params = create_test_quote_params();
1372 let response = successful_quote_response(¶ms.amount_in.to_string());
1373
1374 let quote = NativeClient::process_quote_response(response, ¶ms).unwrap();
1375
1376 assert_eq!(quote.amount_in, params.amount_in);
1377 assert_eq!(
1378 quote
1379 .quote_attributes
1380 .get("deadline_timestamp")
1381 .unwrap()
1382 .as_ref(),
1383 u64::MAX.to_be_bytes()
1384 );
1385 }
1386
1387 #[test]
1388 fn rejects_quote_with_mismatched_input_amount() {
1389 let params = create_test_quote_params();
1390 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1391 response.amount_in = "2".to_string();
1393
1394 let result = NativeClient::process_quote_response(response, ¶ms);
1395
1396 assert!(matches!(
1397 result,
1398 Err(RFQError::ParsingError(message)) if message.contains("Native Relay quote input amount mismatch")
1399 ));
1400 }
1401
1402 #[test]
1403 fn rejects_quote_with_mismatched_signed_input_amount() {
1404 let params = create_test_quote_params();
1405 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1406 response.orders[0].seller_token_amount = "2".to_string();
1407
1408 let result = NativeClient::process_quote_response(response, ¶ms);
1409
1410 assert!(matches!(
1411 result,
1412 Err(RFQError::ParsingError(message)) if message.contains("signed input amount mismatch")
1413 ));
1414 }
1415
1416 #[test]
1417 fn accepts_quote_with_different_effective_input_amount() {
1418 let mut params = create_test_quote_params();
1419 params.amount_in = BigUint::from(100u64);
1420 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1421 response.orders[0].effective_seller_token_amount = "99".to_string();
1422
1423 let quote = NativeClient::process_quote_response(response, ¶ms).unwrap();
1424
1425 assert_eq!(quote.amount_in, params.amount_in);
1426 }
1427
1428 #[rstest]
1429 #[case::seller_token(true)]
1430 #[case::buyer_token(false)]
1431 fn rejects_quote_with_mismatched_token(#[case] mutate_seller_token: bool) {
1432 let params = create_test_quote_params();
1433 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1434 let mismatched_token = "0x5555555555555555555555555555555555555555".to_string();
1435 if mutate_seller_token {
1436 response.orders[0].seller_token = mismatched_token;
1437 } else {
1438 response.orders[0].buyer_token = mismatched_token;
1439 }
1440
1441 let result = NativeClient::process_quote_response(response, ¶ms);
1442
1443 assert!(matches!(
1444 result,
1445 Err(RFQError::ParsingError(message)) if message.contains("quote token mismatch")
1446 ));
1447 }
1448
1449 #[test]
1450 fn rejects_quote_with_mismatched_recipient() {
1451 let params = create_test_quote_params();
1452 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1453 response.orders[0].recipient = "0x5555555555555555555555555555555555555555".to_string();
1454
1455 let result = NativeClient::process_quote_response(response, ¶ms);
1456
1457 assert!(matches!(
1458 result,
1459 Err(RFQError::ParsingError(message)) if message.contains("recipient mismatch")
1460 ));
1461 }
1462
1463 #[test]
1464 fn rejects_quote_without_calldata() {
1465 let params = create_test_quote_params();
1466 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1467 response.tx_request.calldata.clear();
1468
1469 let result = NativeClient::process_quote_response(response, ¶ms);
1470
1471 assert!(matches!(
1472 result,
1473 Err(RFQError::QuoteNotFound(message)) if message.contains("did not include calldata")
1474 ));
1475 }
1476
1477 #[test]
1478 fn rejects_truncated_trade_calldata() {
1479 let params = create_test_quote_params();
1480 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1481 response.tx_request.calldata = format!("0x7083527c{}", "00".repeat(95));
1482
1483 let result = NativeClient::process_quote_response(response, ¶ms);
1484
1485 assert!(matches!(
1486 result,
1487 Err(RFQError::ParsingError(message)) if message.contains("calldata too short")
1488 ));
1489 }
1490
1491 #[test]
1492 fn rejects_quote_with_wrong_trade_rfqt_selector() {
1493 let params = create_test_quote_params();
1494 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1495 let mut calldata = hex::decode(
1496 response
1497 .tx_request
1498 .calldata
1499 .trim_start_matches("0x"),
1500 )
1501 .unwrap();
1502 calldata[..4].copy_from_slice(&[0x09, 0x47, 0xc2, 0xd9]);
1504 response.tx_request.calldata = format!("0x{}", hex::encode(calldata));
1505
1506 let result = NativeClient::process_quote_response(response, ¶ms);
1507
1508 assert!(matches!(
1509 result,
1510 Err(RFQError::ParsingError(message)) if message.contains("selector")
1511 ));
1512 }
1513
1514 #[rstest]
1515 #[case::v4("4")]
1516 #[case::unknown("7")]
1517 fn rejects_quote_with_wrong_router_version(#[case] version: &str) {
1518 let params = create_test_quote_params();
1519 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1520 response.router_version = version.to_string();
1521
1522 assert!(matches!(
1523 NativeClient::process_quote_response(response, ¶ms),
1524 Err(RFQError::ParsingError(message)) if message.contains("Unexpected Native router version")
1525 ));
1526 }
1527
1528 #[rstest]
1529 #[case::seller_offset(68, 68)]
1530 #[case::minimum_offset(36, 36)]
1531 fn rejects_quote_with_noncanonical_override_offsets(
1532 #[case] amount_in_offset: u32,
1533 #[case] amount_out_minimum_offset: u32,
1534 ) {
1535 let params = create_test_quote_params();
1536 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1537 response.amount_in_offset = amount_in_offset;
1538 response.amount_out_minimum_offset = amount_out_minimum_offset;
1539
1540 let result = NativeClient::process_quote_response(response, ¶ms);
1541
1542 assert!(matches!(
1543 result,
1544 Err(RFQError::ParsingError(message)) if message.contains(
1545 "Unexpected Native V6 override offsets"
1546 )
1547 ));
1548 }
1549
1550 #[rstest]
1551 #[case::seller(ACTUAL_SELLER_AMOUNT_OFFSET, "actualSellerAmount")]
1552 #[case::minimum(ACTUAL_MIN_OUTPUT_AMOUNT_OFFSET, "actualMinOutputAmount")]
1553 fn rejects_quote_with_preset_override(
1554 #[case] override_offset: usize,
1555 #[case] field_name: &str,
1556 ) {
1557 let params = create_test_quote_params();
1558 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1559 let mut calldata = hex::decode(
1560 response
1561 .tx_request
1562 .calldata
1563 .trim_start_matches("0x"),
1564 )
1565 .unwrap();
1566 calldata[override_offset + 31] = 1;
1567 response.tx_request.calldata = format!("0x{}", hex::encode(calldata));
1568
1569 let result = NativeClient::process_quote_response(response, ¶ms);
1570
1571 assert!(matches!(
1572 result,
1573 Err(RFQError::ParsingError(message)) if message.contains(field_name)
1574 ));
1575 }
1576
1577 #[test]
1578 fn accepts_native_eth_response_using_zero_address() {
1579 let mut params = create_test_quote_params();
1580 params.token_in = Bytes::zero(20);
1581 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1582 response.orders[0].seller_token = "0x0000000000000000000000000000000000000000".to_string();
1583 response.tx_request.value = params.amount_in.to_string();
1584
1585 let quote = NativeClient::process_quote_response(response, ¶ms).unwrap();
1586
1587 assert_eq!(quote.amount_in, params.amount_in);
1588 }
1589
1590 #[test]
1591 fn rejects_native_quote_with_mismatched_payable_value() {
1592 let mut params = create_test_quote_params();
1593 params.token_in = Bytes::zero(20);
1594 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1595 response.orders[0].seller_token = "0x0000000000000000000000000000000000000000".to_string();
1596 response.tx_request.value = "2".to_string();
1597
1598 let result = NativeClient::process_quote_response(response, ¶ms);
1599
1600 assert!(matches!(
1601 result,
1602 Err(RFQError::ParsingError(message)) if message.contains("payable value mismatch")
1603 ));
1604 }
1605
1606 #[test]
1607 fn rejects_malformed_payable_value() {
1608 let params = create_test_quote_params();
1609 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1610 response.tx_request.value = "not-a-number".to_string();
1611
1612 let result = NativeClient::process_quote_response(response, ¶ms);
1613
1614 assert!(matches!(
1615 result,
1616 Err(RFQError::ParsingError(message)) if message.contains("txRequest.value")
1617 ));
1618 }
1619
1620 #[test]
1621 fn rejects_erc20_quote_with_nonzero_payable_value() {
1622 let params = create_test_quote_params();
1623 let mut response = successful_quote_response(¶ms.amount_in.to_string());
1624 response.tx_request.value = "1".to_string();
1625
1626 let result = NativeClient::process_quote_response(response, ¶ms);
1627
1628 assert!(matches!(
1629 result,
1630 Err(RFQError::ParsingError(message)) if message.contains("payable value mismatch")
1631 ));
1632 }
1633
1634 #[test]
1635 fn accepts_native_eth_orderbook_using_zero_address() {
1636 let tycho_native_eth = Bytes::zero(20);
1637 let usdc = Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap();
1638 let client = NativeClient::new(
1639 Chain::Ethereum,
1640 "test-api-key".to_string(),
1641 HashSet::from([tycho_native_eth.clone(), usdc.clone()]),
1642 0.0,
1643 HashSet::from([usdc.clone()]),
1644 Duration::from_secs(1),
1645 Duration::from_secs(5),
1646 )
1647 .unwrap();
1648
1649 let entry: NativeOrderbookEntry = serde_json::from_value(serde_json::json!({
1650 "base_address": tycho_native_eth.to_string(),
1651 "quote_address": usdc.to_string(),
1652 "minimum_in_base": 1.0,
1653 "side": "bid",
1654 "levels": [[1.0, 3_000.0]]
1655 }))
1656 .unwrap();
1657 let books = client.group_orderbook(vec![entry]);
1658
1659 let book = books.values().next().unwrap();
1660 assert_eq!(book.base_address, tycho_native_eth);
1661 }
1662
1663 fn create_test_client(endpoint: String) -> NativeClient {
1664 let mut client = NativeClient::new(
1665 Chain::Ethereum,
1666 "test-api-key".to_string(),
1667 HashSet::new(),
1668 0.0,
1669 HashSet::new(),
1670 Duration::from_secs(1),
1671 Duration::from_secs(5),
1672 )
1673 .unwrap();
1674 client.endpoint = endpoint;
1675 client
1676 }
1677
1678 #[tokio::test]
1679 async fn requests_v6_firm_quote() {
1680 let listener = TcpListener::bind("127.0.0.1:0")
1681 .await
1682 .unwrap();
1683 let address = listener.local_addr().unwrap();
1684 let server = tokio::spawn(async move {
1685 let (stream, _) = listener.accept().await.unwrap();
1686 let mut reader = BufReader::new(stream);
1687 let mut request_line = String::new();
1688 reader
1689 .read_line(&mut request_line)
1690 .await
1691 .unwrap();
1692 let body = successful_quote_json("1").to_string();
1693 let response = format!(
1694 "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
1695 body.len()
1696 );
1697 reader
1698 .into_inner()
1699 .write_all(response.as_bytes())
1700 .await
1701 .unwrap();
1702 request_line
1703 });
1704 let client = create_test_client(format!("http://{address}"));
1705 let params = create_test_quote_params();
1706
1707 let quote = client
1708 .request_binding_quote(¶ms)
1709 .await
1710 .unwrap();
1711 let request = server.await.unwrap();
1712 let url = reqwest::Url::parse(&format!(
1713 "http://{address}{}",
1714 request
1715 .split_whitespace()
1716 .nth(1)
1717 .unwrap()
1718 ))
1719 .unwrap();
1720 let query: HashMap<_, _> = url.query_pairs().into_owned().collect();
1721
1722 assert_eq!(url.path(), "/firm-quote");
1723 assert_eq!(query.get("version").map(String::as_str), Some("6"));
1724 assert_eq!(
1725 query
1726 .get("allow_multihop")
1727 .map(String::as_str),
1728 Some("false")
1729 );
1730 assert_eq!(quote.amount_in, params.amount_in);
1731 }
1732
1733 #[tokio::test]
1734 async fn requests_native_token_orderbooks() {
1735 let listener = TcpListener::bind("127.0.0.1:0")
1736 .await
1737 .unwrap();
1738 let address = listener.local_addr().unwrap();
1739 let server = tokio::spawn(async move {
1740 let (stream, _) = listener.accept().await.unwrap();
1741 let mut reader = BufReader::new(stream);
1742 let mut request_line = String::new();
1743 reader
1744 .read_line(&mut request_line)
1745 .await
1746 .unwrap();
1747 let mut stream = reader.into_inner();
1748 let response = "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: 2\r\nConnection: close\r\n\r\n[]";
1749 stream
1750 .write_all(response.as_bytes())
1751 .await
1752 .unwrap();
1753 request_line
1754 });
1755 let client = create_test_client(format!("http://{address}"));
1756
1757 let orderbook = client.fetch_orderbook().await.unwrap();
1758 let request = server.await.unwrap();
1759
1760 assert!(orderbook.is_empty());
1761 assert!(request.starts_with("GET /orderbook?"));
1762 assert!(request.contains("chain=ethereum"));
1763 assert!(request.contains("showNative=0x0"));
1764 }
1765
1766 #[rstest]
1767 #[case::with_conversion_helper(true, 300.0, Some(400.0))]
1768 #[case::without_conversion_helper(false, 300.0, None)]
1769 #[case::below_normalized_tvl_threshold(true, 401.0, None)]
1770 #[tokio::test]
1771 async fn stream_uses_unrequested_books_only_for_tvl_conversion(
1772 #[case] include_helper: bool,
1773 #[case] tvl_threshold: f64,
1774 #[case] expected_tvl: Option<f64>,
1775 ) {
1776 let weth = Bytes::from_str("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").unwrap();
1777 let usdt = Bytes::from_str("0xdac17f958d2ee523a2206206994597c13d831ec7").unwrap();
1778 let usdc = Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap();
1779 let mut entries = vec![serde_json::json!({
1780 "base_address": weth.to_string(),
1781 "quote_address": usdt.to_string(),
1782 "minimum_in_base": 0.0,
1783 "side": "bid",
1784 "levels": [[2.0, 100.0]]
1785 })];
1786 if include_helper {
1787 entries.push(serde_json::json!({
1790 "base_address": usdc.to_string(),
1791 "quote_address": usdt.to_string(),
1792 "minimum_in_base": 0.0,
1793 "side": "bid",
1794 "levels": [[1_000.0, 0.5]]
1795 }));
1796 }
1797 let (address, _) =
1798 create_quote_server("200 OK", serde_json::to_string(&entries).unwrap()).await;
1799 let mut client = create_test_client(format!("http://{address}"));
1800 client.tokens = HashSet::from([weth.clone(), usdt.clone()]);
1801 client.quote_tokens = HashSet::from([usdc]);
1802 client.tvl = tvl_threshold;
1803
1804 let (_, update) = timeout(Duration::from_secs(5), client.stream().next())
1805 .await
1806 .expect("orderbook poll timed out")
1807 .expect("stream ended")
1808 .expect("orderbook poll failed");
1809
1810 if let Some(tvl) = expected_tvl {
1811 assert_eq!(update.snapshots.states.len(), 1, "helper must not be emitted");
1812 let component = update
1813 .snapshots
1814 .states
1815 .values()
1816 .next()
1817 .unwrap();
1818 assert_eq!(component.component.tokens, vec![weth, usdt]);
1819 assert_eq!(component.component_tvl, Some(tvl));
1820 } else {
1821 assert!(update.snapshots.states.is_empty());
1822 }
1823 }
1824
1825 async fn create_quote_server(
1826 final_status: impl Into<String>,
1827 final_body: impl Into<String>,
1828 ) -> (std::net::SocketAddr, Arc<AtomicUsize>) {
1829 let final_status = final_status.into();
1830 let final_body = final_body.into();
1831 let listener = TcpListener::bind("127.0.0.1:0")
1832 .await
1833 .unwrap();
1834 let address = listener.local_addr().unwrap();
1835 let request_count = Arc::new(AtomicUsize::new(0));
1836 let server_request_count = request_count.clone();
1837
1838 tokio::spawn(async move {
1839 while let Ok((mut stream, _)) = listener.accept().await {
1840 server_request_count.fetch_add(1, Ordering::SeqCst);
1841 let response = format!(
1842 "HTTP/1.1 {final_status}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{final_body}",
1843 final_body.len()
1844 );
1845 let _ = stream
1846 .write_all(response.as_bytes())
1847 .await;
1848 let _ = stream.shutdown().await;
1849 }
1850 });
1851
1852 (address, request_count)
1853 }
1854
1855 async fn create_hanging_quote_server() -> std::net::SocketAddr {
1856 let listener = TcpListener::bind("127.0.0.1:0")
1857 .await
1858 .unwrap();
1859 let address = listener.local_addr().unwrap();
1860
1861 tokio::spawn(async move {
1862 let (_stream, _) = listener.accept().await.unwrap();
1863 std::future::pending::<()>().await;
1864 });
1865
1866 address
1867 }
1868
1869 #[tokio::test]
1870 async fn handles_orderbook_api_error_with_http_200() {
1871 let (address, request_count) = create_quote_server(
1872 "200 OK",
1873 r#"{"code":171015,"message":"quoted token not available"}"#,
1874 )
1875 .await;
1876 let client = create_test_client(format!("http://{address}"));
1877
1878 let result = client.fetch_orderbook().await;
1879
1880 assert!(matches!(
1881 result,
1882 Err(RFQError::QuoteNotFound(message)) if message.contains("171015")
1883 ));
1884 assert_eq!(request_count.load(Ordering::SeqCst), 1);
1885 }
1886
1887 #[tokio::test]
1888 async fn handles_documented_quote_error_without_retrying() {
1889 let (address, request_count) = create_quote_server(
1890 "200 OK",
1891 r#"{"code":171015,"message":"quoted token not available"}"#,
1892 )
1893 .await;
1894 let client = create_test_client(format!("http://{address}"));
1895
1896 let result = client
1897 .request_binding_quote(&create_test_quote_params())
1898 .await;
1899
1900 match result {
1901 Err(RFQError::QuoteNotFound(message)) => {
1902 assert!(message.contains("171015"));
1903 assert!(message.contains("quoted token not available"));
1904 }
1905 other => panic!("Expected Native API error, got {other:?}"),
1906 }
1907 assert_eq!(request_count.load(Ordering::SeqCst), 1);
1908 }
1909
1910 #[tokio::test]
1911 async fn handles_success_false_as_quote_not_found_without_retrying() {
1912 let mut response = successful_quote_json("1");
1913 response["success"] = serde_json::Value::Bool(false);
1914 response["errorMessage"] = serde_json::Value::String("quote unavailable".to_string());
1915 let (address, request_count) = create_quote_server("200 OK", response.to_string()).await;
1916 let client = create_test_client(format!("http://{address}"));
1917
1918 let result = client
1919 .request_binding_quote(&create_test_quote_params())
1920 .await;
1921
1922 assert!(matches!(
1923 result,
1924 Err(RFQError::QuoteNotFound(message)) if message.contains("quote unavailable")
1925 ));
1926 assert_eq!(request_count.load(Ordering::SeqCst), 1);
1927 }
1928
1929 #[tokio::test]
1930 async fn retries_documented_temporary_api_error() {
1931 let (address, request_count) = create_quote_server(
1932 "200 OK",
1933 r#"{"code":301016,"message":"quote invalid, risk management checks failed"}"#,
1934 )
1935 .await;
1936 let client = create_test_client(format!("http://{address}"));
1937
1938 let result = client
1939 .request_binding_quote(&create_test_quote_params())
1940 .await;
1941
1942 assert!(matches!(result, Err(RFQError::QuoteNotFound(_))));
1943 assert_eq!(request_count.load(Ordering::SeqCst), 3);
1944 }
1945
1946 #[tokio::test]
1947 async fn retries_server_error_without_native_error_envelope() {
1948 let (address, request_count) =
1949 create_quote_server("503 Service Unavailable", "<html>upstream unavailable</html>")
1950 .await;
1951 let client = create_test_client(format!("http://{address}"));
1952
1953 let result = client
1954 .request_binding_quote(&create_test_quote_params())
1955 .await;
1956
1957 assert!(matches!(
1958 result,
1959 Err(RFQError::ConnectionError(message)) if message.contains("503 Service Unavailable")
1960 ));
1961 assert_eq!(request_count.load(Ordering::SeqCst), 3);
1962 }
1963
1964 #[tokio::test]
1965 async fn retries_malformed_success_response() {
1966 let (address, request_count) =
1967 create_quote_server("200 OK", r#"{"unexpected":true}"#).await;
1968 let client = create_test_client(format!("http://{address}"));
1969
1970 let result = client
1971 .request_binding_quote(&create_test_quote_params())
1972 .await;
1973
1974 assert!(matches!(result, Err(RFQError::ParsingError(_))));
1975 assert_eq!(request_count.load(Ordering::SeqCst), 3);
1976 }
1977
1978 #[tokio::test]
1979 async fn times_out_when_quote_response_stalls() {
1980 let address = create_hanging_quote_server().await;
1981 let mut client = create_test_client(format!("http://{address}"));
1982 client.quote_timeout = Duration::from_millis(50);
1983
1984 let result = tokio::time::timeout(
1985 Duration::from_secs(1),
1986 client.request_binding_quote(&create_test_quote_params()),
1987 )
1988 .await
1989 .expect("Native quote timeout did not terminate the request");
1990
1991 assert!(matches!(
1992 result,
1993 Err(RFQError::ConnectionError(message)) if message.contains("timed out after 50ms")
1994 ));
1995 }
1996
1997 #[tokio::test]
1998 async fn shares_quote_timeout_across_retries() {
1999 let listener = TcpListener::bind("127.0.0.1:0")
2000 .await
2001 .unwrap();
2002 let address = listener.local_addr().unwrap();
2003 let request_count = Arc::new(AtomicUsize::new(0));
2004 let server_request_count = request_count.clone();
2005 let params = create_test_quote_params();
2006 let success_body = successful_quote_json(¶ms.amount_in.to_string()).to_string();
2007 let quote_timeout = NATIVE_API_RETRY_DELAY * 2;
2008
2009 tokio::spawn(async move {
2010 let retry_body =
2011 r#"{"code":301016,"message":"quote invalid, risk management checks failed"}"#
2012 .to_string();
2013 for body in [retry_body, success_body] {
2014 let (mut stream, _) = listener.accept().await.unwrap();
2015 let attempt = server_request_count.fetch_add(1, Ordering::SeqCst);
2016 if attempt == 1 {
2017 tokio::time::sleep(quote_timeout * 3 / 4).await;
2019 }
2020 let response = format!(
2021 "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
2022 body.len()
2023 );
2024 let _ = stream
2025 .write_all(response.as_bytes())
2026 .await;
2027 let _ = stream.shutdown().await;
2028 }
2029 });
2030
2031 let mut client = create_test_client(format!("http://{address}"));
2032 client.quote_timeout = quote_timeout;
2033
2034 let result = timeout(Duration::from_secs(5), client.request_binding_quote(¶ms))
2035 .await
2036 .expect("Native quote timeout did not terminate the retries");
2037
2038 assert!(matches!(
2039 result,
2040 Err(RFQError::QuoteNotFound(message)) if message.contains("301016")
2041 ));
2042 assert_eq!(request_count.load(Ordering::SeqCst), 2);
2043 }
2044
2045 #[tokio::test]
2046 async fn does_not_retry_documented_authentication_error() {
2047 let (address, request_count) = create_quote_server(
2048 "200 OK",
2049 r#"{"code":201001,"message":"auth get api key is invalid"}"#,
2050 )
2051 .await;
2052 let client = create_test_client(format!("http://{address}"));
2053
2054 let result = client
2055 .request_binding_quote(&create_test_quote_params())
2056 .await;
2057
2058 assert!(matches!(result, Err(RFQError::FatalError(_))));
2059 assert_eq!(request_count.load(Ordering::SeqCst), 1);
2060 }
2061}