Skip to main content

tycho_simulation/rfq/protocols/liquorice/
client.rs

1use std::{
2    collections::{HashMap, HashSet},
3    str::FromStr,
4    sync::{
5        atomic::{AtomicBool, Ordering},
6        Arc,
7    },
8    time::SystemTime,
9};
10
11use alloy::primitives::utils::keccak256;
12use async_trait::async_trait;
13use base64::{engine::general_purpose::STANDARD as BASE64_STANDARD, Engine};
14use futures::stream::BoxStream;
15use num_bigint::BigUint;
16use reqwest::{Client, RequestBuilder, Response, StatusCode};
17use tokio::time::{interval, timeout, Duration};
18use tracing::{debug, error, info, warn};
19use tycho_common::{
20    models::{protocol::GetAmountOutParams, Chain},
21    simulation::indicatively_priced::SignedQuote,
22    Bytes,
23};
24
25use crate::{
26    evm::protocol::u256_num::biguint_to_u256,
27    rfq::{
28        client::RFQClient,
29        errors::RFQError,
30        models::TimestampHeader,
31        protocols::liquorice::models::{
32            LiquoricePriceLevelsResponse, LiquoriceQuoteRequest, LiquoriceQuoteResponse,
33            LiquoriceTokenPairPrice,
34        },
35    },
36    tycho_client::feed::synchronizer::{ComponentWithState, Snapshot, StateSyncMessage},
37    tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState},
38};
39
40#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
41pub struct LiquoriceClient {
42    chain: Chain,
43    price_levels_endpoint: String,
44    quote_endpoint: String,
45    // Tokens that we want prices for
46    tokens: HashSet<Bytes>,
47    // Min tvl value in the quote token.
48    tvl: f64,
49    // solver header for authentication
50    #[serde(skip_serializing, default)]
51    auth_solver: String,
52    // key header for authentication
53    #[serde(skip_serializing, default)]
54    auth_key: String,
55    quote_tokens: HashSet<Bytes>,
56    poll_time: Duration,
57    quote_timeout: Duration,
58    quote_expiry_secs: u64,
59    /// Liquorice issues credentials under two schemes: newer accounts use
60    /// `Authorization: Basic base64(solver:key)`, legacy accounts use the
61    /// separate `solver` + `authorization` headers. We start with Basic and
62    /// permanently fall back to the legacy scheme if it returns 401.
63    #[serde(skip)]
64    use_legacy_auth: Arc<AtomicBool>,
65}
66
67impl LiquoriceClient {
68    pub const PROTOCOL_SYSTEM: &'static str = "rfq:liquorice";
69
70    #[allow(clippy::too_many_arguments)]
71    pub fn new(
72        chain: Chain,
73        tokens: HashSet<Bytes>,
74        tvl: f64,
75        quote_tokens: HashSet<Bytes>,
76        auth_solver: String,
77        auth_key: String,
78        poll_time: Duration,
79        quote_timeout: Duration,
80        quote_expiry_secs: u64,
81    ) -> Result<Self, RFQError> {
82        Ok(Self {
83            chain,
84            price_levels_endpoint: "https://api.liquorice.tech/v1/solver/price-levels".to_string(),
85            quote_endpoint: "https://api.liquorice.tech/v1/solver/rfq".to_string(),
86            tokens,
87            tvl,
88            auth_solver,
89            auth_key,
90            quote_tokens,
91            poll_time,
92            quote_timeout,
93            quote_expiry_secs,
94            use_legacy_auth: Arc::new(AtomicBool::new(false)),
95        })
96    }
97
98    /// Build the `Authorization: Basic <base64(solver:token)>` header value
99    /// required by the Liquorice REST API.
100    fn basic_auth_header(&self) -> String {
101        let credentials = format!("{}:{}", self.auth_solver, self.auth_key);
102        format!("Basic {}", BASE64_STANDARD.encode(credentials.as_bytes()))
103    }
104
105    /// Applies whichever auth scheme is currently active to the request.
106    fn apply_auth(&self, request: RequestBuilder) -> RequestBuilder {
107        if self
108            .use_legacy_auth
109            .load(Ordering::Relaxed)
110        {
111            request
112                .header("solver", &self.auth_solver)
113                .header("authorization", &self.auth_key)
114        } else {
115            request.header("authorization", self.basic_auth_header())
116        }
117    }
118
119    /// Sends a request, retrying once with the legacy auth scheme if the
120    /// server rejects Basic auth. The working scheme is remembered for the
121    /// lifetime of the client.
122    async fn send_authed(
123        &self,
124        build: impl Fn() -> RequestBuilder,
125    ) -> Result<Response, reqwest::Error> {
126        let response = self.apply_auth(build()).send().await?;
127        if response.status() == StatusCode::UNAUTHORIZED &&
128            !self
129                .use_legacy_auth
130                .swap(true, Ordering::Relaxed)
131        {
132            warn!("Liquorice Basic auth rejected (401); retrying with legacy solver/authorization headers");
133            return self.apply_auth(build()).send().await;
134        }
135        Ok(response)
136    }
137
138    fn normalize_tvl(
139        &self,
140        raw_tvl: f64,
141        quote_token: Bytes,
142        prices_by_mm: &HashMap<String, Vec<LiquoriceTokenPairPrice>>,
143    ) -> Result<f64, RFQError> {
144        if self.quote_tokens.contains(&quote_token) {
145            return Ok(raw_tvl);
146        }
147
148        for approved_quote_token in &self.quote_tokens {
149            for token_prices in prices_by_mm.values() {
150                for token_price in token_prices {
151                    if token_price.base_token == quote_token &&
152                        token_price.quote_token == *approved_quote_token
153                    {
154                        if let Some(price) = token_price.get_price_for_amount(1.0) {
155                            return Ok(raw_tvl * price);
156                        }
157                    }
158                }
159            }
160        }
161
162        Ok(0.0)
163    }
164
165    fn create_component_with_state(
166        &self,
167        component_id: String,
168        tokens: Vec<Bytes>,
169        prices_by_mm: &HashMap<String, LiquoriceTokenPairPrice>,
170        tvl: f64,
171    ) -> ComponentWithState {
172        let protocol_component = ProtocolComponent {
173            id: component_id.clone(),
174            protocol_system: Self::PROTOCOL_SYSTEM.to_string(),
175            protocol_type_name: "liquorice_pool".to_string(),
176            chain: self.chain,
177            tokens,
178            contract_addresses: vec![],
179            ..Default::default()
180        };
181
182        let mut attributes = HashMap::new();
183
184        let prices_json = serde_json::to_string(&prices_by_mm).unwrap_or_default();
185        attributes.insert("prices".to_string(), prices_json.as_bytes().to_vec().into());
186
187        ComponentWithState {
188            state: ProtocolComponentState::new(&component_id, attributes, HashMap::new()),
189            component: protocol_component,
190            component_tvl: Some(tvl),
191            entrypoints: vec![],
192        }
193    }
194
195    fn process_quote_response(
196        quote_response: LiquoriceQuoteResponse,
197        params: &GetAmountOutParams,
198    ) -> Result<SignedQuote, RFQError> {
199        if !quote_response.liquidity_available {
200            debug!(quote_response = ?quote_response, "Liquorice quote response indicates no liquidity");
201            return Err(RFQError::QuoteNotFound(format!(
202                "Liquorice quote not found for {} {} ->{}",
203                params.amount_in, params.token_in, params.token_out,
204            )));
205        }
206
207        info!("Received Liquorice quote response with {} levels", quote_response.levels.len());
208
209        // Find the valid level with the largest quote_token_amount
210        let best_level = quote_response
211            .levels
212            .iter()
213            .filter(|level| level.validate(params).is_ok())
214            .filter_map(|level| {
215                BigUint::from_str(&level.quote_token_amount)
216                    .ok()
217                    .map(|amount| (level, amount))
218            })
219            .max_by(|(_, a), (_, b)| a.cmp(b));
220
221        let (quote_level, _) = best_level.ok_or_else(|| {
222            RFQError::QuoteNotFound(format!(
223                "No valid Liquorice quote levels for {} {} ->{}",
224                params.amount_in, params.token_in, params.token_out,
225            ))
226        })?;
227
228        let mut quote_attributes: HashMap<String, Bytes> = HashMap::new();
229
230        // calldata (pre-encoded by Liquorice API)
231        quote_attributes.insert(
232            "calldata".to_string(),
233            Bytes::from(
234                hex::decode(
235                    quote_level
236                        .tx
237                        .data
238                        .trim_start_matches("0x"),
239                )
240                .map_err(|e| RFQError::ParsingError(format!("Failed to parse calldata: {e}")))?,
241            ),
242        );
243
244        // base_token_amount as U256 (32 bytes big-endian)
245        quote_attributes.insert(
246            "base_token_amount".to_string(),
247            Bytes::from(
248                biguint_to_u256(&BigUint::from_str(&quote_level.base_token_amount).map_err(
249                    |_| {
250                        RFQError::ParsingError(format!(
251                            "Failed to parse base token amount: {}",
252                            quote_level.base_token_amount
253                        ))
254                    },
255                )?)
256                .to_be_bytes::<32>()
257                .to_vec(),
258            ),
259        );
260
261        // partial fill info (if present)
262        if let Some(pf) = &quote_level.partial_fill {
263            quote_attributes.insert(
264                "partial_fill_offset".to_string(),
265                Bytes::from(pf.offset.to_be_bytes().to_vec()),
266            );
267            quote_attributes.insert(
268                "min_base_token_amount".to_string(),
269                Bytes::from(
270                    biguint_to_u256(&BigUint::from_str(&pf.min_base_token_amount).map_err(
271                        |_| {
272                            RFQError::ParsingError(format!(
273                                "Failed to parse min_base_token_amount: {}",
274                                pf.min_base_token_amount
275                            ))
276                        },
277                    )?)
278                    .to_be_bytes::<32>()
279                    .to_vec(),
280                ),
281            );
282        }
283
284        Ok(SignedQuote {
285            base_token: params.token_in.clone(),
286            quote_token: params.token_out.clone(),
287            amount_in: BigUint::from_str(&quote_level.base_token_amount).map_err(|_| {
288                RFQError::ParsingError(format!(
289                    "Failed to parse amount in string: {}",
290                    quote_level.base_token_amount
291                ))
292            })?,
293            amount_out: BigUint::from_str(&quote_level.quote_token_amount).map_err(|_| {
294                RFQError::ParsingError(format!(
295                    "Failed to parse amount out string: {}",
296                    quote_level.quote_token_amount
297                ))
298            })?,
299            quote_attributes,
300        })
301    }
302
303    async fn fetch_price_levels(
304        &self,
305    ) -> Result<HashMap<String, Vec<LiquoriceTokenPairPrice>>, RFQError> {
306        let query_params = vec![("chainId", self.chain.id().to_string())];
307
308        let http_client = Client::new();
309        let response = self
310            .send_authed(|| {
311                http_client
312                    .get(&self.price_levels_endpoint)
313                    .query(&query_params)
314                    .header("accept", "application/json")
315            })
316            .await
317            .map_err(|e| RFQError::ConnectionError(format!("Failed to fetch price levels: {e}")))?;
318
319        if !response.status().is_success() {
320            return Err(RFQError::ConnectionError(format!(
321                "HTTP error {}: {}",
322                response.status(),
323                response
324                    .text()
325                    .await
326                    .unwrap_or_default()
327            )));
328        }
329
330        let price_response: LiquoricePriceLevelsResponse = response.json().await.map_err(|e| {
331            RFQError::ParsingError(format!("Failed to parse price levels response: {e}"))
332        })?;
333
334        Ok(price_response.prices)
335    }
336}
337
338#[async_trait]
339impl RFQClient for LiquoriceClient {
340    fn stream(
341        &self,
342    ) -> BoxStream<'static, Result<(String, StateSyncMessage<TimestampHeader>), RFQError>> {
343        let client = self.clone();
344
345        Box::pin(async_stream::stream! {
346            let mut current_components: HashMap<String, ComponentWithState> = HashMap::new();
347            let mut ticker = interval(client.poll_time);
348
349            info!("Starting Liquorice price levels polling every {} seconds", client.poll_time.as_secs());
350            info!("TVL threshold: {:.2}", client.tvl);
351
352            loop {
353                ticker.tick().await;
354
355                match client.fetch_price_levels().await {
356                    Ok(prices_by_mm) => {
357                        let mut new_components = HashMap::new();
358
359                        // Group qualifying MMs by token pair
360                        struct PricesWithTvl {
361                            // MM name -> price levels for the token pair
362                            mm_prices: HashMap<String, LiquoriceTokenPairPrice>,
363                            // The highest TVL among MMs for this token pair, used as the component's TVL
364                            tvl: f64,
365                        }
366                        let mut pair_mm_prices: HashMap<(Bytes, Bytes), PricesWithTvl> = HashMap::new();
367
368                        info!("Fetched price levels from {} market makers", prices_by_mm.len());
369                        for (mm_name, token_pair_prices) in prices_by_mm.iter() {
370                            for token_pair_price in token_pair_prices {
371                                let base_token = &token_pair_price.base_token;
372                                let quote_token = &token_pair_price.quote_token;
373
374                                if !client.tokens.contains(base_token) || !client.tokens.contains(quote_token) {
375                                    continue;
376                                }
377
378                                let tvl = token_pair_price.calculate_tvl();
379                                let normalized_tvl = client.normalize_tvl(
380                                    tvl,
381                                    token_pair_price.quote_token.clone(),
382                                    &prices_by_mm,
383                                )?;
384
385                                if normalized_tvl < client.tvl {
386                                    info!("Filtering out MM {} for pair {}/{} due to low TVL: {:.2} < {:.2}",
387                                          mm_name, hex::encode(base_token), hex::encode(quote_token),
388                                          normalized_tvl, client.tvl);
389                                    continue;
390                                }
391
392                                let entry = pair_mm_prices
393                                    .entry((base_token.clone(), quote_token.clone()))
394                                    .or_insert_with(|| PricesWithTvl { mm_prices: HashMap::new(), tvl: f64::NEG_INFINITY });
395                                entry.tvl = entry.tvl.max(normalized_tvl);
396                                entry.mm_prices.insert(mm_name.clone(), token_pair_price.clone());
397                            }
398                        }
399
400                        for ((base_token, quote_token), PricesWithTvl { mm_prices, tvl: component_tvl }) in pair_mm_prices {
401                            let pair_str = format!("liquorice_{}/{}", hex::encode(&base_token), hex::encode(&quote_token));
402                            let component_id = format!("{}", keccak256(pair_str.as_bytes()));
403
404                            let tokens = vec![base_token, quote_token];
405
406                            let component_with_state = client.create_component_with_state(
407                                component_id.clone(),
408                                tokens,
409                                &mm_prices,
410                                component_tvl,
411                            );
412                            new_components.insert(component_id, component_with_state);
413                        }
414
415                        let removed_components: HashMap<String, ProtocolComponent> = current_components
416                            .iter()
417                            .filter(|&(id, _)| !new_components.contains_key(id))
418                            .map(|(k, v)| (k.clone(), v.component.clone()))
419                            .collect();
420
421                        current_components = new_components.clone();
422
423                        let snapshot = Snapshot {
424                            states: new_components,
425                            vm_storage: HashMap::new(),
426                        };
427                        let timestamp = SystemTime::now().duration_since(
428                            SystemTime::UNIX_EPOCH
429                        ).map_err(
430                            |_| RFQError::ParsingError("SystemTime before UNIX EPOCH!".into())
431                        )?.as_secs();
432
433                        let msg = StateSyncMessage::<TimestampHeader> {
434                            header: TimestampHeader { timestamp },
435                            snapshots: snapshot,
436                            deltas: None,
437                            removed_components,
438                        };
439
440                        yield Ok(("liquorice".to_string(), msg));
441                    },
442                    Err(e) => {
443                        error!("Failed to fetch price levels from Liquorice API: {}", e);
444                        continue;
445                    }
446                }
447            }
448        })
449    }
450
451    async fn request_binding_quote(
452        &self,
453        params: &GetAmountOutParams,
454    ) -> Result<SignedQuote, RFQError> {
455        let expiry = SystemTime::now()
456            .duration_since(SystemTime::UNIX_EPOCH)
457            .map_err(|_| RFQError::ParsingError("SystemTime before UNIX EPOCH!".into()))?
458            .as_secs() +
459            self.quote_expiry_secs;
460
461        let rfq_id = uuid::Uuid::new_v4().to_string();
462
463        let quote_request = LiquoriceQuoteRequest {
464            chain_id: self.chain.id(),
465            rfq_id: rfq_id.clone(),
466            expiry,
467            base_token: params.token_in.to_string(),
468            quote_token: params.token_out.to_string(),
469            trader: params.receiver.to_string(),
470            effective_trader: Some(params.sender.to_string()),
471            base_token_amount: Some(params.amount_in.to_string()),
472            quote_token_amount: None,
473        };
474
475        debug!(quote_request = ?quote_request, "Sending Liquorice quote request");
476
477        let url = self.quote_endpoint.clone();
478
479        let start_time = std::time::Instant::now();
480        const MAX_RETRIES: u32 = 3;
481        let mut last_error = None;
482
483        for attempt in 0..MAX_RETRIES {
484            let elapsed = start_time.elapsed();
485            if elapsed >= self.quote_timeout {
486                return Err(last_error.unwrap_or_else(|| {
487                    RFQError::ConnectionError(format!(
488                        "Liquorice quote request timed out after {} seconds",
489                        self.quote_timeout.as_secs()
490                    ))
491                }));
492            }
493
494            let remaining_time = self.quote_timeout - elapsed;
495
496            let http_client = Client::new();
497            let response = match timeout(
498                remaining_time,
499                self.send_authed(|| {
500                    http_client
501                        .post(&url)
502                        .json(&quote_request)
503                        .header("accept", "application/json")
504                }),
505            )
506            .await
507            {
508                Ok(Ok(resp)) => resp,
509                Ok(Err(e)) => {
510                    warn!(
511                        "Liquorice quote request failed (attempt {}/{}): {}",
512                        attempt + 1,
513                        MAX_RETRIES,
514                        e
515                    );
516                    last_error = Some(RFQError::ConnectionError(format!(
517                        "Failed to send Liquorice quote request: {e}"
518                    )));
519                    if attempt < MAX_RETRIES - 1 {
520                        tokio::time::sleep(Duration::from_millis(100)).await;
521                        continue;
522                    } else {
523                        return Err(last_error.unwrap());
524                    }
525                }
526                Err(_) => {
527                    return Err(RFQError::ConnectionError(format!(
528                        "Liquorice quote request timed out after {} seconds",
529                        self.quote_timeout.as_secs()
530                    )));
531                }
532            };
533
534            if response.status() != 200 {
535                let err_msg = match response.text().await {
536                    Ok(text) => text,
537                    Err(e) => {
538                        warn!(
539                            "Liquorice error response parsing failed (attempt {}/{}): {}",
540                            attempt + 1,
541                            MAX_RETRIES,
542                            e
543                        );
544                        last_error = Some(RFQError::ParsingError(format!(
545                            "Failed to read response text from Liquorice failed request: {e}"
546                        )));
547                        if attempt < MAX_RETRIES - 1 {
548                            tokio::time::sleep(Duration::from_millis(100)).await;
549                            continue;
550                        } else {
551                            return Err(last_error.unwrap());
552                        }
553                    }
554                };
555                last_error = Some(RFQError::FatalError(format!(
556                    "Failed to send Liquorice quote request: {err_msg}",
557                )));
558                if attempt < MAX_RETRIES - 1 {
559                    warn!(
560                        "Liquorice returned non-200 status (attempt {}/{}): {}",
561                        attempt + 1,
562                        MAX_RETRIES,
563                        err_msg
564                    );
565                    tokio::time::sleep(Duration::from_millis(100)).await;
566                    continue;
567                } else {
568                    return Err(last_error.unwrap());
569                }
570            }
571
572            let quote_response = match response
573                .json::<LiquoriceQuoteResponse>()
574                .await
575            {
576                Ok(resp) => resp,
577                Err(e) => {
578                    warn!(
579                        "Liquorice quote response parsing failed (attempt {}/{}): {}",
580                        attempt + 1,
581                        MAX_RETRIES,
582                        e
583                    );
584                    last_error = Some(RFQError::ParsingError(format!(
585                        "Failed to parse Liquorice quote response: {e}"
586                    )));
587                    if attempt < MAX_RETRIES - 1 {
588                        tokio::time::sleep(Duration::from_millis(100)).await;
589                        continue;
590                    } else {
591                        return Err(last_error.unwrap());
592                    }
593                }
594            };
595
596            return Self::process_quote_response(quote_response, params);
597        }
598
599        Err(last_error.unwrap_or_else(|| {
600            RFQError::ConnectionError("Liquorice quote request failed after retries".to_string())
601        }))
602    }
603}
604
605#[cfg(test)]
606mod tests {
607    use std::{str::FromStr, time::Duration};
608
609    use super::*;
610    use crate::rfq::protocols::liquorice::models::{LiquoricePriceLevel, LiquoriceTokenPairPrice};
611
612    #[test]
613    fn test_normalize_tvl_same_quote_token() {
614        let client = create_test_client();
615        let prices = HashMap::new();
616
617        let result = client.normalize_tvl(
618            1000.0,
619            Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap(),
620            &prices,
621        );
622        assert!(result.is_ok());
623        assert_eq!(result.unwrap(), 1000.0);
624    }
625
626    #[test]
627    fn test_normalize_tvl_different_quote_token() {
628        let client = create_test_client();
629        let mut prices = HashMap::new();
630        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
631        let usdc = Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap();
632
633        let eth_usdc_price = LiquoriceTokenPairPrice {
634            base_token: weth.clone(),
635            quote_token: usdc,
636            levels: vec![LiquoricePriceLevel { quantity: 1.0, price: 3000.0 }],
637            updated_at: None,
638        };
639
640        prices.insert("test_mm".to_string(), vec![eth_usdc_price]);
641
642        let result = client.normalize_tvl(2.0, weth, &prices);
643        assert!(result.is_ok());
644        assert_eq!(result.unwrap(), 6000.0);
645    }
646
647    #[test]
648    fn test_normalize_tvl_no_conversion_available() {
649        let client = create_test_client();
650        let prices = HashMap::new();
651        let result = client.normalize_tvl(
652            1000.0,
653            Bytes::from_str("0x1234567890123456789012345678901234567890").unwrap(),
654            &prices,
655        );
656        assert!(result.is_ok());
657        assert_eq!(result.unwrap(), 0.0);
658    }
659
660    fn create_test_client() -> LiquoriceClient {
661        let quote_tokens = HashSet::from([
662            Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap(), // USDC
663            Bytes::from_str("0xdAC17F958D2ee523a2206206994597C13D831ec7").unwrap(), // USDT
664        ]);
665
666        LiquoriceClient::new(
667            Chain::Ethereum,
668            HashSet::new(),
669            1.0,
670            quote_tokens,
671            "test_solver".to_string(),
672            "test_key".to_string(),
673            Duration::from_secs(5),
674            Duration::from_secs(5),
675            300,
676        )
677        .unwrap()
678    }
679
680    async fn create_delayed_response_server(delay_ms: u64) -> std::net::SocketAddr {
681        use tokio::{io::AsyncWriteExt, net::TcpListener};
682
683        let listener = TcpListener::bind("127.0.0.1:0")
684            .await
685            .unwrap();
686        let addr = listener.local_addr().unwrap();
687
688        let json_response = r#"{"rfqId":"test-rfq-id","liquidityAvailable":true,"levels":[{"makerRfqId":"maker-rfq-1","maker":"test-maker","nonce":"0x0000000000000000000000000000000000000000000000000000000000000001","expiry":1707847360,"tx":{"to":"0x71D9750ECF0c5081FAE4E3EDC4253E52024b0B59","data":"0xdeadbeef"},"baseToken":"0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2","quoteToken":"0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599","baseTokenAmount":"1000000000000000000","quoteTokenAmount":"3329502","partialFill":null,"allowances":[]}]}"#;
689
690        tokio::spawn(async move {
691            while let Ok((mut stream, _)) = listener.accept().await {
692                let json_response_clone = json_response.to_owned();
693                tokio::spawn(async move {
694                    tokio::time::sleep(Duration::from_millis(delay_ms)).await;
695                    let response = format!(
696                        "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
697                        json_response_clone.len(),
698                        json_response_clone
699                    );
700                    let _ = stream
701                        .write_all(response.as_bytes())
702                        .await;
703                    let _ = stream.flush().await;
704                    let _ = stream.shutdown().await;
705                });
706            }
707        });
708
709        tokio::time::sleep(Duration::from_millis(50)).await;
710        addr
711    }
712
713    fn create_test_liquorice_client(
714        quote_endpoint: String,
715        quote_timeout: Duration,
716    ) -> LiquoriceClient {
717        let token_in = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
718        let token_out = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
719
720        LiquoriceClient {
721            chain: Chain::Ethereum,
722            price_levels_endpoint: "http://unused/price-levels".to_string(),
723            quote_endpoint,
724            tokens: HashSet::from([token_in, token_out]),
725            tvl: 10.0,
726            auth_solver: "test_solver".to_string(),
727            auth_key: "test_key".to_string(),
728            quote_tokens: HashSet::new(),
729            poll_time: Duration::from_secs(0),
730            quote_timeout,
731            quote_expiry_secs: 300,
732            use_legacy_auth: Arc::new(AtomicBool::new(false)),
733        }
734    }
735
736    fn make_quote_level(
737        base_token: &str,
738        quote_token: &str,
739        base_token_amount: &str,
740        quote_token_amount: &str,
741        partial_fill: Option<crate::rfq::protocols::liquorice::models::LiquoricePartialFill>,
742    ) -> crate::rfq::protocols::liquorice::models::LiquoriceQuoteLevel {
743        use crate::rfq::protocols::liquorice::models::{LiquoriceQuoteLevel, LiquoriceTx};
744        LiquoriceQuoteLevel {
745            maker_rfq_id: "maker-1".to_string(),
746            maker: "test-maker".to_string(),
747            expiry: 9999999999,
748            tx: LiquoriceTx {
749                to: "0x1111111111111111111111111111111111111111".to_string(),
750                data: "0xdeadbeef".to_string(),
751            },
752            base_token: base_token.to_string(),
753            quote_token: quote_token.to_string(),
754            base_token_amount: base_token_amount.to_string(),
755            quote_token_amount: quote_token_amount.to_string(),
756            partial_fill,
757        }
758    }
759
760    fn make_params(token_in: &str, token_out: &str, amount_in: u64) -> GetAmountOutParams {
761        GetAmountOutParams {
762            amount_in: BigUint::from(amount_in),
763            token_in: Bytes::from_str(token_in).unwrap(),
764            token_out: Bytes::from_str(token_out).unwrap(),
765            sender: Bytes::from_str("0x3333333333333333333333333333333333333333").unwrap(),
766            receiver: Bytes::from_str("0x4444444444444444444444444444444444444444").unwrap(),
767        }
768    }
769
770    #[test]
771    fn test_process_quote_response_no_liquidity() {
772        use crate::rfq::protocols::liquorice::models::LiquoriceQuoteResponse;
773
774        let response = LiquoriceQuoteResponse {
775            rfq_id: "r1".to_string(),
776            liquidity_available: false,
777            levels: vec![],
778        };
779        let params = make_params(
780            "0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2",
781            "0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599",
782            1_000_000_000_000_000_000,
783        );
784        let result = LiquoriceClient::process_quote_response(response, &params);
785        assert!(
786            matches!(result, Err(RFQError::QuoteNotFound(_))),
787            "expected QuoteNotFound, got {:?}",
788            result
789        );
790    }
791
792    #[test]
793    fn test_process_quote_response_partial_fill_attributes() {
794        use crate::rfq::protocols::liquorice::models::{
795            LiquoricePartialFill, LiquoriceQuoteResponse,
796        };
797
798        let token_in = "0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2";
799        let token_out = "0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599";
800        let amount_in = 1_000_000_000_000_000_000u64;
801
802        let level = make_quote_level(
803            token_in,
804            token_out,
805            &amount_in.to_string(),
806            "3329502",
807            Some(LiquoricePartialFill {
808                offset: 68,
809                min_base_token_amount: "500000000000000000".to_string(),
810            }),
811        );
812        let response = LiquoriceQuoteResponse {
813            rfq_id: "r1".to_string(),
814            liquidity_available: true,
815            levels: vec![level],
816        };
817        let params = make_params(token_in, token_out, amount_in);
818
819        let quote = LiquoriceClient::process_quote_response(response, &params).unwrap();
820
821        // partial_fill_offset: 4-byte big-endian encoding of 68
822        let offset_bytes = quote.quote_attributes["partial_fill_offset"].clone();
823        assert_eq!(offset_bytes.as_ref(), &68u32.to_be_bytes());
824
825        // min_base_token_amount: 32-byte big-endian U256 of 500000000000000000
826        let min_amount_bytes = quote.quote_attributes["min_base_token_amount"].clone();
827        assert_eq!(min_amount_bytes.len(), 32);
828        let expected_min = BigUint::from(500_000_000_000_000_000u64);
829        let actual_min = BigUint::from_bytes_be(min_amount_bytes.as_ref());
830        assert_eq!(actual_min, expected_min);
831    }
832
833    #[test]
834    fn test_process_quote_response_selects_best_valid_level() {
835        use crate::rfq::protocols::liquorice::models::LiquoriceQuoteResponse;
836
837        let token_in = "0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2";
838        let token_out = "0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599";
839        let amount_in = 1_000_000_000_000_000_000u64;
840
841        // One level with wrong base_token_amount (will fail validate) and two
842        // valid levels with different quote_token_amounts; expect the higher one
843        // to be chosen.
844        let invalid_level = make_quote_level(token_in, token_out, "999", "9999999", None);
845        let lower_level =
846            make_quote_level(token_in, token_out, &amount_in.to_string(), "3000000", None);
847        let best_level =
848            make_quote_level(token_in, token_out, &amount_in.to_string(), "3500000", None);
849
850        let response = LiquoriceQuoteResponse {
851            rfq_id: "r1".to_string(),
852            liquidity_available: true,
853            levels: vec![invalid_level, lower_level, best_level],
854        };
855        let params = make_params(token_in, token_out, amount_in);
856
857        let quote = LiquoriceClient::process_quote_response(response, &params).unwrap();
858        assert_eq!(quote.amount_out, BigUint::from(3_500_000u64));
859    }
860
861    fn create_test_quote_params() -> GetAmountOutParams {
862        let token_in = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
863        let token_out = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
864        let router = Bytes::from_str("0xfD0b31d2E955fA55e3fa641Fe90e08b677188d35").unwrap();
865
866        GetAmountOutParams {
867            amount_in: BigUint::from(1_000000000000000000u64),
868            token_in,
869            token_out,
870            sender: router.clone(),
871            receiver: router,
872        }
873    }
874
875    #[tokio::test]
876    async fn test_liquorice_quote_timeout() {
877        let addr = create_delayed_response_server(500).await;
878
879        let client_short_timeout = create_test_liquorice_client(
880            format!("http://127.0.0.1:{}/rfq", addr.port()),
881            Duration::from_millis(200),
882        );
883        let params = create_test_quote_params();
884
885        let start = std::time::Instant::now();
886        let result = client_short_timeout
887            .request_binding_quote(&params)
888            .await;
889        let elapsed = start.elapsed();
890
891        assert!(result.is_err());
892        let err = result.unwrap_err();
893        match err {
894            RFQError::ConnectionError(msg) => {
895                assert!(msg.contains("timed out"), "Expected timeout error, got: {}", msg);
896            }
897            _ => panic!("Expected ConnectionError, got: {:?}", err),
898        }
899        assert!(
900            elapsed.as_millis() >= 200 && elapsed.as_millis() < 400,
901            "Expected timeout around 200ms, got: {:?}",
902            elapsed
903        );
904
905        let client_long_timeout = create_test_liquorice_client(
906            format!("http://127.0.0.1:{}/rfq", addr.port()),
907            Duration::from_secs(1),
908        );
909
910        let result = client_long_timeout
911            .request_binding_quote(&params)
912            .await;
913        assert!(result.is_ok(), "Expected success, got: {:?}", result);
914    }
915
916    async fn create_retry_server() -> (std::net::SocketAddr, std::sync::Arc<std::sync::Mutex<u32>>)
917    {
918        use std::sync::{Arc, Mutex};
919
920        use tokio::{io::AsyncWriteExt, net::TcpListener};
921
922        let request_count = Arc::new(Mutex::new(0u32));
923        let request_count_clone = request_count.clone();
924
925        let listener = TcpListener::bind("127.0.0.1:0")
926            .await
927            .unwrap();
928        let addr = listener.local_addr().unwrap();
929
930        let json_response = r#"{"rfqId":"test-rfq-id","liquidityAvailable":true,"levels":[{"makerRfqId":"maker-rfq-1","maker":"test-maker","nonce":"0x0000000000000000000000000000000000000000000000000000000000000001","expiry":1707847360,"tx":{"to":"0x71D9750ECF0c5081FAE4E3EDC4253E52024b0B59","data":"0xdeadbeef"},"baseToken":"0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2","quoteToken":"0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599","baseTokenAmount":"1000000000000000000","quoteTokenAmount":"3329502","partialFill":null,"allowances":[]}]}"#;
931
932        tokio::spawn(async move {
933            while let Ok((mut stream, _)) = listener.accept().await {
934                let count_clone = request_count_clone.clone();
935                let json_response_clone = json_response.to_owned();
936                tokio::spawn(async move {
937                    *count_clone.lock().unwrap() += 1;
938                    let count = *count_clone.lock().unwrap();
939                    println!("Mock server: Received request #{count}");
940
941                    if count <= 2 {
942                        let response = "HTTP/1.1 500 Internal Server Error\r\nContent-Length: 21\r\n\r\nInternal Server Error";
943                        let _ = stream
944                            .write_all(response.as_bytes())
945                            .await;
946                    } else {
947                        let response = format!(
948                            "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
949                            json_response_clone.len(),
950                            json_response_clone
951                        );
952                        let _ = stream
953                            .write_all(response.as_bytes())
954                            .await;
955                    }
956                    let _ = stream.flush().await;
957                    let _ = stream.shutdown().await;
958                });
959            }
960        });
961
962        tokio::time::sleep(Duration::from_millis(50)).await;
963        (addr, request_count)
964    }
965
966    #[tokio::test]
967    async fn test_liquorice_quote_retry_on_bad_response() {
968        let (addr, request_count) = create_retry_server().await;
969
970        let client = create_test_liquorice_client(
971            format!("http://127.0.0.1:{}/rfq", addr.port()),
972            Duration::from_secs(5),
973        );
974        let params = create_test_quote_params();
975        let result = client
976            .request_binding_quote(&params)
977            .await;
978
979        assert!(result.is_ok(), "Expected success after retries, got: {:?}", result);
980        let quote = result.unwrap();
981
982        assert_eq!(quote.amount_in, BigUint::from(1_000000000000000000u64));
983        assert_eq!(quote.amount_out, BigUint::from(3329502u64));
984
985        let final_count = *request_count.lock().unwrap();
986        assert_eq!(final_count, 3, "Expected 3 requests, got {}", final_count);
987    }
988}