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