1use crate::rate_limiter::{AddressRateLimiter, RestRateLimiter};
2use crate::ws::manager::{
3 managed_stream, HyperliquidSubscription, HyperliquidWsManager, WsSendRateLimiter,
4};
5use crate::ws::{HyperliquidWsBook, HyperliquidWsInboundMessage};
6use async_trait::async_trait;
7use futures_util::{stream, StreamExt};
8use guilder_abstraction::{
9 self, AssetContext, BoxStream, Deposit, EcdsaSignature, ExternalSigner, Fill, FundingPayment,
10 L2Level, L2Snapshot, L2Update, Liquidation, OpenOrder, OrderPlacement, OrderSide, OrderType,
11 OrderUpdate, Position, PredictedFunding, TimeInForce, UserFill, Withdrawal,
12};
13use reqwest::Client;
14use rust_decimal::Decimal;
15use serde::Deserialize;
16use serde_json::Value;
17use std::collections::HashMap;
18use std::str::FromStr;
19use std::sync::{Arc, RwLock};
20const HYPERLIQUID_INFO_URL: &str = "https://api.hyperliquid.xyz/info";
21const HYPERLIQUID_EXCHANGE_URL: &str = "https://api.hyperliquid.xyz/exchange";
22
23async fn parse_response<T: for<'de> serde::Deserialize<'de>>(
24 resp: reqwest::Response,
25) -> Result<T, String> {
26 let status = resp.status();
27 let text = resp
28 .text()
29 .await
30 .map_err(|e| format!("failed to read response body (status {status}): {e}"))?;
31
32 if text.is_empty() {
33 return Err(format!(
34 "empty response body from Hyperliquid (HTTP {status})"
35 ));
36 }
37
38 serde_json::from_str(&text).map_err(|e| {
39 let snippet = if text.len() > 512 {
40 format!("{}... ({} bytes total)", &text[..256], text.len())
41 } else {
42 text.clone()
43 };
44 format!("deserialize error (HTTP {status}): {e}: {snippet}")
45 })
46}
47
48pub struct HyperliquidClient {
49 client: Client,
50 user_address: Option<String>,
51 private_key: Option<String>,
52 external_signer: Option<Arc<dyn ExternalSigner>>,
53 rest_limiter: Arc<RestRateLimiter>,
54 address_limiter: Arc<AddressRateLimiter>,
55 market_ws_manager: HyperliquidWsManager,
56 user_ws_managers: Arc<RwLock<HashMap<String, HyperliquidWsManager>>>,
57 ws_send_limiter: WsSendRateLimiter,
58}
59
60impl Default for HyperliquidClient {
61 fn default() -> Self {
62 Self::new()
63 }
64}
65
66impl HyperliquidClient {
67 pub fn new() -> Self {
68 let ws_send_limiter = WsSendRateLimiter::new();
69 HyperliquidClient {
70 client: Client::new(),
71 user_address: None,
72 private_key: None,
73 external_signer: None,
74 rest_limiter: Arc::new(RestRateLimiter::new()),
75 address_limiter: Arc::new(AddressRateLimiter::new()),
76 market_ws_manager: HyperliquidWsManager::new(None, ws_send_limiter.clone()),
77 user_ws_managers: Arc::new(RwLock::new(HashMap::new())),
78 ws_send_limiter,
79 }
80 }
81
82 pub fn with_auth(user_address: impl Into<String>, private_key: String) -> Self {
83 let ws_send_limiter = WsSendRateLimiter::new();
84 HyperliquidClient {
85 client: Client::new(),
86 user_address: Some(user_address.into()),
87 private_key: Some(private_key),
88 external_signer: None,
89 rest_limiter: Arc::new(RestRateLimiter::new()),
90 address_limiter: Arc::new(AddressRateLimiter::new()),
91 market_ws_manager: HyperliquidWsManager::new(None, ws_send_limiter.clone()),
92 user_ws_managers: Arc::new(RwLock::new(HashMap::new())),
93 ws_send_limiter,
94 }
95 }
96
97 pub fn with_external_signer(
103 user_address: impl Into<String>,
104 signer: Arc<dyn ExternalSigner>,
105 ) -> Self {
106 let ws_send_limiter = WsSendRateLimiter::new();
107 HyperliquidClient {
108 client: Client::new(),
109 user_address: Some(user_address.into()),
110 private_key: None,
111 external_signer: Some(signer),
112 rest_limiter: Arc::new(RestRateLimiter::new()),
113 address_limiter: Arc::new(AddressRateLimiter::new()),
114 market_ws_manager: HyperliquidWsManager::new(None, ws_send_limiter.clone()),
115 user_ws_managers: Arc::new(RwLock::new(HashMap::new())),
116 ws_send_limiter,
117 }
118 }
119
120 pub fn with_budgets(mut self, rest_weight: u32, addr_budget: u64) -> Self {
123 self.rest_limiter = Arc::new(RestRateLimiter::new_with_budget(rest_weight));
124 self.address_limiter = Arc::new(AddressRateLimiter::new_with_budget(addr_budget));
125 self
126 }
127
128 async fn info_post(
132 &self,
133 body: Value,
134 weight: u32,
135 call: &str,
136 ) -> Result<reqwest::Response, String> {
137 self.rest_limiter.acquire(weight).await.map_err(|e| {
138 format!(
139 "rate_limited: info_post ({call}) budget exhausted, retry_after_ms={}",
140 e.retry_after.as_millis()
141 )
142 })?;
143 self.client
144 .post(HYPERLIQUID_INFO_URL)
145 .json(&body)
146 .send()
147 .await
148 .map_err(|e| e.to_string())
149 }
150 fn require_user_address(&self) -> Result<String, String> {
151 self.user_address
152 .clone()
153 .ok_or_else(|| "user address required: use HyperliquidClient::with_auth".to_string())
154 }
155
156 fn require_private_key(&self) -> Result<Option<&str>, String> {
157 if self.external_signer.is_some() {
158 return Ok(None);
159 }
160 self.private_key
161 .as_deref()
162 .map(Some)
163 .ok_or_else(|| "private key or external signer required: use with_auth or with_external_signer".to_string())
164 }
165
166 async fn get_asset_index(&self, symbol: &str) -> Result<usize, String> {
167 let resp = self
169 .info_post(serde_json::json!({"type": "meta"}), 20, "get_asset_index")
170 .await?;
171 let meta: MetaResponse = parse_response(resp).await?;
172 meta.universe
173 .iter()
174 .position(|a| a.name == symbol)
175 .ok_or_else(|| format!("symbol {} not found", symbol))
176 }
177
178 async fn submit_signed_action(
179 &self,
180 action: Value,
181 vault_address: Option<&str>,
182 ) -> Result<Value, String> {
183 let private_key = self.require_private_key()?;
184 let nonce = std::time::SystemTime::now()
185 .duration_since(std::time::UNIX_EPOCH)
186 .unwrap()
187 .as_millis() as u64;
188
189 let (r, s, v) = sign_action(
190 private_key,
191 self.external_signer.as_ref(),
192 &action,
193 vault_address,
194 nonce,
195 )
196 .await?;
197
198 let payload = serde_json::json!({
199 "action": action,
200 "nonce": nonce,
201 "signature": {"r": r, "s": s, "v": v},
202 "vaultAddress": null,
203 "expiresAfter": null
204 });
205
206 self.rest_limiter.acquire(1).await.map_err(|e| {
208 format!(
209 "rate_limited: rest_weight exhausted, retry_after_ms={}",
210 e.retry_after.as_millis()
211 )
212 })?;
213 self.address_limiter.acquire(1, false).await.map_err(|e| {
214 format!(
215 "rate_limited: address quota exhausted, retry_after_ms={}",
216 e.retry_after.as_millis()
217 )
218 })?;
219
220 let resp = self
221 .client
222 .post(HYPERLIQUID_EXCHANGE_URL)
223 .json(&payload)
224 .send()
225 .await
226 .map_err(|e| e.to_string())?;
227
228 let status = resp.status();
229 if !status.is_success() {
230 let text = resp.text().await.map_err(|e| e.to_string())?;
231 return Err(format!("HTTP {status}: {text}"));
232 }
233
234 let body: Value = parse_response(resp).await?;
235 if body["status"].as_str() == Some("err") {
236 return Err(body["response"]
237 .as_str()
238 .unwrap_or("unknown error")
239 .to_string());
240 }
241 Ok(body)
242 }
243}
244
245#[derive(Deserialize)]
248struct MetaResponse {
249 universe: Vec<AssetInfo>,
250}
251
252#[derive(Deserialize)]
253struct AssetInfo {
254 name: String,
255 #[serde(rename = "szDecimals")]
256 sz_decimals: i32,
257 #[serde(rename = "isDelisted", default)]
258 is_delisted: bool,
259}
260
261type MetaAndAssetCtxsResponse = (MetaResponse, Vec<RestAssetCtx>);
262
263#[derive(Deserialize)]
264#[serde(rename_all = "camelCase")]
265#[allow(dead_code)]
266struct RestAssetCtx {
267 open_interest: String,
268 funding: String,
269 mark_px: String,
270 day_ntl_vlm: String,
271 mid_px: Option<String>,
272 oracle_px: Option<String>,
273 premium: Option<String>,
274 prev_day_px: Option<String>,
275}
276
277#[derive(Deserialize)]
278#[serde(rename_all = "camelCase")]
279#[allow(dead_code)]
280struct ClearinghouseStateResponse {
281 margin_summary: MarginSummary,
282 asset_positions: Vec<AssetPosition>,
283}
284
285#[derive(Deserialize)]
287#[serde(rename_all = "camelCase")]
288#[allow(dead_code)]
289struct MarginSummary {
290 account_value: String,
291 #[serde(default)]
292 total_ntl_pos: Option<String>,
293 #[serde(default)]
294 total_raw_usd: Option<String>,
295 #[serde(default)]
296 total_margin_used: Option<String>,
297}
298
299#[derive(Deserialize)]
300struct AssetPosition {
301 position: PositionDetail,
302}
303
304#[derive(Deserialize)]
305#[serde(rename_all = "camelCase")]
306struct PositionDetail {
307 coin: String,
308 szi: String,
310 entry_px: Option<String>,
311}
312
313#[derive(Deserialize)]
314#[serde(rename_all = "camelCase")]
315struct RestOpenOrder {
316 coin: String,
317 side: String,
318 limit_px: String,
319 sz: String,
320 oid: i64,
321 orig_sz: String,
322 cloid: Option<String>,
323}
324
325type PredictedFundingsResponse = Vec<(String, Vec<(String, Option<PredictedFundingEntry>)>)>;
328
329#[derive(Deserialize)]
330#[serde(rename_all = "camelCase")]
331struct PredictedFundingEntry {
332 funding_rate: String,
333 next_funding_time: i64,
334}
335
336fn parse_decimal(s: &str) -> Option<Decimal> {
339 Decimal::from_str(s).ok()
340}
341
342fn keccak256(data: &[u8]) -> [u8; 32] {
343 use sha3::{Digest, Keccak256};
344 Keccak256::digest(data).into()
345}
346
347fn hyperliquid_domain_separator() -> [u8; 32] {
349 let type_hash = keccak256(
350 b"EIP712Domain(string name,string version,uint256 chainId,address verifyingContract)",
351 );
352 let name_hash = keccak256(b"Exchange");
353 let version_hash = keccak256(b"1");
354 let mut chain_id = [0u8; 32];
355 chain_id[28..32].copy_from_slice(&1337u32.to_be_bytes());
356 let verifying_contract = [0u8; 32];
357
358 let mut data = [0u8; 160];
359 data[..32].copy_from_slice(&type_hash);
360 data[32..64].copy_from_slice(&name_hash);
361 data[64..96].copy_from_slice(&version_hash);
362 data[96..128].copy_from_slice(&chain_id);
363 data[128..160].copy_from_slice(&verifying_contract);
364 keccak256(&data)
365}
366
367fn value_to_msgpack(val: &Value) -> Vec<u8> {
370 match val {
371 Value::Null => vec![0xc0],
372 Value::Bool(true) => vec![0xc3],
373 Value::Bool(false) => vec![0xc2],
374 Value::Number(n) => {
375 if let Some(i) = n.as_i64() {
376 if i >= 0 {
377 if i <= 127 {
378 vec![i as u8]
379 } else if i <= 255 {
380 vec![0xcc, i as u8]
381 } else if i <= 65535 {
382 let mut buf = vec![0xcd];
383 buf.extend_from_slice(&(i as u16).to_be_bytes());
384 buf
385 } else if i <= 4294967295 {
386 let mut buf = vec![0xce];
387 buf.extend_from_slice(&(i as u32).to_be_bytes());
388 buf
389 } else {
390 let mut buf = vec![0xcf];
391 buf.extend_from_slice(&(i as u64).to_be_bytes());
392 buf
393 }
394 } else if i >= -32 {
395 vec![0xe0 | (i as u8)]
396 } else if i >= -128 {
397 vec![0xd0, i as i8 as u8]
398 } else if i >= -32768 {
399 let mut buf = vec![0xd1];
400 buf.extend_from_slice(&(i as i16).to_be_bytes());
401 buf
402 } else if i >= -2147483648 {
403 let mut buf = vec![0xd2];
404 buf.extend_from_slice(&(i as i32).to_be_bytes());
405 buf
406 } else {
407 let mut buf = vec![0xd3];
408 buf.extend_from_slice(&i.to_be_bytes());
409 buf
410 }
411 } else if let Some(f) = n.as_f64() {
412 let mut buf = vec![0xcb];
413 buf.extend_from_slice(&f.to_be_bytes());
414 buf
415 } else {
416 let u = n.as_u64().unwrap();
417 if u <= 127 {
418 vec![u as u8]
419 } else if u <= 255 {
420 vec![0xcc, u as u8]
421 } else if u <= 65535 {
422 let mut buf = vec![0xcd];
423 buf.extend_from_slice(&(u as u16).to_be_bytes());
424 buf
425 } else if u <= 4294967295 {
426 let mut buf = vec![0xce];
427 buf.extend_from_slice(&(u as u32).to_be_bytes());
428 buf
429 } else {
430 let mut buf = vec![0xcf];
431 buf.extend_from_slice(&u.to_be_bytes());
432 buf
433 }
434 }
435 }
436 Value::String(s) => {
437 let bytes = s.as_bytes();
438 let len = bytes.len();
439 let mut buf = Vec::new();
440 if len <= 31 {
441 buf.push(0xa0 | (len as u8));
442 } else if len <= 255 {
443 buf.push(0xd9);
444 buf.push(len as u8);
445 } else if len <= 65535 {
446 buf.push(0xda);
447 buf.extend_from_slice(&(len as u16).to_be_bytes());
448 } else {
449 buf.push(0xdb);
450 buf.extend_from_slice(&(len as u32).to_be_bytes());
451 }
452 buf.extend_from_slice(bytes);
453 buf
454 }
455 Value::Array(arr) => {
456 let len = arr.len();
457 let mut buf = Vec::new();
458 if len <= 15 {
459 buf.push(0x90 | (len as u8));
460 } else if len <= 65535 {
461 buf.push(0xdc);
462 buf.extend_from_slice(&(len as u16).to_be_bytes());
463 } else {
464 buf.push(0xdd);
465 buf.extend_from_slice(&(len as u32).to_be_bytes());
466 }
467 for item in arr {
468 buf.extend_from_slice(&value_to_msgpack(item));
469 }
470 buf
471 }
472 Value::Object(map) => {
473 let len = map.len();
474 let mut buf = Vec::new();
475 if len <= 15 {
476 buf.push(0x80 | (len as u8));
477 } else if len <= 65535 {
478 buf.push(0xde);
479 buf.extend_from_slice(&(len as u16).to_be_bytes());
480 } else {
481 buf.push(0xdf);
482 buf.extend_from_slice(&(len as u32).to_be_bytes());
483 }
484 for (key, value) in map {
485 buf.extend_from_slice(&value_to_msgpack(&Value::String(key.clone())));
486 buf.extend_from_slice(&value_to_msgpack(value));
487 }
488 buf
489 }
490 }
491}
492
493fn action_to_canonical_msgpack(action: &Value) -> Result<Vec<u8>, String> {
496 Ok(value_to_msgpack(action))
497}
498
499fn build_order_msgpack(
502 asset_idx: usize,
503 is_buy: bool,
504 price: &str,
505 size: &str,
506 reduce_only: bool,
507 order_kind: &str,
508 tif: &[u8],
509 cloid: Option<&str>,
510) -> Vec<u8> {
511 let field_count = if cloid.is_some() { 7 } else { 6 };
512 let mut buf = Vec::new();
513 buf.push(0x80 | (field_count as u8)); buf.extend_from_slice(&value_to_msgpack(&Value::String("a".to_string())));
517 buf.extend_from_slice(&value_to_msgpack(&Value::Number(serde_json::Number::from(
518 asset_idx,
519 ))));
520
521 buf.extend_from_slice(&value_to_msgpack(&Value::String("b".to_string())));
523 buf.push(if is_buy { 0xc3 } else { 0xc2 });
524
525 buf.extend_from_slice(&value_to_msgpack(&Value::String("p".to_string())));
527 buf.extend_from_slice(&value_to_msgpack(&Value::String(price.to_string())));
528
529 buf.extend_from_slice(&value_to_msgpack(&Value::String("s".to_string())));
531 buf.extend_from_slice(&value_to_msgpack(&Value::String(size.to_string())));
532
533 buf.extend_from_slice(&value_to_msgpack(&Value::String("r".to_string())));
535 buf.push(if reduce_only { 0xc3 } else { 0xc2 });
536
537 buf.extend_from_slice(&value_to_msgpack(&Value::String("t".to_string())));
539 buf.push(0x81);
541 buf.extend_from_slice(&value_to_msgpack(&Value::String(order_kind.to_string())));
542 buf.push(0x81);
544 buf.extend_from_slice(&value_to_msgpack(&Value::String("tif".to_string())));
545 buf.extend_from_slice(&value_to_msgpack(&Value::String(
546 String::from_utf8_lossy(tif).to_string(),
547 )));
548
549 if let Some(c) = cloid {
551 buf.extend_from_slice(&value_to_msgpack(&Value::String("c".to_string())));
552 buf.extend_from_slice(&value_to_msgpack(&Value::String(c.to_string())));
553 }
554
555 buf
556}
557
558fn build_trigger_order_msgpack(
561 asset_idx: usize,
562 is_buy: bool,
563 price: &str,
564 size: &str,
565 reduce_only: bool,
566 trigger_px: &str,
567 is_market: bool,
568 tpsl: &str,
569 cloid: Option<&str>,
570) -> Vec<u8> {
571 let field_count = if cloid.is_some() { 7 } else { 6 };
572 let mut buf = Vec::new();
573 buf.push(0x80 | (field_count as u8)); buf.extend_from_slice(&value_to_msgpack(&Value::String("a".to_string())));
577 buf.extend_from_slice(&value_to_msgpack(&Value::Number(serde_json::Number::from(
578 asset_idx,
579 ))));
580
581 buf.extend_from_slice(&value_to_msgpack(&Value::String("b".to_string())));
583 buf.push(if is_buy { 0xc3 } else { 0xc2 });
584
585 buf.extend_from_slice(&value_to_msgpack(&Value::String("p".to_string())));
587 buf.extend_from_slice(&value_to_msgpack(&Value::String(price.to_string())));
588
589 buf.extend_from_slice(&value_to_msgpack(&Value::String("s".to_string())));
591 buf.extend_from_slice(&value_to_msgpack(&Value::String(size.to_string())));
592
593 buf.extend_from_slice(&value_to_msgpack(&Value::String("r".to_string())));
595 buf.push(if reduce_only { 0xc3 } else { 0xc2 });
596
597 buf.extend_from_slice(&value_to_msgpack(&Value::String("t".to_string())));
599 buf.push(0x81); buf.extend_from_slice(&value_to_msgpack(&Value::String("trigger".to_string())));
601 buf.push(0x83);
603 buf.extend_from_slice(&value_to_msgpack(&Value::String("isMarket".to_string())));
604 buf.push(if is_market { 0xc3 } else { 0xc2 });
605 buf.extend_from_slice(&value_to_msgpack(&Value::String("triggerPx".to_string())));
606 buf.extend_from_slice(&value_to_msgpack(&Value::String(trigger_px.to_string())));
607 buf.extend_from_slice(&value_to_msgpack(&Value::String("tpsl".to_string())));
608 buf.extend_from_slice(&value_to_msgpack(&Value::String(tpsl.to_string())));
609
610 if let Some(c) = cloid {
612 buf.extend_from_slice(&value_to_msgpack(&Value::String("c".to_string())));
613 buf.extend_from_slice(&value_to_msgpack(&Value::String(c.to_string())));
614 }
615
616 buf
617}
618
619fn compute_eip712_digest(
624 msgpack: &[u8],
625 nonce: u64,
626 vault_address: Option<&str>,
627) -> Result<[u8; 32], String> {
628 let mut data = msgpack.to_vec();
629 data.extend_from_slice(&nonce.to_be_bytes());
630 match vault_address {
631 None => data.push(0u8),
632 Some(addr) => {
633 data.push(1u8);
634 let addr_bytes = hex::decode(addr.trim_start_matches("0x"))
635 .map_err(|e| format!("invalid vault address: {}", e))?;
636 data.extend_from_slice(&addr_bytes);
637 }
638 }
639
640 let connection_id = keccak256(&data);
641 let agent_type_hash = keccak256(b"Agent(string source,bytes32 connectionId)");
642 let source_hash = keccak256(b"a");
643 let mut struct_data = [0u8; 96];
644 struct_data[..32].copy_from_slice(&agent_type_hash);
645 struct_data[32..64].copy_from_slice(&source_hash);
646 struct_data[64..96].copy_from_slice(&connection_id);
647 let struct_hash = keccak256(&struct_data);
648
649 let domain_sep = hyperliquid_domain_separator();
650 let mut final_data = Vec::with_capacity(66);
651 final_data.extend_from_slice(b"\x19\x01");
652 final_data.extend_from_slice(&domain_sep);
653 final_data.extend_from_slice(&struct_hash);
654 Ok(keccak256(&final_data))
655}
656
657fn format_ecdsa_signature(sig: &EcdsaSignature) -> (String, String, u8) {
659 let r = format!("0x{}", hex::encode(sig.r));
660 let s = format!("0x{}", hex::encode(sig.s));
661 let v = 27u8 + sig.v;
662 (r, s, v)
663}
664
665fn sign_digest_with_key(
667 digest: &[u8; 32],
668 private_key: &str,
669) -> Result<(String, String, u8), String> {
670 use k256::ecdsa::SigningKey;
671
672 let key_bytes = hex::decode(private_key.trim_start_matches("0x"))
673 .map_err(|e| format!("invalid private key: {}", e))?;
674 let signing_key =
675 SigningKey::from_bytes(key_bytes.as_slice().into()).map_err(|e| e.to_string())?;
676 let (sig, recovery_id) = signing_key
677 .sign_prehash_recoverable(digest)
678 .map_err(|e| e.to_string())?;
679
680 let sig_bytes = sig.to_bytes();
681 let r = format!("0x{}", hex::encode(&sig_bytes[..32]));
682 let s = format!("0x{}", hex::encode(&sig_bytes[32..64]));
683 let v = 27u8 + recovery_id.to_byte();
684
685 Ok((r, s, v))
686}
687
688async fn sign_with_msgpack(
691 msgpack: &[u8],
692 private_key: Option<&str>,
693 external_signer: Option<&Arc<dyn ExternalSigner>>,
694 nonce: u64,
695 vault_address: Option<&str>,
696) -> Result<(String, String, u8), String> {
697 let digest = compute_eip712_digest(msgpack, nonce, vault_address)?;
698
699 if let Some(signer) = external_signer {
700 let sig = signer.sign_prehash(&digest).await?;
701 return Ok(format_ecdsa_signature(&sig));
702 }
703
704 let key = private_key.ok_or_else(|| {
705 "no signing method available: provide private_key or external_signer".to_string()
706 })?;
707 sign_digest_with_key(&digest, key)
708}
709
710async fn sign_action(
714 private_key: Option<&str>,
715 external_signer: Option<&Arc<dyn ExternalSigner>>,
716 action: &Value,
717 vault_address: Option<&str>,
718 nonce: u64,
719) -> Result<(String, String, u8), String> {
720 let msgpack_bytes = action_to_canonical_msgpack(action)?;
721 let digest = compute_eip712_digest(&msgpack_bytes, nonce, vault_address)?;
722
723 if let Some(signer) = external_signer {
724 let sig = signer.sign_prehash(&digest).await?;
725 return Ok(format_ecdsa_signature(&sig));
726 }
727
728 let key = private_key.ok_or_else(|| {
729 "no signing method available: provide private_key or external_signer".to_string()
730 })?;
731 sign_digest_with_key(&digest, key)
732}
733
734#[async_trait]
737impl guilder_abstraction::TestServer for HyperliquidClient {
738 async fn ping(&self) -> Result<bool, String> {
740 self.info_post(serde_json::json!({"type": "allMids"}), 2, "ping")
742 .await
743 .map(|r| r.status().is_success())
744 }
745
746 async fn get_server_time(&self) -> Result<i64, String> {
748 Ok(std::time::SystemTime::now()
749 .duration_since(std::time::UNIX_EPOCH)
750 .map(|d| d.as_millis() as i64)
751 .unwrap_or(0))
752 }
753}
754
755#[async_trait]
756impl guilder_abstraction::GetMarketData for HyperliquidClient {
757 async fn get_symbol(&self) -> Result<Vec<String>, String> {
759 let resp = self
761 .info_post(serde_json::json!({"type": "meta"}), 20, "get_symbol")
762 .await?;
763 parse_response::<MetaResponse>(resp)
764 .await
765 .map(|r| r.universe.into_iter().map(|a| a.name).collect())
766 }
767
768 async fn get_open_interest(&self, symbol: String) -> Result<Decimal, String> {
770 let resp = self
772 .info_post(
773 serde_json::json!({"type": "metaAndAssetCtxs"}),
774 20,
775 "get_open_interest",
776 )
777 .await?;
778 let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
779 .await?
780 .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
781 meta.universe
782 .iter()
783 .position(|a| a.name == symbol)
784 .and_then(|i| ctxs.get(i))
785 .and_then(|ctx| parse_decimal(&ctx.open_interest))
786 .ok_or_else(|| format!("symbol {} not found", symbol))
787 }
788
789 async fn get_asset_context(&self, symbol: String) -> Result<AssetContext, String> {
791 let resp = self
793 .info_post(
794 serde_json::json!({"type": "metaAndAssetCtxs"}),
795 20,
796 "get_asset_context",
797 )
798 .await?;
799 let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
800 .await?
801 .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
802 let idx = meta
803 .universe
804 .iter()
805 .position(|a| a.name == symbol)
806 .ok_or_else(|| format!("symbol {} not found", symbol))?;
807 let ctx = ctxs
808 .get(idx)
809 .ok_or_else(|| format!("symbol {} not found", symbol))?;
810 Ok(AssetContext {
811 symbol,
812 open_interest: parse_decimal(&ctx.open_interest).ok_or("invalid open_interest")?,
813 funding_rate: parse_decimal(&ctx.funding).ok_or("invalid funding")?,
814 mark_price: parse_decimal(&ctx.mark_px).ok_or("invalid mark_px")?,
815 day_volume: parse_decimal(&ctx.day_ntl_vlm).ok_or("invalid day_ntl_vlm")?,
816 mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
817 oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
818 premium: ctx.premium.as_deref().and_then(parse_decimal),
819 prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
820 sz_decimals: meta.universe.get(idx).map(|a| a.sz_decimals).unwrap_or(0),
821 })
822 }
823
824 async fn get_all_asset_contexts(&self) -> Result<Vec<AssetContext>, String> {
827 let resp = self
829 .info_post(
830 serde_json::json!({"type": "metaAndAssetCtxs"}),
831 20,
832 "get_all_asset_contexts",
833 )
834 .await?;
835 let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
836 .await?
837 .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
838 let mut result = Vec::with_capacity(meta.universe.len());
839 for (asset, ctx) in meta.universe.iter().zip(ctxs.iter()) {
840 let Some(open_interest) = parse_decimal(&ctx.open_interest) else {
841 continue;
842 };
843 let Some(funding_rate) = parse_decimal(&ctx.funding) else {
844 continue;
845 };
846 let Some(mark_price) = parse_decimal(&ctx.mark_px) else {
847 continue;
848 };
849 let Some(day_volume) = parse_decimal(&ctx.day_ntl_vlm) else {
850 continue;
851 };
852 result.push(AssetContext {
853 symbol: asset.name.clone(),
854 open_interest,
855 funding_rate,
856 mark_price,
857 day_volume,
858 mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
859 oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
860 premium: ctx.premium.as_deref().and_then(parse_decimal),
861 prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
862 sz_decimals: asset.sz_decimals,
863 });
864 }
865 Ok(result)
866 }
867
868 async fn get_sz_decimals(&self, symbol: String) -> Result<i32, String> {
870 let all = self.get_all_sz_decimals().await?;
871 all.get(&symbol)
872 .copied()
873 .ok_or_else(|| format!("symbol {} not found", symbol))
874 }
875
876 async fn get_all_sz_decimals(&self) -> Result<HashMap<String, i32>, String> {
878 let resp = self
880 .info_post(
881 serde_json::json!({"type": "metaAndAssetCtxs"}),
882 20,
883 "get_all_sz_decimals",
884 )
885 .await?;
886 let (meta, _) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
887 .await?
888 .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
889 Ok(meta
890 .universe
891 .into_iter()
892 .map(|a| (a.name, a.sz_decimals))
893 .collect())
894 }
895
896 async fn get_l2_orderbook(&self, symbol: String) -> Result<L2Snapshot, String> {
898 let resp = self
900 .info_post(
901 serde_json::json!({"type": "l2Book", "coin": symbol}),
902 2,
903 "get_l2_orderbook",
904 )
905 .await?;
906 let book: Option<HyperliquidWsBook> = parse_response(resp).await?;
907 let book = match book {
908 Some(b) => b,
909 None => {
910 return Ok(L2Snapshot {
911 symbol,
912 bids: vec![],
913 asks: vec![],
914 sequence: 0,
915 })
916 }
917 };
918 let bids = book
919 .levels
920 .first()
921 .into_iter()
922 .flatten()
923 .filter_map(|level| {
924 Some(L2Level {
925 price: parse_decimal(&level.px)?,
926 volume: parse_decimal(&level.sz)?,
927 })
928 })
929 .collect();
930 let asks = book
931 .levels
932 .get(1)
933 .into_iter()
934 .flatten()
935 .filter_map(|level| {
936 Some(L2Level {
937 price: parse_decimal(&level.px)?,
938 volume: parse_decimal(&level.sz)?,
939 })
940 })
941 .collect();
942 Ok(L2Snapshot {
943 symbol: book.coin,
944 bids,
945 asks,
946 sequence: book.time,
947 })
948 }
949
950 async fn get_price(&self, symbol: String) -> Result<Decimal, String> {
952 let resp = self
954 .info_post(serde_json::json!({"type": "allMids"}), 2, "get_price")
955 .await?;
956 parse_response::<HashMap<String, String>>(resp)
957 .await?
958 .get(&symbol)
959 .and_then(|s| parse_decimal(s))
960 .ok_or_else(|| format!("symbol {} not found", symbol))
961 }
962
963 async fn get_predicted_fundings(&self) -> Result<Vec<PredictedFunding>, String> {
966 let resp = self
968 .info_post(
969 serde_json::json!({"type": "predictedFundings"}),
970 20,
971 "get_predicted_fundings",
972 )
973 .await?;
974 let data: PredictedFundingsResponse = parse_response(resp).await?;
975 let mut result = Vec::new();
976 for (symbol, venues) in data {
977 for (venue, entry) in venues {
978 let Some(entry) = entry else { continue };
979 if let Some(funding_rate) = parse_decimal(&entry.funding_rate) {
980 result.push(PredictedFunding {
981 symbol: symbol.clone(),
982 venue,
983 funding_rate,
984 next_funding_time_ms: entry.next_funding_time,
985 });
986 }
987 }
988 }
989 Ok(result)
990 }
991}
992
993#[async_trait]
994impl guilder_abstraction::ManageOrder for HyperliquidClient {
995 async fn place_order(
1006 &self,
1007 symbol: String,
1008 side: OrderSide,
1009 price: Decimal,
1010 volume: Decimal,
1011 order_type: OrderType,
1012 time_in_force: TimeInForce,
1013 trigger_price: Option<Decimal>,
1014 reduce_only: bool,
1015 cloid: Option<String>,
1016 ) -> Result<OrderPlacement, String> {
1017 let asset_idx = self.get_asset_index(&symbol).await?;
1019 let is_buy = matches!(side, OrderSide::Buy);
1020
1021 let tif_str = match time_in_force {
1022 TimeInForce::Gtc => "Gtc",
1023 TimeInForce::Ioc => "Ioc",
1024 TimeInForce::Fok => "Fok",
1025 TimeInForce::Alo => "Alo",
1026 };
1027
1028 let cloid_hex = cloid.clone();
1029
1030 let is_trigger = matches!(order_type, OrderType::TakeProfit | OrderType::StopLoss);
1032
1033 let (order_msgpack, order_type_json) = if is_trigger {
1034 let trigger_px = trigger_price
1036 .ok_or_else(|| format!("{:?} order requires trigger_price to be set", order_type))?
1037 .normalize()
1038 .to_string();
1039
1040 let tpsl = match order_type {
1041 OrderType::TakeProfit => "tp",
1042 OrderType::StopLoss => "sl",
1043 _ => unreachable!(),
1044 };
1045
1046 let is_market = matches!(time_in_force, TimeInForce::Ioc);
1049
1050 let price_str = if is_market {
1053 trigger_px.clone()
1054 } else {
1055 price.normalize().to_string()
1056 };
1057
1058 let msgpack = build_trigger_order_msgpack(
1059 asset_idx,
1060 is_buy,
1061 &price_str,
1062 &volume.normalize().to_string(),
1063 reduce_only,
1064 &trigger_px,
1065 is_market,
1066 tpsl,
1067 cloid_hex.as_deref(),
1068 );
1069
1070 let json_type = if is_market {
1071 format!(
1072 r#"{{"trigger":{{"isMarket":true,"triggerPx":"{}","tpsl":"{}"}}}}"#,
1073 trigger_px, tpsl
1074 )
1075 } else {
1076 format!(
1077 r#"{{"trigger":{{"isMarket":false,"triggerPx":"{}","tpsl":"{}"}}}}"#,
1078 trigger_px, tpsl
1079 )
1080 };
1081
1082 (msgpack, json_type)
1083 } else {
1084 let (order_kind, tif_bytes) = match order_type {
1086 OrderType::Limit => ("limit", tif_str.as_bytes()),
1087 OrderType::Market => ("limit", b"Ioc".as_slice()),
1088 _ => unreachable!(),
1089 };
1090
1091 let price_str = price.normalize().to_string();
1092 let size_str = volume.normalize().to_string();
1093
1094 let msgpack = build_order_msgpack(
1095 asset_idx,
1096 is_buy,
1097 &price_str,
1098 &size_str,
1099 reduce_only,
1100 order_kind,
1101 tif_bytes,
1102 cloid_hex.as_deref(),
1103 );
1104
1105 let json_type = match order_type {
1106 OrderType::Limit => format!(r#"{{"limit":{{"tif":"{tif_str}"}}}}"#),
1107 OrderType::Market => r#"{"limit":{"tif":"Ioc"}}"#.to_string(),
1108 _ => unreachable!(),
1109 };
1110
1111 (msgpack, json_type)
1112 };
1113
1114 let mut action_msgpack = Vec::new();
1117 action_msgpack.push(0x83); action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("type".to_string())));
1119 action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("order".to_string())));
1120 action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("orders".to_string())));
1121 action_msgpack.push(0x91); action_msgpack.extend_from_slice(&order_msgpack);
1123 action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("grouping".to_string())));
1124 action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("na".to_string())));
1125
1126 let cloid_json = if let Some(ref c) = cloid_hex {
1127 format!(r#","c":"{c}""#)
1128 } else {
1129 String::new()
1130 };
1131
1132 let reduce_json = if reduce_only { "true" } else { "false" };
1133
1134 let price_for_json = if is_trigger && matches!(time_in_force, TimeInForce::Ioc) {
1135 if let Some(ref tp) = trigger_price {
1137 tp.normalize().to_string()
1138 } else {
1139 price.normalize().to_string()
1140 }
1141 } else {
1142 price.normalize().to_string()
1143 };
1144
1145 let action_json_str = format!(
1146 r#"{{"type":"order","orders":[{{"a":{asset_idx},"b":{is_buy},"p":"{price}","s":"{size}","r":{reduce_json},"t":{order_type_json}{cloid_json}}}],"grouping":"na"}}"#,
1147 price = price_for_json,
1148 size = volume.normalize().to_string(),
1149 );
1150
1151 let private_key = self.require_private_key()?;
1153 let nonce = std::time::SystemTime::now()
1154 .duration_since(std::time::UNIX_EPOCH)
1155 .unwrap()
1156 .as_millis() as u64;
1157
1158 let (r, s, v) = sign_with_msgpack(
1159 &action_msgpack,
1160 private_key,
1161 self.external_signer.as_ref(),
1162 nonce,
1163 None,
1164 )
1165 .await?;
1166
1167 let payload_str = format!(
1168 r#"{{"action":{},"nonce":{},"signature":{{"r":"{}","s":"{}","v":{}}},"vaultAddress":null,"expiresAfter":null}}"#,
1169 action_json_str, nonce, r, s, v
1170 );
1171
1172 self.rest_limiter.acquire(1).await.map_err(|e| {
1173 format!(
1174 "rate_limited: rest_weight exhausted, retry_after_ms={}",
1175 e.retry_after.as_millis()
1176 )
1177 })?;
1178 self.address_limiter.acquire(1, false).await.map_err(|e| {
1179 format!(
1180 "rate_limited: address quota exhausted, retry_after_ms={}",
1181 e.retry_after.as_millis()
1182 )
1183 })?;
1184
1185 let resp = self
1186 .client
1187 .post(HYPERLIQUID_EXCHANGE_URL)
1188 .header("Content-Type", "application/json")
1189 .body(payload_str)
1190 .send()
1191 .await
1192 .map_err(|e| e.to_string())?;
1193
1194 let status = resp.status();
1195 if !status.is_success() {
1196 let text = resp.text().await.map_err(|e| e.to_string())?;
1197 return Err(format!("HTTP {status}: {text}"));
1198 }
1199
1200 let body: Value = parse_response(resp).await?;
1201 if body["status"].as_str() == Some("err") {
1202 return Err(body["response"]
1203 .as_str()
1204 .unwrap_or("unknown error")
1205 .to_string());
1206 }
1207 let statuses = &body["response"]["data"]["statuses"][0];
1208
1209 let (oid, returned_cloid, timestamp_ms) = if let Some(resting) = statuses.get("resting") {
1210 let oid = resting["oid"]
1211 .as_i64()
1212 .ok_or_else(|| format!("resting status missing oid: {}", body))?;
1213 let returned_cloid = resting["cloid"].as_str().map(|s: &str| s.to_string());
1214 let ts = std::time::SystemTime::now()
1216 .duration_since(std::time::UNIX_EPOCH)
1217 .unwrap()
1218 .as_millis() as i64;
1219 (oid, returned_cloid, ts)
1220 } else if let Some(filled) = statuses.get("filled") {
1221 let oid = filled["oid"]
1222 .as_i64()
1223 .ok_or_else(|| format!("filled status missing oid: {}", body))?;
1224 let ts = std::time::SystemTime::now()
1226 .duration_since(std::time::UNIX_EPOCH)
1227 .unwrap()
1228 .as_millis() as i64;
1229 (oid, None, ts)
1230 } else if let Some(error) = statuses.get("error") {
1231 return Err(error
1232 .as_str()
1233 .unwrap_or("order rejected with unknown error")
1234 .to_string());
1235 } else {
1236 return Err(format!("unexpected order status: {}", body));
1237 };
1238
1239 Ok(OrderPlacement {
1240 order_id: oid,
1241 symbol,
1242 side,
1243 price,
1244 quantity: volume,
1245 timestamp_ms,
1246 cloid: returned_cloid.or(cloid),
1247 order_type,
1248 trigger_price,
1249 reduce_only,
1250 })
1251 }
1252
1253 async fn change_order_by_cloid(
1256 &self,
1257 cloid: i64,
1258 price: Decimal,
1259 volume: Decimal,
1260 ) -> Result<i64, String> {
1261 let user = self.require_user_address()?;
1262
1263 let resp = self
1265 .info_post(
1266 serde_json::json!({"type": "openOrders", "user": user}),
1267 20,
1268 "change_order_by_cloid",
1269 )
1270 .await?;
1271 let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1272 let order = orders
1273 .iter()
1274 .find(|o| o.oid == cloid)
1275 .ok_or_else(|| format!("order {} not found", cloid))?;
1276
1277 let asset_idx = self.get_asset_index(&order.coin).await?;
1278 let is_buy = order.side == "B";
1279
1280 let action = serde_json::json!({
1281 "type": "batchModify",
1282 "modifies": [{
1283 "oid": cloid,
1284 "order": {
1285 "a": asset_idx,
1286 "b": is_buy,
1287 "p": price.to_string(),
1288 "s": volume.to_string(),
1289 "r": false,
1290 "t": {"limit": {"tif": "Gtc"}}
1291 }
1292 }]
1293 });
1294
1295 self.submit_signed_action(action, None).await?;
1296 Ok(cloid)
1297 }
1298
1299 async fn cancel_order_by_cloid(&self, cloid: String) -> Result<(), String> {
1304 let user = self.require_user_address()?;
1305
1306 let resp = self
1308 .info_post(
1309 serde_json::json!({"type": "openOrders", "user": user}),
1310 20,
1311 "cancel_order_by_cloid",
1312 )
1313 .await?;
1314 let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1315 let order = orders
1316 .iter()
1317 .find(|o| o.cloid.as_ref() == Some(&cloid))
1318 .ok_or_else(|| format!("order with cloid {} not found", cloid))?;
1319
1320 let meta_resp = self
1322 .info_post(
1323 serde_json::json!({"type": "meta"}),
1324 20,
1325 "cancel_order_by_cloid",
1326 )
1327 .await?;
1328 let meta: MetaResponse = parse_response(meta_resp).await?;
1329
1330 let asset_idx = meta
1331 .universe
1332 .iter()
1333 .position(|a| a.name == order.coin)
1334 .ok_or_else(|| format!("asset {} not found in meta", order.coin))?;
1335
1336 let action = serde_json::json!({
1339 "type": "cancel",
1340 "cancels": [{"a": asset_idx, "o": order.oid}]
1341 });
1342
1343 self.submit_signed_action(action, None).await?;
1344 Ok(())
1345 }
1346
1347 async fn cancel_all_order(&self) -> Result<bool, String> {
1350 let user = self.require_user_address()?;
1351
1352 let resp = self
1354 .info_post(
1355 serde_json::json!({"type": "openOrders", "user": user}),
1356 20,
1357 "cancel_all_order",
1358 )
1359 .await?;
1360 let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1361 if orders.is_empty() {
1362 return Ok(true);
1363 }
1364
1365 let meta_resp = self
1367 .info_post(serde_json::json!({"type": "meta"}), 20, "cancel_all_order")
1368 .await?;
1369 let meta: MetaResponse = parse_response(meta_resp).await?;
1370
1371 let cancels: Vec<Value> = orders
1372 .iter()
1373 .filter_map(|o| {
1374 let asset_idx = meta.universe.iter().position(|a| a.name == o.coin)?;
1375 Some(serde_json::json!({"a": asset_idx, "o": o.oid}))
1376 })
1377 .collect();
1378
1379 let action = serde_json::json!({"type": "cancel", "cancels": cancels});
1380 self.submit_signed_action(action, None).await?;
1381 Ok(true)
1382 }
1383}
1384
1385#[async_trait]
1386impl guilder_abstraction::SubscribeMarketData for HyperliquidClient {
1387 fn subscribe_l2_update(&self, symbol: String) -> BoxStream<Result<L2Update, String>> {
1388 Box::pin(stream::iter(vec![Err(format!(
1389 "subscribe_l2_update is unsupported for {symbol}; use subscribe_l2_snapshot"
1390 ))]))
1391 }
1392
1393 fn subscribe_l2_snapshot(&self, symbol: String) -> BoxStream<Result<L2Snapshot, String>> {
1394 let subscription = HyperliquidSubscription::L2Book { coin: symbol };
1395 Box::pin(managed_stream(
1396 self.market_ws_manager.clone(),
1397 subscription,
1398 |msg: HyperliquidWsInboundMessage| {
1399 if let Some(snapshot) = msg.as_l2_snapshot() {
1400 vec![Ok(snapshot)]
1401 } else {
1402 vec![]
1403 }
1404 },
1405 ))
1406 }
1407
1408 fn subscribe_asset_context(&self, symbol: String) -> BoxStream<Result<AssetContext, String>> {
1409 let subscription = HyperliquidSubscription::ActiveAssetCtx { coin: symbol };
1410 Box::pin(managed_stream(
1411 self.market_ws_manager.clone(),
1412 subscription,
1413 |msg: HyperliquidWsInboundMessage| {
1414 if let Some(ctx) = msg.as_asset_context() {
1415 vec![Ok(ctx)]
1416 } else {
1417 vec![]
1418 }
1419 }))
1420 }
1421
1422 fn subscribe_liquidation(&self, user: String) -> BoxStream<Result<Liquidation, String>> {
1423 subscribe_user_stream(
1424 self,
1425 user.clone(),
1426 HyperliquidSubscription::UserEvents { user_addr: user },
1427 |msg: HyperliquidWsInboundMessage| {
1428 if let Some(liq) = msg.as_liquidation() {
1429 vec![Ok(liq)]
1430 } else {
1431 vec![]
1432 }
1433 },
1434 )
1435 }
1436
1437 fn subscribe_fill(&self, symbol: String) -> BoxStream<Result<Fill, String>> {
1438 let subscription = HyperliquidSubscription::Trades { coin: symbol };
1439 Box::pin(managed_stream(
1440 self.market_ws_manager.clone(),
1441 subscription,
1442 |msg: HyperliquidWsInboundMessage| {
1443 if let Some(fills) = msg.as_trades() {
1444 fills.into_iter().map(Ok).collect()
1445 } else {
1446 vec![]
1447 }
1448 }))
1449 }
1450
1451 async fn unsubscribe_all(&self) {
1454 self.market_ws_manager.shutdown();
1455 }
1456}
1457
1458fn subscribe_user_stream<T, F>(
1459 client: &HyperliquidClient,
1460 user_addr: String,
1461 subscription: HyperliquidSubscription,
1462 parse: F,
1463) -> BoxStream<Result<T, String>>
1464where
1465 T: Send + 'static,
1466 F: Fn(HyperliquidWsInboundMessage) -> Vec<Result<T, String>> + Send + Sync + 'static,
1467{
1468 let manager = get_or_create_user_manager(
1469 &client.user_ws_managers,
1470 client.ws_send_limiter.clone(),
1471 user_addr,
1472 );
1473 Box::pin(async_stream::stream! {
1474 let stream = managed_stream(manager, subscription, parse);
1475 tokio::pin!(stream);
1476 while let Some(item) = stream.next().await {
1477 yield item;
1478 }
1479 })
1480}
1481
1482fn get_or_create_user_manager(
1483 user_ws_managers: &RwLock<HashMap<String, HyperliquidWsManager>>,
1484 ws_send_limiter: WsSendRateLimiter,
1485 user_addr: String,
1486) -> HyperliquidWsManager {
1487 {
1488 let managers = user_ws_managers.read().unwrap_or_else(|e| e.into_inner());
1489 if let Some(manager) = managers.get(&user_addr) {
1490 return manager.clone();
1491 }
1492 }
1493
1494 let mut managers = user_ws_managers.write().unwrap_or_else(|e| e.into_inner());
1495 managers
1496 .entry(user_addr.clone())
1497 .or_insert_with(|| HyperliquidWsManager::new(Some(user_addr), ws_send_limiter))
1498 .clone()
1499}
1500
1501#[async_trait]
1502impl guilder_abstraction::GetAccountSnapshot for HyperliquidClient {
1503 async fn get_positions(&self) -> Result<Vec<Position>, String> {
1506 let user = self.require_user_address()?;
1507 let resp = self
1509 .info_post(
1510 serde_json::json!({"type": "clearinghouseState", "user": user}),
1511 2,
1512 "get_positions",
1513 )
1514 .await?;
1515 let state: ClearinghouseStateResponse = parse_response(resp).await?;
1516
1517 Ok(state
1518 .asset_positions
1519 .into_iter()
1520 .filter_map(|ap| {
1521 let p = ap.position;
1522 let size = parse_decimal(&p.szi)?;
1523 if size.is_zero() {
1524 return None;
1525 }
1526 let entry_price = p
1527 .entry_px
1528 .as_deref()
1529 .and_then(parse_decimal)
1530 .unwrap_or_default();
1531 let side = if size > Decimal::ZERO {
1532 OrderSide::Buy
1533 } else {
1534 OrderSide::Sell
1535 };
1536 Some(Position {
1537 symbol: p.coin,
1538 side,
1539 size: size.abs(),
1540 entry_price,
1541 })
1542 })
1543 .collect())
1544 }
1545
1546 async fn get_open_orders(&self) -> Result<Vec<OpenOrder>, String> {
1549 let user = self.require_user_address()?;
1550 let resp = self
1552 .info_post(
1553 serde_json::json!({"type": "openOrders", "user": user}),
1554 20,
1555 "get_open_orders",
1556 )
1557 .await?;
1558 let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1559
1560 Ok(orders
1561 .into_iter()
1562 .filter_map(|o| {
1563 let price = parse_decimal(&o.limit_px)?;
1564 let quantity = parse_decimal(&o.orig_sz)?;
1565 let remaining = parse_decimal(&o.sz)?;
1566 let filled_quantity = quantity - remaining;
1567 let side = if o.side == "B" {
1568 OrderSide::Buy
1569 } else {
1570 OrderSide::Sell
1571 };
1572 Some(OpenOrder {
1573 order_id: o.oid,
1574 symbol: o.coin,
1575 side,
1576 price,
1577 quantity,
1578 filled_quantity,
1579 order_type: None, trigger_price: None, reduce_only: false, })
1583 })
1584 .collect())
1585 }
1586
1587 async fn get_balance(&self) -> Result<Vec<guilder_abstraction::AccountBalance>, String> {
1591 let user = self.require_user_address()?;
1592 let resp = self
1594 .info_post(
1595 serde_json::json!({"type": "spotClearinghouseState", "user": user}),
1596 15,
1597 "get_balance",
1598 )
1599 .await?;
1600
1601 #[derive(Deserialize)]
1602 struct SpotStateResponse {
1603 balances: Vec<SpotBalance>,
1604 #[serde(rename = "tokenToAvailableAfterMaintenance")]
1605 token_to_available_after_maintenance: Vec<(i32, String)>,
1606 }
1607
1608 #[allow(dead_code)]
1609 #[derive(Deserialize)]
1610 struct SpotBalance {
1611 coin: String,
1612 total: String,
1613 hold: String,
1614 #[serde(default)]
1615 token: Option<i32>,
1616 #[serde(default)]
1617 #[serde(rename = "entryNtl")]
1618 entry_ntl: Option<String>,
1619 }
1620
1621 let state: SpotStateResponse = parse_response(resp).await?;
1622
1623 let safe_map: HashMap<i32, Decimal> = state
1626 .token_to_available_after_maintenance
1627 .into_iter()
1628 .filter_map(|(token_id, value)| {
1629 parse_decimal(&value).map(|v| (token_id, v))
1630 })
1631 .collect();
1632
1633 let perp_margin_used: Option<Decimal> = match self
1637 .info_post(
1638 serde_json::json!({"type": "clearinghouseState", "user": user}),
1639 2,
1640 "get_balance_margin",
1641 )
1642 .await
1643 {
1644 Ok(resp) => parse_response::<ClearinghouseStateResponse>(resp)
1645 .await
1646 .ok()
1647 .and_then(|ch| {
1648 ch.margin_summary
1649 .total_margin_used
1650 .as_deref()
1651 .and_then(parse_decimal)
1652 }),
1653 Err(_) => None,
1654 };
1655
1656 state
1657 .balances
1658 .into_iter()
1659 .map(|balance| {
1660 let equity = parse_decimal(&balance.total)
1661 .ok_or_else(|| "invalid total balance".to_string())?;
1662 let hold = parse_decimal(&balance.hold)
1663 .ok_or_else(|| "invalid hold balance".to_string())?;
1664 let free = equity - hold;
1665
1666 let safe = balance.token.and_then(|token_id| safe_map.get(&token_id).copied());
1667 let usable = match safe {
1668 Some(s) => std::cmp::min(free, s),
1669 None => free,
1670 };
1671 let maintenance = safe.map(|s| equity - s);
1672
1673 let margin_used = if balance.coin == "USDC" {
1675 perp_margin_used
1676 } else {
1677 None
1678 };
1679
1680 Ok(guilder_abstraction::AccountBalance {
1681 token: balance.coin,
1682 equity,
1683 free,
1684 safe,
1685 usable,
1686 hold,
1687 margin_used,
1688 maintenance,
1689 })
1690 })
1691 .collect()
1692 }
1693
1694 async fn get_user_rate_limit(&self) -> Result<guilder_abstraction::UserRateLimit, String> {
1697 let user = self.require_user_address()?;
1698 let resp = self
1699 .info_post(
1700 serde_json::json!({"type": "userRateLimit", "user": user}),
1701 20,
1702 "get_user_rate_limit",
1703 )
1704 .await?;
1705 let val = parse_response::<Value>(resp).await?;
1706
1707 let cumulative_volume = val["cumVlm"]
1708 .as_str()
1709 .and_then(parse_decimal)
1710 .ok_or_else(|| "missing or invalid cumVlm".to_string())?;
1711 let requests_used = val["nRequestsUsed"]
1712 .as_i64()
1713 .ok_or_else(|| "missing or invalid nRequestsUsed".to_string())?;
1714 let requests_cap = val["nRequestsCap"]
1715 .as_i64()
1716 .ok_or_else(|| "missing or invalid nRequestsCap".to_string())?;
1717 let requests_surplus = val["nRequestsSurplus"]
1718 .as_i64()
1719 .ok_or_else(|| "missing or invalid nRequestsSurplus".to_string())?;
1720
1721 Ok(guilder_abstraction::UserRateLimit {
1722 cumulative_volume,
1723 requests_used,
1724 requests_cap,
1725 requests_surplus,
1726 })
1727 }
1728}
1729
1730#[async_trait]
1731impl guilder_abstraction::SubscribeUserEvents for HyperliquidClient {
1732 fn subscribe_user_fills(&self) -> BoxStream<Result<UserFill, String>> {
1733 let Some(addr) = self.user_address.as_ref() else {
1734 return Box::pin(stream::iter(vec![Err(
1735 "user address not registered".to_string()
1736 )]));
1737 };
1738 subscribe_user_stream(
1739 self,
1740 addr.clone(),
1741 HyperliquidSubscription::UserEvents {
1742 user_addr: addr.clone(),
1743 },
1744 |msg: HyperliquidWsInboundMessage| {
1745 if let Some(fills) = msg.as_user_fills() {
1746 fills.into_iter().map(Ok).collect()
1747 } else {
1748 vec![]
1749 }
1750 },
1751 )
1752 }
1753
1754 fn subscribe_order_updates(&self) -> BoxStream<Result<OrderUpdate, String>> {
1755 let Some(addr) = self.user_address.as_ref() else {
1756 return Box::pin(stream::iter(vec![Err(
1757 "user address not registered".to_string()
1758 )]));
1759 };
1760 subscribe_user_stream(
1761 self,
1762 addr.clone(),
1763 HyperliquidSubscription::OrderUpdates {
1764 user_addr: addr.clone(),
1765 },
1766 |msg: HyperliquidWsInboundMessage| {
1767 if let Some(updates) = msg.as_order_updates() {
1768 updates.into_iter().map(Ok).collect()
1769 } else {
1770 vec![]
1771 }
1772 },
1773 )
1774 }
1775
1776 fn subscribe_funding_payments(&self) -> BoxStream<Result<FundingPayment, String>> {
1777 let Some(addr) = self.user_address.as_ref() else {
1778 return Box::pin(stream::iter(vec![Err(
1779 "user address not registered".to_string()
1780 )]));
1781 };
1782 subscribe_user_stream(
1783 self,
1784 addr.clone(),
1785 HyperliquidSubscription::UserEvents {
1786 user_addr: addr.clone(),
1787 },
1788 |msg: HyperliquidWsInboundMessage| {
1789 if let Some(p) = msg.as_funding_payment() {
1790 vec![Ok(p)]
1791 } else {
1792 vec![]
1793 }
1794 },
1795 )
1796 }
1797
1798 fn subscribe_deposits(&self) -> BoxStream<Result<Deposit, String>> {
1799 let Some(addr) = self.user_address.as_ref() else {
1800 return Box::pin(stream::iter(vec![Err(
1801 "user address not registered".to_string()
1802 )]));
1803 };
1804 subscribe_user_stream(
1805 self,
1806 addr.clone(),
1807 HyperliquidSubscription::NonFundingLedger {
1808 user_addr: addr.clone(),
1809 },
1810 |msg: HyperliquidWsInboundMessage| {
1811 if let Some(deps) = msg.as_deposits() {
1812 deps.into_iter().map(Ok).collect()
1813 } else {
1814 vec![]
1815 }
1816 },
1817 )
1818 }
1819
1820 fn subscribe_withdrawals(&self) -> BoxStream<Result<Withdrawal, String>> {
1821 let Some(addr) = self.user_address.as_ref() else {
1822 return Box::pin(stream::iter(vec![Err(
1823 "user address not registered".to_string()
1824 )]));
1825 };
1826 subscribe_user_stream(
1827 self,
1828 addr.clone(),
1829 HyperliquidSubscription::NonFundingLedger {
1830 user_addr: addr.clone(),
1831 },
1832 |msg: HyperliquidWsInboundMessage| {
1833 if let Some(wds) = msg.as_withdrawals() {
1834 wds.into_iter().map(Ok).collect()
1835 } else {
1836 vec![]
1837 }
1838 },
1839 )
1840 }
1841
1842 fn subscribe_spot_balance(
1844 &self,
1845 ) -> BoxStream<Result<Vec<guilder_abstraction::AccountBalance>, String>> {
1846 let Some(addr) = self.user_address.as_ref() else {
1847 return Box::pin(stream::iter(vec![Err(
1848 "user address not registered".to_string()
1849 )]));
1850 };
1851 self.subscribe_spot_balance_with_address(addr.clone())
1852 }
1853
1854 fn subscribe_spot_balance_with_address(
1856 &self,
1857 address: String,
1858 ) -> BoxStream<Result<Vec<guilder_abstraction::AccountBalance>, String>> {
1859 subscribe_user_stream(
1860 self,
1861 address.clone(),
1862 HyperliquidSubscription::UserEvents { user_addr: address },
1863 |msg: HyperliquidWsInboundMessage| {
1864 if let Some(balances) = msg.as_spot_balance() {
1865 vec![Ok(balances)]
1866 } else {
1867 vec![]
1868 }
1869 },
1870 )
1871 }
1872
1873 async fn unsubscribe_user_events(&self) {
1874 if let Some(addr) = &self.user_address {
1876 self.market_ws_manager.unsubscribe_user(addr);
1877 }
1878 let managers = self.user_ws_managers.read().unwrap_or_else(|e| e.into_inner());
1880 for (addr, manager) in managers.iter() {
1881 manager.unsubscribe_user(addr);
1882 }
1883 }
1884}
1885
1886#[async_trait]
1887impl guilder_abstraction::SubscribeMarketDataOps for HyperliquidClient {
1888 async fn unsubscribe_market_data(&self, symbol: String) {
1889 self.market_ws_manager.unsubscribe_by_coin(&symbol);
1890 }
1891}
1892
1893#[cfg(test)]
1894mod msgpack_tests {
1895 use super::*;
1896 use serde_json::json;
1897
1898 #[test]
1899 fn test_msgpack_null() {
1900 let result = value_to_msgpack(&Value::Null);
1901 assert_eq!(result, vec![0xc0]);
1902 }
1903
1904 #[test]
1905 fn test_msgpack_bool() {
1906 assert_eq!(value_to_msgpack(&Value::Bool(true)), vec![0xc3]);
1907 assert_eq!(value_to_msgpack(&Value::Bool(false)), vec![0xc2]);
1908 }
1909
1910 #[test]
1911 fn test_msgpack_positive_fixint() {
1912 assert_eq!(value_to_msgpack(&json!(0)), vec![0x00]);
1914 assert_eq!(value_to_msgpack(&json!(1)), vec![0x01]);
1915 assert_eq!(value_to_msgpack(&json!(127)), vec![0x7f]);
1916 }
1917
1918 #[test]
1919 fn test_msgpack_uint8() {
1920 assert_eq!(value_to_msgpack(&json!(128)), vec![0xcc, 0x80]);
1922 assert_eq!(value_to_msgpack(&json!(255)), vec![0xcc, 0xff]);
1923 }
1924
1925 #[test]
1926 fn test_msgpack_uint16() {
1927 assert_eq!(value_to_msgpack(&json!(256)), vec![0xcd, 0x01, 0x00]);
1929 assert_eq!(value_to_msgpack(&json!(65535)), vec![0xcd, 0xff, 0xff]);
1930 }
1931
1932 #[test]
1933 fn test_msgpack_uint32() {
1934 assert_eq!(
1936 value_to_msgpack(&json!(65536)),
1937 vec![0xce, 0x00, 0x01, 0x00, 0x00]
1938 );
1939 assert_eq!(
1940 value_to_msgpack(&json!(4294967295u64)),
1941 vec![0xce, 0xff, 0xff, 0xff, 0xff]
1942 );
1943 }
1944
1945 #[test]
1946 fn test_msgpack_uint64() {
1947 let big: u64 = 4294967296;
1949 assert_eq!(
1950 value_to_msgpack(&json!(big)),
1951 vec![0xcf, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00]
1952 );
1953 }
1954
1955 #[test]
1956 fn test_msgpack_negative_fixint() {
1957 assert_eq!(value_to_msgpack(&json!(-1)), vec![0xff]);
1959 assert_eq!(value_to_msgpack(&json!(-32)), vec![0xe0]);
1960 }
1961
1962 #[test]
1963 fn test_msgpack_int8() {
1964 assert_eq!(value_to_msgpack(&json!(-33)), vec![0xd0, 0xdf]);
1966 assert_eq!(value_to_msgpack(&json!(-128)), vec![0xd0, 0x80]);
1967 }
1968
1969 #[test]
1970 fn test_msgpack_int16() {
1971 assert_eq!(value_to_msgpack(&json!(-129)), vec![0xd1, 0xff, 0x7f]);
1973 assert_eq!(value_to_msgpack(&json!(-32768)), vec![0xd1, 0x80, 0x00]);
1974 }
1975
1976 #[test]
1977 fn test_msgpack_int32() {
1978 assert_eq!(
1980 value_to_msgpack(&json!(-32769)),
1981 vec![0xd2, 0xff, 0xff, 0x7f, 0xff]
1982 );
1983 assert_eq!(
1984 value_to_msgpack(&json!(-2147483648i64)),
1985 vec![0xd2, 0x80, 0x00, 0x00, 0x00]
1986 );
1987 }
1988
1989 #[test]
1990 fn test_msgpack_int64() {
1991 let val: i64 = -2147483649;
1992 let result = value_to_msgpack(&json!(val));
1993 assert_eq!(result[0], 0xd3); assert_eq!(result.len(), 9);
1995 }
1996
1997 #[test]
1998 fn test_msgpack_float() {
1999 let result = value_to_msgpack(&json!(3.14));
2000 assert_eq!(result[0], 0xcb); assert_eq!(result.len(), 9);
2002 }
2003
2004 #[test]
2005 fn test_msgpack_fixstr() {
2006 assert_eq!(value_to_msgpack(&json!("")), vec![0xa0]);
2008 assert_eq!(value_to_msgpack(&json!("hello")), {
2009 let mut expected = vec![0xa5];
2010 expected.extend_from_slice(b"hello");
2011 expected
2012 });
2013 let s = "a".repeat(31);
2014 let result = value_to_msgpack(&json!(s));
2015 assert_eq!(result[0], 0xbf); assert_eq!(result.len(), 32);
2017 }
2018
2019 #[test]
2020 fn test_msgpack_str8() {
2021 let s = "a".repeat(32);
2022 let result = value_to_msgpack(&json!(s));
2023 assert_eq!(result[0], 0xd9); assert_eq!(result[1], 32);
2025 assert_eq!(result.len(), 34);
2026 }
2027
2028 #[test]
2029 fn test_msgpack_fixarray() {
2030 assert_eq!(value_to_msgpack(&json!([])), vec![0x90]);
2032 let result = value_to_msgpack(&json!([1, 2, 3]));
2033 assert_eq!(result[0], 0x93);
2034 assert_eq!(result, vec![0x93, 0x01, 0x02, 0x03]);
2035 }
2036
2037 #[test]
2038 fn test_msgpack_fixmap() {
2039 assert_eq!(value_to_msgpack(&json!({})), vec![0x80]);
2041 let result = value_to_msgpack(&json!({"a": 1}));
2042 assert_eq!(result[0], 0x81); assert_eq!(result, {
2044 let mut expected = vec![0x81];
2045 expected.extend_from_slice(&value_to_msgpack(&json!("a")));
2046 expected.extend_from_slice(&value_to_msgpack(&json!(1)));
2047 expected
2048 });
2049 }
2050
2051 #[test]
2052 fn test_msgmap_preserves_insertion_order() {
2053 let val = json!({
2055 "z": 1,
2056 "a": 2,
2057 "m": 3
2058 });
2059 let result = value_to_msgpack(&val);
2060 assert_eq!(result[0], 0x83);
2062 assert_eq!(result[1], 0xa1); assert_eq!(result[2], b'z');
2065 }
2066
2067 #[test]
2068 fn test_msgpack_mixed_array() {
2069 let val = json!([null, true, false, 42, "hi", [1, 2]]);
2070 let result = value_to_msgpack(&val);
2071 assert_eq!(result[0], 0x96); assert_eq!(result[1], 0xc0); assert_eq!(result[2], 0xc3); assert_eq!(result[3], 0xc2); assert_eq!(result[4], 0x2a); assert_eq!(result[5], 0xa2);
2078 assert_eq!(result[6], b'h');
2079 assert_eq!(result[7], b'i');
2080 }
2081
2082 #[test]
2083 fn test_build_order_msgpack_without_cloid() {
2084 let result = build_order_msgpack(
2085 0, true, "1000", "0.1", false, "limit", b"gtc", None, );
2094 assert_eq!(result[0], 0x86);
2096 }
2097
2098 #[test]
2099 fn test_build_order_msgpack_with_cloid() {
2100 let result = build_order_msgpack(
2101 0, true, "1000", "0.1", false, "limit", b"gtc", Some("my-cloid"), );
2110 assert_eq!(result[0], 0x87);
2112 }
2113
2114 #[test]
2115 fn test_action_to_canonical_msgpack() {
2116 let action = json!({
2117 "type": "order",
2118 "orders": [{"a": 0, "b": true, "p": "1000", "s": "0.1", "r": false, "t": {"limit": {"tif": "gtc"}}}],
2119 "grouping": "na"
2120 });
2121 let result = action_to_canonical_msgpack(&action).unwrap();
2122 assert_eq!(result[0], 0x83);
2124 }
2125
2126 #[test]
2127 fn test_msgpack_matches_rmp_serde_for_simple_values() {
2128 use rmp_serde::to_vec;
2130
2131 for val in [
2132 json!(0),
2133 json!(127),
2134 json!(255),
2135 json!(1000),
2136 json!(-1),
2137 json!(-32),
2138 json!(-128),
2139 ] {
2140 let ours = value_to_msgpack(&val);
2141 let theirs = to_vec(&val).unwrap();
2142 assert_eq!(
2143 ours, theirs,
2144 "mismatch for {}: ours={:?}, rmp={:?}",
2145 val, ours, theirs
2146 );
2147 }
2148 }
2149
2150 #[test]
2151 fn test_msgpack_string_encoding() {
2152 use rmp_serde::to_vec;
2153 for val in [
2154 json!(""),
2155 json!("a"),
2156 json!("hello world"),
2157 json!("BTC-USD"),
2158 ] {
2159 let ours = value_to_msgpack(&val);
2160 let theirs = to_vec(&val).unwrap();
2161 assert_eq!(
2162 ours, theirs,
2163 "mismatch for {}: ours={:?}, rmp={:?}",
2164 val, ours, theirs
2165 );
2166 }
2167 }
2168
2169 #[test]
2170 fn test_msgpack_bool_encoding() {
2171 use rmp_serde::to_vec;
2172 let theirs = to_vec(&json!(true)).unwrap();
2173 assert_eq!(value_to_msgpack(&json!(true)), theirs);
2174 let theirs = to_vec(&json!(false)).unwrap();
2175 assert_eq!(value_to_msgpack(&json!(false)), theirs);
2176 }
2177
2178 #[test]
2179 fn test_msgpack_null_encoding() {
2180 use rmp_serde::to_vec;
2181 let theirs = to_vec(&Value::Null).unwrap();
2182 assert_eq!(value_to_msgpack(&Value::Null), theirs);
2183 }
2184
2185 #[test]
2186 fn test_msgpack_nested_object() {
2187 let val = json!({
2188 "outer": {
2189 "inner": 42
2190 }
2191 });
2192 let result = value_to_msgpack(&val);
2193 assert_eq!(result[0], 0x81); }
2195
2196 #[test]
2197 fn test_msgpack_empty_containers() {
2198 assert_eq!(value_to_msgpack(&json!([])), vec![0x90]);
2199 assert_eq!(value_to_msgpack(&json!({})), vec![0x80]);
2200 }
2201}
2202
2203
2204#[async_trait::async_trait]
2205impl guilder_abstraction::ListingEventSource for HyperliquidClient {
2206 async fn get_listing_events(
2214 &self,
2215 ) -> Result<Vec<guilder_abstraction::ListingEvent>, String> {
2216 let resp = self
2218 .info_post(serde_json::json!({"type": "meta"}), 20, "get_listing_events")
2219 .await?;
2220 let meta = parse_response::<MetaResponse>(resp).await?;
2221 Ok(meta
2222 .universe
2223 .into_iter()
2224 .map(|a| guilder_abstraction::ListingEvent {
2225 ticker: a.name,
2226 exchange: "hyperliquid_perp".to_string(),
2227 event: if a.is_delisted { "delist" } else { "list" }.to_string(),
2228 event_time: String::new(),
2229 })
2230 .collect())
2231 }
2232
2233 async fn get_current_universe(
2234 &self,
2235 ) -> Result<Vec<guilder_abstraction::SymbolStatus>, String> {
2236 let resp = self
2238 .info_post(serde_json::json!({"type": "meta"}), 20, "get_current_universe")
2239 .await?;
2240 let meta = parse_response::<MetaResponse>(resp).await?;
2241 Ok(meta
2242 .universe
2243 .into_iter()
2244 .map(|a| guilder_abstraction::SymbolStatus {
2245 ticker: a.name,
2246 exchange: "hyperliquid_perp".to_string(),
2247 is_delisted: a.is_delisted,
2248 })
2249 .collect())
2250 }
2251}