Skip to main content

nautilus_kraken/websocket/spot_v2/
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//! WebSocket message parsers for converting Kraken streaming data to Nautilus domain models.
17
18use anyhow::Context;
19use jiff::Timestamp;
20use nautilus_core::{UUID4, nanos::UnixNanos};
21use nautilus_model::{
22    data::{Bar, BarSpecification, BarType, BookOrder, OrderBookDelta, QuoteTick, TradeTick},
23    enums::{
24        AggregationSource, AggressorSide, BarAggregation, BookAction, LiquiditySide, OrderSide,
25        OrderStatus, OrderType, PriceType, RecordFlag, TimeInForce, TriggerType,
26    },
27    identifiers::{AccountId, ClientOrderId, InstrumentId, TradeId, VenueOrderId},
28    instruments::{Instrument, any::InstrumentAny},
29    reports::{FillReport, OrderStatusReport},
30    types::{Currency, Money, Price, Quantity},
31};
32use rust_decimal::Decimal;
33
34use super::{
35    enums::{KrakenExecType, KrakenLiquidityInd, KrakenWsOrderStatus},
36    messages::{
37        KrakenSpotWsMessage, KrakenWsBookData, KrakenWsBookLevel, KrakenWsExecutionData,
38        KrakenWsOhlcData, KrakenWsOrderResponse, KrakenWsTickerData, KrakenWsTradeData,
39    },
40};
41use crate::common::enums::{KrakenOrderSide, KrakenOrderType, KrakenTimeInForce};
42
43/// Parses Kraken WebSocket ticker data into a Nautilus quote tick.
44///
45/// # Errors
46///
47/// Returns an error if:
48/// - Bid or ask price/quantity cannot be parsed.
49pub fn parse_quote_tick(
50    ticker: &KrakenWsTickerData,
51    instrument: &InstrumentAny,
52    ts_init: UnixNanos,
53) -> anyhow::Result<QuoteTick> {
54    let instrument_id = instrument.id();
55    let price_precision = instrument.price_precision();
56    let size_precision = instrument.size_precision();
57
58    let bid_price = Price::from_decimal_dp(ticker.bid, price_precision).with_context(|| {
59        format!("Failed to construct bid Price with precision {price_precision}")
60    })?;
61    let bid_size =
62        Quantity::from_decimal_dp(ticker.bid_qty, size_precision).with_context(|| {
63            format!("Failed to construct bid Quantity with precision {size_precision}")
64        })?;
65
66    let ask_price = Price::from_decimal_dp(ticker.ask, price_precision).with_context(|| {
67        format!("Failed to construct ask Price with precision {price_precision}")
68    })?;
69    let ask_size =
70        Quantity::from_decimal_dp(ticker.ask_qty, size_precision).with_context(|| {
71            format!("Failed to construct ask Quantity with precision {size_precision}")
72        })?;
73
74    let ts_event = datetime_to_nanos(ticker.timestamp, "ticker.timestamp")?;
75
76    Ok(QuoteTick::new(
77        instrument_id,
78        bid_price,
79        ask_price,
80        bid_size,
81        ask_size,
82        ts_event,
83        ts_init,
84    ))
85}
86
87/// Parses Kraken WebSocket trade data into a Nautilus trade tick.
88///
89/// # Errors
90///
91/// Returns an error if:
92/// - Price or quantity cannot be parsed.
93/// - Timestamp is invalid.
94pub fn parse_trade_tick(
95    trade: &KrakenWsTradeData,
96    instrument: &InstrumentAny,
97    ts_init: UnixNanos,
98) -> anyhow::Result<TradeTick> {
99    let instrument_id = instrument.id();
100    let price_precision = instrument.price_precision();
101    let size_precision = instrument.size_precision();
102
103    let price = Price::from_decimal_dp(trade.price, price_precision)
104        .with_context(|| format!("Failed to construct Price with precision {price_precision}"))?;
105    let size = Quantity::from_decimal_dp(trade.qty, size_precision)
106        .with_context(|| format!("Failed to construct Quantity with precision {size_precision}"))?;
107
108    let aggressor = match trade.side {
109        KrakenOrderSide::Buy => AggressorSide::Buy,
110        KrakenOrderSide::Sell => AggressorSide::Sell,
111    };
112
113    let trade_id = TradeId::new_checked(trade.trade_id.to_string())?;
114    let ts_event = datetime_to_nanos(trade.timestamp, "trade.timestamp")?;
115
116    TradeTick::new_checked(
117        instrument_id,
118        price,
119        size,
120        aggressor,
121        trade_id,
122        ts_event,
123        ts_init,
124    )
125    .context("Failed to construct TradeTick from Kraken WebSocket trade")
126}
127
128/// Parses Kraken WebSocket book data into Nautilus order book deltas.
129///
130/// Returns a vector of deltas, one for each bid and ask level.
131///
132/// # Errors
133///
134/// Returns an error if:
135/// - Price or quantity cannot be parsed.
136/// - Timestamp is invalid.
137pub fn parse_book_deltas(
138    book: &KrakenWsBookData,
139    instrument: &InstrumentAny,
140    sequence: u64,
141    is_snapshot: bool,
142    ts_init: UnixNanos,
143) -> anyhow::Result<Vec<OrderBookDelta>> {
144    let instrument_id = instrument.id();
145    let price_precision = instrument.price_precision();
146    let size_precision = instrument.size_precision();
147
148    let ts_event = datetime_to_nanos(book.timestamp, "book.timestamp")?;
149
150    let mut current_sequence = sequence;
151    let mut deltas = Vec::new();
152
153    if is_snapshot {
154        deltas.push(OrderBookDelta::clear(
155            instrument_id,
156            current_sequence,
157            ts_event,
158            ts_init,
159        ));
160        current_sequence += 1;
161    }
162
163    if let Some(ref bids) = book.bids {
164        parse_book_side_levels(
165            bids,
166            OrderSide::Buy,
167            is_snapshot,
168            instrument_id,
169            price_precision,
170            size_precision,
171            ts_event,
172            ts_init,
173            &mut current_sequence,
174            &mut deltas,
175        )?;
176    }
177
178    if let Some(ref asks) = book.asks {
179        parse_book_side_levels(
180            asks,
181            OrderSide::Sell,
182            is_snapshot,
183            instrument_id,
184            price_precision,
185            size_precision,
186            ts_event,
187            ts_init,
188            &mut current_sequence,
189            &mut deltas,
190        )?;
191    }
192
193    if let Some(last) = deltas.last_mut() {
194        last.flags |= RecordFlag::F_LAST as u8;
195    }
196
197    Ok(deltas)
198}
199
200#[expect(clippy::too_many_arguments)]
201fn parse_book_side_levels(
202    levels: &[KrakenWsBookLevel],
203    side: OrderSide,
204    is_snapshot: bool,
205    instrument_id: InstrumentId,
206    price_precision: u8,
207    size_precision: u8,
208    ts_event: UnixNanos,
209    ts_init: UnixNanos,
210    current_sequence: &mut u64,
211    deltas: &mut Vec<OrderBookDelta>,
212) -> anyhow::Result<()> {
213    for level in levels {
214        let Some(delta) = parse_book_level(
215            level,
216            side,
217            is_snapshot,
218            instrument_id,
219            price_precision,
220            size_precision,
221            *current_sequence,
222            ts_event,
223            ts_init,
224        )?
225        else {
226            continue;
227        };
228        deltas.push(delta);
229        *current_sequence += 1;
230    }
231
232    Ok(())
233}
234
235#[expect(clippy::too_many_arguments)]
236fn parse_book_level(
237    level: &KrakenWsBookLevel,
238    side: OrderSide,
239    is_snapshot: bool,
240    instrument_id: InstrumentId,
241    price_precision: u8,
242    size_precision: u8,
243    sequence: u64,
244    ts_event: UnixNanos,
245    ts_init: UnixNanos,
246) -> anyhow::Result<Option<OrderBookDelta>> {
247    let price = Price::from_decimal_dp(level.price, price_precision)
248        .with_context(|| format!("Failed to construct Price with precision {price_precision}"))?;
249    let size = Quantity::from_decimal_dp(level.qty, size_precision)
250        .with_context(|| format!("Failed to construct Quantity with precision {size_precision}"))?;
251
252    let action = if is_snapshot {
253        if size.raw == 0 {
254            return Ok(None);
255        }
256        BookAction::Add
257    } else if size.raw == 0 {
258        BookAction::Delete
259    } else {
260        BookAction::Update
261    };
262
263    let order_id = price.raw as u64;
264    let order = BookOrder::new(side, price, size, order_id);
265    let mut flags = RecordFlag::F_MBP as u8;
266    if is_snapshot {
267        flags |= RecordFlag::F_SNAPSHOT as u8;
268    }
269
270    Ok(Some(OrderBookDelta::new(
271        instrument_id,
272        action,
273        order,
274        flags,
275        sequence,
276        ts_event,
277        ts_init,
278    )))
279}
280
281pub(super) fn datetime_to_nanos(value: Timestamp, field: &str) -> anyhow::Result<UnixNanos> {
282    let nanos = u64::try_from(value.as_nanosecond())
283        .with_context(|| format!("Timestamp predates Unix epoch: {field}='{value}'"))?;
284    Ok(UnixNanos::from(nanos))
285}
286
287/// Parses Kraken WebSocket OHLC data into a Nautilus bar.
288///
289/// The bar's `ts_event` is computed as `interval_begin` + `interval` minutes.
290///
291/// # Errors
292///
293/// Returns an error if:
294/// - Price or quantity values cannot be parsed.
295/// - The interval cannot be converted to a valid bar specification.
296pub fn parse_ws_bar(
297    ohlc: &KrakenWsOhlcData,
298    instrument: &InstrumentAny,
299    ts_init: UnixNanos,
300) -> anyhow::Result<Bar> {
301    let instrument_id = instrument.id();
302    let price_precision = instrument.price_precision();
303    let size_precision = instrument.size_precision();
304
305    let open = Price::from_decimal_dp(ohlc.open, price_precision)?;
306    let high = Price::from_decimal_dp(ohlc.high, price_precision)?;
307    let low = Price::from_decimal_dp(ohlc.low, price_precision)?;
308    let close = Price::from_decimal_dp(ohlc.close, price_precision)?;
309    let volume = Quantity::from_decimal_dp(ohlc.volume, size_precision)?;
310
311    let bar_spec = interval_to_bar_spec(ohlc.interval)?;
312    let bar_type = BarType::new(instrument_id, bar_spec, AggregationSource::External);
313
314    // Compute bar close time: interval_begin + interval minutes
315    let interval_secs = i64::from(ohlc.interval) * 60;
316    let close_time = ohlc.interval_begin + jiff::SignedDuration::from_secs(interval_secs);
317    let ts_event = UnixNanos::from(u64::try_from(close_time.as_nanosecond()).unwrap_or(0));
318
319    Bar::new_checked(bar_type, open, high, low, close, volume, ts_event, ts_init)
320}
321
322/// Converts a Kraken OHLC interval (minutes) to a Nautilus bar specification.
323fn interval_to_bar_spec(interval: u32) -> anyhow::Result<BarSpecification> {
324    let (step, aggregation) = match interval {
325        1 => (1, BarAggregation::Minute),
326        5 => (5, BarAggregation::Minute),
327        15 => (15, BarAggregation::Minute),
328        30 => (30, BarAggregation::Minute),
329        60 => (1, BarAggregation::Hour),
330        240 => (4, BarAggregation::Hour),
331        1440 => (1, BarAggregation::Day),
332        10080 => (1, BarAggregation::Week),
333        21600 => (15, BarAggregation::Day), // 21600 min = 360 hours = 15 days
334        _ => anyhow::bail!("Unsupported Kraken OHLC interval: {interval}"),
335    };
336
337    Ok(BarSpecification::new(step, aggregation, PriceType::Last))
338}
339
340/// Parses Kraken execution type and order status to Nautilus order status.
341fn parse_order_status(
342    exec_type: KrakenExecType,
343    order_status: Option<KrakenWsOrderStatus>,
344) -> OrderStatus {
345    match exec_type {
346        KrakenExecType::Canceled => return OrderStatus::Canceled,
347        KrakenExecType::Expired => return OrderStatus::Expired,
348        KrakenExecType::Filled => return OrderStatus::Filled,
349        KrakenExecType::Trade => {
350            return match order_status {
351                Some(KrakenWsOrderStatus::Filled) => OrderStatus::Filled,
352                Some(KrakenWsOrderStatus::PartiallyFilled) | None => OrderStatus::PartiallyFilled,
353                Some(status) => status.into(),
354            };
355        }
356        _ => {}
357    }
358
359    match order_status {
360        Some(status) => status.into(),
361        None => OrderStatus::Accepted,
362    }
363}
364
365/// Parses Kraken order type to Nautilus order type.
366fn parse_order_type(order_type: Option<KrakenOrderType>) -> OrderType {
367    match order_type {
368        Some(KrakenOrderType::Market) => OrderType::Market,
369        Some(KrakenOrderType::Limit) => OrderType::Limit,
370        Some(KrakenOrderType::StopLoss) => OrderType::StopMarket,
371        Some(KrakenOrderType::TakeProfit) => OrderType::MarketIfTouched,
372        Some(KrakenOrderType::StopLossLimit) => OrderType::StopLimit,
373        Some(KrakenOrderType::TakeProfitLimit) => OrderType::LimitIfTouched,
374        // Trailing stops lack offset fields in WS reports, map to non-trailing equivalents
375        Some(KrakenOrderType::TrailingStop) => OrderType::StopMarket,
376        Some(KrakenOrderType::TrailingStopLimit) => OrderType::StopLimit,
377        Some(KrakenOrderType::SettlePosition) => OrderType::Market,
378        None => OrderType::Limit,
379    }
380}
381
382/// Parses Kraken time-in-force to Nautilus time-in-force.
383fn parse_time_in_force(
384    time_in_force: Option<KrakenTimeInForce>,
385    post_only: Option<bool>,
386) -> TimeInForce {
387    // Handle post_only flag
388    if post_only == Some(true) {
389        return TimeInForce::Gtc;
390    }
391
392    match time_in_force {
393        Some(KrakenTimeInForce::GoodTilCancelled) => TimeInForce::Gtc,
394        Some(KrakenTimeInForce::ImmediateOrCancel) => TimeInForce::Ioc,
395        Some(KrakenTimeInForce::GoodTilDate) => TimeInForce::Gtd,
396        Some(KrakenTimeInForce::FillOrKill) => TimeInForce::Fok,
397        None => TimeInForce::Gtc,
398    }
399}
400
401fn parse_liquidity_side(liquidity_ind: Option<KrakenLiquidityInd>) -> LiquiditySide {
402    liquidity_ind.map_or(LiquiditySide::NoLiquiditySide, Into::into)
403}
404
405/// Parses a Kraken WebSocket execution message into an [`OrderStatusReport`].
406///
407/// # Errors
408///
409/// Returns an error if required fields are missing or cannot be parsed.
410pub fn parse_ws_order_status_report(
411    exec: &KrakenWsExecutionData,
412    instrument: &InstrumentAny,
413    account_id: AccountId,
414    cached_order_qty: Option<Decimal>,
415    ts_init: UnixNanos,
416) -> anyhow::Result<OrderStatusReport> {
417    let instrument_id = instrument.id();
418    let venue_order_id = VenueOrderId::new(&exec.order_id);
419    let order_side = exec.side.map(Into::into);
420    let order_type = parse_order_type(exec.order_type);
421    let time_in_force = parse_time_in_force(exec.time_in_force, exec.post_only);
422    let order_status = parse_order_status(exec.exec_type, exec.order_status);
423
424    let price_precision = instrument.price_precision();
425    let size_precision = instrument.size_precision();
426
427    // Quantity fallback: order_qty -> cached -> cum_qty -> last_qty (for trade snapshots)
428    let last_qty = exec
429        .last_qty
430        .map(|qty| Quantity::from_decimal_dp(qty, size_precision))
431        .transpose()
432        .context("Failed to parse last_qty")?;
433
434    let filled_qty = exec
435        .cum_qty
436        .map(|qty| Quantity::from_decimal_dp(qty, size_precision))
437        .transpose()
438        .context("Failed to parse cum_qty")?
439        .or(last_qty)
440        .unwrap_or_else(|| Quantity::zero(size_precision));
441
442    let quantity = exec
443        .order_qty
444        .or(cached_order_qty)
445        .map(|qty| Quantity::from_decimal_dp(qty, size_precision))
446        .transpose()
447        .context("Failed to parse order_qty")?
448        .unwrap_or(filled_qty);
449
450    let ts_event = datetime_to_nanos(exec.timestamp, "execution.timestamp")?;
451
452    let mut report = OrderStatusReport::new(
453        account_id,
454        instrument_id,
455        None, // client_order_id set below if present
456        venue_order_id,
457        order_side,
458        order_type,
459        time_in_force,
460        order_status,
461        quantity,
462        filled_qty,
463        ts_event,
464        ts_event,
465        ts_init,
466        Some(UUID4::new()),
467    );
468
469    if let Some(ref cl_ord_id) = exec.cl_ord_id
470        && !cl_ord_id.is_empty()
471    {
472        report = report.with_client_order_id(ClientOrderId::new(cl_ord_id));
473    }
474
475    // Price fallback: limit_price -> avg_price -> last_price
476    // Note: pending_new messages may not include any price fields, which is fine for
477    // orders we submitted (engine already has the price from submission)
478    let price_value = exec
479        .limit_price
480        .filter(|p| *p > Decimal::ZERO)
481        .or(exec.avg_price.filter(|p| *p > Decimal::ZERO))
482        .or(exec.last_price.filter(|p| *p > Decimal::ZERO));
483
484    if let Some(px) = price_value {
485        let price =
486            Price::from_decimal_dp(px, price_precision).context("Failed to parse order price")?;
487        report = report.with_price(price);
488    }
489
490    // avg_px fallback: avg_price -> cum_cost / cum_qty -> last_price (for single trades/snapshots)
491    let avg_px = exec
492        .avg_price
493        .filter(|p| *p > Decimal::ZERO)
494        .or_else(|| match (exec.cum_cost, exec.cum_qty) {
495            (Some(cost), Some(qty)) if qty > Decimal::ZERO => Some(cost / qty),
496            _ => None,
497        })
498        .or_else(|| exec.last_price.filter(|p| *p > Decimal::ZERO));
499
500    if let Some(avg_price) = avg_px {
501        report.avg_px = Some(avg_price);
502    }
503
504    if exec.post_only == Some(true) {
505        report = report.with_post_only(true);
506    }
507
508    if exec.reduce_only == Some(true) {
509        report = report.with_reduce_only(true);
510    }
511
512    if let Some(ref reason) = exec.reason
513        && !reason.is_empty()
514    {
515        report = report.with_cancel_reason(reason.clone());
516    }
517
518    // Set trigger type for conditional orders (WebSocket doesn't provide trigger field)
519    let is_conditional = matches!(
520        order_type,
521        OrderType::StopMarket
522            | OrderType::StopLimit
523            | OrderType::MarketIfTouched
524            | OrderType::LimitIfTouched
525    );
526
527    if is_conditional {
528        report = report.with_trigger_type(TriggerType::Default);
529    }
530
531    Ok(report)
532}
533
534/// Parses a Kraken WebSocket trade execution into a [`FillReport`].
535///
536/// This should only be called when exec_type is "trade".
537///
538/// # Errors
539///
540/// Returns an error if required fields are missing or cannot be parsed.
541pub fn parse_ws_fill_report(
542    exec: &KrakenWsExecutionData,
543    instrument: &InstrumentAny,
544    account_id: AccountId,
545    ts_init: UnixNanos,
546) -> anyhow::Result<FillReport> {
547    let instrument_id = instrument.id();
548    let venue_order_id = VenueOrderId::new(&exec.order_id);
549
550    let exec_id = exec
551        .exec_id
552        .as_ref()
553        .context("Missing exec_id for trade execution")?;
554    let trade_id =
555        TradeId::new_checked(exec_id).context("Invalid exec_id in Kraken trade execution")?;
556
557    let order_side = exec
558        .side
559        .map(Into::into)
560        .context("Missing side for trade execution")?;
561
562    let price_precision = instrument.price_precision();
563    let size_precision = instrument.size_precision();
564
565    let last_qty = exec
566        .last_qty
567        .map(|qty| Quantity::from_decimal_dp(qty, size_precision))
568        .transpose()
569        .context("Failed to parse last_qty")?
570        .context("Missing last_qty for trade execution")?;
571
572    let last_px = exec
573        .last_price
574        .map(|px| Price::from_decimal_dp(px, price_precision))
575        .transpose()
576        .context("Failed to parse last_price")?
577        .context("Missing last_price for trade execution")?;
578
579    let liquidity_side = parse_liquidity_side(exec.liquidity_ind);
580
581    // Calculate commission from fees array
582    let commission = if let Some(ref fees) = exec.fees {
583        if let Some(fee) = fees.first() {
584            let currency = Currency::get_or_create_crypto(&fee.asset);
585            Money::from_decimal(fee.qty.abs(), currency).context("Failed to parse fill fee")?
586        } else {
587            Money::zero(instrument.quote_currency())
588        }
589    } else {
590        Money::zero(instrument.quote_currency())
591    };
592
593    let ts_event = datetime_to_nanos(exec.timestamp, "execution.timestamp")?;
594
595    let client_order_id = exec
596        .cl_ord_id
597        .as_ref()
598        .filter(|s| !s.is_empty())
599        .map(ClientOrderId::new);
600
601    Ok(FillReport::new(
602        account_id,
603        instrument_id,
604        venue_order_id,
605        trade_id,
606        order_side,
607        last_qty,
608        last_px,
609        commission,
610        liquidity_side,
611        client_order_id,
612        None, // venue_position_id
613        ts_event,
614        ts_init,
615        None, // report_id
616    ))
617}
618
619/// Parses a raw WebSocket JSON string and returns [`KrakenSpotWsMessage::OrderResponse`] if the
620/// message is an order-method response envelope, or `Ok(None)` for unrecognised messages.
621///
622/// # Errors
623///
624/// Returns an error if the message appears to be an order response but cannot be deserialized.
625pub fn parse_order_response(text: &str) -> anyhow::Result<Option<KrakenSpotWsMessage>> {
626    let value: serde_json::Value =
627        serde_json::from_str(text).with_context(|| format!("Failed to parse JSON: {text}"))?;
628
629    let method_str = match value.get("method").and_then(|m| m.as_str()) {
630        Some(s) => s.to_owned(),
631        None => return Ok(None),
632    };
633
634    if !matches!(
635        method_str.as_str(),
636        "add_order" | "amend_order" | "cancel_order" | "batch_add"
637    ) {
638        return Ok(None);
639    }
640
641    let response: KrakenWsOrderResponse = serde_json::from_value(value).with_context(|| {
642        format!("Failed to deserialize order response for method '{method_str}'")
643    })?;
644    Ok(Some(KrakenSpotWsMessage::OrderResponse(response)))
645}
646
647#[cfg(test)]
648mod tests {
649    use nautilus_model::{identifiers::Symbol, types::Currency};
650    use rstest::rstest;
651    use rust_decimal_macros::dec;
652    use ustr::Ustr;
653
654    use super::*;
655    use crate::{common::consts::KRAKEN_VENUE, websocket::spot_v2::messages::KrakenWsRawMessage};
656
657    const TS: UnixNanos = UnixNanos::new(1_700_000_000_000_000_000);
658
659    #[rstest]
660    fn test_parse_time_in_force_fok() {
661        assert_eq!(
662            parse_time_in_force(Some(KrakenTimeInForce::FillOrKill), None),
663            TimeInForce::Fok
664        );
665    }
666
667    fn load_test_json(filename: &str) -> String {
668        let path = format!("test_data/{filename}");
669        std::fs::read_to_string(&path)
670            .unwrap_or_else(|e| panic!("Failed to load test data from {path}: {e}"))
671    }
672
673    fn create_mock_instrument() -> InstrumentAny {
674        use nautilus_model::instruments::currency_pair::CurrencyPair;
675
676        let instrument_id = InstrumentId::new(Symbol::new("BTC/USD"), *KRAKEN_VENUE);
677        InstrumentAny::CurrencyPair(
678            CurrencyPair::builder()
679                .instrument_id(instrument_id)
680                .raw_symbol(Symbol::new("XBTUSDT"))
681                .base_currency(Currency::BTC())
682                .quote_currency(Currency::USDT())
683                .price_precision(1)
684                .size_precision(8)
685                .price_increment(Price::from("0.1"))
686                .size_increment(Quantity::from("0.00000001"))
687                .ts_event(TS)
688                .ts_init(TS)
689                .build()
690                .unwrap(),
691        )
692    }
693
694    #[rstest]
695    fn test_parse_quote_tick() {
696        let json = load_test_json("ws_ticker_snapshot.json");
697        let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
698        let ticker: KrakenWsTickerData = serde_json::from_str(message.data[0].get()).unwrap();
699
700        let instrument = create_mock_instrument();
701        let quote_tick = parse_quote_tick(&ticker, &instrument, TS).unwrap();
702
703        assert_eq!(quote_tick.instrument_id, instrument.id());
704        assert_eq!(quote_tick.bid_price, Price::from("105944.20"));
705        assert_eq!(quote_tick.ask_price, Price::from("105944.30"));
706        assert_eq!(quote_tick.bid_size, Quantity::from("2.5"));
707        assert_eq!(quote_tick.ask_size, Quantity::from("3.2"));
708        assert_eq!(
709            quote_tick.ts_event,
710            UnixNanos::from(1_671_960_659_123_456_000)
711        );
712        assert_eq!(quote_tick.ts_init, TS);
713    }
714
715    #[rstest]
716    fn test_parse_trade_tick() {
717        let json = load_test_json("ws_trade_update.json");
718        let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
719        let trade: KrakenWsTradeData = serde_json::from_str(message.data[0].get()).unwrap();
720
721        let instrument = create_mock_instrument();
722        let trade_tick = parse_trade_tick(&trade, &instrument, TS).unwrap();
723
724        assert_eq!(trade_tick.instrument_id, instrument.id());
725        assert_eq!(trade_tick.price, Price::from("105944.20"));
726        assert_eq!(trade_tick.size, Quantity::from("0.00027625"));
727        assert!(matches!(
728            trade_tick.aggressor_side,
729            AggressorSide::Buy | AggressorSide::Sell
730        ));
731        assert_eq!(
732            trade_tick.ts_event,
733            UnixNanos::from(1_696_613_755_440_295_000)
734        );
735        assert_eq!(trade_tick.ts_init, TS);
736    }
737
738    #[rstest]
739    fn test_parse_book_deltas_snapshot() {
740        let json = load_test_json("ws_book_snapshot.json");
741        let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
742        let book: KrakenWsBookData = serde_json::from_str(message.data[0].get()).unwrap();
743
744        let instrument = create_mock_instrument();
745        let deltas = parse_book_deltas(&book, &instrument, 1, true, TS).unwrap();
746
747        assert!(!deltas.is_empty());
748
749        let bid_count = deltas
750            .iter()
751            .filter(|d| d.order.side == OrderSide::Buy.into())
752            .count();
753        let ask_count = deltas
754            .iter()
755            .filter(|d| d.order.side == OrderSide::Sell.into())
756            .count();
757
758        assert!(bid_count > 0);
759        assert!(ask_count > 0);
760
761        let first_delta = &deltas[0];
762        assert_eq!(first_delta.instrument_id, instrument.id());
763        assert_eq!(first_delta.action, BookAction::Clear);
764        assert!(RecordFlag::F_SNAPSHOT.matches(first_delta.flags));
765        assert!(!RecordFlag::F_LAST.matches(first_delta.flags));
766
767        assert!(deltas[1..].iter().all(|d| d.action == BookAction::Add));
768        assert!(
769            deltas[1..]
770                .iter()
771                .all(|d| RecordFlag::F_MBP.matches(d.flags))
772        );
773        assert!(
774            deltas[1..]
775                .iter()
776                .all(|d| RecordFlag::F_SNAPSHOT.matches(d.flags))
777        );
778        assert!(RecordFlag::F_LAST.matches(deltas.last().unwrap().flags));
779
780        let expected_ts_event = UnixNanos::from(1_696_613_755_440_295_000);
781        assert!(deltas.iter().all(|d| d.ts_event == expected_ts_event));
782        assert!(deltas.iter().all(|d| d.ts_init == TS));
783    }
784
785    #[rstest]
786    fn test_parse_book_deltas_update() {
787        let json = load_test_json("ws_book_update.json");
788        let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
789        let book: KrakenWsBookData = serde_json::from_str(message.data[0].get()).unwrap();
790
791        let instrument = create_mock_instrument();
792        let deltas = parse_book_deltas(&book, &instrument, 1, false, TS).unwrap();
793
794        assert!(!deltas.is_empty());
795
796        let first_delta = &deltas[0];
797        assert_eq!(first_delta.instrument_id, instrument.id());
798        assert_eq!(first_delta.action, BookAction::Update);
799        assert_eq!(first_delta.order.side, OrderSide::Buy.into());
800        assert_eq!(first_delta.order.price, Price::from("105944.20"));
801        assert!(RecordFlag::F_MBP.matches(first_delta.flags));
802        assert!(RecordFlag::F_LAST.matches(first_delta.flags));
803        assert!(!RecordFlag::F_SNAPSHOT.matches(first_delta.flags));
804
805        let expected_ts_event = UnixNanos::from(1_696_613_755_440_295_000);
806        assert!(deltas.iter().all(|d| d.ts_event == expected_ts_event));
807        assert!(deltas.iter().all(|d| d.ts_init == TS));
808    }
809
810    #[rstest]
811    fn test_parse_book_deltas_snapshot_skips_zero_qty_levels() {
812        let book = KrakenWsBookData {
813            symbol: Ustr::from("BTC/USD"),
814            bids: Some(vec![KrakenWsBookLevel {
815                price: dec!(100),
816                qty: Decimal::ZERO,
817            }]),
818            asks: Some(vec![KrakenWsBookLevel {
819                price: dec!(101),
820                qty: dec!(2),
821            }]),
822            checksum: Some(0),
823            timestamp: "2024-01-01T00:00:00Z".parse().unwrap(),
824        };
825
826        let instrument = create_mock_instrument();
827        let deltas = parse_book_deltas(&book, &instrument, 7, true, TS).unwrap();
828
829        assert_eq!(deltas.len(), 2);
830        assert_eq!(deltas[0].action, BookAction::Clear);
831        assert_eq!(deltas[0].sequence, 7);
832        assert!(RecordFlag::F_SNAPSHOT.matches(deltas[0].flags));
833        assert!(!RecordFlag::F_LAST.matches(deltas[0].flags));
834
835        let add = &deltas[1];
836        assert_eq!(add.action, BookAction::Add);
837        assert_eq!(add.sequence, 8);
838        assert_eq!(add.order.side, OrderSide::Sell.into());
839        assert_eq!(add.order.price, Price::from("101.0"));
840        assert!(RecordFlag::F_MBP.matches(add.flags));
841        assert!(RecordFlag::F_SNAPSHOT.matches(add.flags));
842        assert!(RecordFlag::F_LAST.matches(add.flags));
843    }
844
845    #[rstest]
846    fn test_parse_book_deltas_update_zero_qty_deletes_level() {
847        let book = KrakenWsBookData {
848            symbol: Ustr::from("BTC/USD"),
849            bids: Some(vec![KrakenWsBookLevel {
850                price: dec!(100),
851                qty: Decimal::ZERO,
852            }]),
853            asks: Some(vec![]),
854            checksum: Some(0),
855            timestamp: "2024-01-01T00:00:00Z".parse().unwrap(),
856        };
857
858        let instrument = create_mock_instrument();
859        let deltas = parse_book_deltas(&book, &instrument, 11, false, TS).unwrap();
860
861        assert_eq!(deltas.len(), 1);
862        let delete = &deltas[0];
863        assert_eq!(delete.action, BookAction::Delete);
864        assert_eq!(delete.sequence, 11);
865        assert_eq!(delete.order.side, OrderSide::Buy.into());
866        assert_eq!(delete.order.price, Price::from("100.0"));
867        assert_eq!(delete.order.size.raw, 0);
868        assert!(RecordFlag::F_MBP.matches(delete.flags));
869        assert!(RecordFlag::F_LAST.matches(delete.flags));
870        assert!(!RecordFlag::F_SNAPSHOT.matches(delete.flags));
871    }
872
873    #[rstest]
874    fn test_parse_ws_order_status_report_preserves_decimal_avg_px() {
875        let execution = ws_execution_data(Some(KrakenOrderSide::Buy));
876
877        let report = parse_ws_order_status_report(
878            &execution,
879            &create_mock_instrument(),
880            AccountId::from("KRAKEN-001"),
881            None,
882            TS,
883        )
884        .unwrap();
885
886        assert_eq!(report.avg_px, Some(dec!(0.1234567890123456789012345678)));
887    }
888
889    #[rstest]
890    fn test_parse_ws_order_status_report_preserves_missing_side() {
891        let execution = ws_execution_data(None);
892
893        let report = parse_ws_order_status_report(
894            &execution,
895            &create_mock_instrument(),
896            AccountId::from("KRAKEN-001"),
897            None,
898            TS,
899        )
900        .unwrap();
901
902        assert_eq!(report.order_side, None);
903    }
904
905    #[rstest]
906    fn test_parse_ws_fill_report_rejects_missing_side() {
907        let mut execution = ws_execution_data(None);
908        execution.exec_type = KrakenExecType::Trade;
909        execution.exec_id = Some("TRADE-1".to_string());
910        execution.last_qty = Some(dec!(1));
911        execution.last_price = Some(dec!(100));
912
913        let error = parse_ws_fill_report(
914            &execution,
915            &create_mock_instrument(),
916            AccountId::from("KRAKEN-001"),
917            TS,
918        )
919        .expect_err("a trade execution without a side must be rejected");
920
921        assert_eq!(error.to_string(), "Missing side for trade execution");
922    }
923
924    fn ws_execution_data(side: Option<KrakenOrderSide>) -> KrakenWsExecutionData {
925        KrakenWsExecutionData {
926            exec_type: KrakenExecType::Status,
927            order_id: "ORDER-1".to_string(),
928            cl_ord_id: Some("CLIENT-1".to_string()),
929            symbol: Some("BTC/USD".to_string()),
930            side,
931            order_type: Some(KrakenOrderType::Limit),
932            order_qty: Some(dec!(3)),
933            limit_price: None,
934            order_status: Some(KrakenWsOrderStatus::PartiallyFilled),
935            cum_qty: Some(dec!(3)),
936            cum_cost: Some(dec!(0.3703703670370370367037037034)),
937            avg_price: None,
938            time_in_force: Some(KrakenTimeInForce::GoodTilCancelled),
939            post_only: Some(false),
940            reduce_only: Some(false),
941            timestamp: "2024-01-01T00:00:00Z".parse().unwrap(),
942            exec_id: None,
943            last_qty: None,
944            last_price: None,
945            cost: None,
946            liquidity_ind: None,
947            fees: None,
948            fee_usd_equiv: None,
949            reason: None,
950        }
951    }
952
953    #[rstest]
954    fn test_datetime_to_nanos() {
955        let dt = "2023-10-06T17:35:55.440295Z".parse::<Timestamp>().unwrap();
956        let result = datetime_to_nanos(dt, "test").unwrap();
957        assert_eq!(result, UnixNanos::from(1_696_613_755_440_295_000));
958    }
959
960    #[rstest]
961    fn test_datetime_to_nanos_out_of_range_errors() {
962        let dt = "1500-01-01T00:00:00Z".parse::<Timestamp>().unwrap();
963        let result = datetime_to_nanos(dt, "test");
964        assert!(result.is_err());
965        let err = result.unwrap_err().to_string();
966        assert!(err.contains("test"));
967    }
968
969    #[rstest]
970    fn test_parse_ws_bar() {
971        let json = load_test_json("ws_ohlc_update.json");
972        let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
973        let ohlc: KrakenWsOhlcData = serde_json::from_str(message.data[0].get()).unwrap();
974
975        let instrument = create_mock_instrument();
976        let bar = parse_ws_bar(&ohlc, &instrument, TS).unwrap();
977
978        assert_eq!(bar.bar_type.instrument_id(), instrument.id());
979        assert_eq!(bar.open, Price::from("106038.2"));
980        assert_eq!(bar.high, Price::from("106044.3"));
981        assert_eq!(bar.low, Price::from("106038.1"));
982        assert_eq!(bar.close, Price::from("106040.1"));
983        assert_eq!(bar.volume, Quantity::from("30927.68066226"));
984
985        let spec = bar.bar_type.spec();
986        assert_eq!(spec.step.get(), 1);
987        assert_eq!(spec.aggregation, BarAggregation::Minute);
988        assert_eq!(spec.price_type, PriceType::Last);
989
990        // Verify ts_event is computed as interval_begin + interval (close time)
991        // interval_begin is 2023-10-04T16:25:00Z, interval is 1 minute, so close is 16:26:00Z
992        let expected_close = ohlc.interval_begin + jiff::SignedDuration::from_mins(1);
993        let expected_ts_event =
994            UnixNanos::from(u64::try_from(expected_close.as_nanosecond()).unwrap());
995        assert_eq!(bar.ts_event, expected_ts_event);
996    }
997
998    #[rstest]
999    fn test_interval_to_bar_spec() {
1000        let test_cases = [
1001            (1, 1, BarAggregation::Minute),
1002            (5, 5, BarAggregation::Minute),
1003            (15, 15, BarAggregation::Minute),
1004            (30, 30, BarAggregation::Minute),
1005            (60, 1, BarAggregation::Hour),
1006            (240, 4, BarAggregation::Hour),
1007            (1440, 1, BarAggregation::Day),
1008            (10080, 1, BarAggregation::Week),
1009            (21600, 15, BarAggregation::Day), // 21600 min = 15 days
1010        ];
1011
1012        for (interval, expected_step, expected_aggregation) in test_cases {
1013            let spec = interval_to_bar_spec(interval).unwrap();
1014            assert_eq!(
1015                spec.step.get(),
1016                expected_step,
1017                "Failed for interval {interval}"
1018            );
1019            assert_eq!(
1020                spec.aggregation, expected_aggregation,
1021                "Failed for interval {interval}"
1022            );
1023            assert_eq!(spec.price_type, PriceType::Last);
1024        }
1025    }
1026
1027    #[rstest]
1028    fn test_interval_to_bar_spec_invalid() {
1029        let result = interval_to_bar_spec(999);
1030        assert!(result.is_err());
1031    }
1032
1033    #[rstest]
1034    fn test_parse_order_response_envelope_returns_order_response_variant() {
1035        use crate::websocket::spot_v2::enums::KrakenWsMethod;
1036
1037        let raw = load_test_json("ws_add_order_response_success.json");
1038        let parsed = parse_order_response(&raw).expect("parse ok");
1039        match parsed {
1040            Some(KrakenSpotWsMessage::OrderResponse(resp)) => {
1041                assert_eq!(resp.method, KrakenWsMethod::AddOrder);
1042                assert_eq!(resp.req_id, Some(42));
1043                assert!(resp.success);
1044            }
1045            other => panic!("expected OrderResponse, was {other:?}"),
1046        }
1047    }
1048
1049    #[rstest]
1050    fn test_parse_order_response_returns_none_for_non_order_method() {
1051        let json = r#"{"method":"subscribe","req_id":1,"success":true}"#;
1052        let result = parse_order_response(json).expect("parse ok");
1053        assert!(result.is_none());
1054    }
1055
1056    #[rstest]
1057    fn test_parse_order_response_returns_none_for_data_message() {
1058        let json = r#"{"channel":"ticker","type":"snapshot","data":[]}"#;
1059        let result = parse_order_response(json).expect("parse ok");
1060        assert!(result.is_none());
1061    }
1062}