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