Skip to main content

nautilus_okx/websocket/
messages.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Data structures modelling OKX WebSocket request and response payloads.
17
18use derive_builder::Builder;
19use nautilus_model::{
20    data::{Data, FundingRateUpdate, InstrumentStatus, OrderBookDeltas},
21    events::{
22        AccountState, OrderAccepted, OrderCancelRejected, OrderCanceled, OrderExpired,
23        OrderModifyRejected, OrderRejected, OrderTriggered, OrderUpdated,
24    },
25    identifiers::ClientOrderId,
26    instruments::InstrumentAny,
27    reports::{FillReport, OrderStatusReport, PositionStatusReport},
28};
29use serde::{Deserialize, Serialize};
30use ustr::Ustr;
31
32use super::enums::{OKXWsChannel, OKXWsOperation};
33use crate::{
34    common::{
35        enums::{
36            OKXAlgoOrderStatus, OKXAlgoOrderType, OKXBookAction, OKXCandleConfirm, OKXExecType,
37            OKXInstrumentType, OKXMarginMode, OKXOrderCategory, OKXOrderStatus, OKXOrderType,
38            OKXPositionSide, OKXPriceType, OKXQuickMarginType, OKXSelfTradePreventionMode,
39            OKXSettlementState, OKXSide, OKXTargetCurrency, OKXTradeMode, OKXTriggerType,
40        },
41        models::{OKXInstrument, OKXRpiBookLevel},
42        parse::{
43            deserialize_empty_string_as_none, deserialize_empty_ustr_as_none,
44            deserialize_string_to_u64, deserialize_target_currency_as_none,
45        },
46    },
47    http::models::OKXSpreadOrder,
48    websocket::enums::OKXSubscriptionEvent,
49};
50
51#[derive(Debug, Clone)]
52pub enum NautilusWsMessage {
53    Data(Vec<Data>),
54    Deltas(OrderBookDeltas),
55    FundingRates(Vec<FundingRateUpdate>),
56    Instrument(Box<InstrumentAny>, Option<InstrumentStatus>),
57    InstrumentStatus(InstrumentStatus),
58    AccountUpdate(AccountState),
59    PositionUpdate(PositionStatusReport),
60    OrderAccepted(OrderAccepted),
61    OrderCanceled(OrderCanceled),
62    OrderExpired(OrderExpired),
63    OrderRejected(OrderRejected),
64    OrderCancelRejected(OrderCancelRejected),
65    OrderModifyRejected(OrderModifyRejected),
66    OrderTriggered(OrderTriggered),
67    OrderUpdated(OrderUpdated),
68    ExecutionReports(Vec<ExecutionReport>),
69    Error(OKXWebSocketError),
70    Raw(serde_json::Value), // Unhandled channels
71    Reconnected,
72    Authenticated,
73}
74
75/// Represents an OKX WebSocket error.
76#[derive(Debug, Clone, Serialize, Deserialize)]
77pub struct OKXWebSocketError {
78    /// Error code from OKX (e.g., "50101").
79    pub code: String,
80    /// Error message from OKX.
81    pub message: String,
82    /// Connection ID if available.
83    pub conn_id: Option<String>,
84    /// Timestamp when the error occurred.
85    pub timestamp: u64,
86}
87
88#[derive(Debug, Clone)]
89#[allow(
90    clippy::large_enum_variant,
91    reason = "the variant size gap only crosses the threshold when high-precision widens the raw types"
92)]
93pub enum ExecutionReport {
94    Order(OrderStatusReport),
95    Fill(FillReport),
96}
97
98/// Output from the OKX WebSocket handler.
99///
100/// Contains venue-specific types only. Data parsing occurs in `PyOKXWebSocketClient`
101/// (using an instruments cache), and execution parsing occurs in `execution.rs`
102/// (using the system Cache for order lookups).
103#[derive(Debug)]
104pub enum OKXWsMessage {
105    /// Order book snapshot or update.
106    BookData {
107        arg: OKXWebSocketArg,
108        action: OKXBookAction,
109        data: Vec<OKXBookMsg>,
110    },
111    /// Retail Price Improvement order book snapshot or update.
112    RpiBookData {
113        arg: OKXWebSocketArg,
114        action: OKXBookAction,
115        data: Vec<OKXRpiBookMsg>,
116    },
117    /// Data from a non-book channel (trades, tickers, mark price, funding, candles, etc.).
118    ChannelData {
119        channel: OKXWsChannel,
120        inst_id: Option<Ustr>,
121        data: serde_json::Value,
122    },
123    /// Response to a WebSocket order command (place, cancel, amend, mass-cancel).
124    OrderResponse {
125        id: Option<String>,
126        op: OKXWsOperation,
127        code: String,
128        msg: String,
129        data: Vec<serde_json::Value>,
130    },
131    /// Order push channel updates.
132    Orders(Vec<OKXOrderMsg>),
133    /// Nitro spread order push channel updates.
134    SpreadOrders(Vec<OKXSpreadOrder>),
135    /// Algo order push channel updates.
136    AlgoOrders(Vec<OKXAlgoOrderMsg>),
137    /// Account channel update (raw JSON).
138    Account(serde_json::Value),
139    /// Positions channel update (raw JSON).
140    Positions(serde_json::Value),
141    /// Liquidation risk warnings for account positions.
142    LiquidationWarnings(Vec<OKXLiquidationWarningMsg>),
143    /// Instrument definition updates.
144    Instruments(Vec<OKXInstrument>),
145    /// A WebSocket send failed without a structured venue response.
146    SendFailed {
147        request_id: String,
148        client_order_ids: Vec<ClientOrderId>,
149        op: Option<OKXWsOperation>,
150        error: super::error::OKXWsError,
151    },
152    /// The venue rejected a subscribe request, so no data will flow for it.
153    SubscriptionFailed {
154        channel: OKXWsChannel,
155        inst_id: Option<Ustr>,
156        code: String,
157        msg: String,
158    },
159    /// Error received from OKX.
160    Error(OKXWebSocketError),
161    /// WebSocket reconnected.
162    Reconnected,
163    /// WebSocket authenticated.
164    Authenticated,
165}
166
167/// Generic WebSocket request for OKX trading commands.
168#[derive(Debug, Serialize)]
169#[serde(rename_all = "camelCase")]
170pub struct OKXWsRequest<T> {
171    /// Client request ID (required for order operations).
172    #[serde(skip_serializing_if = "Option::is_none")]
173    pub id: Option<String>,
174    /// Operation type (order, cancel-order, amend-order).
175    pub op: OKXWsOperation,
176    /// Request effective deadline. Unix timestamp format in milliseconds.
177    /// This is when the request itself expires, not related to order expiration.
178    #[serde(skip_serializing_if = "Option::is_none")]
179    pub exp_time: Option<String>,
180    /// Arguments payload for the operation.
181    pub args: Vec<T>,
182}
183
184/// OKX WebSocket authentication message.
185#[derive(Debug, Serialize)]
186pub struct OKXAuthentication {
187    pub op: &'static str,
188    pub args: Vec<OKXAuthenticationArg>,
189}
190
191/// OKX WebSocket authentication arguments.
192#[derive(Debug, Serialize)]
193#[serde(rename_all = "camelCase")]
194pub struct OKXAuthenticationArg {
195    pub api_key: String,
196    pub passphrase: String,
197    pub timestamp: String,
198    pub sign: String,
199}
200
201#[derive(Debug, Serialize)]
202pub struct OKXSubscription {
203    pub op: OKXWsOperation,
204    pub args: Vec<OKXSubscriptionArg>,
205}
206
207#[derive(Clone, Debug)]
208pub struct OKXSubscriptionArg {
209    pub channel: OKXWsChannel,
210    pub inst_type: Option<OKXInstrumentType>,
211    pub inst_family: Option<Ustr>,
212    pub inst_id: Option<Ustr>,
213}
214
215impl Serialize for OKXSubscriptionArg {
216    fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
217        use serde::ser::SerializeMap;
218
219        let mut map = serializer.serialize_map(None)?;
220        map.serialize_entry("channel", &self.channel)?;
221
222        if let Some(inst_type) = &self.inst_type {
223            map.serialize_entry("instType", inst_type)?;
224        }
225
226        if let Some(inst_family) = &self.inst_family {
227            map.serialize_entry("instFamily", inst_family)?;
228        }
229
230        if let Some(inst_id) = &self.inst_id {
231            let key = if self.channel.is_spread() {
232                "sprdId"
233            } else {
234                "instId"
235            };
236            map.serialize_entry(key, inst_id)?;
237        }
238
239        map.end()
240    }
241}
242
243/// OKX WebSocket message variants.
244///
245/// Uses custom deserialization that checks discriminant fields (event, op, action)
246/// to determine the correct variant.
247#[derive(Debug)]
248pub enum OKXWsFrame {
249    Login {
250        event: String,
251        code: String,
252        msg: String,
253        conn_id: String,
254    },
255    Subscription {
256        event: OKXSubscriptionEvent,
257        arg: OKXWebSocketArg,
258        conn_id: String,
259        code: Option<String>,
260        msg: Option<String>,
261    },
262    ChannelConnCount {
263        event: String,
264        channel: OKXWsChannel,
265        conn_count: String,
266        conn_id: String,
267    },
268    OrderResponse {
269        id: Option<String>,
270        op: OKXWsOperation,
271        code: String,
272        msg: String,
273        data: Vec<serde_json::Value>,
274    },
275    BookData {
276        arg: OKXWebSocketArg,
277        action: OKXBookAction,
278        data: Vec<OKXBookMsg>,
279    },
280    RpiBookData {
281        arg: OKXWebSocketArg,
282        action: OKXBookAction,
283        data: Vec<OKXRpiBookMsg>,
284    },
285    Data {
286        arg: OKXWebSocketArg,
287        data: serde_json::Value,
288    },
289    Error {
290        arg: Option<OKXWebSocketArg>,
291        code: String,
292        msg: String,
293    },
294    Ping,
295    Reconnected,
296}
297
298impl<'de> Deserialize<'de> for OKXWsFrame {
299    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
300    where
301        D: serde::Deserializer<'de>,
302    {
303        use serde::de::Error;
304
305        // Buffer once via serde_json::Value, then take ownership of the
306        // typed subtrees with `.remove(...)` instead of `.cloned()`: the
307        // latter deep-cloned every level for L2 books and dominated the
308        // inbound decode cost.
309        let mut value = serde_json::Value::deserialize(deserializer)?;
310        let obj = value
311            .as_object_mut()
312            .ok_or_else(|| D::Error::custom("expected JSON object for OKXWsFrame"))?;
313
314        // Check discriminant fields in priority order. Discriminants stay
315        // borrowed via `.get(...).as_str()`; only the structured payloads
316        // (`arg`, `data`, `channel`, `op`, `action`) are moved out.
317
318        // 1. "event" field - Login, Subscription, ChannelConnCount, or Error
319        if let Some(event) = obj.get("event").and_then(|v| v.as_str()) {
320            match event {
321                "login" => return parse_login(obj),
322                "subscribe" | "unsubscribe" => return parse_subscription(obj),
323                "error" => return parse_error(obj),
324                _ if obj.contains_key("channel") && obj.contains_key("connCount") => {
325                    return parse_channel_conn_count(obj);
326                }
327                _ => {}
328            }
329        }
330
331        // 2. "op" field - OrderResponse
332        if obj.contains_key("op") {
333            return parse_order_response(obj);
334        }
335
336        // 3. "action" + "arg" - BookData
337        if obj.contains_key("action") && obj.contains_key("arg") {
338            return parse_book_data(obj);
339        }
340
341        // 4. "arg" + "data" without "action" - Data
342        if obj.contains_key("arg") && obj.contains_key("data") {
343            return parse_data(obj);
344        }
345
346        // 5. Fallback to Error if it has "code" and "msg"
347        if obj.contains_key("code") && obj.contains_key("msg") {
348            return parse_error(obj);
349        }
350
351        // No variant matched; no `remove` happened above, so `value` is still
352        // intact. Serialize it back into the error message to preserve the
353        // original diagnostic shape.
354        Err(D::Error::custom(format!(
355            "cannot determine OKXWsFrame variant from: {}",
356            serde_json::to_string(&value).unwrap_or_default()
357        )))
358    }
359}
360
361#[inline]
362fn take_str<E: serde::de::Error>(
363    obj: &mut serde_json::Map<String, serde_json::Value>,
364    key: &'static str,
365) -> Result<String, E> {
366    match obj.remove(key) {
367        Some(serde_json::Value::String(s)) => Ok(s),
368        Some(_) => Err(E::custom(format!("field `{key}` is not a string"))),
369        None => Err(E::missing_field(key)),
370    }
371}
372
373#[inline]
374fn take_optional_str(
375    obj: &mut serde_json::Map<String, serde_json::Value>,
376    key: &'static str,
377) -> Option<String> {
378    match obj.remove(key) {
379        Some(serde_json::Value::String(s)) => Some(s),
380        _ => None,
381    }
382}
383
384fn parse_login<E: serde::de::Error>(
385    obj: &mut serde_json::Map<String, serde_json::Value>,
386) -> Result<OKXWsFrame, E> {
387    Ok(OKXWsFrame::Login {
388        event: take_str(obj, "event")?,
389        code: take_str(obj, "code")?,
390        msg: take_str(obj, "msg")?,
391        conn_id: take_str(obj, "connId")?,
392    })
393}
394
395fn parse_subscription<E: serde::de::Error>(
396    obj: &mut serde_json::Map<String, serde_json::Value>,
397) -> Result<OKXWsFrame, E> {
398    let event_val = obj
399        .remove("event")
400        .ok_or_else(|| E::missing_field("event"))?;
401    let event: OKXSubscriptionEvent =
402        serde_json::from_value(event_val).map_err(|e| E::custom(format!("invalid event: {e}")))?;
403
404    let arg_val = obj.remove("arg").ok_or_else(|| E::missing_field("arg"))?;
405    let arg: OKXWebSocketArg =
406        serde_json::from_value(arg_val).map_err(|e| E::custom(format!("invalid arg: {e}")))?;
407
408    Ok(OKXWsFrame::Subscription {
409        event,
410        arg,
411        conn_id: take_str(obj, "connId")?,
412        code: take_optional_str(obj, "code"),
413        msg: take_optional_str(obj, "msg"),
414    })
415}
416
417fn parse_channel_conn_count<E: serde::de::Error>(
418    obj: &mut serde_json::Map<String, serde_json::Value>,
419) -> Result<OKXWsFrame, E> {
420    let channel_val = obj
421        .remove("channel")
422        .ok_or_else(|| E::missing_field("channel"))?;
423    let channel: OKXWsChannel = serde_json::from_value(channel_val)
424        .map_err(|e| E::custom(format!("invalid channel: {e}")))?;
425
426    Ok(OKXWsFrame::ChannelConnCount {
427        event: take_str(obj, "event")?,
428        channel,
429        conn_count: take_str(obj, "connCount")?,
430        conn_id: take_str(obj, "connId")?,
431    })
432}
433
434fn parse_order_response<E: serde::de::Error>(
435    obj: &mut serde_json::Map<String, serde_json::Value>,
436) -> Result<OKXWsFrame, E> {
437    let op_val = obj.remove("op").ok_or_else(|| E::missing_field("op"))?;
438    let op: OKXWsOperation =
439        serde_json::from_value(op_val).map_err(|e| E::custom(format!("invalid op: {e}")))?;
440
441    let data: Vec<serde_json::Value> = match obj.remove("data") {
442        Some(v) => {
443            serde_json::from_value(v).map_err(|e| E::custom(format!("invalid data: {e}")))?
444        }
445        None => Vec::new(),
446    };
447
448    Ok(OKXWsFrame::OrderResponse {
449        id: take_optional_str(obj, "id"),
450        op,
451        code: take_str(obj, "code")?,
452        msg: take_str(obj, "msg")?,
453        data,
454    })
455}
456
457fn parse_book_data<E: serde::de::Error>(
458    obj: &mut serde_json::Map<String, serde_json::Value>,
459) -> Result<OKXWsFrame, E> {
460    let arg_val = obj.remove("arg").ok_or_else(|| E::missing_field("arg"))?;
461    let arg: OKXWebSocketArg =
462        serde_json::from_value(arg_val).map_err(|e| E::custom(format!("invalid arg: {e}")))?;
463
464    let action_val = obj
465        .remove("action")
466        .ok_or_else(|| E::missing_field("action"))?;
467    let action: OKXBookAction = serde_json::from_value(action_val)
468        .map_err(|e| E::custom(format!("invalid action: {e}")))?;
469
470    let data_val = obj.remove("data").ok_or_else(|| E::missing_field("data"))?;
471    if arg.channel == OKXWsChannel::BooksRpi {
472        let data: Vec<OKXRpiBookMsg> = serde_json::from_value(data_val)
473            .map_err(|e| E::custom(format!("invalid data: {e}")))?;
474        return Ok(OKXWsFrame::RpiBookData { arg, action, data });
475    }
476
477    let data: Vec<OKXBookMsg> =
478        serde_json::from_value(data_val).map_err(|e| E::custom(format!("invalid data: {e}")))?;
479    Ok(OKXWsFrame::BookData { arg, action, data })
480}
481
482fn parse_data<E: serde::de::Error>(
483    obj: &mut serde_json::Map<String, serde_json::Value>,
484) -> Result<OKXWsFrame, E> {
485    let arg_val = obj.remove("arg").ok_or_else(|| E::missing_field("arg"))?;
486    let arg: OKXWebSocketArg =
487        serde_json::from_value(arg_val).map_err(|e| E::custom(format!("invalid arg: {e}")))?;
488
489    let data = obj.remove("data").ok_or_else(|| E::missing_field("data"))?;
490
491    Ok(OKXWsFrame::Data { arg, data })
492}
493
494fn parse_error<E: serde::de::Error>(
495    obj: &mut serde_json::Map<String, serde_json::Value>,
496) -> Result<OKXWsFrame, E> {
497    let arg = obj
498        .remove("arg")
499        .map(serde_json::from_value)
500        .transpose()
501        .map_err(|e| E::custom(format!("invalid arg: {e}")))?;
502
503    Ok(OKXWsFrame::Error {
504        arg,
505        code: take_str(obj, "code")?,
506        msg: take_str(obj, "msg")?,
507    })
508}
509
510#[derive(Debug, Serialize, Deserialize)]
511#[serde(rename_all = "camelCase")]
512pub struct OKXWebSocketArg {
513    /// Channel name that pushed the data.
514    pub channel: OKXWsChannel,
515    // Spread channels identify the instrument by `sprdId`; a spread's symbol equals
516    // its `sprdId`, and a message carries `instId` xor `sprdId`, so the alias resolves
517    // both to one field without collision.
518    #[serde(default, alias = "sprdId")]
519    pub inst_id: Option<Ustr>,
520    #[serde(default)]
521    pub inst_type: Option<OKXInstrumentType>,
522    #[serde(default)]
523    pub inst_family: Option<Ustr>,
524    #[serde(default)]
525    pub bar: Option<Ustr>,
526}
527
528/// Ticker data for an instrument.
529#[derive(Debug, Serialize, Deserialize)]
530#[serde(rename_all = "camelCase")]
531pub struct OKXTickerMsg {
532    /// Instrument type, e.g. "SPOT", "SWAP".
533    pub inst_type: OKXInstrumentType,
534    /// Instrument ID, e.g. "BTC-USDT".
535    pub inst_id: Ustr,
536    /// Last traded price.
537    #[serde(rename = "last")]
538    pub last_px: String,
539    /// Last traded size.
540    pub last_sz: String,
541    /// Best ask price.
542    pub ask_px: String,
543    /// Best ask size.
544    pub ask_sz: String,
545    /// Best bid price.
546    pub bid_px: String,
547    /// Best bid size.
548    pub bid_sz: String,
549    /// 24-hour opening price.
550    pub open24h: String,
551    /// 24-hour highest price.
552    pub high24h: String,
553    /// 24-hour lowest price.
554    pub low24h: String,
555    /// 24-hour trading volume in quote currency.
556    pub vol_ccy_24h: String,
557    /// 24-hour trading volume.
558    pub vol24h: String,
559    /// The opening price of the day (UTC 0).
560    pub sod_utc0: String,
561    /// The opening price of the day (UTC 8).
562    pub sod_utc8: String,
563    /// Timestamp of the data generation, Unix timestamp format in milliseconds.
564    #[serde(deserialize_with = "deserialize_string_to_u64")]
565    pub ts: u64,
566    /// Order source for RPI liquidity identification.
567    #[serde(default)]
568    pub source: Option<String>,
569}
570
571/// Represents a single order in the order book.
572#[derive(Debug, Serialize, Deserialize)]
573pub struct OrderBookEntry {
574    /// Price of the order.
575    pub price: String,
576    /// Size of the order.
577    pub size: String,
578    // Spread book levels (`sprd-books5`) are 3-element `[price, size, count]`,
579    // omitting the liquidated-orders field standard books carry; default the
580    // trailing counts so both array shapes deserialize. Only price/size are used.
581    /// Number of liquidated orders.
582    #[serde(default)]
583    pub liquidated_orders_count: String,
584    /// Total number of orders at this price.
585    #[serde(default)]
586    pub orders_count: String,
587}
588
589/// Order book data for an instrument.
590#[derive(Debug, Serialize, Deserialize)]
591#[serde(rename_all = "camelCase")]
592pub struct OKXBookMsg {
593    /// Order book asks [price, size, liquidated orders count, orders count].
594    pub asks: Vec<OrderBookEntry>,
595    /// Order book bids [price, size, liquidated orders count, orders count].
596    pub bids: Vec<OrderBookEntry>,
597    /// Checksum value.
598    pub checksum: Option<i64>,
599    /// Sequence ID of the last sent message. Only applicable to books, books-l2-tbt, books50-l2-tbt.
600    pub prev_seq_id: Option<i64>,
601    /// Sequence ID of the current message, implementation details below.
602    pub seq_id: u64,
603    /// Order book generation time, Unix timestamp format in milliseconds, e.g. 1597026383085.
604    #[serde(deserialize_with = "deserialize_string_to_u64")]
605    pub ts: u64,
606}
607
608/// Retail Price Improvement order book data for an instrument.
609#[derive(Debug, Serialize, Deserialize)]
610#[serde(rename_all = "camelCase", deny_unknown_fields)]
611pub struct OKXRpiBookMsg {
612    /// Ask levels [price, total quantity, non-RPI quantity, order count].
613    pub asks: Vec<OKXRpiBookLevel>,
614    /// Bid levels [price, total quantity, non-RPI quantity, order count].
615    pub bids: Vec<OKXRpiBookLevel>,
616    /// Sequence ID of the previous message.
617    pub prev_seq_id: i64,
618    /// Sequence ID of the current message.
619    pub seq_id: u64,
620    /// Order book generation time, Unix timestamp in milliseconds.
621    #[serde(deserialize_with = "deserialize_string_to_u64")]
622    pub ts: u64,
623}
624
625/// Trade data for an instrument.
626#[derive(Debug, Serialize, Deserialize)]
627#[serde(rename_all = "camelCase")]
628pub struct OKXTradeMsg {
629    // Spread public trades (`sprd-public-trades`) key the instrument as `sprdId`
630    // and omit `count`; the actual instrument is resolved from the channel arg, so
631    // both fields are tolerated here and unused by parsing.
632    /// Instrument ID (`instId`, or `sprdId` for spread public trades).
633    #[serde(default, alias = "sprdId")]
634    pub inst_id: Ustr,
635    /// Trade ID.
636    pub trade_id: String,
637    /// Trade price.
638    pub px: String,
639    /// Trade size.
640    pub sz: String,
641    /// Trade direction (buy or sell).
642    pub side: OKXSide,
643    /// Count (absent on spread public trades).
644    #[serde(default)]
645    pub count: String,
646    /// Trade timestamp, Unix timestamp format in milliseconds.
647    #[serde(deserialize_with = "deserialize_string_to_u64")]
648    pub ts: u64,
649    /// Order source (0: normal, 1: RPI).
650    #[serde(default)]
651    pub source: Option<String>,
652    /// Sequence ID for trade events.
653    #[serde(default)]
654    pub seq_id: Option<u64>,
655}
656
657/// Funding rate data for perpetual swaps.
658#[derive(Debug, Serialize, Deserialize)]
659#[serde(rename_all = "camelCase")]
660pub struct OKXFundingRateMsg {
661    /// Instrument type.
662    #[serde(default)]
663    pub inst_type: Option<OKXInstrumentType>,
664    /// Instrument ID.
665    pub inst_id: Ustr,
666    /// Current funding rate.
667    pub funding_rate: Ustr,
668    /// Predicted next funding rate.
669    pub next_funding_rate: Ustr,
670    /// Minimum funding rate.
671    #[serde(default)]
672    pub min_funding_rate: Option<String>,
673    /// Maximum funding rate.
674    #[serde(default)]
675    pub max_funding_rate: Option<String>,
676    /// Settlement state.
677    #[serde(default)]
678    pub sett_state: OKXSettlementState,
679    /// Settlement funding rate.
680    #[serde(default)]
681    pub sett_funding_rate: Option<String>,
682    /// Current premium.
683    #[serde(default)]
684    pub premium: Option<String>,
685    /// Funding rate calculation method.
686    #[serde(default)]
687    pub method: Option<String>,
688    /// Funding time, Unix timestamp format in milliseconds.
689    #[serde(deserialize_with = "deserialize_string_to_u64")]
690    pub funding_time: u64,
691    /// Next funding time, Unix timestamp format in milliseconds (used to determine funding interval).
692    #[serde(deserialize_with = "deserialize_string_to_u64")]
693    pub next_funding_time: u64,
694    /// Message timestamp, Unix timestamp format in milliseconds.
695    #[serde(deserialize_with = "deserialize_string_to_u64")]
696    pub ts: u64,
697}
698
699/// Mark price data for perpetual swaps.
700#[derive(Debug, Serialize, Deserialize)]
701#[serde(rename_all = "camelCase")]
702pub struct OKXMarkPriceMsg {
703    /// Instrument ID.
704    pub inst_id: Ustr,
705    /// Current mark price.
706    pub mark_px: String,
707    /// Timestamp of the data generation, Unix timestamp format in milliseconds.
708    #[serde(deserialize_with = "deserialize_string_to_u64")]
709    pub ts: u64,
710}
711
712/// Index price data.
713#[derive(Debug, Serialize, Deserialize)]
714#[serde(rename_all = "camelCase")]
715pub struct OKXIndexPriceMsg {
716    /// Index name, e.g. "BTC-USD".
717    pub inst_id: Ustr,
718    /// Latest index price.
719    pub idx_px: String,
720    /// 24-hour highest price.
721    pub high24h: String,
722    /// 24-hour lowest price.
723    pub low24h: String,
724    /// 24-hour opening price.
725    pub open24h: String,
726    /// The opening price of the day (UTC 0).
727    pub sod_utc0: String,
728    /// The opening price of the day (UTC 8).
729    pub sod_utc8: String,
730    /// Timestamp of the data generation, Unix timestamp format in milliseconds.
731    #[serde(deserialize_with = "deserialize_string_to_u64")]
732    pub ts: u64,
733}
734
735/// Price limit data (upper and lower limits).
736#[derive(Debug, Serialize, Deserialize)]
737#[serde(rename_all = "camelCase")]
738pub struct OKXPriceLimitMsg {
739    /// Instrument ID.
740    pub inst_id: Ustr,
741    /// Buy limit price.
742    pub buy_lmt: String,
743    /// Sell limit price.
744    pub sell_lmt: String,
745    /// Timestamp of the data generation, Unix timestamp format in milliseconds.
746    #[serde(deserialize_with = "deserialize_string_to_u64")]
747    pub ts: u64,
748}
749
750/// Candlestick data for an instrument.
751#[derive(Debug, Serialize, Deserialize)]
752#[serde(rename_all = "camelCase")]
753pub struct OKXCandleMsg {
754    /// Candlestick timestamp, Unix timestamp format in milliseconds.
755    #[serde(deserialize_with = "deserialize_string_to_u64")]
756    pub ts: u64,
757    /// Opening price.
758    pub o: String,
759    /// Highest price.
760    pub h: String,
761    /// Lowest price.
762    pub l: String,
763    /// Closing price.
764    pub c: String,
765    /// Trading volume in contracts.
766    pub vol: String,
767    /// Trading volume in quote currency.
768    pub vol_ccy: String,
769    pub vol_ccy_quote: String,
770    /// Whether this is a completed candle.
771    pub confirm: OKXCandleConfirm,
772}
773
774/// Open interest data.
775#[derive(Debug, Serialize, Deserialize)]
776#[serde(rename_all = "camelCase")]
777pub struct OKXOpenInterestMsg {
778    /// Instrument ID.
779    pub inst_id: Ustr,
780    /// Open interest in contracts.
781    pub oi: String,
782    /// Open interest in quote currency.
783    pub oi_ccy: String,
784    /// Timestamp of the data generation, Unix timestamp format in milliseconds.
785    #[serde(deserialize_with = "deserialize_string_to_u64")]
786    pub ts: u64,
787}
788
789/// Option market data summary.
790#[derive(Debug, Serialize, Deserialize)]
791#[serde(rename_all = "camelCase")]
792pub struct OKXOptionSummaryMsg {
793    /// Instrument type.
794    #[serde(default)]
795    pub inst_type: Option<OKXInstrumentType>,
796    /// Instrument ID.
797    pub inst_id: Ustr,
798    /// Underlying.
799    pub uly: String,
800    /// Delta.
801    pub delta: String,
802    /// Gamma.
803    pub gamma: String,
804    /// Theta.
805    pub theta: String,
806    /// Vega.
807    pub vega: String,
808    /// Black-Scholes delta.
809    #[serde(alias = "deltaBS")]
810    pub delta_bs: String,
811    /// Black-Scholes gamma.
812    #[serde(alias = "gammaBS")]
813    pub gamma_bs: String,
814    /// Black-Scholes theta.
815    #[serde(alias = "thetaBS")]
816    pub theta_bs: String,
817    /// Black-Scholes vega.
818    #[serde(alias = "vegaBS")]
819    pub vega_bs: String,
820    /// Realized volatility.
821    pub real_vol: String,
822    /// Bid volatility.
823    pub bid_vol: String,
824    /// Ask volatility.
825    pub ask_vol: String,
826    /// Mark volatility.
827    pub mark_vol: String,
828    /// Leverage.
829    pub lever: String,
830    /// Forward price.
831    #[serde(default)]
832    pub fwd_px: Option<String>,
833    /// Mark price.
834    #[serde(default)]
835    pub mark_px: Option<String>,
836    /// Volatility level.
837    #[serde(default)]
838    pub vol_lv: Option<String>,
839    /// Timestamp of the data generation, Unix timestamp format in milliseconds.
840    #[serde(deserialize_with = "deserialize_string_to_u64")]
841    pub ts: u64,
842}
843
844/// Estimated delivery/exercise price data.
845#[derive(Debug, Serialize, Deserialize)]
846#[serde(rename_all = "camelCase")]
847pub struct OKXEstimatedPriceMsg {
848    /// Instrument ID.
849    pub inst_id: Ustr,
850    /// Estimated settlement price.
851    pub settle_px: String,
852    /// Timestamp of the data generation, Unix timestamp format in milliseconds.
853    #[serde(deserialize_with = "deserialize_string_to_u64")]
854    pub ts: u64,
855}
856
857/// Platform status updates.
858#[derive(Debug, Serialize, Deserialize)]
859#[serde(rename_all = "camelCase")]
860pub struct OKXStatusMsg {
861    /// System maintenance status.
862    pub title: Ustr,
863    /// Status type: planned or scheduled.
864    #[serde(rename = "type")]
865    pub status_type: Ustr,
866    /// System maintenance state: canceled, completed, pending, ongoing.
867    pub state: Ustr,
868    /// Expected completion timestamp.
869    pub end_time: Option<String>,
870    /// Planned start timestamp.
871    pub begin_time: Option<String>,
872    /// Service involved.
873    pub service_type: Option<Ustr>,
874    /// Reason for status change.
875    pub reason: Option<String>,
876    /// Timestamp of the data generation, Unix timestamp format in milliseconds.
877    #[serde(deserialize_with = "deserialize_string_to_u64")]
878    pub ts: u64,
879}
880
881pub use crate::common::models::OKXAttachedAlgoOrd;
882
883/// Liquidation risk warning pushed by the `liquidation-warning` channel.
884///
885/// OKX sends this when an isolated position, or all positions under cross
886/// margin, approach liquidation. It is a risk warning only: the position may
887/// already be liquidated by the time the message arrives.
888#[derive(Clone, Debug, Serialize, Deserialize)]
889#[serde(rename_all = "camelCase")]
890pub struct OKXLiquidationWarningMsg {
891    /// Instrument type.
892    pub inst_type: OKXInstrumentType,
893    /// Instrument family.
894    #[serde(default)]
895    pub inst_family: Option<Ustr>,
896    /// Instrument ID.
897    pub inst_id: Ustr,
898    /// Margin mode.
899    pub mgn_mode: OKXMarginMode,
900    /// Position ID.
901    #[serde(default)]
902    pub pos_id: Option<Ustr>,
903    /// Position side.
904    pub pos_side: OKXPositionSide,
905    /// Position quantity.
906    pub pos: String,
907    /// Position currency (margin positions only).
908    #[serde(default)]
909    pub pos_ccy: Option<Ustr>,
910    /// Leverage.
911    pub lever: String,
912    /// Mark price.
913    pub mark_px: String,
914    /// Maintenance margin ratio.
915    pub mgn_ratio: String,
916    /// Margin currency.
917    pub ccy: Ustr,
918    /// Creation time, Unix timestamp in milliseconds.
919    #[serde(deserialize_with = "deserialize_string_to_u64")]
920    pub c_time: u64,
921    /// Last update time, Unix timestamp in milliseconds.
922    #[serde(deserialize_with = "deserialize_string_to_u64")]
923    pub u_time: u64,
924    /// Push time, Unix timestamp in milliseconds.
925    #[serde(default)]
926    pub p_time: Option<String>,
927}
928
929/// Linked algo order metadata from order push updates.
930#[derive(Clone, Debug, Default, Serialize, Deserialize)]
931#[serde(rename_all = "camelCase")]
932pub struct OKXLinkedAlgoOrd {
933    /// Parent algo order ID.
934    #[serde(default)]
935    pub algo_id: String,
936}
937
938/// Order update message from WebSocket orders channel.
939#[derive(Clone, Debug, Serialize, Deserialize)]
940#[serde(rename_all = "camelCase")]
941pub struct OKXOrderMsg {
942    /// Accumulated filled size.
943    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
944    pub acc_fill_sz: Option<String>,
945    /// Algo order ID.
946    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
947    pub algo_id: Option<String>,
948    /// Average price.
949    pub avg_px: String,
950    /// Creation time, Unix timestamp in milliseconds.
951    #[serde(deserialize_with = "deserialize_string_to_u64")]
952    pub c_time: u64,
953    /// Cancel source.
954    #[serde(default)]
955    pub cancel_source: Option<String>,
956    /// Cancel source reason.
957    #[serde(default)]
958    pub cancel_source_reason: Option<String>,
959    /// Order category (normal, liquidation, ADL, etc.).
960    pub category: OKXOrderCategory,
961    /// Currency.
962    pub ccy: Ustr,
963    /// Client order ID.
964    pub cl_ord_id: String,
965    /// Parent algo client order ID if present.
966    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
967    pub algo_cl_ord_id: Option<String>,
968    /// Attached child client order ID if surfaced at the top level.
969    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
970    pub attach_algo_cl_ord_id: Option<String>,
971    /// Attached TP/SL child order metadata.
972    #[serde(default)]
973    pub attach_algo_ords: Vec<OKXAttachedAlgoOrd>,
974    /// Event contract market outcome, if applicable.
975    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
976    pub outcome: Option<String>,
977    /// Fee (cumulative).
978    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
979    pub fee: Option<String>,
980    /// Fee currency.
981    pub fee_ccy: Ustr,
982    /// Fee for this fill.
983    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
984    pub fill_fee: Option<String>,
985    /// Fill fee currency.
986    #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
987    pub fill_fee_ccy: Option<Ustr>,
988    /// Mark price at fill time.
989    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
990    pub fill_mark_px: Option<String>,
991    /// Mark volatility at fill time (options).
992    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
993    pub fill_mark_vol: Option<String>,
994    /// Implied volatility at fill time (options).
995    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
996    pub fill_px_vol: Option<String>,
997    /// Fill price in USD (options).
998    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
999    pub fill_px_usd: Option<String>,
1000    /// Forward price at fill time (options).
1001    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1002    pub fill_fwd_px: Option<String>,
1003    /// Fill notional in USD.
1004    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1005    pub fill_notional_usd: Option<String>,
1006    /// PnL for this fill.
1007    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1008    pub fill_pnl: Option<String>,
1009    /// Fill price.
1010    pub fill_px: String,
1011    /// Fill size.
1012    pub fill_sz: String,
1013    /// Fill time, Unix timestamp in milliseconds.
1014    #[serde(deserialize_with = "deserialize_string_to_u64")]
1015    pub fill_time: u64,
1016    /// Instrument ID.
1017    pub inst_id: Ustr,
1018    /// Instrument type.
1019    pub inst_type: OKXInstrumentType,
1020    /// Whether the TP order is a limit order.
1021    #[serde(default)]
1022    pub is_tp_limit: Option<String>,
1023    /// Leverage.
1024    pub lever: String,
1025    /// Linked algo order metadata.
1026    #[serde(default)]
1027    pub linked_algo_ord: Option<OKXLinkedAlgoOrd>,
1028    /// Notional value in USD.
1029    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1030    pub notional_usd: Option<String>,
1031    /// Order ID.
1032    pub ord_id: Ustr,
1033    /// Order type.
1034    pub ord_type: OKXOrderType,
1035    /// Profit and loss.
1036    pub pnl: String,
1037    /// Position side.
1038    pub pos_side: OKXPositionSide,
1039    /// Price (algo orders use ordPx instead).
1040    #[serde(default)]
1041    pub px: String,
1042    /// Price type (options).
1043    #[serde(default)]
1044    pub px_type: OKXPriceType,
1045    /// Price in USD (options).
1046    #[serde(default)]
1047    pub px_usd: Option<String>,
1048    /// Price in volatility (options).
1049    #[serde(default)]
1050    pub px_vol: Option<String>,
1051    /// Quick margin type.
1052    #[serde(default)]
1053    pub quick_mgn_type: OKXQuickMarginType,
1054    /// Rebate amount.
1055    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1056    pub rebate: Option<String>,
1057    /// Rebate currency.
1058    #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
1059    pub rebate_ccy: Option<Ustr>,
1060    /// Reduce only flag.
1061    pub reduce_only: String,
1062    /// Side.
1063    pub side: OKXSide,
1064    /// Stop-loss order price.
1065    #[serde(default)]
1066    pub sl_ord_px: Option<String>,
1067    /// Stop-loss trigger price.
1068    #[serde(default)]
1069    pub sl_trigger_px: Option<String>,
1070    /// Stop-loss trigger price type (last, mark, index).
1071    #[serde(default)]
1072    pub sl_trigger_px_type: Option<OKXTriggerType>,
1073    /// Order source.
1074    #[serde(default)]
1075    pub source: Option<String>,
1076    /// Order state.
1077    pub state: OKXOrderStatus,
1078    /// Self-trade prevention ID.
1079    #[serde(default)]
1080    pub stp_id: Option<String>,
1081    /// Self-trade prevention mode.
1082    #[serde(default)]
1083    pub stp_mode: OKXSelfTradePreventionMode,
1084    /// Execution type.
1085    pub exec_type: OKXExecType,
1086    /// Size.
1087    pub sz: String,
1088    /// Order tag.
1089    #[serde(default)]
1090    pub tag: Option<String>,
1091    /// Trade mode.
1092    pub td_mode: OKXTradeMode,
1093    /// Target currency (base_ccy or quote_ccy). Empty for margin modes.
1094    #[serde(default, deserialize_with = "deserialize_target_currency_as_none")]
1095    pub tgt_ccy: Option<OKXTargetCurrency>,
1096    /// Take-profit order price.
1097    #[serde(default)]
1098    pub tp_ord_px: Option<String>,
1099    /// Take-profit trigger price.
1100    #[serde(default)]
1101    pub tp_trigger_px: Option<String>,
1102    /// Take-profit trigger price type (last, mark, index).
1103    #[serde(default)]
1104    pub tp_trigger_px_type: Option<OKXTriggerType>,
1105    /// Trade ID.
1106    pub trade_id: String,
1107    /// Last update time, Unix timestamp in milliseconds.
1108    #[serde(deserialize_with = "deserialize_string_to_u64")]
1109    pub u_time: u64,
1110    /// Amend result code.
1111    #[serde(default)]
1112    pub amend_result: Option<String>,
1113    /// Request ID (for amend responses).
1114    #[serde(default)]
1115    pub req_id: Option<String>,
1116    /// Error code.
1117    #[serde(default)]
1118    pub code: Option<String>,
1119    /// Error message.
1120    #[serde(default)]
1121    pub msg: Option<String>,
1122}
1123
1124/// Represents an algo order message from WebSocket updates.
1125#[derive(Clone, Debug, Deserialize, Serialize)]
1126#[serde(rename_all = "camelCase")]
1127pub struct OKXAlgoOrderMsg {
1128    /// Algorithm ID.
1129    pub algo_id: String,
1130    /// Algorithm client order ID.
1131    #[serde(default)]
1132    pub algo_cl_ord_id: String,
1133    /// Client order ID (empty for algo orders until triggered).
1134    pub cl_ord_id: String,
1135    /// Order ID (empty until algo order is triggered).
1136    pub ord_id: String,
1137    /// Triggered child order IDs.
1138    #[serde(default)]
1139    pub ord_id_list: Vec<String>,
1140    /// Instrument ID.
1141    pub inst_id: Ustr,
1142    /// Instrument type.
1143    pub inst_type: OKXInstrumentType,
1144    /// Algo order type (trigger, move_order_stop, oco, iceberg, twap).
1145    pub ord_type: OKXAlgoOrderType,
1146    /// Order state.
1147    pub state: OKXAlgoOrderStatus,
1148    /// Side.
1149    pub side: OKXSide,
1150    /// Position side.
1151    pub pos_side: OKXPositionSide,
1152    /// Size.
1153    #[serde(default)]
1154    pub sz: String,
1155    /// Trigger price.
1156    #[serde(default)]
1157    pub trigger_px: String,
1158    /// Trigger price type (last, mark, index).
1159    #[serde(default)]
1160    pub trigger_px_type: OKXTriggerType,
1161    /// Stop-loss trigger price for conditional close orders.
1162    #[serde(default)]
1163    pub sl_trigger_px: String,
1164    /// Stop-loss order price for conditional close orders.
1165    #[serde(default)]
1166    pub sl_ord_px: String,
1167    /// Stop-loss trigger price type (last, mark, index).
1168    #[serde(default)]
1169    pub sl_trigger_px_type: OKXTriggerType,
1170    /// Take-profit trigger price for conditional close orders.
1171    #[serde(default)]
1172    pub tp_trigger_px: String,
1173    /// Take-profit order price for conditional close orders.
1174    #[serde(default)]
1175    pub tp_ord_px: String,
1176    /// Take-profit trigger price type (last, mark, index).
1177    #[serde(default)]
1178    pub tp_trigger_px_type: OKXTriggerType,
1179    /// Order price (-1 for market orders).
1180    #[serde(default)]
1181    pub ord_px: String,
1182    /// Trade mode.
1183    pub td_mode: OKXTradeMode,
1184    /// Leverage.
1185    pub lever: String,
1186    /// Reduce only flag.
1187    #[serde(default)]
1188    pub reduce_only: String,
1189    /// Fraction of the position to close for close-order algos.
1190    #[serde(default)]
1191    pub close_fraction: String,
1192    /// Actual filled price.
1193    #[serde(default)]
1194    pub actual_px: String,
1195    /// Actual filled size.
1196    #[serde(default)]
1197    pub actual_sz: String,
1198    /// Notional USD value.
1199    #[serde(default)]
1200    pub notional_usd: String,
1201    /// Creation time, Unix timestamp in milliseconds.
1202    #[serde(deserialize_with = "deserialize_string_to_u64")]
1203    pub c_time: u64,
1204    /// Update time, Unix timestamp in milliseconds.
1205    #[serde(deserialize_with = "deserialize_string_to_u64")]
1206    pub u_time: u64,
1207    /// Trigger time (empty until triggered).
1208    #[serde(default)]
1209    pub trigger_time: String,
1210    /// Failure code for rejected algo orders.
1211    #[serde(default)]
1212    pub fail_code: String,
1213    /// Tag.
1214    #[serde(default)]
1215    pub tag: String,
1216    /// Callback price ratio for trailing stop (e.g. "0.01" for 1%).
1217    #[serde(default)]
1218    pub callback_ratio: String,
1219    /// Callback price spread for trailing stop (absolute distance).
1220    #[serde(default)]
1221    pub callback_spread: String,
1222    /// Activation price for trailing stop.
1223    #[serde(default)]
1224    pub active_px: String,
1225    /// Currency.
1226    #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
1227    pub ccy: Option<Ustr>,
1228    /// Target currency (base_ccy or quote_ccy).
1229    #[serde(default, deserialize_with = "deserialize_target_currency_as_none")]
1230    pub tgt_ccy: Option<OKXTargetCurrency>,
1231    /// Fee amount.
1232    #[serde(default)]
1233    pub fee: Option<String>,
1234    /// Fee currency.
1235    #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
1236    pub fee_ccy: Option<Ustr>,
1237    /// Trigger order type (fok, ioc).
1238    #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1239    pub advance_ord_type: Option<String>,
1240}
1241
1242/// Parameters for WebSocket place order operation.
1243#[derive(Clone, Debug, Default, Deserialize, Serialize, Builder)]
1244#[builder(default)]
1245#[builder(setter(into, strip_option))]
1246#[serde(rename_all = "camelCase")]
1247pub struct WsAttachAlgoOrdParams {
1248    /// Attached algo client order ID.
1249    #[serde(skip_serializing_if = "Option::is_none")]
1250    pub attach_algo_cl_ord_id: Option<String>,
1251    /// Stop-loss trigger price.
1252    #[serde(skip_serializing_if = "Option::is_none")]
1253    pub sl_trigger_px: Option<String>,
1254    /// Stop-loss order price (`-1` for market).
1255    #[serde(skip_serializing_if = "Option::is_none")]
1256    pub sl_ord_px: Option<String>,
1257    /// Stop-loss trigger price type (last, mark, index).
1258    #[serde(skip_serializing_if = "Option::is_none")]
1259    pub sl_trigger_px_type: Option<OKXTriggerType>,
1260    /// Take-profit trigger price.
1261    #[serde(skip_serializing_if = "Option::is_none")]
1262    pub tp_trigger_px: Option<String>,
1263    /// Take-profit order price (`-1` for market).
1264    #[serde(skip_serializing_if = "Option::is_none")]
1265    pub tp_ord_px: Option<String>,
1266    /// Take-profit trigger price type (last, mark, index).
1267    #[serde(skip_serializing_if = "Option::is_none")]
1268    pub tp_trigger_px_type: Option<OKXTriggerType>,
1269    /// Callback ratio for attached trailing stop orders.
1270    #[serde(skip_serializing_if = "Option::is_none")]
1271    pub callback_ratio: Option<String>,
1272    /// Callback spread for attached trailing stop orders.
1273    #[serde(skip_serializing_if = "Option::is_none")]
1274    pub callback_spread: Option<String>,
1275    /// Activation price for attached trailing stop orders.
1276    #[serde(skip_serializing_if = "Option::is_none")]
1277    pub active_px: Option<String>,
1278    /// New callback ratio for amended attached trailing stop orders.
1279    #[serde(skip_serializing_if = "Option::is_none")]
1280    pub new_callback_ratio: Option<String>,
1281    /// New callback spread for amended attached trailing stop orders.
1282    #[serde(skip_serializing_if = "Option::is_none")]
1283    pub new_callback_spread: Option<String>,
1284    /// New activation price for amended attached trailing stop orders.
1285    #[serde(skip_serializing_if = "Option::is_none")]
1286    pub new_active_px: Option<String>,
1287}
1288
1289/// Parameters for WebSocket place order operation.
1290#[derive(Clone, Debug, Deserialize, Serialize, Builder)]
1291#[builder(setter(into, strip_option))]
1292#[serde(rename_all = "camelCase")]
1293pub struct WsPostOrderParams {
1294    /// Instrument type: SPOT, MARGIN, SWAP, FUTURES, OPTION (optional for WebSocket).
1295    #[builder(default)]
1296    #[serde(skip_serializing_if = "Option::is_none")]
1297    pub inst_type: Option<OKXInstrumentType>,
1298    /// Instrument ID code (numeric). Replaced `instId` for WebSocket order operations.
1299    pub inst_id_code: u64,
1300    /// Trading mode: cash, isolated, cross.
1301    pub td_mode: OKXTradeMode,
1302    /// Margin currency (only for isolated margin).
1303    #[builder(default)]
1304    #[serde(skip_serializing_if = "Option::is_none")]
1305    pub ccy: Option<Ustr>,
1306    /// Unique client order ID.
1307    #[builder(default)]
1308    #[serde(skip_serializing_if = "Option::is_none")]
1309    pub cl_ord_id: Option<String>,
1310    /// Order side: buy or sell.
1311    pub side: OKXSide,
1312    /// Position side: long, short, net (optional).
1313    #[builder(default)]
1314    #[serde(skip_serializing_if = "Option::is_none")]
1315    pub pos_side: Option<OKXPositionSide>,
1316    /// Order type: limit, market, post_only, fok, ioc, etc.
1317    pub ord_type: OKXOrderType,
1318    /// Order size.
1319    pub sz: String,
1320    /// Order price (required for limit orders).
1321    #[builder(default)]
1322    #[serde(skip_serializing_if = "Option::is_none")]
1323    pub px: Option<String>,
1324    /// Price in USD, only applicable to options. Mutually exclusive with `px` and `px_vol`.
1325    #[builder(default)]
1326    #[serde(rename = "pxUsd", skip_serializing_if = "Option::is_none")]
1327    pub px_usd: Option<String>,
1328    /// Price in implied volatility (1 = 100%), only applicable to options.
1329    /// Mutually exclusive with `px` and `px_usd`.
1330    #[builder(default)]
1331    #[serde(rename = "pxVol", skip_serializing_if = "Option::is_none")]
1332    pub px_vol: Option<String>,
1333    /// Reduce-only flag.
1334    #[builder(default)]
1335    #[serde(skip_serializing_if = "Option::is_none")]
1336    pub reduce_only: Option<bool>,
1337    /// Whether to close the entire position.
1338    #[builder(default)]
1339    #[serde(rename = "closePosition", skip_serializing_if = "Option::is_none")]
1340    pub close_position: Option<bool>,
1341    /// Target currency for net orders.
1342    #[builder(default)]
1343    #[serde(skip_serializing_if = "Option::is_none")]
1344    pub tgt_ccy: Option<OKXTargetCurrency>,
1345    /// Order tag for categorization.
1346    #[builder(default)]
1347    #[serde(skip_serializing_if = "Option::is_none")]
1348    pub tag: Option<String>,
1349    /// Attached TP/SL orders submitted with the parent order.
1350    #[builder(default)]
1351    #[serde(skip_serializing_if = "Option::is_none")]
1352    pub attach_algo_ords: Option<Vec<WsAttachAlgoOrdParams>>,
1353    /// Event contract speed bump flag. Use "1" for non-post-only EVENTS orders.
1354    #[builder(default)]
1355    #[serde(skip_serializing_if = "Option::is_none")]
1356    pub speed_bump: Option<String>,
1357    /// Event contract market outcome: yes or no.
1358    #[builder(default)]
1359    #[serde(skip_serializing_if = "Option::is_none")]
1360    pub outcome: Option<String>,
1361    /// Slippage tolerance for market orders, expressed as a decimal fraction
1362    /// (e.g., "0.005" for 0.5%). Supported instrument/order-type scope is
1363    /// venue-controlled; rejected with `54084`/`54085` if exceeded or out of
1364    /// the venue's accepted range. See the OKX v5 docs for the current matrix.
1365    #[builder(default)]
1366    #[serde(skip_serializing_if = "Option::is_none")]
1367    pub slippage_pct: Option<String>,
1368    /// Whether the order may take RPI liquidity.
1369    #[builder(default)]
1370    #[serde(skip_serializing_if = "Option::is_none")]
1371    pub rpi_taker_access: Option<bool>,
1372    /// Whether OKX may round the order price to an eligible RPI price.
1373    #[builder(default)]
1374    #[serde(skip_serializing_if = "Option::is_none")]
1375    pub rpi_px_round: Option<bool>,
1376}
1377
1378/// Parameters for WebSocket cancel order operation (instType not included).
1379#[derive(Clone, Debug, Default, Deserialize, Serialize, Builder)]
1380#[builder(default)]
1381#[builder(setter(into, strip_option))]
1382#[serde(rename_all = "camelCase")]
1383pub struct WsCancelOrderParams {
1384    /// Instrument ID code (numeric). Replaced `instId` for WebSocket order operations.
1385    pub inst_id_code: u64,
1386    /// Exchange-assigned order ID.
1387    #[serde(skip_serializing_if = "Option::is_none")]
1388    pub ord_id: Option<String>,
1389    /// User-assigned client order ID.
1390    #[serde(skip_serializing_if = "Option::is_none")]
1391    pub cl_ord_id: Option<String>,
1392}
1393
1394/// Parameters for WebSocket mass cancel operation.
1395#[derive(Clone, Debug, Default, Deserialize, Serialize, Builder)]
1396#[builder(default)]
1397#[builder(setter(into, strip_option))]
1398#[serde(rename_all = "camelCase")]
1399pub struct WsMassCancelParams {
1400    /// Instrument type.
1401    pub inst_type: OKXInstrumentType,
1402    /// Instrument family, e.g. "BTC-USD", "BTC-USDT".
1403    pub inst_family: Ustr,
1404}
1405
1406/// Parameters for WebSocket amend order operation (instType not included).
1407#[derive(Clone, Debug, Default, Deserialize, Serialize, Builder)]
1408#[builder(default)]
1409#[builder(setter(into, strip_option))]
1410#[serde(rename_all = "camelCase")]
1411pub struct WsAmendOrderParams {
1412    /// Instrument ID code (numeric). Replaced `instId` for WebSocket order operations.
1413    pub inst_id_code: u64,
1414    /// Exchange-assigned order ID (optional if using clOrdId).
1415    #[serde(skip_serializing_if = "Option::is_none")]
1416    pub ord_id: Option<String>,
1417    /// User-assigned client order ID (optional if using ordId).
1418    #[serde(skip_serializing_if = "Option::is_none")]
1419    pub cl_ord_id: Option<String>,
1420    /// Client request ID for correlating the amendment response.
1421    #[serde(skip_serializing_if = "Option::is_none")]
1422    pub req_id: Option<String>,
1423    /// New order price (optional).
1424    #[serde(skip_serializing_if = "Option::is_none")]
1425    pub new_px: Option<String>,
1426    /// New price in USD, only applicable to options. Must match the pricing mode used at placement.
1427    #[serde(rename = "newPxUsd", skip_serializing_if = "Option::is_none")]
1428    pub new_px_usd: Option<String>,
1429    /// New price in implied volatility, only applicable to options.
1430    /// Must match the pricing mode used at placement.
1431    #[serde(rename = "newPxVol", skip_serializing_if = "Option::is_none")]
1432    pub new_px_vol: Option<String>,
1433    /// New order size (optional).
1434    #[serde(skip_serializing_if = "Option::is_none")]
1435    pub new_sz: Option<String>,
1436    /// Whether the order may take RPI liquidity after amendment.
1437    #[serde(skip_serializing_if = "Option::is_none")]
1438    pub rpi_taker_access: Option<bool>,
1439    /// Whether OKX may round the amended price to an eligible RPI price.
1440    #[serde(skip_serializing_if = "Option::is_none")]
1441    pub rpi_px_round: Option<bool>,
1442    /// Event contract speed bump flag. Use "1" for non-post-only EVENTS orders.
1443    #[serde(skip_serializing_if = "Option::is_none")]
1444    pub speed_bump: Option<String>,
1445}
1446
1447/// Parameters for WebSocket algo order placement.
1448#[derive(Clone, Debug, Deserialize, Serialize, Builder)]
1449#[builder(setter(into, strip_option))]
1450#[serde(rename_all = "camelCase")]
1451pub struct WsPostAlgoOrderParams {
1452    /// Instrument ID code (numeric). Replaced `instId` for WebSocket order operations.
1453    pub inst_id_code: u64,
1454    /// Trading mode: cash, isolated, cross.
1455    pub td_mode: OKXTradeMode,
1456    /// Order side: buy or sell.
1457    pub side: OKXSide,
1458    /// Order type: trigger (for stop orders).
1459    pub ord_type: OKXAlgoOrderType,
1460    /// Order size.
1461    pub sz: String,
1462    /// Client order ID (optional).
1463    #[builder(default)]
1464    #[serde(skip_serializing_if = "Option::is_none")]
1465    pub cl_ord_id: Option<String>,
1466    /// Position side: long, short, net (optional).
1467    #[builder(default)]
1468    #[serde(skip_serializing_if = "Option::is_none")]
1469    pub pos_side: Option<OKXPositionSide>,
1470    /// Trigger price for stop/conditional orders.
1471    #[serde(skip_serializing_if = "Option::is_none")]
1472    pub trigger_px: Option<String>,
1473    /// Trigger price type: last, index, mark.
1474    #[builder(default)]
1475    #[serde(skip_serializing_if = "Option::is_none")]
1476    pub trigger_px_type: Option<OKXTriggerType>,
1477    /// Order price (for limit orders after trigger).
1478    #[builder(default)]
1479    #[serde(skip_serializing_if = "Option::is_none")]
1480    pub order_px: Option<String>,
1481    /// Reduce-only flag.
1482    #[builder(default)]
1483    #[serde(skip_serializing_if = "Option::is_none")]
1484    pub reduce_only: Option<bool>,
1485    /// Order tag for categorization.
1486    #[builder(default)]
1487    #[serde(skip_serializing_if = "Option::is_none")]
1488    pub tag: Option<String>,
1489    /// Callback rate for trailing stop (e.g., "0.01" for 1%).
1490    #[builder(default)]
1491    #[serde(skip_serializing_if = "Option::is_none")]
1492    pub callback_ratio: Option<String>,
1493    /// Callback spread for trailing stop (fixed price distance).
1494    #[builder(default)]
1495    #[serde(skip_serializing_if = "Option::is_none")]
1496    pub callback_spread: Option<String>,
1497    /// Activation price for trailing stop.
1498    #[builder(default)]
1499    #[serde(skip_serializing_if = "Option::is_none")]
1500    pub active_px: Option<String>,
1501}
1502
1503/// Parameters for WebSocket cancel algo order operation.
1504#[derive(Clone, Debug, Deserialize, Serialize, Builder)]
1505#[builder(setter(into, strip_option))]
1506#[serde(rename_all = "camelCase")]
1507pub struct WsCancelAlgoOrderParams {
1508    /// Instrument ID code (numeric). Replaced `instId` for WebSocket order operations.
1509    pub inst_id_code: u64,
1510    /// Algo order ID.
1511    #[serde(skip_serializing_if = "Option::is_none")]
1512    pub algo_id: Option<String>,
1513    /// Client algo order ID.
1514    #[serde(skip_serializing_if = "Option::is_none")]
1515    pub algo_cl_ord_id: Option<String>,
1516}
1517
1518#[cfg(test)]
1519mod tests {
1520    use nautilus_core::time::get_atomic_clock_realtime;
1521    use rstest::rstest;
1522    use rust_decimal::Decimal;
1523
1524    use super::*;
1525    use crate::common::testing::load_test_json;
1526
1527    #[rstest]
1528    fn test_deserialize_websocket_arg() {
1529        let json_str = r#"{"channel":"instruments","instType":"SPOT"}"#;
1530
1531        let result: Result<OKXWebSocketArg, _> = serde_json::from_str(json_str);
1532        match result {
1533            Ok(arg) => {
1534                assert_eq!(arg.channel, OKXWsChannel::Instruments);
1535                assert_eq!(arg.inst_type, Some(OKXInstrumentType::Spot));
1536                assert_eq!(arg.inst_id, None);
1537            }
1538            Err(e) => {
1539                panic!("Failed to deserialize WebSocket arg: {e}");
1540            }
1541        }
1542    }
1543
1544    #[rstest]
1545    fn test_deserialize_subscribe_variant_direct() {
1546        #[derive(Debug, Deserialize)]
1547        #[serde(rename_all = "camelCase")]
1548        struct SubscribeMsg {
1549            event: String,
1550            arg: OKXWebSocketArg,
1551            conn_id: String,
1552        }
1553
1554        let json_str = r#"{"event":"subscribe","arg":{"channel":"instruments","instType":"SPOT"},"connId":"380cfa6a"}"#;
1555
1556        let result: Result<SubscribeMsg, _> = serde_json::from_str(json_str);
1557        match result {
1558            Ok(msg) => {
1559                assert_eq!(msg.event, "subscribe");
1560                assert_eq!(msg.arg.channel, OKXWsChannel::Instruments);
1561                assert_eq!(msg.conn_id, "380cfa6a");
1562            }
1563            Err(e) => {
1564                panic!("Failed to deserialize subscribe message directly: {e}");
1565            }
1566        }
1567    }
1568
1569    #[rstest]
1570    fn test_deserialize_subscribe_confirmation() {
1571        let json_str = r#"{"event":"subscribe","arg":{"channel":"instruments","instType":"SPOT"},"connId":"380cfa6a"}"#;
1572
1573        let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1574        match result {
1575            Ok(msg) => {
1576                if let OKXWsFrame::Subscription {
1577                    event,
1578                    arg,
1579                    conn_id,
1580                    ..
1581                } = msg
1582                {
1583                    assert_eq!(event, OKXSubscriptionEvent::Subscribe);
1584                    assert_eq!(arg.channel, OKXWsChannel::Instruments);
1585                    assert_eq!(conn_id, "380cfa6a");
1586                } else {
1587                    panic!("Expected Subscribe variant, was: {msg:?}");
1588                }
1589            }
1590            Err(e) => {
1591                panic!("Failed to deserialize subscription confirmation: {e}");
1592            }
1593        }
1594    }
1595
1596    #[rstest]
1597    fn test_deserialize_subscribe_with_inst_id() {
1598        let json_str = r#"{"event":"subscribe","arg":{"channel":"candle1m","instId":"ETH-USDT"},"connId":"358602f5"}"#;
1599
1600        let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1601        match result {
1602            Ok(msg) => {
1603                if let OKXWsFrame::Subscription {
1604                    event,
1605                    arg,
1606                    conn_id,
1607                    ..
1608                } = msg
1609                {
1610                    assert_eq!(event, OKXSubscriptionEvent::Subscribe);
1611                    assert_eq!(arg.channel, OKXWsChannel::Candle1Minute);
1612                    assert_eq!(conn_id, "358602f5");
1613                } else {
1614                    panic!("Expected Subscribe variant, was: {msg:?}");
1615                }
1616            }
1617            Err(e) => {
1618                panic!("Failed to deserialize subscription confirmation: {e}");
1619            }
1620        }
1621    }
1622
1623    #[rstest]
1624    fn test_channel_serialization_for_logging() {
1625        let channel = OKXWsChannel::Candle1Minute;
1626        let serialized = serde_json::to_string(&channel).unwrap();
1627        let cleaned = serialized.trim_matches('"').to_string();
1628        assert_eq!(cleaned, "candle1m");
1629
1630        let channel = OKXWsChannel::BboTbt;
1631        let serialized = serde_json::to_string(&channel).unwrap();
1632        let cleaned = serialized.trim_matches('"').to_string();
1633        assert_eq!(cleaned, "bbo-tbt");
1634
1635        let channel = OKXWsChannel::Trades;
1636        let serialized = serde_json::to_string(&channel).unwrap();
1637        let cleaned = serialized.trim_matches('"').to_string();
1638        assert_eq!(cleaned, "trades");
1639    }
1640
1641    #[rstest]
1642    fn test_order_response_with_enum_operation() {
1643        let json_str = r#"{"id":"req-123","op":"order","code":"0","msg":"","data":[]}"#;
1644        let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1645        match result {
1646            Ok(OKXWsFrame::OrderResponse {
1647                id,
1648                op,
1649                code,
1650                msg,
1651                data,
1652            }) => {
1653                assert_eq!(id, Some("req-123".to_string()));
1654                assert_eq!(op, OKXWsOperation::Order);
1655                assert_eq!(code, "0");
1656                assert_eq!(msg, "");
1657                assert!(data.is_empty());
1658            }
1659            Ok(other) => panic!("Expected OrderResponse, was: {other:?}"),
1660            Err(e) => panic!("Failed to deserialize: {e}"),
1661        }
1662
1663        let json_str = r#"{"id":"cancel-456","op":"cancel-order","code":"50001","msg":"Order not found","data":[]}"#;
1664        let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1665        match result {
1666            Ok(OKXWsFrame::OrderResponse {
1667                id,
1668                op,
1669                code,
1670                msg,
1671                data,
1672            }) => {
1673                assert_eq!(id, Some("cancel-456".to_string()));
1674                assert_eq!(op, OKXWsOperation::CancelOrder);
1675                assert_eq!(code, "50001");
1676                assert_eq!(msg, "Order not found");
1677                assert!(data.is_empty());
1678            }
1679            Ok(other) => panic!("Expected OrderResponse, was: {other:?}"),
1680            Err(e) => panic!("Failed to deserialize: {e}"),
1681        }
1682
1683        let json_str = r#"{"id":"amend-789","op":"amend-order","code":"50002","msg":"Invalid price","data":[]}"#;
1684        let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1685        match result {
1686            Ok(OKXWsFrame::OrderResponse {
1687                id,
1688                op,
1689                code,
1690                msg,
1691                data,
1692            }) => {
1693                assert_eq!(id, Some("amend-789".to_string()));
1694                assert_eq!(op, OKXWsOperation::AmendOrder);
1695                assert_eq!(code, "50002");
1696                assert_eq!(msg, "Invalid price");
1697                assert!(data.is_empty());
1698            }
1699            Ok(other) => panic!("Expected OrderResponse, was: {other:?}"),
1700            Err(e) => panic!("Failed to deserialize: {e}"),
1701        }
1702    }
1703
1704    #[rstest]
1705    fn test_operation_enum_serialization() {
1706        let op = OKXWsOperation::Order;
1707        let serialized = serde_json::to_string(&op).unwrap();
1708        assert_eq!(serialized, "\"order\"");
1709
1710        let op = OKXWsOperation::CancelOrder;
1711        let serialized = serde_json::to_string(&op).unwrap();
1712        assert_eq!(serialized, "\"cancel-order\"");
1713
1714        let op = OKXWsOperation::AmendOrder;
1715        let serialized = serde_json::to_string(&op).unwrap();
1716        assert_eq!(serialized, "\"amend-order\"");
1717
1718        let op = OKXWsOperation::Subscribe;
1719        let serialized = serde_json::to_string(&op).unwrap();
1720        assert_eq!(serialized, "\"subscribe\"");
1721    }
1722
1723    #[rstest]
1724    fn test_order_response_parsing() {
1725        let success_response = r#"{
1726            "id": "req-123",
1727            "op": "order",
1728            "code": "0",
1729            "msg": "",
1730            "data": [{"sMsg": "Order placed successfully"}]
1731        }"#;
1732
1733        let parsed: OKXWsFrame = serde_json::from_str(success_response).unwrap();
1734
1735        match parsed {
1736            OKXWsFrame::OrderResponse {
1737                id,
1738                op,
1739                code,
1740                msg,
1741                data,
1742            } => {
1743                assert_eq!(id, Some("req-123".to_string()));
1744                assert_eq!(op, OKXWsOperation::Order);
1745                assert_eq!(code, "0");
1746                assert_eq!(msg, "");
1747                assert_eq!(data.len(), 1);
1748            }
1749            _ => panic!("Expected OrderResponse variant"),
1750        }
1751
1752        let failure_response = r#"{
1753            "id": "req-456",
1754            "op": "cancel-order",
1755            "code": "50001",
1756            "msg": "Order not found",
1757            "data": [{"sMsg": "Order with client order ID not found"}]
1758        }"#;
1759
1760        let parsed: OKXWsFrame = serde_json::from_str(failure_response).unwrap();
1761
1762        match parsed {
1763            OKXWsFrame::OrderResponse {
1764                id,
1765                op,
1766                code,
1767                msg,
1768                data,
1769            } => {
1770                assert_eq!(id, Some("req-456".to_string()));
1771                assert_eq!(op, OKXWsOperation::CancelOrder);
1772                assert_eq!(code, "50001");
1773                assert_eq!(msg, "Order not found");
1774                assert_eq!(data.len(), 1);
1775            }
1776            _ => panic!("Expected OrderResponse variant"),
1777        }
1778    }
1779
1780    #[rstest]
1781    fn test_subscription_event_parsing() {
1782        let subscription_json = r#"{
1783            "event": "subscribe",
1784            "arg": {
1785                "channel": "tickers",
1786                "instId": "BTC-USDT"
1787            },
1788            "connId": "a4d3ae55"
1789        }"#;
1790
1791        let parsed: OKXWsFrame = serde_json::from_str(subscription_json).unwrap();
1792
1793        match parsed {
1794            OKXWsFrame::Subscription {
1795                event,
1796                arg,
1797                conn_id,
1798                ..
1799            } => {
1800                assert_eq!(
1801                    event,
1802                    crate::websocket::enums::OKXSubscriptionEvent::Subscribe
1803                );
1804                assert_eq!(arg.channel, OKXWsChannel::Tickers);
1805                assert_eq!(arg.inst_id, Some(Ustr::from("BTC-USDT")));
1806                assert_eq!(conn_id, "a4d3ae55");
1807            }
1808            _ => panic!("Expected Subscription variant"),
1809        }
1810    }
1811
1812    #[rstest]
1813    fn test_login_event_parsing() {
1814        let login_success = r#"{
1815            "event": "login",
1816            "code": "0",
1817            "msg": "Login successful",
1818            "connId": "a4d3ae55"
1819        }"#;
1820
1821        let parsed: OKXWsFrame = serde_json::from_str(login_success).unwrap();
1822
1823        match parsed {
1824            OKXWsFrame::Login {
1825                event,
1826                code,
1827                msg,
1828                conn_id,
1829            } => {
1830                assert_eq!(event, "login");
1831                assert_eq!(code, "0");
1832                assert_eq!(msg, "Login successful");
1833                assert_eq!(conn_id, "a4d3ae55");
1834            }
1835            _ => panic!("Expected Login variant, was: {parsed:?}"),
1836        }
1837    }
1838
1839    #[rstest]
1840    fn test_error_event_parsing() {
1841        let error_json = r#"{
1842            "code": "60012",
1843            "msg": "Invalid request"
1844        }"#;
1845
1846        let parsed: OKXWsFrame = serde_json::from_str(error_json).unwrap();
1847
1848        match parsed {
1849            OKXWsFrame::Error { arg, code, msg } => {
1850                assert!(arg.is_none());
1851                assert_eq!(code, "60012");
1852                assert_eq!(msg, "Invalid request");
1853            }
1854            _ => panic!("Expected Error variant"),
1855        }
1856    }
1857
1858    #[rstest]
1859    fn test_error_event_with_event_field_parsing() {
1860        // OKX sends error events with "event":"error" field (e.g., login failures)
1861        let error_json = r#"{
1862            "event": "error",
1863            "code": "60018",
1864            "msg": "Invalid sign"
1865        }"#;
1866
1867        let parsed: OKXWsFrame = serde_json::from_str(error_json).unwrap();
1868
1869        match parsed {
1870            OKXWsFrame::Error { arg, code, msg } => {
1871                assert!(arg.is_none());
1872                assert_eq!(code, "60018");
1873                assert_eq!(msg, "Invalid sign");
1874            }
1875            _ => panic!("Expected Error variant, was: {parsed:?}"),
1876        }
1877    }
1878
1879    #[rstest]
1880    fn test_subscription_error_with_arg_field_parsing() {
1881        // OKX sends subscription errors with arg field (channel subscription failures)
1882        let error_json = r#"{
1883            "event": "error",
1884            "arg": {"channel": "tickers", "instId": "INVALID-INST"},
1885            "code": "60012",
1886            "msg": "Invalid request: channel not found",
1887            "connId": "a4d3ae55"
1888        }"#;
1889
1890        let parsed: OKXWsFrame = serde_json::from_str(error_json).unwrap();
1891
1892        match parsed {
1893            OKXWsFrame::Error { arg, code, msg } => {
1894                let arg = arg.expect("subscription error arg");
1895                assert_eq!(arg.channel, OKXWsChannel::Tickers);
1896                assert_eq!(arg.inst_id, Some(Ustr::from("INVALID-INST")));
1897                assert_eq!(code, "60012");
1898                assert_eq!(msg, "Invalid request: channel not found");
1899            }
1900            _ => panic!("Expected Error variant, was: {parsed:?}"),
1901        }
1902    }
1903
1904    #[rstest]
1905    fn test_websocket_request_serialization() {
1906        let request = OKXWsRequest {
1907            id: Some("req-123".to_string()),
1908            op: OKXWsOperation::Order,
1909            args: vec![serde_json::json!({
1910                "instId": "BTC-USDT",
1911                "tdMode": "cash",
1912                "side": "buy",
1913                "ordType": "market",
1914                "sz": "0.1"
1915            })],
1916            exp_time: None,
1917        };
1918
1919        let serialized = serde_json::to_string(&request).unwrap();
1920        let parsed: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1921
1922        assert_eq!(parsed["id"], "req-123");
1923        assert_eq!(parsed["op"], "order");
1924        assert!(parsed["args"].is_array());
1925        assert_eq!(parsed["args"].as_array().unwrap().len(), 1);
1926    }
1927
1928    #[rstest]
1929    fn test_subscription_request_serialization() {
1930        let subscription = OKXSubscription {
1931            op: OKXWsOperation::Subscribe,
1932            args: vec![OKXSubscriptionArg {
1933                channel: OKXWsChannel::Tickers,
1934                inst_type: Some(OKXInstrumentType::Spot),
1935                inst_family: None,
1936                inst_id: Some(Ustr::from("BTC-USDT")),
1937            }],
1938        };
1939
1940        let serialized = serde_json::to_string(&subscription).unwrap();
1941        let parsed: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1942
1943        assert_eq!(parsed["op"], "subscribe");
1944        assert!(parsed["args"].is_array());
1945        assert_eq!(parsed["args"][0]["channel"], "tickers");
1946        assert_eq!(parsed["args"][0]["instType"], "SPOT");
1947        assert_eq!(parsed["args"][0]["instId"], "BTC-USDT");
1948    }
1949
1950    #[rstest]
1951    fn test_error_message_extraction() {
1952        let responses = vec![
1953            (
1954                r#"{
1955                "id": "req-123",
1956                "op": "order",
1957                "code": "50001",
1958                "msg": "Order failed",
1959                "data": [{"sMsg": "Insufficient balance"}]
1960            }"#,
1961                "Insufficient balance",
1962            ),
1963            (
1964                r#"{
1965                "id": "req-456",
1966                "op": "cancel-order",
1967                "code": "50002",
1968                "msg": "Cancel failed",
1969                "data": [{}]
1970            }"#,
1971                "Cancel failed",
1972            ),
1973        ];
1974
1975        for (response_json, expected_msg) in responses {
1976            let parsed: OKXWsFrame = serde_json::from_str(response_json).unwrap();
1977
1978            match parsed {
1979                OKXWsFrame::OrderResponse {
1980                    id: _,
1981                    op: _,
1982                    code,
1983                    msg,
1984                    data,
1985                } => {
1986                    assert_ne!(code, "0"); // Error response
1987
1988                    // Extract error message with fallback logic
1989                    let error_msg = data
1990                        .first()
1991                        .and_then(|d| d.get("sMsg"))
1992                        .and_then(|s| s.as_str())
1993                        .filter(|s| !s.is_empty())
1994                        .unwrap_or(&msg);
1995
1996                    assert_eq!(error_msg, expected_msg);
1997                }
1998                _ => panic!("Expected OrderResponse variant"),
1999            }
2000        }
2001    }
2002
2003    #[rstest]
2004    fn test_book_data_parsing() {
2005        let book_data_json = r#"{
2006            "arg": {
2007                "channel": "books",
2008                "instId": "BTC-USDT"
2009            },
2010            "action": "snapshot",
2011            "data": [{
2012                "asks": [["50000.0", "0.1", "0", "1"]],
2013                "bids": [["49999.0", "0.2", "0", "1"]],
2014                "ts": "1640995200000",
2015                "checksum": 123456789,
2016                "seqId": 1000
2017            }]
2018        }"#;
2019
2020        let parsed: OKXWsFrame = serde_json::from_str(book_data_json).unwrap();
2021
2022        match parsed {
2023            OKXWsFrame::BookData { arg, action, data } => {
2024                assert_eq!(arg.channel, OKXWsChannel::Books);
2025                assert_eq!(arg.inst_id, Some(Ustr::from("BTC-USDT")));
2026                assert_eq!(
2027                    action,
2028                    super::super::super::common::enums::OKXBookAction::Snapshot
2029                );
2030                assert_eq!(data.len(), 1);
2031            }
2032            _ => panic!("Expected BookData variant"),
2033        }
2034    }
2035
2036    #[rstest]
2037    fn test_rpi_book_fixtures_preserve_depth_types_and_sequence() {
2038        let snapshot: OKXWsFrame =
2039            serde_json::from_str(&load_test_json("ws_books_rpi_snapshot.json")).unwrap();
2040        let update: OKXWsFrame =
2041            serde_json::from_str(&load_test_json("ws_books_rpi_update.json")).unwrap();
2042
2043        let OKXWsFrame::RpiBookData { arg, action, data } = snapshot else {
2044            panic!("Expected RPI book snapshot");
2045        };
2046        let snapshot = &data[0];
2047        assert_eq!(arg.channel, OKXWsChannel::BooksRpi);
2048        assert_eq!(arg.inst_id, Some(Ustr::from("OMI-USD")));
2049        assert_eq!(action, OKXBookAction::Snapshot);
2050        assert_eq!(data.len(), 1);
2051        assert_eq!(snapshot.asks.len(), 4);
2052        assert_eq!(snapshot.bids.len(), 10);
2053        assert_eq!(
2054            snapshot.asks[0],
2055            OKXRpiBookLevel(
2056                Decimal::from_str_exact("0.0001617").unwrap(),
2057                Decimal::from_str_exact("12325166.992").unwrap(),
2058                Decimal::from(1000),
2059                2,
2060            )
2061        );
2062        assert_eq!(snapshot.prev_seq_id, -1);
2063        assert_eq!(snapshot.seq_id, 1_082_831_226);
2064        assert_eq!(snapshot.ts, 1_785_406_442_403);
2065
2066        let OKXWsFrame::RpiBookData { arg, action, data } = update else {
2067            panic!("Expected RPI book update");
2068        };
2069        let update = &data[0];
2070        assert_eq!(arg.channel, OKXWsChannel::BooksRpi);
2071        assert_eq!(arg.inst_id, Some(Ustr::from("OMI-USD")));
2072        assert_eq!(action, OKXBookAction::Update);
2073        assert_eq!(data.len(), 1);
2074        assert_eq!(update.asks.len(), 2);
2075        assert!(update.bids.is_empty());
2076        assert_eq!(
2077            update.asks[1],
2078            OKXRpiBookLevel(
2079                Decimal::from_str_exact("0.0001625").unwrap(),
2080                Decimal::from_str_exact("12324367.786").unwrap(),
2081                Decimal::from(1000),
2082                2,
2083            )
2084        );
2085        assert_eq!(update.prev_seq_id, snapshot.seq_id as i64);
2086        assert_eq!(update.seq_id, 1_082_831_230);
2087        assert_eq!(update.ts, 1_785_406_443_903);
2088    }
2089
2090    #[rstest]
2091    fn test_rpi_book_rejects_checksum_field() {
2092        let mut payload: serde_json::Value =
2093            serde_json::from_str(&load_test_json("ws_books_rpi_update.json")).unwrap();
2094        payload["data"][0]["checksum"] = serde_json::json!(0);
2095
2096        let error = serde_json::from_value::<OKXWsFrame>(payload).unwrap_err();
2097
2098        assert!(error.to_string().contains("checksum"));
2099    }
2100
2101    #[rstest]
2102    fn test_data_event_parsing() {
2103        let data_json = r#"{
2104            "arg": {
2105                "channel": "trades",
2106                "instId": "BTC-USDT"
2107            },
2108            "data": [{
2109                "instId": "BTC-USDT",
2110                "tradeId": "12345",
2111                "px": "50000.0",
2112                "sz": "0.1",
2113                "side": "buy",
2114                "ts": "1640995200000"
2115            }]
2116        }"#;
2117
2118        let parsed: OKXWsFrame = serde_json::from_str(data_json).unwrap();
2119
2120        match parsed {
2121            OKXWsFrame::Data { arg, data } => {
2122                assert_eq!(arg.channel, OKXWsChannel::Trades);
2123                assert_eq!(arg.inst_id, Some(Ustr::from("BTC-USDT")));
2124                assert!(data.is_array());
2125            }
2126            _ => panic!("Expected Data variant"),
2127        }
2128    }
2129
2130    #[rstest]
2131    fn test_nautilus_message_variants() {
2132        let clock = get_atomic_clock_realtime();
2133        let ts_init = clock.get_time_ns();
2134
2135        let error = OKXWebSocketError {
2136            code: "60012".to_string(),
2137            message: "Invalid request".to_string(),
2138            conn_id: None,
2139            timestamp: ts_init.as_u64(),
2140        };
2141        let error_msg = NautilusWsMessage::Error(error);
2142
2143        match error_msg {
2144            NautilusWsMessage::Error(e) => {
2145                assert_eq!(e.code, "60012");
2146                assert_eq!(e.message, "Invalid request");
2147            }
2148            _ => panic!("Expected Error variant"),
2149        }
2150
2151        let raw_scenarios = vec![
2152            ::serde_json::json!({"unknown": "data"}),
2153            ::serde_json::json!({"channel": "unsupported", "data": [1, 2, 3]}),
2154            ::serde_json::json!({"complex": {"nested": {"structure": true}}}),
2155        ];
2156
2157        for raw_data in raw_scenarios {
2158            let raw_msg = NautilusWsMessage::Raw(raw_data.clone());
2159
2160            match raw_msg {
2161                NautilusWsMessage::Raw(data) => {
2162                    assert_eq!(data, raw_data);
2163                }
2164                _ => panic!("Expected Raw variant"),
2165            }
2166        }
2167    }
2168
2169    #[rstest]
2170    fn test_order_response_parsing_success() {
2171        let order_response_json = r#"{
2172            "id": "req-123",
2173            "op": "order",
2174            "code": "0",
2175            "msg": "",
2176            "data": [{"sMsg": "Order placed successfully"}]
2177        }"#;
2178
2179        let parsed: OKXWsFrame = serde_json::from_str(order_response_json).unwrap();
2180
2181        match parsed {
2182            OKXWsFrame::OrderResponse {
2183                id,
2184                op,
2185                code,
2186                msg,
2187                data,
2188            } => {
2189                assert_eq!(id, Some("req-123".to_string()));
2190                assert_eq!(op, OKXWsOperation::Order);
2191                assert_eq!(code, "0");
2192                assert_eq!(msg, "");
2193                assert_eq!(data.len(), 1);
2194            }
2195            _ => panic!("Expected OrderResponse variant"),
2196        }
2197    }
2198
2199    #[rstest]
2200    fn test_order_response_parsing_failure() {
2201        let order_response_json = r#"{
2202            "id": "req-456",
2203            "op": "cancel-order",
2204            "code": "50001",
2205            "msg": "Order not found",
2206            "data": [{"sMsg": "Order with client order ID not found"}]
2207        }"#;
2208
2209        let parsed: OKXWsFrame = serde_json::from_str(order_response_json).unwrap();
2210
2211        match parsed {
2212            OKXWsFrame::OrderResponse {
2213                id,
2214                op,
2215                code,
2216                msg,
2217                data,
2218            } => {
2219                assert_eq!(id, Some("req-456".to_string()));
2220                assert_eq!(op, OKXWsOperation::CancelOrder);
2221                assert_eq!(code, "50001");
2222                assert_eq!(msg, "Order not found");
2223                assert_eq!(data.len(), 1);
2224            }
2225            _ => panic!("Expected OrderResponse variant"),
2226        }
2227    }
2228
2229    #[rstest]
2230    fn test_message_request_serialization() {
2231        let request = OKXWsRequest {
2232            id: Some("req-123".to_string()),
2233            op: OKXWsOperation::Order,
2234            args: vec![::serde_json::json!({
2235                "instId": "BTC-USDT",
2236                "tdMode": "cash",
2237                "side": "buy",
2238                "ordType": "market",
2239                "sz": "0.1"
2240            })],
2241            exp_time: None,
2242        };
2243
2244        let serialized = serde_json::to_string(&request).unwrap();
2245        let parsed: serde_json::Value = serde_json::from_str(&serialized).unwrap();
2246
2247        assert_eq!(parsed["id"], "req-123");
2248        assert_eq!(parsed["op"], "order");
2249        assert!(parsed["args"].is_array());
2250        assert_eq!(parsed["args"].as_array().unwrap().len(), 1);
2251    }
2252
2253    #[rstest]
2254    fn test_ws_post_order_params_serializes_inst_id_code() {
2255        use super::WsPostOrderParamsBuilder;
2256        use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2257
2258        let params = WsPostOrderParamsBuilder::default()
2259            .inst_id_code(10459u64)
2260            .td_mode(OKXTradeMode::Cross)
2261            .side(OKXSide::Buy)
2262            .ord_type(OKXOrderType::Limit)
2263            .sz("0.01".to_string())
2264            .px("50000".to_string())
2265            .build()
2266            .unwrap();
2267
2268        let json = serde_json::to_string(&params).unwrap();
2269
2270        assert!(json.contains("\"instIdCode\":10459"));
2271        assert!(!json.contains("\"instId\""));
2272    }
2273
2274    #[rstest]
2275    fn test_ws_post_order_params_serializes_slippage_pct() {
2276        use super::WsPostOrderParamsBuilder;
2277        use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2278
2279        let params = WsPostOrderParamsBuilder::default()
2280            .inst_id_code(10459u64)
2281            .td_mode(OKXTradeMode::Cross)
2282            .side(OKXSide::Buy)
2283            .ord_type(OKXOrderType::Market)
2284            .sz("0.01".to_string())
2285            .slippage_pct("0.005".to_string())
2286            .build()
2287            .unwrap();
2288
2289        let json: serde_json::Value = serde_json::to_value(&params).unwrap();
2290        assert_eq!(json["slippagePct"], "0.005");
2291    }
2292
2293    #[rstest]
2294    fn test_ws_post_order_params_omits_slippage_pct_when_unset() {
2295        use super::WsPostOrderParamsBuilder;
2296        use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2297
2298        let params = WsPostOrderParamsBuilder::default()
2299            .inst_id_code(10459u64)
2300            .td_mode(OKXTradeMode::Cross)
2301            .side(OKXSide::Buy)
2302            .ord_type(OKXOrderType::Market)
2303            .sz("0.01".to_string())
2304            .build()
2305            .unwrap();
2306
2307        let json = serde_json::to_string(&params).unwrap();
2308        assert!(!json.contains("slippagePct"));
2309    }
2310
2311    #[rstest]
2312    fn test_ws_post_order_params_serializes_attached_tp_sl() {
2313        use super::{WsAttachAlgoOrdParamsBuilder, WsPostOrderParamsBuilder};
2314        use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode, OKXTriggerType};
2315
2316        let params = WsPostOrderParamsBuilder::default()
2317            .inst_id_code(10459u64)
2318            .td_mode(OKXTradeMode::Cross)
2319            .side(OKXSide::Buy)
2320            .ord_type(OKXOrderType::Limit)
2321            .sz("0.01".to_string())
2322            .px("50000".to_string())
2323            .attach_algo_ords(vec![
2324                WsAttachAlgoOrdParamsBuilder::default()
2325                    .attach_algo_cl_ord_id("O-bracket-sl")
2326                    .sl_trigger_px("39000")
2327                    .sl_ord_px("-1")
2328                    .sl_trigger_px_type(OKXTriggerType::Last)
2329                    .build()
2330                    .unwrap(),
2331                WsAttachAlgoOrdParamsBuilder::default()
2332                    .attach_algo_cl_ord_id("O-bracket-tp")
2333                    .tp_trigger_px("41000")
2334                    .tp_ord_px("-1")
2335                    .tp_trigger_px_type(OKXTriggerType::Last)
2336                    .build()
2337                    .unwrap(),
2338            ])
2339            .build()
2340            .unwrap();
2341
2342        let json = serde_json::to_string(&params).unwrap();
2343
2344        assert!(json.contains("\"attachAlgoOrds\""));
2345        assert!(json.contains("\"attachAlgoClOrdId\":\"O-bracket-sl\""));
2346        assert!(json.contains("\"slTriggerPx\":\"39000\""));
2347        assert!(json.contains("\"slOrdPx\":\"-1\""));
2348        assert!(json.contains("\"attachAlgoClOrdId\":\"O-bracket-tp\""));
2349        assert!(json.contains("\"tpTriggerPx\":\"41000\""));
2350        assert!(json.contains("\"tpOrdPx\":\"-1\""));
2351    }
2352
2353    #[rstest]
2354    fn test_ws_cancel_order_params_serializes_inst_id_code() {
2355        use super::WsCancelOrderParamsBuilder;
2356
2357        let params = WsCancelOrderParamsBuilder::default()
2358            .inst_id_code(10461u64)
2359            .ord_id("12345678".to_string())
2360            .build()
2361            .unwrap();
2362
2363        let json = serde_json::to_string(&params).unwrap();
2364
2365        assert!(json.contains("\"instIdCode\":10461"));
2366        assert!(!json.contains("\"instId\""));
2367        assert!(json.contains("\"ordId\":\"12345678\""));
2368    }
2369
2370    #[rstest]
2371    fn test_ws_amend_order_params_serializes_inst_id_code() {
2372        use super::WsAmendOrderParamsBuilder;
2373
2374        let params = WsAmendOrderParamsBuilder::default()
2375            .inst_id_code(10459u64)
2376            .cl_ord_id("client123".to_string())
2377            .new_px("51000".to_string())
2378            .build()
2379            .unwrap();
2380
2381        let json = serde_json::to_string(&params).unwrap();
2382
2383        assert!(json.contains("\"instIdCode\":10459"));
2384        assert!(!json.contains("\"instId\""));
2385        assert!(json.contains("\"newPx\":\"51000\""));
2386    }
2387
2388    #[rstest]
2389    fn test_ws_post_algo_order_params_serializes_inst_id_code() {
2390        use super::WsPostAlgoOrderParamsBuilder;
2391        use crate::common::enums::{OKXAlgoOrderType, OKXSide, OKXTradeMode, OKXTriggerType};
2392
2393        let params = WsPostAlgoOrderParamsBuilder::default()
2394            .inst_id_code(10459u64)
2395            .td_mode(OKXTradeMode::Cross)
2396            .side(OKXSide::Buy)
2397            .ord_type(OKXAlgoOrderType::Trigger)
2398            .sz("0.01".to_string())
2399            .trigger_px("48000".to_string())
2400            .trigger_px_type(OKXTriggerType::Last)
2401            .build()
2402            .unwrap();
2403
2404        let json = serde_json::to_string(&params).unwrap();
2405
2406        assert!(json.contains("\"instIdCode\":10459"));
2407        assert!(!json.contains("\"instId\""));
2408        assert!(json.contains("\"triggerPx\":\"48000\""));
2409    }
2410
2411    #[rstest]
2412    fn test_ws_cancel_algo_order_params_serializes_inst_id_code() {
2413        let params = WsCancelAlgoOrderParams {
2414            inst_id_code: 10459,
2415            algo_id: Some("987654321".to_string()),
2416            algo_cl_ord_id: None,
2417        };
2418
2419        let json = serde_json::to_string(&params).unwrap();
2420
2421        assert!(json.contains("\"instIdCode\":10459"));
2422        assert!(!json.contains("\"instId\""));
2423        assert!(json.contains("\"algoId\":\"987654321\""));
2424    }
2425
2426    #[rstest]
2427    fn test_ws_post_order_params_serializes_px_usd() {
2428        use super::WsPostOrderParamsBuilder;
2429        use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2430
2431        let params = WsPostOrderParamsBuilder::default()
2432            .inst_id_code(10459u64)
2433            .td_mode(OKXTradeMode::Cross)
2434            .side(OKXSide::Buy)
2435            .ord_type(OKXOrderType::Limit)
2436            .sz("1".to_string())
2437            .px_usd("100.5".to_string())
2438            .build()
2439            .unwrap();
2440
2441        let json = serde_json::to_string(&params).unwrap();
2442        assert!(json.contains("\"pxUsd\":\"100.5\""));
2443        assert!(!json.contains("\"pxVol\""));
2444        assert!(!json.contains("\"px\":"));
2445    }
2446
2447    #[rstest]
2448    fn test_ws_post_order_params_serializes_px_vol() {
2449        use super::WsPostOrderParamsBuilder;
2450        use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2451
2452        let params = WsPostOrderParamsBuilder::default()
2453            .inst_id_code(10459u64)
2454            .td_mode(OKXTradeMode::Cross)
2455            .side(OKXSide::Buy)
2456            .ord_type(OKXOrderType::Limit)
2457            .sz("1".to_string())
2458            .px_vol("0.55".to_string())
2459            .build()
2460            .unwrap();
2461
2462        let json = serde_json::to_string(&params).unwrap();
2463        assert!(json.contains("\"pxVol\":\"0.55\""));
2464        assert!(!json.contains("\"pxUsd\""));
2465        assert!(!json.contains("\"px\":"));
2466    }
2467
2468    #[rstest]
2469    fn test_ws_amend_order_params_serializes_new_px_usd() {
2470        use super::WsAmendOrderParamsBuilder;
2471
2472        let params = WsAmendOrderParamsBuilder::default()
2473            .inst_id_code(10459u64)
2474            .cl_ord_id("client123".to_string())
2475            .new_px_usd("105.0".to_string())
2476            .build()
2477            .unwrap();
2478
2479        let json = serde_json::to_string(&params).unwrap();
2480        assert!(json.contains("\"newPxUsd\":\"105.0\""));
2481        assert!(!json.contains("\"newPx\":"));
2482        assert!(!json.contains("\"newPxVol\""));
2483    }
2484
2485    #[rstest]
2486    fn test_ws_amend_order_params_serializes_new_px_vol() {
2487        use super::WsAmendOrderParamsBuilder;
2488
2489        let params = WsAmendOrderParamsBuilder::default()
2490            .inst_id_code(10459u64)
2491            .cl_ord_id("client123".to_string())
2492            .new_px_vol("0.60".to_string())
2493            .build()
2494            .unwrap();
2495
2496        let json = serde_json::to_string(&params).unwrap();
2497        assert!(json.contains("\"newPxVol\":\"0.60\""));
2498        assert!(!json.contains("\"newPx\":"));
2499        assert!(!json.contains("\"newPxUsd\""));
2500    }
2501
2502    #[rstest]
2503    fn test_ws_event_contract_markets_channel_serialization() {
2504        let json = serde_json::to_string(&OKXWsChannel::EventContractMarkets).unwrap();
2505        let channel: OKXWsChannel = serde_json::from_str(&json).unwrap();
2506
2507        assert_eq!(json, "\"event-contract-markets\"");
2508        assert_eq!(channel, OKXWsChannel::EventContractMarkets);
2509    }
2510
2511    #[rstest]
2512    fn test_ws_post_order_params_serializes_event_contract_fields() {
2513        use super::WsPostOrderParamsBuilder;
2514        use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2515
2516        let params = WsPostOrderParamsBuilder::default()
2517            .inst_id_code(10459u64)
2518            .td_mode(OKXTradeMode::Cash)
2519            .side(OKXSide::Buy)
2520            .ord_type(OKXOrderType::Limit)
2521            .sz("10".to_string())
2522            .px("0.42".to_string())
2523            .speed_bump("1")
2524            .outcome("yes")
2525            .build()
2526            .unwrap();
2527
2528        let json: serde_json::Value = serde_json::to_value(&params).unwrap();
2529
2530        assert_eq!(json["speedBump"], "1");
2531        assert_eq!(json["outcome"], "yes");
2532    }
2533
2534    #[rstest]
2535    fn test_ws_amend_order_params_serializes_speed_bump() {
2536        use super::WsAmendOrderParamsBuilder;
2537
2538        let params = WsAmendOrderParamsBuilder::default()
2539            .inst_id_code(10459u64)
2540            .cl_ord_id("event-1".to_string())
2541            .new_px("0.43".to_string())
2542            .speed_bump("1")
2543            .build()
2544            .unwrap();
2545
2546        let json: serde_json::Value = serde_json::to_value(&params).unwrap();
2547
2548        assert_eq!(json["speedBump"], "1");
2549    }
2550
2551    #[rstest]
2552    fn test_ws_attach_algo_ord_params_serializes_trailing_fields() {
2553        use super::WsAttachAlgoOrdParamsBuilder;
2554
2555        let params = WsAttachAlgoOrdParamsBuilder::default()
2556            .attach_algo_cl_ord_id("trail-1")
2557            .callback_ratio("0.01")
2558            .active_px("64000")
2559            .new_callback_ratio("0.02")
2560            .new_callback_spread("25")
2561            .new_active_px("65000")
2562            .build()
2563            .unwrap();
2564
2565        let json: serde_json::Value = serde_json::to_value(&params).unwrap();
2566
2567        assert_eq!(json["callbackRatio"], "0.01");
2568        assert_eq!(json["activePx"], "64000");
2569        assert_eq!(json["newCallbackRatio"], "0.02");
2570        assert_eq!(json["newCallbackSpread"], "25");
2571        assert_eq!(json["newActivePx"], "65000");
2572        assert!(json.get("callbackSpread").is_none());
2573    }
2574
2575    #[rstest]
2576    fn test_subscription_arg_serializes_sprd_id_for_spread_channels() {
2577        let arg = OKXSubscriptionArg {
2578            channel: OKXWsChannel::SprdBooks5,
2579            inst_type: None,
2580            inst_family: None,
2581            inst_id: Some(Ustr::from("ETH-USD-260925_ETH-USD-261225")),
2582        };
2583        let json = serde_json::to_value(&arg).unwrap();
2584        assert_eq!(json["channel"], "sprd-books5");
2585        assert_eq!(json["sprdId"], "ETH-USD-260925_ETH-USD-261225");
2586        assert!(json.get("instId").is_none());
2587    }
2588
2589    #[rstest]
2590    fn test_subscription_arg_serializes_inst_id_for_standard_channels() {
2591        let arg = OKXSubscriptionArg {
2592            channel: OKXWsChannel::BboTbt,
2593            inst_type: None,
2594            inst_family: None,
2595            inst_id: Some(Ustr::from("BTC-USDT")),
2596        };
2597        let json = serde_json::to_value(&arg).unwrap();
2598        assert_eq!(json["instId"], "BTC-USDT");
2599        assert!(json.get("sprdId").is_none());
2600    }
2601
2602    #[rstest]
2603    fn test_websocket_arg_resolves_sprd_id_into_inst_id() {
2604        let arg: OKXWebSocketArg = serde_json::from_value(serde_json::json!({
2605            "channel": "sprd-bbo-tbt",
2606            "sprdId": "ETH-USD-260925_ETH-USD-261225",
2607        }))
2608        .unwrap();
2609        assert_eq!(arg.channel, OKXWsChannel::SprdBboTbt);
2610        assert_eq!(
2611            arg.inst_id,
2612            Some(Ustr::from("ETH-USD-260925_ETH-USD-261225"))
2613        );
2614    }
2615
2616    #[rstest]
2617    fn test_book_msg_parses_three_element_spread_levels() {
2618        // sprd-books5 levels are [price, size, count] (3 elements), unlike the
2619        // 4-element standard book levels.
2620        let msg: OKXBookMsg = serde_json::from_value(serde_json::json!({
2621            "asks": [["16.7", "100", "1"]],
2622            "bids": [["16.65", "100", "1"]],
2623            "ts": "1780044924909",
2624            "seqId": 1779935772619784_u64,
2625        }))
2626        .unwrap();
2627        assert_eq!(msg.asks[0].price, "16.7");
2628        assert_eq!(msg.asks[0].size, "100");
2629        assert_eq!(msg.bids[0].price, "16.65");
2630    }
2631
2632    #[rstest]
2633    fn test_trade_msg_parses_spread_public_trade() {
2634        // sprd-public-trades keys the instrument as `sprdId` and omits `count`.
2635        let msg: OKXTradeMsg = serde_json::from_value(serde_json::json!({
2636            "sprdId": "ETH-USD-260925_ETH-USD-261225",
2637            "tradeId": "3392538740127301632",
2638            "px": "16.9",
2639            "sz": "100",
2640            "side": "sell",
2641            "ts": "1780047866507",
2642        }))
2643        .unwrap();
2644        assert_eq!(msg.inst_id, Ustr::from("ETH-USD-260925_ETH-USD-261225"));
2645        assert_eq!(msg.px, "16.9");
2646        assert_eq!(msg.side, OKXSide::Sell);
2647        assert!(msg.count.is_empty());
2648    }
2649}