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