Skip to main content

guilder_client_hyperliquid/
client.rs

1use crate::rate_limiter::{AddressRateLimiter, RestRateLimiter};
2use futures_util::{stream, StreamExt};
3use guilder_abstraction::{
4    self, AssetContext, BoxStream, Deposit, Fill, FundingPayment, L2Update, Liquidation, OpenOrder,
5    OrderPlacement, OrderSide, OrderStatus, OrderType, OrderUpdate, Position, PredictedFunding,
6    Side, TimeInForce, UserFill, Withdrawal,
7};
8use reqwest::Client;
9use rust_decimal::Decimal;
10use serde::Deserialize;
11use serde_json::Value;
12use std::collections::HashMap;
13use std::str::FromStr;
14use std::sync::Arc;
15const HYPERLIQUID_INFO_URL: &str = "https://api.hyperliquid.xyz/info";
16const HYPERLIQUID_EXCHANGE_URL: &str = "https://api.hyperliquid.xyz/exchange";
17
18async fn parse_response<T: for<'de> serde::Deserialize<'de>>(
19    resp: reqwest::Response,
20) -> Result<T, String> {
21    let status = resp.status();
22    let text = resp.text().await.map_err(|e| {
23        format!("failed to read response body (status {status}): {e}")
24    })?;
25
26    if text.is_empty() {
27        return Err(format!(
28            "empty response body from Hyperliquid (HTTP {status})"
29        ));
30    }
31
32    serde_json::from_str(&text).map_err(|e| {
33        let snippet = if text.len() > 512 {
34            format!("{}... ({} bytes total)", &text[..256], text.len())
35        } else {
36            text.clone()
37        };
38        format!("deserialize error (HTTP {status}): {e}: {snippet}")
39    })
40}
41
42pub struct HyperliquidClient {
43    client: Client,
44    user_address: Option<String>,
45    private_key: Option<String>,
46    rest_limiter: Arc<RestRateLimiter>,
47    address_limiter: Arc<AddressRateLimiter>,
48    ws_mux: crate::ws::WsMux,
49}
50
51impl Default for HyperliquidClient {
52    fn default() -> Self {
53        Self::new()
54    }
55}
56
57impl HyperliquidClient {
58    pub fn new() -> Self {
59        HyperliquidClient {
60            client: Client::new(),
61            user_address: None,
62            private_key: None,
63            rest_limiter: Arc::new(RestRateLimiter::new()),
64            address_limiter: Arc::new(AddressRateLimiter::new()),
65            ws_mux: crate::ws::WsMux::new(),
66        }
67    }
68
69    pub fn with_auth(user_address: impl Into<String>, private_key: String) -> Self {
70        HyperliquidClient {
71            client: Client::new(),
72            user_address: Some(user_address.into()),
73            private_key: Some(private_key),
74            rest_limiter: Arc::new(RestRateLimiter::new()),
75            address_limiter: Arc::new(AddressRateLimiter::new()),
76            ws_mux: crate::ws::WsMux::new(),
77        }
78    }
79
80    /// Configure rate limit budgets (rest_weight/min, address_requests).
81    /// Defaults: 1200 rest weight/min, 10000 address requests.
82    pub fn with_budgets(mut self, rest_weight: u32, addr_budget: u64) -> Self {
83        self.rest_limiter = Arc::new(RestRateLimiter::new_with_budget(rest_weight));
84        self.address_limiter = Arc::new(AddressRateLimiter::new_with_budget(addr_budget));
85        self
86    }
87
88    /// POST to the info endpoint, consuming `weight` from the REST rate-limit budget.
89    async fn info_post(
90        &self,
91        body: Value,
92        weight: u32,
93        call: &str,
94    ) -> Result<reqwest::Response, String> {
95        self.rest_limiter.acquire_blocking(weight, call).await;
96        self.client
97            .post(HYPERLIQUID_INFO_URL)
98            .json(&body)
99            .send()
100            .await
101            .map_err(|e| e.to_string())
102    }
103    fn require_user_address(&self) -> Result<String, String> {
104        self.user_address
105            .clone()
106            .ok_or_else(|| "user address required: use HyperliquidClient::with_auth".to_string())
107    }
108
109    fn require_private_key(&self) -> Result<&str, String> {
110        self.private_key
111            .as_deref()
112            .ok_or_else(|| "private key required: use HyperliquidClient::with_auth".to_string())
113    }
114
115    async fn get_asset_index(&self, symbol: &str) -> Result<usize, String> {
116        // `meta` is an "all other info" request → weight 20
117        let resp = self
118            .info_post(serde_json::json!({"type": "meta"}), 20, "get_asset_index")
119            .await?;
120        let meta: MetaResponse = parse_response(resp).await?;
121        meta.universe
122            .iter()
123            .position(|a| a.name == symbol)
124            .ok_or_else(|| format!("symbol {} not found", symbol))
125    }
126
127    async fn submit_signed_action(
128        &self,
129        action: Value,
130        vault_address: Option<&str>,
131    ) -> Result<Value, String> {
132        let private_key = self.require_private_key()?;
133        let nonce = std::time::SystemTime::now()
134            .duration_since(std::time::UNIX_EPOCH)
135            .unwrap()
136            .as_millis() as u64;
137
138        let (r, s, v) = sign_action(private_key, &action, vault_address, nonce)?;
139
140        let payload = serde_json::json!({
141            "action": action,
142            "nonce": nonce,
143            "signature": {"r": r, "s": s, "v": v},
144            "vaultAddress": null
145        });
146
147        // Check both rate limiters non-blocking — fail fast, no retry.
148        self.rest_limiter.acquire(1).await.map_err(|e| {
149            format!(
150                "rate_limited: rest_weight exhausted, retry_after_ms={}",
151                e.retry_after.as_millis()
152            )
153        })?;
154        self.address_limiter.acquire(1, false).await.map_err(|e| {
155            format!(
156                "rate_limited: address quota exhausted, retry_after_ms={}",
157                e.retry_after.as_millis()
158            )
159        })?;
160
161        let resp = self
162            .client
163            .post(HYPERLIQUID_EXCHANGE_URL)
164            .json(&payload)
165            .send()
166            .await
167            .map_err(|e| e.to_string())?;
168
169        let body: Value = parse_response(resp).await?;
170        if body["status"].as_str() == Some("err") {
171            return Err(body["response"]
172                .as_str()
173                .unwrap_or("unknown error")
174                .to_string());
175        }
176        Ok(body)
177    }
178}
179
180// --- REST deserialization types ---
181
182#[derive(Deserialize)]
183struct MetaResponse {
184    universe: Vec<AssetInfo>,
185}
186
187#[derive(Deserialize)]
188struct AssetInfo {
189    name: String,
190}
191
192type MetaAndAssetCtxsResponse = (MetaResponse, Vec<RestAssetCtx>);
193
194#[derive(Deserialize)]
195#[serde(rename_all = "camelCase")]
196#[allow(dead_code)]
197struct RestAssetCtx {
198    open_interest: String,
199    funding: String,
200    mark_px: String,
201    day_ntl_vlm: String,
202    mid_px: Option<String>,
203    oracle_px: Option<String>,
204    premium: Option<String>,
205    prev_day_px: Option<String>,
206}
207
208#[derive(Deserialize)]
209#[serde(rename_all = "camelCase")]
210struct ClearinghouseStateResponse {
211    margin_summary: MarginSummary,
212    asset_positions: Vec<AssetPosition>,
213}
214
215#[derive(Deserialize)]
216#[serde(rename_all = "camelCase")]
217struct MarginSummary {
218    account_value: String,
219}
220
221#[derive(Deserialize)]
222struct AssetPosition {
223    position: PositionDetail,
224}
225
226#[derive(Deserialize)]
227#[serde(rename_all = "camelCase")]
228struct PositionDetail {
229    coin: String,
230    /// positive = long, negative = short
231    szi: String,
232    entry_px: Option<String>,
233}
234
235#[derive(Deserialize)]
236#[serde(rename_all = "camelCase")]
237struct RestOpenOrder {
238    coin: String,
239    side: String,
240    limit_px: String,
241    sz: String,
242    oid: i64,
243    orig_sz: String,
244}
245
246// predictedFundings response: Vec<(coin, Vec<(venue, entry_or_null)>)>
247// The API returns null for venues that don't list the coin.
248type PredictedFundingsResponse = Vec<(String, Vec<(String, Option<PredictedFundingEntry>)>)>;
249
250#[derive(Deserialize)]
251#[serde(rename_all = "camelCase")]
252struct PredictedFundingEntry {
253    funding_rate: String,
254    next_funding_time: i64,
255}
256
257// --- WebSocket envelope and data shapes ---
258
259#[derive(Deserialize)]
260struct WsEnvelope {
261    channel: String,
262    #[serde(default)]
263    data: Value,
264}
265
266#[derive(Deserialize)]
267struct WsBook {
268    coin: String,
269    levels: Vec<Vec<WsLevel>>,
270    time: i64,
271}
272
273#[derive(Deserialize)]
274struct WsLevel {
275    px: String,
276    sz: String,
277}
278
279#[derive(Deserialize)]
280#[serde(rename_all = "camelCase")]
281struct WsAssetCtx {
282    coin: String,
283    ctx: WsPerpsCtx,
284}
285
286#[derive(Deserialize)]
287#[serde(rename_all = "camelCase")]
288struct WsPerpsCtx {
289    open_interest: String,
290    funding: String,
291    mark_px: String,
292    day_ntl_vlm: String,
293    mid_px: Option<String>,
294    oracle_px: Option<String>,
295    premium: Option<String>,
296    prev_day_px: Option<String>,
297}
298
299#[derive(Deserialize)]
300struct WsUserEvent {
301    liquidation: Option<WsLiquidation>,
302    fills: Option<Vec<WsUserFill>>,
303    funding: Option<WsFunding>,
304    spot_state: Option<WsSpotState>,
305}
306
307#[derive(Deserialize)]
308struct WsSpotState {
309    balances: Option<Vec<WsSpotBalance>>,
310}
311
312#[derive(Deserialize)]
313struct WsSpotBalance {
314    coin: String,
315    total: String,
316    hold: String,
317}
318
319#[derive(Deserialize)]
320struct WsLiquidation {
321    liquidated_user: String,
322    liquidated_ntl_pos: String,
323    liquidated_account_value: String,
324}
325
326#[derive(Deserialize)]
327struct WsUserFill {
328    coin: String,
329    px: String,
330    sz: String,
331    side: String,
332    time: i64,
333    oid: i64,
334    fee: String,
335    /// Client order ID assigned at placement, if any.
336    #[serde(default)]
337    cloid: Option<String>,
338}
339
340#[derive(Deserialize)]
341struct WsFunding {
342    time: i64,
343    coin: String,
344    usdc: String,
345}
346
347#[derive(Deserialize)]
348struct WsTrade {
349    coin: String,
350    side: String,
351    px: String,
352    sz: String,
353    time: i64,
354    tid: i64,
355}
356
357#[derive(Deserialize)]
358struct WsOrderUpdate {
359    order: WsOrderInfo,
360    status: String,
361    #[serde(rename = "statusTimestamp")]
362    status_timestamp: i64,
363}
364
365#[derive(Deserialize)]
366#[serde(rename_all = "camelCase")]
367struct WsOrderInfo {
368    coin: String,
369    side: String,
370    limit_px: String,
371    sz: String,
372    oid: i64,
373    orig_sz: String,
374    /// Client order ID assigned at placement, if any.
375    #[serde(default)]
376    cloid: Option<String>,
377}
378
379// --- WebSocket ledger update shapes (deposits / withdrawals) ---
380
381#[derive(Deserialize)]
382struct WsLedgerUpdates {
383    updates: Vec<WsLedgerEntry>,
384}
385
386#[derive(Deserialize)]
387struct WsLedgerEntry {
388    time: i64,
389    delta: WsLedgerDelta,
390}
391
392#[derive(Deserialize)]
393struct WsLedgerDelta {
394    #[serde(rename = "type")]
395    kind: String,
396    usdc: Option<String>,
397}
398
399// --- Helpers ---
400
401fn parse_decimal(s: &str) -> Option<Decimal> {
402    Decimal::from_str(s).ok()
403}
404
405fn keccak256(data: &[u8]) -> [u8; 32] {
406    use sha3::{Digest, Keccak256};
407    Keccak256::digest(data).into()
408}
409
410/// EIP-712 domain separator for Hyperliquid mainnet (Arbitrum, chainId=42161).
411fn hyperliquid_domain_separator() -> [u8; 32] {
412    let type_hash = keccak256(
413        b"EIP712Domain(string name,string version,uint256 chainId,address verifyingContract)",
414    );
415    let name_hash = keccak256(b"Exchange");
416    let version_hash = keccak256(b"1");
417    let mut chain_id = [0u8; 32];
418    chain_id[28..32].copy_from_slice(&42161u32.to_be_bytes());
419    let verifying_contract = [0u8; 32];
420
421    let mut data = [0u8; 160];
422    data[..32].copy_from_slice(&type_hash);
423    data[32..64].copy_from_slice(&name_hash);
424    data[64..96].copy_from_slice(&version_hash);
425    data[96..128].copy_from_slice(&chain_id);
426    data[128..160].copy_from_slice(&verifying_contract);
427    keccak256(&data)
428}
429
430/// Signs a Hyperliquid exchange action using EIP-712.
431/// Returns (r, s, v) where r and s are "0x"-prefixed hex strings and v is 27 or 28.
432fn sign_action(
433    private_key: &str,
434    action: &Value,
435    vault_address: Option<&str>,
436    nonce: u64,
437) -> Result<(String, String, u8), String> {
438    use k256::ecdsa::SigningKey;
439
440    // Step 1: msgpack-encode the action, append nonce + vault flag
441    let msgpack_bytes = rmp_serde::to_vec(action).map_err(|e| e.to_string())?;
442    let mut data = msgpack_bytes;
443    data.extend_from_slice(&nonce.to_be_bytes());
444    match vault_address {
445        None => data.push(0u8),
446        Some(addr) => {
447            data.push(1u8);
448            let addr_bytes = hex::decode(addr.trim_start_matches("0x"))
449                .map_err(|e| format!("invalid vault address: {}", e))?;
450            data.extend_from_slice(&addr_bytes);
451        }
452    }
453    let connection_id = keccak256(&data);
454
455    // Step 2: hash the Agent struct
456    let agent_type_hash = keccak256(b"Agent(string source,bytes32 connectionId)");
457    let source_hash = keccak256(b"a"); // "a" = mainnet
458    let mut struct_data = [0u8; 96];
459    struct_data[..32].copy_from_slice(&agent_type_hash);
460    struct_data[32..64].copy_from_slice(&source_hash);
461    struct_data[64..96].copy_from_slice(&connection_id);
462    let struct_hash = keccak256(&struct_data);
463
464    // Step 3: EIP-712 final hash
465    let domain_sep = hyperliquid_domain_separator();
466    let mut final_data = Vec::with_capacity(66);
467    final_data.extend_from_slice(b"\x19\x01");
468    final_data.extend_from_slice(&domain_sep);
469    final_data.extend_from_slice(&struct_hash);
470    let final_hash = keccak256(&final_data);
471
472    // Step 4: sign with secp256k1
473    let key_bytes = hex::decode(private_key.trim_start_matches("0x"))
474        .map_err(|e| format!("invalid private key: {}", e))?;
475    let signing_key =
476        SigningKey::from_bytes(key_bytes.as_slice().into()).map_err(|e| e.to_string())?;
477    let (sig, recovery_id) = signing_key
478        .sign_prehash_recoverable(&final_hash)
479        .map_err(|e| e.to_string())?;
480
481    let sig_bytes = sig.to_bytes();
482    let r = format!("0x{}", hex::encode(&sig_bytes[..32]));
483    let s = format!("0x{}", hex::encode(&sig_bytes[32..64]));
484    let v = 27u8 + recovery_id.to_byte();
485
486    Ok((r, s, v))
487}
488
489// --- Trait implementations ---
490
491#[allow(async_fn_in_trait)]
492impl guilder_abstraction::TestServer for HyperliquidClient {
493    /// Sends a lightweight allMids request; returns true if the server responds 200 OK.
494    async fn ping(&self) -> Result<bool, String> {
495        // allMids → weight 2
496        self.info_post(serde_json::json!({"type": "allMids"}), 2, "ping")
497            .await
498            .map(|r| r.status().is_success())
499    }
500
501    /// Hyperliquid has no dedicated server-time endpoint; returns local UTC ms.
502    async fn get_server_time(&self) -> Result<i64, String> {
503        Ok(std::time::SystemTime::now()
504            .duration_since(std::time::UNIX_EPOCH)
505            .map(|d| d.as_millis() as i64)
506            .unwrap_or(0))
507    }
508}
509
510#[allow(async_fn_in_trait)]
511impl guilder_abstraction::GetMarketData for HyperliquidClient {
512    /// Returns all perpetual asset names from Hyperliquid's meta endpoint.
513    async fn get_symbol(&self) -> Result<Vec<String>, String> {
514        // meta → weight 20
515        let resp = self
516            .info_post(serde_json::json!({"type": "meta"}), 20, "get_symbol")
517            .await?;
518        parse_response::<MetaResponse>(resp)
519            .await
520            .map(|r| r.universe.into_iter().map(|a| a.name).collect())
521    }
522
523    /// Returns the current open interest for `symbol` from metaAndAssetCtxs.
524    async fn get_open_interest(&self, symbol: String) -> Result<Decimal, String> {
525        // metaAndAssetCtxs → weight 20
526        let resp = self
527            .info_post(
528                serde_json::json!({"type": "metaAndAssetCtxs"}),
529                20,
530                "get_open_interest",
531            )
532            .await?;
533        let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
534            .await?
535            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
536        meta.universe
537            .iter()
538            .position(|a| a.name == symbol)
539            .and_then(|i| ctxs.get(i))
540            .and_then(|ctx| parse_decimal(&ctx.open_interest))
541            .ok_or_else(|| format!("symbol {} not found", symbol))
542    }
543
544    /// Returns a full AssetContext snapshot for `symbol` from metaAndAssetCtxs.
545    async fn get_asset_context(&self, symbol: String) -> Result<AssetContext, String> {
546        // metaAndAssetCtxs → weight 20
547        let resp = self
548            .info_post(
549                serde_json::json!({"type": "metaAndAssetCtxs"}),
550                20,
551                "get_asset_context",
552            )
553            .await?;
554        let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
555            .await?
556            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
557        let idx = meta
558            .universe
559            .iter()
560            .position(|a| a.name == symbol)
561            .ok_or_else(|| format!("symbol {} not found", symbol))?;
562        let ctx = ctxs
563            .get(idx)
564            .ok_or_else(|| format!("symbol {} not found", symbol))?;
565        Ok(AssetContext {
566            symbol,
567            open_interest: parse_decimal(&ctx.open_interest).ok_or("invalid open_interest")?,
568            funding_rate: parse_decimal(&ctx.funding).ok_or("invalid funding")?,
569            mark_price: parse_decimal(&ctx.mark_px).ok_or("invalid mark_px")?,
570            day_volume: parse_decimal(&ctx.day_ntl_vlm).ok_or("invalid day_ntl_vlm")?,
571            mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
572            oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
573            premium: ctx.premium.as_deref().and_then(parse_decimal),
574            prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
575        })
576    }
577
578    /// Fetches metaAndAssetCtxs once and returns all asset contexts in universe order.
579    /// Prefer this over repeated `get_asset_context` calls to avoid rate-limiting.
580    async fn get_all_asset_contexts(&self) -> Result<Vec<AssetContext>, String> {
581        // metaAndAssetCtxs → weight 20
582        let resp = self
583            .info_post(
584                serde_json::json!({"type": "metaAndAssetCtxs"}),
585                20,
586                "get_all_asset_contexts",
587            )
588            .await?;
589        let (meta, ctxs) = parse_response::<Option<MetaAndAssetCtxsResponse>>(resp)
590            .await?
591            .ok_or_else(|| "metaAndAssetCtxs returned null".to_string())?;
592        let mut result = Vec::with_capacity(meta.universe.len());
593        for (asset, ctx) in meta.universe.iter().zip(ctxs.iter()) {
594            let Some(open_interest) = parse_decimal(&ctx.open_interest) else {
595                continue;
596            };
597            let Some(funding_rate) = parse_decimal(&ctx.funding) else {
598                continue;
599            };
600            let Some(mark_price) = parse_decimal(&ctx.mark_px) else {
601                continue;
602            };
603            let Some(day_volume) = parse_decimal(&ctx.day_ntl_vlm) else {
604                continue;
605            };
606            result.push(AssetContext {
607                symbol: asset.name.clone(),
608                open_interest,
609                funding_rate,
610                mark_price,
611                day_volume,
612                mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
613                oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
614                premium: ctx.premium.as_deref().and_then(parse_decimal),
615                prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
616            });
617        }
618        Ok(result)
619    }
620
621    /// Returns a full L2 orderbook snapshot for `symbol` from the l2Book REST endpoint.
622    /// Levels are returned as individual `L2Update` items; all share the same `sequence` (timestamp).
623    async fn get_l2_orderbook(&self, symbol: String) -> Result<Vec<L2Update>, String> {
624        // l2Book → weight 2
625        let resp = self
626            .info_post(
627                serde_json::json!({"type": "l2Book", "coin": symbol}),
628                2,
629                "get_l2_orderbook",
630            )
631            .await?;
632        let book: Option<WsBook> = parse_response(resp).await?;
633        let book = match book {
634            Some(b) => b,
635            None => return Ok(vec![]),
636        };
637        let mut levels = Vec::new();
638        for level in book.levels.first().into_iter().flatten() {
639            if let (Some(price), Some(volume)) =
640                (parse_decimal(&level.px), parse_decimal(&level.sz))
641            {
642                levels.push(L2Update {
643                    symbol: book.coin.clone(),
644                    price,
645                    volume,
646                    side: Side::Ask,
647                    sequence: book.time,
648                });
649            }
650        }
651        for level in book.levels.get(1).into_iter().flatten() {
652            if let (Some(price), Some(volume)) =
653                (parse_decimal(&level.px), parse_decimal(&level.sz))
654            {
655                levels.push(L2Update {
656                    symbol: book.coin.clone(),
657                    price,
658                    volume,
659                    side: Side::Bid,
660                    sequence: book.time,
661                });
662            }
663        }
664        Ok(levels)
665    }
666
667    /// Returns the mid-price of `symbol` (e.g. "BTC") from allMids.
668    async fn get_price(&self, symbol: String) -> Result<Decimal, String> {
669        // allMids → weight 2
670        let resp = self
671            .info_post(serde_json::json!({"type": "allMids"}), 2, "get_price")
672            .await?;
673        parse_response::<HashMap<String, String>>(resp)
674            .await?
675            .get(&symbol)
676            .and_then(|s| parse_decimal(s))
677            .ok_or_else(|| format!("symbol {} not found", symbol))
678    }
679
680    /// Returns predicted funding rates for all symbols across all venues.
681    /// Null venue entries (unsupported coins) are silently skipped.
682    async fn get_predicted_fundings(&self) -> Result<Vec<PredictedFunding>, String> {
683        // predictedFundings → weight 20
684        let resp = self
685            .info_post(
686                serde_json::json!({"type": "predictedFundings"}),
687                20,
688                "get_predicted_fundings",
689            )
690            .await?;
691        let data: PredictedFundingsResponse = parse_response(resp).await?;
692        let mut result = Vec::new();
693        for (symbol, venues) in data {
694            for (venue, entry) in venues {
695                let Some(entry) = entry else { continue };
696                if let Some(funding_rate) = parse_decimal(&entry.funding_rate) {
697                    result.push(PredictedFunding {
698                        symbol: symbol.clone(),
699                        venue,
700                        funding_rate,
701                        next_funding_time_ms: entry.next_funding_time,
702                    });
703                }
704            }
705        }
706        Ok(result)
707    }
708}
709
710#[allow(async_fn_in_trait)]
711impl guilder_abstraction::ManageOrder for HyperliquidClient {
712    /// Places an order on Hyperliquid. Requires `with_auth`. Returns an `OrderPlacement` with
713    /// the exchange-assigned order ID. Market orders are submitted as aggressive limit orders (IOC).
714    ///
715    /// If `cloid` is provided, Hyperliquid attaches it to the order lifecycle — fills and order
716    /// updates will carry the same cloid back, enabling end-to-end intent tracing without a
717    /// separate order_id mapping.
718    async fn place_order(
719        &self,
720        symbol: String,
721        side: OrderSide,
722        price: Decimal,
723        volume: Decimal,
724        order_type: OrderType,
725        time_in_force: TimeInForce,
726        cloid: Option<String>,
727    ) -> Result<OrderPlacement, String> {
728        // Rate limiting is handled in submit_signed_action (non-blocking).
729        let asset_idx = self.get_asset_index(&symbol).await?;
730        let is_buy = matches!(side, OrderSide::Buy);
731
732        let tif_str = match time_in_force {
733            TimeInForce::Gtc => "Gtc",
734            TimeInForce::Ioc => "Ioc",
735            TimeInForce::Fok => "Fok",
736        };
737        // Market orders are IOC limit orders at a wide price
738        let order_type_val = match order_type {
739            OrderType::Limit => serde_json::json!({"limit": {"tif": tif_str}}),
740            OrderType::Market => serde_json::json!({"limit": {"tif": "Ioc"}}),
741        };
742
743        let mut order_json = serde_json::json!({
744            "a": asset_idx,
745            "b": is_buy,
746            "p": price.to_string(),
747            "s": volume.to_string(),
748            "r": false,
749            "t": order_type_val
750        });
751        if let Some(ref c) = cloid {
752            order_json["c"] = serde_json::json!(c);
753        }
754
755        let action = serde_json::json!({
756            "type": "order",
757            "orders": [order_json],
758            "grouping": "na"
759        });
760
761        let resp = self.submit_signed_action(action, None).await?;
762        let oid = resp["response"]["data"]["statuses"][0]["resting"]["oid"]
763            .as_i64()
764            .or_else(|| resp["response"]["data"]["statuses"][0]["filled"]["oid"].as_i64())
765            .ok_or_else(|| format!("unexpected response: {}", resp))?;
766
767        let timestamp_ms = std::time::SystemTime::now()
768            .duration_since(std::time::UNIX_EPOCH)
769            .unwrap()
770            .as_millis() as i64;
771
772        Ok(OrderPlacement {
773            order_id: oid,
774            symbol,
775            side,
776            price,
777            quantity: volume,
778            timestamp_ms,
779            cloid,
780        })
781    }
782
783    /// Modifies price and size of an existing order by its order ID. Requires `with_auth`.
784    /// Fetches the order's current coin and side before submitting the modify action.
785    async fn change_order_by_cloid(
786        &self,
787        cloid: i64,
788        price: Decimal,
789        volume: Decimal,
790    ) -> Result<i64, String> {
791        let user = self.require_user_address()?;
792
793        // openOrders → weight 20; get_asset_index → meta weight 20
794        let resp = self
795            .info_post(
796                serde_json::json!({"type": "openOrders", "user": user}),
797                20,
798                "change_order_by_cloid",
799            )
800            .await?;
801        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
802        let order = orders
803            .iter()
804            .find(|o| o.oid == cloid)
805            .ok_or_else(|| format!("order {} not found", cloid))?;
806
807        let asset_idx = self.get_asset_index(&order.coin).await?;
808        let is_buy = order.side == "B";
809
810        let action = serde_json::json!({
811            "type": "batchModify",
812            "modifies": [{
813                "oid": cloid,
814                "order": {
815                    "a": asset_idx,
816                    "b": is_buy,
817                    "p": price.to_string(),
818                    "s": volume.to_string(),
819                    "r": false,
820                    "t": {"limit": {"tif": "Gtc"}}
821                }
822            }]
823        });
824
825        self.submit_signed_action(action, None).await?;
826        Ok(cloid)
827    }
828
829    /// Cancels a single order by its order ID. Requires `with_auth`.
830    /// Fetches open orders to resolve the coin/asset before cancelling.
831    async fn cancel_order(&self, cloid: i64) -> Result<i64, String> {
832        let user = self.require_user_address()?;
833
834        // openOrders → weight 20
835        let resp = self
836            .info_post(
837                serde_json::json!({"type": "openOrders", "user": user}),
838                20,
839                "cancel_order",
840            )
841            .await?;
842        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
843        let order = orders
844            .iter()
845            .find(|o| o.oid == cloid)
846            .ok_or_else(|| format!("order {} not found", cloid))?;
847
848        let asset_idx = self.get_asset_index(&order.coin).await?;
849        let action = serde_json::json!({
850            "type": "cancel",
851            "cancels": [{"a": asset_idx, "o": cloid}]
852        });
853
854        self.submit_signed_action(action, None).await?;
855        Ok(cloid)
856    }
857
858    /// Cancels all open orders. Requires `with_auth`.
859    /// Fetches all open orders and submits a batch cancel in a single signed request.
860    async fn cancel_all_order(&self) -> Result<bool, String> {
861        let user = self.require_user_address()?;
862
863        // openOrders → weight 20
864        let resp = self
865            .info_post(
866                serde_json::json!({"type": "openOrders", "user": user}),
867                20,
868                "cancel_all_order",
869            )
870            .await?;
871        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
872        if orders.is_empty() {
873            return Ok(true);
874        }
875
876        // meta → weight 20
877        let meta_resp = self
878            .info_post(serde_json::json!({"type": "meta"}), 20, "cancel_all_order")
879            .await?;
880        let meta: MetaResponse = parse_response(meta_resp).await?;
881
882        let cancels: Vec<Value> = orders
883            .iter()
884            .filter_map(|o| {
885                let asset_idx = meta.universe.iter().position(|a| a.name == o.coin)?;
886                Some(serde_json::json!({"a": asset_idx, "o": o.oid}))
887            })
888            .collect();
889
890        let action = serde_json::json!({"type": "cancel", "cancels": cancels});
891        self.submit_signed_action(action, None).await?;
892        Ok(true)
893    }
894}
895
896#[allow(async_fn_in_trait)]
897impl guilder_abstraction::SubscribeMarketData for HyperliquidClient {
898    fn subscribe_l2_update(&self, symbol: String) -> BoxStream<Result<L2Update, String>> {
899        let sub = serde_json::json!({
900            "method": "subscribe",
901            "subscription": {"type": "l2Book", "coin": symbol.clone()}
902        });
903        let key = crate::ws::SubKey {
904            channel: "l2Book".to_string(),
905            routing_key: symbol,
906        };
907        let stream = self.ws_mux.subscribe(key, sub);
908        Box::pin(async_stream::stream! {
909            for await msg in stream {
910                let Ok(env) = serde_json::from_str::<WsEnvelope>(&msg) else {
911                    continue;
912                };
913                if env.channel != "l2Book" {
914                    continue;
915                }
916                let Ok(book) = serde_json::from_value::<WsBook>(env.data) else {
917                    continue;
918                };
919                for level in book.levels.first().into_iter().flatten() {
920                    if let (Some(price), Some(volume)) =
921                        (parse_decimal(&level.px), parse_decimal(&level.sz))
922                    {
923                        yield Ok(L2Update {
924                            symbol: book.coin.clone(),
925                            price,
926                            volume,
927                            side: Side::Ask,
928                            sequence: book.time,
929                        });
930                    }
931                }
932                for level in book.levels.get(1).into_iter().flatten() {
933                    if let (Some(price), Some(volume)) =
934                        (parse_decimal(&level.px), parse_decimal(&level.sz))
935                    {
936                        yield Ok(L2Update {
937                            symbol: book.coin.clone(),
938                            price,
939                            volume,
940                            side: Side::Bid,
941                            sequence: book.time,
942                        });
943                    }
944                }
945            }
946        })
947    }
948
949    fn subscribe_asset_context(&self, symbol: String) -> BoxStream<Result<AssetContext, String>> {
950        let sub = serde_json::json!({
951            "method": "subscribe",
952            "subscription": {"type": "activeAssetCtx", "coin": symbol.clone()}
953        });
954        let key = crate::ws::SubKey {
955            channel: "activeAssetCtx".to_string(),
956            routing_key: symbol,
957        };
958        let stream = self.ws_mux.subscribe(key, sub);
959        Box::pin(async_stream::stream! {
960            for await msg in stream {
961                let Ok(env) = serde_json::from_str::<WsEnvelope>(&msg) else {
962                    continue;
963                };
964                if env.channel != "activeAssetCtx" {
965                    continue;
966                }
967                let Ok(update) = serde_json::from_value::<WsAssetCtx>(env.data) else {
968                    continue;
969                };
970                let ctx = &update.ctx;
971                let (Some(open_interest), Some(funding_rate), Some(mark_price), Some(day_volume)) = (
972                    parse_decimal(&ctx.open_interest),
973                    parse_decimal(&ctx.funding),
974                    parse_decimal(&ctx.mark_px),
975                    parse_decimal(&ctx.day_ntl_vlm),
976                ) else {
977                    continue;
978                };
979                yield Ok(AssetContext {
980                    symbol: update.coin,
981                    open_interest,
982                    funding_rate,
983                    mark_price,
984                    day_volume,
985                    mid_price: ctx.mid_px.as_deref().and_then(parse_decimal),
986                    oracle_price: ctx.oracle_px.as_deref().and_then(parse_decimal),
987                    premium: ctx.premium.as_deref().and_then(parse_decimal),
988                    prev_day_price: ctx.prev_day_px.as_deref().and_then(parse_decimal),
989                });
990            }
991        })
992    }
993
994    fn subscribe_liquidation(&self, user: String) -> BoxStream<Result<Liquidation, String>> {
995        let sub = serde_json::json!({
996            "method": "subscribe",
997            "subscription": {"type": "userEvents", "user": user.clone()}
998        });
999        let key = crate::ws::SubKey {
1000            channel: "userEvents".to_string(),
1001            routing_key: user,
1002        };
1003        let raw_stream = self.ws_mux.subscribe(key, sub);
1004        Box::pin(
1005            raw_stream
1006                .filter_map(|text| async move {
1007                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1008                        return None;
1009                    };
1010                    if env.channel != "userEvents" {
1011                        return None;
1012                    }
1013                    let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1014                        return None;
1015                    };
1016                    let liq = event.liquidation?;
1017                    let (Some(notional_position), Some(account_value)) = (
1018                        parse_decimal(&liq.liquidated_ntl_pos),
1019                        parse_decimal(&liq.liquidated_account_value),
1020                    ) else {
1021                        return None;
1022                    };
1023                    let item = Liquidation {
1024                        symbol: String::new(),
1025                        side: OrderSide::Sell,
1026                        liquidated_user: liq.liquidated_user,
1027                        notional_position,
1028                        account_value,
1029                    };
1030                    Some(stream::iter(vec![Ok(item)]))
1031                })
1032                .flatten(),
1033        )
1034    }
1035
1036    fn subscribe_fill(&self, symbol: String) -> BoxStream<Result<Fill, String>> {
1037        let sub = serde_json::json!({
1038            "method": "subscribe",
1039            "subscription": {"type": "trades", "coin": symbol.clone()}
1040        });
1041        let key = crate::ws::SubKey {
1042            channel: "trades".to_string(),
1043            routing_key: symbol,
1044        };
1045        let stream = self.ws_mux.subscribe(key, sub);
1046        Box::pin(async_stream::stream! {
1047            for await msg in stream {
1048                let Ok(env) = serde_json::from_str::<WsEnvelope>(&msg) else {
1049                    continue;
1050                };
1051                if env.channel != "trades" {
1052                    continue;
1053                }
1054                let Ok(trades) = serde_json::from_value::<Vec<WsTrade>>(env.data) else {
1055                    continue;
1056                };
1057                for trade in trades {
1058                    let side = if trade.side == "B" {
1059                        OrderSide::Buy
1060                    } else {
1061                        OrderSide::Sell
1062                    };
1063                    let price = parse_decimal(&trade.px);
1064                    let volume = parse_decimal(&trade.sz);
1065                    if let (Some(price), Some(volume)) = (price, volume) {
1066                        yield Ok(Fill {
1067                            symbol: trade.coin,
1068                            price,
1069                            volume,
1070                            side,
1071                            timestamp_ms: trade.time,
1072                            trade_id: trade.tid,
1073                        });
1074                    }
1075                }
1076            }
1077        })
1078    }
1079}
1080
1081#[allow(async_fn_in_trait)]
1082impl guilder_abstraction::GetAccountSnapshot for HyperliquidClient {
1083    /// Returns open positions from `clearinghouseState`. Requires `with_auth`.
1084    /// Zero-size positions are filtered out. Positive `szi` = long, negative = short.
1085    async fn get_positions(&self) -> Result<Vec<Position>, String> {
1086        let user = self.require_user_address()?;
1087        // clearinghouseState → weight 2
1088        let resp = self
1089            .info_post(
1090                serde_json::json!({"type": "clearinghouseState", "user": user}),
1091                2,
1092                "get_positions",
1093            )
1094            .await?;
1095        let state: ClearinghouseStateResponse = parse_response(resp).await?;
1096
1097        Ok(state
1098            .asset_positions
1099            .into_iter()
1100            .filter_map(|ap| {
1101                let p = ap.position;
1102                let size = parse_decimal(&p.szi)?;
1103                if size.is_zero() {
1104                    return None;
1105                }
1106                let entry_price = p
1107                    .entry_px
1108                    .as_deref()
1109                    .and_then(parse_decimal)
1110                    .unwrap_or_default();
1111                let side = if size > Decimal::ZERO {
1112                    OrderSide::Buy
1113                } else {
1114                    OrderSide::Sell
1115                };
1116                Some(Position {
1117                    symbol: p.coin,
1118                    side,
1119                    size: size.abs(),
1120                    entry_price,
1121                })
1122            })
1123            .collect())
1124    }
1125
1126    /// Returns resting orders from Hyperliquid's `openOrders` endpoint. Requires `with_auth`.
1127    /// `filled_quantity` is derived as `origSz - sz` (original size minus remaining size).
1128    async fn get_open_orders(&self) -> Result<Vec<OpenOrder>, String> {
1129        let user = self.require_user_address()?;
1130        // openOrders → weight 20
1131        let resp = self
1132            .info_post(
1133                serde_json::json!({"type": "openOrders", "user": user}),
1134                20,
1135                "get_open_orders",
1136            )
1137            .await?;
1138        let orders: Vec<RestOpenOrder> = parse_response(resp).await?;
1139
1140        Ok(orders
1141            .into_iter()
1142            .filter_map(|o| {
1143                let price = parse_decimal(&o.limit_px)?;
1144                let quantity = parse_decimal(&o.orig_sz)?;
1145                let remaining = parse_decimal(&o.sz)?;
1146                let filled_quantity = quantity - remaining;
1147                let side = if o.side == "B" {
1148                    OrderSide::Buy
1149                } else {
1150                    OrderSide::Sell
1151                };
1152                Some(OpenOrder {
1153                    order_id: o.oid,
1154                    symbol: o.coin,
1155                    side,
1156                    price,
1157                    quantity,
1158                    filled_quantity,
1159                })
1160            })
1161            .collect())
1162    }
1163
1164    /// Returns total account value (collateral) from `clearinghouseState`. Requires `with_auth`.
1165    async fn get_collateral(&self) -> Result<Decimal, String> {
1166        let user = self.require_user_address()?;
1167        // clearinghouseState → weight 2
1168        let resp = self
1169            .info_post(
1170                serde_json::json!({"type": "clearinghouseState", "user": user}),
1171                2,
1172                "get_collateral",
1173            )
1174            .await?;
1175        let state: ClearinghouseStateResponse = parse_response(resp).await?;
1176        parse_decimal(&state.margin_summary.account_value)
1177            .ok_or_else(|| "invalid account value".to_string())
1178    }
1179
1180    /// Returns all spot wallet balances from `spotState`. Requires `with_auth`.
1181    async fn get_spot_balance(&self) -> Result<Vec<guilder_abstraction::Balance>, String> {
1182        let user = self.require_user_address()?;
1183        // spotClearinghouseState → weight 15
1184        let resp = self
1185            .info_post(
1186                serde_json::json!({"type": "spotClearinghouseState", "user": user}),
1187                15,
1188                "get_spot_balance",
1189            )
1190            .await?;
1191
1192        #[derive(Deserialize)]
1193        struct SpotStateResponse {
1194            balances: Vec<SpotBalance>,
1195        }
1196
1197        #[allow(dead_code)]
1198        #[derive(Deserialize)]
1199        struct SpotBalance {
1200            coin: String,
1201            total: String,
1202            hold: String,
1203            #[serde(default)]
1204            token: Option<i32>,
1205            #[serde(default)]
1206            #[serde(rename = "entryNtl")]
1207            entry_ntl: Option<String>,
1208        }
1209
1210        let state: SpotStateResponse = parse_response(resp).await?;
1211
1212        state
1213            .balances
1214            .into_iter()
1215            .map(|balance| {
1216                let total = parse_decimal(&balance.total)
1217                    .ok_or_else(|| "invalid total balance".to_string())?;
1218                let locked = parse_decimal(&balance.hold)
1219                    .ok_or_else(|| "invalid hold balance".to_string())?;
1220                let available = total - locked;
1221
1222                Ok(guilder_abstraction::Balance {
1223                    coin: balance.coin,
1224                    total,
1225                    available,
1226                    locked,
1227                })
1228            })
1229            .collect()
1230    }
1231
1232    /// Returns clearing house collateral balance for an asset. Currently returns the total collateral.
1233    /// Requires `with_auth`.
1234    async fn get_collateral_balance(
1235        &self,
1236        asset: String,
1237    ) -> Result<guilder_abstraction::Balance, String> {
1238        if asset.to_uppercase() != "USDC" {
1239            return Err(format!("only USDC collateral is supported, got {}", asset));
1240        }
1241
1242        let total = self.get_collateral().await?;
1243        Ok(guilder_abstraction::Balance {
1244            coin: "USDC".to_string(),
1245            total,
1246            available: total,
1247            locked: Decimal::ZERO,
1248        })
1249    }
1250
1251    /// Returns the user's address-level API rate limit budget.
1252    /// Queries Hyperliquid's `userRateLimit` info endpoint for authoritative server-side counts.
1253    async fn get_user_rate_limit(&self) -> Result<guilder_abstraction::UserRateLimit, String> {
1254        let user = self.require_user_address()?;
1255        let resp = self
1256            .info_post(
1257                serde_json::json!({"type": "userRateLimit", "user": user}),
1258                20,
1259                "get_user_rate_limit",
1260            )
1261            .await?;
1262        let val = parse_response::<Value>(resp).await?;
1263
1264        let cumulative_volume = val["cumVlm"]
1265            .as_str()
1266            .and_then(parse_decimal)
1267            .ok_or_else(|| "missing or invalid cumVlm".to_string())?;
1268        let requests_used = val["nRequestsUsed"]
1269            .as_i64()
1270            .ok_or_else(|| "missing or invalid nRequestsUsed".to_string())?;
1271        let requests_cap = val["nRequestsCap"]
1272            .as_i64()
1273            .ok_or_else(|| "missing or invalid nRequestsCap".to_string())?;
1274        let requests_surplus = val["nRequestsSurplus"]
1275            .as_i64()
1276            .ok_or_else(|| "missing or invalid nRequestsSurplus".to_string())?;
1277
1278        Ok(guilder_abstraction::UserRateLimit {
1279            cumulative_volume,
1280            requests_used,
1281            requests_cap,
1282            requests_surplus,
1283        })
1284    }
1285}
1286
1287#[allow(async_fn_in_trait)]
1288impl guilder_abstraction::SubscribeUserEvents for HyperliquidClient {
1289    fn subscribe_user_fills(&self) -> BoxStream<Result<UserFill, String>> {
1290        let Some(addr) = self.user_address.as_ref() else {
1291            return Box::pin(stream::empty());
1292        };
1293        let addr_str = addr.clone();
1294        let sub = serde_json::json!({
1295            "method": "subscribe",
1296            "subscription": {"type": "userEvents", "user": addr_str.clone()}
1297        });
1298        let key = crate::ws::SubKey {
1299            channel: "userEvents".to_string(),
1300            routing_key: addr_str,
1301        };
1302        let raw_stream = self.ws_mux.subscribe(key, sub);
1303        Box::pin(
1304            raw_stream
1305                .filter_map(|text| async move {
1306                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1307                        return None;
1308                    };
1309                    if env.channel != "userEvents" {
1310                        return None;
1311                    }
1312                    let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1313                        return None;
1314                    };
1315                    let items: Vec<_> = event
1316                        .fills
1317                        .unwrap_or_default()
1318                        .into_iter()
1319                        .filter_map(|fill| {
1320                            let side = if fill.side == "B" {
1321                                OrderSide::Buy
1322                            } else {
1323                                OrderSide::Sell
1324                            };
1325                            let price = parse_decimal(&fill.px)?;
1326                            let quantity = parse_decimal(&fill.sz)?;
1327                            let fee_usd = parse_decimal(&fill.fee)?;
1328                            Some(UserFill {
1329                                order_id: fill.oid,
1330                                symbol: fill.coin,
1331                                side,
1332                                price,
1333                                quantity,
1334                                fee_usd,
1335                                timestamp_ms: fill.time,
1336                                cloid: fill.cloid,
1337                            })
1338                        })
1339                        .collect();
1340                    if items.is_empty() {
1341                        None
1342                    } else {
1343                        Some(stream::iter(items.into_iter().map(Ok)))
1344                    }
1345                })
1346                .flatten(),
1347        )
1348    }
1349
1350    fn subscribe_order_updates(&self) -> BoxStream<Result<OrderUpdate, String>> {
1351        let Some(addr) = self.user_address.as_ref() else {
1352            return Box::pin(stream::empty());
1353        };
1354        let addr_str = addr.clone();
1355        let sub = serde_json::json!({
1356            "method": "subscribe",
1357            "subscription": {"type": "orderUpdates", "user": addr_str.clone()}
1358        });
1359        let key = crate::ws::SubKey {
1360            channel: "orderUpdates".to_string(),
1361            routing_key: addr_str,
1362        };
1363        let raw_stream = self.ws_mux.subscribe(key, sub);
1364        Box::pin(
1365            raw_stream
1366                .filter_map(|text| async move {
1367                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1368                        return None;
1369                    };
1370                    if env.channel != "orderUpdates" {
1371                        return None;
1372                    }
1373                    let Ok(updates) = serde_json::from_value::<Vec<WsOrderUpdate>>(env.data) else {
1374                        return None;
1375                    };
1376                    let items: Vec<_> = updates
1377                        .into_iter()
1378                        .map(|upd| {
1379                            let status = match upd.status.as_str() {
1380                                "open" => OrderStatus::Placed,
1381                                "filled" => OrderStatus::Filled,
1382                                "canceled" | "cancelled" => OrderStatus::Cancelled,
1383                                _ => OrderStatus::PartiallyFilled,
1384                            };
1385                            let side = if upd.order.side == "B" {
1386                                OrderSide::Buy
1387                            } else {
1388                                OrderSide::Sell
1389                            };
1390                            OrderUpdate {
1391                                order_id: upd.order.oid,
1392                                symbol: upd.order.coin,
1393                                status,
1394                                side: Some(side),
1395                                price: parse_decimal(&upd.order.limit_px),
1396                                quantity: parse_decimal(&upd.order.orig_sz),
1397                                remaining_quantity: parse_decimal(&upd.order.sz),
1398                                timestamp_ms: upd.status_timestamp,
1399                                cloid: upd.order.cloid,
1400                            }
1401                        })
1402                        .collect();
1403                    if items.is_empty() {
1404                        None
1405                    } else {
1406                        Some(stream::iter(items.into_iter().map(Ok)))
1407                    }
1408                })
1409                .flatten(),
1410        )
1411    }
1412
1413    fn subscribe_funding_payments(&self) -> BoxStream<Result<FundingPayment, String>> {
1414        let Some(addr) = self.user_address.as_ref() else {
1415            return Box::pin(stream::empty());
1416        };
1417        let addr_str = addr.clone();
1418        let sub = serde_json::json!({
1419            "method": "subscribe",
1420            "subscription": {"type": "userEvents", "user": addr_str.clone()}
1421        });
1422        let key = crate::ws::SubKey {
1423            channel: "userEvents".to_string(),
1424            routing_key: addr_str,
1425        };
1426        let raw_stream = self.ws_mux.subscribe(key, sub);
1427        Box::pin(
1428            raw_stream
1429                .filter_map(|text| async move {
1430                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1431                        return None;
1432                    };
1433                    if env.channel != "userEvents" {
1434                        return None;
1435                    }
1436                    let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1437                        return None;
1438                    };
1439                    let funding = event.funding?;
1440                    let amount_usd = parse_decimal(&funding.usdc)?;
1441                    let item = FundingPayment {
1442                        symbol: funding.coin,
1443                        amount_usd,
1444                        timestamp_ms: funding.time,
1445                    };
1446                    Some(stream::iter(vec![Ok(item)]))
1447                })
1448                .flatten(),
1449        )
1450    }
1451
1452    fn subscribe_deposits(&self) -> BoxStream<Result<Deposit, String>> {
1453        let Some(addr) = self.user_address.as_ref() else {
1454            return Box::pin(stream::empty());
1455        };
1456        let addr_str = addr.clone();
1457        let sub = serde_json::json!({
1458            "method": "subscribe",
1459            "subscription": {"type": "userNonFundingLedgerUpdates", "user": addr_str.clone()}
1460        });
1461        let key = crate::ws::SubKey {
1462            channel: "userNonFundingLedgerUpdates".to_string(),
1463            routing_key: addr_str,
1464        };
1465        let raw_stream = self.ws_mux.subscribe(key, sub);
1466        Box::pin(
1467            raw_stream
1468                .filter_map(|text| async move {
1469                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1470                        return None;
1471                    };
1472                    if env.channel != "userNonFundingLedgerUpdates" {
1473                        return None;
1474                    }
1475                    let Ok(ledger) = serde_json::from_value::<WsLedgerUpdates>(env.data) else {
1476                        return None;
1477                    };
1478                    let items: Vec<_> = ledger
1479                        .updates
1480                        .into_iter()
1481                        .filter_map(|e| {
1482                            if e.delta.kind != "deposit" {
1483                                return None;
1484                            }
1485                            let amount_usd = e.delta.usdc.as_deref().and_then(parse_decimal)?;
1486                            Some(Deposit {
1487                                asset: "USDC".to_string(),
1488                                amount_usd,
1489                                timestamp_ms: e.time,
1490                            })
1491                        })
1492                        .collect();
1493                    if items.is_empty() {
1494                        None
1495                    } else {
1496                        Some(stream::iter(items.into_iter().map(Ok)))
1497                    }
1498                })
1499                .flatten(),
1500        )
1501    }
1502
1503    fn subscribe_withdrawals(&self) -> BoxStream<Result<Withdrawal, String>> {
1504        let Some(addr) = self.user_address.as_ref() else {
1505            return Box::pin(stream::empty());
1506        };
1507        let addr_str = addr.clone();
1508        let sub = serde_json::json!({
1509            "method": "subscribe",
1510            "subscription": {"type": "userNonFundingLedgerUpdates", "user": addr_str.clone()}
1511        });
1512        let key = crate::ws::SubKey {
1513            channel: "userNonFundingLedgerUpdates".to_string(),
1514            routing_key: addr_str,
1515        };
1516        let raw_stream = self.ws_mux.subscribe(key, sub);
1517        Box::pin(
1518            raw_stream
1519                .filter_map(|text| async move {
1520                    let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1521                        return None;
1522                    };
1523                    if env.channel != "userNonFundingLedgerUpdates" {
1524                        return None;
1525                    }
1526                    let Ok(ledger) = serde_json::from_value::<WsLedgerUpdates>(env.data) else {
1527                        return None;
1528                    };
1529                    let items: Vec<_> = ledger
1530                        .updates
1531                        .into_iter()
1532                        .filter_map(|e| {
1533                            if e.delta.kind != "withdraw" {
1534                                return None;
1535                            }
1536                            let amount_usd = e.delta.usdc.as_deref().and_then(parse_decimal)?;
1537                            Some(Withdrawal {
1538                                asset: "USDC".to_string(),
1539                                amount_usd,
1540                                timestamp_ms: e.time,
1541                            })
1542                        })
1543                        .collect();
1544                    if items.is_empty() {
1545                        None
1546                    } else {
1547                        Some(stream::iter(items.into_iter().map(Ok)))
1548                    }
1549                })
1550                .flatten(),
1551        )
1552    }
1553
1554    /// Subscribe to spot wallet balance updates for the registered user address.
1555    /// Requires authentication (address must be set). Returns error if address not registered.
1556    fn subscribe_spot_balance(
1557        &self,
1558    ) -> BoxStream<Result<Vec<guilder_abstraction::Balance>, String>> {
1559        let Some(addr) = self.user_address.as_ref() else {
1560            return Box::pin(stream::iter(vec![Err(
1561                "user address not registered".to_string()
1562            )]));
1563        };
1564        let addr_str = addr.clone();
1565        self.subscribe_spot_balance_with_address(addr_str)
1566    }
1567
1568    /// Subscribe to spot wallet balance updates for a specific address.
1569    fn subscribe_spot_balance_with_address(
1570        &self,
1571        address: String,
1572    ) -> BoxStream<Result<Vec<guilder_abstraction::Balance>, String>> {
1573        let addr_str = address;
1574        let sub = serde_json::json!({
1575            "method": "subscribe",
1576            "subscription": {"type": "userEvents", "user": addr_str.clone()}
1577        });
1578        let key = crate::ws::SubKey {
1579            channel: "userEvents".to_string(),
1580            routing_key: addr_str,
1581        };
1582        let raw_stream = self.ws_mux.subscribe(key, sub);
1583        Box::pin(raw_stream.filter_map(|text| async move {
1584            let Ok(env) = serde_json::from_str::<WsEnvelope>(&text) else {
1585                return None;
1586            };
1587            if env.channel != "userEvents" {
1588                return None;
1589            }
1590            let Ok(event) = serde_json::from_value::<WsUserEvent>(env.data) else {
1591                return None;
1592            };
1593
1594            // Extract spot balances from the event
1595            let spot_state = event.spot_state?;
1596            let balances = spot_state.balances?;
1597
1598            let items: Vec<_> = balances
1599                .into_iter()
1600                .filter_map(|b| {
1601                    let total = parse_decimal(&b.total)?;
1602                    let locked = parse_decimal(&b.hold)?;
1603                    let available = total - locked;
1604                    Some(guilder_abstraction::Balance {
1605                        coin: b.coin,
1606                        total,
1607                        available,
1608                        locked,
1609                    })
1610                })
1611                .collect();
1612
1613            if items.is_empty() {
1614                None
1615            } else {
1616                Some(Ok(items))
1617            }
1618        }))
1619    }
1620}