1use 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), Reconnected,
72 Authenticated,
73}
74
75#[derive(Debug, Clone, Serialize, Deserialize)]
77pub struct OKXWebSocketError {
78 pub code: String,
80 pub message: String,
82 pub conn_id: Option<String>,
84 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#[derive(Debug)]
104pub enum OKXWsMessage {
105 BookData {
107 arg: OKXWebSocketArg,
108 action: OKXBookAction,
109 data: Vec<OKXBookMsg>,
110 },
111 RpiBookData {
113 arg: OKXWebSocketArg,
114 action: OKXBookAction,
115 data: Vec<OKXRpiBookMsg>,
116 },
117 ChannelData {
119 channel: OKXWsChannel,
120 inst_id: Option<Ustr>,
121 data: serde_json::Value,
122 },
123 OrderResponse {
125 id: Option<String>,
126 op: OKXWsOperation,
127 code: String,
128 msg: String,
129 data: Vec<serde_json::Value>,
130 },
131 Orders(Vec<OKXOrderMsg>),
133 SpreadOrders(Vec<OKXSpreadOrder>),
135 AlgoOrders(Vec<OKXAlgoOrderMsg>),
137 Account(serde_json::Value),
139 Positions(serde_json::Value),
141 LiquidationWarnings(Vec<OKXLiquidationWarningMsg>),
143 Instruments(Vec<OKXInstrument>),
145 SendFailed {
147 request_id: String,
148 client_order_ids: Vec<ClientOrderId>,
149 op: Option<OKXWsOperation>,
150 error: super::error::OKXWsError,
151 },
152 SubscriptionFailed {
154 channel: OKXWsChannel,
155 inst_id: Option<Ustr>,
156 code: String,
157 msg: String,
158 },
159 Error(OKXWebSocketError),
161 Reconnected,
163 Authenticated,
165}
166
167#[derive(Debug, Serialize)]
169#[serde(rename_all = "camelCase")]
170pub struct OKXWsRequest<T> {
171 #[serde(skip_serializing_if = "Option::is_none")]
173 pub id: Option<String>,
174 pub op: OKXWsOperation,
176 #[serde(skip_serializing_if = "Option::is_none")]
179 pub exp_time: Option<String>,
180 pub args: Vec<T>,
182}
183
184#[derive(Debug, Serialize)]
186pub struct OKXAuthentication {
187 pub op: &'static str,
188 pub args: Vec<OKXAuthenticationArg>,
189}
190
191#[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#[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 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 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 if obj.contains_key("op") {
333 return parse_order_response(obj);
334 }
335
336 if obj.contains_key("action") && obj.contains_key("arg") {
338 return parse_book_data(obj);
339 }
340
341 if obj.contains_key("arg") && obj.contains_key("data") {
343 return parse_data(obj);
344 }
345
346 if obj.contains_key("code") && obj.contains_key("msg") {
348 return parse_error(obj);
349 }
350
351 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 pub channel: OKXWsChannel,
515 #[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#[derive(Debug, Serialize, Deserialize)]
530#[serde(rename_all = "camelCase")]
531pub struct OKXTickerMsg {
532 pub inst_type: OKXInstrumentType,
534 pub inst_id: Ustr,
536 #[serde(rename = "last")]
538 pub last_px: String,
539 pub last_sz: String,
541 pub ask_px: String,
543 pub ask_sz: String,
545 pub bid_px: String,
547 pub bid_sz: String,
549 pub open24h: String,
551 pub high24h: String,
553 pub low24h: String,
555 pub vol_ccy_24h: String,
557 pub vol24h: String,
559 pub sod_utc0: String,
561 pub sod_utc8: String,
563 #[serde(deserialize_with = "deserialize_string_to_u64")]
565 pub ts: u64,
566 #[serde(default)]
568 pub source: Option<String>,
569}
570
571#[derive(Debug, Serialize, Deserialize)]
573pub struct OrderBookEntry {
574 pub price: String,
576 pub size: String,
578 #[serde(default)]
583 pub liquidated_orders_count: String,
584 #[serde(default)]
586 pub orders_count: String,
587}
588
589#[derive(Debug, Serialize, Deserialize)]
591#[serde(rename_all = "camelCase")]
592pub struct OKXBookMsg {
593 pub asks: Vec<OrderBookEntry>,
595 pub bids: Vec<OrderBookEntry>,
597 pub checksum: Option<i64>,
599 pub prev_seq_id: Option<i64>,
601 pub seq_id: u64,
603 #[serde(deserialize_with = "deserialize_string_to_u64")]
605 pub ts: u64,
606}
607
608#[derive(Debug, Serialize, Deserialize)]
610#[serde(rename_all = "camelCase", deny_unknown_fields)]
611pub struct OKXRpiBookMsg {
612 pub asks: Vec<OKXRpiBookLevel>,
614 pub bids: Vec<OKXRpiBookLevel>,
616 pub prev_seq_id: i64,
618 pub seq_id: u64,
620 #[serde(deserialize_with = "deserialize_string_to_u64")]
622 pub ts: u64,
623}
624
625#[derive(Debug, Serialize, Deserialize)]
627#[serde(rename_all = "camelCase")]
628pub struct OKXTradeMsg {
629 #[serde(default, alias = "sprdId")]
634 pub inst_id: Ustr,
635 pub trade_id: String,
637 pub px: String,
639 pub sz: String,
641 pub side: OKXSide,
643 #[serde(default)]
645 pub count: String,
646 #[serde(deserialize_with = "deserialize_string_to_u64")]
648 pub ts: u64,
649 #[serde(default)]
651 pub source: Option<String>,
652 #[serde(default)]
654 pub seq_id: Option<u64>,
655}
656
657#[derive(Debug, Serialize, Deserialize)]
659#[serde(rename_all = "camelCase")]
660pub struct OKXFundingRateMsg {
661 #[serde(default)]
663 pub inst_type: Option<OKXInstrumentType>,
664 pub inst_id: Ustr,
666 pub funding_rate: Ustr,
668 pub next_funding_rate: Ustr,
670 #[serde(default)]
672 pub min_funding_rate: Option<String>,
673 #[serde(default)]
675 pub max_funding_rate: Option<String>,
676 #[serde(default)]
678 pub sett_state: OKXSettlementState,
679 #[serde(default)]
681 pub sett_funding_rate: Option<String>,
682 #[serde(default)]
684 pub premium: Option<String>,
685 #[serde(default)]
687 pub method: Option<String>,
688 #[serde(deserialize_with = "deserialize_string_to_u64")]
690 pub funding_time: u64,
691 #[serde(deserialize_with = "deserialize_string_to_u64")]
693 pub next_funding_time: u64,
694 #[serde(deserialize_with = "deserialize_string_to_u64")]
696 pub ts: u64,
697}
698
699#[derive(Debug, Serialize, Deserialize)]
701#[serde(rename_all = "camelCase")]
702pub struct OKXMarkPriceMsg {
703 pub inst_id: Ustr,
705 pub mark_px: String,
707 #[serde(deserialize_with = "deserialize_string_to_u64")]
709 pub ts: u64,
710}
711
712#[derive(Debug, Serialize, Deserialize)]
714#[serde(rename_all = "camelCase")]
715pub struct OKXIndexPriceMsg {
716 pub inst_id: Ustr,
718 pub idx_px: String,
720 pub high24h: String,
722 pub low24h: String,
724 pub open24h: String,
726 pub sod_utc0: String,
728 pub sod_utc8: String,
730 #[serde(deserialize_with = "deserialize_string_to_u64")]
732 pub ts: u64,
733}
734
735#[derive(Debug, Serialize, Deserialize)]
737#[serde(rename_all = "camelCase")]
738pub struct OKXPriceLimitMsg {
739 pub inst_id: Ustr,
741 pub buy_lmt: String,
743 pub sell_lmt: String,
745 #[serde(deserialize_with = "deserialize_string_to_u64")]
747 pub ts: u64,
748}
749
750#[derive(Debug, Serialize, Deserialize)]
752#[serde(rename_all = "camelCase")]
753pub struct OKXCandleMsg {
754 #[serde(deserialize_with = "deserialize_string_to_u64")]
756 pub ts: u64,
757 pub o: String,
759 pub h: String,
761 pub l: String,
763 pub c: String,
765 pub vol: String,
767 pub vol_ccy: String,
769 pub vol_ccy_quote: String,
770 pub confirm: OKXCandleConfirm,
772}
773
774#[derive(Debug, Serialize, Deserialize)]
776#[serde(rename_all = "camelCase")]
777pub struct OKXOpenInterestMsg {
778 pub inst_id: Ustr,
780 pub oi: String,
782 pub oi_ccy: String,
784 #[serde(deserialize_with = "deserialize_string_to_u64")]
786 pub ts: u64,
787}
788
789#[derive(Debug, Serialize, Deserialize)]
791#[serde(rename_all = "camelCase")]
792pub struct OKXOptionSummaryMsg {
793 #[serde(default)]
795 pub inst_type: Option<OKXInstrumentType>,
796 pub inst_id: Ustr,
798 pub uly: String,
800 pub delta: String,
802 pub gamma: String,
804 pub theta: String,
806 pub vega: String,
808 #[serde(alias = "deltaBS")]
810 pub delta_bs: String,
811 #[serde(alias = "gammaBS")]
813 pub gamma_bs: String,
814 #[serde(alias = "thetaBS")]
816 pub theta_bs: String,
817 #[serde(alias = "vegaBS")]
819 pub vega_bs: String,
820 pub real_vol: String,
822 pub bid_vol: String,
824 pub ask_vol: String,
826 pub mark_vol: String,
828 pub lever: String,
830 #[serde(default)]
832 pub fwd_px: Option<String>,
833 #[serde(default)]
835 pub mark_px: Option<String>,
836 #[serde(default)]
838 pub vol_lv: Option<String>,
839 #[serde(deserialize_with = "deserialize_string_to_u64")]
841 pub ts: u64,
842}
843
844#[derive(Debug, Serialize, Deserialize)]
846#[serde(rename_all = "camelCase")]
847pub struct OKXEstimatedPriceMsg {
848 pub inst_id: Ustr,
850 pub settle_px: String,
852 #[serde(deserialize_with = "deserialize_string_to_u64")]
854 pub ts: u64,
855}
856
857#[derive(Debug, Serialize, Deserialize)]
859#[serde(rename_all = "camelCase")]
860pub struct OKXStatusMsg {
861 pub title: Ustr,
863 #[serde(rename = "type")]
865 pub status_type: Ustr,
866 pub state: Ustr,
868 pub end_time: Option<String>,
870 pub begin_time: Option<String>,
872 pub service_type: Option<Ustr>,
874 pub reason: Option<String>,
876 #[serde(deserialize_with = "deserialize_string_to_u64")]
878 pub ts: u64,
879}
880
881pub use crate::common::models::OKXAttachedAlgoOrd;
882
883#[derive(Clone, Debug, Serialize, Deserialize)]
889#[serde(rename_all = "camelCase")]
890pub struct OKXLiquidationWarningMsg {
891 pub inst_type: OKXInstrumentType,
893 #[serde(default)]
895 pub inst_family: Option<Ustr>,
896 pub inst_id: Ustr,
898 pub mgn_mode: OKXMarginMode,
900 #[serde(default)]
902 pub pos_id: Option<Ustr>,
903 pub pos_side: OKXPositionSide,
905 pub pos: String,
907 #[serde(default)]
909 pub pos_ccy: Option<Ustr>,
910 pub lever: String,
912 pub mark_px: String,
914 pub mgn_ratio: String,
916 pub ccy: Ustr,
918 #[serde(deserialize_with = "deserialize_string_to_u64")]
920 pub c_time: u64,
921 #[serde(deserialize_with = "deserialize_string_to_u64")]
923 pub u_time: u64,
924 #[serde(default)]
926 pub p_time: Option<String>,
927}
928
929#[derive(Clone, Debug, Default, Serialize, Deserialize)]
931#[serde(rename_all = "camelCase")]
932pub struct OKXLinkedAlgoOrd {
933 #[serde(default)]
935 pub algo_id: String,
936}
937
938#[derive(Clone, Debug, Serialize, Deserialize)]
940#[serde(rename_all = "camelCase")]
941pub struct OKXOrderMsg {
942 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
944 pub acc_fill_sz: Option<String>,
945 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
947 pub algo_id: Option<String>,
948 pub avg_px: String,
950 #[serde(deserialize_with = "deserialize_string_to_u64")]
952 pub c_time: u64,
953 #[serde(default)]
955 pub cancel_source: Option<String>,
956 #[serde(default)]
958 pub cancel_source_reason: Option<String>,
959 pub category: OKXOrderCategory,
961 pub ccy: Ustr,
963 pub cl_ord_id: String,
965 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
967 pub algo_cl_ord_id: Option<String>,
968 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
970 pub attach_algo_cl_ord_id: Option<String>,
971 #[serde(default)]
973 pub attach_algo_ords: Vec<OKXAttachedAlgoOrd>,
974 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
976 pub outcome: Option<String>,
977 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
979 pub fee: Option<String>,
980 pub fee_ccy: Ustr,
982 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
984 pub fill_fee: Option<String>,
985 #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
987 pub fill_fee_ccy: Option<Ustr>,
988 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
990 pub fill_mark_px: Option<String>,
991 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
993 pub fill_mark_vol: Option<String>,
994 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
996 pub fill_px_vol: Option<String>,
997 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
999 pub fill_px_usd: Option<String>,
1000 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1002 pub fill_fwd_px: Option<String>,
1003 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1005 pub fill_notional_usd: Option<String>,
1006 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1008 pub fill_pnl: Option<String>,
1009 pub fill_px: String,
1011 pub fill_sz: String,
1013 #[serde(deserialize_with = "deserialize_string_to_u64")]
1015 pub fill_time: u64,
1016 pub inst_id: Ustr,
1018 pub inst_type: OKXInstrumentType,
1020 #[serde(default)]
1022 pub is_tp_limit: Option<String>,
1023 pub lever: String,
1025 #[serde(default)]
1027 pub linked_algo_ord: Option<OKXLinkedAlgoOrd>,
1028 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1030 pub notional_usd: Option<String>,
1031 pub ord_id: Ustr,
1033 pub ord_type: OKXOrderType,
1035 pub pnl: String,
1037 pub pos_side: OKXPositionSide,
1039 #[serde(default)]
1041 pub px: String,
1042 #[serde(default)]
1044 pub px_type: OKXPriceType,
1045 #[serde(default)]
1047 pub px_usd: Option<String>,
1048 #[serde(default)]
1050 pub px_vol: Option<String>,
1051 #[serde(default)]
1053 pub quick_mgn_type: OKXQuickMarginType,
1054 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1056 pub rebate: Option<String>,
1057 #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
1059 pub rebate_ccy: Option<Ustr>,
1060 pub reduce_only: String,
1062 pub side: OKXSide,
1064 #[serde(default)]
1066 pub sl_ord_px: Option<String>,
1067 #[serde(default)]
1069 pub sl_trigger_px: Option<String>,
1070 #[serde(default)]
1072 pub sl_trigger_px_type: Option<OKXTriggerType>,
1073 #[serde(default)]
1075 pub source: Option<String>,
1076 pub state: OKXOrderStatus,
1078 #[serde(default)]
1080 pub stp_id: Option<String>,
1081 #[serde(default)]
1083 pub stp_mode: OKXSelfTradePreventionMode,
1084 pub exec_type: OKXExecType,
1086 pub sz: String,
1088 #[serde(default)]
1090 pub tag: Option<String>,
1091 pub td_mode: OKXTradeMode,
1093 #[serde(default, deserialize_with = "deserialize_target_currency_as_none")]
1095 pub tgt_ccy: Option<OKXTargetCurrency>,
1096 #[serde(default)]
1098 pub tp_ord_px: Option<String>,
1099 #[serde(default)]
1101 pub tp_trigger_px: Option<String>,
1102 #[serde(default)]
1104 pub tp_trigger_px_type: Option<OKXTriggerType>,
1105 pub trade_id: String,
1107 #[serde(deserialize_with = "deserialize_string_to_u64")]
1109 pub u_time: u64,
1110 #[serde(default)]
1112 pub amend_result: Option<String>,
1113 #[serde(default)]
1115 pub req_id: Option<String>,
1116 #[serde(default)]
1118 pub code: Option<String>,
1119 #[serde(default)]
1121 pub msg: Option<String>,
1122}
1123
1124#[derive(Clone, Debug, Deserialize, Serialize)]
1126#[serde(rename_all = "camelCase")]
1127pub struct OKXAlgoOrderMsg {
1128 pub algo_id: String,
1130 #[serde(default)]
1132 pub algo_cl_ord_id: String,
1133 pub cl_ord_id: String,
1135 pub ord_id: String,
1137 #[serde(default)]
1139 pub ord_id_list: Vec<String>,
1140 pub inst_id: Ustr,
1142 pub inst_type: OKXInstrumentType,
1144 pub ord_type: OKXAlgoOrderType,
1146 pub state: OKXAlgoOrderStatus,
1148 pub side: OKXSide,
1150 pub pos_side: OKXPositionSide,
1152 #[serde(default)]
1154 pub sz: String,
1155 #[serde(default)]
1157 pub trigger_px: String,
1158 #[serde(default)]
1160 pub trigger_px_type: OKXTriggerType,
1161 #[serde(default)]
1163 pub sl_trigger_px: String,
1164 #[serde(default)]
1166 pub sl_ord_px: String,
1167 #[serde(default)]
1169 pub sl_trigger_px_type: OKXTriggerType,
1170 #[serde(default)]
1172 pub tp_trigger_px: String,
1173 #[serde(default)]
1175 pub tp_ord_px: String,
1176 #[serde(default)]
1178 pub tp_trigger_px_type: OKXTriggerType,
1179 #[serde(default)]
1181 pub ord_px: String,
1182 pub td_mode: OKXTradeMode,
1184 pub lever: String,
1186 #[serde(default)]
1188 pub reduce_only: String,
1189 #[serde(default)]
1191 pub close_fraction: String,
1192 #[serde(default)]
1194 pub actual_px: String,
1195 #[serde(default)]
1197 pub actual_sz: String,
1198 #[serde(default)]
1200 pub notional_usd: String,
1201 #[serde(deserialize_with = "deserialize_string_to_u64")]
1203 pub c_time: u64,
1204 #[serde(deserialize_with = "deserialize_string_to_u64")]
1206 pub u_time: u64,
1207 #[serde(default)]
1209 pub trigger_time: String,
1210 #[serde(default)]
1212 pub fail_code: String,
1213 #[serde(default)]
1215 pub tag: String,
1216 #[serde(default)]
1218 pub callback_ratio: String,
1219 #[serde(default)]
1221 pub callback_spread: String,
1222 #[serde(default)]
1224 pub active_px: String,
1225 #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
1227 pub ccy: Option<Ustr>,
1228 #[serde(default, deserialize_with = "deserialize_target_currency_as_none")]
1230 pub tgt_ccy: Option<OKXTargetCurrency>,
1231 #[serde(default)]
1233 pub fee: Option<String>,
1234 #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
1236 pub fee_ccy: Option<Ustr>,
1237 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1239 pub advance_ord_type: Option<String>,
1240}
1241
1242#[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 #[serde(skip_serializing_if = "Option::is_none")]
1250 pub attach_algo_cl_ord_id: Option<String>,
1251 #[serde(skip_serializing_if = "Option::is_none")]
1253 pub sl_trigger_px: Option<String>,
1254 #[serde(skip_serializing_if = "Option::is_none")]
1256 pub sl_ord_px: Option<String>,
1257 #[serde(skip_serializing_if = "Option::is_none")]
1259 pub sl_trigger_px_type: Option<OKXTriggerType>,
1260 #[serde(skip_serializing_if = "Option::is_none")]
1262 pub tp_trigger_px: Option<String>,
1263 #[serde(skip_serializing_if = "Option::is_none")]
1265 pub tp_ord_px: Option<String>,
1266 #[serde(skip_serializing_if = "Option::is_none")]
1268 pub tp_trigger_px_type: Option<OKXTriggerType>,
1269 #[serde(skip_serializing_if = "Option::is_none")]
1271 pub callback_ratio: Option<String>,
1272 #[serde(skip_serializing_if = "Option::is_none")]
1274 pub callback_spread: Option<String>,
1275 #[serde(skip_serializing_if = "Option::is_none")]
1277 pub active_px: Option<String>,
1278 #[serde(skip_serializing_if = "Option::is_none")]
1280 pub new_callback_ratio: Option<String>,
1281 #[serde(skip_serializing_if = "Option::is_none")]
1283 pub new_callback_spread: Option<String>,
1284 #[serde(skip_serializing_if = "Option::is_none")]
1286 pub new_active_px: Option<String>,
1287}
1288
1289#[derive(Clone, Debug, Deserialize, Serialize, Builder)]
1291#[builder(setter(into, strip_option))]
1292#[serde(rename_all = "camelCase")]
1293pub struct WsPostOrderParams {
1294 #[builder(default)]
1296 #[serde(skip_serializing_if = "Option::is_none")]
1297 pub inst_type: Option<OKXInstrumentType>,
1298 pub inst_id_code: u64,
1300 pub td_mode: OKXTradeMode,
1302 #[builder(default)]
1304 #[serde(skip_serializing_if = "Option::is_none")]
1305 pub ccy: Option<Ustr>,
1306 #[builder(default)]
1308 #[serde(skip_serializing_if = "Option::is_none")]
1309 pub cl_ord_id: Option<String>,
1310 pub side: OKXSide,
1312 #[builder(default)]
1314 #[serde(skip_serializing_if = "Option::is_none")]
1315 pub pos_side: Option<OKXPositionSide>,
1316 pub ord_type: OKXOrderType,
1318 pub sz: String,
1320 #[builder(default)]
1322 #[serde(skip_serializing_if = "Option::is_none")]
1323 pub px: Option<String>,
1324 #[builder(default)]
1326 #[serde(rename = "pxUsd", skip_serializing_if = "Option::is_none")]
1327 pub px_usd: Option<String>,
1328 #[builder(default)]
1331 #[serde(rename = "pxVol", skip_serializing_if = "Option::is_none")]
1332 pub px_vol: Option<String>,
1333 #[builder(default)]
1335 #[serde(skip_serializing_if = "Option::is_none")]
1336 pub reduce_only: Option<bool>,
1337 #[builder(default)]
1339 #[serde(rename = "closePosition", skip_serializing_if = "Option::is_none")]
1340 pub close_position: Option<bool>,
1341 #[builder(default)]
1343 #[serde(skip_serializing_if = "Option::is_none")]
1344 pub tgt_ccy: Option<OKXTargetCurrency>,
1345 #[builder(default)]
1347 #[serde(skip_serializing_if = "Option::is_none")]
1348 pub tag: Option<String>,
1349 #[builder(default)]
1351 #[serde(skip_serializing_if = "Option::is_none")]
1352 pub attach_algo_ords: Option<Vec<WsAttachAlgoOrdParams>>,
1353 #[builder(default)]
1355 #[serde(skip_serializing_if = "Option::is_none")]
1356 pub speed_bump: Option<String>,
1357 #[builder(default)]
1359 #[serde(skip_serializing_if = "Option::is_none")]
1360 pub outcome: Option<String>,
1361 #[builder(default)]
1366 #[serde(skip_serializing_if = "Option::is_none")]
1367 pub slippage_pct: Option<String>,
1368 #[builder(default)]
1370 #[serde(skip_serializing_if = "Option::is_none")]
1371 pub rpi_taker_access: Option<bool>,
1372 #[builder(default)]
1374 #[serde(skip_serializing_if = "Option::is_none")]
1375 pub rpi_px_round: Option<bool>,
1376}
1377
1378#[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 pub inst_id_code: u64,
1386 #[serde(skip_serializing_if = "Option::is_none")]
1388 pub ord_id: Option<String>,
1389 #[serde(skip_serializing_if = "Option::is_none")]
1391 pub cl_ord_id: Option<String>,
1392}
1393
1394#[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 pub inst_type: OKXInstrumentType,
1402 pub inst_family: Ustr,
1404}
1405
1406#[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 pub inst_id_code: u64,
1414 #[serde(skip_serializing_if = "Option::is_none")]
1416 pub ord_id: Option<String>,
1417 #[serde(skip_serializing_if = "Option::is_none")]
1419 pub cl_ord_id: Option<String>,
1420 #[serde(skip_serializing_if = "Option::is_none")]
1422 pub req_id: Option<String>,
1423 #[serde(skip_serializing_if = "Option::is_none")]
1425 pub new_px: Option<String>,
1426 #[serde(rename = "newPxUsd", skip_serializing_if = "Option::is_none")]
1428 pub new_px_usd: Option<String>,
1429 #[serde(rename = "newPxVol", skip_serializing_if = "Option::is_none")]
1432 pub new_px_vol: Option<String>,
1433 #[serde(skip_serializing_if = "Option::is_none")]
1435 pub new_sz: Option<String>,
1436 #[serde(skip_serializing_if = "Option::is_none")]
1438 pub rpi_taker_access: Option<bool>,
1439 #[serde(skip_serializing_if = "Option::is_none")]
1441 pub rpi_px_round: Option<bool>,
1442 #[serde(skip_serializing_if = "Option::is_none")]
1444 pub speed_bump: Option<String>,
1445}
1446
1447#[derive(Clone, Debug, Deserialize, Serialize, Builder)]
1449#[builder(setter(into, strip_option))]
1450#[serde(rename_all = "camelCase")]
1451pub struct WsPostAlgoOrderParams {
1452 pub inst_id_code: u64,
1454 pub td_mode: OKXTradeMode,
1456 pub side: OKXSide,
1458 pub ord_type: OKXAlgoOrderType,
1460 pub sz: String,
1462 #[builder(default)]
1464 #[serde(skip_serializing_if = "Option::is_none")]
1465 pub cl_ord_id: Option<String>,
1466 #[builder(default)]
1468 #[serde(skip_serializing_if = "Option::is_none")]
1469 pub pos_side: Option<OKXPositionSide>,
1470 #[serde(skip_serializing_if = "Option::is_none")]
1472 pub trigger_px: Option<String>,
1473 #[builder(default)]
1475 #[serde(skip_serializing_if = "Option::is_none")]
1476 pub trigger_px_type: Option<OKXTriggerType>,
1477 #[builder(default)]
1479 #[serde(skip_serializing_if = "Option::is_none")]
1480 pub order_px: Option<String>,
1481 #[builder(default)]
1483 #[serde(skip_serializing_if = "Option::is_none")]
1484 pub reduce_only: Option<bool>,
1485 #[builder(default)]
1487 #[serde(skip_serializing_if = "Option::is_none")]
1488 pub tag: Option<String>,
1489 #[builder(default)]
1491 #[serde(skip_serializing_if = "Option::is_none")]
1492 pub callback_ratio: Option<String>,
1493 #[builder(default)]
1495 #[serde(skip_serializing_if = "Option::is_none")]
1496 pub callback_spread: Option<String>,
1497 #[builder(default)]
1499 #[serde(skip_serializing_if = "Option::is_none")]
1500 pub active_px: Option<String>,
1501}
1502
1503#[derive(Clone, Debug, Deserialize, Serialize, Builder)]
1505#[builder(setter(into, strip_option))]
1506#[serde(rename_all = "camelCase")]
1507pub struct WsCancelAlgoOrderParams {
1508 pub inst_id_code: u64,
1510 #[serde(skip_serializing_if = "Option::is_none")]
1512 pub algo_id: Option<String>,
1513 #[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 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 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"); 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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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(¶ms).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 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 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}