Skip to main content

tycho_simulation/rfq/protocols/native/
client.rs

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);
43// V6 tradeRFQT: a dynamic quote tuple and two uint256 overrides, i.e. three ABI head words.
44const 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    // Maximum atomic minimum_in_base among the Native entries contributing these levels.
66    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    // Native API error codes:
94    // <https://docs.native.org/native-dev/build-with-native/swap-aggregators/firmquote-swap-apis/miscellaneous/error-handling#error-codes>
95    fn classify_api_error(error: &NativeApiErrorResponse) -> QuoteAttemptError {
96        let message = format!("Native API error {}: {}", error.code, error.message);
97        match error.code {
98            // Native documents these as temporary risk/rate-limit failures.
99            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            // The requested quote is unavailable for the current orderbook/liquidity.
108            // 171055 is not in the public table, but Native returns it when the seller amount is
109            // below the current book minimum.
110            101010 | 171037 | 171011 | 171015 | 171055 | 101007 => {
111                QuoteAttemptError::Fatal(RFQError::QuoteNotFound(message))
112            }
113            // The request itself must be corrected before another attempt can succeed.
114            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                // `group_orderbook` keeps the configured quote token on the quote side, so every
161                // matching conversion book is valued in comparable approved-token units.
162                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                    // Prefer the smaller token address when liquidity is equal so the result does
175                    // not depend on HashMap or HashSet iteration order.
176                    .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            // `showNative` is not boolean: its value selects the address used for native-token
224            // books. Request address(0) so the response matches Tycho's internal representation.
225            .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        // Native can return an API error envelope with HTTP 200.
239        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            // Keep the book direction stable even when Native publishes only one direction of a
278            // pair.
279            // Prefer exactly one configured quote token; if both or neither are configured, use
280            // the sorted pair order established by the grouping key above.
281            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            // Native's minimum_in_base is always denominated in entry.base_address. For a bid the
294            // taker sells base, so it is an input minimum; for an ask the taker receives base, so
295            // it is an output minimum. Mirroring swaps bid/ask and remaps that minimum into the
296            // canonical direction.
297            for entry in entries {
298                // A zero-only direct entry must not suppress usable mirrored liquidity for the
299                // same side. NativeState also filters zero quantities as a defensive measure for
300                // deserialized states that do not pass through this grouping path.
301                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            // Use mirrored levels only when direct ones are absent to avoid double-counting. Keep
335            // the minimum from the selected representation so discarded levels cannot constrain
336            // the surviving side.
337            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            // Use the sorted pair key so the component ID remains stable if Native returns the
348            // opposite book direction in a later poll.
349            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        // 1. Check API-level success
374        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        // Ensure we actually got an order
389        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        // Prevents silently accepting a mismatched/malicious quote.
400        let seller_token = bytes_to_address(&params.token_in)?;
401        let buyer_token = bytes_to_address(&params.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(&params.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        // Security: reject already-expired quotes before we bother building a SignedQuote
432        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        // Bind Native's top-level amountIn and signed sellerTokenAmount to the requested quote
445        // baseline. The encoder stores this baseline as signedAmountIn; the executor supplies
446        // actualSellerAmount when execution receives a different amount from the preceding hop.
447        let quoted_amount_in = BigUint::from_str(&quote_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        // effectiveSellerTokenAmount may differ from the requested gross input for
474        // fee-on-transfer tokens, so it is not an equality invariant here. amountIn and
475        // sellerTokenAmount still bind the quote to Tycho's requested input.
476        // Native requires txRequest.value for native-token quotes. Validate the quoted value here,
477        // while it still describes the original signed amount. During execution, the preceding hop
478        // may deliver either less or more; the executor handles that through actualSellerAmount.
479        let quoted_value = BigUint::from_str(&quote_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        // Decode calldata (pre-built by Native Relay, ready to submit as-is)
504        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        // These offsets are fixed by V6's tradeRFQT ABI (a dynamic quote tuple and
529        // two uint256 overrides). Rejecting any other value catches an incompatible or malformed
530        // API response before it reaches the encoder; the executor independently
531        // hardcodes the same positions rather than trusting route data.
532        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(&quote_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(&quote_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        // Native returns documented API error codes in the body, including with HTTP 200.
625        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                // Native Relay publishes a complete RFQ orderbook once per request. Polling the
674                // full book keeps component creation/removal deterministic and avoids per-pair REST
675                // fan-out.
676                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                    // Keep unrequested books available for TVL conversion, but only emit requested
688                    // markets as components.
689                    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                        // TVL thresholds are applied in approved quote-token units. If Native
699                        // quotes this market against another token, normalize through the most
700                        // liquid available approved quote-token market before filtering.
701                        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                // Emit removals for markets that disappeared from the Relay orderbook or no longer
731                // pass token/TVL filtering.
732                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                // Native is off-chain and timestamped, not block-based. Downstream decoders use
746                // this wall-clock header to build a normal Tycho state update.
747                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(&params.receiver)?;
769        let token_in = bytes_to_address(&params.token_in)?;
770        let token_out = bytes_to_address(&params.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        // Bind the timeout result before inspecting last_error so the attempts future—and its
816        // mutable borrow—has been dropped.
817        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    #[test]
956    fn rejects_unsupported_chain_at_construction() {
957        let result = NativeClient::new(
958            Chain::Polygon,
959            "test-api-key".to_string(),
960            HashSet::new(),
961            0.0,
962            HashSet::new(),
963            Duration::from_secs(1),
964            Duration::from_secs(5),
965        );
966
967        assert!(matches!(result, Err(RFQError::InvalidInput(_))));
968    }
969
970    #[test]
971    fn builder_rejects_zero_poll_time() {
972        let result = NativeClientBuilder::new(Chain::Ethereum, "test-api-key".to_string())
973            .poll_time(Duration::ZERO)
974            .build();
975
976        assert!(matches!(
977            result,
978            Err(RFQError::InvalidInput(message))
979                if message == "Native polling interval must be greater than zero"
980        ));
981    }
982
983    #[test]
984    fn selects_most_liquid_tvl_conversion_book() {
985        let weth = Bytes::from_str("0x3333333333333333333333333333333333333333").unwrap();
986        let usdc = Bytes::from_str("0x1111111111111111111111111111111111111111").unwrap();
987        let usdt = Bytes::from_str("0x2222222222222222222222222222222222222222").unwrap();
988        let wbtc = Bytes::from_str("0x4444444444444444444444444444444444444444").unwrap();
989        let unapproved = Bytes::from_str("0x5555555555555555555555555555555555555555").unwrap();
990        let client = NativeClient::new(
991            Chain::Ethereum,
992            "test-api-key".to_string(),
993            HashSet::from([
994                weth.clone(),
995                usdc.clone(),
996                usdt.clone(),
997                wbtc.clone(),
998                unapproved.clone(),
999            ]),
1000            0.0,
1001            HashSet::from([usdc.clone(), usdt.clone()]),
1002            Duration::from_secs(1),
1003            Duration::from_secs(5),
1004        )
1005        .unwrap();
1006        let books = HashMap::from([
1007            ("usdc".to_string(), conversion_book(weth.clone(), usdc.clone(), 1.0, 100.0)),
1008            ("usdt".to_string(), conversion_book(weth.clone(), usdt.clone(), 2.0, 100.0)),
1009            ("unrelated".to_string(), conversion_book(wbtc, usdc, 1_000.0, 100.0)),
1010            ("unapproved".to_string(), conversion_book(weth.clone(), unapproved, 2_000.0, 100.0)),
1011        ]);
1012
1013        let selected = client
1014            .select_tvl_conversion_book(&weth, &books)
1015            .expect("one conversion book");
1016
1017        assert_eq!(selected.quote_address, usdt);
1018    }
1019
1020    #[test]
1021    fn selects_lower_quote_address_for_equal_tvl_conversion_liquidity() {
1022        let weth = Bytes::from_str("0x3333333333333333333333333333333333333333").unwrap();
1023        let lower_quote = Bytes::from_str("0x1111111111111111111111111111111111111111").unwrap();
1024        let higher_quote = Bytes::from_str("0x2222222222222222222222222222222222222222").unwrap();
1025        let client = NativeClient::new(
1026            Chain::Ethereum,
1027            "test-api-key".to_string(),
1028            HashSet::from([weth.clone(), lower_quote.clone(), higher_quote.clone()]),
1029            0.0,
1030            HashSet::from([lower_quote.clone(), higher_quote.clone()]),
1031            Duration::from_secs(1),
1032            Duration::from_secs(5),
1033        )
1034        .unwrap();
1035        let books = HashMap::from([
1036            ("higher".to_string(), conversion_book(weth.clone(), higher_quote, 2.0, 100.0)),
1037            ("lower".to_string(), conversion_book(weth.clone(), lower_quote.clone(), 1.0, 200.0)),
1038        ]);
1039
1040        let selected = client
1041            .select_tvl_conversion_book(&weth, &books)
1042            .expect("one conversion book");
1043
1044        assert_eq!(selected.quote_address, lower_quote);
1045    }
1046
1047    #[test]
1048    fn creates_indexer_compatible_component_from_relay_orderbook() {
1049        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1050        let usdt = Bytes::from_str("0xdac17f958d2ee523a2206206994597c13d831ec7").unwrap();
1051        let client = NativeClient::new(
1052            Chain::Ethereum,
1053            "test-api-key".to_string(),
1054            HashSet::from([weth.clone(), usdt.clone()]),
1055            0.0,
1056            HashSet::from([usdt.clone()]),
1057            Duration::from_secs(1),
1058            Duration::from_secs(5),
1059        )
1060        .unwrap();
1061
1062        let books = client.group_orderbook(vec![
1063            NativeOrderbookEntry {
1064                base_address: weth.clone(),
1065                quote_address: usdt.clone(),
1066                minimum_in_base: 0.0,
1067                side: NativeOrderbookSide::Bid,
1068                levels: vec![NativePriceLevel { quantity: 0.0001, price: 3213.12345 }],
1069            },
1070            NativeOrderbookEntry {
1071                base_address: weth.clone(),
1072                quote_address: usdt.clone(),
1073                minimum_in_base: 0.0,
1074                side: NativeOrderbookSide::Ask,
1075                levels: vec![NativePriceLevel { quantity: 2.0, price: 3214.0 }],
1076            },
1077            NativeOrderbookEntry {
1078                base_address: usdt.clone(),
1079                quote_address: weth.clone(),
1080                minimum_in_base: 100.0,
1081                side: NativeOrderbookSide::Bid,
1082                levels: vec![NativePriceLevel { quantity: 6428.0, price: 1.0 / 3214.0 }],
1083            },
1084            NativeOrderbookEntry {
1085                base_address: usdt.clone(),
1086                quote_address: weth.clone(),
1087                minimum_in_base: 100.0,
1088                side: NativeOrderbookSide::Ask,
1089                levels: vec![NativePriceLevel { quantity: 0.321312345, price: 1.0 / 3213.12345 }],
1090            },
1091        ]);
1092
1093        let (component_id, book) = books
1094            .into_iter()
1095            .next()
1096            .expect("one grouped book");
1097        let component = client.create_component_with_state(
1098            component_id.clone(),
1099            vec![book.base_address.clone(), book.quote_address.clone()],
1100            book.clone(),
1101            book.calculate_tvl(None)
1102                .expect("TVL should be finite"),
1103        );
1104
1105        assert_eq!(component.component.id, component_id);
1106        assert_eq!(component.component.protocol_system, NativeClient::PROTOCOL_SYSTEM);
1107        assert_eq!(component.component.protocol_type_name, "native_relay_pool");
1108        assert_eq!(component.component.tokens, vec![weth, usdt]);
1109        assert_eq!(component.state.component_id, component_id);
1110
1111        let encoded_book = component
1112            .state
1113            .attributes
1114            .get("book")
1115            .expect("book attribute");
1116        let decoded_book: NativePriceData = serde_json::from_slice(encoded_book).unwrap();
1117        assert_eq!(decoded_book.bids.len(), 1);
1118        assert_eq!(decoded_book.asks.len(), 1);
1119        assert_eq!(decoded_book.bids[0].quantity, 0.0001);
1120        assert_eq!(decoded_book.bids[0].price, 3213.12345);
1121    }
1122
1123    #[test]
1124    fn uses_stable_component_id_and_direction_when_merging_mirrored_books() {
1125        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1126        let usdc = Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap();
1127        let client = NativeClient::new(
1128            Chain::Ethereum,
1129            "test-api-key".to_string(),
1130            HashSet::from([weth.clone(), usdc.clone()]),
1131            0.0,
1132            HashSet::from([usdc.clone()]),
1133            Duration::from_secs(1),
1134            Duration::from_secs(5),
1135        )
1136        .unwrap();
1137        let entries = vec![
1138            NativeOrderbookEntry {
1139                base_address: weth.clone(),
1140                quote_address: usdc.clone(),
1141                minimum_in_base: 100_000_000_000.0,
1142                side: NativeOrderbookSide::Bid,
1143                levels: vec![NativePriceLevel { quantity: 1.0, price: 2_000.0 }],
1144            },
1145            NativeOrderbookEntry {
1146                base_address: weth.clone(),
1147                quote_address: usdc.clone(),
1148                minimum_in_base: 300_000_000_000.0,
1149                side: NativeOrderbookSide::Ask,
1150                levels: vec![NativePriceLevel { quantity: 1.0, price: 2_100.0 }],
1151            },
1152            NativeOrderbookEntry {
1153                base_address: usdc.clone(),
1154                quote_address: weth.clone(),
1155                minimum_in_base: 100.0,
1156                side: NativeOrderbookSide::Bid,
1157                levels: vec![NativePriceLevel { quantity: 2_000.0, price: 0.0005 }],
1158            },
1159            NativeOrderbookEntry {
1160                base_address: usdc.clone(),
1161                quote_address: weth.clone(),
1162                minimum_in_base: 250.0,
1163                side: NativeOrderbookSide::Bid,
1164                levels: vec![NativePriceLevel { quantity: 2_000.0, price: 0.0005 }],
1165            },
1166            NativeOrderbookEntry {
1167                base_address: usdc.clone(),
1168                quote_address: weth.clone(),
1169                minimum_in_base: 400.0,
1170                side: NativeOrderbookSide::Ask,
1171                levels: vec![NativePriceLevel { quantity: 2.0, price: 0.5 }],
1172            },
1173        ];
1174
1175        let forward_only = client.group_orderbook(entries[..2].to_vec());
1176        let reverse_entries = entries[2..].to_vec();
1177        let reverse_only = client.group_orderbook(reverse_entries.clone());
1178        let reversed_reverse_only = client.group_orderbook(
1179            reverse_entries
1180                .into_iter()
1181                .rev()
1182                .collect(),
1183        );
1184        assert_eq!(reverse_only, reversed_reverse_only);
1185        let pair = format!("native_{}/{}", hex::encode(&usdc), hex::encode(&weth));
1186        let component_id = keccak256(pair.as_bytes()).to_string();
1187
1188        let forward_book = forward_only
1189            .get(&component_id)
1190            .expect("forward-only book uses the stable component ID");
1191        assert_eq!(forward_book.base_address, weth);
1192        assert_eq!(forward_book.quote_address, usdc);
1193        assert_eq!(forward_book.minimum_in_base, 100_000_000_000.0);
1194        assert_eq!(forward_book.minimum_in_quote, 0.0);
1195        assert_eq!(forward_book.minimum_out_base, 300_000_000_000.0);
1196        assert_eq!(forward_book.minimum_out_quote, 0.0);
1197
1198        let reverse_book = reverse_only
1199            .get(&component_id)
1200            .expect("reverse-only book uses the stable component ID");
1201        assert_eq!(reverse_book.base_address, weth);
1202        assert_eq!(reverse_book.quote_address, usdc);
1203        assert_eq!(reverse_book.minimum_in_base, 0.0);
1204        assert_eq!(reverse_book.minimum_in_quote, 250.0);
1205        assert_eq!(reverse_book.minimum_out_base, 0.0);
1206        assert_eq!(reverse_book.minimum_out_quote, 400.0);
1207        assert_eq!(reverse_book.bids, vec![NativePriceLevel { quantity: 1.0, price: 2.0 }]);
1208        assert_eq!(
1209            reverse_book.asks,
1210            vec![
1211                NativePriceLevel { quantity: 1.0, price: 2_000.0 },
1212                NativePriceLevel { quantity: 1.0, price: 2_000.0 },
1213            ]
1214        );
1215
1216        let mixed = client.group_orderbook(vec![entries[0].clone(), entries[2].clone()]);
1217        let mixed_book = mixed.get(&component_id).unwrap();
1218        assert_eq!(mixed_book.minimum_in_base, 100_000_000_000.0);
1219        assert_eq!(mixed_book.minimum_in_quote, 100.0);
1220        assert_eq!(mixed_book.minimum_out_base, 0.0);
1221        assert_eq!(mixed_book.minimum_out_quote, 0.0);
1222        assert_eq!(mixed_book.bids, vec![NativePriceLevel { quantity: 1.0, price: 2_000.0 }]);
1223        assert_eq!(mixed_book.asks, vec![NativePriceLevel { quantity: 1.0, price: 2_000.0 }]);
1224
1225        let books = client.group_orderbook(entries.clone());
1226        let reversed_books = client.group_orderbook(entries.into_iter().rev().collect());
1227
1228        assert_eq!(books, reversed_books);
1229        assert_eq!(books.len(), 1);
1230        let book = books.get(&component_id).unwrap();
1231        assert_eq!(book.base_address, weth);
1232        assert_eq!(book.quote_address, usdc);
1233        assert_eq!(book.minimum_in_base, 100_000_000_000.0);
1234        assert_eq!(book.minimum_in_quote, 0.0);
1235        assert_eq!(book.minimum_out_base, 300_000_000_000.0);
1236        assert_eq!(book.minimum_out_quote, 0.0);
1237        assert_eq!(book.bids, vec![NativePriceLevel { quantity: 1.0, price: 2_000.0 }]);
1238        assert_eq!(book.asks, vec![NativePriceLevel { quantity: 1.0, price: 2_100.0 }]);
1239    }
1240
1241    #[test]
1242    fn zero_only_direct_side_does_not_suppress_mirrored_liquidity() {
1243        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1244        let usdc = Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap();
1245        let client = NativeClient::new(
1246            Chain::Ethereum,
1247            "test-api-key".to_string(),
1248            HashSet::from([weth.clone(), usdc.clone()]),
1249            0.0,
1250            HashSet::from([usdc.clone()]),
1251            Duration::from_secs(1),
1252            Duration::from_secs(5),
1253        )
1254        .unwrap();
1255
1256        let books = client.group_orderbook(vec![
1257            NativeOrderbookEntry {
1258                base_address: weth.clone(),
1259                quote_address: usdc.clone(),
1260                minimum_in_base: 999.0,
1261                side: NativeOrderbookSide::Bid,
1262                levels: vec![NativePriceLevel { quantity: 0.0, price: 2_000.0 }],
1263            },
1264            NativeOrderbookEntry {
1265                base_address: usdc,
1266                quote_address: weth,
1267                minimum_in_base: 250.0,
1268                side: NativeOrderbookSide::Ask,
1269                levels: vec![NativePriceLevel { quantity: 2.0, price: 0.5 }],
1270            },
1271        ]);
1272
1273        let book = books
1274            .values()
1275            .next()
1276            .expect("one grouped book");
1277        assert_eq!(book.bids, vec![NativePriceLevel { quantity: 1.0, price: 2.0 }]);
1278        assert_eq!(book.minimum_in_base, 0.0);
1279        assert_eq!(book.minimum_out_quote, 250.0);
1280    }
1281
1282    #[rstest]
1283    #[case::token0_only(true, false, true)]
1284    #[case::token1_only(false, true, false)]
1285    #[case::both_tokens(true, true, false)]
1286    #[case::neither_token(false, false, false)]
1287    fn selects_stable_direction_for_quote_token_preferences(
1288        #[case] token0_is_quote: bool,
1289        #[case] token1_is_quote: bool,
1290        #[case] expected_quote_is_token0: bool,
1291    ) {
1292        let token0 = Bytes::from_str("0x1111111111111111111111111111111111111111").unwrap();
1293        let token1 = Bytes::from_str("0x2222222222222222222222222222222222222222").unwrap();
1294        assert!(token0.as_ref() < token1.as_ref());
1295
1296        let mut quote_tokens = HashSet::new();
1297        if token0_is_quote {
1298            quote_tokens.insert(token0.clone());
1299        }
1300        if token1_is_quote {
1301            quote_tokens.insert(token1.clone());
1302        }
1303        let client = NativeClient::new(
1304            Chain::Ethereum,
1305            "test-api-key".to_string(),
1306            HashSet::from([token0.clone(), token1.clone()]),
1307            0.0,
1308            quote_tokens,
1309            Duration::from_secs(1),
1310            Duration::from_secs(5),
1311        )
1312        .unwrap();
1313
1314        // Native returned only the direction opposite to the sorted fallback.
1315        let books = client.group_orderbook(vec![NativeOrderbookEntry {
1316            base_address: token1.clone(),
1317            quote_address: token0.clone(),
1318            minimum_in_base: 0.0,
1319            side: NativeOrderbookSide::Bid,
1320            levels: vec![NativePriceLevel { quantity: 1.0, price: 1.0 }],
1321        }]);
1322
1323        let book = books
1324            .values()
1325            .next()
1326            .expect("one grouped book");
1327        let (expected_base, expected_quote) =
1328            if expected_quote_is_token0 { (token1, token0) } else { (token0, token1) };
1329        assert_eq!(book.base_address, expected_base);
1330        assert_eq!(book.quote_address, expected_quote);
1331    }
1332
1333    fn create_test_quote_params() -> GetAmountOutParams {
1334        GetAmountOutParams {
1335            amount_in: BigUint::from(1u64),
1336            token_in: Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap(),
1337            token_out: Bytes::from_str("0xdac17f958d2ee523a2206206994597c13d831ec7").unwrap(),
1338            sender: Bytes::from_str("0x3333333333333333333333333333333333333333").unwrap(),
1339            receiver: Bytes::from_str("0x4444444444444444444444444444444444444444").unwrap(),
1340        }
1341    }
1342
1343    #[test]
1344    fn accepts_quote_with_requested_input_amount() {
1345        let params = create_test_quote_params();
1346        let response = successful_quote_response(&params.amount_in.to_string());
1347
1348        let quote = NativeClient::process_quote_response(response, &params).unwrap();
1349
1350        assert_eq!(quote.amount_in, params.amount_in);
1351        assert_eq!(
1352            quote
1353                .quote_attributes
1354                .get("deadline_timestamp")
1355                .unwrap()
1356                .as_ref(),
1357            u64::MAX.to_be_bytes()
1358        );
1359    }
1360
1361    #[test]
1362    fn rejects_quote_with_mismatched_input_amount() {
1363        let params = create_test_quote_params();
1364        let mut response = successful_quote_response(&params.amount_in.to_string());
1365        // Keep sellerTokenAmount correct so only the top-level amountIn check can reject this.
1366        response.amount_in = "2".to_string();
1367
1368        let result = NativeClient::process_quote_response(response, &params);
1369
1370        assert!(matches!(
1371            result,
1372            Err(RFQError::ParsingError(message)) if message.contains("Native Relay quote input amount mismatch")
1373        ));
1374    }
1375
1376    #[test]
1377    fn rejects_quote_with_mismatched_signed_input_amount() {
1378        let params = create_test_quote_params();
1379        let mut response = successful_quote_response(&params.amount_in.to_string());
1380        response.orders[0].seller_token_amount = "2".to_string();
1381
1382        let result = NativeClient::process_quote_response(response, &params);
1383
1384        assert!(matches!(
1385            result,
1386            Err(RFQError::ParsingError(message)) if message.contains("signed input amount mismatch")
1387        ));
1388    }
1389
1390    #[test]
1391    fn accepts_quote_with_different_effective_input_amount() {
1392        let mut params = create_test_quote_params();
1393        params.amount_in = BigUint::from(100u64);
1394        let mut response = successful_quote_response(&params.amount_in.to_string());
1395        response.orders[0].effective_seller_token_amount = "99".to_string();
1396
1397        let quote = NativeClient::process_quote_response(response, &params).unwrap();
1398
1399        assert_eq!(quote.amount_in, params.amount_in);
1400    }
1401
1402    #[rstest]
1403    #[case::seller_token(true)]
1404    #[case::buyer_token(false)]
1405    fn rejects_quote_with_mismatched_token(#[case] mutate_seller_token: bool) {
1406        let params = create_test_quote_params();
1407        let mut response = successful_quote_response(&params.amount_in.to_string());
1408        let mismatched_token = "0x5555555555555555555555555555555555555555".to_string();
1409        if mutate_seller_token {
1410            response.orders[0].seller_token = mismatched_token;
1411        } else {
1412            response.orders[0].buyer_token = mismatched_token;
1413        }
1414
1415        let result = NativeClient::process_quote_response(response, &params);
1416
1417        assert!(matches!(
1418            result,
1419            Err(RFQError::ParsingError(message)) if message.contains("quote token mismatch")
1420        ));
1421    }
1422
1423    #[test]
1424    fn rejects_quote_with_mismatched_recipient() {
1425        let params = create_test_quote_params();
1426        let mut response = successful_quote_response(&params.amount_in.to_string());
1427        response.orders[0].recipient = "0x5555555555555555555555555555555555555555".to_string();
1428
1429        let result = NativeClient::process_quote_response(response, &params);
1430
1431        assert!(matches!(
1432            result,
1433            Err(RFQError::ParsingError(message)) if message.contains("recipient mismatch")
1434        ));
1435    }
1436
1437    #[test]
1438    fn rejects_quote_without_calldata() {
1439        let params = create_test_quote_params();
1440        let mut response = successful_quote_response(&params.amount_in.to_string());
1441        response.tx_request.calldata.clear();
1442
1443        let result = NativeClient::process_quote_response(response, &params);
1444
1445        assert!(matches!(
1446            result,
1447            Err(RFQError::QuoteNotFound(message)) if message.contains("did not include calldata")
1448        ));
1449    }
1450
1451    #[test]
1452    fn rejects_truncated_trade_calldata() {
1453        let params = create_test_quote_params();
1454        let mut response = successful_quote_response(&params.amount_in.to_string());
1455        response.tx_request.calldata = format!("0x7083527c{}", "00".repeat(95));
1456
1457        let result = NativeClient::process_quote_response(response, &params);
1458
1459        assert!(matches!(
1460            result,
1461            Err(RFQError::ParsingError(message)) if message.contains("calldata too short")
1462        ));
1463    }
1464
1465    #[test]
1466    fn rejects_quote_with_wrong_trade_rfqt_selector() {
1467        let params = create_test_quote_params();
1468        let mut response = successful_quote_response(&params.amount_in.to_string());
1469        let mut calldata = hex::decode(
1470            response
1471                .tx_request
1472                .calldata
1473                .trim_start_matches("0x"),
1474        )
1475        .unwrap();
1476        // V4 also exposes tradeRFQT, but its quote tuple has a different ABI.
1477        calldata[..4].copy_from_slice(&[0x09, 0x47, 0xc2, 0xd9]);
1478        response.tx_request.calldata = format!("0x{}", hex::encode(calldata));
1479
1480        let result = NativeClient::process_quote_response(response, &params);
1481
1482        assert!(matches!(
1483            result,
1484            Err(RFQError::ParsingError(message)) if message.contains("selector")
1485        ));
1486    }
1487
1488    #[rstest]
1489    #[case::v4("4")]
1490    #[case::unknown("7")]
1491    fn rejects_quote_with_wrong_router_version(#[case] version: &str) {
1492        let params = create_test_quote_params();
1493        let mut response = successful_quote_response(&params.amount_in.to_string());
1494        response.router_version = version.to_string();
1495
1496        assert!(matches!(
1497            NativeClient::process_quote_response(response, &params),
1498            Err(RFQError::ParsingError(message)) if message.contains("Unexpected Native router version")
1499        ));
1500    }
1501
1502    #[rstest]
1503    #[case::seller_offset(68, 68)]
1504    #[case::minimum_offset(36, 36)]
1505    fn rejects_quote_with_noncanonical_override_offsets(
1506        #[case] amount_in_offset: u32,
1507        #[case] amount_out_minimum_offset: u32,
1508    ) {
1509        let params = create_test_quote_params();
1510        let mut response = successful_quote_response(&params.amount_in.to_string());
1511        response.amount_in_offset = amount_in_offset;
1512        response.amount_out_minimum_offset = amount_out_minimum_offset;
1513
1514        let result = NativeClient::process_quote_response(response, &params);
1515
1516        assert!(matches!(
1517            result,
1518            Err(RFQError::ParsingError(message)) if message.contains(
1519                "Unexpected Native V6 override offsets"
1520            )
1521        ));
1522    }
1523
1524    #[rstest]
1525    #[case::seller(ACTUAL_SELLER_AMOUNT_OFFSET, "actualSellerAmount")]
1526    #[case::minimum(ACTUAL_MIN_OUTPUT_AMOUNT_OFFSET, "actualMinOutputAmount")]
1527    fn rejects_quote_with_preset_override(
1528        #[case] override_offset: usize,
1529        #[case] field_name: &str,
1530    ) {
1531        let params = create_test_quote_params();
1532        let mut response = successful_quote_response(&params.amount_in.to_string());
1533        let mut calldata = hex::decode(
1534            response
1535                .tx_request
1536                .calldata
1537                .trim_start_matches("0x"),
1538        )
1539        .unwrap();
1540        calldata[override_offset + 31] = 1;
1541        response.tx_request.calldata = format!("0x{}", hex::encode(calldata));
1542
1543        let result = NativeClient::process_quote_response(response, &params);
1544
1545        assert!(matches!(
1546            result,
1547            Err(RFQError::ParsingError(message)) if message.contains(field_name)
1548        ));
1549    }
1550
1551    #[test]
1552    fn accepts_native_eth_response_using_zero_address() {
1553        let mut params = create_test_quote_params();
1554        params.token_in = Bytes::zero(20);
1555        let mut response = successful_quote_response(&params.amount_in.to_string());
1556        response.orders[0].seller_token = "0x0000000000000000000000000000000000000000".to_string();
1557        response.tx_request.value = params.amount_in.to_string();
1558
1559        let quote = NativeClient::process_quote_response(response, &params).unwrap();
1560
1561        assert_eq!(quote.amount_in, params.amount_in);
1562    }
1563
1564    #[test]
1565    fn rejects_native_quote_with_mismatched_payable_value() {
1566        let mut params = create_test_quote_params();
1567        params.token_in = Bytes::zero(20);
1568        let mut response = successful_quote_response(&params.amount_in.to_string());
1569        response.orders[0].seller_token = "0x0000000000000000000000000000000000000000".to_string();
1570        response.tx_request.value = "2".to_string();
1571
1572        let result = NativeClient::process_quote_response(response, &params);
1573
1574        assert!(matches!(
1575            result,
1576            Err(RFQError::ParsingError(message)) if message.contains("payable value mismatch")
1577        ));
1578    }
1579
1580    #[test]
1581    fn rejects_malformed_payable_value() {
1582        let params = create_test_quote_params();
1583        let mut response = successful_quote_response(&params.amount_in.to_string());
1584        response.tx_request.value = "not-a-number".to_string();
1585
1586        let result = NativeClient::process_quote_response(response, &params);
1587
1588        assert!(matches!(
1589            result,
1590            Err(RFQError::ParsingError(message)) if message.contains("txRequest.value")
1591        ));
1592    }
1593
1594    #[test]
1595    fn rejects_erc20_quote_with_nonzero_payable_value() {
1596        let params = create_test_quote_params();
1597        let mut response = successful_quote_response(&params.amount_in.to_string());
1598        response.tx_request.value = "1".to_string();
1599
1600        let result = NativeClient::process_quote_response(response, &params);
1601
1602        assert!(matches!(
1603            result,
1604            Err(RFQError::ParsingError(message)) if message.contains("payable value mismatch")
1605        ));
1606    }
1607
1608    #[test]
1609    fn accepts_native_eth_orderbook_using_zero_address() {
1610        let tycho_native_eth = Bytes::zero(20);
1611        let usdc = Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap();
1612        let client = NativeClient::new(
1613            Chain::Ethereum,
1614            "test-api-key".to_string(),
1615            HashSet::from([tycho_native_eth.clone(), usdc.clone()]),
1616            0.0,
1617            HashSet::from([usdc.clone()]),
1618            Duration::from_secs(1),
1619            Duration::from_secs(5),
1620        )
1621        .unwrap();
1622
1623        let entry: NativeOrderbookEntry = serde_json::from_value(serde_json::json!({
1624            "base_address": tycho_native_eth.to_string(),
1625            "quote_address": usdc.to_string(),
1626            "minimum_in_base": 1.0,
1627            "side": "bid",
1628            "levels": [[1.0, 3_000.0]]
1629        }))
1630        .unwrap();
1631        let books = client.group_orderbook(vec![entry]);
1632
1633        let book = books.values().next().unwrap();
1634        assert_eq!(book.base_address, tycho_native_eth);
1635    }
1636
1637    fn create_test_client(endpoint: String) -> NativeClient {
1638        let mut client = NativeClient::new(
1639            Chain::Ethereum,
1640            "test-api-key".to_string(),
1641            HashSet::new(),
1642            0.0,
1643            HashSet::new(),
1644            Duration::from_secs(1),
1645            Duration::from_secs(5),
1646        )
1647        .unwrap();
1648        client.endpoint = endpoint;
1649        client
1650    }
1651
1652    #[tokio::test]
1653    async fn requests_v6_firm_quote() {
1654        let listener = TcpListener::bind("127.0.0.1:0")
1655            .await
1656            .unwrap();
1657        let address = listener.local_addr().unwrap();
1658        let server = tokio::spawn(async move {
1659            let (stream, _) = listener.accept().await.unwrap();
1660            let mut reader = BufReader::new(stream);
1661            let mut request_line = String::new();
1662            reader
1663                .read_line(&mut request_line)
1664                .await
1665                .unwrap();
1666            let body = successful_quote_json("1").to_string();
1667            let response = format!(
1668                "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
1669                body.len()
1670            );
1671            reader
1672                .into_inner()
1673                .write_all(response.as_bytes())
1674                .await
1675                .unwrap();
1676            request_line
1677        });
1678        let client = create_test_client(format!("http://{address}"));
1679        let params = create_test_quote_params();
1680
1681        let quote = client
1682            .request_binding_quote(&params)
1683            .await
1684            .unwrap();
1685        let request = server.await.unwrap();
1686        let url = reqwest::Url::parse(&format!(
1687            "http://{address}{}",
1688            request
1689                .split_whitespace()
1690                .nth(1)
1691                .unwrap()
1692        ))
1693        .unwrap();
1694        let query: HashMap<_, _> = url.query_pairs().into_owned().collect();
1695
1696        assert_eq!(url.path(), "/firm-quote");
1697        assert_eq!(query.get("version").map(String::as_str), Some("6"));
1698        assert_eq!(
1699            query
1700                .get("allow_multihop")
1701                .map(String::as_str),
1702            Some("false")
1703        );
1704        assert_eq!(quote.amount_in, params.amount_in);
1705    }
1706
1707    #[tokio::test]
1708    async fn requests_native_token_orderbooks() {
1709        let listener = TcpListener::bind("127.0.0.1:0")
1710            .await
1711            .unwrap();
1712        let address = listener.local_addr().unwrap();
1713        let server = tokio::spawn(async move {
1714            let (stream, _) = listener.accept().await.unwrap();
1715            let mut reader = BufReader::new(stream);
1716            let mut request_line = String::new();
1717            reader
1718                .read_line(&mut request_line)
1719                .await
1720                .unwrap();
1721            let mut stream = reader.into_inner();
1722            let response = "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: 2\r\nConnection: close\r\n\r\n[]";
1723            stream
1724                .write_all(response.as_bytes())
1725                .await
1726                .unwrap();
1727            request_line
1728        });
1729        let client = create_test_client(format!("http://{address}"));
1730
1731        let orderbook = client.fetch_orderbook().await.unwrap();
1732        let request = server.await.unwrap();
1733
1734        assert!(orderbook.is_empty());
1735        assert!(request.starts_with("GET /orderbook?"));
1736        assert!(request.contains("chain=ethereum"));
1737        assert!(request.contains("showNative=0x0"));
1738    }
1739
1740    #[rstest]
1741    #[case::with_conversion_helper(true, 300.0, Some(400.0))]
1742    #[case::without_conversion_helper(false, 300.0, None)]
1743    #[case::below_normalized_tvl_threshold(true, 401.0, None)]
1744    #[tokio::test]
1745    async fn stream_uses_unrequested_books_only_for_tvl_conversion(
1746        #[case] include_helper: bool,
1747        #[case] tvl_threshold: f64,
1748        #[case] expected_tvl: Option<f64>,
1749    ) {
1750        let weth = Bytes::from_str("0xc02aaa39b223fe8d0a0e5c4f27ead9083c756cc2").unwrap();
1751        let usdt = Bytes::from_str("0xdac17f958d2ee523a2206206994597c13d831ec7").unwrap();
1752        let usdc = Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap();
1753        let mut entries = vec![serde_json::json!({
1754            "base_address": weth.to_string(),
1755            "quote_address": usdt.to_string(),
1756            "minimum_in_base": 0.0,
1757            "side": "bid",
1758            "levels": [[2.0, 100.0]]
1759        })];
1760        if include_helper {
1761            // Use a reversed helper with a non-unit price so the test distinguishes the market's
1762            // 200 USDT of liquidity from its normalized value of 400 USDC.
1763            entries.push(serde_json::json!({
1764                "base_address": usdc.to_string(),
1765                "quote_address": usdt.to_string(),
1766                "minimum_in_base": 0.0,
1767                "side": "bid",
1768                "levels": [[1_000.0, 0.5]]
1769            }));
1770        }
1771        let (address, _) =
1772            create_quote_server("200 OK", serde_json::to_string(&entries).unwrap()).await;
1773        let mut client = create_test_client(format!("http://{address}"));
1774        client.tokens = HashSet::from([weth.clone(), usdt.clone()]);
1775        client.quote_tokens = HashSet::from([usdc]);
1776        client.tvl = tvl_threshold;
1777
1778        let (_, update) = timeout(Duration::from_secs(5), client.stream().next())
1779            .await
1780            .expect("orderbook poll timed out")
1781            .expect("stream ended")
1782            .expect("orderbook poll failed");
1783
1784        if let Some(tvl) = expected_tvl {
1785            assert_eq!(update.snapshots.states.len(), 1, "helper must not be emitted");
1786            let component = update
1787                .snapshots
1788                .states
1789                .values()
1790                .next()
1791                .unwrap();
1792            assert_eq!(component.component.tokens, vec![weth, usdt]);
1793            assert_eq!(component.component_tvl, Some(tvl));
1794        } else {
1795            assert!(update.snapshots.states.is_empty());
1796        }
1797    }
1798
1799    async fn create_quote_server(
1800        final_status: impl Into<String>,
1801        final_body: impl Into<String>,
1802    ) -> (std::net::SocketAddr, Arc<AtomicUsize>) {
1803        let final_status = final_status.into();
1804        let final_body = final_body.into();
1805        let listener = TcpListener::bind("127.0.0.1:0")
1806            .await
1807            .unwrap();
1808        let address = listener.local_addr().unwrap();
1809        let request_count = Arc::new(AtomicUsize::new(0));
1810        let server_request_count = request_count.clone();
1811
1812        tokio::spawn(async move {
1813            while let Ok((mut stream, _)) = listener.accept().await {
1814                server_request_count.fetch_add(1, Ordering::SeqCst);
1815                let response = format!(
1816                    "HTTP/1.1 {final_status}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{final_body}",
1817                    final_body.len()
1818                );
1819                let _ = stream
1820                    .write_all(response.as_bytes())
1821                    .await;
1822                let _ = stream.shutdown().await;
1823            }
1824        });
1825
1826        (address, request_count)
1827    }
1828
1829    async fn create_hanging_quote_server() -> std::net::SocketAddr {
1830        let listener = TcpListener::bind("127.0.0.1:0")
1831            .await
1832            .unwrap();
1833        let address = listener.local_addr().unwrap();
1834
1835        tokio::spawn(async move {
1836            let (_stream, _) = listener.accept().await.unwrap();
1837            std::future::pending::<()>().await;
1838        });
1839
1840        address
1841    }
1842
1843    #[tokio::test]
1844    async fn handles_orderbook_api_error_with_http_200() {
1845        let (address, request_count) = create_quote_server(
1846            "200 OK",
1847            r#"{"code":171015,"message":"quoted token not available"}"#,
1848        )
1849        .await;
1850        let client = create_test_client(format!("http://{address}"));
1851
1852        let result = client.fetch_orderbook().await;
1853
1854        assert!(matches!(
1855            result,
1856            Err(RFQError::QuoteNotFound(message)) if message.contains("171015")
1857        ));
1858        assert_eq!(request_count.load(Ordering::SeqCst), 1);
1859    }
1860
1861    #[tokio::test]
1862    async fn handles_documented_quote_error_without_retrying() {
1863        let (address, request_count) = create_quote_server(
1864            "200 OK",
1865            r#"{"code":171015,"message":"quoted token not available"}"#,
1866        )
1867        .await;
1868        let client = create_test_client(format!("http://{address}"));
1869
1870        let result = client
1871            .request_binding_quote(&create_test_quote_params())
1872            .await;
1873
1874        match result {
1875            Err(RFQError::QuoteNotFound(message)) => {
1876                assert!(message.contains("171015"));
1877                assert!(message.contains("quoted token not available"));
1878            }
1879            other => panic!("Expected Native API error, got {other:?}"),
1880        }
1881        assert_eq!(request_count.load(Ordering::SeqCst), 1);
1882    }
1883
1884    #[tokio::test]
1885    async fn handles_success_false_as_quote_not_found_without_retrying() {
1886        let mut response = successful_quote_json("1");
1887        response["success"] = serde_json::Value::Bool(false);
1888        response["errorMessage"] = serde_json::Value::String("quote unavailable".to_string());
1889        let (address, request_count) = create_quote_server("200 OK", response.to_string()).await;
1890        let client = create_test_client(format!("http://{address}"));
1891
1892        let result = client
1893            .request_binding_quote(&create_test_quote_params())
1894            .await;
1895
1896        assert!(matches!(
1897            result,
1898            Err(RFQError::QuoteNotFound(message)) if message.contains("quote unavailable")
1899        ));
1900        assert_eq!(request_count.load(Ordering::SeqCst), 1);
1901    }
1902
1903    #[tokio::test]
1904    async fn retries_documented_temporary_api_error() {
1905        let (address, request_count) = create_quote_server(
1906            "200 OK",
1907            r#"{"code":301016,"message":"quote invalid, risk management checks failed"}"#,
1908        )
1909        .await;
1910        let client = create_test_client(format!("http://{address}"));
1911
1912        let result = client
1913            .request_binding_quote(&create_test_quote_params())
1914            .await;
1915
1916        assert!(matches!(result, Err(RFQError::QuoteNotFound(_))));
1917        assert_eq!(request_count.load(Ordering::SeqCst), 3);
1918    }
1919
1920    #[tokio::test]
1921    async fn retries_server_error_without_native_error_envelope() {
1922        let (address, request_count) =
1923            create_quote_server("503 Service Unavailable", "<html>upstream unavailable</html>")
1924                .await;
1925        let client = create_test_client(format!("http://{address}"));
1926
1927        let result = client
1928            .request_binding_quote(&create_test_quote_params())
1929            .await;
1930
1931        assert!(matches!(
1932            result,
1933            Err(RFQError::ConnectionError(message)) if message.contains("503 Service Unavailable")
1934        ));
1935        assert_eq!(request_count.load(Ordering::SeqCst), 3);
1936    }
1937
1938    #[tokio::test]
1939    async fn retries_malformed_success_response() {
1940        let (address, request_count) =
1941            create_quote_server("200 OK", r#"{"unexpected":true}"#).await;
1942        let client = create_test_client(format!("http://{address}"));
1943
1944        let result = client
1945            .request_binding_quote(&create_test_quote_params())
1946            .await;
1947
1948        assert!(matches!(result, Err(RFQError::ParsingError(_))));
1949        assert_eq!(request_count.load(Ordering::SeqCst), 3);
1950    }
1951
1952    #[tokio::test]
1953    async fn times_out_when_quote_response_stalls() {
1954        let address = create_hanging_quote_server().await;
1955        let mut client = create_test_client(format!("http://{address}"));
1956        client.quote_timeout = Duration::from_millis(50);
1957
1958        let result = tokio::time::timeout(
1959            Duration::from_secs(1),
1960            client.request_binding_quote(&create_test_quote_params()),
1961        )
1962        .await
1963        .expect("Native quote timeout did not terminate the request");
1964
1965        assert!(matches!(
1966            result,
1967            Err(RFQError::ConnectionError(message)) if message.contains("timed out after 50ms")
1968        ));
1969    }
1970
1971    #[tokio::test]
1972    async fn shares_quote_timeout_across_retries() {
1973        let listener = TcpListener::bind("127.0.0.1:0")
1974            .await
1975            .unwrap();
1976        let address = listener.local_addr().unwrap();
1977        let request_count = Arc::new(AtomicUsize::new(0));
1978        let server_request_count = request_count.clone();
1979        let params = create_test_quote_params();
1980        let success_body = successful_quote_json(&params.amount_in.to_string()).to_string();
1981        let quote_timeout = NATIVE_API_RETRY_DELAY * 2;
1982
1983        tokio::spawn(async move {
1984            let retry_body =
1985                r#"{"code":301016,"message":"quote invalid, risk management checks failed"}"#
1986                    .to_string();
1987            for body in [retry_body, success_body] {
1988                let (mut stream, _) = listener.accept().await.unwrap();
1989                let attempt = server_request_count.fetch_add(1, Ordering::SeqCst);
1990                if attempt == 1 {
1991                    // This fits a fresh timeout, but not the time left after the retry backoff.
1992                    tokio::time::sleep(quote_timeout * 3 / 4).await;
1993                }
1994                let response = format!(
1995                    "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
1996                    body.len()
1997                );
1998                let _ = stream
1999                    .write_all(response.as_bytes())
2000                    .await;
2001                let _ = stream.shutdown().await;
2002            }
2003        });
2004
2005        let mut client = create_test_client(format!("http://{address}"));
2006        client.quote_timeout = quote_timeout;
2007
2008        let result = timeout(Duration::from_secs(5), client.request_binding_quote(&params))
2009            .await
2010            .expect("Native quote timeout did not terminate the retries");
2011
2012        assert!(matches!(
2013            result,
2014            Err(RFQError::QuoteNotFound(message)) if message.contains("301016")
2015        ));
2016        assert_eq!(request_count.load(Ordering::SeqCst), 2);
2017    }
2018
2019    #[tokio::test]
2020    async fn does_not_retry_documented_authentication_error() {
2021        let (address, request_count) = create_quote_server(
2022            "200 OK",
2023            r#"{"code":201001,"message":"auth get api key is invalid"}"#,
2024        )
2025        .await;
2026        let client = create_test_client(format!("http://{address}"));
2027
2028        let result = client
2029            .request_binding_quote(&create_test_quote_params())
2030            .await;
2031
2032        assert!(matches!(result, Err(RFQError::FatalError(_))));
2033        assert_eq!(request_count.load(Ordering::SeqCst), 1);
2034    }
2035}