Skip to main content

nautilus_hyperliquid/websocket/
messages.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16use ahash::AHashMap;
17use derive_builder::Builder;
18#[cfg(test)]
19use nautilus_core::string::secret::REDACTED;
20use nautilus_core::{
21    serialization::{
22        deserialize_decimal, deserialize_decimal_from_str, deserialize_optional_decimal_from_str,
23        serialize_decimal_as_str,
24    },
25    string::secret::SecretString,
26};
27use nautilus_model::{
28    data::{
29        Bar, Data, FundingRateUpdate, IndexPriceUpdate, MarkPriceUpdate, OrderBookDeltas,
30        OrderBookDepth10, QuoteTick, TradeTick,
31    },
32    reports::{FillReport, OrderStatusReport},
33};
34use rust_decimal::Decimal;
35use serde::{Deserialize, Serialize};
36use ustr::Ustr;
37
38use crate::{
39    common::enums::{
40        HyperliquidBarInterval, HyperliquidFillDirection, HyperliquidLiquidationMethod,
41        HyperliquidOrderStatus as HyperliquidOrderStatusEnum, HyperliquidSide,
42        HyperliquidTimeInForce, HyperliquidTpSl, HyperliquidTwapStatus,
43    },
44    http::models::{HyperliquidExchangeAction, HyperliquidExchangeRequest},
45};
46
47/// Represents an outbound WebSocket message from client to Hyperliquid.
48#[derive(Debug, Clone, Serialize)]
49#[serde(tag = "method")]
50#[serde(rename_all = "lowercase")]
51pub enum HyperliquidWsRequest {
52    /// Subscribe to a data feed.
53    Subscribe {
54        /// Subscription details.
55        subscription: SubscriptionRequest,
56    },
57    /// Unsubscribe from a data feed.
58    Unsubscribe {
59        /// Subscription details to remove.
60        subscription: SubscriptionRequest,
61    },
62    /// Post a request (info or action).
63    Post {
64        /// Request ID for tracking.
65        id: u64,
66        /// Request payload.
67        request: PostRequest,
68    },
69    /// Ping for keepalive.
70    Ping,
71}
72
73/// Represents subscription request types for WebSocket feeds.
74#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
75#[serde(tag = "type")]
76#[serde(rename_all = "camelCase")]
77pub enum SubscriptionRequest {
78    /// All mid prices across markets.
79    AllMids {
80        #[serde(skip_serializing_if = "Option::is_none")]
81        dex: Option<String>,
82    },
83    /// Aggregate asset contexts across all perp dexes.
84    AllDexsAssetCtxs,
85    /// Notifications for a user.
86    Notification { user: String },
87    /// Web data for frontend.
88    WebData2 { user: String },
89    /// Candlestick data.
90    Candle {
91        coin: Ustr,
92        interval: HyperliquidBarInterval,
93    },
94    /// Level 2 order book.
95    L2Book {
96        coin: Ustr,
97        #[serde(skip_serializing_if = "Option::is_none")]
98        #[serde(rename = "nSigFigs")]
99        n_sig_figs: Option<u32>,
100        #[serde(skip_serializing_if = "Option::is_none")]
101        mantissa: Option<u32>,
102    },
103    /// Trade updates.
104    Trades { coin: Ustr },
105    /// Order updates for a user.
106    OrderUpdates { user: String },
107    /// User events (fills, funding, liquidations).
108    UserEvents { user: String },
109    /// User fill history.
110    UserFills {
111        user: String,
112        #[serde(skip_serializing_if = "Option::is_none")]
113        #[serde(rename = "aggregateByTime")]
114        aggregate_by_time: Option<bool>,
115    },
116    /// User funding payments.
117    UserFundings { user: String },
118    /// User ledger updates (non-funding).
119    UserNonFundingLedgerUpdates { user: String },
120    /// Active asset context (for perpetuals).
121    ActiveAssetCtx { coin: Ustr },
122    /// Active spot asset context.
123    ActiveSpotAssetCtx { coin: Ustr },
124    /// Active asset data for user.
125    ActiveAssetData { user: String, coin: String },
126    /// TWAP slice fills.
127    UserTwapSliceFills { user: String },
128    /// TWAP history.
129    UserTwapHistory { user: String },
130    /// Best bid/offer updates.
131    Bbo { coin: Ustr },
132}
133
134/// Post request wrapper for info and action requests.
135#[derive(Debug, Clone, Serialize)]
136#[serde(tag = "type")]
137#[serde(rename_all = "lowercase")]
138pub enum PostRequest {
139    /// Info request (no signature required).
140    Info { payload: serde_json::Value },
141    /// Action request (requires signature).
142    Action {
143        payload: HyperliquidExchangeRequest<HyperliquidExchangeAction>,
144    },
145}
146
147/// Action payload with signature.
148#[derive(Debug, Clone, Serialize)]
149pub struct ActionPayload {
150    pub action: ActionRequest,
151    pub nonce: u64,
152    pub signature: SignatureData,
153    #[serde(skip_serializing_if = "Option::is_none")]
154    #[serde(rename = "vaultAddress")]
155    pub vault_address: Option<String>,
156}
157
158/// Signature data.
159#[derive(Debug, Clone, Serialize)]
160pub struct SignatureData {
161    pub r: SecretString,
162    pub s: SecretString,
163    pub v: SecretString,
164}
165
166/// Action request types.
167#[derive(Debug, Clone, Serialize)]
168#[serde(tag = "type")]
169#[serde(rename_all = "lowercase")]
170pub enum ActionRequest {
171    /// Place orders.
172    Order {
173        orders: Vec<OrderRequest>,
174        grouping: String,
175    },
176    /// Cancel orders.
177    Cancel {
178        cancels: Vec<CancelRequest>,
179        #[serde(rename = "f", skip_serializing_if = "Option::is_none")]
180        fast: Option<bool>,
181    },
182    /// Cancel orders by client order ID.
183    CancelByCloid {
184        cancels: Vec<CancelByCloidRequest>,
185        #[serde(rename = "f", skip_serializing_if = "Option::is_none")]
186        fast: Option<bool>,
187    },
188    /// Modify orders.
189    Modify { modifies: Vec<ModifyRequest> },
190}
191
192impl ActionRequest {
193    /// Create a simple order action with default "na" grouping
194    ///
195    /// # Example
196    /// ```ignore
197    /// let action = ActionRequest::order(vec![order1, order2], "na");
198    /// ```
199    pub fn order(orders: Vec<OrderRequest>, grouping: impl Into<String>) -> Self {
200        Self::Order {
201            orders,
202            grouping: grouping.into(),
203        }
204    }
205
206    /// Create a cancel action for multiple orders
207    ///
208    /// # Example
209    /// ```ignore
210    /// let action = ActionRequest::cancel(vec![
211    ///     CancelRequest { a: 0, o: 12345 },
212    ///     CancelRequest { a: 1, o: 67890 },
213    /// ]);
214    /// ```
215    pub fn cancel(cancels: Vec<CancelRequest>) -> Self {
216        Self::Cancel {
217            cancels,
218            fast: None,
219        }
220    }
221
222    /// Create a cancel-by-cloid action
223    ///
224    /// # Example
225    /// ```ignore
226    /// let action = ActionRequest::cancel_by_cloid(vec![
227    ///     CancelByCloidRequest { asset: 0, cloid: "order-1".to_string() },
228    /// ]);
229    /// ```
230    pub fn cancel_by_cloid(cancels: Vec<CancelByCloidRequest>) -> Self {
231        Self::CancelByCloid {
232            cancels,
233            fast: None,
234        }
235    }
236
237    /// Create a modify action for multiple orders
238    ///
239    /// # Example
240    /// ```ignore
241    /// let action = ActionRequest::modify(vec![
242    ///     ModifyRequest { oid: 12345, order: new_order },
243    /// ]);
244    /// ```
245    pub fn modify(modifies: Vec<ModifyRequest>) -> Self {
246        Self::Modify { modifies }
247    }
248}
249
250/// Order placement request.
251#[derive(Debug, Clone, Serialize, Builder)]
252pub struct OrderRequest {
253    /// Asset ID.
254    pub a: u32,
255    /// Buy side (true = buy, false = sell).
256    pub b: bool,
257    /// Price.
258    pub p: String,
259    /// Size.
260    pub s: String,
261    /// Reduce only.
262    pub r: bool,
263    /// Order type.
264    pub t: OrderTypeRequest,
265    /// Client order ID (optional).
266    #[serde(skip_serializing_if = "Option::is_none")]
267    pub c: Option<String>,
268}
269
270/// Order type in request format.
271#[derive(Debug, Clone, Serialize)]
272#[serde(tag = "type")]
273#[serde(rename_all = "lowercase")]
274pub enum OrderTypeRequest {
275    Limit {
276        tif: TimeInForceRequest,
277    },
278    Trigger {
279        #[serde(rename = "isMarket")]
280        is_market: bool,
281        #[serde(rename = "triggerPx")]
282        trigger_px: String,
283        tpsl: TpSlRequest,
284    },
285}
286
287/// Time in force in request format.
288#[derive(Debug, Clone, Serialize)]
289#[serde(rename_all = "PascalCase")]
290pub enum TimeInForceRequest {
291    Alo,
292    Ioc,
293    Gtc,
294}
295
296/// TP/SL in request format.
297#[derive(Debug, Clone, Serialize)]
298#[serde(rename_all = "lowercase")]
299pub enum TpSlRequest {
300    Tp,
301    Sl,
302}
303
304/// Cancel order request.
305#[derive(Debug, Clone, Serialize)]
306pub struct CancelRequest {
307    /// Asset ID.
308    pub a: u32,
309    /// Order ID.
310    pub o: u64,
311}
312
313/// Cancel by client order ID request.
314#[derive(Debug, Clone, Serialize)]
315pub struct CancelByCloidRequest {
316    /// Asset ID.
317    pub asset: u32,
318    /// Client order ID.
319    pub cloid: String,
320}
321
322/// Modify order request.
323#[derive(Debug, Clone, Serialize)]
324pub struct ModifyRequest {
325    /// Order ID.
326    pub oid: u64,
327    /// New order details.
328    pub order: OrderRequest,
329}
330
331/// Subscription response data wrapper.
332#[derive(Debug, Clone, Deserialize)]
333pub struct SubscriptionResponseData {
334    pub method: String,
335    pub subscription: SubscriptionRequest,
336}
337
338/// Inbound WebSocket message from Hyperliquid server.
339#[derive(Debug, Clone, Deserialize)]
340#[serde(tag = "channel")]
341#[serde(rename_all = "camelCase")]
342pub enum HyperliquidWsMessage {
343    /// Subscription confirmation.
344    SubscriptionResponse { data: SubscriptionResponseData },
345    /// Post request response.
346    Post { data: PostResponse },
347    /// All mid prices.
348    AllMids { data: AllMidsData },
349    /// Aggregate asset contexts across all perp dexes.
350    AllDexsAssetCtxs { data: WsAllDexsAssetCtxsData },
351    /// Notifications.
352    Notification { data: NotificationData },
353    /// Web data.
354    WebData2 { data: serde_json::Value },
355    /// Candlestick data.
356    Candle { data: CandleData },
357    /// Level 2 order book.
358    L2Book { data: WsBookData },
359    /// Trade updates.
360    Trades { data: Vec<WsTradeData> },
361    /// Order updates.
362    OrderUpdates { data: Vec<WsOrderData> },
363    /// User events.
364    UserEvents { data: WsUserEventData },
365    /// Generic user channel (Hyperliquid sends fills/events on this channel).
366    #[serde(rename = "user")]
367    User { data: WsUserEventData },
368    /// User fills.
369    UserFills { data: WsUserFillsData },
370    /// User funding payments.
371    UserFundings { data: WsUserFundingsData },
372    /// User ledger updates.
373    UserNonFundingLedgerUpdates { data: serde_json::Value },
374    /// Active asset context.
375    ActiveAssetCtx { data: WsActiveAssetCtxData },
376    /// Active spot asset context (same data as ActiveAssetCtx, different channel name).
377    ActiveSpotAssetCtx { data: WsActiveAssetCtxData },
378    /// Active asset data.
379    ActiveAssetData { data: WsActiveAssetData },
380    /// TWAP slice fills.
381    UserTwapSliceFills { data: WsUserTwapSliceFillsData },
382    /// TWAP history.
383    UserTwapHistory { data: WsUserTwapHistoryData },
384    /// Best bid/offer.
385    Bbo { data: WsBboData },
386    /// Error response.
387    Error { data: String },
388    /// Pong response.
389    Pong,
390}
391
392/// Post response data.
393#[derive(Debug, Clone, Deserialize)]
394pub struct PostResponse {
395    pub id: u64,
396    pub response: PostResponsePayload,
397}
398
399/// Post response payload.
400#[derive(Debug, Clone, Deserialize)]
401#[serde(tag = "type")]
402#[serde(rename_all = "lowercase")]
403pub enum PostResponsePayload {
404    Info { payload: serde_json::Value },
405    Action { payload: serde_json::Value },
406    Error { payload: String },
407}
408
409/// All mid prices data.
410#[derive(Debug, Clone, Deserialize)]
411pub struct AllMidsData {
412    pub mids: AHashMap<Ustr, String>,
413}
414
415/// `allDexsAssetCtxs` data payload.
416#[derive(Debug, Clone, Deserialize)]
417pub struct WsAllDexsAssetCtxsData {
418    pub ctxs: Vec<(String, Vec<PerpsAssetCtx>)>,
419}
420
421/// Notification data.
422#[derive(Debug, Clone, Deserialize)]
423pub struct NotificationData {
424    pub notification: String,
425}
426
427/// Candlestick data.
428#[derive(Debug, Clone, Deserialize)]
429pub struct CandleData {
430    /// Open time (millis).
431    pub t: u64,
432    /// Close time (millis).
433    #[serde(rename = "T")]
434    pub close_time: u64,
435    /// Symbol.
436    pub s: Ustr,
437    /// Interval.
438    pub i: Ustr,
439    /// Open price.
440    #[serde(deserialize_with = "deserialize_decimal_from_str")]
441    pub o: Decimal,
442    /// Close price.
443    #[serde(deserialize_with = "deserialize_decimal_from_str")]
444    pub c: Decimal,
445    /// High price.
446    #[serde(deserialize_with = "deserialize_decimal_from_str")]
447    pub h: Decimal,
448    /// Low price.
449    #[serde(deserialize_with = "deserialize_decimal_from_str")]
450    pub l: Decimal,
451    /// Volume.
452    #[serde(deserialize_with = "deserialize_decimal_from_str")]
453    pub v: Decimal,
454    /// Number of trades.
455    pub n: u32,
456}
457
458/// WebSocket book data.
459#[derive(Debug, Clone, Serialize, Deserialize)]
460pub struct WsBookData {
461    pub coin: Ustr,
462    pub levels: [Vec<WsLevelData>; 2], // [bids, asks]
463    pub time: u64,
464}
465
466/// WebSocket level data.
467#[derive(Debug, Clone, Serialize, Deserialize)]
468pub struct WsLevelData {
469    /// Price.
470    #[serde(
471        deserialize_with = "deserialize_decimal_from_str",
472        serialize_with = "serialize_decimal_as_str"
473    )]
474    pub px: Decimal,
475    /// Size.
476    #[serde(
477        deserialize_with = "deserialize_decimal_from_str",
478        serialize_with = "serialize_decimal_as_str"
479    )]
480    pub sz: Decimal,
481    /// Number of orders.
482    pub n: u32,
483}
484
485/// WebSocket trade data.
486#[derive(Debug, Clone, Serialize, Deserialize)]
487pub struct WsTradeData {
488    pub coin: Ustr,
489    pub side: HyperliquidSide,
490    #[serde(
491        deserialize_with = "deserialize_decimal_from_str",
492        serialize_with = "serialize_decimal_as_str"
493    )]
494    pub px: Decimal,
495    #[serde(
496        deserialize_with = "deserialize_decimal_from_str",
497        serialize_with = "serialize_decimal_as_str"
498    )]
499    pub sz: Decimal,
500    pub hash: String,
501    pub time: u64,
502    pub tid: u64,
503    pub users: [String; 2], // [buyer, seller]
504}
505
506/// WebSocket order data.
507#[derive(Debug, Clone, Deserialize)]
508pub struct WsOrderData {
509    pub order: WsBasicOrderData,
510    pub status: HyperliquidOrderStatusEnum,
511    #[serde(rename = "statusTimestamp")]
512    pub status_timestamp: u64,
513}
514
515/// Basic order data.
516#[derive(Debug, Clone, Deserialize)]
517pub struct WsBasicOrderData {
518    pub coin: Ustr,
519    pub side: HyperliquidSide,
520    #[serde(rename = "limitPx", deserialize_with = "deserialize_decimal_from_str")]
521    pub limit_px: Decimal,
522    #[serde(deserialize_with = "deserialize_decimal_from_str")]
523    pub sz: Decimal,
524    pub oid: u64,
525    pub timestamp: u64,
526    #[serde(rename = "origSz", deserialize_with = "deserialize_decimal_from_str")]
527    pub orig_sz: Decimal,
528    pub cloid: Option<String>,
529    pub tif: Option<HyperliquidTimeInForce>,
530    #[serde(rename = "reduceOnly")]
531    pub reduce_only: Option<bool>,
532    /// Trigger price for conditional orders (stop/take-profit).
533    #[serde(
534        rename = "triggerPx",
535        default,
536        deserialize_with = "deserialize_optional_decimal_from_str"
537    )]
538    pub trigger_px: Option<Decimal>,
539    /// Whether this is a market or limit trigger order.
540    #[serde(rename = "isMarket")]
541    pub is_market: Option<bool>,
542    /// Take-profit or stop-loss indicator.
543    pub tpsl: Option<HyperliquidTpSl>,
544    /// Whether the trigger has been activated.
545    #[serde(rename = "triggerActivated")]
546    pub trigger_activated: Option<bool>,
547    /// Trailing stop parameters if applicable.
548    #[serde(rename = "trailingStop")]
549    pub trailing_stop: Option<WsTrailingStopData>,
550}
551
552/// Trailing stop offset type.
553#[derive(Debug, Clone, Copy, PartialEq, Eq, Deserialize, Serialize)]
554#[serde(rename_all = "camelCase")]
555pub enum TrailingOffsetType {
556    /// Price offset.
557    Price,
558    /// Percentage offset.
559    Percentage,
560    /// Basis points offset.
561    BasisPoints,
562}
563
564impl TrailingOffsetType {
565    /// Format the offset value with the appropriate unit.
566    pub fn format_offset(&self, offset: &str) -> String {
567        match self {
568            Self::Price => offset.to_string(),
569            Self::Percentage => format!("{offset}%"),
570            Self::BasisPoints => format!("{offset} bps"),
571        }
572    }
573}
574
575/// Trailing stop data from WebSocket.
576#[derive(Debug, Clone, Deserialize)]
577pub struct WsTrailingStopData {
578    /// Trailing offset value.
579    #[serde(deserialize_with = "deserialize_decimal_from_str")]
580    pub offset: Decimal,
581    /// Offset type.
582    #[serde(rename = "offsetType")]
583    pub offset_type: TrailingOffsetType,
584    /// Current callback price (highest/lowest price reached).
585    #[serde(
586        rename = "callbackPrice",
587        default,
588        deserialize_with = "deserialize_optional_decimal_from_str"
589    )]
590    pub callback_price: Option<Decimal>,
591}
592
593/// WebSocket user event data.
594#[derive(Debug, Clone, Deserialize)]
595#[serde(untagged)]
596pub enum WsUserEventData {
597    Fills {
598        fills: Vec<WsFillData>,
599    },
600    Funding {
601        funding: WsUserFundingData,
602    },
603    Liquidation {
604        liquidation: WsLiquidationData,
605    },
606    NonUserCancel {
607        #[serde(rename = "nonUserCancel")]
608        non_user_cancel: Vec<WsNonUserCancelData>,
609    },
610    /// Trigger order activated (moved from pending to active).
611    TriggerActivated {
612        #[serde(rename = "triggerActivated")]
613        trigger_activated: WsTriggerActivatedData,
614    },
615    /// Trigger order executed (trigger price reached, order placed).
616    TriggerTriggered {
617        #[serde(rename = "triggerTriggered")]
618        trigger_triggered: WsTriggerTriggeredData,
619    },
620}
621
622/// WebSocket fill data.
623#[derive(Debug, Clone, Deserialize)]
624pub struct WsFillData {
625    pub coin: Ustr,
626    #[serde(deserialize_with = "deserialize_decimal_from_str")]
627    pub px: Decimal,
628    #[serde(deserialize_with = "deserialize_decimal_from_str")]
629    pub sz: Decimal,
630    pub side: HyperliquidSide,
631    pub time: u64,
632    #[serde(
633        rename = "startPosition",
634        deserialize_with = "deserialize_decimal_from_str"
635    )]
636    pub start_position: Decimal,
637    pub dir: HyperliquidFillDirection,
638    #[serde(
639        rename = "closedPnl",
640        deserialize_with = "deserialize_decimal_from_str"
641    )]
642    pub closed_pnl: Decimal,
643    pub hash: String,
644    pub oid: u64,
645    pub crossed: bool,
646    #[serde(deserialize_with = "deserialize_decimal_from_str")]
647    pub fee: Decimal,
648    pub tid: u64,
649    #[serde(default)]
650    pub liquidation: Option<FillLiquidationData>,
651    #[serde(rename = "feeToken")]
652    pub fee_token: Ustr,
653    #[serde(
654        rename = "builderFee",
655        default,
656        deserialize_with = "deserialize_optional_decimal_from_str"
657    )]
658    pub builder_fee: Option<Decimal>,
659    /// Client order ID (hex string with 0x prefix).
660    pub cloid: Option<String>,
661    /// TWAP order ID if this fill is part of a TWAP order.
662    #[serde(rename = "twapId")]
663    pub twap_id: Option<serde_json::Value>,
664}
665
666/// Fill liquidation data.
667#[derive(Debug, Clone, Deserialize)]
668pub struct FillLiquidationData {
669    #[serde(rename = "liquidatedUser")]
670    pub liquidated_user: Option<String>,
671    #[serde(rename = "markPx", deserialize_with = "deserialize_decimal_from_str")]
672    pub mark_px: Decimal,
673    pub method: HyperliquidLiquidationMethod,
674}
675
676/// WebSocket user funding data.
677#[derive(Debug, Clone, Deserialize)]
678pub struct WsUserFundingData {
679    pub time: u64,
680    pub coin: Ustr,
681    #[serde(deserialize_with = "deserialize_decimal_from_str")]
682    pub usdc: Decimal,
683    #[serde(deserialize_with = "deserialize_decimal_from_str")]
684    pub szi: Decimal,
685    #[serde(
686        rename = "fundingRate",
687        deserialize_with = "deserialize_decimal_from_str"
688    )]
689    pub funding_rate: Decimal,
690}
691
692/// WebSocket liquidation data.
693#[derive(Debug, Clone, Deserialize)]
694pub struct WsLiquidationData {
695    pub lid: u64,
696    pub liquidator: String,
697    pub liquidated_user: String,
698    #[serde(deserialize_with = "deserialize_decimal_from_str")]
699    pub liquidated_ntl_pos: Decimal,
700    #[serde(deserialize_with = "deserialize_decimal_from_str")]
701    pub liquidated_account_value: Decimal,
702}
703
704/// WebSocket non-user cancel data.
705#[derive(Debug, Clone, Deserialize)]
706pub struct WsNonUserCancelData {
707    pub coin: Ustr,
708    pub oid: u64,
709}
710
711/// Trigger order activated event data.
712#[derive(Debug, Clone, Deserialize)]
713pub struct WsTriggerActivatedData {
714    pub coin: Ustr,
715    pub oid: u64,
716    pub time: u64,
717    #[serde(
718        rename = "triggerPx",
719        deserialize_with = "deserialize_decimal_from_str"
720    )]
721    pub trigger_px: Decimal,
722    pub tpsl: HyperliquidTpSl,
723}
724
725/// Trigger order triggered event data.
726#[derive(Debug, Clone, Deserialize)]
727pub struct WsTriggerTriggeredData {
728    pub coin: Ustr,
729    pub oid: u64,
730    pub time: u64,
731    #[serde(
732        rename = "triggerPx",
733        deserialize_with = "deserialize_decimal_from_str"
734    )]
735    pub trigger_px: Decimal,
736    #[serde(rename = "marketPx", deserialize_with = "deserialize_decimal_from_str")]
737    pub market_px: Decimal,
738    pub tpsl: HyperliquidTpSl,
739    /// Order ID of the resulting market/limit order after trigger.
740    #[serde(rename = "resultingOid")]
741    pub resulting_oid: Option<u64>,
742}
743
744/// WebSocket user fills data.
745#[derive(Debug, Clone, Deserialize)]
746pub struct WsUserFillsData {
747    #[serde(rename = "isSnapshot")]
748    pub is_snapshot: Option<bool>,
749    pub user: String,
750    pub fills: Vec<WsFillData>,
751}
752
753/// WebSocket user fundings data.
754#[derive(Debug, Clone, Deserialize)]
755pub struct WsUserFundingsData {
756    #[serde(rename = "isSnapshot")]
757    pub is_snapshot: Option<bool>,
758    pub user: String,
759    pub fundings: Vec<WsUserFundingData>,
760}
761
762/// WebSocket active asset context data.
763#[derive(Debug, Clone, Deserialize)]
764#[serde(untagged)]
765pub enum WsActiveAssetCtxData {
766    Perp { coin: Ustr, ctx: PerpsAssetCtx },
767    Spot { coin: Ustr, ctx: SpotAssetCtx },
768}
769
770/// Shared asset context fields.
771#[derive(Debug, Clone, Deserialize)]
772pub struct SharedAssetCtx {
773    #[serde(
774        rename = "dayNtlVlm",
775        deserialize_with = "deserialize_decimal_from_str"
776    )]
777    pub day_ntl_vlm: Decimal,
778    #[serde(
779        rename = "prevDayPx",
780        deserialize_with = "deserialize_decimal_from_str"
781    )]
782    pub prev_day_px: Decimal,
783    #[serde(rename = "markPx", deserialize_with = "deserialize_decimal_from_str")]
784    pub mark_px: Decimal,
785    #[serde(
786        rename = "midPx",
787        default,
788        deserialize_with = "deserialize_optional_decimal_from_str"
789    )]
790    pub mid_px: Option<Decimal>,
791    #[serde(rename = "impactPxs")]
792    pub impact_pxs: Option<Vec<String>>,
793    #[serde(
794        rename = "dayBaseVlm",
795        default,
796        deserialize_with = "deserialize_optional_decimal_from_str"
797    )]
798    pub day_base_vlm: Option<Decimal>,
799}
800
801/// Perps asset context.
802#[derive(Debug, Clone, Deserialize)]
803pub struct PerpsAssetCtx {
804    #[serde(flatten)]
805    pub shared: SharedAssetCtx,
806    #[serde(deserialize_with = "deserialize_decimal_from_str")]
807    pub funding: Decimal,
808    #[serde(
809        rename = "openInterest",
810        deserialize_with = "deserialize_decimal_from_str"
811    )]
812    pub open_interest: Decimal,
813    #[serde(rename = "oraclePx", deserialize_with = "deserialize_decimal_from_str")]
814    pub oracle_px: Decimal,
815    #[serde(default, deserialize_with = "deserialize_optional_decimal_from_str")]
816    pub premium: Option<Decimal>,
817}
818
819/// Spot asset context.
820#[derive(Debug, Clone, Deserialize)]
821pub struct SpotAssetCtx {
822    #[serde(flatten)]
823    pub shared: SharedAssetCtx,
824    #[serde(
825        rename = "circulatingSupply",
826        deserialize_with = "deserialize_decimal_from_str"
827    )]
828    pub circulating_supply: Decimal,
829}
830
831/// WebSocket active asset data.
832#[derive(Debug, Clone, Deserialize)]
833pub struct WsActiveAssetData {
834    pub user: String,
835    pub coin: Ustr,
836    pub leverage: LeverageData,
837    #[serde(rename = "maxTradeSzs")]
838    pub max_trade_szs: [f64; 2],
839    #[serde(rename = "availableToTrade")]
840    pub available_to_trade: [f64; 2],
841}
842
843/// Leverage data.
844#[derive(Debug, Clone, Deserialize)]
845pub struct LeverageData {
846    pub value: f64,
847    pub type_: String,
848}
849
850/// WebSocket TWAP slice fills data.
851#[derive(Debug, Clone, Deserialize)]
852pub struct WsUserTwapSliceFillsData {
853    #[serde(rename = "isSnapshot")]
854    pub is_snapshot: Option<bool>,
855    pub user: String,
856    #[serde(rename = "twapSliceFills")]
857    pub twap_slice_fills: Vec<WsTwapSliceFillData>,
858}
859
860/// TWAP slice fill data.
861#[derive(Debug, Clone, Deserialize)]
862pub struct WsTwapSliceFillData {
863    pub fill: WsFillData,
864    #[serde(rename = "twapId")]
865    pub twap_id: u64,
866}
867
868/// WebSocket TWAP history data.
869#[derive(Debug, Clone, Deserialize)]
870pub struct WsUserTwapHistoryData {
871    #[serde(rename = "isSnapshot")]
872    pub is_snapshot: Option<bool>,
873    pub user: String,
874    pub history: Vec<WsTwapHistoryData>,
875}
876
877/// TWAP history data.
878#[derive(Debug, Clone, Deserialize)]
879pub struct WsTwapHistoryData {
880    pub state: TwapStateData,
881    pub status: TwapStatusData,
882    pub time: u64,
883    #[serde(default, rename = "twapId")]
884    pub twap_id: Option<u64>,
885}
886
887/// TWAP state data.
888#[derive(Debug, Clone, Deserialize)]
889pub struct TwapStateData {
890    pub coin: Ustr,
891    pub user: String,
892    pub side: HyperliquidSide,
893    /// Venue may send a JSON string or number.
894    #[serde(deserialize_with = "deserialize_decimal")]
895    pub sz: Decimal,
896    #[serde(rename = "executedSz", deserialize_with = "deserialize_decimal")]
897    pub executed_sz: Decimal,
898    #[serde(rename = "executedNtl", deserialize_with = "deserialize_decimal")]
899    pub executed_ntl: Decimal,
900    pub minutes: u32,
901    #[serde(rename = "reduceOnly")]
902    pub reduce_only: bool,
903    pub randomize: bool,
904    pub timestamp: u64,
905}
906
907/// TWAP status data.
908#[derive(Debug, Clone, Deserialize)]
909pub struct TwapStatusData {
910    pub status: HyperliquidTwapStatus,
911    /// Present when `status` is `error`; otherwise often omitted.
912    #[serde(default)]
913    pub description: String,
914}
915
916/// WebSocket BBO data.
917#[derive(Debug, Clone, Deserialize)]
918pub struct WsBboData {
919    pub coin: Ustr,
920    pub time: u64,
921    pub bbo: [Option<WsLevelData>; 2], // [bid, ask]
922}
923
924#[cfg(test)]
925mod tests {
926    use rstest::rstest;
927    use rust_decimal_macros::dec;
928    use serde_json;
929
930    use super::*;
931
932    #[rstest]
933    fn test_signature_data_serialization_and_debug_redaction() {
934        let signature = SignatureData {
935            r: SecretString::from("0xsignature-r"),
936            s: SecretString::from("0xsignature-s"),
937            v: SecretString::from("0x1b"),
938        };
939        let wire = serde_json::to_value(&signature).unwrap();
940        let debug = format!("{signature:?}");
941
942        assert_eq!(wire["r"], "0xsignature-r");
943        assert_eq!(wire["s"], "0xsignature-s");
944        assert_eq!(wire["v"], "0x1b");
945        assert_eq!(debug.matches(REDACTED).count(), 3);
946        assert!(!debug.contains("0xsignature-r"));
947        assert!(!debug.contains("0xsignature-s"));
948        assert!(!debug.contains("0x1b"));
949    }
950
951    #[rstest]
952    fn test_subscription_request_serialization() {
953        let sub = SubscriptionRequest::L2Book {
954            coin: Ustr::from("BTC"),
955            n_sig_figs: Some(5),
956            mantissa: None,
957        };
958
959        let json = serde_json::to_string(&sub).unwrap();
960        assert!(json.contains(r#""type":"l2Book""#));
961        assert!(json.contains(r#""coin":"BTC""#));
962    }
963
964    #[rstest]
965    fn test_hyperliquid_ws_request_serialization() {
966        let req = HyperliquidWsRequest::Subscribe {
967            subscription: SubscriptionRequest::Trades {
968                coin: Ustr::from("ETH"),
969            },
970        };
971
972        let json = serde_json::to_string(&req).unwrap();
973        assert!(json.contains(r#""method":"subscribe""#));
974        assert!(json.contains(r#""type":"trades""#));
975    }
976
977    #[rstest]
978    fn test_order_request_serialization() {
979        let order = OrderRequest {
980            a: 0,    // BTC asset ID
981            b: true, // buy
982            p: "50000.0".to_string(),
983            s: "0.1".to_string(),
984            r: false,
985            t: OrderTypeRequest::Limit {
986                tif: TimeInForceRequest::Gtc,
987            },
988            c: Some("client-123".to_string()),
989        };
990
991        let json = serde_json::to_string(&order).unwrap();
992        assert!(json.contains(r#""a":0"#));
993        assert!(json.contains(r#""b":true"#));
994        assert!(json.contains(r#""p":"50000.0""#));
995    }
996
997    #[rstest]
998    fn test_ws_trade_data_deserialization() {
999        let json = r#"{
1000            "coin": "BTC",
1001            "side": "B",
1002            "px": "50000.0",
1003            "sz": "0.1",
1004            "hash": "0x123",
1005            "time": 1234567890,
1006            "tid": 12345,
1007            "users": ["0xabc", "0xdef"]
1008        }"#;
1009
1010        let trade: WsTradeData = serde_json::from_str(json).unwrap();
1011        assert_eq!(trade.coin, "BTC");
1012        assert_eq!(trade.side, HyperliquidSide::Buy);
1013        assert_eq!(trade.px, dec!(50000.0));
1014    }
1015
1016    #[rstest]
1017    fn test_ws_book_data_deserialization() {
1018        let json = r#"{
1019            "coin": "ETH",
1020            "levels": [
1021                [{"px": "3000.0", "sz": "1.0", "n": 1}],
1022                [{"px": "3001.0", "sz": "2.0", "n": 2}]
1023            ],
1024            "time": 1234567890
1025        }"#;
1026
1027        let book: WsBookData = serde_json::from_str(json).unwrap();
1028        assert_eq!(book.coin, "ETH");
1029        assert_eq!(book.levels[0].len(), 1);
1030        assert_eq!(book.levels[1].len(), 1);
1031    }
1032
1033    #[rstest]
1034    fn test_ws_trailing_stop_data_deserialization() {
1035        let json = r#"{
1036            "offset": "100.0",
1037            "offsetType": "price",
1038            "callbackPrice": "50000.0"
1039        }"#;
1040
1041        let data: WsTrailingStopData = serde_json::from_str(json).unwrap();
1042        assert_eq!(data.offset, dec!(100.0));
1043        assert_eq!(data.offset_type, TrailingOffsetType::Price);
1044        assert_eq!(data.callback_price.unwrap(), dec!(50000.0));
1045    }
1046
1047    #[rstest]
1048    fn test_ws_trigger_activated_data_deserialization() {
1049        let json = r#"{
1050            "coin": "BTC",
1051            "oid": 12345,
1052            "time": 1704470400000,
1053            "triggerPx": "50000.0",
1054            "tpsl": "sl"
1055        }"#;
1056
1057        let data: WsTriggerActivatedData = serde_json::from_str(json).unwrap();
1058        assert_eq!(data.coin, Ustr::from("BTC"));
1059        assert_eq!(data.oid, 12345);
1060        assert_eq!(data.trigger_px, dec!(50000.0));
1061        assert_eq!(data.tpsl, HyperliquidTpSl::Sl);
1062        assert_eq!(data.time, 1704470400000);
1063    }
1064
1065    #[rstest]
1066    fn test_ws_trigger_triggered_data_deserialization() {
1067        let json = r#"{
1068            "coin": "ETH",
1069            "oid": 67890,
1070            "time": 1704470500000,
1071            "triggerPx": "3000.0",
1072            "marketPx": "3001.0",
1073            "tpsl": "tp",
1074            "resultingOid": 99999
1075        }"#;
1076
1077        let data: WsTriggerTriggeredData = serde_json::from_str(json).unwrap();
1078        assert_eq!(data.coin, Ustr::from("ETH"));
1079        assert_eq!(data.oid, 67890);
1080        assert_eq!(data.trigger_px, dec!(3000.0));
1081        assert_eq!(data.market_px, dec!(3001.0));
1082        assert_eq!(data.tpsl, HyperliquidTpSl::Tp);
1083        assert_eq!(data.resulting_oid, Some(99999));
1084    }
1085
1086    #[rstest]
1087    fn test_ws_fill_data_deserialization_with_cloid_and_twap() {
1088        let json = r#"{
1089            "coin": "@107",
1090            "px": "31.737",
1091            "sz": "0.31",
1092            "side": "B",
1093            "time": 1769920606068,
1094            "startPosition": "0.0",
1095            "dir": "Buy",
1096            "closedPnl": "0.0",
1097            "hash": "0xc731e7561e5334a0c8ab043472ce7d01d400ff3bb95653726afa92a8dd570e8b",
1098            "oid": 308086083674,
1099            "crossed": true,
1100            "fee": "0.00021699",
1101            "tid": 812806034449156,
1102            "cloid": "0xd211f1c27288259290850338d22132a0",
1103            "feeToken": "HYPE",
1104            "twapId": null
1105        }"#;
1106
1107        let fill: WsFillData = serde_json::from_str(json).unwrap();
1108        assert_eq!(fill.coin, "@107");
1109        assert_eq!(fill.px, dec!(31.737));
1110        assert_eq!(fill.sz, dec!(0.31));
1111        assert_eq!(fill.side, HyperliquidSide::Buy);
1112        assert_eq!(fill.oid, 308086083674);
1113        assert!(fill.crossed);
1114        assert_eq!(fill.fee, dec!(0.00021699));
1115        assert_eq!(fill.fee_token, "HYPE");
1116        assert_eq!(
1117            fill.cloid,
1118            Some("0xd211f1c27288259290850338d22132a0".to_string())
1119        );
1120        assert!(fill.twap_id.is_none() || fill.twap_id == Some(serde_json::Value::Null));
1121    }
1122
1123    #[rstest]
1124    fn test_ws_user_fills_message_deserialization() {
1125        let json = r#"{"channel":"user","data":{"fills":[{"coin":"@107","px":"31.737","sz":"0.31","side":"B","time":1769920606068,"startPosition":"0.0","dir":"Buy","closedPnl":"0.0","hash":"0xc731e7561e5334a0c8ab043472ce7d01d400ff3bb95653726afa92a8dd570e8b","oid":308086083674,"crossed":true,"fee":"0.00021699","tid":812806034449156,"cloid":"0xd211f1c27288259290850338d22132a0","feeToken":"HYPE","twapId":null}]}}"#;
1126
1127        let msg: HyperliquidWsMessage = serde_json::from_str(json).unwrap();
1128
1129        match msg {
1130            HyperliquidWsMessage::User { data } => match data {
1131                WsUserEventData::Fills { fills } => {
1132                    assert_eq!(fills.len(), 1);
1133                    let fill = &fills[0];
1134                    assert_eq!(fill.coin, "@107");
1135                    assert_eq!(fill.px, dec!(31.737));
1136                    assert_eq!(
1137                        fill.cloid,
1138                        Some("0xd211f1c27288259290850338d22132a0".to_string())
1139                    );
1140                }
1141                _ => panic!("Expected Fills variant"),
1142            },
1143            _ => panic!("Expected User channel message"),
1144        }
1145    }
1146
1147    #[rstest]
1148    fn test_ws_user_fills_message_with_builder_fee() {
1149        // Real message from production that was failing
1150        let json = r#"{"channel":"user","data":{"fills":[{"coin":"BTC","px":"79146.0","sz":"0.001","side":"A","time":1769940855551,"startPosition":"0.00093","dir":"Long > Short","closedPnl":"0.046128","hash":"0x5f8b9c337a197c4061050434769793020e020019151c9b1203544786391d562b","oid":308254271324,"crossed":false,"fee":"0.019785","builderFee":"0.007914","tid":404237815023429,"cloid":"0x50663504b0f4fedea00080176229d94f","feeToken":"USDC","twapId":null}]}}"#;
1151
1152        let msg: HyperliquidWsMessage = serde_json::from_str(json).unwrap();
1153
1154        match msg {
1155            HyperliquidWsMessage::User { data } => match data {
1156                WsUserEventData::Fills { fills } => {
1157                    assert_eq!(fills.len(), 1);
1158                    let fill = &fills[0];
1159                    assert_eq!(fill.coin, "BTC");
1160                    assert_eq!(fill.px, dec!(79146.0));
1161                    assert_eq!(fill.side, HyperliquidSide::Sell);
1162                    assert_eq!(fill.builder_fee, Some(dec!(0.007914)));
1163                    assert_eq!(fill.fee_token, "USDC");
1164                }
1165                _ => panic!("Expected Fills variant"),
1166            },
1167            _ => panic!("Expected User channel message"),
1168        }
1169    }
1170
1171    #[rstest]
1172    fn test_ws_user_fills_message_with_liquidation() {
1173        // Real message from production that failed to parse: the liquidation
1174        // block carries `markPx` as a quoted string like every other decimal.
1175        let json = include_str!("../../test_data/ws_user_fill_liquidation.json");
1176
1177        let msg: HyperliquidWsMessage = serde_json::from_str(json).unwrap();
1178
1179        match msg {
1180            HyperliquidWsMessage::User { data } => match data {
1181                WsUserEventData::Fills { fills } => {
1182                    assert_eq!(fills.len(), 1);
1183                    let fill = &fills[0];
1184                    let liquidation = fill.liquidation.as_ref().expect("expected liquidation");
1185                    assert_eq!(fill.coin, "BTC");
1186                    assert_eq!(fill.side, HyperliquidSide::Sell);
1187                    assert_eq!(liquidation.mark_px, dec!(66607.0));
1188                    assert_eq!(liquidation.method, HyperliquidLiquidationMethod::Market);
1189                    assert_eq!(
1190                        liquidation.liquidated_user.as_deref(),
1191                        Some("0x360878d351f05975e25f1807a27895e1e5e004fb"),
1192                    );
1193                }
1194                _ => panic!("Expected Fills variant"),
1195            },
1196            _ => panic!("Expected User channel message"),
1197        }
1198    }
1199
1200    #[rstest]
1201    fn test_ws_trade_data_round_trips_decimals_as_strings() {
1202        // Deserializing into Decimal then serializing must reproduce the
1203        // string wire form (with scale preserved), not emit a JSON number.
1204        let json = r#"{"coin":"BTC","side":"B","px":"66653.0","sz":"0.001","hash":"0xabc","time":1,"tid":2,"users":["0xa","0xb"]}"#;
1205
1206        let trade: WsTradeData = serde_json::from_str(json).unwrap();
1207        assert_eq!(trade.px, dec!(66653.0));
1208        assert_eq!(trade.sz, dec!(0.001));
1209
1210        let value = serde_json::to_value(&trade).unwrap();
1211        assert_eq!(value["px"], serde_json::Value::from("66653.0"));
1212        assert_eq!(value["sz"], serde_json::Value::from("0.001"));
1213    }
1214}
1215
1216/// Nautilus WebSocket message wrapper for routing to execution engine.
1217///
1218/// Wraps parsed messages from the handler.
1219///
1220/// All parsing happens in the handler layer, with parsed Nautilus domain objects.
1221/// passed through to the Python layer.
1222#[derive(Debug, Clone)]
1223pub enum NautilusWsMessage {
1224    /// Execution reports (order status and fills).
1225    ExecutionReports(Vec<ExecutionReport>),
1226    /// Parsed trade ticks.
1227    Trades(Vec<TradeTick>),
1228    /// Parsed quote tick (from BBO).
1229    Quote(QuoteTick),
1230    /// Parsed order book deltas.
1231    Deltas(OrderBookDeltas),
1232    /// Parsed order book depth-10 snapshot.
1233    Depth10(Box<OrderBookDepth10>),
1234    /// Parsed candle/bar.
1235    Candle(Bar),
1236    /// Mark price update.
1237    MarkPrice(MarkPriceUpdate),
1238    /// Index price update.
1239    IndexPrice(IndexPriceUpdate),
1240    /// Funding rate update.
1241    FundingRate(FundingRateUpdate),
1242    /// Custom data (e.g. allMids).
1243    CustomData(Data),
1244    /// Error occurred.
1245    Error(String),
1246    /// WebSocket reconnected.
1247    Reconnected,
1248}
1249
1250/// Execution report wrapper for order status and fill reports.
1251///
1252/// This enum allows both order status updates and fill reports.
1253/// to be sent through the execution engine.
1254#[derive(Debug, Clone)]
1255#[allow(
1256    clippy::large_enum_variant,
1257    reason = "the variant size gap only crosses the threshold when high-precision widens the raw types"
1258)]
1259pub enum ExecutionReport {
1260    /// Order status report.
1261    Order(OrderStatusReport),
1262    /// Fill report.
1263    Fill(FillReport),
1264}