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.as_ref().map(|c| {
1035            let hash = keccak256(c.as_bytes());
1036            format!("0x{}", hex::encode(&hash[..16]))
1037        });
1038
1039        // Build msgpack with Python SDK field order (matching the Python SDK's
1040        // msgpack output). The server hashes the msgpack for signature verification,
1041        // and Python preserves dict insertion order.
1042        let order_msgpack = build_order_msgpack(
1043            asset_idx,
1044            is_buy,
1045            &price_str,
1046            &size_str,
1047            false, // reduce_only
1048            order_kind,
1049            tif_bytes,
1050            cloid_hex.as_deref(),
1051        );
1052
1053        // Build the action-level msgpack with Python SDK field order (insertion order):
1054        // type → orders → grouping
1055        let mut action_msgpack = Vec::new();
1056        action_msgpack.push(0x83); // fixmap(3)
1057        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("type".to_string())));
1058        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("order".to_string())));
1059        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("orders".to_string())));
1060        action_msgpack.push(0x91); // fixarray(1)
1061        action_msgpack.extend_from_slice(&order_msgpack);
1062        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("grouping".to_string())));
1063        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("na".to_string())));
1064
1065        // Build JSON with official SDK field order (alphabetical):
1066        // grouping → orders → type
1067        let order_type_json = match order_type {
1068            OrderType::Limit => format!(r#"{{"limit":{{"tif":"{tif_str}"}}}}"#),
1069            OrderType::Market => r#"{"limit":{"tif":"Ioc"}}"#.to_string(),
1070        };
1071
1072        let cloid_json = if let Some(ref c) = cloid_hex {
1073            format!(r#","c":"{c}""#)
1074        } else {
1075            String::new()
1076        };
1077
1078        let action_json_str = format!(
1079            r#"{{"type":"order","orders":[{{"a":{asset_idx},"b":{is_buy},"p":"{price}","s":"{size}","r":false,"t":{order_type_json}{cloid_json}}}],"grouping":"na"}}"#,
1080            price = price_str,
1081            size = size_str,
1082        );
1083
1084        // Sign using the canonical msgpack (matching Python's field order)
1085        let private_key = self.require_private_key()?;
1086        let nonce = std::time::SystemTime::now()
1087            .duration_since(std::time::UNIX_EPOCH)
1088            .unwrap()
1089            .as_millis() as u64;
1090
1091        let (r, s, v) = sign_with_msgpack(&action_msgpack, private_key, nonce, None)?;
1092
1093        let payload_str = format!(
1094            r#"{{"action":{},"nonce":{},"signature":{{"r":"{}","s":"{}","v":{}}},"vaultAddress":null,"expiresAfter":null}}"#,
1095            action_json_str,
1096            nonce,
1097            r, s, v
1098        );
1099
1100        self.rest_limiter.acquire(1).await.map_err(|e| {
1101            format!(
1102                "rate_limited: rest_weight exhausted, retry_after_ms={}",
1103                e.retry_after.as_millis()
1104            )
1105        })?;
1106        self.address_limiter.acquire(1, false).await.map_err(|e| {
1107            format!(
1108                "rate_limited: address quota exhausted, retry_after_ms={}",
1109                e.retry_after.as_millis()
1110            )
1111        })?;
1112
1113        let resp = self
1114            .client
1115            .post(HYPERLIQUID_EXCHANGE_URL)
1116            .header("Content-Type", "application/json")
1117            .body(payload_str)
1118            .send()
1119            .await
1120            .map_err(|e| e.to_string())?;
1121
1122        let status = resp.status();
1123        if !status.is_success() {
1124            let text = resp.text().await.map_err(|e| e.to_string())?;
1125            return Err(format!("HTTP {status}: {text}"));
1126        }
1127
1128        let body: Value = parse_response(resp).await?;
1129        if body["status"].as_str() == Some("err") {
1130            return Err(body["response"]
1131                .as_str()
1132                .unwrap_or("unknown error")
1133                .to_string());
1134        }
1135        let statuses = &body["response"]["data"]["statuses"][0];
1136
1137        let (oid, returned_cloid, timestamp_ms) = if let Some(resting) = statuses.get("resting") {
1138            let oid = resting["oid"]
1139                .as_i64()
1140                .ok_or_else(|| format!("resting status missing oid: {}", body))?;
1141            let returned_cloid = resting["cloid"].as_str().map(|s: &str| s.to_string());
1142            // resting doesn't include a timestamp
1143            let ts = std::time::SystemTime::now()
1144                .duration_since(std::time::UNIX_EPOCH)
1145                .unwrap()
1146                .as_millis() as i64;
1147            (oid, returned_cloid, ts)
1148        } else if let Some(filled) = statuses.get("filled") {
1149            let oid = filled["oid"]
1150                .as_i64()
1151                .ok_or_else(|| format!("filled status missing oid: {}", body))?;
1152            // filled doesn't include cloid
1153            let ts = std::time::SystemTime::now()
1154                .duration_since(std::time::UNIX_EPOCH)
1155                .unwrap()
1156                .as_millis() as i64;
1157            (oid, None, ts)
1158        } else if let Some(error) = statuses.get("error") {
1159            return Err(error
1160                .as_str()
1161                .unwrap_or("order rejected with unknown error")
1162                .to_string());
1163        } else {
1164            return Err(format!("unexpected order status: {}", body));
1165        };
1166
1167        Ok(OrderPlacement {
1168            order_id: oid,
1169            symbol,
1170            side,
1171            price,
1172            quantity: volume,
1173            timestamp_ms,
1174            cloid: returned_cloid.or(cloid),
1175        })
1176    }
1177
1178    /// Modifies price and size of an existing order by its order ID. Requires `with_auth`.
1179    /// Fetches the order's current coin and side before submitting the modify action.
1180    async fn change_order_by_cloid(
1181        &self,
1182        cloid: i64,
1183        price: Decimal,
1184        volume: Decimal,
1185    ) -> Result<i64, String> {
1186        let user = self.require_user_address()?;
1187
1188        // openOrders → weight 20; get_asset_index → meta weight 20
1189        let resp = self
1190            .info_post(
1191                serde_json::json!({"type": "openOrders", "user": user}),
1192                20,
1193                "change_order_by_cloid",
1194            )
1195            .await?;
1196        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1197        let order = orders
1198            .iter()
1199            .find(|o| o.oid == cloid)
1200            .ok_or_else(|| format!("order {} not found", cloid))?;
1201
1202        let asset_idx = self.get_asset_index(&order.coin).await?;
1203        let is_buy = order.side == "B";
1204
1205        let action = serde_json::json!({
1206            "type": "batchModify",
1207            "modifies": [{
1208                "oid": cloid,
1209                "order": {
1210                    "a": asset_idx,
1211                    "b": is_buy,
1212                    "p": price.to_string(),
1213                    "s": volume.to_string(),
1214                    "r": false,
1215                    "t": {"limit": {"tif": "Gtc"}}
1216                }
1217            }]
1218        });
1219
1220        self.submit_signed_action(action, None).await?;
1221        Ok(cloid)
1222    }
1223
1224    /// Cancels a single order by its order ID. Requires `with_auth`.
1225    /// Fetches open orders to resolve the coin/asset before cancelling.
1226    async fn cancel_order(&self, cloid: i64) -> Result<i64, String> {
1227        let user = self.require_user_address()?;
1228
1229        // openOrders → weight 20
1230        let resp = self
1231            .info_post(
1232                serde_json::json!({"type": "openOrders", "user": user}),
1233                20,
1234                "cancel_order",
1235            )
1236            .await?;
1237        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1238        let order = orders
1239            .iter()
1240            .find(|o| o.oid == cloid)
1241            .ok_or_else(|| format!("order {} not found", cloid))?;
1242
1243        let asset_idx = self.get_asset_index(&order.coin).await?;
1244        let action = serde_json::json!({
1245            "type": "cancel",
1246            "cancels": [{"a": asset_idx, "o": cloid}]
1247        });
1248
1249        self.submit_signed_action(action, None).await?;
1250        Ok(cloid)
1251    }
1252
1253    /// Cancels all open orders. Requires `with_auth`.
1254    /// Fetches all open orders and submits a batch cancel in a single signed request.
1255    async fn cancel_all_order(&self) -> Result<bool, String> {
1256        let user = self.require_user_address()?;
1257
1258        // openOrders → weight 20
1259        let resp = self
1260            .info_post(
1261                serde_json::json!({"type": "openOrders", "user": user}),
1262                20,
1263                "cancel_all_order",
1264            )
1265            .await?;
1266        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1267        if orders.is_empty() {
1268            return Ok(true);
1269        }
1270
1271        // meta → weight 20
1272        let meta_resp = self
1273            .info_post(serde_json::json!({"type": "meta"}), 20, "cancel_all_order")
1274            .await?;
1275        let meta: MetaResponse = parse_response(meta_resp).await?;
1276
1277        let cancels: Vec<Value> = orders
1278            .iter()
1279            .filter_map(|o| {
1280                let asset_idx = meta.universe.iter().position(|a| a.name == o.coin)?;
1281                Some(serde_json::json!({"a": asset_idx, "o": o.oid}))
1282            })
1283            .collect();
1284
1285        let action = serde_json::json!({"type": "cancel", "cancels": cancels});
1286        self.submit_signed_action(action, None).await?;
1287        Ok(true)
1288    }
1289}
1290
1291#[allow(async_fn_in_trait)]
1292impl guilder_abstraction::SubscribeMarketData for HyperliquidClient {
1293    fn subscribe_l2_update(&self, symbol: String) -> BoxStream<Result<L2Update, String>> {
1294        let sub = serde_json::json!({
1295            "method": "subscribe",
1296            "subscription": {"type": "l2Book", "coin": symbol.clone()}
1297        });
1298        let key = crate::ws::SubKey {
1299            channel: "l2Book".to_string(),
1300            routing_key: symbol,
1301        };
1302        let stream = self.ws_mux.subscribe(key, sub);
1303        Box::pin(async_stream::stream! {
1304            for await msg in stream {
1305                let Ok(env) = serde_json::from_str::<WsEnvelope>(&msg) else {
1306                    continue;
1307                };
1308                if env.channel != "l2Book" {
1309                    continue;
1310                }
1311                let Ok(book) = serde_json::from_value::<WsBook>(env.data) else {
1312                    continue;
1313                };
1314                for level in book.levels.first().into_iter().flatten() {
1315                    if let (Some(price), Some(volume)) =
1316                        (parse_decimal(&level.px), parse_decimal(&level.sz))
1317                    {
1318                        yield Ok(L2Update {
1319                            symbol: book.coin.clone(),
1320                            price,
1321                            volume,
1322                            side: Side::Ask,
1323                            sequence: book.time,
1324                        });
1325                    }
1326                }
1327                for level in book.levels.get(1).into_iter().flatten() {
1328                    if let (Some(price), Some(volume)) =
1329                        (parse_decimal(&level.px), parse_decimal(&level.sz))
1330                    {
1331                        yield Ok(L2Update {
1332                            symbol: book.coin.clone(),
1333                            price,
1334                            volume,
1335                            side: Side::Bid,
1336                            sequence: book.time,
1337                        });
1338                    }
1339                }
1340            }
1341        })
1342    }
1343
1344    fn subscribe_asset_context(&self, symbol: String) -> BoxStream<Result<AssetContext, String>> {
1345        let sub = serde_json::json!({
1346            "method": "subscribe",
1347            "subscription": {"type": "activeAssetCtx", "coin": symbol.clone()}
1348        });
1349        let key = crate::ws::SubKey {
1350            channel: "activeAssetCtx".to_string(),
1351            routing_key: symbol,
1352        };
1353        let stream = self.ws_mux.subscribe(key, sub);
1354        Box::pin(async_stream::stream! {
1355            for await msg in stream {
1356                let Ok(env) = serde_json::from_str::<WsEnvelope>(&msg) else {
1357                    continue;
1358                };
1359                if env.channel != "activeAssetCtx" {
1360                    continue;
1361                }
1362                let Ok(update) = serde_json::from_value::<WsAssetCtx>(env.data) else {
1363                    continue;
1364                };
1365                let ctx = &update.ctx;
1366                let (Some(open_interest), Some(funding_rate), Some(mark_price), Some(day_volume)) = (
1367                    parse_decimal(&ctx.open_interest),
1368                    parse_decimal(&ctx.funding),
1369                    parse_decimal(&ctx.mark_px),
1370                    parse_decimal(&ctx.day_ntl_vlm),
1371                ) else {
1372                    continue;
1373                };
1374                yield Ok(AssetContext {
1375                    symbol: update.coin,
1376                    open_interest,
1377                    funding_rate,
1378                    mark_price,
1379                    day_volume,
1380                    mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
1381                    oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
1382                    premium: ctx.premium.as_deref().and_then(parse_decimal),
1383                    prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
1384                    // szDecimals is static metadata, not streamed via activeAssetCtx WS
1385                    sz_decimals: 0,
1386                });
1387            }
1388        })
1389    }
1390
1391    fn subscribe_liquidation(&self, user: String) -> BoxStream<Result<Liquidation, String>> {
1392        let sub = serde_json::json!({
1393            "method": "subscribe",
1394            "subscription": {"type": "userEvents", "user": user.clone()}
1395        });
1396        let key = crate::ws::SubKey {
1397            channel: "userEvents".to_string(),
1398            routing_key: user,
1399        };
1400        let raw_stream = self.ws_mux.subscribe(key, sub);
1401        Box::pin(
1402            raw_stream
1403                .filter_map(|text| async move {
1404                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1405                        return None;
1406                    };
1407                    if env.channel != "userEvents" {
1408                        return None;
1409                    }
1410                    let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1411                        return None;
1412                    };
1413                    let liq = event.liquidation?;
1414                    let (Some(notional_position), Some(account_value)) = (
1415                        parse_decimal(&liq.liquidated_ntl_pos),
1416                        parse_decimal(&liq.liquidated_account_value),
1417                    ) else {
1418                        return None;
1419                    };
1420                    let item = Liquidation {
1421                        symbol: String::new(),
1422                        side: OrderSide::Sell,
1423                        liquidated_user: liq.liquidated_user,
1424                        notional_position,
1425                        account_value,
1426                    };
1427                    Some(stream::iter(vec![Ok(item)]))
1428                })
1429                .flatten(),
1430        )
1431    }
1432
1433    fn subscribe_fill(&self, symbol: String) -> BoxStream<Result<Fill, String>> {
1434        let sub = serde_json::json!({
1435            "method": "subscribe",
1436            "subscription": {"type": "trades", "coin": symbol.clone()}
1437        });
1438        let key = crate::ws::SubKey {
1439            channel: "trades".to_string(),
1440            routing_key: symbol,
1441        };
1442        let stream = self.ws_mux.subscribe(key, sub);
1443        Box::pin(async_stream::stream! {
1444            for await msg in stream {
1445                let Ok(env) = serde_json::from_str::<WsEnvelope>(&msg) else {
1446                    continue;
1447                };
1448                if env.channel != "trades" {
1449                    continue;
1450                }
1451                let Ok(trades) = serde_json::from_value::<Vec<WsTrade>>(env.data) else {
1452                    continue;
1453                };
1454                for trade in trades {
1455                    let side = if trade.side == "B" {
1456                        OrderSide::Buy
1457                    } else {
1458                        OrderSide::Sell
1459                    };
1460                    let price = parse_decimal(&trade.px);
1461                    let volume = parse_decimal(&trade.sz);
1462                    if let (Some(price), Some(volume)) = (price, volume) {
1463                        yield Ok(Fill {
1464                            symbol: trade.coin,
1465                            price,
1466                            volume,
1467                            side,
1468                            timestamp_ms: trade.time,
1469                            trade_id: trade.tid,
1470                        });
1471                    }
1472                }
1473            }
1474        })
1475    }
1476}
1477
1478#[allow(async_fn_in_trait)]
1479impl guilder_abstraction::GetAccountSnapshot for HyperliquidClient {
1480    /// Returns open positions from `clearinghouseState`. Requires `with_auth`.
1481    /// Zero-size positions are filtered out. Positive `szi` = long, negative = short.
1482    async fn get_positions(&self) -> Result<Vec<Position>, String> {
1483        let user = self.require_user_address()?;
1484        // clearinghouseState → weight 2
1485        let resp = self
1486            .info_post(
1487                serde_json::json!({"type": "clearinghouseState", "user": user}),
1488                2,
1489                "get_positions",
1490            )
1491            .await?;
1492        let state: ClearinghouseStateResponse = parse_response(resp).await?;
1493
1494        Ok(state
1495            .asset_positions
1496            .into_iter()
1497            .filter_map(|ap| {
1498                let p = ap.position;
1499                let size = parse_decimal(&p.szi)?;
1500                if size.is_zero() {
1501                    return None;
1502                }
1503                let entry_price = p
1504                    .entry_px
1505                    .as_deref()
1506                    .and_then(parse_decimal)
1507                    .unwrap_or_default();
1508                let side = if size > Decimal::ZERO {
1509                    OrderSide::Buy
1510                } else {
1511                    OrderSide::Sell
1512                };
1513                Some(Position {
1514                    symbol: p.coin,
1515                    side,
1516                    size: size.abs(),
1517                    entry_price,
1518                })
1519            })
1520            .collect())
1521    }
1522
1523    /// Returns resting orders from Hyperliquid's `openOrders` endpoint. Requires `with_auth`.
1524    /// `filled_quantity` is derived as `origSz - sz` (original size minus remaining size).
1525    async fn get_open_orders(&self) -> Result<Vec<OpenOrder>, String> {
1526        let user = self.require_user_address()?;
1527        // openOrders → weight 20
1528        let resp = self
1529            .info_post(
1530                serde_json::json!({"type": "openOrders", "user": user}),
1531                20,
1532                "get_open_orders",
1533            )
1534            .await?;
1535        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1536
1537        Ok(orders
1538            .into_iter()
1539            .filter_map(|o| {
1540                let price = parse_decimal(&o.limit_px)?;
1541                let quantity = parse_decimal(&o.orig_sz)?;
1542                let remaining = parse_decimal(&o.sz)?;
1543                let filled_quantity = quantity - remaining;
1544                let side = if o.side == "B" {
1545                    OrderSide::Buy
1546                } else {
1547                    OrderSide::Sell
1548                };
1549                Some(OpenOrder {
1550                    order_id: o.oid,
1551                    symbol: o.coin,
1552                    side,
1553                    price,
1554                    quantity,
1555                    filled_quantity,
1556                })
1557            })
1558            .collect())
1559    }
1560
1561    /// Returns total account value (collateral) from `clearinghouseState`. Requires `with_auth`.
1562    async fn get_collateral(&self) -> Result<Decimal, String> {
1563        let user = self.require_user_address()?;
1564        // clearinghouseState → weight 2
1565        let resp = self
1566            .info_post(
1567                serde_json::json!({"type": "clearinghouseState", "user": user}),
1568                2,
1569                "get_collateral",
1570            )
1571            .await?;
1572        let state: ClearinghouseStateResponse = parse_response(resp).await?;
1573        parse_decimal(&state.margin_summary.account_value)
1574            .ok_or_else(|| "invalid account value".to_string())
1575    }
1576
1577    /// Returns all spot wallet balances from `spotState`. Requires `with_auth`.
1578    async fn get_spot_balance(&self) -> Result<Vec<guilder_abstraction::Balance>, String> {
1579        let user = self.require_user_address()?;
1580        // spotClearinghouseState → weight 15
1581        let resp = self
1582            .info_post(
1583                serde_json::json!({"type": "spotClearinghouseState", "user": user}),
1584                15,
1585                "get_spot_balance",
1586            )
1587            .await?;
1588
1589        #[derive(Deserialize)]
1590        struct SpotStateResponse {
1591            balances: Vec<SpotBalance>,
1592        }
1593
1594        #[allow(dead_code)]
1595        #[derive(Deserialize)]
1596        struct SpotBalance {
1597            coin: String,
1598            total: String,
1599            hold: String,
1600            #[serde(default)]
1601            token: Option<i32>,
1602            #[serde(default)]
1603            #[serde(rename = "entryNtl")]
1604            entry_ntl: Option<String>,
1605        }
1606
1607        let state: SpotStateResponse = parse_response(resp).await?;
1608
1609        state
1610            .balances
1611            .into_iter()
1612            .map(|balance| {
1613                let total = parse_decimal(&balance.total)
1614                    .ok_or_else(|| "invalid total balance".to_string())?;
1615                let locked = parse_decimal(&balance.hold)
1616                    .ok_or_else(|| "invalid hold balance".to_string())?;
1617                let available = total - locked;
1618
1619                Ok(guilder_abstraction::Balance {
1620                    coin: balance.coin,
1621                    total,
1622                    available,
1623                    locked,
1624                })
1625            })
1626            .collect()
1627    }
1628
1629    /// Returns clearing house collateral balance for an asset. Currently returns the total collateral.
1630    /// Requires `with_auth`.
1631    async fn get_collateral_balance(
1632        &self,
1633        asset: String,
1634    ) -> Result<guilder_abstraction::Balance, String> {
1635        if asset.to_uppercase() != "USDC" {
1636            return Err(format!("only USDC collateral is supported, got {}", asset));
1637        }
1638
1639        let total = self.get_collateral().await?;
1640        Ok(guilder_abstraction::Balance {
1641            coin: "USDC".to_string(),
1642            total,
1643            available: total,
1644            locked: Decimal::ZERO,
1645        })
1646    }
1647
1648    /// Returns the user's address-level API rate limit budget.
1649    /// Queries Hyperliquid's `userRateLimit` info endpoint for authoritative server-side counts.
1650    async fn get_user_rate_limit(&self) -> Result<guilder_abstraction::UserRateLimit, String> {
1651        let user = self.require_user_address()?;
1652        let resp = self
1653            .info_post(
1654                serde_json::json!({"type": "userRateLimit", "user": user}),
1655                20,
1656                "get_user_rate_limit",
1657            )
1658            .await?;
1659        let val = parse_response::<Value>(resp).await?;
1660
1661        let cumulative_volume = val["cumVlm"]
1662            .as_str()
1663            .and_then(parse_decimal)
1664            .ok_or_else(|| "missing or invalid cumVlm".to_string())?;
1665        let requests_used = val["nRequestsUsed"]
1666            .as_i64()
1667            .ok_or_else(|| "missing or invalid nRequestsUsed".to_string())?;
1668        let requests_cap = val["nRequestsCap"]
1669            .as_i64()
1670            .ok_or_else(|| "missing or invalid nRequestsCap".to_string())?;
1671        let requests_surplus = val["nRequestsSurplus"]
1672            .as_i64()
1673            .ok_or_else(|| "missing or invalid nRequestsSurplus".to_string())?;
1674
1675        Ok(guilder_abstraction::UserRateLimit {
1676            cumulative_volume,
1677            requests_used,
1678            requests_cap,
1679            requests_surplus,
1680        })
1681    }
1682}
1683
1684#[allow(async_fn_in_trait)]
1685impl guilder_abstraction::SubscribeUserEvents for HyperliquidClient {
1686    fn subscribe_user_fills(&self) -> BoxStream<Result<UserFill, String>> {
1687        let Some(addr) = self.user_address.as_ref() else {
1688            return Box::pin(stream::empty());
1689        };
1690        let addr_str = addr.clone();
1691        let sub = serde_json::json!({
1692            "method": "subscribe",
1693            "subscription": {"type": "userEvents", "user": addr_str.clone()}
1694        });
1695        let key = crate::ws::SubKey {
1696            channel: "userEvents".to_string(),
1697            routing_key: addr_str,
1698        };
1699        let raw_stream = self.ws_mux.subscribe(key, sub);
1700        Box::pin(
1701            raw_stream
1702                .filter_map(|text| async move {
1703                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1704                        return None;
1705                    };
1706                    if env.channel != "userEvents" {
1707                        return None;
1708                    }
1709                    let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1710                        return None;
1711                    };
1712                    let items: Vec<_> = event
1713                        .fills
1714                        .unwrap_or_default()
1715                        .into_iter()
1716                        .filter_map(|fill| {
1717                            let side = if fill.side == "B" {
1718                                OrderSide::Buy
1719                            } else {
1720                                OrderSide::Sell
1721                            };
1722                            let price = parse_decimal(&fill.px)?;
1723                            let quantity = parse_decimal(&fill.sz)?;
1724                            let fee_usd = parse_decimal(&fill.fee)?;
1725                            Some(UserFill {
1726                                order_id: fill.oid,
1727                                symbol: fill.coin,
1728                                side,
1729                                price,
1730                                quantity,
1731                                fee_usd,
1732                                timestamp_ms: fill.time,
1733                                cloid: fill.cloid,
1734                            })
1735                        })
1736                        .collect();
1737                    if items.is_empty() {
1738                        None
1739                    } else {
1740                        Some(stream::iter(items.into_iter().map(Ok)))
1741                    }
1742                })
1743                .flatten(),
1744        )
1745    }
1746
1747    fn subscribe_order_updates(&self) -> BoxStream<Result<OrderUpdate, String>> {
1748        let Some(addr) = self.user_address.as_ref() else {
1749            return Box::pin(stream::empty());
1750        };
1751        let addr_str = addr.clone();
1752        let sub = serde_json::json!({
1753            "method": "subscribe",
1754            "subscription": {"type": "orderUpdates", "user": addr_str.clone()}
1755        });
1756        let key = crate::ws::SubKey {
1757            channel: "orderUpdates".to_string(),
1758            routing_key: addr_str,
1759        };
1760        let raw_stream = self.ws_mux.subscribe(key, sub);
1761        Box::pin(
1762            raw_stream
1763                .filter_map(|text| async move {
1764                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1765                        return None;
1766                    };
1767                    if env.channel != "orderUpdates" {
1768                        return None;
1769                    }
1770                    let Ok(updates) = serde_json::from_value::<Vec<WsOrderUpdate>>(env.data) else {
1771                        return None;
1772                    };
1773                    let items: Vec<_> = updates
1774                        .into_iter()
1775                        .map(|upd| {
1776                            let status = match upd.status.as_str() {
1777                                "open" => OrderStatus::Placed,
1778                                "filled" => OrderStatus::Filled,
1779                                "canceled" | "cancelled" => OrderStatus::Cancelled,
1780                                _ => OrderStatus::PartiallyFilled,
1781                            };
1782                            let side = if upd.order.side == "B" {
1783                                OrderSide::Buy
1784                            } else {
1785                                OrderSide::Sell
1786                            };
1787                            OrderUpdate {
1788                                order_id: upd.order.oid,
1789                                symbol: upd.order.coin,
1790                                status,
1791                                side: Some(side),
1792                                price: parse_decimal(&upd.order.limit_px),
1793                                quantity: parse_decimal(&upd.order.orig_sz),
1794                                remaining_quantity: parse_decimal(&upd.order.sz),
1795                                timestamp_ms: upd.status_timestamp,
1796                                cloid: upd.order.cloid,
1797                            }
1798                        })
1799                        .collect();
1800                    if items.is_empty() {
1801                        None
1802                    } else {
1803                        Some(stream::iter(items.into_iter().map(Ok)))
1804                    }
1805                })
1806                .flatten(),
1807        )
1808    }
1809
1810    fn subscribe_funding_payments(&self) -> BoxStream<Result<FundingPayment, String>> {
1811        let Some(addr) = self.user_address.as_ref() else {
1812            return Box::pin(stream::empty());
1813        };
1814        let addr_str = addr.clone();
1815        let sub = serde_json::json!({
1816            "method": "subscribe",
1817            "subscription": {"type": "userEvents", "user": addr_str.clone()}
1818        });
1819        let key = crate::ws::SubKey {
1820            channel: "userEvents".to_string(),
1821            routing_key: addr_str,
1822        };
1823        let raw_stream = self.ws_mux.subscribe(key, sub);
1824        Box::pin(
1825            raw_stream
1826                .filter_map(|text| async move {
1827                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1828                        return None;
1829                    };
1830                    if env.channel != "userEvents" {
1831                        return None;
1832                    }
1833                    let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1834                        return None;
1835                    };
1836                    let funding = event.funding?;
1837                    let amount_usd = parse_decimal(&funding.usdc)?;
1838                    let item = FundingPayment {
1839                        symbol: funding.coin,
1840                        amount_usd,
1841                        timestamp_ms: funding.time,
1842                    };
1843                    Some(stream::iter(vec![Ok(item)]))
1844                })
1845                .flatten(),
1846        )
1847    }
1848
1849    fn subscribe_deposits(&self) -> BoxStream<Result<Deposit, String>> {
1850        let Some(addr) = self.user_address.as_ref() else {
1851            return Box::pin(stream::empty());
1852        };
1853        let addr_str = addr.clone();
1854        let sub = serde_json::json!({
1855            "method": "subscribe",
1856            "subscription": {"type": "userNonFundingLedgerUpdates", "user": addr_str.clone()}
1857        });
1858        let key = crate::ws::SubKey {
1859            channel: "userNonFundingLedgerUpdates".to_string(),
1860            routing_key: addr_str,
1861        };
1862        let raw_stream = self.ws_mux.subscribe(key, sub);
1863        Box::pin(
1864            raw_stream
1865                .filter_map(|text| async move {
1866                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1867                        return None;
1868                    };
1869                    if env.channel != "userNonFundingLedgerUpdates" {
1870                        return None;
1871                    }
1872                    let Ok(ledger) = serde_json::from_value::<WsLedgerUpdates>(env.data) else {
1873                        return None;
1874                    };
1875                    let items: Vec<_> = ledger
1876                        .updates
1877                        .into_iter()
1878                        .filter_map(|e| {
1879                            if e.delta.kind != "deposit" {
1880                                return None;
1881                            }
1882                            let amount_usd = e.delta.usdc.as_deref().and_then(parse_decimal)?;
1883                            Some(Deposit {
1884                                asset: "USDC".to_string(),
1885                                amount_usd,
1886                                timestamp_ms: e.time,
1887                            })
1888                        })
1889                        .collect();
1890                    if items.is_empty() {
1891                        None
1892                    } else {
1893                        Some(stream::iter(items.into_iter().map(Ok)))
1894                    }
1895                })
1896                .flatten(),
1897        )
1898    }
1899
1900    fn subscribe_withdrawals(&self) -> BoxStream<Result<Withdrawal, String>> {
1901        let Some(addr) = self.user_address.as_ref() else {
1902            return Box::pin(stream::empty());
1903        };
1904        let addr_str = addr.clone();
1905        let sub = serde_json::json!({
1906            "method": "subscribe",
1907            "subscription": {"type": "userNonFundingLedgerUpdates", "user": addr_str.clone()}
1908        });
1909        let key = crate::ws::SubKey {
1910            channel: "userNonFundingLedgerUpdates".to_string(),
1911            routing_key: addr_str,
1912        };
1913        let raw_stream = self.ws_mux.subscribe(key, sub);
1914        Box::pin(
1915            raw_stream
1916                .filter_map(|text| async move {
1917                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1918                        return None;
1919                    };
1920                    if env.channel != "userNonFundingLedgerUpdates" {
1921                        return None;
1922                    }
1923                    let Ok(ledger) = serde_json::from_value::<WsLedgerUpdates>(env.data) else {
1924                        return None;
1925                    };
1926                    let items: Vec<_> = ledger
1927                        .updates
1928                        .into_iter()
1929                        .filter_map(|e| {
1930                            if e.delta.kind != "withdraw" {
1931                                return None;
1932                            }
1933                            let amount_usd = e.delta.usdc.as_deref().and_then(parse_decimal)?;
1934                            Some(Withdrawal {
1935                                asset: "USDC".to_string(),
1936                                amount_usd,
1937                                timestamp_ms: e.time,
1938                            })
1939                        })
1940                        .collect();
1941                    if items.is_empty() {
1942                        None
1943                    } else {
1944                        Some(stream::iter(items.into_iter().map(Ok)))
1945                    }
1946                })
1947                .flatten(),
1948        )
1949    }
1950
1951    /// Subscribe to spot wallet balance updates for the registered user address.
1952    /// Requires authentication (address must be set). Returns error if address not registered.
1953    fn subscribe_spot_balance(
1954        &self,
1955    ) -> BoxStream<Result<Vec<guilder_abstraction::Balance>, String>> {
1956        let Some(addr) = self.user_address.as_ref() else {
1957            return Box::pin(stream::iter(vec![Err(
1958                "user address not registered".to_string()
1959            )]));
1960        };
1961        let addr_str = addr.clone();
1962        self.subscribe_spot_balance_with_address(addr_str)
1963    }
1964
1965    /// Subscribe to spot wallet balance updates for a specific address.
1966    fn subscribe_spot_balance_with_address(
1967        &self,
1968        address: String,
1969    ) -> BoxStream<Result<Vec<guilder_abstraction::Balance>, String>> {
1970        let addr_str = address;
1971        let sub = serde_json::json!({
1972            "method": "subscribe",
1973            "subscription": {"type": "userEvents", "user": addr_str.clone()}
1974        });
1975        let key = crate::ws::SubKey {
1976            channel: "userEvents".to_string(),
1977            routing_key: addr_str,
1978        };
1979        let raw_stream = self.ws_mux.subscribe(key, sub);
1980        Box::pin(raw_stream.filter_map(|text| async move {
1981            let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1982                return None;
1983            };
1984            if env.channel != "userEvents" {
1985                return None;
1986            }
1987            let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1988                return None;
1989            };
1990
1991            // Extract spot balances from the event
1992            let spot_state = event.spot_state?;
1993            let balances = spot_state.balances?;
1994
1995            let items: Vec<_> = balances
1996                .into_iter()
1997                .filter_map(|b| {
1998                    let total = parse_decimal(&b.total)?;
1999                    let locked = parse_decimal(&b.hold)?;
2000                    let available = total - locked;
2001                    Some(guilder_abstraction::Balance {
2002                        coin: b.coin,
2003                        total,
2004                        available,
2005                        locked,
2006                    })
2007                })
2008                .collect();
2009
2010            if items.is_empty() {
2011                None
2012            } else {
2013                Some(Ok(items))
2014            }
2015        }))
2016    }
2017}
2018
2019#[cfg(test)]
2020mod msgpack_tests {
2021    use super::*;
2022    use serde_json::json;
2023
2024    #[test]
2025    fn test_msgpack_null() {
2026        let result = value_to_msgpack(&Value::Null);
2027        assert_eq!(result, vec![0xc0]);
2028    }
2029
2030    #[test]
2031    fn test_msgpack_bool() {
2032        assert_eq!(value_to_msgpack(&Value::Bool(true)), vec![0xc3]);
2033        assert_eq!(value_to_msgpack(&Value::Bool(false)), vec![0xc2]);
2034    }
2035
2036    #[test]
2037    fn test_msgpack_positive_fixint() {
2038        // 0–127: positive fixint
2039        assert_eq!(value_to_msgpack(&json!(0)), vec![0x00]);
2040        assert_eq!(value_to_msgpack(&json!(1)), vec![0x01]);
2041        assert_eq!(value_to_msgpack(&json!(127)), vec![0x7f]);
2042    }
2043
2044    #[test]
2045    fn test_msgpack_uint8() {
2046        // 128–255: uint8
2047        assert_eq!(value_to_msgpack(&json!(128)), vec![0xcc, 0x80]);
2048        assert_eq!(value_to_msgpack(&json!(255)), vec![0xcc, 0xff]);
2049    }
2050
2051    #[test]
2052    fn test_msgpack_uint16() {
2053        // 256–65535: uint16
2054        assert_eq!(value_to_msgpack(&json!(256)), vec![0xcd, 0x01, 0x00]);
2055        assert_eq!(value_to_msgpack(&json!(65535)), vec![0xcd, 0xff, 0xff]);
2056    }
2057
2058    #[test]
2059    fn test_msgpack_uint32() {
2060        // 65536–4294967295: uint32
2061        assert_eq!(
2062            value_to_msgpack(&json!(65536)),
2063            vec![0xce, 0x00, 0x01, 0x00, 0x00]
2064        );
2065        assert_eq!(
2066            value_to_msgpack(&json!(4294967295u64)),
2067            vec![0xce, 0xff, 0xff, 0xff, 0xff]
2068        );
2069    }
2070
2071    #[test]
2072    fn test_msgpack_uint64() {
2073        // >4294967295: uint64
2074        let big: u64 = 4294967296;
2075        assert_eq!(
2076            value_to_msgpack(&json!(big)),
2077            vec![0xcf, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00]
2078        );
2079    }
2080
2081    #[test]
2082    fn test_msgpack_negative_fixint() {
2083        // -1 to -32: negative fixint
2084        assert_eq!(value_to_msgpack(&json!(-1)), vec![0xff]);
2085        assert_eq!(value_to_msgpack(&json!(-32)), vec![0xe0]);
2086    }
2087
2088    #[test]
2089    fn test_msgpack_int8() {
2090        // -33 to -128: int8
2091        assert_eq!(value_to_msgpack(&json!(-33)), vec![0xd0, 0xdf]);
2092        assert_eq!(value_to_msgpack(&json!(-128)), vec![0xd0, 0x80]);
2093    }
2094
2095    #[test]
2096    fn test_msgpack_int16() {
2097        // -129 to -32768: int16
2098        assert_eq!(value_to_msgpack(&json!(-129)), vec![0xd1, 0xff, 0x7f]);
2099        assert_eq!(value_to_msgpack(&json!(-32768)), vec![0xd1, 0x80, 0x00]);
2100    }
2101
2102    #[test]
2103    fn test_msgpack_int32() {
2104        // -32769 to -2147483648: int32
2105        assert_eq!(
2106            value_to_msgpack(&json!(-32769)),
2107            vec![0xd2, 0xff, 0xff, 0x7f, 0xff]
2108        );
2109        assert_eq!(
2110            value_to_msgpack(&json!(-2147483648i64)),
2111            vec![0xd2, 0x80, 0x00, 0x00, 0x00]
2112        );
2113    }
2114
2115    #[test]
2116    fn test_msgpack_int64() {
2117        let val: i64 = -2147483649;
2118        let result = value_to_msgpack(&json!(val));
2119        assert_eq!(result[0], 0xd3); // int64 marker
2120        assert_eq!(result.len(), 9);
2121    }
2122
2123    #[test]
2124    fn test_msgpack_float() {
2125        let result = value_to_msgpack(&json!(3.14));
2126        assert_eq!(result[0], 0xcb); // float64 marker
2127        assert_eq!(result.len(), 9);
2128    }
2129
2130    #[test]
2131    fn test_msgpack_fixstr() {
2132        // 0–31 bytes: fixstr
2133        assert_eq!(value_to_msgpack(&json!("")), vec![0xa0]);
2134        assert_eq!(
2135            value_to_msgpack(&json!("hello")),
2136            {
2137                let mut expected = vec![0xa5];
2138                expected.extend_from_slice(b"hello");
2139                expected
2140            }
2141        );
2142        let s = "a".repeat(31);
2143        let result = value_to_msgpack(&json!(s));
2144        assert_eq!(result[0], 0xbf); // 0xa0 | 31
2145        assert_eq!(result.len(), 32);
2146    }
2147
2148    #[test]
2149    fn test_msgpack_str8() {
2150        let s = "a".repeat(32);
2151        let result = value_to_msgpack(&json!(s));
2152        assert_eq!(result[0], 0xd9); // str8 marker
2153        assert_eq!(result[1], 32);
2154        assert_eq!(result.len(), 34);
2155    }
2156
2157    #[test]
2158    fn test_msgpack_fixarray() {
2159        // 0–15 elements: fixarray
2160        assert_eq!(value_to_msgpack(&json!([])), vec![0x90]);
2161        let result = value_to_msgpack(&json!([1, 2, 3]));
2162        assert_eq!(result[0], 0x93);
2163        assert_eq!(result, vec![0x93, 0x01, 0x02, 0x03]);
2164    }
2165
2166    #[test]
2167    fn test_msgpack_fixmap() {
2168        // 0–15 entries: fixmap
2169        assert_eq!(value_to_msgpack(&json!({})), vec![0x80]);
2170        let result = value_to_msgpack(&json!({"a": 1}));
2171        assert_eq!(result[0], 0x81); // fixmap(1)
2172        assert_eq!(result, {
2173            let mut expected = vec![0x81];
2174            expected.extend_from_slice(&value_to_msgpack(&json!("a")));
2175            expected.extend_from_slice(&value_to_msgpack(&json!(1)));
2176            expected
2177        });
2178    }
2179
2180    #[test]
2181    fn test_msgmap_preserves_insertion_order() {
2182        // Verify keys are serialized in JSON insertion order, not sorted
2183        let val = json!({
2184            "z": 1,
2185            "a": 2,
2186            "m": 3
2187        });
2188        let result = value_to_msgpack(&val);
2189        // fixmap(3)
2190        assert_eq!(result[0], 0x83);
2191        // First key should be "z" (insertion order), not "a" (sorted)
2192        assert_eq!(result[1], 0xa1); // fixstr(1)
2193        assert_eq!(result[2], b'z');
2194    }
2195
2196    #[test]
2197    fn test_msgpack_mixed_array() {
2198        let val = json!([null, true, false, 42, "hi", [1, 2]]);
2199        let result = value_to_msgpack(&val);
2200        assert_eq!(result[0], 0x96); // fixarray(6)
2201        assert_eq!(result[1], 0xc0); // null
2202        assert_eq!(result[2], 0xc3); // true
2203        assert_eq!(result[3], 0xc2); // false
2204        assert_eq!(result[4], 0x2a); // 42
2205        // "hi" = fixstr(2) + "hi"
2206        assert_eq!(result[5], 0xa2);
2207        assert_eq!(result[6], b'h');
2208        assert_eq!(result[7], b'i');
2209    }
2210
2211    #[test]
2212    fn test_build_order_msgpack_without_cloid() {
2213        let result = build_order_msgpack(
2214            0,       // asset index
2215            true,    // is_buy
2216            "1000",  // price
2217            "0.1",   // size
2218            false,   // reduce_only
2219            "limit", // order_kind
2220            b"gtc",  // tif
2221            None,    // cloid
2222        );
2223        // fixmap(6)
2224        assert_eq!(result[0], 0x86);
2225    }
2226
2227    #[test]
2228    fn test_build_order_msgpack_with_cloid() {
2229        let result = build_order_msgpack(
2230            0,       // asset index
2231            true,    // is_buy
2232            "1000",  // price
2233            "0.1",   // size
2234            false,   // reduce_only
2235            "limit", // order_kind
2236            b"gtc",  // tif
2237            Some("my-cloid"), // cloid
2238        );
2239        // fixmap(7)
2240        assert_eq!(result[0], 0x87);
2241    }
2242
2243    #[test]
2244    fn test_action_to_canonical_msgpack() {
2245        let action = json!({
2246            "type": "order",
2247            "orders": [{"a": 0, "b": true, "p": "1000", "s": "0.1", "r": false, "t": {"limit": {"tif": "gtc"}}}],
2248            "grouping": "na"
2249        });
2250        let result = action_to_canonical_msgpack(&action).unwrap();
2251        // fixmap(3)
2252        assert_eq!(result[0], 0x83);
2253    }
2254
2255    #[test]
2256    fn test_msgpack_matches_rmp_serde_for_simple_values() {
2257        // Verify our encoding matches rmp_serde for simple scalar values
2258        use rmp_serde::to_vec;
2259
2260        for val in [json!(0), json!(127), json!(255), json!(1000), json!(-1), json!(-32), json!(-128)] {
2261            let ours = value_to_msgpack(&val);
2262            let theirs = to_vec(&val).unwrap();
2263            assert_eq!(
2264                ours, theirs,
2265                "mismatch for {}: ours={:?}, rmp={:?}",
2266                val, ours, theirs
2267            );
2268        }
2269    }
2270
2271    #[test]
2272    fn test_msgpack_string_encoding() {
2273        use rmp_serde::to_vec;
2274        for val in [json!(""), json!("a"), json!("hello world"), json!("BTC-USD")] {
2275            let ours = value_to_msgpack(&val);
2276            let theirs = to_vec(&val).unwrap();
2277            assert_eq!(
2278                ours, theirs,
2279                "mismatch for {}: ours={:?}, rmp={:?}",
2280                val, ours, theirs
2281            );
2282        }
2283    }
2284
2285    #[test]
2286    fn test_msgpack_bool_encoding() {
2287        use rmp_serde::to_vec;
2288        let theirs = to_vec(&json!(true)).unwrap();
2289        assert_eq!(value_to_msgpack(&json!(true)), theirs);
2290        let theirs = to_vec(&json!(false)).unwrap();
2291        assert_eq!(value_to_msgpack(&json!(false)), theirs);
2292    }
2293
2294    #[test]
2295    fn test_msgpack_null_encoding() {
2296        use rmp_serde::to_vec;
2297        let theirs = to_vec(&Value::Null).unwrap();
2298        assert_eq!(value_to_msgpack(&Value::Null), theirs);
2299    }
2300
2301    #[test]
2302    fn test_msgpack_nested_object() {
2303        let val = json!({
2304            "outer": {
2305                "inner": 42
2306            }
2307        });
2308        let result = value_to_msgpack(&val);
2309        assert_eq!(result[0], 0x81); // fixmap(1)
2310    }
2311
2312    #[test]
2313    fn test_msgpack_empty_containers() {
2314        assert_eq!(value_to_msgpack(&json!([])), vec![0x90]);
2315        assert_eq!(value_to_msgpack(&json!({})), vec![0x80]);
2316    }
2317}