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