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