Skip to main content

tycho_simulation/rfq/protocols/bebop/
client.rs

1use std::{
2    collections::{HashMap, HashSet},
3    str::FromStr,
4    time::SystemTime,
5};
6
7use alloy::primitives::{utils::keccak256, Address};
8use async_trait::async_trait;
9use futures::{stream::BoxStream, StreamExt};
10use http::Request;
11use num_bigint::BigUint;
12use prost::Message as ProstMessage;
13use reqwest::Client;
14use serde::{Deserialize, Serialize};
15use tokio::time::{sleep, timeout, Duration};
16use tokio_tungstenite::{
17    connect_async_with_config,
18    tungstenite::{handshake::client::generate_key, Message},
19};
20use tracing::{error, info, warn};
21use tycho_common::{
22    models::{protocol::GetAmountOutParams, Chain},
23    simulation::indicatively_priced::SignedQuote,
24    Bytes,
25};
26
27use crate::{
28    rfq::{
29        client::RFQClient,
30        errors::RFQError,
31        models::TimestampHeader,
32        protocols::bebop::models::{
33            BebopOrderToSign, BebopPriceData, BebopPricingUpdate, BebopQuoteResponse,
34        },
35    },
36    tycho_client::feed::synchronizer::{ComponentWithState, Snapshot, StateSyncMessage},
37    tycho_common::models::protocol::{ProtocolComponent, ProtocolComponentState},
38};
39
40fn bytes_to_address(address: &Bytes) -> Result<Address, RFQError> {
41    if address.len() == 20 {
42        Ok(Address::from_slice(address))
43    } else {
44        Err(RFQError::InvalidInput(format!("Invalid ERC20 token address: {address:?}")))
45    }
46}
47
48/// Maps a Chain to its corresponding Bebop WebSocket URL
49fn chain_to_bebop_url(chain: Chain) -> Result<String, RFQError> {
50    let chain_path = match chain {
51        Chain::Ethereum => "ethereum",
52        Chain::Base => "base",
53        _ => return Err(RFQError::FatalError(format!("Unsupported chain: {chain:?}"))),
54    };
55    let url = format!("api.bebop.xyz/pmm/{chain_path}/v3");
56    Ok(url)
57}
58
59#[derive(Clone, Debug, Serialize, Deserialize)]
60pub struct BebopClient {
61    chain: Chain,
62    price_ws: String,
63    quote_endpoint: String,
64    // Tokens that we want prices for
65    tokens: HashSet<Bytes>,
66    // Min tvl value in the quote token.
67    tvl: f64,
68    // key header for authentication
69    #[serde(skip_serializing, default)]
70    ws_key: String,
71    // quote tokens to normalize to for TVL purposes. Should have the same prices.
72    quote_tokens: HashSet<Bytes>,
73    quote_timeout: Duration,
74    /// The real end-user's EOA when the taker is not the end-user's own wallet.
75    origin_address: Option<Bytes>,
76    /// The `to` address of the resulting transaction when a contract executes the swap.
77    origin_target: Option<Bytes>,
78    /// Stable identifier for the upstream flow source when aggregating multiple sources.
79    origin_source: Option<String>,
80    #[serde(default = "default_protocol_system")]
81    protocol_system: String,
82}
83
84fn default_protocol_system() -> String {
85    BebopClient::PROTOCOL_SYSTEM.to_string()
86}
87
88impl BebopClient {
89    pub const PROTOCOL_SYSTEM: &'static str = "rfq:bebop";
90    /// Components executed through Tycho's `BebopFallbackRouter`.
91    pub const FALLBACK_PROTOCOL_SYSTEM: &'static str = "fallback:rfq:bebop";
92
93    pub(super) fn via_fallback_router(mut self) -> Self {
94        self.protocol_system = Self::FALLBACK_PROTOCOL_SYSTEM.to_string();
95        self
96    }
97
98    /// Creates a fully configured client. Prefer constructing through
99    /// [`BebopClientBuilder`](super::client_builder::BebopClientBuilder).
100    #[allow(clippy::too_many_arguments)]
101    pub fn new(
102        chain: Chain,
103        tokens: HashSet<Bytes>,
104        tvl: f64,
105        ws_key: String,
106        quote_tokens: HashSet<Bytes>,
107        quote_timeout: Duration,
108        origin_address: Option<Bytes>,
109        origin_target: Option<Bytes>,
110        origin_source: Option<String>,
111    ) -> Result<Self, RFQError> {
112        let url = chain_to_bebop_url(chain)?;
113        Ok(Self {
114            price_ws: "wss://".to_string() + &url + "/pricing?format=protobuf",
115            quote_endpoint: "https://".to_string() + &url + "/quote",
116            tokens,
117            chain,
118            tvl,
119            ws_key,
120            quote_tokens,
121            quote_timeout,
122            origin_address,
123            origin_target,
124            origin_source,
125            protocol_system: Self::PROTOCOL_SYSTEM.to_string(),
126        })
127    }
128
129    fn create_component_with_state(
130        &self,
131        component_id: String,
132        tokens: Vec<tycho_common::Bytes>,
133        price_data: &BebopPriceData,
134        tvl: f64,
135    ) -> ComponentWithState {
136        let protocol_component = ProtocolComponent {
137            id: component_id.clone(),
138            protocol_system: self.protocol_system.clone(),
139            protocol_type_name: "bebop_pool".to_string(),
140            chain: self.chain,
141            tokens,
142            contract_addresses: vec![], // empty for RFQ
143            static_attributes: Default::default(),
144            change: Default::default(),
145            creation_tx: Default::default(),
146            created_at: Default::default(),
147        };
148
149        let mut attributes = HashMap::new();
150
151        // Store all bids and asks as JSON strings, since we cannot store arrays
152        // Convert flat arrays [price1, size1, price2, size2, ...] to pairs [(price1, size1),
153        // (price2, size2), ...]
154        if !price_data.bids.is_empty() {
155            let bids_pairs: Vec<(f32, f32)> = price_data
156                .bids
157                .as_chunks::<2>()
158                .0
159                .iter()
160                .map(|chunk| (chunk[0], chunk[1]))
161                .collect();
162            let bids_json = serde_json::to_string(&bids_pairs).unwrap_or_default();
163            attributes.insert("bids".to_string(), bids_json.as_bytes().to_vec().into());
164        }
165        if !price_data.asks.is_empty() {
166            let asks_pairs: Vec<(f32, f32)> = price_data
167                .asks
168                .as_chunks::<2>()
169                .0
170                .iter()
171                .map(|chunk| (chunk[0], chunk[1]))
172                .collect();
173            let asks_json = serde_json::to_string(&asks_pairs).unwrap_or_default();
174            attributes.insert("asks".to_string(), asks_json.as_bytes().to_vec().into());
175        }
176
177        ComponentWithState {
178            state: ProtocolComponentState::new(&component_id, attributes, HashMap::new()),
179            component: protocol_component,
180            component_tvl: Some(tvl),
181            entrypoints: vec![],
182        }
183    }
184
185    fn process_quote_response(
186        quote_response: BebopQuoteResponse,
187        params: &GetAmountOutParams,
188    ) -> Result<SignedQuote, RFQError> {
189        match quote_response {
190            BebopQuoteResponse::Success(quote) => {
191                quote.validate(params)?;
192
193                let mut quote_attributes: HashMap<String, Bytes> = HashMap::new();
194                // The contract the calldata targets: either the Bebop settlement
195                // or the Bebop router.
196                quote_attributes.insert("tx_to".into(), quote.tx.to);
197                quote_attributes.insert("calldata".into(), quote.tx.data);
198                quote_attributes.insert(
199                    "partial_fill_offset".into(),
200                    Bytes::from(
201                        quote
202                            .partial_fill_offset
203                            .to_be_bytes()
204                            .to_vec(),
205                    ),
206                );
207                let signed_quote = match quote.to_sign {
208                    BebopOrderToSign::Single(ref single) => SignedQuote {
209                        base_token: params.token_in.clone(),
210                        quote_token: params.token_out.clone(),
211                        amount_in: BigUint::from_str(&single.taker_amount).map_err(|_| {
212                            RFQError::ParsingError(format!(
213                                "Failed to parse amount in string: {}",
214                                single.taker_amount
215                            ))
216                        })?,
217                        amount_out: BigUint::from_str(&single.maker_amount).map_err(|_| {
218                            RFQError::ParsingError(format!(
219                                "Failed to parse amount out string: {}",
220                                single.maker_amount
221                            ))
222                        })?,
223                        quote_attributes,
224                    },
225                    BebopOrderToSign::Aggregate(aggregate) => {
226                        // Sum taker_amounts for taker_tokens matching the token_in
227                        let amount_in: BigUint = aggregate
228                            .taker_tokens
229                            .iter()
230                            .zip(&aggregate.taker_amounts)
231                            .flat_map(|(tokens, amounts)| {
232                                tokens
233                                    .iter()
234                                    .zip(amounts)
235                                    .filter_map(|(token, amount)| {
236                                        if token == &params.token_in {
237                                            BigUint::from_str(amount).ok()
238                                        } else {
239                                            None
240                                        }
241                                    })
242                            })
243                            .sum();
244
245                        // Sum maker_amounts for maker_tokens matching the token_out
246                        let amount_out: BigUint = aggregate
247                            .maker_tokens
248                            .iter()
249                            .zip(&aggregate.maker_amounts)
250                            .flat_map(|(tokens, amounts)| {
251                                tokens
252                                    .iter()
253                                    .zip(amounts)
254                                    .filter_map(|(token, amount)| {
255                                        if token == &params.token_out {
256                                            BigUint::from_str(amount).ok()
257                                        } else {
258                                            None
259                                        }
260                                    })
261                            })
262                            .sum();
263
264                        SignedQuote {
265                            base_token: params.token_in.clone(),
266                            quote_token: params.token_out.clone(),
267                            amount_in,
268                            amount_out,
269                            quote_attributes,
270                        }
271                    }
272                };
273
274                Ok(signed_quote)
275            }
276            BebopQuoteResponse::Error(err) => Err(RFQError::FatalError(format!(
277                "Bebop API error: code {} - {} (requestId: {})",
278                err.error.error_code, err.error.message, err.error.request_id
279            ))),
280        }
281    }
282}
283
284#[async_trait]
285impl RFQClient for BebopClient {
286    fn stream(
287        &self,
288    ) -> BoxStream<'static, Result<(String, StateSyncMessage<TimestampHeader>), RFQError>> {
289        let tokens = self.tokens.clone();
290        let url = self.price_ws.clone();
291        let tvl_threshold = self.tvl;
292        let authorization = format!("Bearer {}", self.ws_key);
293        let client = self.clone();
294
295        Box::pin(async_stream::stream! {
296            let mut current_components: HashMap<String, ComponentWithState> = HashMap::new();
297            let mut consecutive_failures = 0;
298            const MAX_CONSECUTIVE_FAILURES: u32 = 10;
299
300            loop {
301                let request = Request::builder()
302                    .method("GET")
303                    .uri(&url)
304                    .header("Host", "api.bebop.xyz")
305                    .header("Upgrade", "websocket")
306                    .header("Connection", "Upgrade")
307                    .header("Sec-WebSocket-Key", generate_key())
308                    .header("Sec-WebSocket-Version", "13")
309                    .header("Authorization", &authorization)
310                    .body(())
311                    .map_err(|_| RFQError::FatalError("Failed to build request".into()))?;
312
313                // Connect to Bebop WebSocket with custom headers
314                let (ws_stream, _) = match connect_async_with_config(request, None, false).await {
315                    Ok(connection) => {
316                        info!("Successfully connected to Bebop WebSocket");
317                        connection
318                    },
319                    Err(e) => {
320                        consecutive_failures += 1;
321                        error!("Failed to connect to Bebop WebSocket (consecutive failure {}): {}", consecutive_failures, e);
322
323                        if consecutive_failures >= MAX_CONSECUTIVE_FAILURES {
324                            yield Err(RFQError::ConnectionError(format!("Failed to connect after {MAX_CONSECUTIVE_FAILURES} consecutive failures: {e}")));
325                            return;
326                        }
327
328                        let backoff_duration = Duration::from_secs(2_u64.pow(consecutive_failures.min(5)));
329                        info!("Retrying connection in {} seconds...", backoff_duration.as_secs());
330                        sleep(backoff_duration).await;
331                        continue;
332                    }
333                };
334
335                let (_, mut ws_receiver) = ws_stream.split();
336
337                // Message processing loop
338                while let Some(msg) = ws_receiver.next().await {
339                    match msg {
340                        Ok(Message::Binary(data)) => {
341                            match BebopPricingUpdate::decode(&data[..]) {
342                                Ok(protobuf_update) => {
343                                    // A completed handshake says nothing about whether the
344                                    // connection works, so only pricing data clears the counter.
345                                    consecutive_failures = 0;
346
347                                    let mut new_components = HashMap::new();
348
349                                    // Process all pairs directly from protobuf
350                                    for price_data in &protobuf_update.pairs {
351                                        let base_bytes = Bytes::from(price_data.base.clone());
352                                        let quote_bytes = Bytes::from(price_data.quote.clone());
353                                        if tokens.contains(&base_bytes) && tokens.contains(&quote_bytes) {
354                                            let pair_tokens = vec![
355                                                base_bytes.clone(), quote_bytes.clone()
356                                            ];
357
358                                            let mut quote_price_data: Option<&BebopPriceData> = None;
359                                            // The quote token is not one of the approved quote tokens
360                                            // Get the price, so we can normalize our TVL calculation
361                                            if !client.quote_tokens.contains(&quote_bytes) {
362                                                for approved_quote_token in &client.quote_tokens {
363                                                    // Look for a pair containing both our quote token and an approved token
364                                                    // Can be either QUOTE/APPROVED or APPROVED/QUOTE
365                                                    if let Some(quote_data) = protobuf_update.pairs.iter()
366                                                        .find(|p| {
367                                                            (p.base == quote_bytes.as_ref() && p.quote == approved_quote_token.as_ref()) ||
368                                                            (p.quote == quote_bytes.as_ref() && p.base == approved_quote_token.as_ref())
369                                                        }) {
370                                                        quote_price_data = Some(quote_data);
371                                                        break;
372                                                    }
373                                                }
374
375                                                // Quote token doesn't have price levels in approved quote tokens.
376                                                // Skip.
377                                                if quote_price_data.is_none() {
378                                                    warn!("Quote token {} does not have price levels in approved quote token. Skipping.", hex::encode(&quote_bytes));
379                                                    continue;
380                                                }
381                                            }
382
383                                            let tvl = price_data.calculate_tvl(quote_price_data);
384                                            if tvl < tvl_threshold {
385                                                continue;
386                                            }
387
388                                            let pair_str = format!("bebop_{}/{}", hex::encode(&base_bytes), hex::encode(&quote_bytes));
389                                            let component_id = format!("{}", keccak256(pair_str.as_bytes()));
390                                            let component_with_state = client.create_component_with_state(
391                                                component_id.clone(),
392                                                pair_tokens,
393                                                price_data,
394                                                tvl
395                                            );
396                                            new_components.insert(component_id, component_with_state);
397                                        }
398                                    }
399
400                                    // Find components that were removed (existed before but not in this update)
401                                    // This includes components with no bids or asks, since they are filtered
402                                    // out by the tvl threshold.
403                                    let removed_components: HashMap<String, ProtocolComponent> = current_components
404                                        .iter()
405                                        .filter(|&(id, _)| !new_components.contains_key(id))
406                                        .map(|(k, v)| (k.clone(), v.component.clone()))
407                                        .collect();
408
409                                    // Update our current state
410                                    current_components = new_components.clone();
411
412                                    let snapshot = Snapshot {
413                                        states: new_components,
414                                        vm_storage: HashMap::new(),
415                                    };
416                                    let timestamp = SystemTime::now().duration_since(
417                                        SystemTime::UNIX_EPOCH
418                                    ).map_err(
419                                        |_| RFQError::ParsingError("SystemTime before UNIX EPOCH!".into())
420                                    )?.as_secs();
421
422                                    let msg = StateSyncMessage::<TimestampHeader> {
423                                        header: TimestampHeader { timestamp },
424                                        snapshots: snapshot,
425                                        deltas: None, // Deltas are always None - all the changes are absolute
426                                        removed_components,
427                                    };
428
429                                    // Yield one message containing all updated pairs
430                                    yield Ok(("bebop".to_string(), msg));
431                                },
432                                Err(e) => {
433                                    error!("Failed to parse protobuf message: {}", e);
434                                    break;
435                                }
436                            }
437                        }
438                        Ok(Message::Close(frame)) => {
439                            match frame {
440                                Some(frame) => warn!("WebSocket closed by server: {frame}"),
441                                None => warn!("WebSocket closed by server without a close frame"),
442                            }
443                            break;
444                        }
445                        Err(e) => {
446                            error!("WebSocket error: {}", e);
447                            break;
448                        }
449                        _ => {} // Ignore other message types
450                    }
451                }
452
453                // If we're here, the message loop exited - always attempt to reconnect.
454                // Pricing data resets this, so it only grows while the feed stays unusable.
455                consecutive_failures += 1;
456                if consecutive_failures >= MAX_CONSECUTIVE_FAILURES {
457                    yield Err(RFQError::ConnectionError(format!("No pricing data received after {MAX_CONSECUTIVE_FAILURES} consecutive failures")));
458                    return;
459                }
460
461                let backoff_duration = Duration::from_secs(2_u64.pow(consecutive_failures.min(5)));
462                info!("Reconnecting in {} seconds (consecutive failure {})...", backoff_duration.as_secs(), consecutive_failures);
463                sleep(backoff_duration).await;
464                // Continue to the next iteration of the main loop
465            }
466        })
467    }
468
469    async fn request_binding_quote(
470        &self,
471        params: &GetAmountOutParams,
472    ) -> Result<SignedQuote, RFQError> {
473        let sell_token = bytes_to_address(&params.token_in)?.to_string();
474        let buy_token = bytes_to_address(&params.token_out)?.to_string();
475        let sell_amount = params.amount_in.to_string();
476        let sender = bytes_to_address(&params.sender)?.to_string();
477        let receiver = bytes_to_address(&params.receiver)?.to_string();
478
479        let url = self.quote_endpoint.clone();
480
481        let mut query = vec![
482            ("sell_tokens", sell_token),
483            ("buy_tokens", buy_token),
484            ("sell_amounts", sell_amount),
485            ("taker_address", sender),
486            ("receiver_address", receiver),
487            ("approval_type", "Standard".into()),
488            ("skip_validation", "true".into()),
489            ("skip_taker_checks", "true".into()),
490            ("gasless", "false".into()),
491            ("expiry_type", "standard".into()),
492            ("fee", "0".into()),
493            ("is_ui", "false".into()),
494        ];
495        if let Some(origin_address) = &self.origin_address {
496            query.push(("origin_address", bytes_to_address(origin_address)?.to_string()));
497        }
498        if let Some(origin_target) = &self.origin_target {
499            query.push(("origin_target", bytes_to_address(origin_target)?.to_string()));
500        }
501        if let Some(origin_source) = &self.origin_source {
502            query.push(("origin_source", origin_source.clone()));
503        }
504
505        let client = Client::new();
506
507        let start_time = std::time::Instant::now();
508        const MAX_RETRIES: u32 = 3;
509        let mut last_error = None;
510
511        for attempt in 0..MAX_RETRIES {
512            // Check if we have time remaining for this attempt
513            let elapsed = start_time.elapsed();
514            if elapsed >= self.quote_timeout {
515                return Err(last_error.unwrap_or_else(|| {
516                    RFQError::ConnectionError(format!(
517                        "Bebop quote request timed out after {} seconds",
518                        self.quote_timeout.as_secs()
519                    ))
520                }));
521            }
522
523            let remaining_time = self.quote_timeout - elapsed;
524
525            let request = client
526                .get(&url)
527                .query(&query)
528                .header("accept", "application/json")
529                .bearer_auth(&self.ws_key);
530
531            let response = match timeout(remaining_time, request.send()).await {
532                Ok(Ok(resp)) => resp,
533                Ok(Err(e)) => {
534                    warn!(
535                        "Bebop quote request failed (attempt {}/{}): {}",
536                        attempt + 1,
537                        MAX_RETRIES,
538                        e
539                    );
540                    last_error = Some(RFQError::ConnectionError(format!(
541                        "Failed to send Bebop quote request: {e}"
542                    )));
543                    if attempt < MAX_RETRIES - 1 {
544                        continue;
545                    } else {
546                        return Err(last_error.unwrap());
547                    }
548                }
549                Err(_) => {
550                    return Err(RFQError::ConnectionError(format!(
551                        "Bebop quote request timed out after {} seconds",
552                        self.quote_timeout.as_secs()
553                    )));
554                }
555            };
556
557            let quote_response = match response
558                .json::<BebopQuoteResponse>()
559                .await
560            {
561                Ok(resp) => resp,
562                Err(e) => {
563                    warn!(
564                        "Bebop quote response parsing failed (attempt {}/{}): {}",
565                        attempt + 1,
566                        MAX_RETRIES,
567                        e
568                    );
569                    last_error = Some(RFQError::ParsingError(format!(
570                        "Failed to parse Bebop quote response: {e}"
571                    )));
572                    if attempt < MAX_RETRIES - 1 {
573                        sleep(Duration::from_millis(100)).await;
574                        continue;
575                    } else {
576                        return Err(last_error.unwrap());
577                    }
578                }
579            };
580
581            return Self::process_quote_response(quote_response, params);
582        }
583
584        Err(last_error.unwrap_or_else(|| {
585            RFQError::ConnectionError("Bebop quote request failed after retries".to_string())
586        }))
587    }
588}
589
590#[cfg(test)]
591mod tests {
592    use std::{
593        sync::{Arc, Mutex},
594        time::Duration,
595    };
596
597    use dotenv::dotenv;
598    use futures::SinkExt;
599    use tokio::{net::TcpListener, time::timeout};
600    use tokio_tungstenite::accept_async;
601
602    use super::*;
603    use crate::rfq::{
604        constants::get_bebop_auth, protocols::bebop::client_builder::BebopClientBuilder,
605    };
606
607    /// BebopSettlement.swapSingle
608    const SWAP_SINGLE_SELECTOR: [u8; 4] = [0x4d, 0xce, 0xbc, 0xba];
609    /// BebopSettlement.swapAggregate
610    const SWAP_AGGREGATE_SELECTOR: [u8; 4] = [0xa2, 0xf7, 0x48, 0x93];
611    /// BebopRouter.swap
612    const ROUTER_SWAP_SELECTOR: [u8; 4] = [0x95, 0x86, 0xd0, 0xe8];
613
614    #[test]
615    fn test_fallback_router_labels_components() {
616        let direct = BebopClientBuilder::new(Chain::Ethereum, String::new())
617            .build()
618            .unwrap();
619        let via_router = BebopClientBuilder::new(Chain::Ethereum, String::new())
620            .with_fallback_router()
621            .build()
622            .unwrap();
623        let price_data = BebopPriceData::default();
624
625        let component = |client: &BebopClient| {
626            client
627                .create_component_with_state(String::from("bebop"), vec![], &price_data, 0.0)
628                .component
629                .protocol_system
630        };
631
632        assert_eq!(component(&direct), BebopClient::PROTOCOL_SYSTEM);
633        assert_eq!(component(&via_router), BebopClient::FALLBACK_PROTOCOL_SYSTEM);
634        assert_eq!(
635            BebopClient::FALLBACK_PROTOCOL_SYSTEM,
636            tycho_execution::encoding::evm::BEBOP_FALLBACK_PROTOCOL_SYSTEM
637        );
638    }
639
640    #[tokio::test]
641    #[ignore] // Requires network access and setting proper env vars
642    async fn test_bebop_websocket_connection() {
643        // We test with quote tokens that are not USDC in order to ensure our normalization works
644        // fine
645        let wbtc = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
646        let weth = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
647
648        dotenv().expect("Missing .env file");
649        let auth = get_bebop_auth().expect("Failed to get Bebop authentication");
650
651        let quote_tokens = HashSet::from([
652            // Use addresses we forgot to checksum (to test checksumming)
653            Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap(), // USDC
654            Bytes::from_str("0xdac17f958d2ee523a2206206994597c13d831ec7").unwrap(), // USDT
655        ]);
656
657        let client = BebopClient::new(
658            Chain::Ethereum,
659            HashSet::from_iter(vec![weth.clone(), wbtc.clone()]),
660            10.0, // $10 minimum TVL
661            auth.key,
662            quote_tokens,
663            Duration::from_secs(30),
664            None,
665            None,
666            None,
667        )
668        .unwrap();
669
670        let mut stream = client.stream();
671
672        // Test connection and message reception with timeout
673        // Receiving a single decodable pricing message is enough to prove the authenticated
674        // handshake and protobuf decoding work. Bebop only pushes on price changes, so
675        // requiring more messages makes the test flaky against market cadence.
676        let result = timeout(Duration::from_secs(10), async {
677            let mut message_count = 0;
678            let max_messages = 1;
679
680            while let Some(result) = stream.next().await {
681                match result {
682                    Ok((component_id, msg)) => {
683                        println!("Received message with ID: {component_id}");
684
685                        assert!(!component_id.is_empty());
686                        assert_eq!(component_id, "bebop");
687                        assert!(msg.header.timestamp > 0);
688                        assert!(!msg.snapshots.states.is_empty());
689
690                        let snapshot = &msg.snapshots;
691
692                        // We got at least one component
693                        assert!(!snapshot.states.is_empty());
694
695                        println!("Received {} components in this message", snapshot.states.len());
696                        for (id, component_with_state) in &snapshot.states {
697                            assert_eq!(
698                                component_with_state
699                                    .component
700                                    .protocol_system,
701                                "rfq:bebop"
702                            );
703                            assert_eq!(
704                                component_with_state
705                                    .component
706                                    .protocol_type_name,
707                                "bebop_pool"
708                            );
709                            assert_eq!(component_with_state.component.chain, Chain::Ethereum);
710
711                            let attributes = &component_with_state.state.attributes;
712
713                            // Check that bids and asks exist and have non-empty byte strings
714                            assert!(attributes.contains_key("bids"));
715                            assert!(attributes.contains_key("asks"));
716                            assert!(!attributes["bids"].is_empty());
717                            assert!(!attributes["asks"].is_empty());
718
719                            if let Some(tvl) = component_with_state.component_tvl {
720                                assert!(tvl >= 0.0);
721                                println!("Component {id} TVL: ${tvl:.2}");
722                            }
723                        }
724
725                        message_count += 1;
726                        if message_count >= max_messages {
727                            break;
728                        }
729                    }
730                    Err(e) => {
731                        panic!("Stream error: {e}");
732                    }
733                }
734            }
735
736            assert!(message_count > 0, "Should have received at least one message");
737            println!("Successfully received {message_count} messages");
738        })
739        .await;
740
741        match result {
742            Ok(_) => println!("Test completed successfully"),
743            Err(_) => panic!("Test timed out - no messages received within 10 seconds"),
744        }
745    }
746
747    #[tokio::test]
748    async fn test_websocket_reconnection() {
749        // Start a mock WebSocket server that will drop connections intermittently
750        let listener = TcpListener::bind("127.0.0.1:0")
751            .await
752            .unwrap();
753        let addr = listener.local_addr().unwrap();
754
755        // Creates a thread-safe counter.
756        let connection_count = Arc::new(Mutex::new(0u32));
757
758        // We must clone - since we want to read the original value at the end of the test.
759        let connection_count_clone = connection_count.clone();
760
761        tokio::spawn(async move {
762            while let Ok((stream, _)) = listener.accept().await {
763                *connection_count_clone.lock().unwrap() += 1;
764                let count = *connection_count_clone.lock().unwrap();
765                println!("Mock server: Connection #{count} established");
766
767                tokio::spawn(async move {
768                    if let Ok(ws_stream) = accept_async(stream).await {
769                        let (mut ws_sender, _ws_receiver) = ws_stream.split();
770
771                        // Create test protobuf message
772                        let weth_addr =
773                            hex::decode("C02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
774                        let usdc_addr =
775                            hex::decode("A0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap();
776
777                        let test_price_data = BebopPriceData {
778                            base: weth_addr,
779                            quote: usdc_addr,
780                            last_update_ts: 1752617378,
781                            bids: vec![3070.05f32, 0.325717f32],
782                            asks: vec![3070.527f32, 0.325717f32],
783                        };
784
785                        let pricing_update = BebopPricingUpdate { pairs: vec![test_price_data] };
786
787                        let test_message = pricing_update.encode_to_vec();
788
789                        if count == 1 {
790                            // First connection: Send message successfully, then drop
791                            println!("Mock server: Connection #1 - sending message then dropping.");
792                            let _ = ws_sender
793                                .send(Message::Binary(test_message.clone().into()))
794                                .await;
795
796                            // Give time for message to be processed, then drop the connection.
797                            tokio::time::sleep(Duration::from_millis(100)).await;
798                            println!("Mock server: Dropping connection #1");
799                            let _ = ws_sender.close().await;
800                        } else if count == 2 {
801                            // Second connection: Send message successfully and maintain connection
802                            println!("Mock server: Connection #2 - maintaining stable connection.");
803                            let _ = ws_sender
804                                .send(Message::Binary(test_message.clone().into()))
805                                .await;
806                        }
807                    }
808                });
809            }
810        });
811
812        // Wait a moment for the server to start
813        tokio::time::sleep(Duration::from_millis(50)).await;
814
815        let mut test_quote_tokens = HashSet::new();
816        test_quote_tokens
817            .insert(Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap());
818
819        let tokens_formatted = vec![
820            Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap(),
821            Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap(),
822        ];
823
824        // Bypass the new() constructor to mock the URL to point to our mock server.
825        let client = BebopClient {
826            chain: Chain::Ethereum,
827            price_ws: format!("ws://127.0.0.1:{}", addr.port()),
828            tokens: tokens_formatted.into_iter().collect(),
829            tvl: 1000.0,
830            ws_key: "test_key".to_string(),
831            quote_tokens: test_quote_tokens,
832            quote_endpoint: "".to_string(),
833            quote_timeout: Duration::from_secs(5),
834            origin_address: None,
835            origin_target: None,
836            origin_source: None,
837            protocol_system: BebopClient::PROTOCOL_SYSTEM.to_string(),
838        };
839
840        let start_time = std::time::Instant::now();
841        let mut successful_messages = 0;
842        let mut connection_errors = 0;
843        let mut first_message_received = false;
844        let mut second_message_received = false;
845
846        // Expected flow:
847        // 1. Receive first message successfully
848        // 2. Connection drops
849        // 3. Client reconnects
850        // 4. Receive second message successfully
851        // Timeout if two messages are not received within 5 seconds.
852        while start_time.elapsed() < Duration::from_secs(5) && successful_messages < 2 {
853            match timeout(Duration::from_millis(1000), client.stream().next()).await {
854                Ok(Some(result)) => match result {
855                    Ok((_component_id, _message)) => {
856                        successful_messages += 1;
857                        println!("Received successful message {successful_messages}");
858
859                        if successful_messages == 1 {
860                            first_message_received = true;
861                            println!("First message received - connection should drop after this.");
862                        } else if successful_messages == 2 {
863                            second_message_received = true;
864                            println!("Second message received after reconnection.");
865                        }
866                    }
867                    Err(e) => {
868                        connection_errors += 1;
869                        println!("Connection error during reconnection: {e:?}");
870                    }
871                },
872                Ok(None) => {
873                    panic!("Stream ended unexpectedly");
874                }
875                Err(_) => {
876                    println!("Timeout waiting for message (normal during reconnections)");
877                    continue;
878                }
879            }
880        }
881
882        let final_connection_count = *connection_count.lock().unwrap();
883
884        // 1. Exactly 2 connection attempts (initial + reconnect)
885        // 2. Exactly 2 successful messages (one before drop, one after reconnect)
886
887        assert_eq!(final_connection_count, 2);
888        assert!(first_message_received);
889        assert!(second_message_received);
890        assert_eq!(connection_errors, 0);
891        assert_eq!(successful_messages, 2);
892    }
893
894    #[tokio::test]
895    #[ignore] // Requires network access and setting proper env vars
896    async fn test_bebop_quote_single_order() {
897        let token_in = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
898        let token_out = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
899        dotenv().expect("Missing .env file");
900        let auth = get_bebop_auth().expect("Failed to get Bebop authentication");
901
902        let router = Bytes::from_str("0xfD0b31d2E955fA55e3fa641Fe90e08b677188d35").unwrap();
903
904        let client = BebopClient::new(
905            Chain::Ethereum,
906            HashSet::from_iter(vec![token_in.clone(), token_out.clone()]),
907            10.0, // $10 minimum TVL
908            auth.key,
909            HashSet::new(),
910            Duration::from_secs(30),
911            Some(Bytes::from_str("0x00000000219ab540356cBB839Cbe05303d7705Fa").unwrap()),
912            Some(router.clone()),
913            Some("tycho-test".to_string()),
914        )
915        .unwrap();
916
917        let params = GetAmountOutParams {
918            amount_in: BigUint::from(1_000000000000000000u64),
919            token_in: token_in.clone(),
920            token_out: token_out.clone(),
921            sender: router.clone(),
922            receiver: router,
923        };
924        let quote = client
925            .request_binding_quote(&params)
926            .await
927            .unwrap();
928
929        assert_eq!(quote.base_token, token_in);
930        assert_eq!(quote.quote_token, token_out);
931        assert_eq!(quote.amount_in, BigUint::from(1_000000000000000000u64));
932
933        // Conservative sanity bound (0.01 WBTC for 1 WETH) — proves a real, non-dust quote came
934        // back without depending closely on the live WETH/WBTC price.
935        assert!(quote.amount_out > BigUint::from(1_000_000u64));
936
937        // The settlement mode depends on the API account configuration behind BEBOP_KEY:
938        // settlement-mode accounts get BebopSettlement.swapSingle calldata, router-mode
939        // accounts get BebopRouter.swap calldata.
940        let selector = &quote
941            .quote_attributes
942            .get("calldata")
943            .unwrap()[..4];
944        if selector == SWAP_SINGLE_SELECTOR {
945            let partial_fill_offset_slice = quote
946                .quote_attributes
947                .get("partial_fill_offset")
948                .unwrap()
949                .as_ref();
950            let mut partial_fill_offset_array = [0u8; 8];
951            partial_fill_offset_array.copy_from_slice(partial_fill_offset_slice);
952
953            assert_eq!(u64::from_be_bytes(partial_fill_offset_array), 12);
954        } else {
955            assert_eq!(selector, ROUTER_SWAP_SELECTOR);
956        }
957    }
958
959    #[tokio::test]
960    #[ignore] // Requires network access and setting proper env vars
961    async fn test_bebop_quote_aggregate_order() {
962        // This will make a quote request similar to the previous test but with a very big amount
963        // We expect the Bebop Quote to have an aggregate order (split between different mms)
964        let token_in = Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap();
965        let token_out = Bytes::from_str("0xfAbA6f8e4a5E8Ab82F62fe7C39859FA577269BE3").unwrap();
966        dotenv().expect("Missing .env file");
967        let auth = get_bebop_auth().expect("Failed to get Bebop authentication");
968
969        let router = Bytes::from_str("0xfD0b31d2E955fA55e3fa641Fe90e08b677188d35").unwrap();
970
971        let client = BebopClient::new(
972            Chain::Ethereum,
973            HashSet::from_iter(vec![token_in.clone(), token_out.clone()]),
974            10.0, // $10 minimum TVL
975            auth.key,
976            HashSet::new(),
977            Duration::from_secs(30),
978            Some(Bytes::from_str("0x00000000219ab540356cBB839Cbe05303d7705Fa").unwrap()),
979            Some(router.clone()),
980            Some("tycho-test".to_string()),
981        )
982        .unwrap();
983
984        let amount_in = BigUint::from_str("20_000_000_000").unwrap(); // 20k USDC
985        let params = GetAmountOutParams {
986            amount_in: amount_in.clone(),
987            token_in: token_in.clone(),
988            token_out: token_out.clone(),
989            sender: router.clone(),
990            receiver: router,
991        };
992        let quote = client
993            .request_binding_quote(&params)
994            .await
995            .unwrap();
996
997        assert_eq!(quote.base_token, token_in);
998        assert_eq!(quote.quote_token, token_out);
999        assert_eq!(quote.amount_in, amount_in);
1000
1001        // Assuming the USDC - ONDO price doesn't change too much at the time of running this
1002        assert!(quote.amount_out > BigUint::from_str("18000000000000000000000").unwrap()); // ~19k ONDO
1003
1004        // The settlement mode depends on the API account configuration behind BEBOP_KEY:
1005        // settlement-mode accounts get BebopSettlement.swapAggregate calldata, router-mode
1006        // accounts get BebopRouter.swap calldata.
1007        let selector = &quote
1008            .quote_attributes
1009            .get("calldata")
1010            .unwrap()[..4];
1011        if selector == SWAP_AGGREGATE_SELECTOR {
1012            let partial_fill_offset_slice = quote
1013                .quote_attributes
1014                .get("partial_fill_offset")
1015                .unwrap()
1016                .as_ref();
1017            let mut partial_fill_offset_array = [0u8; 8];
1018            partial_fill_offset_array.copy_from_slice(partial_fill_offset_slice);
1019
1020            // This is the only attribute that is significantly different for the Single and
1021            // Aggregate Order
1022            assert_eq!(u64::from_be_bytes(partial_fill_offset_array), 2);
1023        } else {
1024            assert_eq!(selector, ROUTER_SWAP_SELECTOR);
1025        }
1026    }
1027
1028    #[test]
1029    fn test_process_bebop_quote_response_aggregate_order() {
1030        let json =
1031            std::fs::read_to_string("src/rfq/protocols/bebop/test_responses/aggregate_order.json")
1032                .unwrap();
1033        let quote_response: BebopQuoteResponse = serde_json::from_str(&json).unwrap();
1034        let params = GetAmountOutParams {
1035            amount_in: BigUint::from_str("20000000000").unwrap(),
1036            token_in: Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap(),
1037            token_out: Bytes::from_str("0xfAbA6f8e4a5E8Ab82F62fe7C39859FA577269BE3").unwrap(),
1038            sender: Bytes::from_str("0xfd0b31d2e955fa55e3fa641fe90e08b677188d35").unwrap(),
1039            receiver: Bytes::from_str("0xfd0b31d2e955fa55e3fa641fe90e08b677188d35").unwrap(),
1040        };
1041        let res = BebopClient::process_quote_response(quote_response, &params).unwrap();
1042        assert_eq!(res.amount_out, BigUint::from_str("52571055094221715780641").unwrap());
1043        assert_eq!(res.amount_in, BigUint::from_str("20000000000").unwrap());
1044        assert_eq!(res.base_token, params.token_in);
1045        assert_eq!(res.quote_token, params.token_out);
1046    }
1047
1048    #[test]
1049    fn test_process_bebop_quote_response_aggregate_order_with_multihop() {
1050        let json = std::fs::read_to_string(
1051            "src/rfq/protocols/bebop/test_responses/aggregate_order_with_multihop.json",
1052        )
1053        .unwrap();
1054        let quote_response: BebopQuoteResponse = serde_json::from_str(&json).unwrap();
1055        let params = GetAmountOutParams {
1056            amount_in: BigUint::from_str("43067495979235520920162").unwrap(),
1057            token_in: Bytes::from_str("0xDEf1CA1fb7FBcDC777520aa7f396b4E015F497aB").unwrap(),
1058            token_out: Bytes::from_str("0xdAC17F958D2ee523a2206206994597C13D831ec7").unwrap(),
1059            sender: Bytes::from_str("0x809305d724B6E79C71e10a097ABadd1274B9C279").unwrap(),
1060            receiver: Bytes::from_str("0x809305d724B6E79C71e10a097ABadd1274B9C279").unwrap(),
1061        };
1062        let res = BebopClient::process_quote_response(quote_response, &params).unwrap();
1063        assert_eq!(res.amount_out, BigUint::from_str("11186653890").unwrap());
1064        assert_eq!(res.amount_in, BigUint::from_str("43067495979235520920162").unwrap());
1065        assert_eq!(res.base_token, params.token_in);
1066        assert_eq!(res.quote_token, params.token_out);
1067    }
1068
1069    #[test]
1070    fn test_process_bebop_quote_response_single_order() {
1071        // Captured from a settlement-mode API account: the signed order's taker and receiver
1072        // are the requested sender/receiver and the calldata targets the settlement contract.
1073        let json =
1074            std::fs::read_to_string("src/rfq/protocols/bebop/test_responses/single_order.json")
1075                .unwrap();
1076        let quote_response: BebopQuoteResponse = serde_json::from_str(&json).unwrap();
1077        let router = Bytes::from_str("0xfD0b31d2E955fA55e3fa641Fe90e08b677188d35").unwrap();
1078        let params = GetAmountOutParams {
1079            amount_in: BigUint::from_str("1000000000000000000").unwrap(),
1080            token_in: Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap(),
1081            token_out: Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap(),
1082            sender: router.clone(),
1083            receiver: router,
1084        };
1085        let res = BebopClient::process_quote_response(quote_response, &params).unwrap();
1086        assert_eq!(res.amount_in, BigUint::from_str("1000000000000000000").unwrap());
1087        assert_eq!(res.amount_out, BigUint::from_str("2915408").unwrap());
1088        let settlement = Bytes::from_str("0xbbbbbBB520d69a9775E85b458C58c648259FAD5F").unwrap();
1089        assert_eq!(
1090            res.quote_attributes
1091                .get("tx_to")
1092                .unwrap(),
1093            &settlement
1094        );
1095        assert_eq!(
1096            res.quote_attributes
1097                .get("calldata")
1098                .unwrap()[..4],
1099            SWAP_SINGLE_SELECTOR
1100        );
1101    }
1102
1103    #[test]
1104    fn test_process_bebop_quote_response_single_order_router_mode() {
1105        // Captured from an API account configured for router-mode settlement: the signed
1106        // order's taker and receiver are the Bebop router contract (= tx.to), not the
1107        // requested sender/receiver.
1108        let json = std::fs::read_to_string(
1109            "src/rfq/protocols/bebop/test_responses/single_order_router_mode.json",
1110        )
1111        .unwrap();
1112        let quote_response: BebopQuoteResponse = serde_json::from_str(&json).unwrap();
1113        let router = Bytes::from_str("0xfD0b31d2E955fA55e3fa641Fe90e08b677188d35").unwrap();
1114        let params = GetAmountOutParams {
1115            amount_in: BigUint::from_str("1000000000000000000").unwrap(),
1116            token_in: Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap(),
1117            token_out: Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap(),
1118            sender: router.clone(),
1119            receiver: router,
1120        };
1121        let res = BebopClient::process_quote_response(quote_response, &params).unwrap();
1122        assert_eq!(res.amount_in, BigUint::from_str("1000000000000000000").unwrap());
1123        assert_eq!(res.amount_out, BigUint::from_str("2926296").unwrap());
1124        let bebop_router = Bytes::from_str("0xBeb0009ACa35087ce7cCF11637E24dd1Aad3bf2A").unwrap();
1125        assert_eq!(
1126            res.quote_attributes
1127                .get("tx_to")
1128                .unwrap(),
1129            &bebop_router
1130        );
1131        assert_eq!(
1132            res.quote_attributes
1133                .get("calldata")
1134                .unwrap()[..4],
1135            ROUTER_SWAP_SELECTOR
1136        );
1137    }
1138
1139    /// Helper function to create a mock server that responds after a delay
1140    async fn create_delayed_response_server(delay_ms: u64) -> std::net::SocketAddr {
1141        use tokio::io::AsyncWriteExt;
1142
1143        let listener = TcpListener::bind("127.0.0.1:0")
1144            .await
1145            .unwrap();
1146        let addr = listener.local_addr().unwrap();
1147
1148        let json_response =
1149            std::fs::read_to_string("src/rfq/protocols/bebop/test_responses/aggregate_order.json")
1150                .unwrap();
1151
1152        tokio::spawn(async move {
1153            while let Ok((mut stream, _)) = listener.accept().await {
1154                let json_response_clone = json_response.clone();
1155                tokio::spawn(async move {
1156                    sleep(Duration::from_millis(delay_ms)).await;
1157
1158                    let response = format!(
1159                        "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
1160                        json_response_clone.len(),
1161                        json_response_clone
1162                    );
1163                    let _ = stream
1164                        .write_all(response.as_bytes())
1165                        .await;
1166                    let _ = stream.flush().await;
1167                    let _ = stream.shutdown().await;
1168                });
1169            }
1170        });
1171
1172        addr
1173    }
1174
1175    fn create_test_bebop_client(quote_endpoint: String, quote_timeout: Duration) -> BebopClient {
1176        let token_in = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1177        let token_out = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
1178
1179        BebopClient {
1180            chain: Chain::Ethereum,
1181            price_ws: "ws://example.com".to_string(),
1182            quote_endpoint,
1183            tokens: HashSet::from([token_in, token_out]),
1184            tvl: 10.0,
1185            ws_key: "test_key".to_string(),
1186            quote_tokens: HashSet::new(),
1187            quote_timeout,
1188            origin_address: None,
1189            origin_target: None,
1190            origin_source: None,
1191            protocol_system: BebopClient::PROTOCOL_SYSTEM.to_string(),
1192        }
1193    }
1194
1195    /// Helper function to create test quote params matching aggregate_order.json
1196    fn create_test_quote_params() -> GetAmountOutParams {
1197        let token_in = Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap();
1198        let token_out = Bytes::from_str("0xfAbA6f8e4a5E8Ab82F62fe7C39859FA577269BE3").unwrap();
1199        let router = Bytes::from_str("0xfD0b31d2E955fA55e3fa641Fe90e08b677188d35").unwrap();
1200
1201        GetAmountOutParams {
1202            amount_in: BigUint::from_str("20000000000").unwrap(),
1203            token_in,
1204            token_out,
1205            sender: router.clone(),
1206            receiver: router,
1207        }
1208    }
1209
1210    #[tokio::test]
1211    async fn test_bebop_quote_timeout() {
1212        let addr = create_delayed_response_server(500).await;
1213
1214        // Test 1: Client with short timeout (200ms) - should timeout
1215        let client_short_timeout = create_test_bebop_client(
1216            format!("http://127.0.0.1:{}/quote", addr.port()),
1217            Duration::from_millis(200),
1218        );
1219        let params = create_test_quote_params();
1220
1221        let start = std::time::Instant::now();
1222        let result = client_short_timeout
1223            .request_binding_quote(&params)
1224            .await;
1225        let elapsed = start.elapsed();
1226
1227        assert!(result.is_err());
1228        let err = result.unwrap_err();
1229        match err {
1230            RFQError::ConnectionError(msg) => {
1231                assert!(msg.contains("timed out"), "Expected timeout error, got: {}", msg);
1232            }
1233            _ => panic!("Expected ConnectionError, got: {:?}", err),
1234        }
1235        assert!(
1236            elapsed.as_millis() >= 200 && elapsed.as_millis() < 400,
1237            "Expected timeout around 200ms, got: {:?}",
1238            elapsed
1239        );
1240
1241        // Test 2: Client with long timeout (1 seconds) - should wait and receive response
1242        // Note: With retry logic, we may need multiple attempts if the response is malformed,
1243        // so we need a longer timeout to account for retries
1244        let client_long_timeout = create_test_bebop_client(
1245            format!("http://127.0.0.1:{}/quote", addr.port()),
1246            Duration::from_secs(1),
1247        );
1248
1249        let result = client_long_timeout
1250            .request_binding_quote(&params)
1251            .await;
1252
1253        // Should succeed - the server waits 500ms which is within the 1s timeout
1254        assert!(result.is_ok(), "Expected success, got: {:?}", result);
1255        let quote = result.unwrap();
1256
1257        // Verify the quote matches what we expect from aggregate_order.json
1258        assert_eq!(quote.base_token, params.token_in);
1259        assert_eq!(quote.quote_token, params.token_out);
1260    }
1261
1262    /// Helper function to create a mock server that fails twice, then succeeds with
1263    /// aggregate_order.json
1264    async fn create_retry_server() -> (std::net::SocketAddr, Arc<Mutex<u32>>) {
1265        use std::sync::{Arc, Mutex};
1266
1267        use tokio::io::AsyncWriteExt;
1268
1269        let request_count = Arc::new(Mutex::new(0u32));
1270        let request_count_clone = request_count.clone();
1271
1272        let listener = TcpListener::bind("127.0.0.1:0")
1273            .await
1274            .unwrap();
1275        let addr = listener.local_addr().unwrap();
1276
1277        let json_response =
1278            std::fs::read_to_string("src/rfq/protocols/bebop/test_responses/aggregate_order.json")
1279                .unwrap();
1280
1281        tokio::spawn(async move {
1282            while let Ok((mut stream, _)) = listener.accept().await {
1283                let count_clone = request_count_clone.clone();
1284                let json_response_clone = json_response.clone();
1285                tokio::spawn(async move {
1286                    *count_clone.lock().unwrap() += 1;
1287                    let count = *count_clone.lock().unwrap();
1288                    println!("Mock server: Received request #{count}");
1289
1290                    if count <= 2 {
1291                        let response = "HTTP/1.1 500 Internal Server Error\r\nContent-Length: 21\r\n\r\nInternal Server Error";
1292                        let _ = stream
1293                            .write_all(response.as_bytes())
1294                            .await;
1295                    } else {
1296                        let response = format!(
1297                            "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
1298                            json_response_clone.len(),
1299                            json_response_clone
1300                        );
1301                        let _ = stream
1302                            .write_all(response.as_bytes())
1303                            .await;
1304                    }
1305                    let _ = stream.flush().await;
1306                    let _ = stream.shutdown().await;
1307                });
1308            }
1309        });
1310        (addr, request_count)
1311    }
1312
1313    #[tokio::test]
1314    async fn test_bebop_quote_retry_on_bad_response() {
1315        let (addr, request_count) = create_retry_server().await;
1316
1317        let client = create_test_bebop_client(
1318            format!("http://127.0.0.1:{}/quote", addr.port()),
1319            Duration::from_secs(5),
1320        );
1321        let params = create_test_quote_params();
1322        let result = client
1323            .request_binding_quote(&params)
1324            .await;
1325
1326        assert!(result.is_ok(), "Expected success after retries, got: {:?}", result);
1327        let quote = result.unwrap();
1328
1329        // Verify the quote (amounts from aggregate_order.json)
1330        assert_eq!(quote.amount_in, BigUint::from_str("20000000000").unwrap());
1331        assert_eq!(quote.amount_out, BigUint::from_str("52571055094221715780641").unwrap());
1332
1333        // Verify exactly 3 requests were made (2 failures + 1 success)
1334        let final_count = *request_count.lock().unwrap();
1335        assert_eq!(final_count, 3, "Expected 3 requests, got {}", final_count);
1336    }
1337
1338    #[test]
1339    fn test_bebop_client_serialize_deserialize_roundtrip() {
1340        let token_in = Bytes::from_str("0xC02aaA39b223FE8D0A0e5C4F27eAD9083C756Cc2").unwrap();
1341        let token_out = Bytes::from_str("0x2260FAC5E5542a773Aa44fBCfeDf7C193bc2C599").unwrap();
1342        let quote_token = Bytes::from_str("0xa0b86991c6218b36c1d19d4a2e9eb0ce3606eb48").unwrap();
1343
1344        let original = BebopClient {
1345            chain: Chain::Ethereum,
1346            price_ws: "wss://api.bebop.xyz/pricing".to_string(),
1347            quote_endpoint: "https://api.bebop.xyz/quote".to_string(),
1348            tokens: HashSet::from([token_in.clone(), token_out.clone()]),
1349            tvl: 50.5,
1350            ws_key: "secret_key".to_string(),
1351            quote_tokens: HashSet::from([quote_token.clone()]),
1352            quote_timeout: Duration::from_millis(5500),
1353            origin_address: Some(
1354                Bytes::from_str("0x00000000219ab540356cBB839Cbe05303d7705Fa").unwrap(),
1355            ),
1356            origin_target: Some(
1357                Bytes::from_str("0xdA892C989d07A18B5DD3F392d949f00dF15C5736").unwrap(),
1358            ),
1359            origin_source: Some("tycho".to_string()),
1360            protocol_system: BebopClient::PROTOCOL_SYSTEM.to_string(),
1361        };
1362
1363        let serialized = serde_json::to_string(&original).unwrap();
1364        let deserialized: BebopClient = serde_json::from_str(&serialized).unwrap();
1365
1366        // Fields that should round-trip correctly
1367        assert_eq!(deserialized.chain, original.chain);
1368        assert_eq!(deserialized.price_ws, original.price_ws);
1369        assert_eq!(deserialized.quote_endpoint, original.quote_endpoint);
1370        assert_eq!(deserialized.tokens, original.tokens);
1371        assert_eq!(deserialized.tvl, original.tvl);
1372        assert_eq!(deserialized.quote_tokens, original.quote_tokens);
1373        assert_eq!(deserialized.quote_timeout, original.quote_timeout);
1374        assert_eq!(deserialized.origin_address, original.origin_address);
1375        assert_eq!(deserialized.origin_target, original.origin_target);
1376        assert_eq!(deserialized.origin_source, original.origin_source);
1377
1378        // ws_key should NOT round-trip (skip_serializing + default)
1379        assert_eq!(deserialized.ws_key, "");
1380        assert_ne!(deserialized.ws_key, original.ws_key);
1381    }
1382
1383    #[test]
1384    fn test_bebop_client_deserialize_with_credentials() {
1385        // When ws_key is provided in JSON, it should be deserialized
1386        // (skip_serializing only affects serialization, not deserialization)
1387        let json = r#"{
1388            "chain": "ethereum",
1389            "price_ws": "wss://api.bebop.xyz/pricing",
1390            "quote_endpoint": "https://api.bebop.xyz/quote",
1391            "tokens": [],
1392            "tvl": 10.0,
1393            "ws_key": "provided_key",
1394            "quote_tokens": [],
1395            "quote_timeout": {"secs": 30, "nanos": 0}
1396        }"#;
1397
1398        let client: BebopClient = serde_json::from_str(json).unwrap();
1399
1400        // Credentials should be deserialized from JSON
1401        assert_eq!(client.ws_key, "provided_key");
1402    }
1403
1404    #[test]
1405    fn test_process_bebop_quote_response_aggregate_order_router_mode() {
1406        // Captured from a router-mode API account: an aggregate order split across three
1407        // makers where the signed order's taker and receiver are the Bebop router (= tx.to).
1408        let json = std::fs::read_to_string(
1409            "src/rfq/protocols/bebop/test_responses/aggregate_order_router_mode.json",
1410        )
1411        .unwrap();
1412        let quote_response: BebopQuoteResponse = serde_json::from_str(&json).unwrap();
1413        let router = Bytes::from_str("0xfD0b31d2E955fA55e3fa641Fe90e08b677188d35").unwrap();
1414        let params = GetAmountOutParams {
1415            amount_in: BigUint::from_str("20000000000").unwrap(),
1416            token_in: Bytes::from_str("0xA0b86991c6218b36c1d19D4a2e9Eb0cE3606eB48").unwrap(),
1417            token_out: Bytes::from_str("0xfAbA6f8e4a5E8Ab82F62fe7C39859FA577269BE3").unwrap(),
1418            sender: router.clone(),
1419            receiver: router,
1420        };
1421        let res = BebopClient::process_quote_response(quote_response, &params).unwrap();
1422        assert_eq!(res.amount_in, BigUint::from_str("20000000000").unwrap());
1423        assert_eq!(res.amount_out, BigUint::from_str("52577858553072299423490").unwrap());
1424        let bebop_router = Bytes::from_str("0xBeb0009ACa35087ce7cCF11637E24dd1Aad3bf2A").unwrap();
1425        assert_eq!(
1426            res.quote_attributes
1427                .get("tx_to")
1428                .unwrap(),
1429            &bebop_router
1430        );
1431        assert_eq!(
1432            res.quote_attributes
1433                .get("calldata")
1434                .unwrap()[..4],
1435            ROUTER_SWAP_SELECTOR
1436        );
1437    }
1438}