Skip to main content

guilder_client_hyperliquid/
client.rs

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    /// Create a client authenticated via an external signer (TPM, Secure Enclave, etc.).
98    ///
99    /// The signer handles the actual ECDSA signing; the private key never leaves
100    /// the hardware. The client computes the EIP-712 digest and delegates signing
101    /// to the external signer.
102    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    /// Configure rate limit budgets (rest_weight/min, address_requests).
121    /// Defaults: 1200 rest weight/min, 10000 address requests.
122    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    /// POST to the info endpoint, consuming `weight` from the REST rate-limit budget.
129    /// Returns `Err("rate_limited: ...")` immediately if budget is exhausted — no retry.
130    /// Callers should handle gracefully (skip cycle, retry later, etc.).
131    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        // `meta` is an "all other info" request → weight 20
168        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        // Check both rate limiters non-blocking — fail fast, no retry.
207        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// --- REST deserialization types ---
246
247#[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/// Kept for get_positions compatibility; margin_summary fields are unused since get_collateral was removed.
286#[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    /// positive = long, negative = short
309    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
325// predictedFundings response: Vec<(coin, Vec<(venue, entry_or_null)>)>
326// The API returns null for venues that don't list the coin.
327type 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
336// --- Helpers ---
337
338fn 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
347/// EIP-712 domain separator for Hyperliquid L1 actions (chainId=1337).
348fn 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
367/// Convert a `serde_json::Value` to msgpack bytes, preserving the JSON map key order.
368/// This avoids rmp_serde's HashMap-based serialization which reorders map keys.
369fn 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
493/// Convert action to msgpack bytes preserving JSON map key insertion order
494/// (matching Python's msgpack dict ordering).
495fn action_to_canonical_msgpack(action: &Value) -> Result<Vec<u8>, String> {
496    Ok(value_to_msgpack(action))
497}
498
499/// Build msgpack for a single order with Python SDK field order:
500/// a, b, p, s, r, t, c(opt)
501fn 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)); // fixmap
514
515    // "a": asset_idx
516    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    // "b": is_buy
522    buf.extend_from_slice(&value_to_msgpack(&Value::String("b".to_string())));
523    buf.push(if is_buy { 0xc3 } else { 0xc2 });
524
525    // "p": price
526    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    // "s": size (Python SDK puts s before r)
530    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    // "r": reduce_only
534    buf.extend_from_slice(&value_to_msgpack(&Value::String("r".to_string())));
535    buf.push(if reduce_only { 0xc3 } else { 0xc2 });
536
537    // "t": { order_kind: { "tif": tif_str } }
538    buf.extend_from_slice(&value_to_msgpack(&Value::String("t".to_string())));
539    // Inner: fixmap(1) with order_kind key
540    buf.push(0x81);
541    buf.extend_from_slice(&value_to_msgpack(&Value::String(order_kind.to_string())));
542    // Inner-inner: fixmap(1) with "tif" key
543    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    // "c": cloid (optional, appended at END per Python SDK)
550    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
558/// Build msgpack for a trigger order (take profit / stop loss) with Python SDK field order.
559/// Trigger orders use: t = { "trigger": { "isMarket": bool, "triggerPx": str, "tpsl": "tp"|"sl" } }
560fn 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)); // fixmap
574
575    // "a": asset_idx
576    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    // "b": is_buy
582    buf.extend_from_slice(&value_to_msgpack(&Value::String("b".to_string())));
583    buf.push(if is_buy { 0xc3 } else { 0xc2 });
584
585    // "p": price (use trigger_px as price for resting, or "0" for market-on-trigger)
586    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    // "s": size
590    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    // "r": reduce_only
594    buf.extend_from_slice(&value_to_msgpack(&Value::String("r".to_string())));
595    buf.push(if reduce_only { 0xc3 } else { 0xc2 });
596
597    // "t": { "trigger": { "isMarket": bool, "triggerPx": str, "tpsl": str } }
598    buf.extend_from_slice(&value_to_msgpack(&Value::String("t".to_string())));
599    buf.push(0x81); // fixmap(1): "trigger"
600    buf.extend_from_slice(&value_to_msgpack(&Value::String("trigger".to_string())));
601    // trigger object: fixmap(3)
602    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    // "c": cloid (optional, appended at END)
611    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
619/// Compute the EIP-712 digest for a Hyperliquid exchange action.
620///
621/// This is the hash that gets signed — extracted so both direct key and
622/// external signer paths can share it.
623fn 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
657/// Format an EcdsaSignature into the (r, s, v) hex strings expected by Hyperliquid.
658fn 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
665/// Sign a digest using a raw private key (direct key mode).
666fn 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
688/// Sign using pre-built msgpack bytes (bypassing serde_json field ordering).
689/// Supports both direct private key and external signer.
690async 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
710/// Signs a Hyperliquid exchange action using EIP-712.
711/// Returns (r, s, v) where r and s are "0x"-prefixed hex strings and v is 27 or 28.
712/// Supports both direct private key and external signer.
713async 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// --- Trait implementations ---
735
736#[async_trait]
737impl guilder_abstraction::TestServer for HyperliquidClient {
738    /// Sends a lightweight allMids request; returns true if the server responds 200 OK.
739    async fn ping(&self) -> Result<bool, String> {
740        // allMids → weight 2
741        self.info_post(serde_json::json!({"type": "allMids"}), 2, "ping")
742            .await
743            .map(|r| r.status().is_success())
744    }
745
746    /// Hyperliquid has no dedicated server-time endpoint; returns local UTC ms.
747    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    /// Returns all perpetual asset names from Hyperliquid's meta endpoint.
758    async fn get_symbol(&self) -> Result<Vec<String>, String> {
759        // meta → weight 20
760        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    /// Returns the current open interest for `symbol` from metaAndAssetCtxs.
769    async fn get_open_interest(&self, symbol: String) -> Result<Decimal, String> {
770        // metaAndAssetCtxs → weight 20
771        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    /// Returns a full AssetContext snapshot for `symbol` from metaAndAssetCtxs.
790    async fn get_asset_context(&self, symbol: String) -> Result<AssetContext, String> {
791        // metaAndAssetCtxs → weight 20
792        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    /// Fetches metaAndAssetCtxs once and returns all asset contexts in universe order.
825    /// Prefer this over repeated `get_asset_context` calls to avoid rate-limiting.
826    async fn get_all_asset_contexts(&self) -> Result<Vec<AssetContext>, String> {
827        // metaAndAssetCtxs → weight 20
828        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    /// Returns the number of decimal places for order size for a symbol.
869    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    /// Returns sz_decimals for all symbols from the meta universe.
877    async fn get_all_sz_decimals(&self) -> Result<HashMap<String, i32>, String> {
878        // metaAndAssetCtxs → weight 20
879        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    /// Returns a full L2 orderbook snapshot for `symbol` from the l2Book REST endpoint.
897    async fn get_l2_orderbook(&self, symbol: String) -> Result<L2Snapshot, String> {
898        // l2Book → weight 2
899        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    /// Returns the mid-price of `symbol` (e.g. "BTC") from allMids.
951    async fn get_price(&self, symbol: String) -> Result<Decimal, String> {
952        // allMids → weight 2
953        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    /// Returns predicted funding rates for all symbols across all venues.
964    /// Null venue entries (unsupported coins) are silently skipped.
965    async fn get_predicted_fundings(&self) -> Result<Vec<PredictedFunding>, String> {
966        // predictedFundings → weight 20
967        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    /// Places an order on Hyperliquid. Requires `with_auth`. Returns an `OrderPlacement` with
996    /// the exchange-assigned order ID. Market orders are submitted as aggressive limit orders (IOC).
997    ///
998    /// If `cloid` is provided, Hyperliquid attaches it to the order lifecycle — fills and order
999    /// updates will carry the same cloid back, enabling end-to-end intent tracing without a
1000    /// separate order_id mapping.
1001    ///
1002    /// Trigger orders (`TakeProfit` / `StopLoss`) require `trigger_price` to be set. The order
1003    /// activates when the mark price reaches `triggerPx`, then executes as a market or limit
1004    /// order depending on `time_in_force` (`Ioc` = market, `Gtc` = limit).
1005    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        // Rate limiting is handled in submit_signed_action (non-blocking).
1018        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        // Determine if this is a trigger order
1031        let is_trigger = matches!(order_type, OrderType::TakeProfit | OrderType::StopLoss);
1032
1033        let (order_msgpack, order_type_json) = if is_trigger {
1034            // --- Trigger order ---
1035            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            // TimeInForce determines market vs limit on trigger:
1047            // Ioc = market execution (isMarket: true), Gtc/Alo/Fok = limit (isMarket: false)
1048            let is_market = matches!(time_in_force, TimeInForce::Ioc);
1049
1050            // For trigger orders, `p` must be set to the trigger price (not "0"),
1051            // even for market-on-trigger. Hyperliquid validates this field.
1052            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            // --- Regular limit/market order ---
1085            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        // Build the action-level msgpack with Python SDK field order (insertion order):
1115        // type → orders → grouping
1116        let mut action_msgpack = Vec::new();
1117        action_msgpack.push(0x83); // fixmap(3)
1118        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); // fixarray(1)
1122        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            // For market-on-trigger, use trigger price (same as msgpack)
1136            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        // Sign using the canonical msgpack (matching Python's field order)
1152        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            // resting doesn't include a timestamp
1215            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            // filled doesn't include cloid
1225            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    /// Modifies price and size of an existing order by its order ID. Requires `with_auth`.
1254    /// Fetches the order's current coin and side before submitting the modify action.
1255    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        // openOrders → weight 20; get_asset_index → meta weight 20
1264        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    /// Cancels a single order by its client order ID (cloid). Requires `with_auth`.
1300    /// Fetches open orders to resolve the order ID for the matching cloid, then
1301    /// submits a cancel action using the order ID — this works for all order types
1302    /// including trigger orders (TakeProfit/StopLoss).
1303    async fn cancel_order_by_cloid(&self, cloid: String) -> Result<(), String> {
1304        let user = self.require_user_address()?;
1305
1306        // openOrders → weight 20
1307        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        // meta → weight 20
1321        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        // Use the same "cancel" action type as cancel_all_order, which is
1337        // proven to work for all order types including trigger orders.
1338        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    /// Cancels all open orders. Requires `with_auth`.
1348    /// Fetches all open orders and submits a batch cancel in a single signed request.
1349    async fn cancel_all_order(&self) -> Result<bool, String> {
1350        let user = self.require_user_address()?;
1351
1352        // openOrders → weight 20
1353        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        // meta → weight 20
1366        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    /// Gracefully shut down all market data subscriptions.
1452    /// Closes all broadcast channels so subscribers exit without reconnecting.
1453    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    /// Returns open positions from `clearinghouseState`. Requires `with_auth`.
1504    /// Zero-size positions are filtered out. Positive `szi` = long, negative = short.
1505    async fn get_positions(&self) -> Result<Vec<Position>, String> {
1506        let user = self.require_user_address()?;
1507        // clearinghouseState → weight 2
1508        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    /// Returns resting orders from Hyperliquid's `openOrders` endpoint. Requires `with_auth`.
1547    /// `filled_quantity` is derived as `origSz - sz` (original size minus remaining size).
1548    async fn get_open_orders(&self) -> Result<Vec<OpenOrder>, String> {
1549        let user = self.require_user_address()?;
1550        // openOrders → weight 20
1551        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, // openOrders REST endpoint doesn't return order type
1580                    trigger_price: None, // trigger info not included in openOrders response
1581                    reduce_only: false, // default; REST doesn't expose this field
1582                })
1583            })
1584            .collect())
1585    }
1586
1587    /// Returns all per-asset balances from `spotClearinghouseState` with margin health.
1588    /// Also fetches `clearinghouseState` to populate `margin_used` on the USDC entry.
1589    /// Requires `with_auth`.
1590    async fn get_balance(&self) -> Result<Vec<guilder_abstraction::AccountBalance>, String> {
1591        let user = self.require_user_address()?;
1592        // spotClearinghouseState → weight 15
1593        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        // Build a safe lookup map from token ID → available after maintenance.
1624        // Not every token appears in the map — absence means zero maintenance impact.
1625        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        // Fetch perp-side margin summary for margin_used.
1634        // clearinghouseState → weight 2. Failure is non-fatal — margin_used
1635        // stays None for all entries.
1636        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                // margin_used is account-level (perp side); only populate on USDC.
1674                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    /// Returns the user's address-level API rate limit budget.
1695    /// Queries Hyperliquid's `userRateLimit` info endpoint for authoritative server-side counts.
1696    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    /// Subscribe to spot wallet balance updates for the registered user address.
1843    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    /// Subscribe to spot wallet balance updates for a specific address.
1855    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        // Unsubscribe from the market-level user manager if we have a user address.
1875        if let Some(addr) = &self.user_address {
1876            self.market_ws_manager.unsubscribe_user(addr);
1877        }
1878        // Also unsubscribe from all per-user managers.
1879        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        // 0–127: positive fixint
1913        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        // 128–255: uint8
1921        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        // 256–65535: uint16
1928        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        // 65536–4294967295: uint32
1935        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        // >4294967295: uint64
1948        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        // -1 to -32: negative fixint
1958        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        // -33 to -128: int8
1965        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        // -129 to -32768: int16
1972        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        // -32769 to -2147483648: int32
1979        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); // int64 marker
1994        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); // float64 marker
2001        assert_eq!(result.len(), 9);
2002    }
2003
2004    #[test]
2005    fn test_msgpack_fixstr() {
2006        // 0–31 bytes: fixstr
2007        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); // 0xa0 | 31
2016        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); // str8 marker
2024        assert_eq!(result[1], 32);
2025        assert_eq!(result.len(), 34);
2026    }
2027
2028    #[test]
2029    fn test_msgpack_fixarray() {
2030        // 0–15 elements: fixarray
2031        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        // 0–15 entries: fixmap
2040        assert_eq!(value_to_msgpack(&json!({})), vec![0x80]);
2041        let result = value_to_msgpack(&json!({"a": 1}));
2042        assert_eq!(result[0], 0x81); // fixmap(1)
2043        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        // Verify keys are serialized in JSON insertion order, not sorted
2054        let val = json!({
2055            "z": 1,
2056            "a": 2,
2057            "m": 3
2058        });
2059        let result = value_to_msgpack(&val);
2060        // fixmap(3)
2061        assert_eq!(result[0], 0x83);
2062        // First key should be "z" (insertion order), not "a" (sorted)
2063        assert_eq!(result[1], 0xa1); // fixstr(1)
2064        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); // fixarray(6)
2072        assert_eq!(result[1], 0xc0); // null
2073        assert_eq!(result[2], 0xc3); // true
2074        assert_eq!(result[3], 0xc2); // false
2075        assert_eq!(result[4], 0x2a); // 42
2076                                     // "hi" = fixstr(2) + "hi"
2077        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,       // asset index
2086            true,    // is_buy
2087            "1000",  // price
2088            "0.1",   // size
2089            false,   // reduce_only
2090            "limit", // order_kind
2091            b"gtc",  // tif
2092            None,    // cloid
2093        );
2094        // fixmap(6)
2095        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,                // asset index
2102            true,             // is_buy
2103            "1000",           // price
2104            "0.1",            // size
2105            false,            // reduce_only
2106            "limit",          // order_kind
2107            b"gtc",           // tif
2108            Some("my-cloid"), // cloid
2109        );
2110        // fixmap(7)
2111        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        // fixmap(3)
2123        assert_eq!(result[0], 0x83);
2124    }
2125
2126    #[test]
2127    fn test_msgpack_matches_rmp_serde_for_simple_values() {
2128        // Verify our encoding matches rmp_serde for simple scalar values
2129        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); // fixmap(1)
2194    }
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    /// Authoritative lifecycle events from the venue. Hyperliquid's info API
2207    /// exposes only the CURRENT universe (+ isDelisted flags) — it has no
2208    /// historical listing stream, so this returns one synthetic `list` event
2209    /// per never-delisted symbol (event_time unknown → epoch placeholder) and
2210    /// a `delist` event per isDelisted symbol. The QDB sync layer upserts
2211    /// these into token_registry_events; true FIRST-LISTING times for
2212    /// pre-history coins come from earlier sync runs, not from this call.
2213    async fn get_listing_events(
2214        &self,
2215    ) -> Result<Vec<guilder_abstraction::ListingEvent>, String> {
2216        // meta → weight 20
2217        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        // meta → weight 20
2237        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}