Skip to main content

nautilus_hyperliquid/websocket/
parse.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
16//! Parsers for Hyperliquid WebSocket payloads.
17
18use anyhow::Context;
19use nautilus_core::{nanos::UnixNanos, uuid::UUID4};
20use nautilus_model::{
21    data::{
22        Bar, BarType, BookOrder, FundingRateUpdate, IndexPriceUpdate, MarkPriceUpdate,
23        OrderBookDelta, OrderBookDeltas, OrderBookDepth10, QuoteTick, TradeTick,
24        depth::DEPTH10_LEN,
25    },
26    enums::{
27        AggressorSide, BookAction, LiquiditySide, OrderSide, OrderStatus, OrderType, RecordFlag,
28        TimeInForce,
29    },
30    identifiers::{AccountId, ClientOrderId, TradeId, VenueOrderId},
31    instruments::{Instrument, InstrumentAny},
32    reports::{FillReport, OrderStatusReport},
33    types::{Money, Price, Quantity},
34};
35use rust_decimal::Decimal;
36
37use super::messages::{
38    CandleData, TwapStateData, WsActiveAssetCtxData, WsBboData, WsBookData, WsFillData,
39    WsOrderData, WsTradeData, WsTwapHistoryData, WsTwapSliceFillData,
40};
41use crate::{
42    common::{
43        converters::hyperliquid_time_in_force_to_nautilus,
44        enums::{HyperliquidFillDirection, HyperliquidTimeInForce},
45        parse::{
46            is_conditional_order_data, make_fill_trade_id, millis_to_nanos,
47            parse_trigger_order_type,
48        },
49    },
50    data_types::{
51        HyperliquidOpenInterest, HyperliquidPublicTrade, HyperliquidTwapHistory,
52        HyperliquidTwapSliceFill,
53    },
54};
55
56fn parse_price(
57    value: Decimal,
58    instrument: &InstrumentAny,
59    field_name: &str,
60) -> anyhow::Result<Price> {
61    Price::from_decimal_dp(value, instrument.price_precision())
62        .with_context(|| format!("Failed to create price from '{value}' for {field_name}"))
63}
64
65fn parse_quantity(
66    value: Decimal,
67    instrument: &InstrumentAny,
68    field_name: &str,
69) -> anyhow::Result<Quantity> {
70    Quantity::from_decimal_dp(value.abs(), instrument.size_precision())
71        .with_context(|| format!("Failed to create quantity from '{value}' for {field_name}"))
72}
73
74/// Parses a WebSocket trade frame into a [`TradeTick`].
75pub fn parse_ws_trade_tick(
76    trade: &WsTradeData,
77    instrument: &InstrumentAny,
78    ts_init: UnixNanos,
79) -> anyhow::Result<TradeTick> {
80    let price = parse_price(trade.px, instrument, "trade.px")?;
81    let size = parse_quantity(trade.sz, instrument, "trade.sz")?;
82    let aggressor = AggressorSide::from(trade.side);
83    let trade_id = TradeId::new_checked(trade.tid.to_string())
84        .context("invalid trade identifier in Hyperliquid trade message")?;
85    let ts_event = millis_to_nanos(trade.time)?;
86
87    TradeTick::new_checked(
88        instrument.id(),
89        price,
90        size,
91        aggressor,
92        trade_id,
93        ts_event,
94        ts_init,
95    )
96    .context("failed to construct TradeTick from Hyperliquid trade message")
97}
98
99/// Parses a WebSocket trade frame into a complete public Hyperliquid trade.
100pub fn parse_ws_public_trade(
101    trade: &WsTradeData,
102    instrument: &InstrumentAny,
103    ts_init: UnixNanos,
104) -> anyhow::Result<HyperliquidPublicTrade> {
105    let price = parse_price(trade.px, instrument, "trade.px")?;
106    let size = parse_quantity(trade.sz, instrument, "trade.sz")?;
107    let ts_event = millis_to_nanos(trade.time)?;
108
109    Ok(HyperliquidPublicTrade::new(
110        instrument.id(),
111        price,
112        size,
113        AggressorSide::from(trade.side),
114        trade.tid.to_string(),
115        trade.users[0].clone(),
116        trade.users[1].clone(),
117        trade.hash.clone(),
118        ts_event,
119        ts_init,
120    ))
121}
122
123/// Parses a WebSocket L2 order book message into [`OrderBookDeltas`].
124pub fn parse_ws_order_book_deltas(
125    book: &WsBookData,
126    instrument: &InstrumentAny,
127    ts_init: UnixNanos,
128) -> anyhow::Result<OrderBookDeltas> {
129    let ts_event = millis_to_nanos(book.time)?;
130    let bids = &book.levels[0];
131    let asks = &book.levels[1];
132    let mut deltas = Vec::with_capacity(1 + bids.len() + asks.len());
133
134    // Treat every book payload as a snapshot: clear existing depth and rebuild it
135    deltas.push(OrderBookDelta::clear(instrument.id(), 0, ts_event, ts_init));
136
137    for level in bids {
138        let price = parse_price(level.px, instrument, "book.bid.px")?;
139        let size = parse_quantity(level.sz, instrument, "book.bid.sz")?;
140
141        if !size.is_positive() {
142            continue;
143        }
144
145        let order = BookOrder::new(OrderSide::Buy, price, size, 0);
146
147        let delta = OrderBookDelta::new(
148            instrument.id(),
149            BookAction::Add,
150            order,
151            RecordFlag::F_LAST as u8,
152            0, // sequence
153            ts_event,
154            ts_init,
155        );
156
157        deltas.push(delta);
158    }
159
160    for level in asks {
161        let price = parse_price(level.px, instrument, "book.ask.px")?;
162        let size = parse_quantity(level.sz, instrument, "book.ask.sz")?;
163
164        if !size.is_positive() {
165            continue;
166        }
167
168        let order = BookOrder::new(OrderSide::Sell, price, size, 0);
169
170        let delta = OrderBookDelta::new(
171            instrument.id(),
172            BookAction::Add,
173            order,
174            RecordFlag::F_LAST as u8,
175            0, // sequence
176            ts_event,
177            ts_init,
178        );
179
180        deltas.push(delta);
181    }
182
183    Ok(OrderBookDeltas::new(instrument.id(), deltas))
184}
185
186/// Parses a WebSocket L2 order book snapshot into [`OrderBookDepth10`].
187///
188/// Hyperliquid's `l2Book` subscription emits snapshots of bid/ask levels.
189/// Fills any missing levels past the venue-provided depth with zero-size
190/// placeholder orders so the fixed-size `[BookOrder; 10]` arrays are
191/// always fully populated.
192pub fn parse_ws_order_book_depth10(
193    book: &WsBookData,
194    instrument: &InstrumentAny,
195    ts_init: UnixNanos,
196) -> anyhow::Result<OrderBookDepth10> {
197    let ts_event = millis_to_nanos(book.time)?;
198    let price_precision = instrument.price_precision();
199    let size_precision = instrument.size_precision();
200
201    let mut bids: [BookOrder; DEPTH10_LEN] = [BookOrder::default(); DEPTH10_LEN];
202    let mut asks: [BookOrder; DEPTH10_LEN] = [BookOrder::default(); DEPTH10_LEN];
203    let mut bid_counts: [u32; DEPTH10_LEN] = [0; DEPTH10_LEN];
204    let mut ask_counts: [u32; DEPTH10_LEN] = [0; DEPTH10_LEN];
205
206    let raw_bids = book.levels.first().map_or(&[][..], |v| v.as_slice());
207    let raw_asks = book.levels.get(1).map_or(&[][..], |v| v.as_slice());
208
209    for (i, level) in raw_bids.iter().take(DEPTH10_LEN).enumerate() {
210        let price = parse_price(level.px, instrument, "book.bid.px")?;
211        let size = parse_quantity(level.sz, instrument, "book.bid.sz")?;
212        bids[i] = BookOrder::new(OrderSide::Buy, price, size, 0);
213        bid_counts[i] = level.n;
214    }
215
216    for bid in bids.iter_mut().skip(raw_bids.len().min(DEPTH10_LEN)) {
217        *bid = BookOrder::new(
218            OrderSide::Buy,
219            Price::zero(price_precision),
220            Quantity::zero(size_precision),
221            0,
222        );
223    }
224
225    for (i, level) in raw_asks.iter().take(DEPTH10_LEN).enumerate() {
226        let price = parse_price(level.px, instrument, "book.ask.px")?;
227        let size = parse_quantity(level.sz, instrument, "book.ask.sz")?;
228        asks[i] = BookOrder::new(OrderSide::Sell, price, size, 0);
229        ask_counts[i] = level.n;
230    }
231
232    for ask in asks.iter_mut().skip(raw_asks.len().min(DEPTH10_LEN)) {
233        *ask = BookOrder::new(
234            OrderSide::Sell,
235            Price::zero(price_precision),
236            Quantity::zero(size_precision),
237            0,
238        );
239    }
240
241    Ok(OrderBookDepth10::new(
242        instrument.id(),
243        bids,
244        asks,
245        bid_counts,
246        ask_counts,
247        RecordFlag::F_SNAPSHOT as u8,
248        0,
249        ts_event,
250        ts_init,
251    ))
252}
253
254/// Parses a WebSocket BBO (best bid/offer) message into a [`QuoteTick`].
255pub fn parse_ws_quote_tick(
256    bbo: &WsBboData,
257    instrument: &InstrumentAny,
258    ts_init: UnixNanos,
259) -> anyhow::Result<QuoteTick> {
260    let bid_level = bbo.bbo[0]
261        .as_ref()
262        .context("BBO message missing bid level")?;
263    let ask_level = bbo.bbo[1]
264        .as_ref()
265        .context("BBO message missing ask level")?;
266
267    let bid_price = parse_price(bid_level.px, instrument, "bbo.bid.px")?;
268    let ask_price = parse_price(ask_level.px, instrument, "bbo.ask.px")?;
269    let bid_size = parse_quantity(bid_level.sz, instrument, "bbo.bid.sz")?;
270    let ask_size = parse_quantity(ask_level.sz, instrument, "bbo.ask.sz")?;
271
272    let ts_event = millis_to_nanos(bbo.time)?;
273
274    QuoteTick::new_checked(
275        instrument.id(),
276        bid_price,
277        ask_price,
278        bid_size,
279        ask_size,
280        ts_event,
281        ts_init,
282    )
283    .context("failed to construct QuoteTick from Hyperliquid BBO message")
284}
285
286/// Parses a WebSocket candle message into a [`Bar`].
287pub fn parse_ws_candle(
288    candle: &CandleData,
289    instrument: &InstrumentAny,
290    bar_type: &BarType,
291    ts_init: UnixNanos,
292) -> anyhow::Result<Bar> {
293    let open = parse_price(candle.o, instrument, "candle.o")?;
294    let high = parse_price(candle.h, instrument, "candle.h")?;
295    let low = parse_price(candle.l, instrument, "candle.l")?;
296    let close = parse_price(candle.c, instrument, "candle.c")?;
297    let volume = parse_quantity(candle.v, instrument, "candle.v")?;
298
299    let ts_event = millis_to_nanos(candle.t)?;
300
301    Ok(Bar::new(
302        *bar_type, open, high, low, close, volume, ts_event, ts_init,
303    ))
304}
305
306/// Parses a WebSocket order update message into an [`OrderStatusReport`].
307///
308/// This converts Hyperliquid order data from WebSocket into Nautilus order status reports.
309/// Handles both regular and conditional orders (stop/limit-if-touched).
310pub fn parse_ws_order_status_report(
311    order: &WsOrderData,
312    instrument: &InstrumentAny,
313    account_id: AccountId,
314    ts_init: UnixNanos,
315) -> anyhow::Result<OrderStatusReport> {
316    let instrument_id = instrument.id();
317    let venue_order_id = VenueOrderId::new(order.order.oid.to_string());
318    let order_side = OrderSide::from(order.order.side);
319
320    // Determine order type based on trigger info
321    let is_conditional =
322        is_conditional_order_data(order.order.trigger_px, order.order.tpsl.as_ref());
323    let order_type = if is_conditional {
324        if let (Some(is_market), Some(tpsl)) = (order.order.is_market, order.order.tpsl.as_ref()) {
325            parse_trigger_order_type(is_market, tpsl)
326        } else {
327            OrderType::Limit // fallback
328        }
329    } else {
330        OrderType::Limit // Regular limit order
331    };
332
333    let time_in_force = order
334        .order
335        .tif
336        .map_or(TimeInForce::Gtc, hyperliquid_time_in_force_to_nautilus);
337    let order_status = OrderStatus::from(order.status);
338
339    // orig_sz is the original order quantity, sz is the remaining quantity
340    let orig_qty = parse_quantity(order.order.orig_sz, instrument, "order.orig_sz")?;
341    let remaining_qty = parse_quantity(order.order.sz, instrument, "order.sz")?;
342    let filled_qty = orig_qty - orig_qty.min(remaining_qty);
343
344    let price = parse_price(order.order.limit_px, instrument, "order.limitPx")?;
345
346    let ts_accepted = millis_to_nanos(order.order.timestamp)?;
347    let ts_last = millis_to_nanos(order.status_timestamp)?;
348
349    let mut report = OrderStatusReport::new(
350        account_id,
351        instrument_id,
352        None, // venue_order_id_modified
353        venue_order_id,
354        order_side.into(),
355        order_type,
356        time_in_force,
357        order_status,
358        orig_qty, // Use original quantity, not remaining
359        filled_qty,
360        ts_accepted,
361        ts_last,
362        ts_init,
363        Some(UUID4::new()),
364    );
365
366    if let Some(ref cloid) = order.order.cloid {
367        report = report.with_client_order_id(ClientOrderId::new(cloid.as_str()));
368    }
369
370    if matches!(order.order.tif, Some(HyperliquidTimeInForce::Alo)) {
371        report = report.with_post_only(true);
372    }
373
374    if let Some(reduce_only) = order.order.reduce_only {
375        report = report.with_reduce_only(reduce_only);
376    }
377
378    if let Some(reason) = order.status.rejection_reason() {
379        report = report.with_cancel_reason(reason.to_string());
380    }
381
382    report = report.with_price(price);
383
384    if is_conditional && let Some(trigger_px) = order.order.trigger_px {
385        let trigger_price = parse_price(trigger_px, instrument, "order.triggerPx")?;
386        report = report.with_trigger_price(trigger_price);
387    }
388
389    Ok(report)
390}
391
392/// Parses a WebSocket fill message into a [`FillReport`].
393///
394/// This converts Hyperliquid fill data from WebSocket user events into Nautilus fill reports.
395pub fn parse_ws_fill_report(
396    fill: &WsFillData,
397    instrument: &InstrumentAny,
398    account_id: AccountId,
399    ts_init: UnixNanos,
400) -> anyhow::Result<FillReport> {
401    let instrument_id = instrument.id();
402
403    if let Some(liquidation) = fill.liquidation.as_ref() {
404        log::warn!(
405            "Liquidation fill: {} oid={} method={:?} mark_px={} liquidated_user={}",
406            instrument_id,
407            fill.oid,
408            liquidation.method,
409            liquidation.mark_px,
410            liquidation
411                .liquidated_user
412                .as_deref()
413                .unwrap_or("<unknown>"),
414        );
415    } else if matches!(fill.dir, HyperliquidFillDirection::AutoDeleveraging) {
416        log::warn!(
417            "Auto-deleveraging fill: {instrument_id} oid={} px={} sz={}",
418            fill.oid,
419            fill.px,
420            fill.sz,
421        );
422    }
423
424    let venue_order_id = VenueOrderId::new(fill.oid.to_string());
425    let trade_id = make_fill_trade_id(
426        &fill.hash,
427        fill.oid,
428        fill.px,
429        fill.sz,
430        fill.time,
431        fill.start_position,
432    );
433
434    let order_side = OrderSide::from(fill.side);
435    let last_qty = parse_quantity(fill.sz, instrument, "fill.sz")?;
436    let last_px = parse_price(fill.px, instrument, "fill.px")?;
437    let liquidity_side = if fill.crossed {
438        LiquiditySide::Taker
439    } else {
440        LiquiditySide::Maker
441    };
442
443    let fee_amount = fill.fee;
444
445    let commission_currency =
446        crate::http::parse::resolve_fee_currency(fill.fee_token.as_str(), fee_amount, instrument)?;
447
448    let commission = Money::from_decimal(fee_amount, commission_currency)
449        .with_context(|| format!("Failed to create commission from fee='{}'", fill.fee))?;
450    let ts_event = millis_to_nanos(fill.time)?;
451
452    // No client order ID available in fill data directly
453    let client_order_id = None;
454
455    Ok(FillReport::new(
456        account_id,
457        instrument_id,
458        venue_order_id,
459        trade_id,
460        order_side,
461        last_qty,
462        last_px,
463        commission,
464        liquidity_side,
465        client_order_id,
466        None, // venue_position_id
467        ts_event,
468        ts_init,
469        None, // report_id
470    ))
471}
472
473/// Parses a WebSocket ActiveAssetCtx message into mark price, index price, and funding rate updates.
474///
475/// This converts Hyperliquid asset context data into Nautilus price and funding rate updates.
476/// Returns a tuple of (`MarkPriceUpdate`, `Option<IndexPriceUpdate>`, `Option<FundingRateUpdate>`).
477/// Index price and funding rate are only present for perpetual contracts.
478pub fn parse_ws_asset_context(
479    ctx: &WsActiveAssetCtxData,
480    instrument: &InstrumentAny,
481    ts_init: UnixNanos,
482) -> anyhow::Result<(
483    MarkPriceUpdate,
484    Option<IndexPriceUpdate>,
485    Option<FundingRateUpdate>,
486)> {
487    let instrument_id = instrument.id();
488
489    match ctx {
490        WsActiveAssetCtxData::Perp { coin: _, ctx } => {
491            let mark_price = parse_price(ctx.shared.mark_px, instrument, "ctx.mark_px")?;
492            let mark_price_update =
493                MarkPriceUpdate::new(instrument_id, mark_price, ts_init, ts_init);
494
495            let index_price = parse_price(ctx.oracle_px, instrument, "ctx.oracle_px")?;
496            let index_price_update =
497                IndexPriceUpdate::new(instrument_id, index_price, ts_init, ts_init);
498
499            let funding_rate_update = FundingRateUpdate::new(
500                instrument_id,
501                ctx.funding,
502                Some(60), // Hyperliquid exchanges funding hourly
503                None,     // Hyperliquid doesn't provide next funding time in this message
504                ts_init,
505                ts_init,
506            );
507
508            Ok((
509                mark_price_update,
510                Some(index_price_update),
511                Some(funding_rate_update),
512            ))
513        }
514        WsActiveAssetCtxData::Spot { coin: _, ctx } => {
515            let mark_price = parse_price(ctx.shared.mark_px, instrument, "ctx.mark_px")?;
516            let mark_price_update =
517                MarkPriceUpdate::new(instrument_id, mark_price, ts_init, ts_init);
518
519            Ok((mark_price_update, None, None))
520        }
521    }
522}
523
524/// Parses an `activeAssetCtx` open interest string into an open interest custom data update.
525///
526/// The caller is responsible for restricting this to the perpetual branch; spot
527/// `activeAssetCtx` payloads carry no open interest field.
528pub fn parse_ws_open_interest(
529    open_interest: Decimal,
530    instrument: &InstrumentAny,
531    ts_init: UnixNanos,
532) -> anyhow::Result<HyperliquidOpenInterest> {
533    Ok(HyperliquidOpenInterest::new(
534        instrument.id(),
535        open_interest,
536        ts_init,
537        ts_init,
538    ))
539}
540
541/// Converts Hyperliquid TWAP times to nanos.
542///
543/// History row `time` is seconds; `state.timestamp` and fill `time` are milliseconds.
544fn venue_time_to_nanos(value: u64) -> anyhow::Result<UnixNanos> {
545    if value < 100_000_000_000 {
546        Ok(UnixNanos::from(value.checked_mul(1_000_000_000).context(
547            "venue time seconds overflow converting to nanos",
548        )?))
549    } else {
550        millis_to_nanos(value)
551    }
552}
553
554/// Parses one `userTwapHistory` row into custom data.
555///
556/// Unknown coins leave `instrument_id` unset and do not fail the parse.
557pub fn parse_ws_twap_history_row(
558    row: &WsTwapHistoryData,
559    user: &str,
560    is_snapshot: bool,
561    instrument: Option<&InstrumentAny>,
562    ts_init: UnixNanos,
563) -> anyhow::Result<HyperliquidTwapHistory> {
564    let state: &TwapStateData = &row.state;
565    let ts_event = venue_time_to_nanos(row.time)?;
566    let state_timestamp = venue_time_to_nanos(state.timestamp)?;
567    let envelope_user = if user.is_empty() {
568        state.user.as_str()
569    } else {
570        user
571    };
572
573    Ok(HyperliquidTwapHistory::new(
574        envelope_user.to_string(),
575        row.twap_id,
576        state.coin.to_string(),
577        instrument.map(Instrument::id),
578        OrderSide::from(state.side),
579        state.sz,
580        state.executed_sz,
581        state.executed_ntl,
582        state.minutes,
583        state.reduce_only,
584        state.randomize,
585        row.status.status,
586        row.status.description.clone(),
587        state_timestamp,
588        is_snapshot,
589        ts_event,
590        ts_init,
591    ))
592}
593
594/// Parses one `userTwapSliceFills` item into custom data.
595///
596/// Unknown coins leave `instrument_id` unset and do not fail the parse.
597pub fn parse_ws_twap_slice_fill(
598    item: &WsTwapSliceFillData,
599    user: &str,
600    is_snapshot: bool,
601    instrument: Option<&InstrumentAny>,
602    ts_init: UnixNanos,
603) -> anyhow::Result<HyperliquidTwapSliceFill> {
604    let fill = &item.fill;
605    let ts_event = millis_to_nanos(fill.time)?;
606
607    Ok(HyperliquidTwapSliceFill::new(
608        user.to_string(),
609        item.twap_id,
610        fill.coin.to_string(),
611        instrument.map(Instrument::id),
612        fill.px,
613        fill.sz,
614        OrderSide::from(fill.side),
615        fill.hash.clone(),
616        fill.oid,
617        fill.tid,
618        fill.crossed,
619        fill.fee,
620        fill.fee_token.to_string(),
621        fill.dir.to_string(),
622        fill.closed_pnl,
623        is_snapshot,
624        ts_event,
625        ts_init,
626    ))
627}
628
629#[cfg(test)]
630mod tests {
631    use std::str::FromStr;
632
633    use nautilus_model::{
634        identifiers::{InstrumentId, Symbol},
635        instruments::CryptoPerpetual,
636        types::currency::Currency,
637    };
638    use rstest::rstest;
639    use rust_decimal_macros::dec;
640    use ustr::Ustr;
641
642    use super::*;
643    use crate::{
644        common::{
645            consts::HYPERLIQUID_VENUE,
646            enums::{
647                HyperliquidFillDirection, HyperliquidLiquidationMethod,
648                HyperliquidOrderStatus as HyperliquidOrderStatusEnum, HyperliquidSide,
649                HyperliquidTimeInForce,
650            },
651        },
652        websocket::messages::{
653            CandleData, FillLiquidationData, PerpsAssetCtx, SharedAssetCtx, SpotAssetCtx,
654            WsBasicOrderData, WsBookData, WsLevelData,
655        },
656    };
657
658    fn create_test_instrument() -> InstrumentAny {
659        let instrument_id = InstrumentId::new(Symbol::new("BTC-PERP"), *HYPERLIQUID_VENUE);
660
661        InstrumentAny::CryptoPerpetual(
662            CryptoPerpetual::builder()
663                .instrument_id(instrument_id)
664                .raw_symbol(Symbol::new("BTC-PERP"))
665                .base_currency(Currency::from("BTC"))
666                .quote_currency(Currency::from("USDC"))
667                .settlement_currency(Currency::from("USDC"))
668                .is_inverse(false)
669                .price_precision(2)
670                .size_precision(3)
671                .price_increment(Price::from("0.01"))
672                .size_increment(Quantity::from("0.001"))
673                .ts_event(UnixNanos::default())
674                .ts_init(UnixNanos::default())
675                .build()
676                .unwrap(),
677        )
678    }
679
680    #[rstest]
681    fn test_parse_ws_candle_preserves_open_event_and_receipt_initialization_timestamps() {
682        let instrument = create_test_instrument();
683        let bar_type = BarType::from("BTC-PERP.HYPERLIQUID-1-MINUTE-LAST-EXTERNAL");
684        let candle = CandleData {
685            t: 1_700_000_000_000,
686            close_time: 1_700_000_059_999,
687            s: Ustr::from("BTC"),
688            i: Ustr::from("1m"),
689            o: dec!(100.0),
690            c: dec!(100.5),
691            h: dec!(101.0),
692            l: dec!(99.0),
693            v: dec!(10.0),
694            n: 42,
695        };
696        let receipt_timestamp = UnixNanos::from(1_700_000_060_123_000_000);
697
698        let bar = parse_ws_candle(&candle, &instrument, &bar_type, receipt_timestamp).unwrap();
699
700        assert_eq!(bar.ts_event, millis_to_nanos(candle.t).unwrap());
701        assert_eq!(bar.ts_init, receipt_timestamp);
702    }
703
704    #[rstest]
705    fn test_parse_ws_order_status_report_basic() {
706        let instrument = create_test_instrument();
707        let account_id = AccountId::new("HYPERLIQUID-001");
708        let ts_init = UnixNanos::default();
709
710        let order_data = WsOrderData {
711            order: WsBasicOrderData {
712                coin: Ustr::from("BTC"),
713                side: HyperliquidSide::Buy,
714                limit_px: dec!(50000.0),
715                sz: dec!(0.5),
716                oid: 12345,
717                timestamp: 1704470400000,
718                orig_sz: dec!(1.0),
719                cloid: Some("test-order-1".to_string()),
720                tif: Some(HyperliquidTimeInForce::Alo),
721                reduce_only: Some(true),
722                trigger_px: Some(dec!(0.0)),
723                is_market: None,
724                tpsl: None,
725                trigger_activated: None,
726                trailing_stop: None,
727            },
728            status: HyperliquidOrderStatusEnum::Open,
729            status_timestamp: 1704470400000,
730        };
731
732        let result = parse_ws_order_status_report(&order_data, &instrument, account_id, ts_init);
733        assert!(result.is_ok());
734
735        let report = result.unwrap();
736        assert_eq!(report.order_side, OrderSide::Buy.into());
737        assert_eq!(report.order_type, OrderType::Limit);
738        assert_eq!(report.order_status, OrderStatus::Accepted);
739        assert_eq!(report.time_in_force, TimeInForce::Gtc);
740        assert!(report.post_only);
741        assert!(report.reduce_only);
742        assert!(report.trigger_price.is_none());
743    }
744
745    #[rstest]
746    #[case(
747        HyperliquidOrderStatusEnum::BadAloPxRejected,
748        "Post only order would have immediately matched"
749    )]
750    #[case(
751        HyperliquidOrderStatusEnum::ReduceOnlyRejected,
752        "Reduce only order would increase position."
753    )]
754    #[case(
755        HyperliquidOrderStatusEnum::IocCancelRejected,
756        "Order could not immediately match against any resting orders"
757    )]
758    fn test_parse_ws_rejection_preserves_venue_reason(
759        #[case] status: HyperliquidOrderStatusEnum,
760        #[case] expected_reason: &str,
761    ) {
762        let instrument = create_test_instrument();
763        let order_data = WsOrderData {
764            order: WsBasicOrderData {
765                coin: Ustr::from("BTC"),
766                side: HyperliquidSide::Buy,
767                limit_px: dec!(50000.0),
768                sz: dec!(1.0),
769                oid: 12345,
770                timestamp: 1704470400000,
771                orig_sz: dec!(1.0),
772                cloid: Some("test-rejection".to_string()),
773                tif: Some(HyperliquidTimeInForce::Alo),
774                reduce_only: Some(false),
775                trigger_px: None,
776                is_market: None,
777                tpsl: None,
778                trigger_activated: None,
779                trailing_stop: None,
780            },
781            status,
782            status_timestamp: 1704470400000,
783        };
784
785        let report = parse_ws_order_status_report(
786            &order_data,
787            &instrument,
788            AccountId::new("HYPERLIQUID-001"),
789            UnixNanos::default(),
790        )
791        .unwrap();
792
793        assert_eq!(report.order_status, OrderStatus::Rejected);
794        assert_eq!(report.cancel_reason.as_deref(), Some(expected_reason));
795    }
796
797    #[rstest]
798    fn test_parse_ws_fill_report_basic() {
799        let instrument = create_test_instrument();
800        let account_id = AccountId::new("HYPERLIQUID-001");
801        let ts_init = UnixNanos::default();
802
803        let fill_data = WsFillData {
804            coin: Ustr::from("BTC"),
805            px: dec!(50000.0),
806            sz: dec!(0.1),
807            side: HyperliquidSide::Buy,
808            time: 1704470400000,
809            start_position: dec!(0.0),
810            dir: HyperliquidFillDirection::OpenLong,
811            closed_pnl: dec!(0.0),
812            hash: "0xabc123".to_string(),
813            oid: 12345,
814            crossed: true,
815            fee: dec!(0.05),
816            tid: 98765,
817            liquidation: None,
818            fee_token: Ustr::from("USDC"),
819            builder_fee: None,
820            cloid: Some("0xd211f1c27288259290850338d22132a0".to_string()),
821            twap_id: None,
822        };
823
824        let result = parse_ws_fill_report(&fill_data, &instrument, account_id, ts_init);
825        assert!(result.is_ok());
826
827        let report = result.unwrap();
828        assert_eq!(report.order_side, OrderSide::Buy);
829        assert_eq!(report.liquidity_side, LiquiditySide::Taker);
830    }
831
832    #[rstest]
833    fn test_parse_ws_fill_report_with_liquidation() {
834        let instrument = create_test_instrument();
835        let account_id = AccountId::new("HYPERLIQUID-001");
836        let ts_init = UnixNanos::default();
837
838        let fill_data = WsFillData {
839            coin: Ustr::from("BTC"),
840            px: dec!(50000.0),
841            sz: dec!(0.1),
842            side: HyperliquidSide::Sell,
843            time: 1704470400000,
844            start_position: dec!(0.1),
845            dir: HyperliquidFillDirection::CloseLong,
846            closed_pnl: dec!(-25.0),
847            hash: "0xdef456".to_string(),
848            oid: 54321,
849            crossed: true,
850            fee: dec!(0.0),
851            tid: 12345,
852            liquidation: Some(FillLiquidationData {
853                liquidated_user: Some("0xuser".to_string()),
854                mark_px: dec!(50000.0),
855                method: HyperliquidLiquidationMethod::Market,
856            }),
857            fee_token: Ustr::from("USDC"),
858            builder_fee: None,
859            cloid: None,
860            twap_id: None,
861        };
862
863        let report = parse_ws_fill_report(&fill_data, &instrument, account_id, ts_init).unwrap();
864
865        // The fill is still emitted through the standard path; the liquidation
866        // metadata is logged for observability rather than encoded on the report.
867        assert_eq!(report.order_side, OrderSide::Sell);
868        assert_eq!(report.liquidity_side, LiquiditySide::Taker);
869        assert_eq!(report.venue_order_id.to_string(), "54321");
870    }
871
872    #[rstest]
873    fn test_parse_ws_fill_report_outcome_round_trip() {
874        use crate::http::{
875            models::{OutcomeMarket, OutcomeMeta},
876            parse::{create_instrument_from_def, parse_outcome_instruments},
877        };
878
879        let meta = OutcomeMeta {
880            outcomes: vec![OutcomeMarket {
881                outcome: 99,
882                name: "BTC daily".to_string(),
883                description: String::new(),
884                side_specs: vec![],
885            }],
886            questions: vec![],
887        };
888
889        let defs = parse_outcome_instruments(&meta).unwrap();
890        let instrument = create_instrument_from_def(&defs[0], UnixNanos::default()).unwrap();
891        assert_eq!(instrument.id().symbol.as_str(), "99-YES-OUTCOME");
892
893        let fill_data = WsFillData {
894            coin: Ustr::from("#990"),
895            px: dec!(0.4500),
896            sz: dec!(1500.00),
897            side: HyperliquidSide::Buy,
898            time: 1_704_470_400_000,
899            start_position: dec!(0.00),
900            dir: HyperliquidFillDirection::OpenLong,
901            closed_pnl: dec!(0.0),
902            hash: "0xabc789".to_string(),
903            oid: 42_42,
904            crossed: true,
905            fee: dec!(0.0),
906            tid: 7777,
907            liquidation: None,
908            fee_token: Ustr::from("+990"),
909            builder_fee: None,
910            cloid: None,
911            twap_id: None,
912        };
913
914        let report = parse_ws_fill_report(
915            &fill_data,
916            &instrument,
917            AccountId::new("HYPERLIQUID-001"),
918            UnixNanos::default(),
919        )
920        .unwrap();
921
922        // Zero-fee outcome fills fall back to the instrument's quote currency
923        // (USDH) instead of the unregistered side token, keeping downstream
924        // OrderFilled events and persistence on a registered currency.
925        assert_eq!(report.commission.currency.code, "USDH");
926        assert!(report.commission.as_decimal().is_zero());
927        assert_eq!(report.order_side, OrderSide::Buy);
928    }
929
930    #[rstest]
931    fn test_parse_ws_order_book_deltas_snapshot_behavior() {
932        let instrument = create_test_instrument();
933        let ts_init = UnixNanos::default();
934
935        let book = WsBookData {
936            coin: Ustr::from("BTC"),
937            levels: [
938                vec![WsLevelData {
939                    px: dec!(50000.0),
940                    sz: dec!(1.0),
941                    n: 1,
942                }],
943                vec![WsLevelData {
944                    px: dec!(50001.0),
945                    sz: dec!(2.0),
946                    n: 1,
947                }],
948            ],
949            time: 1_704_470_400_000,
950        };
951
952        let deltas = parse_ws_order_book_deltas(&book, &instrument, ts_init).unwrap();
953
954        assert_eq!(deltas.deltas.len(), 3); // clear + bid + ask
955        assert_eq!(deltas.deltas[0].action, BookAction::Clear);
956
957        let bid_delta = &deltas.deltas[1];
958        assert_eq!(bid_delta.action, BookAction::Add);
959        assert_eq!(bid_delta.order.side, OrderSide::Buy.into());
960        assert!(bid_delta.order.size.is_positive());
961        assert_eq!(bid_delta.order.order_id, 0);
962
963        let ask_delta = &deltas.deltas[2];
964        assert_eq!(ask_delta.action, BookAction::Add);
965        assert_eq!(ask_delta.order.side, OrderSide::Sell.into());
966        assert!(ask_delta.order.size.is_positive());
967        assert_eq!(ask_delta.order.order_id, 0);
968    }
969
970    #[rstest]
971    fn test_parse_ws_order_book_depth10_pads_sparse_book() {
972        let instrument = create_test_instrument();
973        let ts_init = UnixNanos::default();
974
975        // 3 bids, 2 asks - Depth10 must pad the remaining 7/8 slots with zero orders
976        let book = WsBookData {
977            coin: Ustr::from("BTC"),
978            levels: [
979                vec![
980                    WsLevelData {
981                        px: dec!(100.00),
982                        sz: dec!(1.0),
983                        n: 2,
984                    },
985                    WsLevelData {
986                        px: dec!(99.99),
987                        sz: dec!(2.0),
988                        n: 3,
989                    },
990                    WsLevelData {
991                        px: dec!(99.98),
992                        sz: dec!(3.0),
993                        n: 1,
994                    },
995                ],
996                vec![
997                    WsLevelData {
998                        px: dec!(100.01),
999                        sz: dec!(1.5),
1000                        n: 1,
1001                    },
1002                    WsLevelData {
1003                        px: dec!(100.02),
1004                        sz: dec!(2.5),
1005                        n: 4,
1006                    },
1007                ],
1008            ],
1009            time: 1_704_470_400_000,
1010        };
1011
1012        let depth = parse_ws_order_book_depth10(&book, &instrument, ts_init).unwrap();
1013
1014        assert_eq!(depth.instrument_id, instrument.id());
1015        assert_eq!(depth.bids.len(), 10);
1016        assert_eq!(depth.asks.len(), 10);
1017
1018        assert_eq!(depth.bids[0].price.as_f64(), 100.00);
1019        assert_eq!(depth.bids[0].side, OrderSide::Buy.into());
1020        assert_eq!(depth.bid_counts[0], 2);
1021        assert_eq!(depth.bids[2].price.as_f64(), 99.98);
1022        assert_eq!(depth.bid_counts[2], 1);
1023
1024        // Padded bid slots
1025        for i in 3..10 {
1026            assert_eq!(depth.bids[i].side, OrderSide::Buy.into());
1027            assert!(depth.bids[i].size.is_zero());
1028            assert_eq!(depth.bid_counts[i], 0);
1029        }
1030
1031        assert_eq!(depth.asks[0].price.as_f64(), 100.01);
1032        assert_eq!(depth.asks[0].side, OrderSide::Sell.into());
1033        assert_eq!(depth.ask_counts[0], 1);
1034        assert_eq!(depth.asks[1].price.as_f64(), 100.02);
1035        assert_eq!(depth.ask_counts[1], 4);
1036
1037        for i in 2..10 {
1038            assert_eq!(depth.asks[i].side, OrderSide::Sell.into());
1039            assert!(depth.asks[i].size.is_zero());
1040            assert_eq!(depth.ask_counts[i], 0);
1041        }
1042
1043        // Snapshot flag set
1044        assert_eq!(depth.flags, RecordFlag::F_SNAPSHOT as u8);
1045        assert_eq!(
1046            depth.ts_event,
1047            UnixNanos::from(1_704_470_400_000 * 1_000_000)
1048        );
1049    }
1050
1051    #[rstest]
1052    fn test_parse_ws_order_book_depth10_truncates_beyond_10() {
1053        let instrument = create_test_instrument();
1054        let ts_init = UnixNanos::default();
1055
1056        let mk_levels = |base: f64, n: usize| -> Vec<WsLevelData> {
1057            (0..n)
1058                .map(|i| WsLevelData {
1059                    px: Decimal::from_str(&format!("{:.2}", base - i as f64 * 0.01)).unwrap(),
1060                    sz: dec!(1.0),
1061                    n: 1,
1062                })
1063                .collect()
1064        };
1065
1066        let book = WsBookData {
1067            coin: Ustr::from("BTC"),
1068            levels: [mk_levels(100.00, 15), mk_levels(100.50, 12)],
1069            time: 1_704_470_400_000,
1070        };
1071
1072        let depth = parse_ws_order_book_depth10(&book, &instrument, ts_init).unwrap();
1073
1074        // Only first 10 on each side retained
1075        for i in 0..10 {
1076            assert!(
1077                !depth.bids[i].size.is_zero(),
1078                "bid slot {i} unexpectedly empty"
1079            );
1080            assert!(
1081                !depth.asks[i].size.is_zero(),
1082                "ask slot {i} unexpectedly empty"
1083            );
1084        }
1085    }
1086
1087    #[rstest]
1088    fn test_parse_ws_asset_context_perp() {
1089        let instrument = create_test_instrument();
1090        let ts_init = UnixNanos::default();
1091
1092        let ctx_data = WsActiveAssetCtxData::Perp {
1093            coin: Ustr::from("BTC"),
1094            ctx: PerpsAssetCtx {
1095                shared: SharedAssetCtx {
1096                    day_ntl_vlm: dec!(1000000.0),
1097                    prev_day_px: dec!(49000.0),
1098                    mark_px: dec!(50000.0),
1099                    mid_px: Some(dec!(50001.0)),
1100                    impact_pxs: Some(vec!["50000.0".to_string(), "50002.0".to_string()]),
1101                    day_base_vlm: Some(dec!(100.0)),
1102                },
1103                funding: dec!(0.0001),
1104                open_interest: dec!(100000.0),
1105                oracle_px: dec!(50005.0),
1106                premium: Some(dec!(-0.0001)),
1107            },
1108        };
1109
1110        let result = parse_ws_asset_context(&ctx_data, &instrument, ts_init);
1111        assert!(result.is_ok());
1112
1113        let (mark_price, index_price, funding_rate) = result.unwrap();
1114
1115        assert_eq!(mark_price.instrument_id, instrument.id());
1116        assert_eq!(mark_price.value.as_f64(), 50_000.0);
1117
1118        assert!(index_price.is_some());
1119        let index = index_price.unwrap();
1120        assert_eq!(index.instrument_id, instrument.id());
1121        assert_eq!(index.value.as_f64(), 50_005.0);
1122
1123        assert!(funding_rate.is_some());
1124        let funding = funding_rate.unwrap();
1125        assert_eq!(funding.instrument_id, instrument.id());
1126        assert_eq!(funding.rate.to_string(), "0.0001");
1127        assert_eq!(funding.interval, Some(60));
1128    }
1129
1130    #[rstest]
1131    fn test_parse_ws_asset_context_spot() {
1132        let instrument = create_test_instrument();
1133        let ts_init = UnixNanos::default();
1134
1135        let ctx_data = WsActiveAssetCtxData::Spot {
1136            coin: Ustr::from("BTC"),
1137            ctx: SpotAssetCtx {
1138                shared: SharedAssetCtx {
1139                    day_ntl_vlm: dec!(1000000.0),
1140                    prev_day_px: dec!(49000.0),
1141                    mark_px: dec!(50000.0),
1142                    mid_px: Some(dec!(50001.0)),
1143                    impact_pxs: Some(vec!["50000.0".to_string(), "50002.0".to_string()]),
1144                    day_base_vlm: Some(dec!(100.0)),
1145                },
1146                circulating_supply: dec!(19000000.0),
1147            },
1148        };
1149
1150        let result = parse_ws_asset_context(&ctx_data, &instrument, ts_init);
1151        assert!(result.is_ok());
1152
1153        let (mark_price, index_price, funding_rate) = result.unwrap();
1154
1155        assert_eq!(mark_price.instrument_id, instrument.id());
1156        assert_eq!(mark_price.value.as_f64(), 50_000.0);
1157        assert!(index_price.is_none());
1158        assert!(funding_rate.is_none());
1159    }
1160
1161    /// Pins the direct `Decimal::from_str` path for the funding rate. An f64
1162    /// round-trip (the prior implementation) cannot represent these values
1163    /// exactly, so the parsed Decimal would diverge from the input string.
1164    #[rstest]
1165    #[case::positive_high_precision("0.0001234567890123456")]
1166    #[case::negative_high_precision("-0.0001234567890123456")]
1167    fn test_parse_ws_asset_context_perp_preserves_funding_precision(#[case] funding_str: &str) {
1168        let instrument = create_test_instrument();
1169        let ts_init = UnixNanos::default();
1170
1171        let expected = Decimal::from_str(funding_str).unwrap();
1172
1173        let ctx_data = WsActiveAssetCtxData::Perp {
1174            coin: Ustr::from("BTC"),
1175            ctx: PerpsAssetCtx {
1176                shared: SharedAssetCtx {
1177                    day_ntl_vlm: dec!(1000000.0),
1178                    prev_day_px: dec!(49000.0),
1179                    mark_px: dec!(50000.0),
1180                    mid_px: None,
1181                    impact_pxs: None,
1182                    day_base_vlm: None,
1183                },
1184                funding: Decimal::from_str(funding_str).unwrap(),
1185                open_interest: dec!(100000.0),
1186                oracle_px: dec!(50005.0),
1187                premium: None,
1188            },
1189        };
1190
1191        let (_, _, funding_rate) = parse_ws_asset_context(&ctx_data, &instrument, ts_init).unwrap();
1192
1193        let funding = funding_rate.expect("perp ctx must yield funding rate");
1194        assert_eq!(funding.rate, expected);
1195    }
1196
1197    #[rstest]
1198    fn test_parse_ws_open_interest_perp() {
1199        let instrument = create_test_instrument();
1200        let ts_init = UnixNanos::default();
1201
1202        let open_interest = parse_ws_open_interest(dec!(100000.0), &instrument, ts_init).unwrap();
1203
1204        assert_eq!(open_interest.instrument_id, instrument.id());
1205        assert_eq!(open_interest.open_interest.to_string(), "100000.0");
1206        assert_eq!(open_interest.ts_event, ts_init);
1207        assert_eq!(open_interest.ts_init, ts_init);
1208    }
1209
1210    #[rstest]
1211    #[case::round("100000.0")]
1212    #[case::precise("100000.123456789")]
1213    fn test_parse_ws_open_interest_preserves_precision(#[case] open_interest_str: &str) {
1214        let instrument = create_test_instrument();
1215        let ts_init = UnixNanos::default();
1216
1217        let expected = Decimal::from_str(open_interest_str).unwrap();
1218
1219        let open_interest = parse_ws_open_interest(
1220            Decimal::from_str(open_interest_str).unwrap(),
1221            &instrument,
1222            ts_init,
1223        )
1224        .unwrap();
1225
1226        assert_eq!(open_interest.open_interest, expected);
1227    }
1228
1229    #[rstest]
1230    fn test_parse_ws_twap_history_row_from_live_mainnet_fixture() {
1231        let fixture = include_str!("../../test_data/ws_user_twap_history.json");
1232        let msg: crate::websocket::messages::HyperliquidWsMessage =
1233            serde_json::from_str(fixture).expect("fixture should deserialize");
1234        let crate::websocket::messages::HyperliquidWsMessage::UserTwapHistory { data } = msg else {
1235            panic!("expected UserTwapHistory");
1236        };
1237        let ts_init = UnixNanos::from(99);
1238        let is_snapshot = data.is_snapshot.unwrap_or(false);
1239
1240        let row =
1241            parse_ws_twap_history_row(&data.history[0], &data.user, is_snapshot, None, ts_init)
1242                .unwrap();
1243
1244        assert!(row.is_snapshot);
1245        assert_eq!(row.user, data.user);
1246        assert_eq!(row.coin, "xyz:HOOD");
1247        assert_eq!(row.twap_id, Some(2081397));
1248        assert!(row.instrument_id.is_none());
1249        assert_eq!(row.side, OrderSide::Buy);
1250        assert_eq!(row.size.to_string(), "100.0");
1251        assert_eq!(row.executed_size.to_string(), "100.0");
1252        assert_eq!(row.minutes, 240);
1253        assert!(!row.randomize);
1254        assert!(!row.reduce_only);
1255        assert_eq!(
1256            row.status,
1257            crate::common::enums::HyperliquidTwapStatus::Finished
1258        );
1259        assert!(row.status_description.is_empty());
1260        // Live mainnet history.time is seconds (not milliseconds).
1261        assert_eq!(row.ts_event, UnixNanos::from(1_785_848_057_000_000_000));
1262        assert_eq!(row.ts_init, ts_init);
1263    }
1264
1265    #[rstest]
1266    fn test_parse_ws_twap_slice_fill_from_live_mainnet_fixture() {
1267        let fixture = include_str!("../../test_data/ws_user_twap_slice_fills.json");
1268        let msg: crate::websocket::messages::HyperliquidWsMessage =
1269            serde_json::from_str(fixture).expect("fixture should deserialize");
1270        let crate::websocket::messages::HyperliquidWsMessage::UserTwapSliceFills { data } = msg
1271        else {
1272            panic!("expected UserTwapSliceFills");
1273        };
1274        let instrument = create_test_instrument();
1275        let ts_init = UnixNanos::from(99);
1276        let is_snapshot = data.is_snapshot.unwrap_or(false);
1277
1278        let fill = parse_ws_twap_slice_fill(
1279            &data.twap_slice_fills[0],
1280            &data.user,
1281            is_snapshot,
1282            Some(&instrument),
1283            ts_init,
1284        )
1285        .unwrap();
1286
1287        assert!(fill.is_snapshot);
1288        assert_eq!(fill.twap_id, 2_087_225);
1289        assert_eq!(
1290            fill.hash,
1291            "0x0000000000000000000000000000000000000000000000000000000000000000"
1292        );
1293        assert_eq!(fill.coin, "BTC");
1294        assert_eq!(fill.instrument_id, Some(instrument.id()));
1295        assert_eq!(fill.side, OrderSide::Buy);
1296        assert_eq!(fill.price.to_string(), "64597.0");
1297        assert!(fill.crossed);
1298    }
1299}