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#[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], }
160
161#[allow(dead_code)]
162#[derive(Deserialize, Debug)]
163struct WsBook {
164 coin: String,
165 time: u64,
166 levels: [Vec<WsLevel>; 2], }
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 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, T: u64, s: String, i: String, o: Decimal, c: Decimal, h: Decimal, l: Decimal, v: Decimal, n: u64, }
227
228#[derive(Deserialize, Debug)]
229pub struct ActiveAssetCtxData {
230 pub coin: String, pub ctx: PerpsAssetCtx, }
233
234#[allow(dead_code, non_snake_case)]
235#[derive(Deserialize, Debug)]
236pub struct PerpsAssetCtx {
237 pub dayNtlVlm: Decimal, pub prevDayPx: Decimal, pub markPx: Decimal, pub midPx: Option<Decimal>, pub funding: Decimal, pub openInterest: Decimal, pub oraclePx: Decimal, }
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 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() } else {
944 format!("{}-USD", coin) };
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; }
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 let delta_vol = if cur_vol >= prev_vol {
1243 cur_vol - prev_vol
1244 } else {
1245 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 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 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 base.to_string()
1830 } else {
1831 sym.to_string()
1832 }
1833}