Skip to main content

guilder_client_hyperliquid/
client.rs

1use crate::rate_limiter::{AddressRateLimiter, RestRateLimiter};
2use crate::ws::manager::{
3    managed_stream, HyperliquidSubscription, HyperliquidWsManager, WsSendRateLimiter,
4};
5use crate::ws::{HyperliquidWsBook, HyperliquidWsInboundMessage};
6use async_trait::async_trait;
7use futures_util::{stream, StreamExt};
8use guilder_abstraction::{
9    self, AssetContext, BoxStream, Deposit, EcdsaSignature, ExternalSigner, Fill, FundingPayment,
10    L2Level, L2Snapshot, L2Update, Liquidation, OpenOrder, OrderPlacement, OrderSide, OrderType,
11    OrderUpdate, Position, PredictedFunding, TimeInForce, UserFill, Withdrawal,
12};
13use reqwest::Client;
14use rust_decimal::Decimal;
15use serde::Deserialize;
16use serde_json::Value;
17use std::collections::HashMap;
18use std::str::FromStr;
19use std::sync::{Arc, RwLock};
20const HYPERLIQUID_INFO_URL: &str = "https://api.hyperliquid.xyz/info";
21const HYPERLIQUID_EXCHANGE_URL: &str = "https://api.hyperliquid.xyz/exchange";
22
23/// REST/WS endpoints, selectable at construction time. `Mainnet` is the
24/// default; `Testnet` points at the official testnet replica
25/// (api.hyperliquid-testnet.xyz) — real order-book mechanics, simulated
26/// balances (faucet-funded). Configurable per Sho 2026-09-21: cleanest at
27/// client construction, not via env or compile-time flags.
28#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
29pub enum HyperliquidNetwork {
30    /// Production: api.hyperliquid.xyz
31    #[default]
32    Mainnet,
33    /// Official testnet replica: api.hyperliquid-testnet.xyz
34    Testnet,
35}
36
37impl HyperliquidNetwork {
38    pub fn info_url(&self) -> &'static str {
39        match self {
40            Self::Mainnet => HYPERLIQUID_INFO_URL,
41            Self::Testnet => "https://api.hyperliquid-testnet.xyz/info",
42        }
43    }
44    pub fn exchange_url(&self) -> &'static str {
45        match self {
46            Self::Mainnet => HYPERLIQUID_EXCHANGE_URL,
47            Self::Testnet => "https://api.hyperliquid-testnet.xyz/exchange",
48        }
49    }
50    pub fn ws_url(&self) -> &'static str {
51        match self {
52            Self::Mainnet => "wss://api.hyperliquid.xyz/ws",
53            Self::Testnet => "wss://api.hyperliquid-testnet.xyz/ws",
54        }
55    }
56    /// EIP-712 phantom-agent `source` for L1 action signing.
57    /// Hyperliquid mainnet uses "a", testnet uses "b" — matches the official
58    /// Python SDK (`construct_phantom_agent`).
59    pub fn eip712_source(&self) -> &'static str {
60        match self {
61            Self::Mainnet => "a",
62            Self::Testnet => "b",
63        }
64    }
65}
66
67async fn parse_response<T: for<'de> serde::Deserialize<'de>>(
68    resp: reqwest::Response,
69) -> Result<T, String> {
70    let status = resp.status();
71    let text = resp
72        .text()
73        .await
74        .map_err(|e| format!("failed to read response body (status {status}): {e}"))?;
75
76    if text.is_empty() {
77        return Err(format!(
78            "empty response body from Hyperliquid (HTTP {status})"
79        ));
80    }
81
82    serde_json::from_str(&text).map_err(|e| {
83        let snippet = if text.len() > 512 {
84            format!("{}... ({} bytes total)", &text[..256], text.len())
85        } else {
86            text.clone()
87        };
88        format!("deserialize error (HTTP {status}): {e}: {snippet}")
89    })
90}
91
92pub struct HyperliquidClient {
93    client: Client,
94    network: HyperliquidNetwork,
95    user_address: Option<String>,
96    private_key: Option<String>,
97    external_signer: Option<Arc<dyn ExternalSigner>>,
98    rest_limiter: Arc<RestRateLimiter>,
99    address_limiter: Arc<AddressRateLimiter>,
100    market_ws_manager: HyperliquidWsManager,
101    user_ws_managers: Arc<RwLock<HashMap<String, HyperliquidWsManager>>>,
102    ws_send_limiter: WsSendRateLimiter,
103    /// Monotonic nonce source (2026-09-28 z4dbg incident): the last nonce this
104    /// client handed out. HL rejects a signed action whose nonce was already
105    /// used (`Invalid nonce: duplicate nonce N`), and the raw wall-clock
106    /// millisecond value repeats whenever two signed actions are built within
107    /// the same millisecond — exactly the emergency close fan-out shape.
108    /// The provider hands out `max(now_ms, last + 1)` under a compare-and-swap
109    /// so concurrent signers can never observe the same nonce twice.
110    last_nonce: Arc<std::sync::atomic::AtomicU64>,
111}
112
113impl Default for HyperliquidClient {
114    fn default() -> Self {
115        Self::new()
116    }
117}
118
119impl HyperliquidClient {
120    pub fn new() -> Self {
121        Self::with_network(HyperliquidNetwork::Mainnet)
122    }
123
124    /// M2 (Sho 2026-09-21): select the network (mainnet/testnet) at
125    /// construction. All REST + WS endpoints follow.
126    pub fn with_network(network: HyperliquidNetwork) -> Self {
127        let ws_send_limiter = WsSendRateLimiter::new();
128        HyperliquidClient {
129            client: Client::new(),
130            network,
131            user_address: None,
132            private_key: None,
133            external_signer: None,
134            rest_limiter: Arc::new(RestRateLimiter::new()),
135            address_limiter: Arc::new(AddressRateLimiter::new()),
136            market_ws_manager: HyperliquidWsManager::new(
137                None,
138                ws_send_limiter.clone(),
139                network.ws_url(),
140            ),
141            user_ws_managers: Arc::new(RwLock::new(HashMap::new())),
142            ws_send_limiter,
143            last_nonce: Arc::new(std::sync::atomic::AtomicU64::new(0)),
144        }
145    }
146
147    /// Monotonic nonce source (2026-09-28 z4dbg fix): HL rejects a signed
148    /// action whose nonce repeats (`Invalid nonce: duplicate nonce N`), and the
149    /// raw wall-clock millisecond repeats whenever two actions are signed in
150    /// the same millisecond. Returns `max(now_ms, last + 1)` via CAS so
151    /// concurrent signers can never observe the same nonce twice.
152    pub(crate) fn next_nonce(&self) -> u64 {
153        use std::sync::atomic::Ordering;
154        let now_ms = std::time::SystemTime::now()
155            .duration_since(std::time::UNIX_EPOCH)
156            .unwrap()
157            .as_millis() as u64;
158        let mut last = self.last_nonce.load(Ordering::Relaxed);
159        loop {
160            let candidate = now_ms.max(last + 1);
161            match self
162                .last_nonce
163                .compare_exchange(last, candidate, Ordering::AcqRel, Ordering::Relaxed)
164            {
165                Ok(_) => return candidate,
166                Err(observed) => last = observed,
167            }
168        }
169    }
170
171    pub fn with_auth(user_address: impl Into<String>, private_key: String) -> Self {
172        Self::with_network_and_auth(HyperliquidNetwork::Mainnet, user_address, private_key)
173    }
174
175    pub fn with_network_and_auth(
176        network: HyperliquidNetwork,
177        user_address: impl Into<String>,
178        private_key: String,
179    ) -> Self {
180        let ws_send_limiter = WsSendRateLimiter::new();
181        HyperliquidClient {
182            client: Client::new(),
183            network,
184            user_address: Some(user_address.into()),
185            private_key: Some(private_key),
186            external_signer: None,
187            rest_limiter: Arc::new(RestRateLimiter::new()),
188            address_limiter: Arc::new(AddressRateLimiter::new()),
189            market_ws_manager: HyperliquidWsManager::new(
190                None,
191                ws_send_limiter.clone(),
192                network.ws_url(),
193            ),
194            user_ws_managers: Arc::new(RwLock::new(HashMap::new())),
195            ws_send_limiter,
196            last_nonce: Arc::new(std::sync::atomic::AtomicU64::new(0)),
197        }
198    }
199
200    /// Create a client authenticated via an external signer (TPM, Secure Enclave, etc.).
201    ///
202    /// The signer handles the actual ECDSA signing; the private key never leaves
203    /// the hardware. The client computes the EIP-712 digest and delegates signing
204    /// to the external signer.
205    pub fn with_external_signer(
206        user_address: impl Into<String>,
207        signer: Arc<dyn ExternalSigner>,
208    ) -> Self {
209        let ws_send_limiter = WsSendRateLimiter::new();
210        HyperliquidClient {
211            client: Client::new(),
212            network: HyperliquidNetwork::Mainnet,
213            user_address: Some(user_address.into()),
214            private_key: None,
215            external_signer: Some(signer),
216            rest_limiter: Arc::new(RestRateLimiter::new()),
217            address_limiter: Arc::new(AddressRateLimiter::new()),
218            market_ws_manager: HyperliquidWsManager::new(
219                None,
220                ws_send_limiter.clone(),
221                HyperliquidNetwork::Mainnet.ws_url(),
222            ),
223            user_ws_managers: Arc::new(RwLock::new(HashMap::new())),
224            ws_send_limiter,
225            last_nonce: Arc::new(std::sync::atomic::AtomicU64::new(0)),
226        }
227    }
228
229    pub fn with_network_and_external_signer(
230        network: HyperliquidNetwork,
231        user_address: impl Into<String>,
232        signer: Arc<dyn ExternalSigner>,
233    ) -> Self {
234        let ws_send_limiter = WsSendRateLimiter::new();
235        HyperliquidClient {
236            client: Client::new(),
237            network,
238            user_address: Some(user_address.into()),
239            private_key: None,
240            external_signer: Some(signer),
241            rest_limiter: Arc::new(RestRateLimiter::new()),
242            address_limiter: Arc::new(AddressRateLimiter::new()),
243            market_ws_manager: HyperliquidWsManager::new(
244                None,
245                ws_send_limiter.clone(),
246                network.ws_url(),
247            ),
248            user_ws_managers: Arc::new(RwLock::new(HashMap::new())),
249            ws_send_limiter,
250            last_nonce: Arc::new(std::sync::atomic::AtomicU64::new(0)),
251        }
252    }
253
254    /// Configure rate limit budgets (rest_weight/min, address_requests).
255    /// Defaults: 1200 rest weight/min, 10000 address requests.
256    pub fn with_budgets(mut self, rest_weight: u32, addr_budget: u64) -> Self {
257        self.rest_limiter = Arc::new(RestRateLimiter::new_with_budget(rest_weight));
258        self.address_limiter = Arc::new(AddressRateLimiter::new_with_budget(addr_budget));
259        self
260    }
261
262    /// POST to the info endpoint, consuming `weight` from the REST rate-limit budget.
263    /// Returns `Err("rate_limited: ...")` immediately if budget is exhausted — no retry.
264    /// Callers should handle gracefully (skip cycle, retry later, etc.).
265    async fn info_post(
266        &self,
267        body: Value,
268        weight: u32,
269        call: &str,
270    ) -> Result<reqwest::Response, String> {
271        self.rest_limiter.acquire(weight).await.map_err(|e| {
272            format!(
273                "rate_limited: info_post ({call}) budget exhausted, retry_after_ms={}",
274                e.retry_after.as_millis()
275            )
276        })?;
277        self.client
278            .post(self.network.info_url())
279            .json(&body)
280            .send()
281            .await
282            .map_err(|e| e.to_string())
283    }
284    fn require_user_address(&self) -> Result<String, String> {
285        self.user_address
286            .clone()
287            .ok_or_else(|| "user address required: use HyperliquidClient::with_auth".to_string())
288    }
289
290    fn require_private_key(&self) -> Result<Option<&str>, String> {
291        if self.external_signer.is_some() {
292            return Ok(None);
293        }
294        self.private_key.as_deref().map(Some).ok_or_else(|| {
295            "private key or external signer required: use with_auth or with_external_signer"
296                .to_string()
297        })
298    }
299
300    async fn get_asset_index(&self, symbol: &str) -> Result<usize, String> {
301        // `meta` is an "all other info" request → weight 20
302        let resp = self
303            .info_post(serde_json::json!({"type": "meta"}), 20, "get_asset_index")
304            .await?;
305        let meta: MetaResponse = parse_response(resp).await?;
306        meta.universe
307            .iter()
308            .position(|a| a.name == symbol)
309            .ok_or_else(|| format!("symbol {} not found", symbol))
310    }
311
312    async fn submit_signed_action(
313        &self,
314        action: Value,
315        vault_address: Option<&str>,
316    ) -> Result<Value, String> {
317        let private_key = self.require_private_key()?;
318        let nonce = self.next_nonce();
319
320        let (r, s, v) = sign_action(
321            private_key,
322            self.external_signer.as_ref(),
323            &action,
324            vault_address,
325            nonce,
326            self.network.eip712_source(),
327        )
328        .await?;
329
330        let payload = serde_json::json!({
331            "action": action,
332            "nonce": nonce,
333            "signature": {"r": r, "s": s, "v": v},
334            "vaultAddress": null,
335            "expiresAfter": null
336        });
337
338        // Check both rate limiters non-blocking — fail fast, no retry.
339        self.rest_limiter.acquire(1).await.map_err(|e| {
340            format!(
341                "rate_limited: rest_weight exhausted, retry_after_ms={}",
342                e.retry_after.as_millis()
343            )
344        })?;
345        self.address_limiter.acquire(1, false).await.map_err(|e| {
346            format!(
347                "rate_limited: address quota exhausted, retry_after_ms={}",
348                e.retry_after.as_millis()
349            )
350        })?;
351
352        let resp = self
353            .client
354            .post(self.network.exchange_url())
355            .json(&payload)
356            .send()
357            .await
358            .map_err(|e| e.to_string())?;
359
360        let status = resp.status();
361        if !status.is_success() {
362            let text = resp.text().await.map_err(|e| e.to_string())?;
363            return Err(format!("HTTP {status}: {text}"));
364        }
365
366        let body: Value = parse_response(resp).await?;
367        if body["status"].as_str() == Some("err") {
368            return Err(body["response"]
369                .as_str()
370                .unwrap_or("unknown error")
371                .to_string());
372        }
373        Ok(body)
374    }
375}
376
377// --- REST deserialization types ---
378
379#[derive(Deserialize)]
380struct MetaResponse {
381    universe: Vec<AssetInfo>,
382}
383
384#[derive(Deserialize)]
385struct AssetInfo {
386    name: String,
387    #[serde(rename = "szDecimals")]
388    sz_decimals: i32,
389    #[serde(rename = "isDelisted", default)]
390    is_delisted: bool,
391}
392
393type MetaAndAssetCtxsResponse = (MetaResponse, Vec<RestAssetCtx>);
394
395#[derive(Deserialize)]
396#[serde(rename_all = "camelCase")]
397#[allow(dead_code)]
398struct RestAssetCtx {
399    open_interest: String,
400    funding: String,
401    mark_px: String,
402    day_ntl_vlm: String,
403    mid_px: Option<String>,
404    oracle_px: Option<String>,
405    premium: Option<String>,
406    prev_day_px: Option<String>,
407}
408
409#[derive(Deserialize)]
410#[serde(rename_all = "camelCase")]
411#[allow(dead_code)]
412struct ClearinghouseStateResponse {
413    margin_summary: MarginSummary,
414    asset_positions: Vec<AssetPosition>,
415}
416
417/// Kept for get_positions compatibility; margin_summary fields are unused since get_collateral was removed.
418#[derive(Deserialize)]
419#[serde(rename_all = "camelCase")]
420#[allow(dead_code)]
421struct MarginSummary {
422    account_value: String,
423    #[serde(default)]
424    total_ntl_pos: Option<String>,
425    #[serde(default)]
426    total_raw_usd: Option<String>,
427    #[serde(default)]
428    total_margin_used: Option<String>,
429}
430
431#[derive(Deserialize)]
432struct AssetPosition {
433    position: PositionDetail,
434}
435
436#[derive(Deserialize)]
437#[serde(rename_all = "camelCase")]
438struct PositionDetail {
439    coin: String,
440    /// positive = long, negative = short
441    szi: String,
442    entry_px: Option<String>,
443    /// Venue-authoritative unrealized PnL (HL reports it per position).
444    #[serde(default)]
445    unrealized_pnl: Option<String>,
446}
447
448#[derive(Deserialize)]
449#[serde(rename_all = "camelCase")]
450struct RestOpenOrder {
451    coin: String,
452    side: String,
453    limit_px: String,
454    sz: String,
455    oid: i64,
456    orig_sz: String,
457    cloid: Option<String>,
458}
459
460// predictedFundings response: Vec<(coin, Vec<(venue, entry_or_null)>)>
461// The API returns null for venues that don't list the coin.
462type PredictedFundingsResponse = Vec<(String, Vec<(String, Option<PredictedFundingEntry>)>)>;
463
464#[derive(Deserialize)]
465#[serde(rename_all = "camelCase")]
466struct PredictedFundingEntry {
467    funding_rate: String,
468    next_funding_time: i64,
469}
470
471// --- Helpers ---
472
473/// spotClearinghouseState response. The maintenance map is ABSENT on
474/// testnet (field only exists on mainnet) — default = empty map = zero
475/// maintenance impact on every token (caught live 2026-09-26: every
476/// testnet get_balance failed `missing field` and equity read as 0).
477#[derive(Deserialize)]
478pub(crate) struct SpotStateResponse {
479    pub balances: Vec<SpotBalance>,
480    #[serde(default, rename = "tokenToAvailableAfterMaintenance")]
481    pub token_to_available_after_maintenance: Vec<(i32, String)>,
482}
483
484#[derive(Deserialize)]
485pub(crate) struct SpotBalance {
486    pub coin: String,
487    pub total: String,
488    pub hold: String,
489    #[serde(default)]
490    pub token: Option<i32>,
491    #[serde(default)]
492    #[serde(rename = "entryNtl")]
493    pub entry_ntl: Option<String>,
494}
495
496/// Map a parsed SPOT state (spotClearinghouseState) into account balances.
497/// Pure — unit-tested against both the mainnet shape (maintenance map
498/// present) and the testnet shape (absent).
499pub(crate) fn map_spot_state(
500    state: SpotStateResponse,
501    perp_margin_used: Option<Decimal>,
502) -> Result<Vec<guilder_abstraction::AccountBalance>, String> {
503    // Build a safe lookup map from token ID → available after maintenance.
504    // Not every token appears in the map — absence means zero maintenance impact.
505    let safe_map: HashMap<i32, Decimal> = state
506        .token_to_available_after_maintenance
507        .into_iter()
508        .filter_map(|(token_id, value)| parse_decimal(&value).map(|v| (token_id, v)))
509        .collect();
510
511    state
512        .balances
513        .into_iter()
514        .map(|balance| {
515            let equity =
516                parse_decimal(&balance.total).ok_or_else(|| "invalid total balance".to_string())?;
517            let hold =
518                parse_decimal(&balance.hold).ok_or_else(|| "invalid hold balance".to_string())?;
519            let free = equity - hold;
520
521            let safe = balance
522                .token
523                .and_then(|token_id| safe_map.get(&token_id).copied());
524            let usable = match safe {
525                Some(s) => std::cmp::min(free, s),
526                None => free,
527            };
528            let maintenance = safe.map(|s| equity - s);
529
530            // margin_used is account-level (perp side); only populate on USDC.
531            let margin_used = if balance.coin == "USDC" {
532                perp_margin_used
533            } else {
534                None
535            };
536
537            Ok(guilder_abstraction::AccountBalance {
538                token: balance.coin,
539                equity,
540                free,
541                safe,
542                usable,
543                hold,
544                margin_used,
545                maintenance,
546                settled_usd: None,
547            })
548        })
549        .collect()
550}
551
552/// Map the PERP margin account (clearinghouseState) into the trading USDC
553/// balance row. Under Hyperliquid's manual-account model spot and futures
554/// are SEPARATE ledgers — futures trading sizes from THIS account:
555/// equity = accountValue, free/usable = accountValue - totalMarginUsed.
556pub(crate) fn map_perp_state(
557    state: ClearinghouseStateResponse,
558) -> Result<guilder_abstraction::AccountBalance, String> {
559    let equity = parse_decimal(&state.margin_summary.account_value)
560        .ok_or_else(|| "invalid accountValue".to_string())?;
561    let margin_used = state
562        .margin_summary
563        .total_margin_used
564        .as_deref()
565        .and_then(parse_decimal)
566        .unwrap_or(Decimal::ZERO);
567    let free = equity - margin_used;
568    // settled cash = accountValue − Σ uPnL(assetPositions) — the accounting
569    // basis. NOT marginSummary.totalRawUsd: that field is the LIQUIDATION
570    // basis (accountValue − notional), not cash (albatross #114 live:
571    // anchoring on it produced a −119 equity, delta tracking the notional).
572    let venue_uPnL: Decimal = state
573        .asset_positions
574        .iter()
575        .filter_map(|ap| ap.position.unrealized_pnl.as_deref().and_then(parse_decimal))
576        .sum();
577    let settled_usd = Some(equity - venue_uPnL);
578    Ok(guilder_abstraction::AccountBalance {
579        token: crate::PERP_LEDGER_TOKEN.to_string(),
580        equity,
581        free,
582        safe: None,
583        usable: free,
584        hold: Decimal::ZERO,
585        margin_used: Some(margin_used),
586        maintenance: Some(equity - margin_used),
587        settled_usd,
588    })
589}
590
591fn parse_decimal(s: &str) -> Option<Decimal> {
592    Decimal::from_str(s).ok()
593}
594
595fn keccak256(data: &[u8]) -> [u8; 32] {
596    use sha3::{Digest, Keccak256};
597    Keccak256::digest(data).into()
598}
599
600/// EIP-712 domain separator for Hyperliquid L1 actions (chainId=1337).
601fn hyperliquid_domain_separator() -> [u8; 32] {
602    let type_hash = keccak256(
603        b"EIP712Domain(string name,string version,uint256 chainId,address verifyingContract)",
604    );
605    let name_hash = keccak256(b"Exchange");
606    let version_hash = keccak256(b"1");
607    let mut chain_id = [0u8; 32];
608    chain_id[28..32].copy_from_slice(&1337u32.to_be_bytes());
609    let verifying_contract = [0u8; 32];
610
611    let mut data = [0u8; 160];
612    data[..32].copy_from_slice(&type_hash);
613    data[32..64].copy_from_slice(&name_hash);
614    data[64..96].copy_from_slice(&version_hash);
615    data[96..128].copy_from_slice(&chain_id);
616    data[128..160].copy_from_slice(&verifying_contract);
617    keccak256(&data)
618}
619
620/// Convert a `serde_json::Value` to msgpack bytes, preserving the JSON map key order.
621/// This avoids rmp_serde's HashMap-based serialization which reorders map keys.
622fn value_to_msgpack(val: &Value) -> Vec<u8> {
623    match val {
624        Value::Null => vec![0xc0],
625        Value::Bool(true) => vec![0xc3],
626        Value::Bool(false) => vec![0xc2],
627        Value::Number(n) => {
628            if let Some(i) = n.as_i64() {
629                if i >= 0 {
630                    if i <= 127 {
631                        vec![i as u8]
632                    } else if i <= 255 {
633                        vec![0xcc, i as u8]
634                    } else if i <= 65535 {
635                        let mut buf = vec![0xcd];
636                        buf.extend_from_slice(&(i as u16).to_be_bytes());
637                        buf
638                    } else if i <= 4294967295 {
639                        let mut buf = vec![0xce];
640                        buf.extend_from_slice(&(i as u32).to_be_bytes());
641                        buf
642                    } else {
643                        let mut buf = vec![0xcf];
644                        buf.extend_from_slice(&(i as u64).to_be_bytes());
645                        buf
646                    }
647                } else if i >= -32 {
648                    vec![0xe0 | (i as u8)]
649                } else if i >= -128 {
650                    vec![0xd0, i as i8 as u8]
651                } else if i >= -32768 {
652                    let mut buf = vec![0xd1];
653                    buf.extend_from_slice(&(i as i16).to_be_bytes());
654                    buf
655                } else if i >= -2147483648 {
656                    let mut buf = vec![0xd2];
657                    buf.extend_from_slice(&(i as i32).to_be_bytes());
658                    buf
659                } else {
660                    let mut buf = vec![0xd3];
661                    buf.extend_from_slice(&i.to_be_bytes());
662                    buf
663                }
664            } else if let Some(f) = n.as_f64() {
665                let mut buf = vec![0xcb];
666                buf.extend_from_slice(&f.to_be_bytes());
667                buf
668            } else {
669                let u = n.as_u64().unwrap();
670                if u <= 127 {
671                    vec![u as u8]
672                } else if u <= 255 {
673                    vec![0xcc, u as u8]
674                } else if u <= 65535 {
675                    let mut buf = vec![0xcd];
676                    buf.extend_from_slice(&(u as u16).to_be_bytes());
677                    buf
678                } else if u <= 4294967295 {
679                    let mut buf = vec![0xce];
680                    buf.extend_from_slice(&(u as u32).to_be_bytes());
681                    buf
682                } else {
683                    let mut buf = vec![0xcf];
684                    buf.extend_from_slice(&u.to_be_bytes());
685                    buf
686                }
687            }
688        }
689        Value::String(s) => {
690            let bytes = s.as_bytes();
691            let len = bytes.len();
692            let mut buf = Vec::new();
693            if len <= 31 {
694                buf.push(0xa0 | (len as u8));
695            } else if len <= 255 {
696                buf.push(0xd9);
697                buf.push(len as u8);
698            } else if len <= 65535 {
699                buf.push(0xda);
700                buf.extend_from_slice(&(len as u16).to_be_bytes());
701            } else {
702                buf.push(0xdb);
703                buf.extend_from_slice(&(len as u32).to_be_bytes());
704            }
705            buf.extend_from_slice(bytes);
706            buf
707        }
708        Value::Array(arr) => {
709            let len = arr.len();
710            let mut buf = Vec::new();
711            if len <= 15 {
712                buf.push(0x90 | (len as u8));
713            } else if len <= 65535 {
714                buf.push(0xdc);
715                buf.extend_from_slice(&(len as u16).to_be_bytes());
716            } else {
717                buf.push(0xdd);
718                buf.extend_from_slice(&(len as u32).to_be_bytes());
719            }
720            for item in arr {
721                buf.extend_from_slice(&value_to_msgpack(item));
722            }
723            buf
724        }
725        Value::Object(map) => {
726            let len = map.len();
727            let mut buf = Vec::new();
728            if len <= 15 {
729                buf.push(0x80 | (len as u8));
730            } else if len <= 65535 {
731                buf.push(0xde);
732                buf.extend_from_slice(&(len as u16).to_be_bytes());
733            } else {
734                buf.push(0xdf);
735                buf.extend_from_slice(&(len as u32).to_be_bytes());
736            }
737            for (key, value) in map {
738                buf.extend_from_slice(&value_to_msgpack(&Value::String(key.clone())));
739                buf.extend_from_slice(&value_to_msgpack(value));
740            }
741            buf
742        }
743    }
744}
745
746/// Convert action to msgpack bytes preserving JSON map key insertion order
747/// (matching Python's msgpack dict ordering).
748fn action_to_canonical_msgpack(action: &Value) -> Result<Vec<u8>, String> {
749    Ok(value_to_msgpack(action))
750}
751
752/// Build msgpack for a single order with Python SDK field order:
753/// a, b, p, s, r, t, c(opt)
754fn build_order_msgpack(
755    asset_idx: usize,
756    is_buy: bool,
757    price: &str,
758    size: &str,
759    reduce_only: bool,
760    order_kind: &str,
761    tif: &[u8],
762    cloid: Option<&str>,
763) -> Vec<u8> {
764    let field_count = if cloid.is_some() { 7 } else { 6 };
765    let mut buf = Vec::new();
766    buf.push(0x80 | (field_count as u8)); // fixmap
767
768    // "a": asset_idx
769    buf.extend_from_slice(&value_to_msgpack(&Value::String("a".to_string())));
770    buf.extend_from_slice(&value_to_msgpack(&Value::Number(serde_json::Number::from(
771        asset_idx,
772    ))));
773
774    // "b": is_buy
775    buf.extend_from_slice(&value_to_msgpack(&Value::String("b".to_string())));
776    buf.push(if is_buy { 0xc3 } else { 0xc2 });
777
778    // "p": price
779    buf.extend_from_slice(&value_to_msgpack(&Value::String("p".to_string())));
780    buf.extend_from_slice(&value_to_msgpack(&Value::String(price.to_string())));
781
782    // "s": size (Python SDK puts s before r)
783    buf.extend_from_slice(&value_to_msgpack(&Value::String("s".to_string())));
784    buf.extend_from_slice(&value_to_msgpack(&Value::String(size.to_string())));
785
786    // "r": reduce_only
787    buf.extend_from_slice(&value_to_msgpack(&Value::String("r".to_string())));
788    buf.push(if reduce_only { 0xc3 } else { 0xc2 });
789
790    // "t": { order_kind: { "tif": tif_str } }
791    buf.extend_from_slice(&value_to_msgpack(&Value::String("t".to_string())));
792    // Inner: fixmap(1) with order_kind key
793    buf.push(0x81);
794    buf.extend_from_slice(&value_to_msgpack(&Value::String(order_kind.to_string())));
795    // Inner-inner: fixmap(1) with "tif" key
796    buf.push(0x81);
797    buf.extend_from_slice(&value_to_msgpack(&Value::String("tif".to_string())));
798    buf.extend_from_slice(&value_to_msgpack(&Value::String(
799        String::from_utf8_lossy(tif).to_string(),
800    )));
801
802    // "c": cloid (optional, appended at END per Python SDK)
803    if let Some(c) = cloid {
804        buf.extend_from_slice(&value_to_msgpack(&Value::String("c".to_string())));
805        buf.extend_from_slice(&value_to_msgpack(&Value::String(c.to_string())));
806    }
807
808    buf
809}
810
811/// Build msgpack for a trigger order (take profit / stop loss) with Python SDK field order.
812/// Trigger orders use: t = { "trigger": { "isMarket": bool, "triggerPx": str, "tpsl": "tp"|"sl" } }
813fn build_trigger_order_msgpack(
814    asset_idx: usize,
815    is_buy: bool,
816    price: &str,
817    size: &str,
818    reduce_only: bool,
819    trigger_px: &str,
820    is_market: bool,
821    tpsl: &str,
822    cloid: Option<&str>,
823) -> Vec<u8> {
824    let field_count = if cloid.is_some() { 7 } else { 6 };
825    let mut buf = Vec::new();
826    buf.push(0x80 | (field_count as u8)); // fixmap
827
828    // "a": asset_idx
829    buf.extend_from_slice(&value_to_msgpack(&Value::String("a".to_string())));
830    buf.extend_from_slice(&value_to_msgpack(&Value::Number(serde_json::Number::from(
831        asset_idx,
832    ))));
833
834    // "b": is_buy
835    buf.extend_from_slice(&value_to_msgpack(&Value::String("b".to_string())));
836    buf.push(if is_buy { 0xc3 } else { 0xc2 });
837
838    // "p": price (use trigger_px as price for resting, or "0" for market-on-trigger)
839    buf.extend_from_slice(&value_to_msgpack(&Value::String("p".to_string())));
840    buf.extend_from_slice(&value_to_msgpack(&Value::String(price.to_string())));
841
842    // "s": size
843    buf.extend_from_slice(&value_to_msgpack(&Value::String("s".to_string())));
844    buf.extend_from_slice(&value_to_msgpack(&Value::String(size.to_string())));
845
846    // "r": reduce_only
847    buf.extend_from_slice(&value_to_msgpack(&Value::String("r".to_string())));
848    buf.push(if reduce_only { 0xc3 } else { 0xc2 });
849
850    // "t": { "trigger": { "isMarket": bool, "triggerPx": str, "tpsl": str } }
851    buf.extend_from_slice(&value_to_msgpack(&Value::String("t".to_string())));
852    buf.push(0x81); // fixmap(1): "trigger"
853    buf.extend_from_slice(&value_to_msgpack(&Value::String("trigger".to_string())));
854    // trigger object: fixmap(3)
855    buf.push(0x83);
856    buf.extend_from_slice(&value_to_msgpack(&Value::String("isMarket".to_string())));
857    buf.push(if is_market { 0xc3 } else { 0xc2 });
858    buf.extend_from_slice(&value_to_msgpack(&Value::String("triggerPx".to_string())));
859    buf.extend_from_slice(&value_to_msgpack(&Value::String(trigger_px.to_string())));
860    buf.extend_from_slice(&value_to_msgpack(&Value::String("tpsl".to_string())));
861    buf.extend_from_slice(&value_to_msgpack(&Value::String(tpsl.to_string())));
862
863    // "c": cloid (optional, appended at END)
864    if let Some(c) = cloid {
865        buf.extend_from_slice(&value_to_msgpack(&Value::String("c".to_string())));
866        buf.extend_from_slice(&value_to_msgpack(&Value::String(c.to_string())));
867    }
868
869    buf
870}
871
872/// Compute the EIP-712 digest for a Hyperliquid exchange action.
873///
874/// This is the hash that gets signed — extracted so both direct key and
875/// external signer paths can share it.
876fn compute_eip712_digest(
877    msgpack: &[u8],
878    nonce: u64,
879    vault_address: Option<&str>,
880    source: &str,
881) -> Result<[u8; 32], String> {
882    let mut data = msgpack.to_vec();
883    data.extend_from_slice(&nonce.to_be_bytes());
884    match vault_address {
885        None => data.push(0u8),
886        Some(addr) => {
887            data.push(1u8);
888            let addr_bytes = hex::decode(addr.trim_start_matches("0x"))
889                .map_err(|e| format!("invalid vault address: {}", e))?;
890            data.extend_from_slice(&addr_bytes);
891        }
892    }
893
894    let connection_id = keccak256(&data);
895    let agent_type_hash = keccak256(b"Agent(string source,bytes32 connectionId)");
896    let source_hash = keccak256(source.as_bytes());
897    let mut struct_data = [0u8; 96];
898    struct_data[..32].copy_from_slice(&agent_type_hash);
899    struct_data[32..64].copy_from_slice(&source_hash);
900    struct_data[64..96].copy_from_slice(&connection_id);
901    let struct_hash = keccak256(&struct_data);
902
903    let domain_sep = hyperliquid_domain_separator();
904    let mut final_data = Vec::with_capacity(66);
905    final_data.extend_from_slice(b"\x19\x01");
906    final_data.extend_from_slice(&domain_sep);
907    final_data.extend_from_slice(&struct_hash);
908    Ok(keccak256(&final_data))
909}
910
911/// Format an EcdsaSignature into the (r, s, v) hex strings expected by Hyperliquid.
912fn format_ecdsa_signature(sig: &EcdsaSignature) -> (String, String, u8) {
913    let r = format!("0x{}", hex::encode(sig.r));
914    let s = format!("0x{}", hex::encode(sig.s));
915    let v = 27u8 + sig.v;
916    (r, s, v)
917}
918
919/// Sign a digest using a raw private key (direct key mode).
920fn sign_digest_with_key(
921    digest: &[u8; 32],
922    private_key: &str,
923) -> Result<(String, String, u8), String> {
924    use k256::ecdsa::SigningKey;
925
926    let key_bytes = hex::decode(private_key.trim_start_matches("0x"))
927        .map_err(|e| format!("invalid private key: {}", e))?;
928    let signing_key =
929        SigningKey::from_bytes(key_bytes.as_slice().into()).map_err(|e| e.to_string())?;
930    let (sig, recovery_id) = signing_key
931        .sign_prehash_recoverable(digest)
932        .map_err(|e| e.to_string())?;
933
934    let sig_bytes = sig.to_bytes();
935    let r = format!("0x{}", hex::encode(&sig_bytes[..32]));
936    let s = format!("0x{}", hex::encode(&sig_bytes[32..64]));
937    let v = 27u8 + recovery_id.to_byte();
938
939    Ok((r, s, v))
940}
941
942/// Sign using pre-built msgpack bytes (bypassing serde_json field ordering).
943/// Supports both direct private key and external signer.
944async fn sign_with_msgpack(
945    msgpack: &[u8],
946    private_key: Option<&str>,
947    external_signer: Option<&Arc<dyn ExternalSigner>>,
948    nonce: u64,
949    vault_address: Option<&str>,
950    source: &str,
951) -> Result<(String, String, u8), String> {
952    let digest = compute_eip712_digest(msgpack, nonce, vault_address, source)?;
953
954    if let Some(signer) = external_signer {
955        let sig = signer.sign_prehash(&digest).await?;
956        return Ok(format_ecdsa_signature(&sig));
957    }
958
959    let key = private_key.ok_or_else(|| {
960        "no signing method available: provide private_key or external_signer".to_string()
961    })?;
962    sign_digest_with_key(&digest, key)
963}
964
965/// Signs a Hyperliquid exchange action using EIP-712.
966/// Returns (r, s, v) where r and s are "0x"-prefixed hex strings and v is 27 or 28.
967/// Supports both direct private key and external signer.
968async fn sign_action(
969    private_key: Option<&str>,
970    external_signer: Option<&Arc<dyn ExternalSigner>>,
971    action: &Value,
972    vault_address: Option<&str>,
973    nonce: u64,
974    source: &str,
975) -> Result<(String, String, u8), String> {
976    let msgpack_bytes = action_to_canonical_msgpack(action)?;
977    let digest = compute_eip712_digest(&msgpack_bytes, nonce, vault_address, source)?;
978
979    if let Some(signer) = external_signer {
980        let sig = signer.sign_prehash(&digest).await?;
981        return Ok(format_ecdsa_signature(&sig));
982    }
983
984    let key = private_key.ok_or_else(|| {
985        "no signing method available: provide private_key or external_signer".to_string()
986    })?;
987    sign_digest_with_key(&digest, key)
988}
989
990// --- Trait implementations ---
991
992#[async_trait]
993impl guilder_abstraction::TestServer for HyperliquidClient {
994    /// Sends a lightweight allMids request; returns true if the server responds 200 OK.
995    async fn ping(&self) -> Result<bool, String> {
996        // allMids → weight 2
997        self.info_post(serde_json::json!({"type": "allMids"}), 2, "ping")
998            .await
999            .map(|r| r.status().is_success())
1000    }
1001
1002    /// Hyperliquid has no dedicated server-time endpoint; returns local UTC ms.
1003    async fn get_server_time(&self) -> Result<i64, String> {
1004        Ok(std::time::SystemTime::now()
1005            .duration_since(std::time::UNIX_EPOCH)
1006            .map(|d| d.as_millis() as i64)
1007            .unwrap_or(0))
1008    }
1009}
1010
1011#[async_trait]
1012impl guilder_abstraction::GetMarketData for HyperliquidClient {
1013    /// Returns all perpetual asset names from Hyperliquid's meta endpoint.
1014    async fn get_symbol(&self) -> Result<Vec<String>, String> {
1015        // meta → weight 20
1016        let resp = self
1017            .info_post(serde_json::json!({"type": "meta"}), 20, "get_symbol")
1018            .await?;
1019        parse_response::<MetaResponse>(resp)
1020            .await
1021            .map(|r| r.universe.into_iter().map(|a| a.name).collect())
1022    }
1023
1024    /// Returns the current open interest for `symbol` from metaAndAssetCtxs.
1025    async fn get_open_interest(&self, symbol: String) -> Result<Decimal, String> {
1026        // metaAndAssetCtxs → weight 20
1027        let resp = self
1028            .info_post(
1029                serde_json::json!({"type": "metaAndAssetCtxs"}),
1030                20,
1031                "get_open_interest",
1032            )
1033            .await?;
1034        let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
1035            .await?
1036            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
1037        meta.universe
1038            .iter()
1039            .position(|a| a.name == symbol)
1040            .and_then(|i| ctxs.get(i))
1041            .and_then(|ctx| parse_decimal(&ctx.open_interest))
1042            .ok_or_else(|| format!("symbol {} not found", symbol))
1043    }
1044
1045    /// Returns a full AssetContext snapshot for `symbol` from metaAndAssetCtxs.
1046    async fn get_asset_context(&self, symbol: String) -> Result<AssetContext, String> {
1047        // metaAndAssetCtxs → weight 20
1048        let resp = self
1049            .info_post(
1050                serde_json::json!({"type": "metaAndAssetCtxs"}),
1051                20,
1052                "get_asset_context",
1053            )
1054            .await?;
1055        let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
1056            .await?
1057            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
1058        let idx = meta
1059            .universe
1060            .iter()
1061            .position(|a| a.name == symbol)
1062            .ok_or_else(|| format!("symbol {} not found", symbol))?;
1063        let ctx = ctxs
1064            .get(idx)
1065            .ok_or_else(|| format!("symbol {} not found", symbol))?;
1066        Ok(AssetContext {
1067            symbol,
1068            open_interest: parse_decimal(&ctx.open_interest).ok_or("invalid open_interest")?,
1069            funding_rate: parse_decimal(&ctx.funding).ok_or("invalid funding")?,
1070            mark_price: parse_decimal(&ctx.mark_px).ok_or("invalid mark_px")?,
1071            day_volume: parse_decimal(&ctx.day_ntl_vlm).ok_or("invalid day_ntl_vlm")?,
1072            mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
1073            oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
1074            premium: ctx.premium.as_deref().and_then(parse_decimal),
1075            prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
1076            sz_decimals: meta.universe.get(idx).map(|a| a.sz_decimals).unwrap_or(0),
1077        })
1078    }
1079
1080    /// Fetches metaAndAssetCtxs once and returns all asset contexts in universe order.
1081    /// Prefer this over repeated `get_asset_context` calls to avoid rate-limiting.
1082    async fn get_all_asset_contexts(&self) -> Result<Vec<AssetContext>, String> {
1083        // metaAndAssetCtxs → weight 20
1084        let resp = self
1085            .info_post(
1086                serde_json::json!({"type": "metaAndAssetCtxs"}),
1087                20,
1088                "get_all_asset_contexts",
1089            )
1090            .await?;
1091        let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
1092            .await?
1093            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
1094        let mut result = Vec::with_capacity(meta.universe.len());
1095        for (asset, ctx) in meta.universe.iter().zip(ctxs.iter()) {
1096            let Some(open_interest) = parse_decimal(&ctx.open_interest) else {
1097                continue;
1098            };
1099            let Some(funding_rate) = parse_decimal(&ctx.funding) else {
1100                continue;
1101            };
1102            let Some(mark_price) = parse_decimal(&ctx.mark_px) else {
1103                continue;
1104            };
1105            let Some(day_volume) = parse_decimal(&ctx.day_ntl_vlm) else {
1106                continue;
1107            };
1108            result.push(AssetContext {
1109                symbol: asset.name.clone(),
1110                open_interest,
1111                funding_rate,
1112                mark_price,
1113                day_volume,
1114                mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
1115                oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
1116                premium: ctx.premium.as_deref().and_then(parse_decimal),
1117                prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
1118                sz_decimals: asset.sz_decimals,
1119            });
1120        }
1121        Ok(result)
1122    }
1123
1124    /// Returns the number of decimal places for order size for a symbol.
1125    async fn get_sz_decimals(&self, symbol: String) -> Result<i32, String> {
1126        let all = self.get_all_sz_decimals().await?;
1127        all.get(&symbol)
1128            .copied()
1129            .ok_or_else(|| format!("symbol {} not found", symbol))
1130    }
1131
1132    /// Returns sz_decimals for all symbols from the meta universe.
1133    async fn get_all_sz_decimals(&self) -> Result<HashMap<String, i32>, String> {
1134        // metaAndAssetCtxs → weight 20
1135        let resp = self
1136            .info_post(
1137                serde_json::json!({"type": "metaAndAssetCtxs"}),
1138                20,
1139                "get_all_sz_decimals",
1140            )
1141            .await?;
1142        let (meta, _) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
1143            .await?
1144            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
1145        Ok(meta
1146            .universe
1147            .into_iter()
1148            .map(|a| (a.name, a.sz_decimals))
1149            .collect())
1150    }
1151
1152    /// Returns a full L2 orderbook snapshot for `symbol` from the l2Book REST endpoint.
1153    async fn get_l2_orderbook(&self, symbol: String) -> Result<L2Snapshot, String> {
1154        // l2Book → weight 2
1155        let resp = self
1156            .info_post(
1157                serde_json::json!({"type": "l2Book", "coin": symbol}),
1158                2,
1159                "get_l2_orderbook",
1160            )
1161            .await?;
1162        let book: Option<HyperliquidWsBook> = parse_response(resp).await?;
1163        let book = match book {
1164            Some(b) => b,
1165            None => {
1166                return Ok(L2Snapshot {
1167                    symbol,
1168                    bids: vec![],
1169                    asks: vec![],
1170                    sequence: 0,
1171                })
1172            }
1173        };
1174        let bids = book
1175            .levels
1176            .first()
1177            .into_iter()
1178            .flatten()
1179            .filter_map(|level| {
1180                Some(L2Level {
1181                    price: parse_decimal(&level.px)?,
1182                    volume: parse_decimal(&level.sz)?,
1183                })
1184            })
1185            .collect();
1186        let asks = book
1187            .levels
1188            .get(1)
1189            .into_iter()
1190            .flatten()
1191            .filter_map(|level| {
1192                Some(L2Level {
1193                    price: parse_decimal(&level.px)?,
1194                    volume: parse_decimal(&level.sz)?,
1195                })
1196            })
1197            .collect();
1198        Ok(L2Snapshot {
1199            symbol: book.coin,
1200            bids,
1201            asks,
1202            sequence: book.time,
1203        })
1204    }
1205
1206    /// Returns the mid-price of `symbol` (e.g. "BTC") from allMids.
1207    async fn get_price(&self, symbol: String) -> Result<Decimal, String> {
1208        // allMids → weight 2
1209        let resp = self
1210            .info_post(serde_json::json!({"type": "allMids"}), 2, "get_price")
1211            .await?;
1212        parse_response::<HashMap<String, String>>(resp)
1213            .await?
1214            .get(&symbol)
1215            .and_then(|s| parse_decimal(s))
1216            .ok_or_else(|| format!("symbol {} not found", symbol))
1217    }
1218
1219    /// Returns predicted funding rates for all symbols across all venues.
1220    /// Null venue entries (unsupported coins) are silently skipped.
1221    async fn get_predicted_fundings(&self) -> Result<Vec<PredictedFunding>, String> {
1222        // predictedFundings → weight 20
1223        let resp = self
1224            .info_post(
1225                serde_json::json!({"type": "predictedFundings"}),
1226                20,
1227                "get_predicted_fundings",
1228            )
1229            .await?;
1230        let data: PredictedFundingsResponse = parse_response(resp).await?;
1231        let mut result = Vec::new();
1232        for (symbol, venues) in data {
1233            for (venue, entry) in venues {
1234                let Some(entry) = entry else { continue };
1235                if let Some(funding_rate) = parse_decimal(&entry.funding_rate) {
1236                    result.push(PredictedFunding {
1237                        symbol: symbol.clone(),
1238                        venue,
1239                        funding_rate,
1240                        next_funding_time_ms: entry.next_funding_time,
1241                    });
1242                }
1243            }
1244        }
1245        Ok(result)
1246    }
1247}
1248
1249#[async_trait]
1250impl guilder_abstraction::ManageOrder for HyperliquidClient {
1251    /// Places an order on Hyperliquid. Requires `with_auth`. Returns an `OrderPlacement` with
1252    /// the exchange-assigned order ID. Market orders are submitted as aggressive limit orders (IOC).
1253    ///
1254    /// If `cloid` is provided, Hyperliquid attaches it to the order lifecycle — fills and order
1255    /// updates will carry the same cloid back, enabling end-to-end intent tracing without a
1256    /// separate order_id mapping.
1257    ///
1258    /// Trigger orders (`TakeProfit` / `StopLoss`) require `trigger_price` to be set. The order
1259    /// activates when the mark price reaches `triggerPx`, then executes as a market or limit
1260    /// order depending on `time_in_force` (`Ioc` = market, `Gtc` = limit).
1261    async fn place_order(
1262        &self,
1263        symbol: String,
1264        side: OrderSide,
1265        price: Decimal,
1266        volume: Decimal,
1267        order_type: OrderType,
1268        time_in_force: TimeInForce,
1269        trigger_price: Option<Decimal>,
1270        reduce_only: bool,
1271        cloid: Option<String>,
1272    ) -> Result<OrderPlacement, String> {
1273        // Rate limiting is handled in submit_signed_action (non-blocking).
1274        let asset_idx = self.get_asset_index(&symbol).await?;
1275        let is_buy = matches!(side, OrderSide::Buy);
1276
1277        let tif_str = match time_in_force {
1278            TimeInForce::Gtc => "Gtc",
1279            TimeInForce::Ioc => "Ioc",
1280            TimeInForce::Fok => "Fok",
1281            TimeInForce::Alo => "Alo",
1282        };
1283
1284        let cloid_hex = cloid.clone();
1285
1286        // Determine if this is a trigger order
1287        let is_trigger = matches!(order_type, OrderType::TakeProfit | OrderType::StopLoss);
1288
1289        let (order_msgpack, order_type_json) = if is_trigger {
1290            // --- Trigger order ---
1291            let trigger_px = trigger_price
1292                .ok_or_else(|| format!("{:?} order requires trigger_price to be set", order_type))?
1293                .normalize()
1294                .to_string();
1295
1296            let tpsl = match order_type {
1297                OrderType::TakeProfit => "tp",
1298                OrderType::StopLoss => "sl",
1299                _ => unreachable!(),
1300            };
1301
1302            // TimeInForce determines market vs limit on trigger:
1303            // Ioc = market execution (isMarket: true), Gtc/Alo/Fok = limit (isMarket: false)
1304            let is_market = matches!(time_in_force, TimeInForce::Ioc);
1305
1306            // For trigger orders, `p` must be set to the trigger price (not "0"),
1307            // even for market-on-trigger. Hyperliquid validates this field.
1308            let price_str = if is_market {
1309                trigger_px.clone()
1310            } else {
1311                price.normalize().to_string()
1312            };
1313
1314            let msgpack = build_trigger_order_msgpack(
1315                asset_idx,
1316                is_buy,
1317                &price_str,
1318                &volume.normalize().to_string(),
1319                reduce_only,
1320                &trigger_px,
1321                is_market,
1322                tpsl,
1323                cloid_hex.as_deref(),
1324            );
1325
1326            let json_type = if is_market {
1327                format!(
1328                    r#"{{"trigger":{{"isMarket":true,"triggerPx":"{}","tpsl":"{}"}}}}"#,
1329                    trigger_px, tpsl
1330                )
1331            } else {
1332                format!(
1333                    r#"{{"trigger":{{"isMarket":false,"triggerPx":"{}","tpsl":"{}"}}}}"#,
1334                    trigger_px, tpsl
1335                )
1336            };
1337
1338            (msgpack, json_type)
1339        } else {
1340            // --- Regular limit/market order ---
1341            let (order_kind, tif_bytes) = match order_type {
1342                OrderType::Limit => ("limit", tif_str.as_bytes()),
1343                OrderType::Market => ("limit", b"Ioc".as_slice()),
1344                _ => unreachable!(),
1345            };
1346
1347            let price_str = price.normalize().to_string();
1348            let size_str = volume.normalize().to_string();
1349
1350            let msgpack = build_order_msgpack(
1351                asset_idx,
1352                is_buy,
1353                &price_str,
1354                &size_str,
1355                reduce_only,
1356                order_kind,
1357                tif_bytes,
1358                cloid_hex.as_deref(),
1359            );
1360
1361            let json_type = match order_type {
1362                OrderType::Limit => format!(r#"{{"limit":{{"tif":"{tif_str}"}}}}"#),
1363                OrderType::Market => r#"{"limit":{"tif":"Ioc"}}"#.to_string(),
1364                _ => unreachable!(),
1365            };
1366
1367            (msgpack, json_type)
1368        };
1369
1370        // Build the action-level msgpack with Python SDK field order (insertion order):
1371        // type → orders → grouping
1372        let mut action_msgpack = Vec::new();
1373        action_msgpack.push(0x83); // fixmap(3)
1374        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("type".to_string())));
1375        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("order".to_string())));
1376        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("orders".to_string())));
1377        action_msgpack.push(0x91); // fixarray(1)
1378        action_msgpack.extend_from_slice(&order_msgpack);
1379        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("grouping".to_string())));
1380        action_msgpack.extend_from_slice(&value_to_msgpack(&Value::String("na".to_string())));
1381
1382        let cloid_json = if let Some(ref c) = cloid_hex {
1383            format!(r#","c":"{c}""#)
1384        } else {
1385            String::new()
1386        };
1387
1388        let reduce_json = if reduce_only { "true" } else { "false" };
1389
1390        let price_for_json = if is_trigger && matches!(time_in_force, TimeInForce::Ioc) {
1391            // For market-on-trigger, use trigger price (same as msgpack)
1392            if let Some(ref tp) = trigger_price {
1393                tp.normalize().to_string()
1394            } else {
1395                price.normalize().to_string()
1396            }
1397        } else {
1398            price.normalize().to_string()
1399        };
1400
1401        let action_json_str = format!(
1402            r#"{{"type":"order","orders":[{{"a":{asset_idx},"b":{is_buy},"p":"{price}","s":"{size}","r":{reduce_json},"t":{order_type_json}{cloid_json}}}],"grouping":"na"}}"#,
1403            price = price_for_json,
1404            size = volume.normalize().to_string(),
1405        );
1406
1407        // Sign using the canonical msgpack (matching Python's field order)
1408        let private_key = self.require_private_key()?;
1409        let nonce = self.next_nonce();
1410
1411        let (r, s, v) = sign_with_msgpack(
1412            &action_msgpack,
1413            private_key,
1414            self.external_signer.as_ref(),
1415            nonce,
1416            None,
1417            self.network.eip712_source(),
1418        )
1419        .await?;
1420
1421        let payload_str = format!(
1422            r#"{{"action":{},"nonce":{},"signature":{{"r":"{}","s":"{}","v":{}}},"vaultAddress":null,"expiresAfter":null}}"#,
1423            action_json_str, nonce, r, s, v
1424        );
1425
1426        self.rest_limiter.acquire(1).await.map_err(|e| {
1427            format!(
1428                "rate_limited: rest_weight exhausted, retry_after_ms={}",
1429                e.retry_after.as_millis()
1430            )
1431        })?;
1432        self.address_limiter.acquire(1, false).await.map_err(|e| {
1433            format!(
1434                "rate_limited: address quota exhausted, retry_after_ms={}",
1435                e.retry_after.as_millis()
1436            )
1437        })?;
1438
1439        let resp = self
1440            .client
1441            .post(self.network.exchange_url())
1442            .header("Content-Type", "application/json")
1443            .body(payload_str)
1444            .send()
1445            .await
1446            .map_err(|e| e.to_string())?;
1447
1448        let status = resp.status();
1449        if !status.is_success() {
1450            let text = resp.text().await.map_err(|e| e.to_string())?;
1451            return Err(format!("HTTP {status}: {text}"));
1452        }
1453
1454        let body: Value = parse_response(resp).await?;
1455        if body["status"].as_str() == Some("err") {
1456            return Err(body["response"]
1457                .as_str()
1458                .unwrap_or("unknown error")
1459                .to_string());
1460        }
1461        let statuses = &body["response"]["data"]["statuses"][0];
1462
1463        let (oid, returned_cloid, timestamp_ms) = if let Some(resting) = statuses.get("resting") {
1464            let oid = resting["oid"]
1465                .as_i64()
1466                .ok_or_else(|| format!("resting status missing oid: {}", body))?;
1467            let returned_cloid = resting["cloid"].as_str().map(|s: &str| s.to_string());
1468            // resting doesn't include a timestamp
1469            let ts = std::time::SystemTime::now()
1470                .duration_since(std::time::UNIX_EPOCH)
1471                .unwrap()
1472                .as_millis() as i64;
1473            (oid, returned_cloid, ts)
1474        } else if let Some(filled) = statuses.get("filled") {
1475            let oid = filled["oid"]
1476                .as_i64()
1477                .ok_or_else(|| format!("filled status missing oid: {}", body))?;
1478            // filled doesn't include cloid
1479            let ts = std::time::SystemTime::now()
1480                .duration_since(std::time::UNIX_EPOCH)
1481                .unwrap()
1482                .as_millis() as i64;
1483            (oid, None, ts)
1484        } else if let Some(error) = statuses.get("error") {
1485            return Err(error
1486                .as_str()
1487                .unwrap_or("order rejected with unknown error")
1488                .to_string());
1489        } else {
1490            return Err(format!("unexpected order status: {}", body));
1491        };
1492
1493        Ok(OrderPlacement {
1494            order_id: oid,
1495            symbol,
1496            side,
1497            price,
1498            quantity: volume,
1499            timestamp_ms,
1500            cloid: returned_cloid.or(cloid),
1501            order_type,
1502            trigger_price,
1503            reduce_only,
1504        })
1505    }
1506
1507    /// Modifies price and size of an existing order by its order ID. Requires `with_auth`.
1508    /// Fetches the order's current coin and side before submitting the modify action.
1509    async fn change_order_by_cloid(
1510        &self,
1511        cloid: i64,
1512        price: Decimal,
1513        volume: Decimal,
1514    ) -> Result<i64, String> {
1515        let user = self.require_user_address()?;
1516
1517        // openOrders → weight 20; get_asset_index → meta weight 20
1518        let resp = self
1519            .info_post(
1520                serde_json::json!({"type": "openOrders", "user": user}),
1521                20,
1522                "change_order_by_cloid",
1523            )
1524            .await?;
1525        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1526        let order = orders
1527            .iter()
1528            .find(|o| o.oid == cloid)
1529            .ok_or_else(|| format!("order {} not found", cloid))?;
1530
1531        let asset_idx = self.get_asset_index(&order.coin).await?;
1532        let is_buy = order.side == "B";
1533
1534        let action = serde_json::json!({
1535            "type": "batchModify",
1536            "modifies": [{
1537                "oid": cloid,
1538                "order": {
1539                    "a": asset_idx,
1540                    "b": is_buy,
1541                    "p": price.to_string(),
1542                    "s": volume.to_string(),
1543                    "r": false,
1544                    "t": {"limit": {"tif": "Gtc"}}
1545                }
1546            }]
1547        });
1548
1549        self.submit_signed_action(action, None).await?;
1550        Ok(cloid)
1551    }
1552
1553    /// Cancels a single order by its client order ID (cloid). Requires `with_auth`.
1554    /// Fetches open orders to resolve the order ID for the matching cloid, then
1555    /// submits a cancel action using the order ID — this works for all order types
1556    /// including trigger orders (TakeProfit/StopLoss).
1557    async fn cancel_order_by_cloid(&self, cloid: String) -> Result<(), String> {
1558        let user = self.require_user_address()?;
1559
1560        // openOrders → weight 20
1561        let resp = self
1562            .info_post(
1563                serde_json::json!({"type": "openOrders", "user": user}),
1564                20,
1565                "cancel_order_by_cloid",
1566            )
1567            .await?;
1568        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1569        let order = orders
1570            .iter()
1571            .find(|o| o.cloid.as_ref() == Some(&cloid))
1572            .ok_or_else(|| format!("order with cloid {} not found", cloid))?;
1573
1574        // meta → weight 20
1575        let meta_resp = self
1576            .info_post(
1577                serde_json::json!({"type": "meta"}),
1578                20,
1579                "cancel_order_by_cloid",
1580            )
1581            .await?;
1582        let meta: MetaResponse = parse_response(meta_resp).await?;
1583
1584        let asset_idx = meta
1585            .universe
1586            .iter()
1587            .position(|a| a.name == order.coin)
1588            .ok_or_else(|| format!("asset {} not found in meta", order.coin))?;
1589
1590        // Use the same "cancel" action type as cancel_all_order, which is
1591        // proven to work for all order types including trigger orders.
1592        let action = serde_json::json!({
1593            "type": "cancel",
1594            "cancels": [{"a": asset_idx, "o": order.oid}]
1595        });
1596
1597        self.submit_signed_action(action, None).await?;
1598        Ok(())
1599    }
1600
1601    /// Cancels all open orders. Requires `with_auth`.
1602    /// Fetches all open orders and submits a batch cancel in a single signed request.
1603    async fn cancel_all_order(&self) -> Result<bool, String> {
1604        let user = self.require_user_address()?;
1605
1606        // openOrders → weight 20
1607        let resp = self
1608            .info_post(
1609                serde_json::json!({"type": "openOrders", "user": user}),
1610                20,
1611                "cancel_all_order",
1612            )
1613            .await?;
1614        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1615        if orders.is_empty() {
1616            return Ok(true);
1617        }
1618
1619        // meta → weight 20
1620        let meta_resp = self
1621            .info_post(serde_json::json!({"type": "meta"}), 20, "cancel_all_order")
1622            .await?;
1623        let meta: MetaResponse = parse_response(meta_resp).await?;
1624
1625        let cancels: Vec<Value> = orders
1626            .iter()
1627            .filter_map(|o| {
1628                let asset_idx = meta.universe.iter().position(|a| a.name == o.coin)?;
1629                Some(serde_json::json!({"a": asset_idx, "o": o.oid}))
1630            })
1631            .collect();
1632
1633        let action = serde_json::json!({"type": "cancel", "cancels": cancels});
1634        self.submit_signed_action(action, None).await?;
1635        Ok(true)
1636    }
1637}
1638
1639#[async_trait]
1640impl guilder_abstraction::SubscribeMarketData for HyperliquidClient {
1641    fn subscribe_l2_update(&self, symbol: String) -> BoxStream<Result<L2Update, String>> {
1642        Box::pin(stream::iter(vec![Err(format!(
1643            "subscribe_l2_update is unsupported for {symbol}; use subscribe_l2_snapshot"
1644        ))]))
1645    }
1646
1647    fn subscribe_l2_snapshot(&self, symbol: String) -> BoxStream<Result<L2Snapshot, String>> {
1648        let subscription = HyperliquidSubscription::L2Book { coin: symbol };
1649        Box::pin(managed_stream(
1650            self.market_ws_manager.clone(),
1651            subscription,
1652            |msg: HyperliquidWsInboundMessage| {
1653                if let Some(snapshot) = msg.as_l2_snapshot() {
1654                    vec![Ok(snapshot)]
1655                } else {
1656                    vec![]
1657                }
1658            },
1659        ))
1660    }
1661
1662    fn subscribe_asset_context(&self, symbol: String) -> BoxStream<Result<AssetContext, String>> {
1663        let subscription = HyperliquidSubscription::ActiveAssetCtx { coin: symbol };
1664        Box::pin(managed_stream(
1665            self.market_ws_manager.clone(),
1666            subscription,
1667            |msg: HyperliquidWsInboundMessage| {
1668                if let Some(ctx) = msg.as_asset_context() {
1669                    vec![Ok(ctx)]
1670                } else {
1671                    vec![]
1672                }
1673            },
1674        ))
1675    }
1676
1677    fn subscribe_liquidation(&self, user: String) -> BoxStream<Result<Liquidation, String>> {
1678        subscribe_user_stream(
1679            self,
1680            user.clone(),
1681            HyperliquidSubscription::UserEvents { user_addr: user },
1682            |msg: HyperliquidWsInboundMessage| {
1683                if let Some(liq) = msg.as_liquidation() {
1684                    vec![Ok(liq)]
1685                } else {
1686                    vec![]
1687                }
1688            },
1689        )
1690    }
1691
1692    fn subscribe_fill(&self, symbol: String) -> BoxStream<Result<Fill, String>> {
1693        let subscription = HyperliquidSubscription::Trades { coin: symbol };
1694        Box::pin(managed_stream(
1695            self.market_ws_manager.clone(),
1696            subscription,
1697            |msg: HyperliquidWsInboundMessage| {
1698                if let Some(fills) = msg.as_trades() {
1699                    fills.into_iter().map(Ok).collect()
1700                } else {
1701                    vec![]
1702                }
1703            },
1704        ))
1705    }
1706
1707    /// Gracefully shut down all market data subscriptions.
1708    /// Closes all broadcast channels so subscribers exit without reconnecting.
1709    async fn unsubscribe_all(&self) {
1710        self.market_ws_manager.shutdown();
1711    }
1712}
1713
1714fn subscribe_user_stream<T, F>(
1715    client: &HyperliquidClient,
1716    user_addr: String,
1717    subscription: HyperliquidSubscription,
1718    parse: F,
1719) -> BoxStream<Result<T, String>>
1720where
1721    T: Send + 'static,
1722    F: Fn(HyperliquidWsInboundMessage) -> Vec<Result<T, String>> + Send + Sync + 'static,
1723{
1724    let manager = get_or_create_user_manager(
1725        &client.user_ws_managers,
1726        client.ws_send_limiter.clone(),
1727        user_addr,
1728        client.network.ws_url(),
1729    );
1730    Box::pin(async_stream::stream! {
1731        let stream = managed_stream(manager, subscription, parse);
1732        tokio::pin!(stream);
1733        while let Some(item) = stream.next().await {
1734            yield item;
1735        }
1736    })
1737}
1738
1739fn get_or_create_user_manager(
1740    user_ws_managers: &RwLock<HashMap<String, HyperliquidWsManager>>,
1741    ws_send_limiter: WsSendRateLimiter,
1742    user_addr: String,
1743    ws_url: &'static str,
1744) -> HyperliquidWsManager {
1745    {
1746        let managers = user_ws_managers.read().unwrap_or_else(|e| e.into_inner());
1747        if let Some(manager) = managers.get(&user_addr) {
1748            return manager.clone();
1749        }
1750    }
1751
1752    let mut managers = user_ws_managers.write().unwrap_or_else(|e| e.into_inner());
1753    managers
1754        .entry(user_addr.clone())
1755        .or_insert_with(|| HyperliquidWsManager::new(Some(user_addr), ws_send_limiter, ws_url))
1756        .clone()
1757}
1758
1759#[async_trait]
1760impl guilder_abstraction::GetAccountSnapshot for HyperliquidClient {
1761    /// Returns open positions from `clearinghouseState`. Requires `with_auth`.
1762    /// Zero-size positions are filtered out. Positive `szi` = long, negative = short.
1763    async fn get_positions(&self) -> Result<Vec<Position>, String> {
1764        let user = self.require_user_address()?;
1765        // clearinghouseState → weight 2
1766        let resp = self
1767            .info_post(
1768                serde_json::json!({"type": "clearinghouseState", "user": user}),
1769                2,
1770                "get_positions",
1771            )
1772            .await?;
1773        let state: ClearinghouseStateResponse = parse_response(resp).await?;
1774
1775        Ok(state
1776            .asset_positions
1777            .into_iter()
1778            .filter_map(|ap| {
1779                let p = ap.position;
1780                let size = parse_decimal(&p.szi)?;
1781                if size.is_zero() {
1782                    return None;
1783                }
1784                // entryPx is null on TRANSIENT reads (position book race) —
1785                // mapping to 0 poisons downstream uPnL accounting AND the
1786                // reconcile anchor (albatross #114 live: anchor −149 per
1787                // BERA notional). Skip positions without a real entry.
1788                let entry_price = p
1789                    .entry_px
1790                    .as_deref()
1791                    .and_then(parse_decimal);
1792                let Some(entry_price) = entry_price else {
1793                    return None;
1794                };
1795                let side = if size > Decimal::ZERO {
1796                    OrderSide::Buy
1797                } else {
1798                    OrderSide::Sell
1799                };
1800                Some(Position {
1801                    symbol: p.coin,
1802                    side,
1803                    size: size.abs(),
1804                    entry_price,
1805                })
1806            })
1807            .collect())
1808    }
1809
1810    /// Returns resting orders from Hyperliquid's `openOrders` endpoint. Requires `with_auth`.
1811    /// `filled_quantity` is derived as `origSz - sz` (original size minus remaining size).
1812    async fn get_open_orders(&self) -> Result<Vec<OpenOrder>, String> {
1813        let user = self.require_user_address()?;
1814        // openOrders → weight 20
1815        let resp = self
1816            .info_post(
1817                serde_json::json!({"type": "openOrders", "user": user}),
1818                20,
1819                "get_open_orders",
1820            )
1821            .await?;
1822        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1823
1824        Ok(orders
1825            .into_iter()
1826            .filter_map(|o| {
1827                let price = parse_decimal(&o.limit_px)?;
1828                let quantity = parse_decimal(&o.orig_sz)?;
1829                let remaining = parse_decimal(&o.sz)?;
1830                let filled_quantity = quantity - remaining;
1831                let side = if o.side == "B" {
1832                    OrderSide::Buy
1833                } else {
1834                    OrderSide::Sell
1835                };
1836                Some(OpenOrder {
1837                    order_id: o.oid,
1838                    symbol: o.coin,
1839                    side,
1840                    price,
1841                    quantity,
1842                    filled_quantity,
1843                    order_type: None, // openOrders REST endpoint doesn't return order type
1844                    trigger_price: None, // trigger info not included in openOrders response
1845                    reduce_only: false, // default; REST doesn't expose this field
1846                })
1847            })
1848            .collect())
1849    }
1850
1851    /// Returns all per-asset balances from `spotClearinghouseState` with margin health.
1852    /// Also fetches `clearinghouseState` to populate `margin_used` on the USDC entry.
1853    /// Requires `with_auth`.
1854    async fn get_balance(&self) -> Result<Vec<guilder_abstraction::AccountBalance>, String> {
1855        let user = self.require_user_address()?;
1856        // spotClearinghouseState → weight 15
1857        let resp = self
1858            .info_post(
1859                serde_json::json!({"type": "spotClearinghouseState", "user": user}),
1860                15,
1861                "get_balance",
1862            )
1863            .await?;
1864
1865        #[allow(dead_code)]
1866        #[derive(Deserialize)]
1867        struct SpotBalance {
1868            coin: String,
1869            total: String,
1870            hold: String,
1871            #[serde(default)]
1872            token: Option<i32>,
1873            #[serde(default)]
1874            #[serde(rename = "entryNtl")]
1875            entry_ntl: Option<String>,
1876        }
1877
1878        let state: SpotStateResponse = parse_response(resp).await?;
1879
1880        // Fetch the PERP ledger (clearinghouseState). weight 2.
1881        let perp_parsed: Option<ClearinghouseStateResponse> = match self
1882            .info_post(
1883                serde_json::json!({"type": "clearinghouseState", "user": user}),
1884                2,
1885                "get_balance_margin",
1886            )
1887            .await
1888        {
1889            Ok(resp) => parse_response::<ClearinghouseStateResponse>(resp).await.ok(),
1890            Err(_) => None,
1891        };
1892
1893        // Perp-side margin_used lands on the SPOT USDC row for backwards
1894        // compatibility (0.6.x consumers).
1895        let perp_margin_used: Option<Decimal> = perp_parsed.as_ref().and_then(|ch| {
1896            ch.margin_summary
1897                .total_margin_used
1898                .as_deref()
1899                .and_then(parse_decimal)
1900        });
1901
1902        // Spot ledger rows — keep per-asset spot balances (incl. USDC spot
1903        // cash, which is NOT margin under the manual-account model).
1904        let mut rows = map_spot_state(state, perp_margin_used)?;
1905
1906        // Perp ledger row appended LAST (stable position for consumers that
1907        // index rows): token = PERP_LEDGER_TOKEN. equity = accountValue,
1908        // free/usable = accountValue − totalMarginUsed, margin fields filled.
1909        // Failure here is fatal — a trader that cannot read its futures
1910        // ledger must not fall back to spot sizing.
1911        let perp_state: ClearinghouseStateResponse = match perp_parsed {
1912            Some(ch) => ch,
1913            None => {
1914                // The margin probe failed — retry once so a transient
1915                // clearinghouse error cannot silently drop the futures row.
1916                let resp = self
1917                    .info_post(
1918                        serde_json::json!({"type": "clearinghouseState", "user": user}),
1919                        2,
1920                        "get_balance_perp_row",
1921                    )
1922                    .await?;
1923                parse_response::<ClearinghouseStateResponse>(resp).await?
1924            }
1925        };
1926        rows.push(map_perp_state(perp_state)?);
1927
1928        Ok(rows)
1929    }
1930
1931    /// Returns the user's address-level API rate limit budget.
1932    /// Queries Hyperliquid's `userRateLimit` info endpoint for authoritative server-side counts.
1933    async fn get_user_rate_limit(&self) -> Result<guilder_abstraction::UserRateLimit, String> {
1934        let user = self.require_user_address()?;
1935        let resp = self
1936            .info_post(
1937                serde_json::json!({"type": "userRateLimit", "user": user}),
1938                20,
1939                "get_user_rate_limit",
1940            )
1941            .await?;
1942        let val = parse_response::<Value>(resp).await?;
1943
1944        let cumulative_volume = val["cumVlm"]
1945            .as_str()
1946            .and_then(parse_decimal)
1947            .ok_or_else(|| "missing or invalid cumVlm".to_string())?;
1948        let requests_used = val["nRequestsUsed"]
1949            .as_i64()
1950            .ok_or_else(|| "missing or invalid nRequestsUsed".to_string())?;
1951        let requests_cap = val["nRequestsCap"]
1952            .as_i64()
1953            .ok_or_else(|| "missing or invalid nRequestsCap".to_string())?;
1954        let requests_surplus = val["nRequestsSurplus"]
1955            .as_i64()
1956            .ok_or_else(|| "missing or invalid nRequestsSurplus".to_string())?;
1957
1958        Ok(guilder_abstraction::UserRateLimit {
1959            cumulative_volume,
1960            requests_used,
1961            requests_cap,
1962            requests_surplus,
1963        })
1964    }
1965}
1966
1967#[async_trait]
1968impl guilder_abstraction::SubscribeUserEvents for HyperliquidClient {
1969    fn subscribe_user_fills(&self) -> BoxStream<Result<UserFill, String>> {
1970        let Some(addr) = self.user_address.as_ref() else {
1971            return Box::pin(stream::iter(vec![Err(
1972                "user address not registered".to_string()
1973            )]));
1974        };
1975        subscribe_user_stream(
1976            self,
1977            addr.clone(),
1978            HyperliquidSubscription::UserEvents {
1979                user_addr: addr.clone(),
1980            },
1981            |msg: HyperliquidWsInboundMessage| {
1982                if let Some(fills) = msg.as_user_fills() {
1983                    fills.into_iter().map(Ok).collect()
1984                } else {
1985                    vec![]
1986                }
1987            },
1988        )
1989    }
1990
1991    fn subscribe_order_updates(&self) -> BoxStream<Result<OrderUpdate, String>> {
1992        let Some(addr) = self.user_address.as_ref() else {
1993            return Box::pin(stream::iter(vec![Err(
1994                "user address not registered".to_string()
1995            )]));
1996        };
1997        subscribe_user_stream(
1998            self,
1999            addr.clone(),
2000            HyperliquidSubscription::OrderUpdates {
2001                user_addr: addr.clone(),
2002            },
2003            |msg: HyperliquidWsInboundMessage| {
2004                if let Some(updates) = msg.as_order_updates() {
2005                    updates.into_iter().map(Ok).collect()
2006                } else {
2007                    vec![]
2008                }
2009            },
2010        )
2011    }
2012
2013    fn subscribe_funding_payments(&self) -> BoxStream<Result<FundingPayment, String>> {
2014        let Some(addr) = self.user_address.as_ref() else {
2015            return Box::pin(stream::iter(vec![Err(
2016                "user address not registered".to_string()
2017            )]));
2018        };
2019        subscribe_user_stream(
2020            self,
2021            addr.clone(),
2022            HyperliquidSubscription::UserEvents {
2023                user_addr: addr.clone(),
2024            },
2025            |msg: HyperliquidWsInboundMessage| {
2026                if let Some(p) = msg.as_funding_payment() {
2027                    vec![Ok(p)]
2028                } else {
2029                    vec![]
2030                }
2031            },
2032        )
2033    }
2034
2035    fn subscribe_deposits(&self) -> BoxStream<Result<Deposit, String>> {
2036        let Some(addr) = self.user_address.as_ref() else {
2037            return Box::pin(stream::iter(vec![Err(
2038                "user address not registered".to_string()
2039            )]));
2040        };
2041        subscribe_user_stream(
2042            self,
2043            addr.clone(),
2044            HyperliquidSubscription::NonFundingLedger {
2045                user_addr: addr.clone(),
2046            },
2047            |msg: HyperliquidWsInboundMessage| {
2048                if let Some(deps) = msg.as_deposits() {
2049                    deps.into_iter().map(Ok).collect()
2050                } else {
2051                    vec![]
2052                }
2053            },
2054        )
2055    }
2056
2057    fn subscribe_withdrawals(&self) -> BoxStream<Result<Withdrawal, String>> {
2058        let Some(addr) = self.user_address.as_ref() else {
2059            return Box::pin(stream::iter(vec![Err(
2060                "user address not registered".to_string()
2061            )]));
2062        };
2063        subscribe_user_stream(
2064            self,
2065            addr.clone(),
2066            HyperliquidSubscription::NonFundingLedger {
2067                user_addr: addr.clone(),
2068            },
2069            |msg: HyperliquidWsInboundMessage| {
2070                if let Some(wds) = msg.as_withdrawals() {
2071                    wds.into_iter().map(Ok).collect()
2072                } else {
2073                    vec![]
2074                }
2075            },
2076        )
2077    }
2078
2079    /// Subscribe to spot wallet balance updates for the registered user address.
2080    fn subscribe_spot_balance(
2081        &self,
2082    ) -> BoxStream<Result<Vec<guilder_abstraction::AccountBalance>, String>> {
2083        let Some(addr) = self.user_address.as_ref() else {
2084            return Box::pin(stream::iter(vec![Err(
2085                "user address not registered".to_string()
2086            )]));
2087        };
2088        self.subscribe_spot_balance_with_address(addr.clone())
2089    }
2090
2091    /// Subscribe to spot wallet balance updates for a specific address.
2092    fn subscribe_spot_balance_with_address(
2093        &self,
2094        address: String,
2095    ) -> BoxStream<Result<Vec<guilder_abstraction::AccountBalance>, String>> {
2096        subscribe_user_stream(
2097            self,
2098            address.clone(),
2099            HyperliquidSubscription::UserEvents { user_addr: address },
2100            |msg: HyperliquidWsInboundMessage| {
2101                if let Some(balances) = msg.as_spot_balance() {
2102                    vec![Ok(balances)]
2103                } else {
2104                    vec![]
2105                }
2106            },
2107        )
2108    }
2109
2110    async fn unsubscribe_user_events(&self) {
2111        // Unsubscribe from the market-level user manager if we have a user address.
2112        if let Some(addr) = &self.user_address {
2113            self.market_ws_manager.unsubscribe_user(addr);
2114        }
2115        // Also unsubscribe from all per-user managers.
2116        let managers = self
2117            .user_ws_managers
2118            .read()
2119            .unwrap_or_else(|e| e.into_inner());
2120        for (addr, manager) in managers.iter() {
2121            manager.unsubscribe_user(addr);
2122        }
2123    }
2124}
2125
2126#[async_trait]
2127impl guilder_abstraction::SubscribeMarketDataOps for HyperliquidClient {
2128    async fn unsubscribe_market_data(&self, symbol: String) {
2129        self.market_ws_manager.unsubscribe_by_coin(&symbol);
2130    }
2131}
2132
2133#[cfg(test)]
2134mod msgpack_tests {
2135    use super::*;
2136    use serde_json::json;
2137
2138    #[test]
2139    fn test_eip712_digest_known_answers() {
2140        // Reference digests generated with the official hyperliquid-python-sdk
2141        // (msgpack.packb → action_hash → EIP-712 Agent payload → eth_account
2142        // sign + ECDSA address-recovery verified). Regen script logic:
2143        // nonce + action mirrored exactly below.
2144        // 0.6.4 regression guard: source was hardcoded to "a" (mainnet), which
2145        // makes every TESTNET signature verify against the wrong phantom agent
2146        // ("User or API Wallet 0x... does not exist").
2147        let action = json!({
2148            "type": "order",
2149            "orders": [{
2150                "a": 1, "b": true, "p": "1500.5", "s": "0.05",
2151                "r": false, "t": {"limit": {"tif": "Gtc"}}
2152            }],
2153            "grouping": "na"
2154        });
2155        let nonce: u64 = 1758572400123;
2156        let msgpack = action_to_canonical_msgpack(&action).unwrap();
2157
2158        let mainnet = compute_eip712_digest(&msgpack, nonce, None, "a").unwrap();
2159        assert_eq!(
2160            hex::encode(mainnet),
2161            "28cd97dd515629633af463bd7edaab14e61b3941b638d410c317cac7e0aed860"
2162        );
2163
2164        let testnet = compute_eip712_digest(&msgpack, nonce, None, "b").unwrap();
2165        assert_eq!(
2166            hex::encode(testnet),
2167            "d452e53f806773ca6d7ab5d10147518bff0d8b64a4404dcb062e3d2fb3d85c0b"
2168        );
2169
2170        assert_ne!(mainnet, testnet);
2171        assert_eq!(HyperliquidNetwork::Mainnet.eip712_source(), "a");
2172        assert_eq!(HyperliquidNetwork::Testnet.eip712_source(), "b");
2173    }
2174
2175    #[test]
2176    fn test_msgpack_null() {
2177        let result = value_to_msgpack(&Value::Null);
2178        assert_eq!(result, vec![0xc0]);
2179    }
2180
2181    #[test]
2182    fn test_msgpack_bool() {
2183        assert_eq!(value_to_msgpack(&Value::Bool(true)), vec![0xc3]);
2184        assert_eq!(value_to_msgpack(&Value::Bool(false)), vec![0xc2]);
2185    }
2186
2187    #[test]
2188    fn test_msgpack_positive_fixint() {
2189        // 0–127: positive fixint
2190        assert_eq!(value_to_msgpack(&json!(0)), vec![0x00]);
2191        assert_eq!(value_to_msgpack(&json!(1)), vec![0x01]);
2192        assert_eq!(value_to_msgpack(&json!(127)), vec![0x7f]);
2193    }
2194
2195    #[test]
2196    fn test_msgpack_uint8() {
2197        // 128–255: uint8
2198        assert_eq!(value_to_msgpack(&json!(128)), vec![0xcc, 0x80]);
2199        assert_eq!(value_to_msgpack(&json!(255)), vec![0xcc, 0xff]);
2200    }
2201
2202    #[test]
2203    fn test_msgpack_uint16() {
2204        // 256–65535: uint16
2205        assert_eq!(value_to_msgpack(&json!(256)), vec![0xcd, 0x01, 0x00]);
2206        assert_eq!(value_to_msgpack(&json!(65535)), vec![0xcd, 0xff, 0xff]);
2207    }
2208
2209    #[test]
2210    fn test_msgpack_uint32() {
2211        // 65536–4294967295: uint32
2212        assert_eq!(
2213            value_to_msgpack(&json!(65536)),
2214            vec![0xce, 0x00, 0x01, 0x00, 0x00]
2215        );
2216        assert_eq!(
2217            value_to_msgpack(&json!(4294967295u64)),
2218            vec![0xce, 0xff, 0xff, 0xff, 0xff]
2219        );
2220    }
2221
2222    #[test]
2223    fn test_msgpack_uint64() {
2224        // >4294967295: uint64
2225        let big: u64 = 4294967296;
2226        assert_eq!(
2227            value_to_msgpack(&json!(big)),
2228            vec![0xcf, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00]
2229        );
2230    }
2231
2232    #[test]
2233    fn test_msgpack_negative_fixint() {
2234        // -1 to -32: negative fixint
2235        assert_eq!(value_to_msgpack(&json!(-1)), vec![0xff]);
2236        assert_eq!(value_to_msgpack(&json!(-32)), vec![0xe0]);
2237    }
2238
2239    #[test]
2240    fn test_msgpack_int8() {
2241        // -33 to -128: int8
2242        assert_eq!(value_to_msgpack(&json!(-33)), vec![0xd0, 0xdf]);
2243        assert_eq!(value_to_msgpack(&json!(-128)), vec![0xd0, 0x80]);
2244    }
2245
2246    #[test]
2247    fn test_msgpack_int16() {
2248        // -129 to -32768: int16
2249        assert_eq!(value_to_msgpack(&json!(-129)), vec![0xd1, 0xff, 0x7f]);
2250        assert_eq!(value_to_msgpack(&json!(-32768)), vec![0xd1, 0x80, 0x00]);
2251    }
2252
2253    #[test]
2254    fn test_msgpack_int32() {
2255        // -32769 to -2147483648: int32
2256        assert_eq!(
2257            value_to_msgpack(&json!(-32769)),
2258            vec![0xd2, 0xff, 0xff, 0x7f, 0xff]
2259        );
2260        assert_eq!(
2261            value_to_msgpack(&json!(-2147483648i64)),
2262            vec![0xd2, 0x80, 0x00, 0x00, 0x00]
2263        );
2264    }
2265
2266    #[test]
2267    fn test_msgpack_int64() {
2268        let val: i64 = -2147483649;
2269        let result = value_to_msgpack(&json!(val));
2270        assert_eq!(result[0], 0xd3); // int64 marker
2271        assert_eq!(result.len(), 9);
2272    }
2273
2274    #[test]
2275    fn test_msgpack_float() {
2276        let result = value_to_msgpack(&json!(3.14));
2277        assert_eq!(result[0], 0xcb); // float64 marker
2278        assert_eq!(result.len(), 9);
2279    }
2280
2281    #[test]
2282    fn test_msgpack_fixstr() {
2283        // 0–31 bytes: fixstr
2284        assert_eq!(value_to_msgpack(&json!("")), vec![0xa0]);
2285        assert_eq!(value_to_msgpack(&json!("hello")), {
2286            let mut expected = vec![0xa5];
2287            expected.extend_from_slice(b"hello");
2288            expected
2289        });
2290        let s = "a".repeat(31);
2291        let result = value_to_msgpack(&json!(s));
2292        assert_eq!(result[0], 0xbf); // 0xa0 | 31
2293        assert_eq!(result.len(), 32);
2294    }
2295
2296    #[test]
2297    fn test_msgpack_str8() {
2298        let s = "a".repeat(32);
2299        let result = value_to_msgpack(&json!(s));
2300        assert_eq!(result[0], 0xd9); // str8 marker
2301        assert_eq!(result[1], 32);
2302        assert_eq!(result.len(), 34);
2303    }
2304
2305    #[test]
2306    fn test_msgpack_fixarray() {
2307        // 0–15 elements: fixarray
2308        assert_eq!(value_to_msgpack(&json!([])), vec![0x90]);
2309        let result = value_to_msgpack(&json!([1, 2, 3]));
2310        assert_eq!(result[0], 0x93);
2311        assert_eq!(result, vec![0x93, 0x01, 0x02, 0x03]);
2312    }
2313
2314    #[test]
2315    fn test_msgpack_fixmap() {
2316        // 0–15 entries: fixmap
2317        assert_eq!(value_to_msgpack(&json!({})), vec![0x80]);
2318        let result = value_to_msgpack(&json!({"a": 1}));
2319        assert_eq!(result[0], 0x81); // fixmap(1)
2320        assert_eq!(result, {
2321            let mut expected = vec![0x81];
2322            expected.extend_from_slice(&value_to_msgpack(&json!("a")));
2323            expected.extend_from_slice(&value_to_msgpack(&json!(1)));
2324            expected
2325        });
2326    }
2327
2328    #[test]
2329    fn test_msgmap_preserves_insertion_order() {
2330        // Verify keys are serialized in JSON insertion order, not sorted
2331        let val = json!({
2332            "z": 1,
2333            "a": 2,
2334            "m": 3
2335        });
2336        let result = value_to_msgpack(&val);
2337        // fixmap(3)
2338        assert_eq!(result[0], 0x83);
2339        // First key should be "z" (insertion order), not "a" (sorted)
2340        assert_eq!(result[1], 0xa1); // fixstr(1)
2341        assert_eq!(result[2], b'z');
2342    }
2343
2344    #[test]
2345    fn test_msgpack_mixed_array() {
2346        let val = json!([null, true, false, 42, "hi", [1, 2]]);
2347        let result = value_to_msgpack(&val);
2348        assert_eq!(result[0], 0x96); // fixarray(6)
2349        assert_eq!(result[1], 0xc0); // null
2350        assert_eq!(result[2], 0xc3); // true
2351        assert_eq!(result[3], 0xc2); // false
2352        assert_eq!(result[4], 0x2a); // 42
2353                                     // "hi" = fixstr(2) + "hi"
2354        assert_eq!(result[5], 0xa2);
2355        assert_eq!(result[6], b'h');
2356        assert_eq!(result[7], b'i');
2357    }
2358
2359    #[test]
2360    fn test_build_order_msgpack_without_cloid() {
2361        let result = build_order_msgpack(
2362            0,       // asset index
2363            true,    // is_buy
2364            "1000",  // price
2365            "0.1",   // size
2366            false,   // reduce_only
2367            "limit", // order_kind
2368            b"gtc",  // tif
2369            None,    // cloid
2370        );
2371        // fixmap(6)
2372        assert_eq!(result[0], 0x86);
2373    }
2374
2375    #[test]
2376    fn test_build_order_msgpack_with_cloid() {
2377        let result = build_order_msgpack(
2378            0,                // asset index
2379            true,             // is_buy
2380            "1000",           // price
2381            "0.1",            // size
2382            false,            // reduce_only
2383            "limit",          // order_kind
2384            b"gtc",           // tif
2385            Some("my-cloid"), // cloid
2386        );
2387        // fixmap(7)
2388        assert_eq!(result[0], 0x87);
2389    }
2390
2391    #[test]
2392    fn test_action_to_canonical_msgpack() {
2393        let action = json!({
2394            "type": "order",
2395            "orders": [{"a": 0, "b": true, "p": "1000", "s": "0.1", "r": false, "t": {"limit": {"tif": "gtc"}}}],
2396            "grouping": "na"
2397        });
2398        let result = action_to_canonical_msgpack(&action).unwrap();
2399        // fixmap(3)
2400        assert_eq!(result[0], 0x83);
2401    }
2402
2403    #[test]
2404    fn test_msgpack_matches_rmp_serde_for_simple_values() {
2405        // Verify our encoding matches rmp_serde for simple scalar values
2406        use rmp_serde::to_vec;
2407
2408        for val in [
2409            json!(0),
2410            json!(127),
2411            json!(255),
2412            json!(1000),
2413            json!(-1),
2414            json!(-32),
2415            json!(-128),
2416        ] {
2417            let ours = value_to_msgpack(&val);
2418            let theirs = to_vec(&val).unwrap();
2419            assert_eq!(
2420                ours, theirs,
2421                "mismatch for {}: ours={:?}, rmp={:?}",
2422                val, ours, theirs
2423            );
2424        }
2425    }
2426
2427    #[test]
2428    fn test_msgpack_string_encoding() {
2429        use rmp_serde::to_vec;
2430        for val in [
2431            json!(""),
2432            json!("a"),
2433            json!("hello world"),
2434            json!("BTC-USD"),
2435        ] {
2436            let ours = value_to_msgpack(&val);
2437            let theirs = to_vec(&val).unwrap();
2438            assert_eq!(
2439                ours, theirs,
2440                "mismatch for {}: ours={:?}, rmp={:?}",
2441                val, ours, theirs
2442            );
2443        }
2444    }
2445
2446    #[test]
2447    fn test_msgpack_bool_encoding() {
2448        use rmp_serde::to_vec;
2449        let theirs = to_vec(&json!(true)).unwrap();
2450        assert_eq!(value_to_msgpack(&json!(true)), theirs);
2451        let theirs = to_vec(&json!(false)).unwrap();
2452        assert_eq!(value_to_msgpack(&json!(false)), theirs);
2453    }
2454
2455    #[test]
2456    fn test_msgpack_null_encoding() {
2457        use rmp_serde::to_vec;
2458        let theirs = to_vec(&Value::Null).unwrap();
2459        assert_eq!(value_to_msgpack(&Value::Null), theirs);
2460    }
2461
2462    #[test]
2463    fn test_msgpack_nested_object() {
2464        let val = json!({
2465            "outer": {
2466                "inner": 42
2467            }
2468        });
2469        let result = value_to_msgpack(&val);
2470        assert_eq!(result[0], 0x81); // fixmap(1)
2471    }
2472
2473    #[test]
2474    fn test_msgpack_empty_containers() {
2475        assert_eq!(value_to_msgpack(&json!([])), vec![0x90]);
2476        assert_eq!(value_to_msgpack(&json!({})), vec![0x80]);
2477    }
2478}
2479
2480#[async_trait::async_trait]
2481impl guilder_abstraction::ListingEventSource for HyperliquidClient {
2482    /// Authoritative lifecycle events from the venue. Hyperliquid's info API
2483    /// exposes only the CURRENT universe (+ isDelisted flags) — it has no
2484    /// historical listing stream, so this returns one synthetic `list` event
2485    /// per never-delisted symbol (event_time unknown → epoch placeholder) and
2486    /// a `delist` event per isDelisted symbol. The QDB sync layer upserts
2487    /// these into token_registry_events; true FIRST-LISTING times for
2488    /// pre-history coins come from earlier sync runs, not from this call.
2489    async fn get_listing_events(&self) -> Result<Vec<guilder_abstraction::ListingEvent>, String> {
2490        // meta → weight 20
2491        let resp = self
2492            .info_post(
2493                serde_json::json!({"type": "meta"}),
2494                20,
2495                "get_listing_events",
2496            )
2497            .await?;
2498        let meta = parse_response::<MetaResponse>(resp).await?;
2499        Ok(meta
2500            .universe
2501            .into_iter()
2502            .map(|a| guilder_abstraction::ListingEvent {
2503                ticker: a.name,
2504                exchange: "hyperliquid_perp".to_string(),
2505                event: if a.is_delisted { "delist" } else { "list" }.to_string(),
2506                event_time: String::new(),
2507            })
2508            .collect())
2509    }
2510
2511    async fn get_current_universe(&self) -> Result<Vec<guilder_abstraction::SymbolStatus>, String> {
2512        // meta → weight 20
2513        let resp = self
2514            .info_post(
2515                serde_json::json!({"type": "meta"}),
2516                20,
2517                "get_current_universe",
2518            )
2519            .await?;
2520        let meta = parse_response::<MetaResponse>(resp).await?;
2521        Ok(meta
2522            .universe
2523            .into_iter()
2524            .map(|a| guilder_abstraction::SymbolStatus {
2525                ticker: a.name,
2526                exchange: "hyperliquid_perp".to_string(),
2527                is_delisted: a.is_delisted,
2528            })
2529            .collect())
2530    }
2531}
2532
2533#[cfg(test)]
2534mod nonce_tests {
2535    use super::*;
2536
2537    /// Nonce monotonicity (2026-09-28 z4dbg incident): the order lane emits
2538    /// signed bursts (emergency close fan-out) where two actions signed in the
2539    /// same millisecond produced the SAME raw timestamp nonce and HL rejected
2540    /// the loser with `Invalid nonce: duplicate nonce N`. A shared client must
2541    /// therefore hand out strictly increasing nonces across threads, even when
2542    /// the wall clock does not advance between calls.
2543    #[tokio::test]
2544    async fn next_nonce_is_strictly_increasing_across_concurrent_calls() {
2545        let client = HyperliquidClient::new();
2546        let n1 = client.next_nonce();
2547        let n2 = client.next_nonce();
2548        let n3 = client.next_nonce();
2549        assert!(n2 > n1, "nonce must strictly increase: {n1} -> {n2}");
2550        assert!(n3 > n2, "nonce must strictly increase: {n2} -> {n3}");
2551    }
2552
2553    /// The nonce must track the wall clock (HL rejects stale nonces far in the
2554    /// past) while never reusing or going backwards when the clock stalls
2555    /// between two consecutive calls: the floor is last_nonce + 1.
2556    #[tokio::test]
2557    async fn next_nonce_tracks_clock_but_never_repeats_or_regresses() {
2558        let client = HyperliquidClient::new();
2559        let now_ms = std::time::SystemTime::now()
2560            .duration_since(std::time::UNIX_EPOCH)
2561            .unwrap()
2562            .as_millis() as u64;
2563        let n1 = client.next_nonce();
2564        assert!(
2565            n1 >= now_ms,
2566            "first nonce must track the current time: {n1} < {now_ms}"
2567        );
2568        let n2 = client.next_nonce();
2569        assert!(n2 > n1, "even a stalled clock must yield n > last ({n1})");
2570    }
2571
2572    /// Concurrent signed bursts (the emergency fan-out shape): N threads each
2573    /// take one nonce and all N must be distinct AND strictly ordered by
2574    /// acquisition — the exact scenario that produced the duplicate.
2575    #[tokio::test]
2576    async fn concurrent_nonce_take_is_duplicate_free() {
2577        let client = Arc::new(HyperliquidClient::new());
2578        let handles: Vec<_> = (0..8)
2579            .map(|_| {
2580                let c = Arc::clone(&client);
2581                std::thread::spawn(move || c.next_nonce())
2582            })
2583            .collect();
2584        let mut taken: Vec<u64> = handles.into_iter().map(|h| h.join().unwrap()).collect();
2585        taken.sort();
2586        let len = taken.len();
2587        taken.dedup();
2588        assert_eq!(
2589            taken.len(),
2590            len,
2591            "concurrent takes must never produce a duplicate nonce: {taken:?}"
2592        );
2593    }
2594}
2595
2596#[cfg(test)]
2597mod spot_state_tests {
2598    use super::*;
2599
2600    /// EXACT response shape Hyperliquid TESTNET returns for
2601    /// spotClearinghouseState — NO tokenToAvailableAfterMaintenance.
2602    /// (Captured live from trader.daometric.com's deserialize warning,
2603    /// 2026-09-26.) Must parse: USDC total 969, free 969, usable 969.
2604    const TESTNET_SHAPE: &str =
2605        r#"{"balances":[{"coin":"USDC","token":0,"total":"969.0","hold":"0.0","entryNtl":"0.0"}]}"#;
2606
2607    /// MAINNET shape — maintenance map present.
2608    const MAINNET_SHAPE: &str = r#"{"balances":[{"coin":"USDC","token":0,"total":"969.0","hold":"10.0","entryNtl":"0.0"}],"tokenToAvailableAfterMaintenance":[[0,"955.0"]]}"#;
2609
2610    #[test]
2611    fn testnet_spot_state_without_maintenance_map_parses() {
2612        let state: SpotStateResponse =
2613            serde_json::from_str(TESTNET_SHAPE).expect("testnet shape must deserialize");
2614        let balances = map_spot_state(state, None).expect("mapping must succeed");
2615        assert_eq!(balances.len(), 1);
2616        let usdc = &balances[0];
2617        assert_eq!(usdc.token, "USDC");
2618        assert_eq!(usdc.equity.to_string(), "969.0");
2619        assert_eq!(usdc.hold.to_string(), "0.0");
2620        assert_eq!(usdc.free.to_string(), "969.0");
2621        // No maintenance map → usable falls back to free.
2622        assert_eq!(usdc.usable.to_string(), "969.0");
2623        assert_eq!(usdc.safe, None);
2624        assert_eq!(usdc.maintenance, None);
2625    }
2626
2627    #[test]
2628    fn mainnet_spot_state_with_maintenance_map_parses() {
2629        let state: SpotStateResponse =
2630            serde_json::from_str(MAINNET_SHAPE).expect("mainnet shape must deserialize");
2631        let balances = map_spot_state(state, None).expect("mapping must succeed");
2632        let usdc = &balances[0];
2633        assert_eq!(usdc.equity.to_string(), "969.0");
2634        assert_eq!(usdc.free.to_string(), "959.0");
2635        // maintenance map present → usable = min(free, safe) = 955.
2636        assert_eq!(usdc.usable.to_string(), "955.0");
2637        assert_eq!(usdc.safe.map(|d| d.to_string()), Some("955.0".to_string()));
2638        assert_eq!(
2639            usdc.maintenance.map(|d| d.to_string()),
2640            Some("14.0".to_string())
2641        );
2642    }
2643
2644    #[test]
2645    fn usdc_balance_carries_perp_margin_used() {
2646        let state: SpotStateResponse =
2647            serde_json::from_str(TESTNET_SHAPE).expect("testnet shape must deserialize");
2648        let balances = map_spot_state(state, Some(Decimal::from(2))).expect("mapping must succeed");
2649        let usdc = &balances[0];
2650        assert_eq!(
2651            usdc.margin_used.map(|d| d.to_string()),
2652            Some("2".to_string())
2653        );
2654    }
2655
2656    /// 0.7.0 contract: `map_balance_rows` (spot + perp) appends the futures
2657    /// ledger row LAST with the marker token. Locked by the trader's
2658    /// ledger-split consumer (account manager).
2659    #[test]
2660    fn balance_rows_end_with_marked_perp_row() {
2661        let state: SpotStateResponse =
2662            serde_json::from_str(TESTNET_SHAPE).expect("testnet shape must deserialize");
2663        let perp: ClearinghouseStateResponse = serde_json::from_str(
2664            r#"{"marginSummary":{"accountValue":"30.0","totalMarginUsed":"5.0"},"assetPositions":[{"position":{"coin":"BTC","szi":"0.1","entryPx":"100","unrealizedPnl":"2.0"}}]}"#,
2665        )
2666        .expect("perp shape must deserialize");
2667        let mut rows =
2668            map_spot_state(state, Some(Decimal::from(2))).expect("spot rows must map");
2669        rows.push(map_perp_state(perp).expect("perp row must map"));
2670
2671        assert_eq!(rows.len(), 2);
2672        assert_eq!(rows[0].token, "USDC");
2673        assert_eq!(rows[0].margin_used.map(|d| d.to_string()), Some("2".to_string()));
2674        assert_eq!(rows[1].token, crate::PERP_LEDGER_TOKEN);
2675        assert_eq!(rows[1].equity.to_string(), "30.0");
2676        assert_eq!(rows[1].free.to_string(), "25.0");
2677        assert_eq!(rows[1].usable.to_string(), "25.0");
2678        assert_eq!(
2679            rows[1].margin_used.map(|d| d.to_string()),
2680            Some("5.0".to_string())
2681        );
2682        // 0.7.4: settled cash = accountValue − Σ venue uPnL = 30 − 2 = 28
2683        // (NOT totalRawUsd — that's the liquidation basis, ≠ cash)
2684        assert_eq!(rows[1].settled_usd.map(|d| d.to_string()), Some("28.0".into()));
2685        // Marker token must never collide with a real asset symbol.
2686        assert!(rows[..rows.len() - 1]
2687            .iter()
2688            .all(|b| b.token != crate::PERP_LEDGER_TOKEN));
2689    }
2690
2691    /// PERP margin account (clearinghouseState) — the trading money under
2692    /// Hyperliquid's manual-account model (spot and futures are SEPARATE
2693    /// ledgers; SMR trades perps, so sizing reads THIS account).
2694    #[test]
2695    fn perp_margin_account_maps_to_usdc_row() {
2696        let state: ClearinghouseStateResponse = serde_json::from_str(
2697            r#"{"marginSummary":{"accountValue":"30.0","totalNtlPos":"0.0","totalRawUsd":"-70.0","totalMarginUsed":"0.0"},"assetPositions":[]}"#,
2698        )
2699        .expect("perp shape must deserialize");
2700        let usdc = map_perp_state(state).expect("perp mapping must succeed");
2701        // 0.7.0: the perp ledger row is marked, not "USDC" — spot and futures
2702        // are separate ledgers and consumers split rows by token.
2703        assert_eq!(usdc.token, crate::PERP_LEDGER_TOKEN);
2704        // 0.7.2: settled_usd carries totalRawUsd (venue settled cash) when present.
2705        // settled = accountValue − ΣuPnL = 30 − 0 = 30 (totalRawUsd −70 ignored)
2706        assert_eq!(usdc.settled_usd.map(|d| d.to_string()), Some("30.0".into()));
2707        assert_eq!(usdc.equity.to_string(), "30.0"); // accountValue
2708        assert_eq!(usdc.free.to_string(), "30.0"); // accountValue - marginUsed
2709        assert_eq!(usdc.usable.to_string(), "30.0");
2710        assert_eq!(
2711            usdc.margin_used.map(|d| d.to_string()),
2712            Some("0.0".to_string())
2713        );
2714    }
2715}