Skip to main content

guilder_client_hyperliquid/
client.rs

1use crate::rate_limiter::{AddressRateLimiter, RestRateLimiter};
2use futures_util::{stream, StreamExt};
3use guilder_abstraction::{
4    self, AssetContext, BoxStream, Deposit, Fill, FundingPayment, L2Update, Liquidation, OpenOrder,
5    OrderPlacement, OrderSide, OrderStatus, OrderType, OrderUpdate, Position, PredictedFunding,
6    Side, TimeInForce, UserFill, Withdrawal,
7};
8use reqwest::Client;
9use rust_decimal::Decimal;
10use serde::Deserialize;
11use serde_json::Value;
12use std::collections::HashMap;
13use std::str::FromStr;
14use std::sync::Arc;
15const HYPERLIQUID_INFO_URL: &str = "https://api.hyperliquid.xyz/info";
16const HYPERLIQUID_EXCHANGE_URL: &str = "https://api.hyperliquid.xyz/exchange";
17
18async fn parse_response<T: for<'de> serde::Deserialize<'de>>(
19    resp: reqwest::Response,
20) -> Result<T, String> {
21    let status = resp.status();
22    let text = resp.text().await.map_err(|e| {
23        format!("failed to read response body (status {status}): {e}")
24    })?;
25
26    if text.is_empty() {
27        return Err(format!(
28            "empty response body from Hyperliquid (HTTP {status})"
29        ));
30    }
31
32    serde_json::from_str(&text).map_err(|e| {
33        let snippet = if text.len() > 512 {
34            format!("{}... ({} bytes total)", &text[..256], text.len())
35        } else {
36            text.clone()
37        };
38        format!("deserialize error (HTTP {status}): {e}: {snippet}")
39    })
40}
41
42pub struct HyperliquidClient {
43    client: Client,
44    user_address: Option<String>,
45    private_key: Option<String>,
46    rest_limiter: Arc<RestRateLimiter>,
47    address_limiter: Arc<AddressRateLimiter>,
48    ws_mux: crate::ws::WsMux,
49}
50
51impl Default for HyperliquidClient {
52    fn default() -> Self {
53        Self::new()
54    }
55}
56
57impl HyperliquidClient {
58    pub fn new() -> Self {
59        HyperliquidClient {
60            client: Client::new(),
61            user_address: None,
62            private_key: None,
63            rest_limiter: Arc::new(RestRateLimiter::new()),
64            address_limiter: Arc::new(AddressRateLimiter::new()),
65            ws_mux: crate::ws::WsMux::new(),
66        }
67    }
68
69    pub fn with_auth(user_address: impl Into<String>, private_key: String) -> Self {
70        HyperliquidClient {
71            client: Client::new(),
72            user_address: Some(user_address.into()),
73            private_key: Some(private_key),
74            rest_limiter: Arc::new(RestRateLimiter::new()),
75            address_limiter: Arc::new(AddressRateLimiter::new()),
76            ws_mux: crate::ws::WsMux::new(),
77        }
78    }
79
80    /// Configure rate limit budgets (rest_weight/min, address_requests).
81    /// Defaults: 1200 rest weight/min, 10000 address requests.
82    pub fn with_budgets(mut self, rest_weight: u32, addr_budget: u64) -> Self {
83        self.rest_limiter = Arc::new(RestRateLimiter::new_with_budget(rest_weight));
84        self.address_limiter = Arc::new(AddressRateLimiter::new_with_budget(addr_budget));
85        self
86    }
87
88    /// POST to the info endpoint, consuming `weight` from the REST rate-limit budget.
89    /// Returns `Err("rate_limited: ...")` immediately if budget is exhausted — no retry.
90    /// Callers should handle gracefully (skip cycle, retry later, etc.).
91    async fn info_post(
92        &self,
93        body: Value,
94        weight: u32,
95        call: &str,
96    ) -> Result<reqwest::Response, String> {
97        self.rest_limiter.acquire(weight).await.map_err(|e| {
98            format!(
99                "rate_limited: info_post ({call}) budget exhausted, retry_after_ms={}",
100                e.retry_after.as_millis()
101            )
102        })?;
103        self.client
104            .post(HYPERLIQUID_INFO_URL)
105            .json(&body)
106            .send()
107            .await
108            .map_err(|e| e.to_string())
109    }
110    fn require_user_address(&self) -> Result<String, String> {
111        self.user_address
112            .clone()
113            .ok_or_else(|| "user address required: use HyperliquidClient::with_auth".to_string())
114    }
115
116    fn require_private_key(&self) -> Result<&str, String> {
117        self.private_key
118            .as_deref()
119            .ok_or_else(|| "private key required: use HyperliquidClient::with_auth".to_string())
120    }
121
122    async fn get_asset_index(&self, symbol: &str) -> Result<usize, String> {
123        // `meta` is an "all other info" request → weight 20
124        let resp = self
125            .info_post(serde_json::json!({"type": "meta"}), 20, "get_asset_index")
126            .await?;
127        let meta: MetaResponse = parse_response(resp).await?;
128        meta.universe
129            .iter()
130            .position(|a| a.name == symbol)
131            .ok_or_else(|| format!("symbol {} not found", symbol))
132    }
133
134    async fn submit_signed_action(
135        &self,
136        action: Value,
137        vault_address: Option<&str>,
138    ) -> Result<Value, String> {
139        let private_key = self.require_private_key()?;
140        let nonce = std::time::SystemTime::now()
141            .duration_since(std::time::UNIX_EPOCH)
142            .unwrap()
143            .as_millis() as u64;
144
145        let (r, s, v) = sign_action(private_key, &action, vault_address, nonce)?;
146
147        let payload = serde_json::json!({
148            "action": action,
149            "nonce": nonce,
150            "signature": {"r": r, "s": s, "v": v},
151            "vaultAddress": null,
152            "expiresAfter": null
153        });
154
155        // Check both rate limiters non-blocking — fail fast, no retry.
156        self.rest_limiter.acquire(1).await.map_err(|e| {
157            format!(
158                "rate_limited: rest_weight exhausted, retry_after_ms={}",
159                e.retry_after.as_millis()
160            )
161        })?;
162        self.address_limiter.acquire(1, false).await.map_err(|e| {
163            format!(
164                "rate_limited: address quota exhausted, retry_after_ms={}",
165                e.retry_after.as_millis()
166            )
167        })?;
168
169        let resp = self
170            .client
171            .post(HYPERLIQUID_EXCHANGE_URL)
172            .json(&payload)
173            .send()
174            .await
175            .map_err(|e| e.to_string())?;
176
177        let status = resp.status();
178        if !status.is_success() {
179            let text = resp.text().await.map_err(|e| e.to_string())?;
180            return Err(format!("HTTP {status}: {text}"));
181        }
182
183        let body: Value = parse_response(resp).await?;
184        if body["status"].as_str() == Some("err") {
185            return Err(body["response"]
186                .as_str()
187                .unwrap_or("unknown error")
188                .to_string());
189        }
190        Ok(body)
191    }
192}
193
194// --- REST deserialization types ---
195
196#[derive(Deserialize)]
197struct MetaResponse {
198    universe: Vec<AssetInfo>,
199}
200
201#[derive(Deserialize)]
202struct AssetInfo {
203    name: String,
204    #[serde(rename = "szDecimals")]
205    sz_decimals: i32,
206}
207
208type MetaAndAssetCtxsResponse = (MetaResponse, Vec<RestAssetCtx>);
209
210#[derive(Deserialize)]
211#[serde(rename_all = "camelCase")]
212#[allow(dead_code)]
213struct RestAssetCtx {
214    open_interest: String,
215    funding: String,
216    mark_px: String,
217    day_ntl_vlm: String,
218    mid_px: Option<String>,
219    oracle_px: Option<String>,
220    premium: Option<String>,
221    prev_day_px: Option<String>,
222}
223
224#[derive(Deserialize)]
225#[serde(rename_all = "camelCase")]
226struct ClearinghouseStateResponse {
227    margin_summary: MarginSummary,
228    asset_positions: Vec<AssetPosition>,
229}
230
231#[derive(Deserialize)]
232#[serde(rename_all = "camelCase")]
233struct MarginSummary {
234    account_value: String,
235}
236
237#[derive(Deserialize)]
238struct AssetPosition {
239    position: PositionDetail,
240}
241
242#[derive(Deserialize)]
243#[serde(rename_all = "camelCase")]
244struct PositionDetail {
245    coin: String,
246    /// positive = long, negative = short
247    szi: String,
248    entry_px: Option<String>,
249}
250
251#[derive(Deserialize)]
252#[serde(rename_all = "camelCase")]
253struct RestOpenOrder {
254    coin: String,
255    side: String,
256    limit_px: String,
257    sz: String,
258    oid: i64,
259    orig_sz: String,
260    cloid: Option<String>,
261}
262
263// predictedFundings response: Vec<(coin, Vec<(venue, entry_or_null)>)>
264// The API returns null for venues that don't list the coin.
265type PredictedFundingsResponse = Vec<(String, Vec<(String, Option<PredictedFundingEntry>)>)>;
266
267#[derive(Deserialize)]
268#[serde(rename_all = "camelCase")]
269struct PredictedFundingEntry {
270    funding_rate: String,
271    next_funding_time: i64,
272}
273
274// --- WebSocket envelope and data shapes ---
275
276#[derive(Deserialize)]
277struct WsEnvelope {
278    channel: String,
279    #[serde(default)]
280    data: Value,
281}
282
283#[derive(Deserialize)]
284struct WsBook {
285    coin: String,
286    levels: Vec<Vec<WsLevel>>,
287    time: i64,
288}
289
290#[derive(Deserialize)]
291struct WsLevel {
292    px: String,
293    sz: String,
294}
295
296#[derive(Deserialize)]
297#[serde(rename_all = "camelCase")]
298struct WsAssetCtx {
299    coin: String,
300    ctx: WsPerpsCtx,
301}
302
303#[derive(Deserialize)]
304#[serde(rename_all = "camelCase")]
305struct WsPerpsCtx {
306    open_interest: String,
307    funding: String,
308    mark_px: String,
309    day_ntl_vlm: String,
310    mid_px: Option<String>,
311    oracle_px: Option<String>,
312    premium: Option<String>,
313    prev_day_px: Option<String>,
314}
315
316#[derive(Deserialize)]
317struct WsUserEvent {
318    liquidation: Option<WsLiquidation>,
319    fills: Option<Vec<WsUserFill>>,
320    funding: Option<WsFunding>,
321    spot_state: Option<WsSpotState>,
322}
323
324#[derive(Deserialize)]
325struct WsSpotState {
326    balances: Option<Vec<WsSpotBalance>>,
327}
328
329#[derive(Deserialize)]
330struct WsSpotBalance {
331    coin: String,
332    total: String,
333    hold: String,
334}
335
336#[derive(Deserialize)]
337struct WsLiquidation {
338    liquidated_user: String,
339    liquidated_ntl_pos: String,
340    liquidated_account_value: String,
341}
342
343#[derive(Deserialize)]
344struct WsUserFill {
345    coin: String,
346    px: String,
347    sz: String,
348    side: String,
349    time: i64,
350    oid: i64,
351    fee: String,
352    /// Client order ID assigned at placement, if any.
353    #[serde(default)]
354    cloid: Option<String>,
355}
356
357#[derive(Deserialize)]
358struct WsFunding {
359    time: i64,
360    coin: String,
361    usdc: String,
362}
363
364#[derive(Deserialize)]
365struct WsTrade {
366    coin: String,
367    side: String,
368    px: String,
369    sz: String,
370    time: i64,
371    tid: i64,
372}
373
374#[derive(Deserialize)]
375struct WsOrderUpdate {
376    order: WsOrderInfo,
377    status: String,
378    #[serde(rename = "statusTimestamp")]
379    status_timestamp: i64,
380}
381
382#[derive(Deserialize)]
383#[serde(rename_all = "camelCase")]
384struct WsOrderInfo {
385    coin: String,
386    side: String,
387    limit_px: String,
388    sz: String,
389    oid: i64,
390    orig_sz: String,
391    /// Client order ID assigned at placement, if any.
392    #[serde(default)]
393    cloid: Option<String>,
394}
395
396// --- WebSocket ledger update shapes (deposits / withdrawals) ---
397
398#[derive(Deserialize)]
399struct WsLedgerUpdates {
400    updates: Vec<WsLedgerEntry>,
401}
402
403#[derive(Deserialize)]
404struct WsLedgerEntry {
405    time: i64,
406    delta: WsLedgerDelta,
407}
408
409#[derive(Deserialize)]
410struct WsLedgerDelta {
411    #[serde(rename = "type")]
412    kind: String,
413    usdc: Option<String>,
414}
415
416// --- Helpers ---
417
418fn parse_decimal(s: &str) -> Option<Decimal> {
419    Decimal::from_str(s).ok()
420}
421
422fn keccak256(data: &[u8]) -> [u8; 32] {
423    use sha3::{Digest, Keccak256};
424    Keccak256::digest(data).into()
425}
426
427/// EIP-712 domain separator for Hyperliquid L1 actions (chainId=1337).
428fn hyperliquid_domain_separator() -> [u8; 32] {
429    let type_hash = keccak256(
430        b"EIP712Domain(string name,string version,uint256 chainId,address verifyingContract)",
431    );
432    let name_hash = keccak256(b"Exchange");
433    let version_hash = keccak256(b"1");
434    let mut chain_id = [0u8; 32];
435    chain_id[28..32].copy_from_slice(&1337u32.to_be_bytes());
436    let verifying_contract = [0u8; 32];
437
438    let mut data = [0u8; 160];
439    data[..32].copy_from_slice(&type_hash);
440    data[32..64].copy_from_slice(&name_hash);
441    data[64..96].copy_from_slice(&version_hash);
442    data[96..128].copy_from_slice(&chain_id);
443    data[128..160].copy_from_slice(&verifying_contract);
444    keccak256(&data)
445}
446
447/// Convert a `serde_json::Value` to msgpack bytes, preserving the JSON map key order.
448/// This avoids rmp_serde's HashMap-based serialization which reorders map keys.
449fn value_to_msgpack(val: &Value) -> Vec<u8> {
450    match val {
451        Value::Null => vec![0xc0],
452        Value::Bool(true) => vec![0xc3],
453        Value::Bool(false) => vec![0xc2],
454        Value::Number(n) => {
455            if let Some(i) = n.as_i64() {
456                if i >= 0 {
457                    if i <= 127 {
458                        vec![i as u8]
459                    } else if i <= 255 {
460                        vec![0xcc, i as u8]
461                    } else if i <= 65535 {
462                        let mut buf = vec![0xcd];
463                        buf.extend_from_slice(&(i as u16).to_be_bytes());
464                        buf
465                    } else if i <= 4294967295 {
466                        let mut buf = vec![0xce];
467                        buf.extend_from_slice(&(i as u32).to_be_bytes());
468                        buf
469                    } else {
470                        let mut buf = vec![0xcf];
471                        buf.extend_from_slice(&(i as u64).to_be_bytes());
472                        buf
473                    }
474                } else if i >= -32 {
475                    vec![0xe0 | (i as u8)]
476                } else if i >= -128 {
477                    vec![0xd0, i as i8 as u8]
478                } else if i >= -32768 {
479                    let mut buf = vec![0xd1];
480                    buf.extend_from_slice(&(i as i16).to_be_bytes());
481                    buf
482                } else if i >= -2147483648 {
483                    let mut buf = vec![0xd2];
484                    buf.extend_from_slice(&(i as i32).to_be_bytes());
485                    buf
486                } else {
487                    let mut buf = vec![0xd3];
488                    buf.extend_from_slice(&i.to_be_bytes());
489                    buf
490                }
491            } else if let Some(f) = n.as_f64() {
492                let mut buf = vec![0xcb];
493                buf.extend_from_slice(&f.to_be_bytes());
494                buf
495            } else {
496                let u = n.as_u64().unwrap();
497                if u <= 127 {
498                    vec![u as u8]
499                } else if u <= 255 {
500                    vec![0xcc, u as u8]
501                } else if u <= 65535 {
502                    let mut buf = vec![0xcd];
503                    buf.extend_from_slice(&(u as u16).to_be_bytes());
504                    buf
505                } else if u <= 4294967295 {
506                    let mut buf = vec![0xce];
507                    buf.extend_from_slice(&(u as u32).to_be_bytes());
508                    buf
509                } else {
510                    let mut buf = vec![0xcf];
511                    buf.extend_from_slice(&u.to_be_bytes());
512                    buf
513                }
514            }
515        }
516        Value::String(s) => {
517            let bytes = s.as_bytes();
518            let len = bytes.len();
519            let mut buf = Vec::new();
520            if len <= 31 {
521                buf.push(0xa0 | (len as u8));
522            } else if len <= 255 {
523                buf.push(0xd9);
524                buf.push(len as u8);
525            } else if len <= 65535 {
526                buf.push(0xda);
527                buf.extend_from_slice(&(len as u16).to_be_bytes());
528            } else {
529                buf.push(0xdb);
530                buf.extend_from_slice(&(len as u32).to_be_bytes());
531            }
532            buf.extend_from_slice(bytes);
533            buf
534        }
535        Value::Array(arr) => {
536            let len = arr.len();
537            let mut buf = Vec::new();
538            if len <= 15 {
539                buf.push(0x90 | (len as u8));
540            } else if len <= 65535 {
541                buf.push(0xdc);
542                buf.extend_from_slice(&(len as u16).to_be_bytes());
543            } else {
544                buf.push(0xdd);
545                buf.extend_from_slice(&(len as u32).to_be_bytes());
546            }
547            for item in arr {
548                buf.extend_from_slice(&value_to_msgpack(item));
549            }
550            buf
551        }
552        Value::Object(map) => {
553            let len = map.len();
554            let mut buf = Vec::new();
555            if len <= 15 {
556                buf.push(0x80 | (len as u8));
557            } else if len <= 65535 {
558                buf.push(0xde);
559                buf.extend_from_slice(&(len as u16).to_be_bytes());
560            } else {
561                buf.push(0xdf);
562                buf.extend_from_slice(&(len as u32).to_be_bytes());
563            }
564            for (key, value) in map {
565                buf.extend_from_slice(&value_to_msgpack(&Value::String(key.clone())));
566                buf.extend_from_slice(&value_to_msgpack(value));
567            }
568            buf
569        }
570    }
571}
572
573/// Convert action to msgpack bytes preserving JSON map key insertion order
574/// (matching Python's msgpack dict ordering).
575fn action_to_canonical_msgpack(action: &Value) -> Result<Vec<u8>, String> {
576    Ok(value_to_msgpack(action))
577}
578
579/// Build msgpack for a single order with Python SDK field order:
580/// a, b, p, s, r, t, c(opt)
581fn build_order_msgpack(
582    asset_idx: usize,
583    is_buy: bool,
584    price: &str,
585    size: &str,
586    reduce_only: bool,
587    order_kind: &str,
588    tif: &[u8],
589    cloid: Option<&str>,
590) -> Vec<u8> {
591    let field_count = if cloid.is_some() { 7 } else { 6 };
592    let mut buf = Vec::new();
593    buf.push(0x80 | (field_count as u8)); // fixmap
594
595    // "a": asset_idx
596    buf.extend_from_slice(&value_to_msgpack(&Value::String("a".to_string())));
597    buf.extend_from_slice(&value_to_msgpack(&Value::Number(serde_json::Number::from(asset_idx))));
598
599    // "b": is_buy
600    buf.extend_from_slice(&value_to_msgpack(&Value::String("b".to_string())));
601    buf.push(if is_buy { 0xc3 } else { 0xc2 });
602
603    // "p": price
604    buf.extend_from_slice(&value_to_msgpack(&Value::String("p".to_string())));
605    buf.extend_from_slice(&value_to_msgpack(&Value::String(price.to_string())));
606
607    // "s": size (Python SDK puts s before r)
608    buf.extend_from_slice(&value_to_msgpack(&Value::String("s".to_string())));
609    buf.extend_from_slice(&value_to_msgpack(&Value::String(size.to_string())));
610
611    // "r": reduce_only
612    buf.extend_from_slice(&value_to_msgpack(&Value::String("r".to_string())));
613    buf.push(if reduce_only { 0xc3 } else { 0xc2 });
614
615    // "t": { order_kind: { "tif": tif_str } }
616    buf.extend_from_slice(&value_to_msgpack(&Value::String("t".to_string())));
617    // Inner: fixmap(1) with order_kind key
618    buf.push(0x81);
619    buf.extend_from_slice(&value_to_msgpack(&Value::String(order_kind.to_string())));
620    // Inner-inner: fixmap(1) with "tif" key
621    buf.push(0x81);
622    buf.extend_from_slice(&value_to_msgpack(&Value::String("tif".to_string())));
623    buf.extend_from_slice(&value_to_msgpack(&Value::String(String::from_utf8_lossy(tif).to_string())));
624
625    // "c": cloid (optional, appended at END per Python SDK)
626    if let Some(c) = cloid {
627        buf.extend_from_slice(&value_to_msgpack(&Value::String("c".to_string())));
628        buf.extend_from_slice(&value_to_msgpack(&Value::String(c.to_string())));
629    }
630
631    buf
632}
633
634/// Sign using pre-built msgpack bytes (bypassing serde_json field ordering).
635fn sign_with_msgpack(
636    msgpack: &[u8],
637    private_key: &str,
638    nonce: u64,
639    vault_address: Option<&str>,
640) -> Result<(String, String, u8), String> {
641    use k256::ecdsa::SigningKey;
642
643    let mut data = msgpack.to_vec();
644    data.extend_from_slice(&nonce.to_be_bytes());
645    match vault_address {
646        None => data.push(0u8),
647        Some(addr) => {
648            data.push(1u8);
649            let addr_bytes = hex::decode(addr.trim_start_matches("0x"))
650                .map_err(|e| format!("invalid vault address: {}", e))?;
651            data.extend_from_slice(&addr_bytes);
652        }
653    }
654
655    let connection_id = keccak256(&data);
656    let agent_type_hash = keccak256(b"Agent(string source,bytes32 connectionId)");
657    let source_hash = keccak256(b"a");
658    let mut struct_data = [0u8; 96];
659    struct_data[..32].copy_from_slice(&agent_type_hash);
660    struct_data[32..64].copy_from_slice(&source_hash);
661    struct_data[64..96].copy_from_slice(&connection_id);
662    let struct_hash = keccak256(&struct_data);
663
664    let domain_sep = hyperliquid_domain_separator();
665    let mut final_data = Vec::with_capacity(66);
666    final_data.extend_from_slice(b"\x19\x01");
667    final_data.extend_from_slice(&domain_sep);
668    final_data.extend_from_slice(&struct_hash);
669    let final_hash = keccak256(&final_data);
670
671    let key_bytes = hex::decode(private_key.trim_start_matches("0x"))
672        .map_err(|e| format!("invalid private key: {}", e))?;
673    let signing_key =
674        SigningKey::from_bytes(key_bytes.as_slice().into()).map_err(|e| e.to_string())?;
675    let (sig, recovery_id) = signing_key
676        .sign_prehash_recoverable(&final_hash)
677        .map_err(|e| e.to_string())?;
678
679    let sig_bytes = sig.to_bytes();
680    let r = format!("0x{}", hex::encode(&sig_bytes[..32]));
681    let s = format!("0x{}", hex::encode(&sig_bytes[32..64]));
682    let v = 27u8 + recovery_id.to_byte();
683
684    Ok((r, s, v))
685}
686
687/// Signs a Hyperliquid exchange action using EIP-712.
688/// Returns (r, s, v) where r and s are "0x"-prefixed hex strings and v is 27 or 28.
689fn sign_action(
690    private_key: &str,
691    action: &Value,
692    vault_address: Option<&str>,
693    nonce: u64,
694) -> Result<(String, String, u8), String> {
695    use k256::ecdsa::SigningKey;
696
697    // Step 1: msgpack-encode the action preserving Python dict field order,
698    // then append nonce + vault flag.
699    let msgpack_bytes = action_to_canonical_msgpack(action)?;
700    let mut data = msgpack_bytes;
701    data.extend_from_slice(&nonce.to_be_bytes());
702    match vault_address {
703        None => data.push(0u8),
704        Some(addr) => {
705            data.push(1u8);
706            let addr_bytes = hex::decode(addr.trim_start_matches("0x"))
707                .map_err(|e| format!("invalid vault address: {}", e))?;
708            data.extend_from_slice(&addr_bytes);
709        }
710    }
711    let connection_id = keccak256(&data);
712
713    // Step 2: hash the Agent struct
714    let agent_type_hash = keccak256(b"Agent(string source,bytes32 connectionId)");
715    let source_hash = keccak256(b"a"); // "a" = mainnet
716    let mut struct_data = [0u8; 96];
717    struct_data[..32].copy_from_slice(&agent_type_hash);
718    struct_data[32..64].copy_from_slice(&source_hash);
719    struct_data[64..96].copy_from_slice(&connection_id);
720    let struct_hash = keccak256(&struct_data);
721
722    // Step 3: EIP-712 final hash
723    let domain_sep = hyperliquid_domain_separator();
724    let mut final_data = Vec::with_capacity(66);
725    final_data.extend_from_slice(b"\x19\x01");
726    final_data.extend_from_slice(&domain_sep);
727    final_data.extend_from_slice(&struct_hash);
728    let final_hash = keccak256(&final_data);
729
730    // Step 4: sign with secp256k1
731    let key_bytes = hex::decode(private_key.trim_start_matches("0x"))
732        .map_err(|e| format!("invalid private key: {}", e))?;
733    let signing_key =
734        SigningKey::from_bytes(key_bytes.as_slice().into()).map_err(|e| e.to_string())?;
735    let (sig, recovery_id) = signing_key
736        .sign_prehash_recoverable(&final_hash)
737        .map_err(|e| e.to_string())?;
738
739    let sig_bytes = sig.to_bytes();
740    let r = format!("0x{}", hex::encode(&sig_bytes[..32]));
741    let s = format!("0x{}", hex::encode(&sig_bytes[32..64]));
742    let v = 27u8 + recovery_id.to_byte();
743
744    Ok((r, s, v))
745}
746
747// --- Trait implementations ---
748
749#[allow(async_fn_in_trait)]
750impl guilder_abstraction::TestServer for HyperliquidClient {
751    /// Sends a lightweight allMids request; returns true if the server responds 200 OK.
752    async fn ping(&self) -> Result<bool, String> {
753        // allMids → weight 2
754        self.info_post(serde_json::json!({"type": "allMids"}), 2, "ping")
755            .await
756            .map(|r| r.status().is_success())
757    }
758
759    /// Hyperliquid has no dedicated server-time endpoint; returns local UTC ms.
760    async fn get_server_time(&self) -> Result<i64, String> {
761        Ok(std::time::SystemTime::now()
762            .duration_since(std::time::UNIX_EPOCH)
763            .map(|d| d.as_millis() as i64)
764            .unwrap_or(0))
765    }
766}
767
768#[allow(async_fn_in_trait)]
769impl guilder_abstraction::GetMarketData for HyperliquidClient {
770    /// Returns all perpetual asset names from Hyperliquid's meta endpoint.
771    async fn get_symbol(&self) -> Result<Vec<String>, String> {
772        // meta → weight 20
773        let resp = self
774            .info_post(serde_json::json!({"type": "meta"}), 20, "get_symbol")
775            .await?;
776        parse_response::<MetaResponse>(resp)
777            .await
778            .map(|r| r.universe.into_iter().map(|a| a.name).collect())
779    }
780
781    /// Returns the current open interest for `symbol` from metaAndAssetCtxs.
782    async fn get_open_interest(&self, symbol: String) -> Result<Decimal, String> {
783        // metaAndAssetCtxs → weight 20
784        let resp = self
785            .info_post(
786                serde_json::json!({"type": "metaAndAssetCtxs"}),
787                20,
788                "get_open_interest",
789            )
790            .await?;
791        let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
792            .await?
793            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
794        meta.universe
795            .iter()
796            .position(|a| a.name == symbol)
797            .and_then(|i| ctxs.get(i))
798            .and_then(|ctx| parse_decimal(&ctx.open_interest))
799            .ok_or_else(|| format!("symbol {} not found", symbol))
800    }
801
802    /// Returns a full AssetContext snapshot for `symbol` from metaAndAssetCtxs.
803    async fn get_asset_context(&self, symbol: String) -> Result<AssetContext, String> {
804        // metaAndAssetCtxs → weight 20
805        let resp = self
806            .info_post(
807                serde_json::json!({"type": "metaAndAssetCtxs"}),
808                20,
809                "get_asset_context",
810            )
811            .await?;
812        let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
813            .await?
814            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
815        let idx = meta
816            .universe
817            .iter()
818            .position(|a| a.name == symbol)
819            .ok_or_else(|| format!("symbol {} not found", symbol))?;
820        let ctx = ctxs
821            .get(idx)
822            .ok_or_else(|| format!("symbol {} not found", symbol))?;
823        Ok(AssetContext {
824            symbol,
825            open_interest: parse_decimal(&ctx.open_interest).ok_or("invalid open_interest")?,
826            funding_rate: parse_decimal(&ctx.funding).ok_or("invalid funding")?,
827            mark_price: parse_decimal(&ctx.mark_px).ok_or("invalid mark_px")?,
828            day_volume: parse_decimal(&ctx.day_ntl_vlm).ok_or("invalid day_ntl_vlm")?,
829            mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
830            oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
831            premium: ctx.premium.as_deref().and_then(parse_decimal),
832            prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
833            sz_decimals: meta.universe.get(idx).map(|a| a.sz_decimals).unwrap_or(0),
834        })
835    }
836
837    /// Fetches metaAndAssetCtxs once and returns all asset contexts in universe order.
838    /// Prefer this over repeated `get_asset_context` calls to avoid rate-limiting.
839    async fn get_all_asset_contexts(&self) -> Result<Vec<AssetContext>, String> {
840        // metaAndAssetCtxs → weight 20
841        let resp = self
842            .info_post(
843                serde_json::json!({"type": "metaAndAssetCtxs"}),
844                20,
845                "get_all_asset_contexts",
846            )
847            .await?;
848        let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
849            .await?
850            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
851        let mut result = Vec::with_capacity(meta.universe.len());
852        for (asset, ctx) in meta.universe.iter().zip(ctxs.iter()) {
853            let Some(open_interest) = parse_decimal(&ctx.open_interest) else {
854                continue;
855            };
856            let Some(funding_rate) = parse_decimal(&ctx.funding) else {
857                continue;
858            };
859            let Some(mark_price) = parse_decimal(&ctx.mark_px) else {
860                continue;
861            };
862            let Some(day_volume) = parse_decimal(&ctx.day_ntl_vlm) else {
863                continue;
864            };
865            result.push(AssetContext {
866                symbol: asset.name.clone(),
867                open_interest,
868                funding_rate,
869                mark_price,
870                day_volume,
871                mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
872                oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
873                premium: ctx.premium.as_deref().and_then(parse_decimal),
874                prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
875                sz_decimals: asset.sz_decimals,
876            });
877        }
878        Ok(result)
879    }
880
881    /// Returns the number of decimal places for order size for a symbol.
882    async fn get_sz_decimals(&self, symbol: String) -> Result<i32, String> {
883        let all = self.get_all_sz_decimals().await?;
884        all.get(&symbol)
885            .copied()
886            .ok_or_else(|| format!("symbol {} not found", symbol))
887    }
888
889    /// Returns sz_decimals for all symbols from the meta universe.
890    async fn get_all_sz_decimals(&self) -> Result<HashMap<String, i32>, String> {
891        // metaAndAssetCtxs → weight 20
892        let resp = self
893            .info_post(
894                serde_json::json!({"type": "metaAndAssetCtxs"}),
895                20,
896                "get_all_sz_decimals",
897            )
898            .await?;
899        let (meta, _) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
900            .await?
901            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
902        Ok(meta
903            .universe
904            .into_iter()
905            .map(|a| (a.name, a.sz_decimals))
906            .collect())
907    }
908
909    /// Returns a full L2 orderbook snapshot for `symbol` from the l2Book REST endpoint.
910    /// Levels are returned as individual `L2Update` items; all share the same `sequence` (timestamp).
911    async fn get_l2_orderbook(&self, symbol: String) -> Result<Vec<L2Update>, String> {
912        // l2Book → weight 2
913        let resp = self
914            .info_post(
915                serde_json::json!({"type": "l2Book", "coin": symbol}),
916                2,
917                "get_l2_orderbook",
918            )
919            .await?;
920        let book: Option<WsBook> = parse_response(resp).await?;
921        let book = match book {
922            Some(b) => b,
923            None => return Ok(vec![]),
924        };
925        let mut levels = Vec::new();
926        for level in book.levels.first().into_iter().flatten() {
927            if let (Some(price), Some(volume)) =
928                (parse_decimal(&level.px), parse_decimal(&level.sz))
929            {
930                levels.push(L2Update {
931                    symbol: book.coin.clone(),
932                    price,
933                    volume,
934                    side: Side::Ask,
935                    sequence: book.time,
936                });
937            }
938        }
939        for level in book.levels.get(1).into_iter().flatten() {
940            if let (Some(price), Some(volume)) =
941                (parse_decimal(&level.px), parse_decimal(&level.sz))
942            {
943                levels.push(L2Update {
944                    symbol: book.coin.clone(),
945                    price,
946                    volume,
947                    side: Side::Bid,
948                    sequence: book.time,
949                });
950            }
951        }
952        Ok(levels)
953    }
954
955    /// Returns the mid-price of `symbol` (e.g. "BTC") from allMids.
956    async fn get_price(&self, symbol: String) -> Result<Decimal, String> {
957        // allMids → weight 2
958        let resp = self
959            .info_post(serde_json::json!({"type": "allMids"}), 2, "get_price")
960            .await?;
961        parse_response::<HashMap<String, String>>(resp)
962            .await?
963            .get(&symbol)
964            .and_then(|s| parse_decimal(s))
965            .ok_or_else(|| format!("symbol {} not found", symbol))
966    }
967
968    /// Returns predicted funding rates for all symbols across all venues.
969    /// Null venue entries (unsupported coins) are silently skipped.
970    async fn get_predicted_fundings(&self) -> Result<Vec<PredictedFunding>, String> {
971        // predictedFundings → weight 20
972        let resp = self
973            .info_post(
974                serde_json::json!({"type": "predictedFundings"}),
975                20,
976                "get_predicted_fundings",
977            )
978            .await?;
979        let data: PredictedFundingsResponse = parse_response(resp).await?;
980        let mut result = Vec::new();
981        for (symbol, venues) in data {
982            for (venue, entry) in venues {
983                let Some(entry) = entry else { continue };
984                if let Some(funding_rate) = parse_decimal(&entry.funding_rate) {
985                    result.push(PredictedFunding {
986                        symbol: symbol.clone(),
987                        venue,
988                        funding_rate,
989                        next_funding_time_ms: entry.next_funding_time,
990                    });
991                }
992            }
993        }
994        Ok(result)
995    }
996}
997
998#[allow(async_fn_in_trait)]
999impl guilder_abstraction::ManageOrder for HyperliquidClient {
1000    /// Places an order on Hyperliquid. Requires `with_auth`. Returns an `OrderPlacement` with
1001    /// the exchange-assigned order ID. Market orders are submitted as aggressive limit orders (IOC).
1002    ///
1003    /// If `cloid` is provided, Hyperliquid attaches it to the order lifecycle — fills and order
1004    /// updates will carry the same cloid back, enabling end-to-end intent tracing without a
1005    /// separate order_id mapping.
1006    async fn place_order(
1007        &self,
1008        symbol: String,
1009        side: OrderSide,
1010        price: Decimal,
1011        volume: Decimal,
1012        order_type: OrderType,
1013        time_in_force: TimeInForce,
1014        cloid: Option<String>,
1015    ) -> Result<OrderPlacement, String> {
1016        // Rate limiting is handled in submit_signed_action (non-blocking).
1017        let asset_idx = self.get_asset_index(&symbol).await?;
1018        let is_buy = matches!(side, OrderSide::Buy);
1019
1020        let tif_str = match time_in_force {
1021            TimeInForce::Gtc => "Gtc",
1022            TimeInForce::Ioc => "Ioc",
1023            TimeInForce::Fok => "Fok",
1024        };
1025
1026        // Market orders are IOC limit orders at a wide price
1027        let (order_kind, tif_bytes) = match order_type {
1028            OrderType::Limit => ("limit", tif_str.as_bytes()),
1029            OrderType::Market => ("limit", b"Ioc".as_slice()),
1030        };
1031
1032        let price_str = price.normalize().to_string();
1033        let size_str = volume.normalize().to_string();
1034
1035        let cloid_hex = cloid.clone();
1036
1037        // Build msgpack with Python SDK field order (matching the Python SDK's
1038        // msgpack output). The server hashes the msgpack for signature verification,
1039        // and Python preserves dict insertion order.
1040        let order_msgpack = build_order_msgpack(
1041            asset_idx,
1042            is_buy,
1043            &price_str,
1044            &size_str,
1045            false, // reduce_only
1046            order_kind,
1047            tif_bytes,
1048            cloid_hex.as_deref(),
1049        );
1050
1051        // Build the action-level msgpack with Python SDK field order (insertion order):
1052        // type → orders → grouping
1053        let mut action_msgpack = Vec::new();
1054        action_msgpack.push(0x83); // fixmap(3)
1055        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("type".to_string())));
1056        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("order".to_string())));
1057        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("orders".to_string())));
1058        action_msgpack.push(0x91); // fixarray(1)
1059        action_msgpack.extend_from_slice(&order_msgpack);
1060        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("grouping".to_string())));
1061        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("na".to_string())));
1062
1063        // Build JSON with official SDK field order (alphabetical):
1064        // grouping → orders → type
1065        let order_type_json = match order_type {
1066            OrderType::Limit => format!(r#"{{"limit":{{"tif":"{tif_str}"}}}}"#),
1067            OrderType::Market => r#"{"limit":{"tif":"Ioc"}}"#.to_string(),
1068        };
1069
1070        let cloid_json = if let Some(ref c) = cloid_hex {
1071            format!(r#","c":"{c}""#)
1072        } else {
1073            String::new()
1074        };
1075
1076        let action_json_str = format!(
1077            r#"{{"type":"order","orders":[{{"a":{asset_idx},"b":{is_buy},"p":"{price}","s":"{size}","r":false,"t":{order_type_json}{cloid_json}}}],"grouping":"na"}}"#,
1078            price = price_str,
1079            size = size_str,
1080        );
1081
1082        // Sign using the canonical msgpack (matching Python's field order)
1083        let private_key = self.require_private_key()?;
1084        let nonce = std::time::SystemTime::now()
1085            .duration_since(std::time::UNIX_EPOCH)
1086            .unwrap()
1087            .as_millis() as u64;
1088
1089        let (r, s, v) = sign_with_msgpack(&action_msgpack, private_key, nonce, None)?;
1090
1091        let payload_str = format!(
1092            r#"{{"action":{},"nonce":{},"signature":{{"r":"{}","s":"{}","v":{}}},"vaultAddress":null,"expiresAfter":null}}"#,
1093            action_json_str,
1094            nonce,
1095            r, s, v
1096        );
1097
1098        self.rest_limiter.acquire(1).await.map_err(|e| {
1099            format!(
1100                "rate_limited: rest_weight exhausted, retry_after_ms={}",
1101                e.retry_after.as_millis()
1102            )
1103        })?;
1104        self.address_limiter.acquire(1, false).await.map_err(|e| {
1105            format!(
1106                "rate_limited: address quota exhausted, retry_after_ms={}",
1107                e.retry_after.as_millis()
1108            )
1109        })?;
1110
1111        let resp = self
1112            .client
1113            .post(HYPERLIQUID_EXCHANGE_URL)
1114            .header("Content-Type", "application/json")
1115            .body(payload_str)
1116            .send()
1117            .await
1118            .map_err(|e| e.to_string())?;
1119
1120        let status = resp.status();
1121        if !status.is_success() {
1122            let text = resp.text().await.map_err(|e| e.to_string())?;
1123            return Err(format!("HTTP {status}: {text}"));
1124        }
1125
1126        let body: Value = parse_response(resp).await?;
1127        if body["status"].as_str() == Some("err") {
1128            return Err(body["response"]
1129                .as_str()
1130                .unwrap_or("unknown error")
1131                .to_string());
1132        }
1133        let statuses = &body["response"]["data"]["statuses"][0];
1134
1135        let (oid, returned_cloid, timestamp_ms) = if let Some(resting) = statuses.get("resting") {
1136            let oid = resting["oid"]
1137                .as_i64()
1138                .ok_or_else(|| format!("resting status missing oid: {}", body))?;
1139            let returned_cloid = resting["cloid"].as_str().map(|s: &str| s.to_string());
1140            // resting doesn't include a timestamp
1141            let ts = std::time::SystemTime::now()
1142                .duration_since(std::time::UNIX_EPOCH)
1143                .unwrap()
1144                .as_millis() as i64;
1145            (oid, returned_cloid, ts)
1146        } else if let Some(filled) = statuses.get("filled") {
1147            let oid = filled["oid"]
1148                .as_i64()
1149                .ok_or_else(|| format!("filled status missing oid: {}", body))?;
1150            // filled doesn't include cloid
1151            let ts = std::time::SystemTime::now()
1152                .duration_since(std::time::UNIX_EPOCH)
1153                .unwrap()
1154                .as_millis() as i64;
1155            (oid, None, ts)
1156        } else if let Some(error) = statuses.get("error") {
1157            return Err(error
1158                .as_str()
1159                .unwrap_or("order rejected with unknown error")
1160                .to_string());
1161        } else {
1162            return Err(format!("unexpected order status: {}", body));
1163        };
1164
1165        Ok(OrderPlacement {
1166            order_id: oid,
1167            symbol,
1168            side,
1169            price,
1170            quantity: volume,
1171            timestamp_ms,
1172            cloid: returned_cloid.or(cloid),
1173        })
1174    }
1175
1176    /// Modifies price and size of an existing order by its order ID. Requires `with_auth`.
1177    /// Fetches the order's current coin and side before submitting the modify action.
1178    async fn change_order_by_cloid(
1179        &self,
1180        cloid: i64,
1181        price: Decimal,
1182        volume: Decimal,
1183    ) -> Result<i64, String> {
1184        let user = self.require_user_address()?;
1185
1186        // openOrders → weight 20; get_asset_index → meta weight 20
1187        let resp = self
1188            .info_post(
1189                serde_json::json!({"type": "openOrders", "user": user}),
1190                20,
1191                "change_order_by_cloid",
1192            )
1193            .await?;
1194        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1195        let order = orders
1196            .iter()
1197            .find(|o| o.oid == cloid)
1198            .ok_or_else(|| format!("order {} not found", cloid))?;
1199
1200        let asset_idx = self.get_asset_index(&order.coin).await?;
1201        let is_buy = order.side == "B";
1202
1203        let action = serde_json::json!({
1204            "type": "batchModify",
1205            "modifies": [{
1206                "oid": cloid,
1207                "order": {
1208                    "a": asset_idx,
1209                    "b": is_buy,
1210                    "p": price.to_string(),
1211                    "s": volume.to_string(),
1212                    "r": false,
1213                    "t": {"limit": {"tif": "Gtc"}}
1214                }
1215            }]
1216        });
1217
1218        self.submit_signed_action(action, None).await?;
1219        Ok(cloid)
1220    }
1221
1222    /// Cancels a single order by its client order ID (cloid). Requires `with_auth`.
1223    /// Fetches open orders to resolve the coin/asset for the matching cloid before cancelling.
1224    async fn cancel_order_by_cloid(&self, cloid: String) -> Result<(), String> {
1225        let user = self.require_user_address()?;
1226
1227        // openOrders → weight 20
1228        let resp = self
1229            .info_post(
1230                serde_json::json!({"type": "openOrders", "user": user}),
1231                20,
1232                "cancel_order_by_cloid",
1233            )
1234            .await?;
1235        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1236        let order = orders
1237            .iter()
1238            .find(|o| o.cloid.as_ref() == Some(&cloid))
1239            .ok_or_else(|| format!("order with cloid {} not found", cloid))?;
1240
1241        let asset_idx = self.get_asset_index(&order.coin).await?;
1242        let action = serde_json::json!({
1243            "type": "cancel",
1244            "cancels": [{"a": asset_idx, "cloid": cloid}]
1245        });
1246
1247        self.submit_signed_action(action, None).await?;
1248        Ok(())
1249    }
1250
1251    /// Cancels all open orders. Requires `with_auth`.
1252    /// Fetches all open orders and submits a batch cancel in a single signed request.
1253    async fn cancel_all_order(&self) -> Result<bool, String> {
1254        let user = self.require_user_address()?;
1255
1256        // openOrders → weight 20
1257        let resp = self
1258            .info_post(
1259                serde_json::json!({"type": "openOrders", "user": user}),
1260                20,
1261                "cancel_all_order",
1262            )
1263            .await?;
1264        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1265        if orders.is_empty() {
1266            return Ok(true);
1267        }
1268
1269        // meta → weight 20
1270        let meta_resp = self
1271            .info_post(serde_json::json!({"type": "meta"}), 20, "cancel_all_order")
1272            .await?;
1273        let meta: MetaResponse = parse_response(meta_resp).await?;
1274
1275        let cancels: Vec<Value> = orders
1276            .iter()
1277            .filter_map(|o| {
1278                let asset_idx = meta.universe.iter().position(|a| a.name == o.coin)?;
1279                Some(serde_json::json!({"a": asset_idx, "o": o.oid}))
1280            })
1281            .collect();
1282
1283        let action = serde_json::json!({"type": "cancel", "cancels": cancels});
1284        self.submit_signed_action(action, None).await?;
1285        Ok(true)
1286    }
1287}
1288
1289#[allow(async_fn_in_trait)]
1290impl guilder_abstraction::SubscribeMarketData for HyperliquidClient {
1291    fn subscribe_l2_update(&self, symbol: String) -> BoxStream<Result<L2Update, String>> {
1292        let sub = serde_json::json!({
1293            "method": "subscribe",
1294            "subscription": {"type": "l2Book", "coin": symbol.clone()}
1295        });
1296        let key = crate::ws::SubKey {
1297            channel: "l2Book".to_string(),
1298            routing_key: symbol,
1299        };
1300        let stream = self.ws_mux.subscribe(key, sub);
1301        Box::pin(async_stream::stream! {
1302            for await msg in stream {
1303                let Ok(env) = serde_json::from_str::<WsEnvelope>(&msg) else {
1304                    continue;
1305                };
1306                if env.channel != "l2Book" {
1307                    continue;
1308                }
1309                let Ok(book) = serde_json::from_value::<WsBook>(env.data) else {
1310                    continue;
1311                };
1312                for level in book.levels.first().into_iter().flatten() {
1313                    if let (Some(price), Some(volume)) =
1314                        (parse_decimal(&level.px), parse_decimal(&level.sz))
1315                    {
1316                        yield Ok(L2Update {
1317                            symbol: book.coin.clone(),
1318                            price,
1319                            volume,
1320                            side: Side::Ask,
1321                            sequence: book.time,
1322                        });
1323                    }
1324                }
1325                for level in book.levels.get(1).into_iter().flatten() {
1326                    if let (Some(price), Some(volume)) =
1327                        (parse_decimal(&level.px), parse_decimal(&level.sz))
1328                    {
1329                        yield Ok(L2Update {
1330                            symbol: book.coin.clone(),
1331                            price,
1332                            volume,
1333                            side: Side::Bid,
1334                            sequence: book.time,
1335                        });
1336                    }
1337                }
1338            }
1339        })
1340    }
1341
1342    fn subscribe_asset_context(&self, symbol: String) -> BoxStream<Result<AssetContext, String>> {
1343        let sub = serde_json::json!({
1344            "method": "subscribe",
1345            "subscription": {"type": "activeAssetCtx", "coin": symbol.clone()}
1346        });
1347        let key = crate::ws::SubKey {
1348            channel: "activeAssetCtx".to_string(),
1349            routing_key: symbol,
1350        };
1351        let stream = self.ws_mux.subscribe(key, sub);
1352        Box::pin(async_stream::stream! {
1353            for await msg in stream {
1354                let Ok(env) = serde_json::from_str::<WsEnvelope>(&msg) else {
1355                    continue;
1356                };
1357                if env.channel != "activeAssetCtx" {
1358                    continue;
1359                }
1360                let Ok(update) = serde_json::from_value::<WsAssetCtx>(env.data) else {
1361                    continue;
1362                };
1363                let ctx = &update.ctx;
1364                let (Some(open_interest), Some(funding_rate), Some(mark_price), Some(day_volume)) = (
1365                    parse_decimal(&ctx.open_interest),
1366                    parse_decimal(&ctx.funding),
1367                    parse_decimal(&ctx.mark_px),
1368                    parse_decimal(&ctx.day_ntl_vlm),
1369                ) else {
1370                    continue;
1371                };
1372                yield Ok(AssetContext {
1373                    symbol: update.coin,
1374                    open_interest,
1375                    funding_rate,
1376                    mark_price,
1377                    day_volume,
1378                    mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
1379                    oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
1380                    premium: ctx.premium.as_deref().and_then(parse_decimal),
1381                    prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
1382                    // szDecimals is static metadata, not streamed via activeAssetCtx WS
1383                    sz_decimals: 0,
1384                });
1385            }
1386        })
1387    }
1388
1389    fn subscribe_liquidation(&self, user: String) -> BoxStream<Result<Liquidation, String>> {
1390        let sub = serde_json::json!({
1391            "method": "subscribe",
1392            "subscription": {"type": "userEvents", "user": user.clone()}
1393        });
1394        let key = crate::ws::SubKey {
1395            channel: "user".to_string(),
1396            routing_key: user,
1397        };
1398        let raw_stream = self.ws_mux.subscribe(key, sub);
1399        Box::pin(
1400            raw_stream
1401                .filter_map(|text| async move {
1402                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1403                        return None;
1404                    };
1405                    if env.channel != "user" {
1406                        return None;
1407                    }
1408                    let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1409                        return None;
1410                    };
1411                    let liq = event.liquidation?;
1412                    let (Some(notional_position), Some(account_value)) = (
1413                        parse_decimal(&liq.liquidated_ntl_pos),
1414                        parse_decimal(&liq.liquidated_account_value),
1415                    ) else {
1416                        return None;
1417                    };
1418                    let item = Liquidation {
1419                        symbol: String::new(),
1420                        side: OrderSide::Sell,
1421                        liquidated_user: liq.liquidated_user,
1422                        notional_position,
1423                        account_value,
1424                    };
1425                    Some(stream::iter(vec![Ok(item)]))
1426                })
1427                .flatten(),
1428        )
1429    }
1430
1431    fn subscribe_fill(&self, symbol: String) -> BoxStream<Result<Fill, String>> {
1432        let sub = serde_json::json!({
1433            "method": "subscribe",
1434            "subscription": {"type": "trades", "coin": symbol.clone()}
1435        });
1436        let key = crate::ws::SubKey {
1437            channel: "trades".to_string(),
1438            routing_key: symbol,
1439        };
1440        let stream = self.ws_mux.subscribe(key, sub);
1441        Box::pin(async_stream::stream! {
1442            for await msg in stream {
1443                let Ok(env) = serde_json::from_str::<WsEnvelope>(&msg) else {
1444                    continue;
1445                };
1446                if env.channel != "trades" {
1447                    continue;
1448                }
1449                let Ok(trades) = serde_json::from_value::<Vec<WsTrade>>(env.data) else {
1450                    continue;
1451                };
1452                for trade in trades {
1453                    let side = if trade.side == "B" {
1454                        OrderSide::Buy
1455                    } else {
1456                        OrderSide::Sell
1457                    };
1458                    let price = parse_decimal(&trade.px);
1459                    let volume = parse_decimal(&trade.sz);
1460                    if let (Some(price), Some(volume)) = (price, volume) {
1461                        yield Ok(Fill {
1462                            symbol: trade.coin,
1463                            price,
1464                            volume,
1465                            side,
1466                            timestamp_ms: trade.time,
1467                            trade_id: trade.tid,
1468                        });
1469                    }
1470                }
1471            }
1472        })
1473    }
1474}
1475
1476#[allow(async_fn_in_trait)]
1477impl guilder_abstraction::GetAccountSnapshot for HyperliquidClient {
1478    /// Returns open positions from `clearinghouseState`. Requires `with_auth`.
1479    /// Zero-size positions are filtered out. Positive `szi` = long, negative = short.
1480    async fn get_positions(&self) -> Result<Vec<Position>, String> {
1481        let user = self.require_user_address()?;
1482        // clearinghouseState → weight 2
1483        let resp = self
1484            .info_post(
1485                serde_json::json!({"type": "clearinghouseState", "user": user}),
1486                2,
1487                "get_positions",
1488            )
1489            .await?;
1490        let state: ClearinghouseStateResponse = parse_response(resp).await?;
1491
1492        Ok(state
1493            .asset_positions
1494            .into_iter()
1495            .filter_map(|ap| {
1496                let p = ap.position;
1497                let size = parse_decimal(&p.szi)?;
1498                if size.is_zero() {
1499                    return None;
1500                }
1501                let entry_price = p
1502                    .entry_px
1503                    .as_deref()
1504                    .and_then(parse_decimal)
1505                    .unwrap_or_default();
1506                let side = if size > Decimal::ZERO {
1507                    OrderSide::Buy
1508                } else {
1509                    OrderSide::Sell
1510                };
1511                Some(Position {
1512                    symbol: p.coin,
1513                    side,
1514                    size: size.abs(),
1515                    entry_price,
1516                })
1517            })
1518            .collect())
1519    }
1520
1521    /// Returns resting orders from Hyperliquid's `openOrders` endpoint. Requires `with_auth`.
1522    /// `filled_quantity` is derived as `origSz - sz` (original size minus remaining size).
1523    async fn get_open_orders(&self) -> Result<Vec<OpenOrder>, String> {
1524        let user = self.require_user_address()?;
1525        // openOrders → weight 20
1526        let resp = self
1527            .info_post(
1528                serde_json::json!({"type": "openOrders", "user": user}),
1529                20,
1530                "get_open_orders",
1531            )
1532            .await?;
1533        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1534
1535        Ok(orders
1536            .into_iter()
1537            .filter_map(|o| {
1538                let price = parse_decimal(&o.limit_px)?;
1539                let quantity = parse_decimal(&o.orig_sz)?;
1540                let remaining = parse_decimal(&o.sz)?;
1541                let filled_quantity = quantity - remaining;
1542                let side = if o.side == "B" {
1543                    OrderSide::Buy
1544                } else {
1545                    OrderSide::Sell
1546                };
1547                Some(OpenOrder {
1548                    order_id: o.oid,
1549                    symbol: o.coin,
1550                    side,
1551                    price,
1552                    quantity,
1553                    filled_quantity,
1554                })
1555            })
1556            .collect())
1557    }
1558
1559    /// Returns total account value (collateral) from `clearinghouseState`. Requires `with_auth`.
1560    async fn get_collateral(&self) -> Result<Decimal, String> {
1561        let user = self.require_user_address()?;
1562        // clearinghouseState → weight 2
1563        let resp = self
1564            .info_post(
1565                serde_json::json!({"type": "clearinghouseState", "user": user}),
1566                2,
1567                "get_collateral",
1568            )
1569            .await?;
1570        let state: ClearinghouseStateResponse = parse_response(resp).await?;
1571        parse_decimal(&state.margin_summary.account_value)
1572            .ok_or_else(|| "invalid account value".to_string())
1573    }
1574
1575    /// Returns all spot wallet balances from `spotState`. Requires `with_auth`.
1576    async fn get_spot_balance(&self) -> Result<Vec<guilder_abstraction::Balance>, String> {
1577        let user = self.require_user_address()?;
1578        // spotClearinghouseState → weight 15
1579        let resp = self
1580            .info_post(
1581                serde_json::json!({"type": "spotClearinghouseState", "user": user}),
1582                15,
1583                "get_spot_balance",
1584            )
1585            .await?;
1586
1587        #[derive(Deserialize)]
1588        struct SpotStateResponse {
1589            balances: Vec<SpotBalance>,
1590        }
1591
1592        #[allow(dead_code)]
1593        #[derive(Deserialize)]
1594        struct SpotBalance {
1595            coin: String,
1596            total: String,
1597            hold: String,
1598            #[serde(default)]
1599            token: Option<i32>,
1600            #[serde(default)]
1601            #[serde(rename = "entryNtl")]
1602            entry_ntl: Option<String>,
1603        }
1604
1605        let state: SpotStateResponse = parse_response(resp).await?;
1606
1607        state
1608            .balances
1609            .into_iter()
1610            .map(|balance| {
1611                let total = parse_decimal(&balance.total)
1612                    .ok_or_else(|| "invalid total balance".to_string())?;
1613                let locked = parse_decimal(&balance.hold)
1614                    .ok_or_else(|| "invalid hold balance".to_string())?;
1615                let available = total - locked;
1616
1617                Ok(guilder_abstraction::Balance {
1618                    coin: balance.coin,
1619                    total,
1620                    available,
1621                    locked,
1622                })
1623            })
1624            .collect()
1625    }
1626
1627    /// Returns clearing house collateral balance for an asset. Currently returns the total collateral.
1628    /// Requires `with_auth`.
1629    async fn get_collateral_balance(
1630        &self,
1631        asset: String,
1632    ) -> Result<guilder_abstraction::Balance, String> {
1633        if asset.to_uppercase() != "USDC" {
1634            return Err(format!("only USDC collateral is supported, got {}", asset));
1635        }
1636
1637        let total = self.get_collateral().await?;
1638        Ok(guilder_abstraction::Balance {
1639            coin: "USDC".to_string(),
1640            total,
1641            available: total,
1642            locked: Decimal::ZERO,
1643        })
1644    }
1645
1646    /// Returns the user's address-level API rate limit budget.
1647    /// Queries Hyperliquid's `userRateLimit` info endpoint for authoritative server-side counts.
1648    async fn get_user_rate_limit(&self) -> Result<guilder_abstraction::UserRateLimit, String> {
1649        let user = self.require_user_address()?;
1650        let resp = self
1651            .info_post(
1652                serde_json::json!({"type": "userRateLimit", "user": user}),
1653                20,
1654                "get_user_rate_limit",
1655            )
1656            .await?;
1657        let val = parse_response::<Value>(resp).await?;
1658
1659        let cumulative_volume = val["cumVlm"]
1660            .as_str()
1661            .and_then(parse_decimal)
1662            .ok_or_else(|| "missing or invalid cumVlm".to_string())?;
1663        let requests_used = val["nRequestsUsed"]
1664            .as_i64()
1665            .ok_or_else(|| "missing or invalid nRequestsUsed".to_string())?;
1666        let requests_cap = val["nRequestsCap"]
1667            .as_i64()
1668            .ok_or_else(|| "missing or invalid nRequestsCap".to_string())?;
1669        let requests_surplus = val["nRequestsSurplus"]
1670            .as_i64()
1671            .ok_or_else(|| "missing or invalid nRequestsSurplus".to_string())?;
1672
1673        Ok(guilder_abstraction::UserRateLimit {
1674            cumulative_volume,
1675            requests_used,
1676            requests_cap,
1677            requests_surplus,
1678        })
1679    }
1680}
1681
1682#[allow(async_fn_in_trait)]
1683impl guilder_abstraction::SubscribeUserEvents for HyperliquidClient {
1684    fn subscribe_user_fills(&self) -> BoxStream<Result<UserFill, String>> {
1685        let Some(addr) = self.user_address.as_ref() else {
1686            return Box::pin(stream::empty());
1687        };
1688        let addr_str = addr.clone();
1689        let sub = serde_json::json!({
1690            "method": "subscribe",
1691            "subscription": {"type": "userEvents", "user": addr_str.clone()}
1692        });
1693        let key = crate::ws::SubKey {
1694            channel: "user".to_string(),
1695            routing_key: addr_str,
1696        };
1697        let raw_stream = self.ws_mux.subscribe(key, sub);
1698        Box::pin(
1699            raw_stream
1700                .filter_map(|text| async move {
1701                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1702                        return None;
1703                    };
1704                    if env.channel != "user" {
1705                        return None;
1706                    }
1707                    let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1708                        return None;
1709                    };
1710                    let items: Vec<_> = event
1711                        .fills
1712                        .unwrap_or_default()
1713                        .into_iter()
1714                        .filter_map(|fill| {
1715                            let side = if fill.side == "B" {
1716                                OrderSide::Buy
1717                            } else {
1718                                OrderSide::Sell
1719                            };
1720                            let price = parse_decimal(&fill.px)?;
1721                            let quantity = parse_decimal(&fill.sz)?;
1722                            let fee_usd = parse_decimal(&fill.fee)?;
1723                            Some(UserFill {
1724                                order_id: fill.oid,
1725                                symbol: fill.coin,
1726                                side,
1727                                price,
1728                                quantity,
1729                                fee_usd,
1730                                timestamp_ms: fill.time,
1731                                cloid: fill.cloid,
1732                            })
1733                        })
1734                        .collect();
1735                    if items.is_empty() {
1736                        None
1737                    } else {
1738                        Some(stream::iter(items.into_iter().map(Ok)))
1739                    }
1740                })
1741                .flatten(),
1742        )
1743    }
1744
1745    fn subscribe_order_updates(&self) -> BoxStream<Result<OrderUpdate, String>> {
1746        let Some(addr) = self.user_address.as_ref() else {
1747            return Box::pin(stream::empty());
1748        };
1749        let addr_str = addr.clone();
1750        let sub = serde_json::json!({
1751            "method": "subscribe",
1752            "subscription": {"type": "orderUpdates", "user": addr_str.clone()}
1753        });
1754        let key = crate::ws::SubKey {
1755            channel: "orderUpdates".to_string(),
1756            routing_key: addr_str,
1757        };
1758        let raw_stream = self.ws_mux.subscribe(key, sub);
1759        Box::pin(
1760            raw_stream
1761                .filter_map(|text| async move {
1762                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1763                        return None;
1764                    };
1765                    if env.channel != "orderUpdates" {
1766                        return None;
1767                    }
1768                    let Ok(updates) = serde_json::from_value::<Vec<WsOrderUpdate>>(env.data) else {
1769                        return None;
1770                    };
1771                    let items: Vec<_> = updates
1772                        .into_iter()
1773                        .map(|upd| {
1774                            let status = match upd.status.as_str() {
1775                                "open" => OrderStatus::Placed,
1776                                "filled" => OrderStatus::Filled,
1777                                "canceled" | "cancelled" => OrderStatus::Cancelled,
1778                                _ => OrderStatus::PartiallyFilled,
1779                            };
1780                            let side = if upd.order.side == "B" {
1781                                OrderSide::Buy
1782                            } else {
1783                                OrderSide::Sell
1784                            };
1785                            OrderUpdate {
1786                                order_id: upd.order.oid,
1787                                symbol: upd.order.coin,
1788                                status,
1789                                side: Some(side),
1790                                price: parse_decimal(&upd.order.limit_px),
1791                                quantity: parse_decimal(&upd.order.orig_sz),
1792                                remaining_quantity: parse_decimal(&upd.order.sz),
1793                                timestamp_ms: upd.status_timestamp,
1794                                cloid: upd.order.cloid,
1795                            }
1796                        })
1797                        .collect();
1798                    if items.is_empty() {
1799                        None
1800                    } else {
1801                        Some(stream::iter(items.into_iter().map(Ok)))
1802                    }
1803                })
1804                .flatten(),
1805        )
1806    }
1807
1808    fn subscribe_funding_payments(&self) -> BoxStream<Result<FundingPayment, String>> {
1809        let Some(addr) = self.user_address.as_ref() else {
1810            return Box::pin(stream::empty());
1811        };
1812        let addr_str = addr.clone();
1813        let sub = serde_json::json!({
1814            "method": "subscribe",
1815            "subscription": {"type": "userEvents", "user": addr_str.clone()}
1816        });
1817        let key = crate::ws::SubKey {
1818            channel: "user".to_string(),
1819            routing_key: addr_str,
1820        };
1821        let raw_stream = self.ws_mux.subscribe(key, sub);
1822        Box::pin(
1823            raw_stream
1824                .filter_map(|text| async move {
1825                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1826                        return None;
1827                    };
1828                    if env.channel != "user" {
1829                        return None;
1830                    }
1831                    let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1832                        return None;
1833                    };
1834                    let funding = event.funding?;
1835                    let amount_usd = parse_decimal(&funding.usdc)?;
1836                    let item = FundingPayment {
1837                        symbol: funding.coin,
1838                        amount_usd,
1839                        timestamp_ms: funding.time,
1840                    };
1841                    Some(stream::iter(vec![Ok(item)]))
1842                })
1843                .flatten(),
1844        )
1845    }
1846
1847    fn subscribe_deposits(&self) -> BoxStream<Result<Deposit, String>> {
1848        let Some(addr) = self.user_address.as_ref() else {
1849            return Box::pin(stream::empty());
1850        };
1851        let addr_str = addr.clone();
1852        let sub = serde_json::json!({
1853            "method": "subscribe",
1854            "subscription": {"type": "userNonFundingLedgerUpdates", "user": addr_str.clone()}
1855        });
1856        let key = crate::ws::SubKey {
1857            channel: "userNonFundingLedgerUpdates".to_string(),
1858            routing_key: addr_str,
1859        };
1860        let raw_stream = self.ws_mux.subscribe(key, sub);
1861        Box::pin(
1862            raw_stream
1863                .filter_map(|text| async move {
1864                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1865                        return None;
1866                    };
1867                    if env.channel != "userNonFundingLedgerUpdates" {
1868                        return None;
1869                    }
1870                    let Ok(ledger) = serde_json::from_value::<WsLedgerUpdates>(env.data) else {
1871                        return None;
1872                    };
1873                    let items: Vec<_> = ledger
1874                        .updates
1875                        .into_iter()
1876                        .filter_map(|e| {
1877                            if e.delta.kind != "deposit" {
1878                                return None;
1879                            }
1880                            let amount_usd = e.delta.usdc.as_deref().and_then(parse_decimal)?;
1881                            Some(Deposit {
1882                                asset: "USDC".to_string(),
1883                                amount_usd,
1884                                timestamp_ms: e.time,
1885                            })
1886                        })
1887                        .collect();
1888                    if items.is_empty() {
1889                        None
1890                    } else {
1891                        Some(stream::iter(items.into_iter().map(Ok)))
1892                    }
1893                })
1894                .flatten(),
1895        )
1896    }
1897
1898    fn subscribe_withdrawals(&self) -> BoxStream<Result<Withdrawal, String>> {
1899        let Some(addr) = self.user_address.as_ref() else {
1900            return Box::pin(stream::empty());
1901        };
1902        let addr_str = addr.clone();
1903        let sub = serde_json::json!({
1904            "method": "subscribe",
1905            "subscription": {"type": "userNonFundingLedgerUpdates", "user": addr_str.clone()}
1906        });
1907        let key = crate::ws::SubKey {
1908            channel: "userNonFundingLedgerUpdates".to_string(),
1909            routing_key: addr_str,
1910        };
1911        let raw_stream = self.ws_mux.subscribe(key, sub);
1912        Box::pin(
1913            raw_stream
1914                .filter_map(|text| async move {
1915                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1916                        return None;
1917                    };
1918                    if env.channel != "userNonFundingLedgerUpdates" {
1919                        return None;
1920                    }
1921                    let Ok(ledger) = serde_json::from_value::<WsLedgerUpdates>(env.data) else {
1922                        return None;
1923                    };
1924                    let items: Vec<_> = ledger
1925                        .updates
1926                        .into_iter()
1927                        .filter_map(|e| {
1928                            if e.delta.kind != "withdraw" {
1929                                return None;
1930                            }
1931                            let amount_usd = e.delta.usdc.as_deref().and_then(parse_decimal)?;
1932                            Some(Withdrawal {
1933                                asset: "USDC".to_string(),
1934                                amount_usd,
1935                                timestamp_ms: e.time,
1936                            })
1937                        })
1938                        .collect();
1939                    if items.is_empty() {
1940                        None
1941                    } else {
1942                        Some(stream::iter(items.into_iter().map(Ok)))
1943                    }
1944                })
1945                .flatten(),
1946        )
1947    }
1948
1949    /// Subscribe to spot wallet balance updates for the registered user address.
1950    /// Requires authentication (address must be set). Returns error if address not registered.
1951    fn subscribe_spot_balance(
1952        &self,
1953    ) -> BoxStream<Result<Vec<guilder_abstraction::Balance>, String>> {
1954        let Some(addr) = self.user_address.as_ref() else {
1955            return Box::pin(stream::iter(vec![Err(
1956                "user address not registered".to_string()
1957            )]));
1958        };
1959        let addr_str = addr.clone();
1960        self.subscribe_spot_balance_with_address(addr_str)
1961    }
1962
1963    /// Subscribe to spot wallet balance updates for a specific address.
1964    fn subscribe_spot_balance_with_address(
1965        &self,
1966        address: String,
1967    ) -> BoxStream<Result<Vec<guilder_abstraction::Balance>, String>> {
1968        let addr_str = address;
1969        let sub = serde_json::json!({
1970            "method": "subscribe",
1971            "subscription": {"type": "userEvents", "user": addr_str.clone()}
1972        });
1973        let key = crate::ws::SubKey {
1974            channel: "user".to_string(),
1975            routing_key: addr_str,
1976        };
1977        let raw_stream = self.ws_mux.subscribe(key, sub);
1978        Box::pin(raw_stream.filter_map(|text| async move {
1979            let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1980                return None;
1981            };
1982            if env.channel != "user" {
1983                return None;
1984            }
1985            let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1986                return None;
1987            };
1988
1989            // Extract spot balances from the event
1990            let spot_state = event.spot_state?;
1991            let balances = spot_state.balances?;
1992
1993            let items: Vec<_> = balances
1994                .into_iter()
1995                .filter_map(|b| {
1996                    let total = parse_decimal(&b.total)?;
1997                    let locked = parse_decimal(&b.hold)?;
1998                    let available = total - locked;
1999                    Some(guilder_abstraction::Balance {
2000                        coin: b.coin,
2001                        total,
2002                        available,
2003                        locked,
2004                    })
2005                })
2006                .collect();
2007
2008            if items.is_empty() {
2009                None
2010            } else {
2011                Some(Ok(items))
2012            }
2013        }))
2014    }
2015}
2016
2017#[cfg(test)]
2018mod msgpack_tests {
2019    use super::*;
2020    use serde_json::json;
2021
2022    #[test]
2023    fn test_msgpack_null() {
2024        let result = value_to_msgpack(&Value::Null);
2025        assert_eq!(result, vec![0xc0]);
2026    }
2027
2028    #[test]
2029    fn test_msgpack_bool() {
2030        assert_eq!(value_to_msgpack(&Value::Bool(true)), vec![0xc3]);
2031        assert_eq!(value_to_msgpack(&Value::Bool(false)), vec![0xc2]);
2032    }
2033
2034    #[test]
2035    fn test_msgpack_positive_fixint() {
2036        // 0–127: positive fixint
2037        assert_eq!(value_to_msgpack(&json!(0)), vec![0x00]);
2038        assert_eq!(value_to_msgpack(&json!(1)), vec![0x01]);
2039        assert_eq!(value_to_msgpack(&json!(127)), vec![0x7f]);
2040    }
2041
2042    #[test]
2043    fn test_msgpack_uint8() {
2044        // 128–255: uint8
2045        assert_eq!(value_to_msgpack(&json!(128)), vec![0xcc, 0x80]);
2046        assert_eq!(value_to_msgpack(&json!(255)), vec![0xcc, 0xff]);
2047    }
2048
2049    #[test]
2050    fn test_msgpack_uint16() {
2051        // 256–65535: uint16
2052        assert_eq!(value_to_msgpack(&json!(256)), vec![0xcd, 0x01, 0x00]);
2053        assert_eq!(value_to_msgpack(&json!(65535)), vec![0xcd, 0xff, 0xff]);
2054    }
2055
2056    #[test]
2057    fn test_msgpack_uint32() {
2058        // 65536–4294967295: uint32
2059        assert_eq!(
2060            value_to_msgpack(&json!(65536)),
2061            vec![0xce, 0x00, 0x01, 0x00, 0x00]
2062        );
2063        assert_eq!(
2064            value_to_msgpack(&json!(4294967295u64)),
2065            vec![0xce, 0xff, 0xff, 0xff, 0xff]
2066        );
2067    }
2068
2069    #[test]
2070    fn test_msgpack_uint64() {
2071        // >4294967295: uint64
2072        let big: u64 = 4294967296;
2073        assert_eq!(
2074            value_to_msgpack(&json!(big)),
2075            vec![0xcf, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00]
2076        );
2077    }
2078
2079    #[test]
2080    fn test_msgpack_negative_fixint() {
2081        // -1 to -32: negative fixint
2082        assert_eq!(value_to_msgpack(&json!(-1)), vec![0xff]);
2083        assert_eq!(value_to_msgpack(&json!(-32)), vec![0xe0]);
2084    }
2085
2086    #[test]
2087    fn test_msgpack_int8() {
2088        // -33 to -128: int8
2089        assert_eq!(value_to_msgpack(&json!(-33)), vec![0xd0, 0xdf]);
2090        assert_eq!(value_to_msgpack(&json!(-128)), vec![0xd0, 0x80]);
2091    }
2092
2093    #[test]
2094    fn test_msgpack_int16() {
2095        // -129 to -32768: int16
2096        assert_eq!(value_to_msgpack(&json!(-129)), vec![0xd1, 0xff, 0x7f]);
2097        assert_eq!(value_to_msgpack(&json!(-32768)), vec![0xd1, 0x80, 0x00]);
2098    }
2099
2100    #[test]
2101    fn test_msgpack_int32() {
2102        // -32769 to -2147483648: int32
2103        assert_eq!(
2104            value_to_msgpack(&json!(-32769)),
2105            vec![0xd2, 0xff, 0xff, 0x7f, 0xff]
2106        );
2107        assert_eq!(
2108            value_to_msgpack(&json!(-2147483648i64)),
2109            vec![0xd2, 0x80, 0x00, 0x00, 0x00]
2110        );
2111    }
2112
2113    #[test]
2114    fn test_msgpack_int64() {
2115        let val: i64 = -2147483649;
2116        let result = value_to_msgpack(&json!(val));
2117        assert_eq!(result[0], 0xd3); // int64 marker
2118        assert_eq!(result.len(), 9);
2119    }
2120
2121    #[test]
2122    fn test_msgpack_float() {
2123        let result = value_to_msgpack(&json!(3.14));
2124        assert_eq!(result[0], 0xcb); // float64 marker
2125        assert_eq!(result.len(), 9);
2126    }
2127
2128    #[test]
2129    fn test_msgpack_fixstr() {
2130        // 0–31 bytes: fixstr
2131        assert_eq!(value_to_msgpack(&json!("")), vec![0xa0]);
2132        assert_eq!(
2133            value_to_msgpack(&json!("hello")),
2134            {
2135                let mut expected = vec![0xa5];
2136                expected.extend_from_slice(b"hello");
2137                expected
2138            }
2139        );
2140        let s = "a".repeat(31);
2141        let result = value_to_msgpack(&json!(s));
2142        assert_eq!(result[0], 0xbf); // 0xa0 | 31
2143        assert_eq!(result.len(), 32);
2144    }
2145
2146    #[test]
2147    fn test_msgpack_str8() {
2148        let s = "a".repeat(32);
2149        let result = value_to_msgpack(&json!(s));
2150        assert_eq!(result[0], 0xd9); // str8 marker
2151        assert_eq!(result[1], 32);
2152        assert_eq!(result.len(), 34);
2153    }
2154
2155    #[test]
2156    fn test_msgpack_fixarray() {
2157        // 0–15 elements: fixarray
2158        assert_eq!(value_to_msgpack(&json!([])), vec![0x90]);
2159        let result = value_to_msgpack(&json!([1, 2, 3]));
2160        assert_eq!(result[0], 0x93);
2161        assert_eq!(result, vec![0x93, 0x01, 0x02, 0x03]);
2162    }
2163
2164    #[test]
2165    fn test_msgpack_fixmap() {
2166        // 0–15 entries: fixmap
2167        assert_eq!(value_to_msgpack(&json!({})), vec![0x80]);
2168        let result = value_to_msgpack(&json!({"a": 1}));
2169        assert_eq!(result[0], 0x81); // fixmap(1)
2170        assert_eq!(result, {
2171            let mut expected = vec![0x81];
2172            expected.extend_from_slice(&value_to_msgpack(&json!("a")));
2173            expected.extend_from_slice(&value_to_msgpack(&json!(1)));
2174            expected
2175        });
2176    }
2177
2178    #[test]
2179    fn test_msgmap_preserves_insertion_order() {
2180        // Verify keys are serialized in JSON insertion order, not sorted
2181        let val = json!({
2182            "z": 1,
2183            "a": 2,
2184            "m": 3
2185        });
2186        let result = value_to_msgpack(&val);
2187        // fixmap(3)
2188        assert_eq!(result[0], 0x83);
2189        // First key should be "z" (insertion order), not "a" (sorted)
2190        assert_eq!(result[1], 0xa1); // fixstr(1)
2191        assert_eq!(result[2], b'z');
2192    }
2193
2194    #[test]
2195    fn test_msgpack_mixed_array() {
2196        let val = json!([null, true, false, 42, "hi", [1, 2]]);
2197        let result = value_to_msgpack(&val);
2198        assert_eq!(result[0], 0x96); // fixarray(6)
2199        assert_eq!(result[1], 0xc0); // null
2200        assert_eq!(result[2], 0xc3); // true
2201        assert_eq!(result[3], 0xc2); // false
2202        assert_eq!(result[4], 0x2a); // 42
2203        // "hi" = fixstr(2) + "hi"
2204        assert_eq!(result[5], 0xa2);
2205        assert_eq!(result[6], b'h');
2206        assert_eq!(result[7], b'i');
2207    }
2208
2209    #[test]
2210    fn test_build_order_msgpack_without_cloid() {
2211        let result = build_order_msgpack(
2212            0,       // asset index
2213            true,    // is_buy
2214            "1000",  // price
2215            "0.1",   // size
2216            false,   // reduce_only
2217            "limit", // order_kind
2218            b"gtc",  // tif
2219            None,    // cloid
2220        );
2221        // fixmap(6)
2222        assert_eq!(result[0], 0x86);
2223    }
2224
2225    #[test]
2226    fn test_build_order_msgpack_with_cloid() {
2227        let result = build_order_msgpack(
2228            0,       // asset index
2229            true,    // is_buy
2230            "1000",  // price
2231            "0.1",   // size
2232            false,   // reduce_only
2233            "limit", // order_kind
2234            b"gtc",  // tif
2235            Some("my-cloid"), // cloid
2236        );
2237        // fixmap(7)
2238        assert_eq!(result[0], 0x87);
2239    }
2240
2241    #[test]
2242    fn test_action_to_canonical_msgpack() {
2243        let action = json!({
2244            "type": "order",
2245            "orders": [{"a": 0, "b": true, "p": "1000", "s": "0.1", "r": false, "t": {"limit": {"tif": "gtc"}}}],
2246            "grouping": "na"
2247        });
2248        let result = action_to_canonical_msgpack(&action).unwrap();
2249        // fixmap(3)
2250        assert_eq!(result[0], 0x83);
2251    }
2252
2253    #[test]
2254    fn test_msgpack_matches_rmp_serde_for_simple_values() {
2255        // Verify our encoding matches rmp_serde for simple scalar values
2256        use rmp_serde::to_vec;
2257
2258        for val in [json!(0), json!(127), json!(255), json!(1000), json!(-1), json!(-32), json!(-128)] {
2259            let ours = value_to_msgpack(&val);
2260            let theirs = to_vec(&val).unwrap();
2261            assert_eq!(
2262                ours, theirs,
2263                "mismatch for {}: ours={:?}, rmp={:?}",
2264                val, ours, theirs
2265            );
2266        }
2267    }
2268
2269    #[test]
2270    fn test_msgpack_string_encoding() {
2271        use rmp_serde::to_vec;
2272        for val in [json!(""), json!("a"), json!("hello world"), json!("BTC-USD")] {
2273            let ours = value_to_msgpack(&val);
2274            let theirs = to_vec(&val).unwrap();
2275            assert_eq!(
2276                ours, theirs,
2277                "mismatch for {}: ours={:?}, rmp={:?}",
2278                val, ours, theirs
2279            );
2280        }
2281    }
2282
2283    #[test]
2284    fn test_msgpack_bool_encoding() {
2285        use rmp_serde::to_vec;
2286        let theirs = to_vec(&json!(true)).unwrap();
2287        assert_eq!(value_to_msgpack(&json!(true)), theirs);
2288        let theirs = to_vec(&json!(false)).unwrap();
2289        assert_eq!(value_to_msgpack(&json!(false)), theirs);
2290    }
2291
2292    #[test]
2293    fn test_msgpack_null_encoding() {
2294        use rmp_serde::to_vec;
2295        let theirs = to_vec(&Value::Null).unwrap();
2296        assert_eq!(value_to_msgpack(&Value::Null), theirs);
2297    }
2298
2299    #[test]
2300    fn test_msgpack_nested_object() {
2301        let val = json!({
2302            "outer": {
2303                "inner": 42
2304            }
2305        });
2306        let result = value_to_msgpack(&val);
2307        assert_eq!(result[0], 0x81); // fixmap(1)
2308    }
2309
2310    #[test]
2311    fn test_msgpack_empty_containers() {
2312        assert_eq!(value_to_msgpack(&json!([])), vec![0x90]);
2313        assert_eq!(value_to_msgpack(&json!({})), vec![0x80]);
2314    }
2315}