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