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//! Parsing helpers 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 = Quantity::from_raw(
343        orig_qty.raw.saturating_sub(remaining_qty.raw),
344        instrument.size_precision(),
345    );
346
347    let price = parse_price(order.order.limit_px, instrument, "order.limitPx")?;
348
349    let ts_accepted = millis_to_nanos(order.order.timestamp)?;
350    let ts_last = millis_to_nanos(order.status_timestamp)?;
351
352    let mut report = OrderStatusReport::new(
353        account_id,
354        instrument_id,
355        None, // venue_order_id_modified
356        venue_order_id,
357        order_side.into(),
358        order_type,
359        time_in_force,
360        order_status,
361        orig_qty, // Use original quantity, not remaining
362        filled_qty,
363        ts_accepted,
364        ts_last,
365        ts_init,
366        Some(UUID4::new()),
367    );
368
369    if let Some(ref cloid) = order.order.cloid {
370        report = report.with_client_order_id(ClientOrderId::new(cloid.as_str()));
371    }
372
373    if matches!(order.order.tif, Some(HyperliquidTimeInForce::Alo)) {
374        report = report.with_post_only(true);
375    }
376
377    if let Some(reduce_only) = order.order.reduce_only {
378        report = report.with_reduce_only(reduce_only);
379    }
380
381    if let Some(reason) = order.status.rejection_reason() {
382        report = report.with_cancel_reason(reason.to_string());
383    }
384
385    report = report.with_price(price);
386
387    if is_conditional && let Some(trigger_px) = order.order.trigger_px {
388        let trigger_price = parse_price(trigger_px, instrument, "order.triggerPx")?;
389        report = report.with_trigger_price(trigger_price);
390    }
391
392    Ok(report)
393}
394
395/// Parses a WebSocket fill message into a [`FillReport`].
396///
397/// This converts Hyperliquid fill data from WebSocket user events into Nautilus fill reports.
398pub fn parse_ws_fill_report(
399    fill: &WsFillData,
400    instrument: &InstrumentAny,
401    account_id: AccountId,
402    ts_init: UnixNanos,
403) -> anyhow::Result<FillReport> {
404    let instrument_id = instrument.id();
405
406    if let Some(liquidation) = fill.liquidation.as_ref() {
407        log::warn!(
408            "Liquidation fill: {} oid={} method={:?} mark_px={} liquidated_user={}",
409            instrument_id,
410            fill.oid,
411            liquidation.method,
412            liquidation.mark_px,
413            liquidation
414                .liquidated_user
415                .as_deref()
416                .unwrap_or("<unknown>"),
417        );
418    } else if matches!(fill.dir, HyperliquidFillDirection::AutoDeleveraging) {
419        log::warn!(
420            "Auto-deleveraging fill: {instrument_id} oid={} px={} sz={}",
421            fill.oid,
422            fill.px,
423            fill.sz,
424        );
425    }
426
427    let venue_order_id = VenueOrderId::new(fill.oid.to_string());
428    let trade_id = make_fill_trade_id(
429        &fill.hash,
430        fill.oid,
431        fill.px,
432        fill.sz,
433        fill.time,
434        fill.start_position,
435    );
436
437    let order_side = OrderSide::from(fill.side);
438    let last_qty = parse_quantity(fill.sz, instrument, "fill.sz")?;
439    let last_px = parse_price(fill.px, instrument, "fill.px")?;
440    let liquidity_side = if fill.crossed {
441        LiquiditySide::Taker
442    } else {
443        LiquiditySide::Maker
444    };
445
446    let fee_amount = fill.fee;
447
448    let commission_currency =
449        crate::http::parse::resolve_fee_currency(fill.fee_token.as_str(), fee_amount, instrument)?;
450
451    let commission = Money::from_decimal(fee_amount, commission_currency)
452        .with_context(|| format!("Failed to create commission from fee='{}'", fill.fee))?;
453    let ts_event = millis_to_nanos(fill.time)?;
454
455    // No client order ID available in fill data directly
456    let client_order_id = None;
457
458    Ok(FillReport::new(
459        account_id,
460        instrument_id,
461        venue_order_id,
462        trade_id,
463        order_side,
464        last_qty,
465        last_px,
466        commission,
467        liquidity_side,
468        client_order_id,
469        None, // venue_position_id
470        ts_event,
471        ts_init,
472        None, // report_id
473    ))
474}
475
476/// Parses a WebSocket ActiveAssetCtx message into mark price, index price, and funding rate updates.
477///
478/// This converts Hyperliquid asset context data into Nautilus price and funding rate updates.
479/// Returns a tuple of (`MarkPriceUpdate`, `Option<IndexPriceUpdate>`, `Option<FundingRateUpdate>`).
480/// Index price and funding rate are only present for perpetual contracts.
481pub fn parse_ws_asset_context(
482    ctx: &WsActiveAssetCtxData,
483    instrument: &InstrumentAny,
484    ts_init: UnixNanos,
485) -> anyhow::Result<(
486    MarkPriceUpdate,
487    Option<IndexPriceUpdate>,
488    Option<FundingRateUpdate>,
489)> {
490    let instrument_id = instrument.id();
491
492    match ctx {
493        WsActiveAssetCtxData::Perp { coin: _, ctx } => {
494            let mark_price = parse_price(ctx.shared.mark_px, instrument, "ctx.mark_px")?;
495            let mark_price_update =
496                MarkPriceUpdate::new(instrument_id, mark_price, ts_init, ts_init);
497
498            let index_price = parse_price(ctx.oracle_px, instrument, "ctx.oracle_px")?;
499            let index_price_update =
500                IndexPriceUpdate::new(instrument_id, index_price, ts_init, ts_init);
501
502            let funding_rate_update = FundingRateUpdate::new(
503                instrument_id,
504                ctx.funding,
505                Some(60), // Hyperliquid exchanges funding hourly
506                None,     // Hyperliquid doesn't provide next funding time in this message
507                ts_init,
508                ts_init,
509            );
510
511            Ok((
512                mark_price_update,
513                Some(index_price_update),
514                Some(funding_rate_update),
515            ))
516        }
517        WsActiveAssetCtxData::Spot { coin: _, ctx } => {
518            let mark_price = parse_price(ctx.shared.mark_px, instrument, "ctx.mark_px")?;
519            let mark_price_update =
520                MarkPriceUpdate::new(instrument_id, mark_price, ts_init, ts_init);
521
522            Ok((mark_price_update, None, None))
523        }
524    }
525}
526
527/// Parses an `activeAssetCtx` open interest string into an open interest custom data update.
528///
529/// The caller is responsible for restricting this to the perpetual branch; spot
530/// `activeAssetCtx` payloads carry no open interest field.
531pub fn parse_ws_open_interest(
532    open_interest: Decimal,
533    instrument: &InstrumentAny,
534    ts_init: UnixNanos,
535) -> anyhow::Result<HyperliquidOpenInterest> {
536    Ok(HyperliquidOpenInterest::new(
537        instrument.id(),
538        open_interest,
539        ts_init,
540        ts_init,
541    ))
542}
543
544/// Converts Hyperliquid TWAP times to nanos.
545///
546/// History row `time` is seconds; `state.timestamp` and fill `time` are milliseconds.
547fn venue_time_to_nanos(value: u64) -> anyhow::Result<UnixNanos> {
548    if value < 100_000_000_000 {
549        Ok(UnixNanos::from(value.checked_mul(1_000_000_000).context(
550            "venue time seconds overflow converting to nanos",
551        )?))
552    } else {
553        millis_to_nanos(value)
554    }
555}
556
557/// Parses one `userTwapHistory` row into custom data.
558///
559/// Unknown coins leave `instrument_id` unset and do not fail the parse.
560pub fn parse_ws_twap_history_row(
561    row: &WsTwapHistoryData,
562    user: &str,
563    is_snapshot: bool,
564    instrument: Option<&InstrumentAny>,
565    ts_init: UnixNanos,
566) -> anyhow::Result<HyperliquidTwapHistory> {
567    let state: &TwapStateData = &row.state;
568    let ts_event = venue_time_to_nanos(row.time)?;
569    let state_timestamp = venue_time_to_nanos(state.timestamp)?;
570    let envelope_user = if user.is_empty() {
571        state.user.as_str()
572    } else {
573        user
574    };
575
576    Ok(HyperliquidTwapHistory::new(
577        envelope_user.to_string(),
578        row.twap_id,
579        state.coin.to_string(),
580        instrument.map(Instrument::id),
581        OrderSide::from(state.side),
582        state.sz,
583        state.executed_sz,
584        state.executed_ntl,
585        state.minutes,
586        state.reduce_only,
587        state.randomize,
588        row.status.status,
589        row.status.description.clone(),
590        state_timestamp,
591        is_snapshot,
592        ts_event,
593        ts_init,
594    ))
595}
596
597/// Parses one `userTwapSliceFills` item into custom data.
598///
599/// Unknown coins leave `instrument_id` unset and do not fail the parse.
600pub fn parse_ws_twap_slice_fill(
601    item: &WsTwapSliceFillData,
602    user: &str,
603    is_snapshot: bool,
604    instrument: Option<&InstrumentAny>,
605    ts_init: UnixNanos,
606) -> anyhow::Result<HyperliquidTwapSliceFill> {
607    let fill = &item.fill;
608    let ts_event = millis_to_nanos(fill.time)?;
609
610    Ok(HyperliquidTwapSliceFill::new(
611        user.to_string(),
612        item.twap_id,
613        fill.coin.to_string(),
614        instrument.map(Instrument::id),
615        fill.px,
616        fill.sz,
617        OrderSide::from(fill.side),
618        fill.hash.clone(),
619        fill.oid,
620        fill.tid,
621        fill.crossed,
622        fill.fee,
623        fill.fee_token.to_string(),
624        fill.dir.to_string(),
625        fill.closed_pnl,
626        is_snapshot,
627        ts_event,
628        ts_init,
629    ))
630}
631
632#[cfg(test)]
633mod tests {
634    use std::str::FromStr;
635
636    use nautilus_model::{
637        identifiers::{InstrumentId, Symbol},
638        instruments::CryptoPerpetual,
639        types::currency::Currency,
640    };
641    use rstest::rstest;
642    use rust_decimal_macros::dec;
643    use ustr::Ustr;
644
645    use super::*;
646    use crate::{
647        common::{
648            consts::HYPERLIQUID_VENUE,
649            enums::{
650                HyperliquidFillDirection, HyperliquidLiquidationMethod,
651                HyperliquidOrderStatus as HyperliquidOrderStatusEnum, HyperliquidSide,
652                HyperliquidTimeInForce,
653            },
654        },
655        websocket::messages::{
656            CandleData, FillLiquidationData, PerpsAssetCtx, SharedAssetCtx, SpotAssetCtx,
657            WsBasicOrderData, WsBookData, WsLevelData,
658        },
659    };
660
661    fn create_test_instrument() -> InstrumentAny {
662        let instrument_id = InstrumentId::new(Symbol::new("BTC-PERP"), *HYPERLIQUID_VENUE);
663
664        InstrumentAny::CryptoPerpetual(
665            CryptoPerpetual::builder()
666                .instrument_id(instrument_id)
667                .raw_symbol(Symbol::new("BTC-PERP"))
668                .base_currency(Currency::from("BTC"))
669                .quote_currency(Currency::from("USDC"))
670                .settlement_currency(Currency::from("USDC"))
671                .is_inverse(false)
672                .price_precision(2)
673                .size_precision(3)
674                .price_increment(Price::from("0.01"))
675                .size_increment(Quantity::from("0.001"))
676                .ts_event(UnixNanos::default())
677                .ts_init(UnixNanos::default())
678                .build()
679                .unwrap(),
680        )
681    }
682
683    #[rstest]
684    fn test_parse_ws_candle_preserves_open_event_and_receipt_initialization_timestamps() {
685        let instrument = create_test_instrument();
686        let bar_type = BarType::from("BTC-PERP.HYPERLIQUID-1-MINUTE-LAST-EXTERNAL");
687        let candle = CandleData {
688            t: 1_700_000_000_000,
689            close_time: 1_700_000_059_999,
690            s: Ustr::from("BTC"),
691            i: Ustr::from("1m"),
692            o: dec!(100.0),
693            c: dec!(100.5),
694            h: dec!(101.0),
695            l: dec!(99.0),
696            v: dec!(10.0),
697            n: 42,
698        };
699        let receipt_timestamp = UnixNanos::from(1_700_000_060_123_000_000);
700
701        let bar = parse_ws_candle(&candle, &instrument, &bar_type, receipt_timestamp).unwrap();
702
703        assert_eq!(bar.ts_event, millis_to_nanos(candle.t).unwrap());
704        assert_eq!(bar.ts_init, receipt_timestamp);
705    }
706
707    #[rstest]
708    fn test_parse_ws_order_status_report_basic() {
709        let instrument = create_test_instrument();
710        let account_id = AccountId::new("HYPERLIQUID-001");
711        let ts_init = UnixNanos::default();
712
713        let order_data = WsOrderData {
714            order: WsBasicOrderData {
715                coin: Ustr::from("BTC"),
716                side: HyperliquidSide::Buy,
717                limit_px: dec!(50000.0),
718                sz: dec!(0.5),
719                oid: 12345,
720                timestamp: 1704470400000,
721                orig_sz: dec!(1.0),
722                cloid: Some("test-order-1".to_string()),
723                tif: Some(HyperliquidTimeInForce::Alo),
724                reduce_only: Some(true),
725                trigger_px: Some(dec!(0.0)),
726                is_market: None,
727                tpsl: None,
728                trigger_activated: None,
729                trailing_stop: None,
730            },
731            status: HyperliquidOrderStatusEnum::Open,
732            status_timestamp: 1704470400000,
733        };
734
735        let result = parse_ws_order_status_report(&order_data, &instrument, account_id, ts_init);
736        assert!(result.is_ok());
737
738        let report = result.unwrap();
739        assert_eq!(report.order_side, OrderSide::Buy.into());
740        assert_eq!(report.order_type, OrderType::Limit);
741        assert_eq!(report.order_status, OrderStatus::Accepted);
742        assert_eq!(report.time_in_force, TimeInForce::Gtc);
743        assert!(report.post_only);
744        assert!(report.reduce_only);
745        assert!(report.trigger_price.is_none());
746    }
747
748    #[rstest]
749    #[case(
750        HyperliquidOrderStatusEnum::BadAloPxRejected,
751        "Post only order would have immediately matched"
752    )]
753    #[case(
754        HyperliquidOrderStatusEnum::ReduceOnlyRejected,
755        "Reduce only order would increase position."
756    )]
757    #[case(
758        HyperliquidOrderStatusEnum::IocCancelRejected,
759        "Order could not immediately match against any resting orders"
760    )]
761    fn test_parse_ws_rejection_preserves_venue_reason(
762        #[case] status: HyperliquidOrderStatusEnum,
763        #[case] expected_reason: &str,
764    ) {
765        let instrument = create_test_instrument();
766        let order_data = WsOrderData {
767            order: WsBasicOrderData {
768                coin: Ustr::from("BTC"),
769                side: HyperliquidSide::Buy,
770                limit_px: dec!(50000.0),
771                sz: dec!(1.0),
772                oid: 12345,
773                timestamp: 1704470400000,
774                orig_sz: dec!(1.0),
775                cloid: Some("test-rejection".to_string()),
776                tif: Some(HyperliquidTimeInForce::Alo),
777                reduce_only: Some(false),
778                trigger_px: None,
779                is_market: None,
780                tpsl: None,
781                trigger_activated: None,
782                trailing_stop: None,
783            },
784            status,
785            status_timestamp: 1704470400000,
786        };
787
788        let report = parse_ws_order_status_report(
789            &order_data,
790            &instrument,
791            AccountId::new("HYPERLIQUID-001"),
792            UnixNanos::default(),
793        )
794        .unwrap();
795
796        assert_eq!(report.order_status, OrderStatus::Rejected);
797        assert_eq!(report.cancel_reason.as_deref(), Some(expected_reason));
798    }
799
800    #[rstest]
801    fn test_parse_ws_fill_report_basic() {
802        let instrument = create_test_instrument();
803        let account_id = AccountId::new("HYPERLIQUID-001");
804        let ts_init = UnixNanos::default();
805
806        let fill_data = WsFillData {
807            coin: Ustr::from("BTC"),
808            px: dec!(50000.0),
809            sz: dec!(0.1),
810            side: HyperliquidSide::Buy,
811            time: 1704470400000,
812            start_position: dec!(0.0),
813            dir: HyperliquidFillDirection::OpenLong,
814            closed_pnl: dec!(0.0),
815            hash: "0xabc123".to_string(),
816            oid: 12345,
817            crossed: true,
818            fee: dec!(0.05),
819            tid: 98765,
820            liquidation: None,
821            fee_token: Ustr::from("USDC"),
822            builder_fee: None,
823            cloid: Some("0xd211f1c27288259290850338d22132a0".to_string()),
824            twap_id: None,
825        };
826
827        let result = parse_ws_fill_report(&fill_data, &instrument, account_id, ts_init);
828        assert!(result.is_ok());
829
830        let report = result.unwrap();
831        assert_eq!(report.order_side, OrderSide::Buy);
832        assert_eq!(report.liquidity_side, LiquiditySide::Taker);
833    }
834
835    #[rstest]
836    fn test_parse_ws_fill_report_with_liquidation() {
837        let instrument = create_test_instrument();
838        let account_id = AccountId::new("HYPERLIQUID-001");
839        let ts_init = UnixNanos::default();
840
841        let fill_data = WsFillData {
842            coin: Ustr::from("BTC"),
843            px: dec!(50000.0),
844            sz: dec!(0.1),
845            side: HyperliquidSide::Sell,
846            time: 1704470400000,
847            start_position: dec!(0.1),
848            dir: HyperliquidFillDirection::CloseLong,
849            closed_pnl: dec!(-25.0),
850            hash: "0xdef456".to_string(),
851            oid: 54321,
852            crossed: true,
853            fee: dec!(0.0),
854            tid: 12345,
855            liquidation: Some(FillLiquidationData {
856                liquidated_user: Some("0xuser".to_string()),
857                mark_px: dec!(50000.0),
858                method: HyperliquidLiquidationMethod::Market,
859            }),
860            fee_token: Ustr::from("USDC"),
861            builder_fee: None,
862            cloid: None,
863            twap_id: None,
864        };
865
866        let report = parse_ws_fill_report(&fill_data, &instrument, account_id, ts_init).unwrap();
867
868        // The fill is still emitted through the standard path; the liquidation
869        // metadata is logged for observability rather than encoded on the report.
870        assert_eq!(report.order_side, OrderSide::Sell);
871        assert_eq!(report.liquidity_side, LiquiditySide::Taker);
872        assert_eq!(report.venue_order_id.to_string(), "54321");
873    }
874
875    #[rstest]
876    fn test_parse_ws_fill_report_outcome_round_trip() {
877        use crate::http::{
878            models::{OutcomeMarket, OutcomeMeta},
879            parse::{create_instrument_from_def, parse_outcome_instruments},
880        };
881
882        let meta = OutcomeMeta {
883            outcomes: vec![OutcomeMarket {
884                outcome: 99,
885                name: "BTC daily".to_string(),
886                description: String::new(),
887                side_specs: vec![],
888            }],
889            questions: vec![],
890        };
891
892        let defs = parse_outcome_instruments(&meta).unwrap();
893        let instrument = create_instrument_from_def(&defs[0], UnixNanos::default()).unwrap();
894        assert_eq!(instrument.id().symbol.as_str(), "99-YES-OUTCOME");
895
896        let fill_data = WsFillData {
897            coin: Ustr::from("#990"),
898            px: dec!(0.4500),
899            sz: dec!(1500.00),
900            side: HyperliquidSide::Buy,
901            time: 1_704_470_400_000,
902            start_position: dec!(0.00),
903            dir: HyperliquidFillDirection::OpenLong,
904            closed_pnl: dec!(0.0),
905            hash: "0xabc789".to_string(),
906            oid: 42_42,
907            crossed: true,
908            fee: dec!(0.0),
909            tid: 7777,
910            liquidation: None,
911            fee_token: Ustr::from("+990"),
912            builder_fee: None,
913            cloid: None,
914            twap_id: None,
915        };
916
917        let report = parse_ws_fill_report(
918            &fill_data,
919            &instrument,
920            AccountId::new("HYPERLIQUID-001"),
921            UnixNanos::default(),
922        )
923        .unwrap();
924
925        // Zero-fee outcome fills fall back to the instrument's quote currency
926        // (USDH) instead of the unregistered side token, keeping downstream
927        // OrderFilled events and persistence on a registered currency.
928        assert_eq!(report.commission.currency.code.as_str(), "USDH");
929        assert!(report.commission.as_decimal().is_zero());
930        assert_eq!(report.order_side, OrderSide::Buy);
931    }
932
933    #[rstest]
934    fn test_parse_ws_order_book_deltas_snapshot_behavior() {
935        let instrument = create_test_instrument();
936        let ts_init = UnixNanos::default();
937
938        let book = WsBookData {
939            coin: Ustr::from("BTC"),
940            levels: [
941                vec![WsLevelData {
942                    px: dec!(50000.0),
943                    sz: dec!(1.0),
944                    n: 1,
945                }],
946                vec![WsLevelData {
947                    px: dec!(50001.0),
948                    sz: dec!(2.0),
949                    n: 1,
950                }],
951            ],
952            time: 1_704_470_400_000,
953        };
954
955        let deltas = parse_ws_order_book_deltas(&book, &instrument, ts_init).unwrap();
956
957        assert_eq!(deltas.deltas.len(), 3); // clear + bid + ask
958        assert_eq!(deltas.deltas[0].action, BookAction::Clear);
959
960        let bid_delta = &deltas.deltas[1];
961        assert_eq!(bid_delta.action, BookAction::Add);
962        assert_eq!(bid_delta.order.side, OrderSide::Buy.into());
963        assert!(bid_delta.order.size.is_positive());
964        assert_eq!(bid_delta.order.order_id, 0);
965
966        let ask_delta = &deltas.deltas[2];
967        assert_eq!(ask_delta.action, BookAction::Add);
968        assert_eq!(ask_delta.order.side, OrderSide::Sell.into());
969        assert!(ask_delta.order.size.is_positive());
970        assert_eq!(ask_delta.order.order_id, 0);
971    }
972
973    #[rstest]
974    fn test_parse_ws_order_book_depth10_pads_sparse_book() {
975        let instrument = create_test_instrument();
976        let ts_init = UnixNanos::default();
977
978        // 3 bids, 2 asks - Depth10 must pad the remaining 7/8 slots with zero orders
979        let book = WsBookData {
980            coin: Ustr::from("BTC"),
981            levels: [
982                vec![
983                    WsLevelData {
984                        px: dec!(100.00),
985                        sz: dec!(1.0),
986                        n: 2,
987                    },
988                    WsLevelData {
989                        px: dec!(99.99),
990                        sz: dec!(2.0),
991                        n: 3,
992                    },
993                    WsLevelData {
994                        px: dec!(99.98),
995                        sz: dec!(3.0),
996                        n: 1,
997                    },
998                ],
999                vec![
1000                    WsLevelData {
1001                        px: dec!(100.01),
1002                        sz: dec!(1.5),
1003                        n: 1,
1004                    },
1005                    WsLevelData {
1006                        px: dec!(100.02),
1007                        sz: dec!(2.5),
1008                        n: 4,
1009                    },
1010                ],
1011            ],
1012            time: 1_704_470_400_000,
1013        };
1014
1015        let depth = parse_ws_order_book_depth10(&book, &instrument, ts_init).unwrap();
1016
1017        assert_eq!(depth.instrument_id, instrument.id());
1018        assert_eq!(depth.bids.len(), 10);
1019        assert_eq!(depth.asks.len(), 10);
1020
1021        assert_eq!(depth.bids[0].price.as_f64(), 100.00);
1022        assert_eq!(depth.bids[0].side, OrderSide::Buy.into());
1023        assert_eq!(depth.bid_counts[0], 2);
1024        assert_eq!(depth.bids[2].price.as_f64(), 99.98);
1025        assert_eq!(depth.bid_counts[2], 1);
1026
1027        // Padded bid slots
1028        for i in 3..10 {
1029            assert_eq!(depth.bids[i].side, OrderSide::Buy.into());
1030            assert!(depth.bids[i].size.is_zero());
1031            assert_eq!(depth.bid_counts[i], 0);
1032        }
1033
1034        assert_eq!(depth.asks[0].price.as_f64(), 100.01);
1035        assert_eq!(depth.asks[0].side, OrderSide::Sell.into());
1036        assert_eq!(depth.ask_counts[0], 1);
1037        assert_eq!(depth.asks[1].price.as_f64(), 100.02);
1038        assert_eq!(depth.ask_counts[1], 4);
1039
1040        for i in 2..10 {
1041            assert_eq!(depth.asks[i].side, OrderSide::Sell.into());
1042            assert!(depth.asks[i].size.is_zero());
1043            assert_eq!(depth.ask_counts[i], 0);
1044        }
1045
1046        // Snapshot flag set
1047        assert_eq!(depth.flags, RecordFlag::F_SNAPSHOT as u8);
1048        assert_eq!(
1049            depth.ts_event,
1050            UnixNanos::from(1_704_470_400_000 * 1_000_000)
1051        );
1052    }
1053
1054    #[rstest]
1055    fn test_parse_ws_order_book_depth10_truncates_beyond_10() {
1056        let instrument = create_test_instrument();
1057        let ts_init = UnixNanos::default();
1058
1059        let mk_levels = |base: f64, n: usize| -> Vec<WsLevelData> {
1060            (0..n)
1061                .map(|i| WsLevelData {
1062                    px: Decimal::from_str(&format!("{:.2}", base - i as f64 * 0.01)).unwrap(),
1063                    sz: dec!(1.0),
1064                    n: 1,
1065                })
1066                .collect()
1067        };
1068
1069        let book = WsBookData {
1070            coin: Ustr::from("BTC"),
1071            levels: [mk_levels(100.00, 15), mk_levels(100.50, 12)],
1072            time: 1_704_470_400_000,
1073        };
1074
1075        let depth = parse_ws_order_book_depth10(&book, &instrument, ts_init).unwrap();
1076
1077        // Only first 10 on each side retained
1078        for i in 0..10 {
1079            assert!(
1080                !depth.bids[i].size.is_zero(),
1081                "bid slot {i} unexpectedly empty"
1082            );
1083            assert!(
1084                !depth.asks[i].size.is_zero(),
1085                "ask slot {i} unexpectedly empty"
1086            );
1087        }
1088    }
1089
1090    #[rstest]
1091    fn test_parse_ws_asset_context_perp() {
1092        let instrument = create_test_instrument();
1093        let ts_init = UnixNanos::default();
1094
1095        let ctx_data = WsActiveAssetCtxData::Perp {
1096            coin: Ustr::from("BTC"),
1097            ctx: PerpsAssetCtx {
1098                shared: SharedAssetCtx {
1099                    day_ntl_vlm: dec!(1000000.0),
1100                    prev_day_px: dec!(49000.0),
1101                    mark_px: dec!(50000.0),
1102                    mid_px: Some(dec!(50001.0)),
1103                    impact_pxs: Some(vec!["50000.0".to_string(), "50002.0".to_string()]),
1104                    day_base_vlm: Some(dec!(100.0)),
1105                },
1106                funding: dec!(0.0001),
1107                open_interest: dec!(100000.0),
1108                oracle_px: dec!(50005.0),
1109                premium: Some(dec!(-0.0001)),
1110            },
1111        };
1112
1113        let result = parse_ws_asset_context(&ctx_data, &instrument, ts_init);
1114        assert!(result.is_ok());
1115
1116        let (mark_price, index_price, funding_rate) = result.unwrap();
1117
1118        assert_eq!(mark_price.instrument_id, instrument.id());
1119        assert_eq!(mark_price.value.as_f64(), 50_000.0);
1120
1121        assert!(index_price.is_some());
1122        let index = index_price.unwrap();
1123        assert_eq!(index.instrument_id, instrument.id());
1124        assert_eq!(index.value.as_f64(), 50_005.0);
1125
1126        assert!(funding_rate.is_some());
1127        let funding = funding_rate.unwrap();
1128        assert_eq!(funding.instrument_id, instrument.id());
1129        assert_eq!(funding.rate.to_string(), "0.0001");
1130        assert_eq!(funding.interval, Some(60));
1131    }
1132
1133    #[rstest]
1134    fn test_parse_ws_asset_context_spot() {
1135        let instrument = create_test_instrument();
1136        let ts_init = UnixNanos::default();
1137
1138        let ctx_data = WsActiveAssetCtxData::Spot {
1139            coin: Ustr::from("BTC"),
1140            ctx: SpotAssetCtx {
1141                shared: SharedAssetCtx {
1142                    day_ntl_vlm: dec!(1000000.0),
1143                    prev_day_px: dec!(49000.0),
1144                    mark_px: dec!(50000.0),
1145                    mid_px: Some(dec!(50001.0)),
1146                    impact_pxs: Some(vec!["50000.0".to_string(), "50002.0".to_string()]),
1147                    day_base_vlm: Some(dec!(100.0)),
1148                },
1149                circulating_supply: dec!(19000000.0),
1150            },
1151        };
1152
1153        let result = parse_ws_asset_context(&ctx_data, &instrument, ts_init);
1154        assert!(result.is_ok());
1155
1156        let (mark_price, index_price, funding_rate) = result.unwrap();
1157
1158        assert_eq!(mark_price.instrument_id, instrument.id());
1159        assert_eq!(mark_price.value.as_f64(), 50_000.0);
1160        assert!(index_price.is_none());
1161        assert!(funding_rate.is_none());
1162    }
1163
1164    /// Pins the direct `Decimal::from_str` path for the funding rate. An f64
1165    /// round-trip (the prior implementation) cannot represent these values
1166    /// exactly, so the parsed Decimal would diverge from the input string.
1167    #[rstest]
1168    #[case::positive_high_precision("0.0001234567890123456")]
1169    #[case::negative_high_precision("-0.0001234567890123456")]
1170    fn test_parse_ws_asset_context_perp_preserves_funding_precision(#[case] funding_str: &str) {
1171        let instrument = create_test_instrument();
1172        let ts_init = UnixNanos::default();
1173
1174        let expected = Decimal::from_str(funding_str).unwrap();
1175
1176        let ctx_data = WsActiveAssetCtxData::Perp {
1177            coin: Ustr::from("BTC"),
1178            ctx: PerpsAssetCtx {
1179                shared: SharedAssetCtx {
1180                    day_ntl_vlm: dec!(1000000.0),
1181                    prev_day_px: dec!(49000.0),
1182                    mark_px: dec!(50000.0),
1183                    mid_px: None,
1184                    impact_pxs: None,
1185                    day_base_vlm: None,
1186                },
1187                funding: Decimal::from_str(funding_str).unwrap(),
1188                open_interest: dec!(100000.0),
1189                oracle_px: dec!(50005.0),
1190                premium: None,
1191            },
1192        };
1193
1194        let (_, _, funding_rate) = parse_ws_asset_context(&ctx_data, &instrument, ts_init).unwrap();
1195
1196        let funding = funding_rate.expect("perp ctx must yield funding rate");
1197        assert_eq!(funding.rate, expected);
1198    }
1199
1200    #[rstest]
1201    fn test_parse_ws_open_interest_perp() {
1202        let instrument = create_test_instrument();
1203        let ts_init = UnixNanos::default();
1204
1205        let open_interest = parse_ws_open_interest(dec!(100000.0), &instrument, ts_init).unwrap();
1206
1207        assert_eq!(open_interest.instrument_id, instrument.id());
1208        assert_eq!(open_interest.open_interest.to_string(), "100000.0");
1209        assert_eq!(open_interest.ts_event, ts_init);
1210        assert_eq!(open_interest.ts_init, ts_init);
1211    }
1212
1213    #[rstest]
1214    #[case::round("100000.0")]
1215    #[case::precise("100000.123456789")]
1216    fn test_parse_ws_open_interest_preserves_precision(#[case] open_interest_str: &str) {
1217        let instrument = create_test_instrument();
1218        let ts_init = UnixNanos::default();
1219
1220        let expected = Decimal::from_str(open_interest_str).unwrap();
1221
1222        let open_interest = parse_ws_open_interest(
1223            Decimal::from_str(open_interest_str).unwrap(),
1224            &instrument,
1225            ts_init,
1226        )
1227        .unwrap();
1228
1229        assert_eq!(open_interest.open_interest, expected);
1230    }
1231
1232    #[rstest]
1233    fn test_parse_ws_twap_history_row_from_live_mainnet_fixture() {
1234        let fixture = include_str!("../../test_data/ws_user_twap_history.json");
1235        let msg: crate::websocket::messages::HyperliquidWsMessage =
1236            serde_json::from_str(fixture).expect("fixture should deserialize");
1237        let crate::websocket::messages::HyperliquidWsMessage::UserTwapHistory { data } = msg else {
1238            panic!("expected UserTwapHistory");
1239        };
1240        let ts_init = UnixNanos::from(99);
1241        let is_snapshot = data.is_snapshot.unwrap_or(false);
1242
1243        let row =
1244            parse_ws_twap_history_row(&data.history[0], &data.user, is_snapshot, None, ts_init)
1245                .unwrap();
1246
1247        assert!(row.is_snapshot);
1248        assert_eq!(row.user, data.user);
1249        assert_eq!(row.coin, "xyz:HOOD");
1250        assert_eq!(row.twap_id, Some(2081397));
1251        assert!(row.instrument_id.is_none());
1252        assert_eq!(row.side, OrderSide::Buy);
1253        assert_eq!(row.size.to_string(), "100.0");
1254        assert_eq!(row.executed_size.to_string(), "100.0");
1255        assert_eq!(row.minutes, 240);
1256        assert!(!row.randomize);
1257        assert!(!row.reduce_only);
1258        assert_eq!(
1259            row.status,
1260            crate::common::enums::HyperliquidTwapStatus::Finished
1261        );
1262        assert!(row.status_description.is_empty());
1263        // Live mainnet history.time is seconds (not milliseconds).
1264        assert_eq!(row.ts_event, UnixNanos::from(1_785_848_057_000_000_000));
1265        assert_eq!(row.ts_init, ts_init);
1266    }
1267
1268    #[rstest]
1269    fn test_parse_ws_twap_slice_fill_from_live_mainnet_fixture() {
1270        let fixture = include_str!("../../test_data/ws_user_twap_slice_fills.json");
1271        let msg: crate::websocket::messages::HyperliquidWsMessage =
1272            serde_json::from_str(fixture).expect("fixture should deserialize");
1273        let crate::websocket::messages::HyperliquidWsMessage::UserTwapSliceFills { data } = msg
1274        else {
1275            panic!("expected UserTwapSliceFills");
1276        };
1277        let instrument = create_test_instrument();
1278        let ts_init = UnixNanos::from(99);
1279        let is_snapshot = data.is_snapshot.unwrap_or(false);
1280
1281        let fill = parse_ws_twap_slice_fill(
1282            &data.twap_slice_fills[0],
1283            &data.user,
1284            is_snapshot,
1285            Some(&instrument),
1286            ts_init,
1287        )
1288        .unwrap();
1289
1290        assert!(fill.is_snapshot);
1291        assert_eq!(fill.twap_id, 2_087_225);
1292        assert_eq!(
1293            fill.hash,
1294            "0x0000000000000000000000000000000000000000000000000000000000000000"
1295        );
1296        assert_eq!(fill.coin, "BTC");
1297        assert_eq!(fill.instrument_id, Some(instrument.id()));
1298        assert_eq!(fill.side, OrderSide::Buy);
1299        assert_eq!(fill.price.to_string(), "64597.0");
1300        assert!(fill.crossed);
1301    }
1302}