dex_connector/
hyperliquid_connector.rs

1use crate::{
2    dex_connector::{slippage_price, string_to_decimal, DexConnector},
3    dex_request::{DexError, DexRequest, HttpMethod},
4    dex_websocket::DexWebSocket,
5    BalanceResponse, CanceledOrder, CanceledOrdersResponse, CombinedBalanceResponse,
6    CreateOrderResponse, FilledOrder, FilledOrdersResponse, LastTradesResponse, OpenOrdersResponse,
7    OrderSide, TickerResponse, TpSl,
8};
9use async_trait::async_trait;
10use chrono::{DateTime, Duration as ChronoDuration, Utc};
11use debot_utils::parse_to_decimal;
12use ethers::{signers::LocalWallet, types::H160};
13use futures::{
14    stream::{SplitSink, SplitStream},
15    SinkExt, StreamExt,
16};
17use hyperliquid_rust_sdk_fork::{
18    BaseUrl, ClientCancelRequest, ClientLimit, ClientOrder, ClientOrderRequest, ClientTrigger,
19    ExchangeClient, ExchangeDataStatus, ExchangeResponseStatus,
20};
21use reqwest::Client;
22use rust_decimal::prelude::ToPrimitive;
23use rust_decimal::Decimal;
24use serde::{Deserialize, Serialize};
25use serde_json::Value;
26use std::{
27    collections::HashMap,
28    str::FromStr,
29    sync::{
30        atomic::{AtomicBool, Ordering},
31        Arc,
32    },
33    time::Duration,
34};
35use tokio::{
36    net::TcpStream,
37    select,
38    signal::unix::{signal, SignalKind},
39    sync::{Mutex, RwLock},
40    task::JoinHandle,
41    time::sleep,
42};
43use tokio_tungstenite::{tungstenite::protocol::Message, MaybeTlsStream, WebSocketStream};
44
45struct Config {
46    evm_wallet_address: String,
47    symbol_list: Vec<String>,
48}
49
50// --- Spot metadata support ---
51#[derive(Deserialize, Debug)]
52struct SpotMetaToken {
53    #[serde(rename = "name")]
54    _name: String,
55    #[serde(rename = "szDecimals")]
56    _sz_decimals: u32,
57    #[serde(rename = "weiDecimals")]
58    _wei_decimals: u32,
59    #[serde(rename = "index")]
60    _index: usize,
61}
62
63#[derive(Deserialize, Debug, Clone)]
64struct SpotMetaUniverse {
65    #[serde(rename = "name")]
66    name: String,
67    #[serde(rename = "tokens")]
68    _tokens: Vec<usize>,
69    #[serde(rename = "index")]
70    index: usize,
71}
72
73#[derive(Deserialize, Debug)]
74struct SpotMetaResponse {
75    #[serde(rename = "tokens")]
76    _tokens: Vec<SpotMetaToken>,
77    #[serde(rename = "universe")]
78    universe: Vec<SpotMetaUniverse>,
79}
80
81#[derive(Serialize, Debug)]
82struct InfoRequest<'a> {
83    #[serde(rename = "type")]
84    req_type: &'a str,
85    #[serde(skip_serializing_if = "Option::is_none")]
86    user: Option<&'a str>,
87}
88
89#[derive(Debug)]
90struct TradeResult {
91    pub filled_side: OrderSide,
92    pub filled_size: Decimal,
93    pub filled_value: Decimal,
94    pub filled_fee: Decimal,
95    order_id: String,
96    pub is_rejected: bool,
97}
98
99#[derive(Debug, Clone)]
100pub struct CancelEvent {
101    pub order_id: String,
102    pub timestamp: u64,
103}
104
105#[derive(Default)]
106struct DynamicMarketInfo {
107    pub best_bid: Option<Decimal>,
108    pub best_ask: Option<Decimal>,
109    pub market_price: Option<Decimal>,
110    pub min_tick: Option<Decimal>,
111    pub volume: Option<Decimal>,
112    pub num_trades: Option<u64>,
113    pub open_interest: Option<Decimal>,
114    pub funding_rate: Option<Decimal>,
115    pub oracle_price: Option<Decimal>,
116}
117
118#[derive(Clone)]
119struct StaticMarketInfo {
120    pub decimals: u32,
121    pub _max_leverage: u32,
122}
123
124#[derive(Clone)]
125#[allow(dead_code)]
126struct MaintenanceInfo {
127    next_start: Option<DateTime<Utc>>,
128    fetched_at: DateTime<Utc>,
129}
130
131#[derive(Deserialize, Debug)]
132pub struct OrderUpdateDetail {
133    pub coin: String,
134    #[serde(rename = "oid")]
135    pub oid: u64,
136}
137
138#[derive(Deserialize, Debug)]
139pub struct OrderUpdate {
140    pub order: OrderUpdateDetail,
141    pub status: String,
142    #[serde(rename = "statusTimestamp")]
143    pub status_timestamp: u64,
144}
145
146#[allow(dead_code)]
147#[derive(Deserialize, Debug)]
148struct WsLevel {
149    px: String,
150    sz: String,
151    n: u64,
152}
153
154#[allow(dead_code)]
155#[derive(Deserialize, Debug)]
156struct WsBbo {
157    coin: String,
158    time: u64,
159    bbo: [Option<WsLevel>; 2], // [bestBid?, bestAsk?]
160}
161
162#[allow(dead_code)]
163#[derive(Deserialize, Debug)]
164struct WsBook {
165    coin: String,
166    time: u64,
167    levels: [Vec<WsLevel>; 2], // [bids, asks]
168}
169
170pub struct HyperliquidConnector {
171    config: Config,
172    request: DexRequest,
173    web_socket: DexWebSocket,
174    running: Arc<AtomicBool>,
175    read_socket: Arc<Mutex<Option<SplitStream<WebSocketStream<MaybeTlsStream<TcpStream>>>>>>,
176    write_socket:
177        Arc<Mutex<Option<SplitSink<WebSocketStream<MaybeTlsStream<TcpStream>>, Message>>>>,
178    task_handle_read_message: Arc<Mutex<Option<JoinHandle<()>>>>,
179    task_handle_read_sigterm: Arc<Mutex<Option<JoinHandle<()>>>>,
180    // 1st key = symbol, 2nd key = order_id
181    trade_results: Arc<RwLock<HashMap<String, HashMap<String, TradeResult>>>>,
182    canceled_results: Arc<RwLock<HashMap<String, HashMap<String, CancelEvent>>>>,
183    dynamic_market_info: Arc<RwLock<HashMap<String, DynamicMarketInfo>>>,
184    static_market_info: HashMap<String, StaticMarketInfo>,
185    spot_index_map: HashMap<String, usize>,
186    spot_reverse_map: Arc<HashMap<usize, String>>,
187    exchange_client: ExchangeClient,
188    maintenance: Arc<RwLock<MaintenanceInfo>>,
189    last_volumes: Arc<Mutex<HashMap<String, Decimal>>>,
190}
191
192#[derive(Debug)]
193struct WebSocketMessage {
194    _channel: String,
195    data: WebSocketData,
196}
197
198#[derive(Debug)]
199enum WebSocketData {
200    AllMidsData(AllMidsData),
201    UserFillsData(UserFillsData),
202    CandleData(CandleData),
203    ActiveAssetCtxData(ActiveAssetCtxData),
204    OrderUpdatesData(Vec<OrderUpdate>),
205    Bbo(WsBbo),
206    L2Book(WsBook),
207}
208
209#[derive(Deserialize, Debug)]
210struct AllMidsData {
211    mids: HashMap<String, String>,
212}
213
214#[allow(dead_code, non_snake_case)]
215#[derive(Deserialize, Debug)]
216struct CandleData {
217    t: u64,     // Open time (milliseconds)
218    T: u64,     // Close time (milliseconds)
219    s: String,  // Symbol
220    i: String,  // Interval
221    o: Decimal, // Open price
222    c: Decimal, // Close price
223    h: Decimal, // High price
224    l: Decimal, // Low price
225    v: Decimal, // Volume
226    n: u64,     // Number of trades
227}
228
229#[derive(Deserialize, Debug)]
230pub struct ActiveAssetCtxData {
231    pub coin: String,       // The asset symbol (e.g., BTC-USD)
232    pub ctx: PerpsAssetCtx, // The asset context containing market details
233}
234
235#[allow(dead_code, non_snake_case)]
236#[derive(Deserialize, Debug)]
237pub struct PerpsAssetCtx {
238    pub dayNtlVlm: Decimal,     // Daily notional volume
239    pub prevDayPx: Decimal,     // Previous day's price
240    pub markPx: Decimal,        // Mark price
241    pub midPx: Option<Decimal>, // Mid price (optional)
242    pub funding: Decimal,       // Funding rate
243    pub openInterest: Decimal,  // Open interest
244    pub oraclePx: Decimal,      // Oracle price
245}
246
247#[derive(Serialize, Deserialize, Debug)]
248pub struct UserFillsData {
249    pub user: String,
250    pub fills: Vec<Fill>,
251}
252
253#[derive(Serialize, Deserialize, Debug)]
254pub struct Fill {
255    pub coin: String,
256    pub px: Decimal,
257    pub sz: Decimal,
258    pub side: String,
259    pub dir: String,
260    #[serde(rename = "closedPnl")]
261    pub closed_pnl: Decimal,
262    pub oid: u64,
263    pub tid: u64,
264    pub fee: Decimal,
265}
266
267impl<'de> Deserialize<'de> for WebSocketMessage {
268    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
269    where
270        D: serde::Deserializer<'de>,
271    {
272        #[derive(Deserialize)]
273        struct Helper {
274            channel: String,
275            data: serde_json::Value,
276        }
277        let helper = Helper::deserialize(deserializer)?;
278        let data = match helper.channel.as_str() {
279            "allMids" => AllMidsData::deserialize(helper.data)
280                .map(WebSocketData::AllMidsData)
281                .map_err(serde::de::Error::custom)?,
282            "userFills" => UserFillsData::deserialize(helper.data)
283                .map(WebSocketData::UserFillsData)
284                .map_err(serde::de::Error::custom)?,
285            "orderUpdates" => Vec::<OrderUpdate>::deserialize(helper.data)
286                .map(WebSocketData::OrderUpdatesData)
287                .map_err(serde::de::Error::custom)?,
288            "candle" => CandleData::deserialize(helper.data)
289                .map(WebSocketData::CandleData)
290                .map_err(serde::de::Error::custom)?,
291            "activeAssetCtx" => ActiveAssetCtxData::deserialize(helper.data)
292                .map(WebSocketData::ActiveAssetCtxData)
293                .map_err(serde::de::Error::custom)?,
294            "bbo" => WsBbo::deserialize(helper.data)
295                .map(WebSocketData::Bbo)
296                .map_err(serde::de::Error::custom)?,
297            "l2Book" => WsBook::deserialize(helper.data)
298                .map(WebSocketData::L2Book)
299                .map_err(serde::de::Error::custom)?,
300            _ => return Err(serde::de::Error::custom("unknown channel type")),
301        };
302        Ok(WebSocketMessage {
303            _channel: helper.channel,
304            data,
305        })
306    }
307}
308
309impl HyperliquidConnector {
310    pub async fn new(
311        rest_endpoint: &str,
312        web_socket_endpoint: &str,
313        private_key: &str,
314        evm_wallet_address: &str,
315        vault_address: Option<String>,
316        use_agent: bool,
317        agent_name: Option<String>,
318        symbol_list: &[&str],
319    ) -> Result<Self, DexError> {
320        let request = DexRequest::new(rest_endpoint.to_owned()).await?;
321        let web_socket = DexWebSocket::new(web_socket_endpoint.to_owned());
322
323        let evm_wallet_address = vault_address
324            .clone()
325            .unwrap_or_else(|| evm_wallet_address.into());
326        let config = Config {
327            evm_wallet_address,
328            symbol_list: symbol_list.iter().map(|s| s.to_string()).collect(),
329        };
330
331        let vault_address: Option<H160> = vault_address
332            .as_deref()
333            .and_then(|v| H160::from_str(v).ok());
334
335        let mut local_wallet: LocalWallet = private_key.parse().unwrap();
336
337        if use_agent {
338            let ec_tmp =
339                ExchangeClient::new(None, local_wallet, Some(BaseUrl::Mainnet), None, None)
340                    .await
341                    .map_err(|e| DexError::Other(e.to_string()))?;
342
343            let (pk, status) = ec_tmp
344                .approve_agent(None, agent_name.clone())
345                .await
346                .map_err(|e| DexError::Other(e.to_string()))?;
347
348            match status {
349                ExchangeResponseStatus::Ok(_) => {
350                    log::info!("approve_agent succeeded for {:?}", agent_name);
351                    local_wallet = pk
352                        .parse()
353                        .map_err(|e| DexError::Other(format!("Failed to parse agent pk: {e}")))?;
354                }
355                ExchangeResponseStatus::Err(e) => {
356                    log::error!("approve_agent failed: {e}");
357                    return Err(DexError::Other(format!("approve_agent failed: {e}")));
358                }
359            }
360        }
361
362        let exchange_client = ExchangeClient::new(
363            None,
364            local_wallet,
365            Some(BaseUrl::Mainnet),
366            None,
367            vault_address,
368        )
369        .await
370        .map_err(|e| DexError::Other(e.to_string()))?;
371
372        let mut instance = HyperliquidConnector {
373            config,
374            request,
375            web_socket,
376            trade_results: Arc::new(RwLock::new(HashMap::new())),
377            canceled_results: Arc::new(RwLock::new(HashMap::new())),
378            running: Arc::new(AtomicBool::new(false)),
379            read_socket: Arc::new(Mutex::new(None)),
380            write_socket: Arc::new(Mutex::new(None)),
381            task_handle_read_message: Arc::new(Mutex::new(None)),
382            task_handle_read_sigterm: Arc::new(Mutex::new(None)),
383            dynamic_market_info: Arc::new(RwLock::new(HashMap::new())),
384            static_market_info: HashMap::new(),
385            spot_index_map: HashMap::new(),
386            spot_reverse_map: Arc::new(HashMap::new()),
387            exchange_client,
388            maintenance: Arc::new(RwLock::new(MaintenanceInfo {
389                next_start: None,
390                fetched_at: Utc::now() - ChronoDuration::hours(1),
391            })),
392            last_volumes: Arc::new(Mutex::new(HashMap::new())),
393        };
394
395        instance.spawn_maintenance_watcher();
396
397        instance.retrive_market_metadata().await?;
398
399        let info_payload = serde_json::to_string(&InfoRequest {
400            req_type: "spotMeta",
401            user: None,
402        })
403        .map_err(|e| DexError::Other(e.to_string()))?;
404
405        let spot_meta: SpotMetaResponse = instance
406            .request
407            .handle_request::<SpotMetaResponse, InfoRequest<'_>>(
408                HttpMethod::Post,
409                "/info".into(),
410                &HashMap::new(),
411                info_payload,
412            )
413            .await?;
414
415        // index → token_name
416        let token_name_map: HashMap<usize, String> = spot_meta
417            ._tokens
418            .iter()
419            .map(|t| (t._index, t._name.clone()))
420            .collect();
421
422        let mut idx_from_pair = HashMap::<String, usize>::new();
423        let mut pair_from_idx = HashMap::<usize, String>::new();
424
425        for uni in &spot_meta.universe {
426            let pair = if !uni.name.starts_with('@') {
427                uni.name.clone()
428            } else if uni._tokens.len() == 2 {
429                format!(
430                    "{}/{}",
431                    token_name_map.get(&uni._tokens[0]).unwrap_or(&"?".into()),
432                    token_name_map.get(&uni._tokens[1]).unwrap_or(&"?".into())
433                )
434            } else {
435                log::warn!(
436                    "universe idx {} has unexpected token vec {:?}",
437                    uni.index,
438                    uni._tokens
439                );
440                uni.name.clone()
441            };
442
443            idx_from_pair.insert(pair.clone(), uni.index);
444            pair_from_idx.insert(uni.index, pair);
445        }
446
447        instance.spot_index_map = idx_from_pair;
448        instance.spot_reverse_map = Arc::new(pair_from_idx);
449
450        {
451            let token_decimals: HashMap<String, u32> = spot_meta
452                ._tokens
453                .iter()
454                .map(|t| (t._name.clone(), t._sz_decimals))
455                .collect();
456
457            let mut sm = std::mem::take(&mut instance.static_market_info);
458
459            for uni in &spot_meta.universe {
460                let pair = if !uni.name.starts_with('@') {
461                    uni.name.clone()
462                } else if uni._tokens.len() == 2 {
463                    format!(
464                        "{}/{}",
465                        spot_meta._tokens[uni._tokens[0]]._name,
466                        spot_meta._tokens[uni._tokens[1]]._name
467                    )
468                } else {
469                    uni.name.clone()
470                };
471
472                let base = pair.split('/').next().unwrap();
473                let decimals = *token_decimals.get(base).unwrap_or(&0);
474
475                sm.insert(
476                    pair.clone(),
477                    StaticMarketInfo {
478                        decimals,
479                        _max_leverage: 0,
480                    },
481                );
482            }
483
484            instance.static_market_info = sm;
485        }
486
487        Ok(instance)
488    }
489
490    fn spawn_maintenance_watcher(&self) {
491        let cache = self.maintenance.clone();
492        tokio::spawn(async move {
493            let client = Client::builder()
494                .timeout(std::time::Duration::from_secs(2))
495                .build()
496                .expect("reqwest client");
497
498            loop {
499                if let Ok(res) = client
500                    .get("https://hyperliquid.statuspage.io/api/v2/scheduled-maintenances/upcoming.json")
501                    .send()
502                    .await
503                {
504                    if let Ok(json) = res.json::<Value>().await {
505                        let next = json
506                            .get("scheduled_maintenances")
507                            .and_then(|v| v.get(0))
508                            .and_then(|v| v.get("scheduled_for"))
509                            .and_then(|v| v.as_str())
510                            .and_then(|s| DateTime::parse_from_rfc3339(s).ok())
511                            .map(|dt| dt.with_timezone(&Utc));
512
513                        *cache.write().await = MaintenanceInfo {
514                            next_start: next,
515                            fetched_at: Utc::now(),
516                        };
517                    }
518                }
519                sleep(Duration::from_secs(600)).await;
520            }
521        });
522    }
523
524    pub async fn start_web_socket(&self) -> Result<(), DexError> {
525        log::info!("start_web_socket");
526
527        let (write, read) = self
528            .web_socket
529            .clone()
530            .connect()
531            .await
532            .map_err(|_| DexError::Other("Failed to connect to WebSocket".to_string()))?;
533
534        {
535            let mut read_lock = self.read_socket.lock().await;
536            *read_lock = Some(read);
537        }
538        {
539            let mut write_lock = self.write_socket.lock().await;
540            *write_lock = Some(write);
541        }
542
543        self.running.store(true, Ordering::SeqCst);
544        self.subscribe_to_channels(&self.config.evm_wallet_address)
545            .await?;
546
547        let running = self.running.clone();
548        let read_sock = self.read_socket.clone();
549        let write_sock = self.write_socket.clone();
550        let dmi = self.dynamic_market_info.clone();
551        let trs = self.trade_results.clone();
552        let rev_map = self.spot_reverse_map.clone();
553        let crs = self.canceled_results.clone();
554        let static_info = self.static_market_info.clone();
555
556        let reader_handle = tokio::spawn(async move {
557            let mut idle_counter = 0;
558            while running.load(Ordering::SeqCst) {
559                if let Some(stream) = read_sock.lock().await.as_mut() {
560                    tokio::select! {
561                        msg = stream.next() => match msg {
562                            Some(Ok(Message::Text(txt))) => {
563                                idle_counter = 0;
564                                if txt == "{}" {
565                                    if let Some(w) = write_sock.lock().await.as_mut() {
566                                        let _ = w.send(Message::Text(txt)).await;
567                                    }
568                                } else {
569                                    if let Err(e) = HyperliquidConnector::handle_websocket_message(
570                                        Message::Text(txt),
571                                        dmi.clone(),
572                                        trs.clone(),
573                                        rev_map.clone(),
574                                        crs.clone(),
575                                        static_info.clone(),
576                                    ).await {
577                                        log::error!("WebSocket handler error: {:?}", e);
578                                        break;
579                                    }
580                                }
581                            }
582                            Some(Ok(_)) => {
583                            }
584                            Some(Err(err)) => {
585                                log::error!("WebSocket read error: {:?}", err);
586                                break;
587                            }
588                            None => {
589                                log::info!("WebSocket stream closed");
590                                break;
591                            }
592                        },
593                        _ = tokio::time::sleep(Duration::from_secs(10)) => {
594                            idle_counter += 1;
595                            if idle_counter >= 10 {
596                                log::error!("No WebSocket messages for 100s, shutting down reader");
597                                break;
598                            }
599                        }
600                    }
601                }
602            }
603            running.store(false, Ordering::SeqCst);
604            log::info!("WebSocket reader task ended");
605        });
606        *self.task_handle_read_message.lock().await = Some(reader_handle);
607
608        let running_for_sig = self.running.clone();
609        let sig_handle = tokio::spawn(async move {
610            let mut sigterm =
611                signal(SignalKind::terminate()).expect("Failed to bind SIGTERM handler");
612            loop {
613                select! {
614                    _ = sigterm.recv() => {
615                        log::info!("SIGTERM received, stopping WebSocket");
616                        running_for_sig.store(false, Ordering::SeqCst);
617                        break;
618                    }
619                    _ = tokio::time::sleep(Duration::from_secs(1)) => {
620                        if !running_for_sig.load(Ordering::SeqCst) {
621                            break;
622                        }
623                    }
624                }
625            }
626        });
627        *self.task_handle_read_sigterm.lock().await = Some(sig_handle);
628
629        Ok(())
630    }
631
632    pub async fn stop_web_socket(&self) -> Result<(), DexError> {
633        log::info!("stop_web_socket");
634        self.running.store(false, Ordering::SeqCst);
635
636        {
637            let mut write_guard = self.write_socket.lock().await;
638            if let Some(write_socket) = write_guard.as_mut() {
639                if let Err(e) = write_socket.send(Message::Close(None)).await {
640                    log::error!("Failed to send WebSocket close message: {:?}", e);
641                }
642            }
643            *write_guard = None;
644        }
645
646        {
647            let mut read_guard = self.read_socket.lock().await;
648            *read_guard = None;
649        }
650
651        if let Some(handle) = self.task_handle_read_message.lock().await.take() {
652            let _ = handle.await;
653        }
654
655        if let Some(handle) = self.task_handle_read_sigterm.lock().await.take() {
656            let _ = handle.await;
657        }
658
659        drop(self.web_socket.clone());
660
661        Ok(())
662    }
663
664    async fn subscribe_to_channels(&self, user_address: &str) -> Result<(), DexError> {
665        let all_mids_subscription = serde_json::json!({
666            "method": "subscribe",
667            "subscription": {
668                "type": "allMids"
669            }
670        })
671        .to_string();
672
673        let user_fills_subscription = serde_json::json!({
674            "method": "subscribe",
675            "subscription": {
676                "type": "userFills",
677                "user": user_address
678            }
679        })
680        .to_string();
681
682        let order_updates_subscription = serde_json::json!({
683            "method": "subscribe",
684            "subscription": {
685                "type": "orderUpdates",
686                "user": user_address
687            }
688        })
689        .to_string();
690
691        let mut write_socket_lock = self.write_socket.lock().await;
692
693        if let Some(write_socket) = write_socket_lock.as_mut() {
694            if let Err(e) = write_socket
695                .send(Message::Text(all_mids_subscription))
696                .await
697            {
698                return Err(DexError::WebSocketError(format!(
699                    "Failed to subscribe to allMids: {}",
700                    e
701                )));
702            }
703
704            if let Err(e) = write_socket
705                .send(Message::Text(order_updates_subscription))
706                .await
707            {
708                return Err(DexError::WebSocketError(format!(
709                    "Failed to subscribe to userFills: {}",
710                    e
711                )));
712            }
713
714            if let Err(e) = write_socket
715                .send(Message::Text(user_fills_subscription))
716                .await
717            {
718                return Err(DexError::WebSocketError(format!(
719                    "Failed to subscribe to userFills: {}",
720                    e
721                )));
722            }
723
724            for symbol in &self.config.symbol_list {
725                let coin = resolve_coin(symbol, &self.spot_index_map);
726                let candle_subscription = serde_json::json!({
727                    "method": "subscribe",
728                    "subscription": {
729                        "type": "candle",
730                        "coin": coin,
731                        "interval": "1m"
732                    }
733                })
734                .to_string();
735                if let Err(e) = write_socket.send(Message::Text(candle_subscription)).await {
736                    return Err(DexError::WebSocketError(format!(
737                        "Failed to subscribe to candle for {}: {}",
738                        symbol, e
739                    )));
740                }
741
742                let active_asset_ctx_subscription = serde_json::json!({
743                    "method": "subscribe",
744                    "subscription": {
745                        "type": "activeAssetCtx",
746                        "coin": coin,
747                    }
748                })
749                .to_string();
750                if let Err(e) = write_socket
751                    .send(Message::Text(active_asset_ctx_subscription))
752                    .await
753                {
754                    return Err(DexError::WebSocketError(format!(
755                        "Failed to subscribe to activeAssetCtx: {}",
756                        e
757                    )));
758                }
759
760                let bbo_subscription = serde_json::json!({
761                    "method": "subscribe",
762                    "subscription": {
763                        "type": "bbo",
764                        "coin": coin
765                    }
766                })
767                .to_string();
768                if let Err(e) = write_socket.send(Message::Text(bbo_subscription)).await {
769                    return Err(DexError::WebSocketError(format!(
770                        "Failed to subscribe to bbo: {}",
771                        e
772                    )));
773                }
774
775                let l2_subscription = serde_json::json!({
776                    "method": "subscribe",
777                    "subscription": {
778                        "type": "l2Book",
779                        "coin": coin
780                    }
781                })
782                .to_string();
783                if let Err(e) = write_socket.send(Message::Text(l2_subscription)).await {
784                    return Err(DexError::WebSocketError(format!(
785                        "Failed to subscribe to l2: {}",
786                        e
787                    )));
788                }
789            }
790        } else {
791            return Err(DexError::WebSocketError(
792                "Write socket is not available".to_string(),
793            ));
794        }
795
796        Ok(())
797    }
798
799    async fn handle_websocket_message(
800        msg: Message,
801        dynamic_market_info: Arc<RwLock<HashMap<String, DynamicMarketInfo>>>,
802        trade_results: Arc<RwLock<HashMap<String, HashMap<String, TradeResult>>>>,
803        spot_reverse_map: Arc<HashMap<usize, String>>,
804        canceled_results: Arc<RwLock<HashMap<String, HashMap<String, CancelEvent>>>>,
805        static_market_info: HashMap<String, StaticMarketInfo>,
806    ) -> Result<(), DexError> {
807        if let Message::Text(text) = msg {
808            for line in text.split('\n') {
809                if line.is_empty() {
810                    continue;
811                }
812                if let Ok(message) = serde_json::from_str::<WebSocketMessage>(line) {
813                    log::trace!("[WebSocketMessage] channel = {}", message._channel);
814
815                    match &message.data {
816                        WebSocketData::Bbo(bbo) => {
817                            log::trace!("[WebSocketMessage] BBO coin = {}", bbo.coin);
818                        }
819                        WebSocketData::L2Book(book) => {
820                            log::trace!("[WebSocketMessage] L2Book coin = {}", book.coin);
821                        }
822                        WebSocketData::AllMidsData(data) => {
823                            log::trace!("[WebSocketMessage] allMids keys = {:?}", data.mids.keys());
824                        }
825                        _ => {}
826                    }
827
828                    match message.data {
829                        WebSocketData::AllMidsData(ref data) => {
830                            Self::process_all_mids_message(
831                                data,
832                                dynamic_market_info.clone(),
833                                spot_reverse_map.clone(),
834                                &static_market_info,
835                            )
836                            .await;
837                        }
838                        WebSocketData::CandleData(ref data) => {
839                            Self::process_candle_message(
840                                data,
841                                dynamic_market_info.clone(),
842                                spot_reverse_map.clone(),
843                            )
844                            .await;
845                        }
846                        WebSocketData::UserFillsData(ref data) => {
847                            Self::process_account_data(data, trade_results.clone()).await;
848                        }
849                        WebSocketData::ActiveAssetCtxData(ref data) => {
850                            Self::process_active_asset_ctx_message(
851                                data,
852                                dynamic_market_info.clone(),
853                                spot_reverse_map.clone(),
854                            )
855                            .await;
856                        }
857                        WebSocketData::OrderUpdatesData(ref orders) => {
858                            Self::process_order_updates_message(
859                                orders,
860                                canceled_results.clone(),
861                                trade_results.clone(),
862                            )
863                            .await;
864                        }
865                        WebSocketData::Bbo(ref bbo) => {
866                            // "@123" → 123 → "UBTC/USDC"
867                            let idx = bbo
868                                .coin
869                                .strip_prefix('@')
870                                .and_then(|s| s.parse::<usize>().ok());
871                            let coin = idx
872                                .and_then(|i| spot_reverse_map.get(&i).cloned())
873                                .unwrap_or_else(|| bbo.coin.clone());
874                            let market_key = if coin.contains('/') {
875                                coin.clone()
876                            } else {
877                                format!("{}-USD", coin)
878                            };
879                            let mut info_map = dynamic_market_info.write().await;
880                            let info = info_map.entry(market_key).or_default();
881                            info.best_bid = bbo
882                                .bbo
883                                .get(0)
884                                .and_then(|lvl| lvl.as_ref())
885                                .map(|l| Decimal::from_str(&l.px).unwrap());
886                            info.best_ask = bbo
887                                .bbo
888                                .get(1)
889                                .and_then(|lvl| lvl.as_ref())
890                                .map(|l| Decimal::from_str(&l.px).unwrap());
891                        }
892                        WebSocketData::L2Book(ref book) => {
893                            let idx = book
894                                .coin
895                                .strip_prefix('@')
896                                .and_then(|s| s.parse::<usize>().ok());
897                            let coin = idx
898                                .and_then(|i| spot_reverse_map.get(&i).cloned())
899                                .unwrap_or_else(|| book.coin.clone());
900                            let market_key = if coin.contains('/') {
901                                coin.clone()
902                            } else {
903                                format!("{}-USD", coin)
904                            };
905                            let mut info_map = dynamic_market_info.write().await;
906                            let info = info_map.entry(market_key).or_default();
907                            info.best_bid = book.levels[0]
908                                .get(0)
909                                .map(|lvl| Decimal::from_str(&lvl.px).unwrap());
910                            info.best_ask = book.levels[1]
911                                .get(0)
912                                .map(|lvl| Decimal::from_str(&lvl.px).unwrap());
913                        }
914                    }
915                }
916            }
917        }
918        Ok(())
919    }
920
921    async fn process_order_updates_message(
922        orders: &[OrderUpdate],
923        canceled_results: Arc<RwLock<HashMap<String, HashMap<String, CancelEvent>>>>,
924        trade_results: Arc<RwLock<HashMap<String, HashMap<String, TradeResult>>>>,
925    ) {
926        for upd in orders.iter() {
927            let symbol = if upd.order.coin.contains('/') || upd.order.coin.contains('-') {
928                upd.order.coin.clone()
929            } else {
930                format!("{}-USD", upd.order.coin)
931            };
932
933            match upd.status.as_str() {
934                "canceled" => {
935                    let evt = CancelEvent {
936                        order_id: upd.order.oid.to_string(),
937                        timestamp: upd.status_timestamp,
938                    };
939                    canceled_results
940                        .write()
941                        .await
942                        .entry(symbol)
943                        .or_default()
944                        .insert(evt.order_id.clone(), evt);
945                }
946                "rejected" => {
947                    let mut trs = trade_results.write().await;
948                    let entry = trs.entry(symbol).or_default();
949                    entry.insert(
950                        upd.order.oid.to_string(),
951                        TradeResult {
952                            filled_side: OrderSide::Long,
953                            filled_size: Decimal::ZERO,
954                            filled_value: Decimal::ZERO,
955                            filled_fee: Decimal::ZERO,
956                            order_id: upd.order.oid.to_string(),
957                            is_rejected: true,
958                        },
959                    );
960                }
961                _ => {}
962            }
963        }
964    }
965
966    async fn process_all_mids_message(
967        mids_data: &AllMidsData,
968        dynamic_market_info: Arc<RwLock<HashMap<String, DynamicMarketInfo>>>,
969        spot_reverse_map: Arc<HashMap<usize, String>>,
970        static_market_info: &HashMap<String, StaticMarketInfo>,
971    ) {
972        for (raw_coin, mid_price_str) in &mids_data.mids {
973            let coin = if let Some(stripped) = raw_coin.strip_prefix('@') {
974                let idx = stripped.parse::<usize>().unwrap_or_else(|_| {
975                    log::warn!("[resolve_coin] invalid @index: {}", stripped);
976                    0
977                });
978
979                let mapped = spot_reverse_map.get(&idx).cloned();
980                match mapped {
981                    Some(mapped) => mapped,
982                    None => {
983                        log::trace!(
984                            "[resolve_coin] spot_reverse_map missing: {} (index: {})",
985                            raw_coin,
986                            idx
987                        );
988                        raw_coin.clone()
989                    }
990                }
991            } else {
992                raw_coin.clone()
993            };
994
995            let market_key = if coin.contains('/') || coin.contains('-') {
996                coin.clone() // Spot: UBTC/USDC,  etc.
997            } else {
998                format!("{}-USD", coin) // Perp: BTC-USD, etc.
999            };
1000
1001            if let Ok(mid) = string_to_decimal(Some(mid_price_str.clone())) {
1002                let mut guard = dynamic_market_info.write().await;
1003                let info = guard.entry(market_key.clone()).or_default();
1004                let sz_decimals = static_market_info
1005                    .get(&market_key)
1006                    .map(|m| m.decimals)
1007                    .unwrap_or_else(|| {
1008                        log::trace!("no static for {}, default 0", market_key);
1009                        0
1010                    });
1011                let is_spot = market_key.contains('/');
1012
1013                let base_tick = Self::calculate_min_tick(mid, sz_decimals, is_spot);
1014                info.min_tick = Some(base_tick);
1015                info.market_price = Some(mid);
1016            }
1017        }
1018    }
1019
1020    async fn process_candle_message(
1021        candle: &CandleData,
1022        dynamic_market_info: Arc<RwLock<HashMap<String, DynamicMarketInfo>>>,
1023        spot_reverse_map: Arc<HashMap<usize, String>>,
1024    ) {
1025        let coin = if let Some(stripped) = candle.s.strip_prefix('@') {
1026            stripped
1027                .parse::<usize>()
1028                .ok()
1029                .and_then(|idx| spot_reverse_map.get(&idx).cloned())
1030                .unwrap_or_else(|| {
1031                    log::trace!(
1032                        "in spot_reverse_map: {} is missing (@{})",
1033                        candle.s,
1034                        stripped
1035                    );
1036                    candle.s.clone()
1037                })
1038        } else {
1039            candle.s.clone()
1040        };
1041
1042        let market_key = if coin.contains('/') || coin.contains('-') {
1043            coin.clone()
1044        } else {
1045            format!("{}-USD", coin)
1046        };
1047
1048        let mut guard = dynamic_market_info.write().await;
1049        let info = guard.entry(market_key.clone()).or_default();
1050        info.volume = Some(candle.v);
1051        info.num_trades = Some(candle.n);
1052    }
1053
1054    async fn process_active_asset_ctx_message(
1055        asset_data: &ActiveAssetCtxData,
1056        dynamic_market_info: Arc<RwLock<HashMap<String, DynamicMarketInfo>>>,
1057        spot_reverse_map: Arc<HashMap<usize, String>>,
1058    ) {
1059        let coin = if let Some(stripped) = asset_data.coin.strip_prefix('@') {
1060            stripped
1061                .parse::<usize>()
1062                .ok()
1063                .and_then(|idx| spot_reverse_map.get(&idx).cloned())
1064                .unwrap_or_else(|| {
1065                    log::trace!(
1066                        "in spot_reverse_map {} is missing (@{})",
1067                        asset_data.coin,
1068                        stripped
1069                    );
1070                    asset_data.coin.clone()
1071                })
1072        } else {
1073            asset_data.coin.clone()
1074        };
1075
1076        let market_key = if coin.contains('/') || coin.contains('-') {
1077            coin.clone()
1078        } else {
1079            format!("{}-USD", coin)
1080        };
1081
1082        let mut guard = dynamic_market_info.write().await;
1083        let info = guard
1084            .entry(market_key.clone())
1085            .or_insert_with(DynamicMarketInfo::default);
1086        info.funding_rate = Some(asset_data.ctx.funding);
1087        info.open_interest = Some(asset_data.ctx.openInterest);
1088        info.oracle_price = Some(asset_data.ctx.oraclePx);
1089    }
1090
1091    async fn process_account_data(
1092        data: &UserFillsData,
1093        trade_results: Arc<RwLock<HashMap<String, HashMap<String, TradeResult>>>>,
1094    ) {
1095        for fill in &data.fills {
1096            log::debug!("{:?}", fill);
1097
1098            let filled_side = if fill.side == "A" {
1099                OrderSide::Short
1100            } else {
1101                OrderSide::Long
1102            };
1103
1104            let filled_size = fill.sz;
1105            let filled_price = fill.px;
1106            let filled_value = filled_size * filled_price;
1107            let filled_fee = fill.fee;
1108            let order_id = fill.oid;
1109            let trade_id = fill.tid;
1110
1111            let market_id = if fill.coin.contains('/') || fill.coin.contains('-') {
1112                fill.coin.clone()
1113            } else {
1114                format!("{}-USD", fill.coin)
1115            };
1116
1117            let trade_result = TradeResult {
1118                filled_side,
1119                filled_size,
1120                filled_value,
1121                filled_fee,
1122                order_id: order_id.to_string(),
1123                is_rejected: false,
1124            };
1125
1126            let mut trade_results_guard = trade_results.write().await;
1127            let market_map = trade_results_guard.entry(market_id.clone()).or_default();
1128            let key = trade_id.to_string();
1129
1130            if let Some(existing) = market_map.get_mut(&key) {
1131                existing.filled_size += trade_result.filled_size;
1132                existing.filled_value += trade_result.filled_value;
1133                existing.filled_fee += trade_result.filled_fee;
1134            } else {
1135                market_map.insert(key, trade_result);
1136            }
1137        }
1138    }
1139}
1140
1141#[derive(Serialize, Debug, Clone)]
1142struct HyperliquidDefaultPayload {
1143    r#type: String,
1144    #[serde(skip_serializing_if = "Option::is_none")]
1145    user: Option<String>,
1146}
1147
1148#[derive(Deserialize, Debug)]
1149struct HyperliquidRetrieveUserStateResponse {
1150    #[serde(rename = "marginSummary")]
1151    margin_summary: Option<HyperliquidMarginSummary>,
1152}
1153#[derive(Deserialize, Debug)]
1154struct HyperliquidMarginSummary {
1155    #[serde(rename = "accountValue")]
1156    account_value: String,
1157    #[serde(rename = "totalRawUsd")]
1158    total_rawusd: String,
1159}
1160
1161#[derive(Deserialize, Debug)]
1162struct HyperliquidRetriveUserOpenOrder {
1163    coin: String,
1164    oid: u64,
1165}
1166
1167#[derive(Deserialize, Debug)]
1168struct HyperliquidRetriveUserPositionResponse {
1169    #[serde(rename = "assetPositions")]
1170    asset_positions: Vec<HyperliquidRetriveUserPositionResponseBody>,
1171}
1172#[derive(Deserialize, Debug)]
1173struct HyperliquidRetriveUserPositionResponseBody {
1174    position: HyperliquidRetriveUserPosition,
1175}
1176#[derive(Deserialize, Debug)]
1177struct HyperliquidRetriveUserPosition {
1178    coin: String,
1179    szi: Decimal,
1180}
1181
1182#[derive(Deserialize, Debug)]
1183struct HyperliquidRetriveMarketMetadataResponse {
1184    universe: Vec<HyperliquidRetriveMarketMetadata>,
1185}
1186#[derive(Deserialize, Debug)]
1187struct HyperliquidRetriveMarketMetadata {
1188    name: String,
1189    #[serde(rename = "szDecimals")]
1190    decimals: u32,
1191    #[serde(rename = "maxLeverage")]
1192    max_leverage: u32,
1193}
1194
1195#[derive(Deserialize, Debug)]
1196struct HyperliquidSpotBalanceResponse {
1197    balances: Vec<HyperliquidSpotBalance>,
1198}
1199
1200#[derive(Deserialize, Debug)]
1201struct HyperliquidSpotBalance {
1202    coin: String,
1203    total: String,
1204}
1205
1206#[async_trait]
1207impl DexConnector for HyperliquidConnector {
1208    async fn start(&self) -> Result<(), DexError> {
1209        self.start_web_socket().await?;
1210        sleep(Duration::from_secs(5)).await;
1211        self.wait_for_market_ready(60).await?;
1212        Ok(())
1213    }
1214
1215    async fn stop(&self) -> Result<(), DexError> {
1216        self.stop_web_socket().await?;
1217        Ok(())
1218    }
1219
1220    async fn restart(&self, max_retries: i32) -> Result<(), DexError> {
1221        log::info!("Restarting WebSocket connection...");
1222
1223        let mut retry_count = 0;
1224        let mut backoff_delay = Duration::from_secs(1);
1225
1226        while retry_count < max_retries {
1227            if let Err(e) = self.stop_web_socket().await {
1228                log::error!(
1229                    "Failed to stop WebSocket on attempt {}: {:?}",
1230                    retry_count + 1,
1231                    e
1232                );
1233            } else {
1234                log::info!(
1235                    "Successfully stopped WebSocket on attempt {}.",
1236                    retry_count + 1
1237                );
1238            }
1239
1240            sleep(backoff_delay).await;
1241
1242            match self.start_web_socket().await {
1243                Ok(_) => {
1244                    log::info!(
1245                        "Successfully started WebSocket on attempt {}.",
1246                        retry_count + 1
1247                    );
1248                    return Ok(());
1249                }
1250                Err(e) => {
1251                    log::error!(
1252                        "Failed to start WebSocket on attempt {}: {:?}",
1253                        retry_count + 1,
1254                        e
1255                    );
1256                    retry_count += 1;
1257                    backoff_delay *= 2; // Exponential backoff
1258                }
1259            }
1260        }
1261
1262        log::error!(
1263            "Failed to restart WebSocket after {} attempts.",
1264            max_retries
1265        );
1266        Err(DexError::Other(format!(
1267            "Failed to restart WebSocket after {} attempts.",
1268            max_retries
1269        )))
1270    }
1271
1272    async fn set_leverage(&self, symbol: &str, leverage: u32) -> Result<(), DexError> {
1273        let asset = Self::extract_asset_name(symbol);
1274        self.exchange_client
1275            .update_leverage(leverage, asset, false, None)
1276            .await
1277            .map_err(|e| DexError::Other(e.to_string()))?;
1278        Ok(())
1279    }
1280
1281    async fn get_ticker(
1282        &self,
1283        symbol: &str,
1284        _test_price: Option<Decimal>,
1285    ) -> Result<TickerResponse, DexError> {
1286        if !self.running.load(Ordering::SeqCst) {
1287            return Err(DexError::NoConnection);
1288        }
1289
1290        let dynamic_info_guard = self.dynamic_market_info.read().await;
1291        let dynamic_info = dynamic_info_guard
1292            .get(symbol)
1293            .ok_or_else(|| DexError::Other("No dynamic market info available".to_string()))?;
1294        let price = dynamic_info
1295            .market_price
1296            .ok_or_else(|| DexError::Other("No price available".to_string()))?;
1297        let min_tick = dynamic_info.min_tick;
1298        let num_trades = dynamic_info.num_trades;
1299        let funding_rate = dynamic_info.funding_rate;
1300        let open_interest = dynamic_info.open_interest;
1301        let oracle_price = dynamic_info.oracle_price;
1302
1303        let cur_vol = dynamic_info.volume.unwrap_or(Decimal::ZERO);
1304        let mut lv = self.last_volumes.lock().await;
1305        let prev_vol = lv.get(symbol).cloned().unwrap_or(Decimal::ZERO);
1306        // make sure we never return a negative delta if cur_vol resets each candle
1307        let delta_vol = if cur_vol >= prev_vol {
1308            cur_vol - prev_vol
1309        } else {
1310            // volume counter has rolled over/reset at candle boundary
1311            cur_vol
1312        };
1313        lv.insert(symbol.to_string(), cur_vol);
1314
1315        Ok(TickerResponse {
1316            symbol: symbol.to_owned(),
1317            price,
1318            min_tick,
1319            min_order: None,
1320            volume: Some(delta_vol),
1321            num_trades,
1322            funding_rate,
1323            open_interest,
1324            oracle_price,
1325        })
1326    }
1327
1328    async fn get_filled_orders(&self, symbol: &str) -> Result<FilledOrdersResponse, DexError> {
1329        let mut response: Vec<FilledOrder> = vec![];
1330        let trade_results_guard = self.trade_results.read().await;
1331        let orders = match trade_results_guard.get(symbol) {
1332            Some(v) => v,
1333            None => return Ok(FilledOrdersResponse::default()),
1334        };
1335        for (trade_id, order) in orders.iter() {
1336            let filled_order = FilledOrder {
1337                order_id: order.order_id.clone(),
1338                trade_id: trade_id.clone(),
1339                is_rejected: order.is_rejected,
1340                filled_side: Some(order.filled_side.clone()),
1341                filled_size: Some(order.filled_size),
1342                filled_fee: Some(order.filled_fee),
1343                filled_value: Some(order.filled_value),
1344            };
1345            response.push(filled_order);
1346        }
1347
1348        Ok(FilledOrdersResponse { orders: response })
1349    }
1350
1351    async fn get_canceled_orders(&self, symbol: &str) -> Result<CanceledOrdersResponse, DexError> {
1352        let mut resp = Vec::new();
1353        let guard = self.canceled_results.read().await;
1354        if let Some(map) = guard.get(symbol) {
1355            for (_, evt) in map.iter() {
1356                resp.push(CanceledOrder {
1357                    order_id: evt.order_id.clone(),
1358                    canceled_timestamp: evt.timestamp,
1359                });
1360            }
1361        }
1362        Ok(CanceledOrdersResponse { orders: resp })
1363    }
1364
1365    async fn get_open_orders(&self, symbol: &str) -> Result<OpenOrdersResponse, DexError> {
1366        let all_orders = self.get_orders().await?;
1367
1368        // Convert symbol to internal representation if needed
1369        let target_coin = if let Some(internal_id) = self.spot_index_map.get(symbol) {
1370            format!("@{}", internal_id)
1371        } else {
1372            symbol.to_string()
1373        };
1374
1375        let filtered_orders: Vec<crate::OpenOrder> = all_orders
1376            .into_iter()
1377            .filter(|order| order.coin == target_coin || order.coin == symbol)
1378            .map(|hyperliquid_order| crate::OpenOrder {
1379                order_id: hyperliquid_order.oid.to_string(),
1380                symbol: symbol.to_string(),
1381                side: OrderSide::Long, // Default, we need more detailed API to get actual side
1382                size: Decimal::ZERO,   // Default, we need more detailed API to get actual size
1383                price: Decimal::ZERO,  // Default, we need more detailed API to get actual price
1384                status: "open".to_string(),
1385            })
1386            .collect();
1387
1388        Ok(OpenOrdersResponse {
1389            orders: filtered_orders,
1390        })
1391    }
1392
1393    async fn get_balance(&self, symbol: Option<&str>) -> Result<BalanceResponse, DexError> {
1394        // Always get both spot and perp balances
1395        let spot_action = HyperliquidDefaultPayload {
1396            r#type: "spotClearinghouseState".into(),
1397            user: Some(self.config.evm_wallet_address.clone()),
1398        };
1399        let spot_res: HyperliquidSpotBalanceResponse = self
1400            .handle_request_with_action("/info".into(), &spot_action)
1401            .await?;
1402
1403        let perp_action = HyperliquidDefaultPayload {
1404            r#type: "clearinghouseState".into(),
1405            user: Some(self.config.evm_wallet_address.clone()),
1406        };
1407        let perp_res = self
1408            .handle_request_with_action::<HyperliquidRetrieveUserStateResponse, _>(
1409                "/info".into(),
1410                &perp_action,
1411            )
1412            .await?;
1413
1414        log::debug!("spot balances = {:?}", spot_res.balances);
1415        log::debug!("perp margin summary = {:?}", perp_res.margin_summary);
1416
1417        if let Some(pair) = symbol {
1418            // "UBTC/USDC" → "UBTC"
1419            let base_coin = pair.split('/').next().unwrap_or(pair);
1420
1421            let mut usdc_total = Decimal::ZERO;
1422            let mut base_total = Decimal::ZERO;
1423            for b in &spot_res.balances {
1424                match b.coin.as_str() {
1425                    "USDC" => {
1426                        usdc_total = parse_to_decimal(&b.total)?;
1427                        log::debug!("USDC total = {}", usdc_total);
1428                    }
1429                    c if c == base_coin => {
1430                        base_total = parse_to_decimal(&b.total)?;
1431                        log::debug!("{} total = {}", c, base_total);
1432                    }
1433                    _ => {}
1434                }
1435            }
1436
1437            let price_key = pair.to_string();
1438            let px = self
1439                .get_market_price(&price_key)
1440                .await
1441                .unwrap_or(Decimal::ZERO);
1442
1443            log::debug!("price_key = {}, px = {}", price_key, px);
1444
1445            // Get perp balance as well
1446            let (perp_equity, perp_balance) = if let Some(summary) = &perp_res.margin_summary {
1447                let equity = parse_to_decimal(&summary.account_value)?;
1448                let balance = parse_to_decimal(&summary.total_rawusd)?;
1449                log::debug!("perp equity = {}, perp balance = {}", equity, balance);
1450                (equity, balance)
1451            } else {
1452                (Decimal::ZERO, Decimal::ZERO)
1453            };
1454
1455            // Combine spot and perp balances
1456            let spot_equity = base_total * px + usdc_total;
1457            let total_equity = spot_equity + perp_equity;
1458            let total_balance = usdc_total + perp_balance;
1459
1460            log::debug!(
1461                "final equity = {} (spot: {} + perp: {}), balance = {} (spot: {} + perp: {})",
1462                total_equity,
1463                spot_equity,
1464                perp_equity,
1465                total_balance,
1466                usdc_total,
1467                perp_balance
1468            );
1469
1470            return Ok(BalanceResponse {
1471                equity: total_equity,
1472                balance: total_balance,
1473                position_entry_price: None,
1474                position_sign: None,
1475            });
1476        }
1477
1478        // When no symbol is specified, return total balance from both spot and perp
1479        let mut total_equity = Decimal::ZERO;
1480        let mut total_balance = Decimal::ZERO;
1481
1482        // Add spot balances
1483        for b in &spot_res.balances {
1484            let balance = parse_to_decimal(&b.total)?;
1485            if b.coin == "USDC" {
1486                total_balance += balance;
1487                total_equity += balance;
1488            } else {
1489                // For non-USDC tokens, try to get their price and add to equity
1490                let symbol_key = format!("{}/USDC", b.coin);
1491                if let Ok(px) = self.get_market_price(&symbol_key).await {
1492                    total_equity += balance * px;
1493                }
1494            }
1495        }
1496
1497        // Add perp balances
1498        if let Some(summary) = perp_res.margin_summary {
1499            let perp_equity = parse_to_decimal(&summary.account_value)?;
1500            let perp_balance = parse_to_decimal(&summary.total_rawusd)?;
1501            total_equity += perp_equity;
1502            total_balance += perp_balance;
1503            log::debug!(
1504                "perp equity = {}, perp balance = {}",
1505                perp_equity,
1506                perp_balance
1507            );
1508        }
1509
1510        log::debug!(
1511            "final total equity = {}, total balance = {}",
1512            total_equity,
1513            total_balance
1514        );
1515        Ok(BalanceResponse {
1516            equity: total_equity,
1517            balance: total_balance,
1518            position_entry_price: None,
1519            position_sign: None,
1520        })
1521    }
1522
1523    async fn get_combined_balance(&self) -> Result<CombinedBalanceResponse, DexError> {
1524        // Hyperliquid connector: minimal implementation for compilation only
1525        Err(DexError::Other(
1526            "get_combined_balance not implemented for HyperLiquid".to_string(),
1527        ))
1528    }
1529
1530    async fn get_last_trades(&self, _symbol: &str) -> Result<LastTradesResponse, DexError> {
1531        // TODO: Implement HyperLiquid last trades functionality
1532        Err(DexError::Other(
1533            "get_last_trades not implemented for HyperLiquid".to_string(),
1534        ))
1535    }
1536
1537    async fn clear_filled_order(&self, symbol: &str, trade_id: &str) -> Result<(), DexError> {
1538        let mut m = self.trade_results.write().await;
1539        if let Some(map) = m.get_mut(symbol) {
1540            if map.remove(trade_id).is_some() {
1541                Ok(())
1542            } else {
1543                Err(DexError::Other(format!(
1544                    "filled trade(trade_id:{}({})) does not exist",
1545                    trade_id, symbol
1546                )))
1547            }
1548        } else {
1549            Err(DexError::Other(format!(
1550                "filled trade(symbol:{}({})) does not exist",
1551                symbol, trade_id
1552            )))
1553        }
1554    }
1555
1556    async fn clear_all_filled_orders(&self) -> Result<(), DexError> {
1557        let mut trade_results_guard = self.trade_results.write().await;
1558        trade_results_guard.clear();
1559        Ok(())
1560    }
1561
1562    async fn clear_canceled_order(&self, symbol: &str, order_id: &str) -> Result<(), DexError> {
1563        let mut guard = self.canceled_results.write().await;
1564        if let Some(map) = guard.get_mut(symbol) {
1565            if map.remove(order_id).is_some() {
1566                return Ok(());
1567            }
1568        }
1569        Err(DexError::Other(format!(
1570            "canceled order {} for {} not found",
1571            order_id, symbol
1572        )))
1573    }
1574
1575    async fn clear_all_canceled_orders(&self) -> Result<(), DexError> {
1576        self.canceled_results.write().await.clear();
1577        Ok(())
1578    }
1579
1580    async fn create_order(
1581        &self,
1582        symbol: &str,
1583        size: Decimal,
1584        side: OrderSide,
1585        price: Option<Decimal>,
1586        spread: Option<i64>,
1587    ) -> Result<CreateOrderResponse, DexError> {
1588        let (price, time_in_force) = match price {
1589            Some(v) => {
1590                if spread.is_some() {
1591                    let map = self.dynamic_market_info.read().await;
1592                    let info = map
1593                        .get(symbol)
1594                        .ok_or_else(|| DexError::Other(format!("No market info for {}", symbol)))?;
1595                    let bid = info
1596                        .best_bid
1597                        .ok_or_else(|| DexError::Other("No best_bid".into()))?;
1598                    let ask = info
1599                        .best_ask
1600                        .ok_or_else(|| DexError::Other("No best_ask".into()))?;
1601                    let mid = (bid + ask) * Decimal::new(5, 1);
1602                    let tick = info
1603                        .min_tick
1604                        .ok_or_else(|| DexError::Other("No min_tick".into()))?;
1605                    let spread = Decimal::from(spread.unwrap());
1606                    log::debug!(
1607                        "bid = {}, min = {}, ask = {}, tick = {}, spread = {}",
1608                        bid,
1609                        mid,
1610                        ask,
1611                        tick,
1612                        spread
1613                    );
1614                    let calc = if side == OrderSide::Long {
1615                        mid - tick * spread
1616                    } else {
1617                        mid + tick * spread
1618                    };
1619                    (calc, "Alo")
1620                } else {
1621                    (v, "Alo")
1622                }
1623            }
1624            None => {
1625                let price = self.get_worst_price(symbol, &side).await?;
1626                (price, "Ioc")
1627            }
1628        };
1629
1630        let dynamic_market_info_guard = self.dynamic_market_info.read().await;
1631        let market_info = dynamic_market_info_guard
1632            .get(symbol)
1633            .ok_or_else(|| DexError::Other("Market info not found".to_string()))?;
1634        let min_tick = market_info
1635            .min_tick
1636            .ok_or_else(|| DexError::Other("Min tick not set for market".to_string()))?;
1637
1638        let rounded_price = Self::round_price(price, min_tick, side.clone());
1639        let rounded_size = self.floor_size(size, symbol);
1640
1641        log::info!(
1642            "[create_order] sym={} tif={} px={} size={} notional={} min_tick={} sz_decimals={}",
1643            symbol,
1644            time_in_force,
1645            rounded_price,
1646            rounded_size,
1647            rounded_price * rounded_size,
1648            min_tick,
1649            self.static_market_info
1650                .get(symbol)
1651                .map(|m| m.decimals)
1652                .unwrap_or(0),
1653        );
1654        let asset = resolve_coin(symbol, &self.spot_index_map);
1655
1656        let order = ClientOrderRequest {
1657            asset,
1658            is_buy: side == OrderSide::Long,
1659            reduce_only: false,
1660            limit_px: rounded_price
1661                .to_f64()
1662                .ok_or_else(|| DexError::Other("Conversion to f64 failed".to_string()))?,
1663            sz: rounded_size
1664                .to_f64()
1665                .ok_or_else(|| DexError::Other("Conversion to f64 failed".to_string()))?,
1666            cloid: None,
1667            order_type: ClientOrder::Limit(ClientLimit {
1668                tif: time_in_force.to_string(),
1669            }),
1670        };
1671
1672        let res = self.exchange_client.order(order, None).await.map_err(|e| {
1673            log::error!(
1674                "[create_order] order failed: symbol = {}, size = {}, error = {}",
1675                symbol,
1676                rounded_size,
1677                e
1678            );
1679            DexError::Other(e.to_string())
1680        })?;
1681
1682        let res = match res {
1683            ExchangeResponseStatus::Ok(exchange_response) => exchange_response,
1684            ExchangeResponseStatus::Err(e) => return Err(DexError::ServerResponse(e.to_string())),
1685        };
1686        let status = res.data.unwrap().statuses[0].clone();
1687        let order_id = match status {
1688            ExchangeDataStatus::Filled(order) => order.oid,
1689            ExchangeDataStatus::Resting(order) => order.oid,
1690            _ => {
1691                return Err(DexError::ServerResponse(
1692                    "Unknown ExchangeDataStaus".to_owned(),
1693                ))
1694            }
1695        };
1696
1697        Ok(CreateOrderResponse {
1698            order_id: order_id.to_string(),
1699            ordered_price: rounded_price,
1700            ordered_size: rounded_size,
1701        })
1702    }
1703
1704    async fn create_trigger_order(
1705        &self,
1706        symbol: &str,
1707        size: Decimal,
1708        side: OrderSide,
1709        trigger_px: Decimal,
1710        is_market: bool,
1711        tpsl: TpSl,
1712        _reduce_only: bool,
1713    ) -> Result<CreateOrderResponse, DexError> {
1714        // Resolve the exchange asset code
1715        let asset = resolve_coin(symbol, &self.spot_index_map);
1716
1717        // Convert Decimal to f64, failing if out of range
1718        let limit_px = trigger_px
1719            .to_f64()
1720            .ok_or_else(|| DexError::Other("Failed to convert trigger_px to f64".into()))?;
1721        let sz = size
1722            .to_f64()
1723            .ok_or_else(|| DexError::Other("Failed to convert size to f64".into()))?;
1724
1725        // Build the ClientOrderRequest with a Trigger type
1726        let request = ClientOrderRequest {
1727            asset,
1728            is_buy: side == OrderSide::Long,
1729            reduce_only: false,
1730            limit_px,
1731            sz,
1732            cloid: None,
1733            order_type: ClientOrder::Trigger(ClientTrigger {
1734                is_market,
1735                trigger_px: limit_px,
1736                // TpSl enum serialized lowercase: "tp" or "sl"
1737                tpsl: format!("{:?}", tpsl).to_lowercase(),
1738            }),
1739        };
1740
1741        // Send the order through the exchange client
1742        let resp_status = self
1743            .exchange_client
1744            .order(request, None)
1745            .await
1746            .map_err(|e| DexError::Other(format!("Order request failed: {}", e)))?;
1747
1748        // Unwrap the response status
1749        let exchange_response = match resp_status {
1750            ExchangeResponseStatus::Ok(x) => x,
1751            ExchangeResponseStatus::Err(e) => return Err(DexError::ServerResponse(e.to_string())),
1752        };
1753
1754        // Expect at least one status entry
1755        let status = exchange_response
1756            .data
1757            .unwrap()
1758            .statuses
1759            .into_iter()
1760            .next()
1761            .ok_or_else(|| DexError::Other("No order status returned".into()))?;
1762
1763        // Extract the order ID from either Filled or Resting
1764        let oid = match status {
1765            ExchangeDataStatus::Filled(o) => o.oid,
1766            ExchangeDataStatus::Resting(o) => o.oid,
1767            _ => {
1768                return Err(DexError::ServerResponse(
1769                    "Unrecognized exchange status".into(),
1770                ))
1771            }
1772        };
1773
1774        // Return the CreateOrderResponse
1775        Ok(CreateOrderResponse {
1776            order_id: oid.to_string(),
1777            ordered_price: trigger_px,
1778            ordered_size: size,
1779        })
1780    }
1781
1782    async fn cancel_order(&self, symbol: &str, order_id: &str) -> Result<(), DexError> {
1783        let asset = resolve_coin(symbol, &self.spot_index_map);
1784        let cancel = ClientCancelRequest {
1785            asset,
1786            oid: u64::from_str(order_id).unwrap_or_default(),
1787        };
1788
1789        self.exchange_client
1790            .cancel(cancel, None)
1791            .await
1792            .map_err(|e| DexError::Other(e.to_string()))?;
1793
1794        Ok(())
1795    }
1796
1797    async fn cancel_all_orders(&self, symbol: Option<String>) -> Result<(), DexError> {
1798        let open_orders = self.get_orders().await?;
1799        let order_ids: Vec<String> = open_orders
1800            .iter()
1801            .filter_map(|order| {
1802                let idx_opt = order.coin.strip_prefix('@').and_then(|s| s.parse::<usize>().ok());
1803                let external_sym = idx_opt
1804                    .and_then(|idx| self.spot_reverse_map.get(&idx).cloned())
1805                    .unwrap_or_else(|| format!("{}-USD", order.coin));
1806
1807                    log::debug!(
1808                        "cancel_all_orders: raw coin = {}, idx = {:?}, external_sym = {:?}, target = {:?}",
1809                        order.coin, idx_opt, external_sym, symbol
1810                    );
1811
1812                if symbol.as_deref().map_or(true, |s| s == &external_sym) {
1813                    Some(order.oid.to_string())
1814                } else {
1815                    None
1816                }
1817            })
1818            .collect();
1819        self.cancel_orders(symbol, order_ids).await
1820    }
1821
1822    async fn cancel_orders(
1823        &self,
1824        symbol: Option<String>,
1825        order_ids: Vec<String>,
1826    ) -> Result<(), DexError> {
1827        let open_orders = self.get_orders().await?;
1828        let mut cancels = Vec::new();
1829
1830        for order in open_orders {
1831            let idx_opt = order
1832                .coin
1833                .strip_prefix('@')
1834                .and_then(|s| s.parse::<usize>().ok());
1835            let external_sym = idx_opt
1836                .and_then(|idx| self.spot_reverse_map.get(&idx).cloned())
1837                .unwrap_or_else(|| format!("{}-USD", order.coin));
1838
1839            log::debug!(
1840                    "cancel_orders: raw coin = {}, idx = {:?}, external_sym = {:?}, requested_ids = {:?}",
1841                    order.coin, idx_opt, external_sym, order_ids
1842                );
1843
1844            if symbol.as_deref().map_or(true, |s| s == &external_sym)
1845                && order_ids.contains(&order.oid.to_string())
1846            {
1847                let asset = resolve_coin(&external_sym, &self.spot_index_map);
1848                cancels.push(ClientCancelRequest {
1849                    asset,
1850                    oid: order.oid,
1851                });
1852            }
1853        }
1854
1855        if !cancels.is_empty() {
1856            self.exchange_client
1857                .bulk_cancel(cancels, None)
1858                .await
1859                .map_err(|e| DexError::Other(e.to_string()))?;
1860        }
1861        Ok(())
1862    }
1863
1864    async fn close_all_positions(&self, symbol: Option<String>) -> Result<(), DexError> {
1865        let open_positions = self.get_positions().await?;
1866        for p in open_positions {
1867            let position = p.position;
1868            let idx_opt = position
1869                .coin
1870                .strip_prefix('@')
1871                .and_then(|s| s.parse::<usize>().ok());
1872            let external_sym = idx_opt
1873                .and_then(|idx| self.spot_reverse_map.get(&idx).cloned())
1874                .unwrap_or_else(|| format!("{}-USD", position.coin));
1875            if symbol.as_deref().map_or(true, |s| s == &external_sym) {
1876                let reversed_side = if position.szi.is_sign_negative() {
1877                    OrderSide::Long
1878                } else {
1879                    OrderSide::Short
1880                };
1881                let size = position.szi.abs();
1882                let _ = self
1883                    .create_order(&external_sym, size, reversed_side, None, None)
1884                    .await;
1885            }
1886        }
1887        Ok(())
1888    }
1889
1890    async fn clear_last_trades(&self, _symbol: &str) -> Result<(), DexError> {
1891        Ok(())
1892    }
1893
1894    async fn is_upcoming_maintenance(&self) -> bool {
1895        let info = self.maintenance.read().await;
1896        if let Some(start) = info.next_start {
1897            let now = Utc::now();
1898            if now < start && (start - now) <= ChronoDuration::hours(6) {
1899                return true;
1900            }
1901        }
1902        false
1903    }
1904
1905    async fn sign_evm_65b(&self, _message: &str) -> Result<String, DexError> {
1906        Err(DexError::Other(
1907            "65B EVM signature not supported for Hyperliquid".to_string(),
1908        ))
1909    }
1910
1911    async fn sign_evm_65b_with_eip191(&self, _message: &str) -> Result<String, DexError> {
1912        Err(DexError::Other(
1913            "65B EIP-191 signature not supported for Hyperliquid".to_string(),
1914        ))
1915    }
1916}
1917
1918impl HyperliquidConnector {
1919    async fn handle_request_with_action<T, U>(
1920        &self,
1921        request_url: String,
1922        action: &U,
1923    ) -> Result<T, DexError>
1924    where
1925        T: for<'de> Deserialize<'de>,
1926        U: Serialize + std::fmt::Debug + Clone,
1927    {
1928        let json_payload =
1929            serde_json::to_value(action).map_err(|e| DexError::Other(e.to_string()))?;
1930
1931        log::debug!("json_payload = {:?}", json_payload);
1932
1933        self.request
1934            .handle_request::<T, U>(
1935                HttpMethod::Post,
1936                request_url,
1937                &HashMap::new(),
1938                json_payload.to_string(),
1939            )
1940            .await
1941            .map_err(|e| DexError::Other(e.to_string()))
1942    }
1943
1944    async fn get_positions(
1945        &self,
1946    ) -> Result<Vec<HyperliquidRetriveUserPositionResponseBody>, DexError> {
1947        let request_url = "/info";
1948        let action = HyperliquidDefaultPayload {
1949            r#type: "clearinghouseState".to_owned(),
1950            user: Some(self.config.evm_wallet_address.clone()),
1951        };
1952        let res: HyperliquidRetriveUserPositionResponse = self
1953            .handle_request_with_action::<HyperliquidRetriveUserPositionResponse, HyperliquidDefaultPayload>(
1954                request_url.to_string(),
1955                &action,
1956            )
1957            .await?;
1958
1959        Ok(res.asset_positions)
1960    }
1961
1962    async fn get_orders(&self) -> Result<Vec<HyperliquidRetriveUserOpenOrder>, DexError> {
1963        let request_url = "/info";
1964        let action = HyperliquidDefaultPayload {
1965            r#type: "openOrders".to_owned(),
1966            user: Some(self.config.evm_wallet_address.clone()),
1967        };
1968        let res: Vec<HyperliquidRetriveUserOpenOrder> = self
1969            .handle_request_with_action::<Vec<HyperliquidRetriveUserOpenOrder>, HyperliquidDefaultPayload>(
1970                request_url.to_string(),
1971                &action,
1972            )
1973            .await?;
1974
1975        Ok(res)
1976    }
1977
1978    async fn retrive_market_metadata(&mut self) -> Result<(), DexError> {
1979        let request_url = "/info";
1980        let action = HyperliquidDefaultPayload {
1981            r#type: "meta".to_owned(),
1982            user: None,
1983        };
1984        let res = self
1985            .handle_request_with_action::<HyperliquidRetriveMarketMetadataResponse, HyperliquidDefaultPayload>(
1986                request_url.to_string(),
1987                &action,
1988            )
1989            .await?;
1990
1991        let mut static_market_info_update = HashMap::new();
1992        for metadata in res.universe.into_iter() {
1993            let market_id = format!("{}-USD", metadata.name);
1994            static_market_info_update.insert(
1995                market_id,
1996                StaticMarketInfo {
1997                    decimals: metadata.decimals,
1998                    _max_leverage: metadata.max_leverage,
1999                },
2000            );
2001        }
2002
2003        self.static_market_info = static_market_info_update;
2004
2005        Ok(())
2006    }
2007
2008    async fn get_worst_price(&self, symbol: &str, side: &OrderSide) -> Result<Decimal, DexError> {
2009        let market_price = self.get_market_price(symbol).await?;
2010
2011        let worst_price = slippage_price(market_price, *side == OrderSide::Long);
2012        Ok(worst_price)
2013    }
2014
2015    async fn get_market_price(&self, symbol: &str) -> Result<Decimal, DexError> {
2016        let market_info_guard = self.dynamic_market_info.read().await;
2017        match market_info_guard.get(symbol) {
2018            Some(v) => match v.market_price {
2019                Some(price) => Ok(price),
2020                None => Err(DexError::Other("Price is None".to_string())),
2021            },
2022            None => Err(DexError::Other("No price available".to_string())),
2023        }
2024    }
2025
2026    fn calculate_min_tick(price: Decimal, sz_decimals: u32, is_spot: bool) -> Decimal {
2027        log::trace!(
2028            "calculate_min_tick called: price={}, sz_decimals={}, is_spot={}",
2029            price,
2030            sz_decimals,
2031            is_spot
2032        );
2033
2034        let price_str = price.to_string();
2035        let integer_part = price_str.split('.').next().unwrap_or("");
2036        let integer_digits = if integer_part == "0" {
2037            0
2038        } else {
2039            integer_part.len()
2040        };
2041
2042        let scale_by_sig: u32 = if integer_digits >= 5 {
2043            0
2044        } else {
2045            (5 - integer_digits) as u32
2046        };
2047
2048        let max_decimals: u32 = if is_spot { 8u32 } else { 6u32 };
2049        let scale_by_dec: u32 = max_decimals.saturating_sub(sz_decimals);
2050        let scale: u32 = scale_by_sig.min(scale_by_dec);
2051
2052        log::trace!(
2053            "calculate_min_tick internals: integer_digits={}, scale_by_sig={}, max_decimals={}, scale_by_dec={}, scale={}",
2054            integer_digits,
2055            scale_by_sig,
2056            max_decimals,
2057            scale_by_dec,
2058            scale
2059        );
2060
2061        let min_tick = Decimal::new(1, scale);
2062
2063        log::trace!(
2064            "calculate_min_tick result: min_tick={}, (1e-{})",
2065            min_tick,
2066            scale
2067        );
2068
2069        min_tick
2070    }
2071
2072    fn round_price(price: Decimal, min_tick: Decimal, order_side: OrderSide) -> Decimal {
2073        if min_tick.is_zero() {
2074            log::error!("round_price: min_tick is zero");
2075            return price;
2076        }
2077
2078        match order_side {
2079            OrderSide::Long => (price / min_tick).floor() * min_tick,
2080            OrderSide::Short => (price / min_tick).ceil() * min_tick,
2081        }
2082    }
2083
2084    fn floor_size(&self, size: Decimal, symbol: &str) -> Decimal {
2085        let decimals = match self.static_market_info.get(symbol) {
2086            Some(v) => v.decimals,
2087            None => {
2088                log::error!("symbol meta is not available: {}", symbol);
2089                return size;
2090            }
2091        };
2092
2093        size.round_dp(decimals)
2094    }
2095
2096    fn extract_asset_name(symbol: &str) -> &str {
2097        symbol.split('-').next().unwrap_or(symbol)
2098    }
2099
2100    /// Wait until best_bid and best_ask are available for all configured symbols
2101    async fn wait_for_market_ready(&self, timeout_secs: u64) -> Result<(), DexError> {
2102        use tokio::time::{sleep, Instant};
2103
2104        let deadline = Instant::now() + Duration::from_secs(timeout_secs);
2105
2106        loop {
2107            let mut all_ready = true;
2108            {
2109                let map = self.dynamic_market_info.read().await;
2110                for symbol in &self.config.symbol_list {
2111                    if let Some(info) = map.get(symbol) {
2112                        log::warn!(
2113                            "[wait_for_market_ready] symbol = {}, best_bid = {:?}, best_ask = {:?}",
2114                            symbol,
2115                            info.best_bid,
2116                            info.best_ask
2117                        );
2118                        if info.best_bid.is_none() || info.best_ask.is_none() {
2119                            log::info!("Waiting for best_bid/best_ask for symbol: {}", symbol);
2120                            all_ready = false;
2121                            break;
2122                        }
2123                    } else {
2124                        log::info!("Market info not found yet for symbol: {}", symbol);
2125                        all_ready = false;
2126                        break;
2127                    }
2128                }
2129            }
2130
2131            if all_ready {
2132                log::info!("All symbols are market-ready.");
2133                return Ok(());
2134            }
2135
2136            if Instant::now() >= deadline {
2137                return Err(DexError::Other(
2138                    "Timed out waiting for market data".to_string(),
2139                ));
2140            }
2141
2142            sleep(Duration::from_millis(200)).await;
2143        }
2144    }
2145}
2146
2147fn resolve_coin(sym: &str, map: &HashMap<String, usize>) -> String {
2148    if sym.contains('/') {
2149        // ---- Spot ----
2150        match map.get(sym) {
2151            Some(idx) => format!("@{}", idx),
2152            None => {
2153                log::warn!("resolve_coin: {} is not in spot_index_map", sym);
2154                sym.to_string()
2155            }
2156        }
2157    } else if let Some(base) = sym.strip_suffix("-USD") {
2158        // ---- Perp ----
2159        base.to_string()
2160    } else {
2161        sym.to_string()
2162    }
2163}