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