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: HashSet<Bytes>,
47 tvl: f64,
49 #[serde(skip_serializing, default)]
51 auth_solver: String,
52 #[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 #[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 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 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 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("e_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 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 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 quote_attributes.insert(
246 "base_token_amount".to_string(),
247 Bytes::from(
248 biguint_to_u256(&BigUint::from_str("e_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 if let Some(pf) = "e_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("e_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("e_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 struct PricesWithTvl {
361 mm_prices: HashMap<String, LiquoriceTokenPairPrice>,
363 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("e_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("e_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(), Bytes::from_str("0xdAC17F958D2ee523a2206206994597C13D831ec7").unwrap(), ]);
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, ¶ms);
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, ¶ms).unwrap();
820
821 let offset_bytes = quote.quote_attributes["partial_fill_offset"].clone();
823 assert_eq!(offset_bytes.as_ref(), &68u32.to_be_bytes());
824
825 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 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, ¶ms).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(¶ms)
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(¶ms)
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(¶ms)
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}