Skip to main content

tycho_simulation/rfq/protocols/native/
client.rs

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