1use derive_builder::Builder;
19use nautilus_core::string::secret::SecretString;
20use nautilus_model::{
21 data::{Data, FundingRateUpdate, InstrumentStatus, OrderBookDeltas},
22 events::{
23 AccountState, OrderAccepted, OrderCancelRejected, OrderCanceled, OrderExpired,
24 OrderModifyRejected, OrderRejected, OrderTriggered, OrderUpdated,
25 },
26 identifiers::ClientOrderId,
27 instruments::InstrumentAny,
28 reports::{FillReport, OrderStatusReport, PositionStatusReport},
29};
30use serde::{Deserialize, Serialize};
31use ustr::Ustr;
32use zeroize::Zeroize;
33
34use super::enums::{OKXWsChannel, OKXWsOperation};
35use crate::{
36 common::{
37 enums::{
38 OKXAlgoOrderStatus, OKXAlgoOrderType, OKXBookAction, OKXCandleConfirm, OKXExecType,
39 OKXInstrumentType, OKXMarginMode, OKXOrderCategory, OKXOrderStatus, OKXOrderType,
40 OKXPositionSide, OKXPriceType, OKXQuickMarginType, OKXSelfTradePreventionMode,
41 OKXSettlementState, OKXSide, OKXTargetCurrency, OKXTradeMode, OKXTriggerType,
42 },
43 models::{OKXInstrument, OKXRpiBookLevel},
44 parse::{
45 deserialize_empty_string_as_none, deserialize_empty_ustr_as_none,
46 deserialize_string_to_u64, deserialize_target_currency_as_none,
47 },
48 },
49 http::models::OKXSpreadOrder,
50 websocket::enums::OKXSubscriptionEvent,
51};
52
53#[derive(Debug, Clone)]
54pub enum NautilusWsMessage {
55 Data(Vec<Data>),
56 Deltas(OrderBookDeltas),
57 FundingRates(Vec<FundingRateUpdate>),
58 Instrument(Box<InstrumentAny>, Option<InstrumentStatus>),
59 InstrumentStatus(InstrumentStatus),
60 AccountUpdate(AccountState),
61 PositionUpdate(PositionStatusReport),
62 OrderAccepted(OrderAccepted),
63 OrderCanceled(OrderCanceled),
64 OrderExpired(OrderExpired),
65 OrderRejected(OrderRejected),
66 OrderCancelRejected(OrderCancelRejected),
67 OrderModifyRejected(OrderModifyRejected),
68 OrderTriggered(OrderTriggered),
69 OrderUpdated(OrderUpdated),
70 ExecutionReports(Vec<ExecutionReport>),
71 Error(OKXWebSocketError),
72 Raw(serde_json::Value), Reconnected,
74 Authenticated,
75}
76
77#[derive(Debug, Clone, Serialize, Deserialize)]
79pub struct OKXWebSocketError {
80 pub code: String,
82 pub message: String,
84 pub conn_id: Option<String>,
86 pub timestamp: u64,
88}
89
90#[derive(Debug, Clone)]
91#[allow(
92 clippy::large_enum_variant,
93 reason = "the variant size gap only crosses the threshold when high-precision widens the raw types"
94)]
95pub enum ExecutionReport {
96 Order(OrderStatusReport),
97 Fill(FillReport),
98}
99
100#[derive(Debug)]
106pub enum OKXWsMessage {
107 BookData {
109 arg: OKXWebSocketArg,
110 action: OKXBookAction,
111 data: Vec<OKXBookMsg>,
112 },
113 RpiBookData {
115 arg: OKXWebSocketArg,
116 action: OKXBookAction,
117 data: Vec<OKXRpiBookMsg>,
118 },
119 ChannelData {
121 channel: OKXWsChannel,
122 inst_id: Option<Ustr>,
123 data: serde_json::Value,
124 },
125 OrderResponse {
127 id: Option<String>,
128 op: OKXWsOperation,
129 code: String,
130 msg: String,
131 data: Vec<serde_json::Value>,
132 },
133 Orders(Vec<OKXOrderMsg>),
135 SpreadOrders(Vec<OKXSpreadOrder>),
137 AlgoOrders(Vec<OKXAlgoOrderMsg>),
139 Account(serde_json::Value),
141 Positions(serde_json::Value),
143 LiquidationWarnings(Vec<OKXLiquidationWarningMsg>),
145 Instruments(Vec<OKXInstrument>),
147 SendFailed {
149 request_id: String,
150 client_order_ids: Vec<ClientOrderId>,
151 op: Option<OKXWsOperation>,
152 error: super::error::OKXWsError,
153 },
154 SubscriptionFailed {
156 channel: OKXWsChannel,
157 inst_id: Option<Ustr>,
158 code: String,
159 msg: String,
160 },
161 Error(OKXWebSocketError),
163 Reconnected,
165 Authenticated,
167}
168
169#[derive(Debug, Serialize)]
171#[serde(rename_all = "camelCase")]
172pub struct OKXWsRequest<T> {
173 #[serde(skip_serializing_if = "Option::is_none")]
175 pub id: Option<String>,
176 pub op: OKXWsOperation,
178 #[serde(skip_serializing_if = "Option::is_none")]
181 pub exp_time: Option<String>,
182 pub args: Vec<T>,
184}
185
186#[derive(Debug, Serialize, Zeroize)]
188pub struct OKXAuthentication {
189 #[zeroize(skip)]
190 pub op: &'static str,
191 pub args: Vec<OKXAuthenticationArg>,
192}
193
194#[derive(Debug, Serialize, Zeroize)]
196#[serde(rename_all = "camelCase")]
197pub struct OKXAuthenticationArg {
198 pub api_key: SecretString,
199 pub passphrase: SecretString,
200 pub timestamp: String,
201 pub sign: SecretString,
202}
203
204#[derive(Debug, Serialize)]
205pub struct OKXSubscription {
206 pub op: OKXWsOperation,
207 pub args: Vec<OKXSubscriptionArg>,
208}
209
210#[derive(Clone, Debug)]
211pub struct OKXSubscriptionArg {
212 pub channel: OKXWsChannel,
213 pub inst_type: Option<OKXInstrumentType>,
214 pub inst_family: Option<Ustr>,
215 pub inst_id: Option<Ustr>,
216}
217
218impl Serialize for OKXSubscriptionArg {
219 fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
220 use serde::ser::SerializeMap;
221
222 let mut map = serializer.serialize_map(None)?;
223 map.serialize_entry("channel", &self.channel)?;
224
225 if let Some(inst_type) = &self.inst_type {
226 map.serialize_entry("instType", inst_type)?;
227 }
228
229 if let Some(inst_family) = &self.inst_family {
230 map.serialize_entry("instFamily", inst_family)?;
231 }
232
233 if let Some(inst_id) = &self.inst_id {
234 let key = if self.channel.is_spread() {
235 "sprdId"
236 } else {
237 "instId"
238 };
239 map.serialize_entry(key, inst_id)?;
240 }
241
242 map.end()
243 }
244}
245
246#[derive(Debug)]
251pub enum OKXWsFrame {
252 Login {
253 event: String,
254 code: String,
255 msg: String,
256 conn_id: String,
257 },
258 Subscription {
259 event: OKXSubscriptionEvent,
260 arg: OKXWebSocketArg,
261 conn_id: String,
262 code: Option<String>,
263 msg: Option<String>,
264 },
265 ChannelConnCount {
266 event: String,
267 channel: OKXWsChannel,
268 conn_count: String,
269 conn_id: String,
270 },
271 OrderResponse {
272 id: Option<String>,
273 op: OKXWsOperation,
274 code: String,
275 msg: String,
276 data: Vec<serde_json::Value>,
277 },
278 BookData {
279 arg: OKXWebSocketArg,
280 action: OKXBookAction,
281 data: Vec<OKXBookMsg>,
282 },
283 RpiBookData {
284 arg: OKXWebSocketArg,
285 action: OKXBookAction,
286 data: Vec<OKXRpiBookMsg>,
287 },
288 Data {
289 arg: OKXWebSocketArg,
290 data: serde_json::Value,
291 },
292 Error {
293 arg: Option<OKXWebSocketArg>,
294 code: String,
295 msg: String,
296 },
297 Ping,
298 Reconnected,
299}
300
301impl<'de> Deserialize<'de> for OKXWsFrame {
302 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
303 where
304 D: serde::Deserializer<'de>,
305 {
306 use serde::de::Error;
307
308 let mut value = serde_json::Value::deserialize(deserializer)?;
313 let obj = value
314 .as_object_mut()
315 .ok_or_else(|| D::Error::custom("expected JSON object for OKXWsFrame"))?;
316
317 if let Some(event) = obj.get("event").and_then(|v| v.as_str()) {
323 match event {
324 "login" => return parse_login(obj),
325 "subscribe" | "unsubscribe" => return parse_subscription(obj),
326 "error" => return parse_error(obj),
327 _ if obj.contains_key("channel") && obj.contains_key("connCount") => {
328 return parse_channel_conn_count(obj);
329 }
330 _ => {}
331 }
332 }
333
334 if obj.contains_key("op") {
336 return parse_order_response(obj);
337 }
338
339 if obj.contains_key("action") && obj.contains_key("arg") {
341 return parse_book_data(obj);
342 }
343
344 if obj.contains_key("arg") && obj.contains_key("data") {
346 return parse_data(obj);
347 }
348
349 if obj.contains_key("code") && obj.contains_key("msg") {
351 return parse_error(obj);
352 }
353
354 Err(D::Error::custom(format!(
358 "cannot determine OKXWsFrame variant from: {}",
359 serde_json::to_string(&value).unwrap_or_default()
360 )))
361 }
362}
363
364#[inline]
365fn take_str<E: serde::de::Error>(
366 obj: &mut serde_json::Map<String, serde_json::Value>,
367 key: &'static str,
368) -> Result<String, E> {
369 match obj.remove(key) {
370 Some(serde_json::Value::String(s)) => Ok(s),
371 Some(_) => Err(E::custom(format!("field `{key}` is not a string"))),
372 None => Err(E::missing_field(key)),
373 }
374}
375
376#[inline]
377fn take_optional_str(
378 obj: &mut serde_json::Map<String, serde_json::Value>,
379 key: &'static str,
380) -> Option<String> {
381 match obj.remove(key) {
382 Some(serde_json::Value::String(s)) => Some(s),
383 _ => None,
384 }
385}
386
387fn parse_login<E: serde::de::Error>(
388 obj: &mut serde_json::Map<String, serde_json::Value>,
389) -> Result<OKXWsFrame, E> {
390 Ok(OKXWsFrame::Login {
391 event: take_str(obj, "event")?,
392 code: take_str(obj, "code")?,
393 msg: take_str(obj, "msg")?,
394 conn_id: take_str(obj, "connId")?,
395 })
396}
397
398fn parse_subscription<E: serde::de::Error>(
399 obj: &mut serde_json::Map<String, serde_json::Value>,
400) -> Result<OKXWsFrame, E> {
401 let event_val = obj
402 .remove("event")
403 .ok_or_else(|| E::missing_field("event"))?;
404 let event: OKXSubscriptionEvent =
405 serde_json::from_value(event_val).map_err(|e| E::custom(format!("invalid event: {e}")))?;
406
407 let arg_val = obj.remove("arg").ok_or_else(|| E::missing_field("arg"))?;
408 let arg: OKXWebSocketArg =
409 serde_json::from_value(arg_val).map_err(|e| E::custom(format!("invalid arg: {e}")))?;
410
411 Ok(OKXWsFrame::Subscription {
412 event,
413 arg,
414 conn_id: take_str(obj, "connId")?,
415 code: take_optional_str(obj, "code"),
416 msg: take_optional_str(obj, "msg"),
417 })
418}
419
420fn parse_channel_conn_count<E: serde::de::Error>(
421 obj: &mut serde_json::Map<String, serde_json::Value>,
422) -> Result<OKXWsFrame, E> {
423 let channel_val = obj
424 .remove("channel")
425 .ok_or_else(|| E::missing_field("channel"))?;
426 let channel: OKXWsChannel = serde_json::from_value(channel_val)
427 .map_err(|e| E::custom(format!("invalid channel: {e}")))?;
428
429 Ok(OKXWsFrame::ChannelConnCount {
430 event: take_str(obj, "event")?,
431 channel,
432 conn_count: take_str(obj, "connCount")?,
433 conn_id: take_str(obj, "connId")?,
434 })
435}
436
437fn parse_order_response<E: serde::de::Error>(
438 obj: &mut serde_json::Map<String, serde_json::Value>,
439) -> Result<OKXWsFrame, E> {
440 let op_val = obj.remove("op").ok_or_else(|| E::missing_field("op"))?;
441 let op: OKXWsOperation =
442 serde_json::from_value(op_val).map_err(|e| E::custom(format!("invalid op: {e}")))?;
443
444 let data: Vec<serde_json::Value> = match obj.remove("data") {
445 Some(v) => {
446 serde_json::from_value(v).map_err(|e| E::custom(format!("invalid data: {e}")))?
447 }
448 None => Vec::new(),
449 };
450
451 Ok(OKXWsFrame::OrderResponse {
452 id: take_optional_str(obj, "id"),
453 op,
454 code: take_str(obj, "code")?,
455 msg: take_str(obj, "msg")?,
456 data,
457 })
458}
459
460fn parse_book_data<E: serde::de::Error>(
461 obj: &mut serde_json::Map<String, serde_json::Value>,
462) -> Result<OKXWsFrame, E> {
463 let arg_val = obj.remove("arg").ok_or_else(|| E::missing_field("arg"))?;
464 let arg: OKXWebSocketArg =
465 serde_json::from_value(arg_val).map_err(|e| E::custom(format!("invalid arg: {e}")))?;
466
467 let action_val = obj
468 .remove("action")
469 .ok_or_else(|| E::missing_field("action"))?;
470 let action: OKXBookAction = serde_json::from_value(action_val)
471 .map_err(|e| E::custom(format!("invalid action: {e}")))?;
472
473 let data_val = obj.remove("data").ok_or_else(|| E::missing_field("data"))?;
474 if arg.channel == OKXWsChannel::BooksRpi {
475 let data: Vec<OKXRpiBookMsg> = serde_json::from_value(data_val)
476 .map_err(|e| E::custom(format!("invalid data: {e}")))?;
477 return Ok(OKXWsFrame::RpiBookData { arg, action, data });
478 }
479
480 let data: Vec<OKXBookMsg> =
481 serde_json::from_value(data_val).map_err(|e| E::custom(format!("invalid data: {e}")))?;
482 Ok(OKXWsFrame::BookData { arg, action, data })
483}
484
485fn parse_data<E: serde::de::Error>(
486 obj: &mut serde_json::Map<String, serde_json::Value>,
487) -> Result<OKXWsFrame, E> {
488 let arg_val = obj.remove("arg").ok_or_else(|| E::missing_field("arg"))?;
489 let arg: OKXWebSocketArg =
490 serde_json::from_value(arg_val).map_err(|e| E::custom(format!("invalid arg: {e}")))?;
491
492 let data = obj.remove("data").ok_or_else(|| E::missing_field("data"))?;
493
494 Ok(OKXWsFrame::Data { arg, data })
495}
496
497fn parse_error<E: serde::de::Error>(
498 obj: &mut serde_json::Map<String, serde_json::Value>,
499) -> Result<OKXWsFrame, E> {
500 let arg = obj
501 .remove("arg")
502 .map(serde_json::from_value)
503 .transpose()
504 .map_err(|e| E::custom(format!("invalid arg: {e}")))?;
505
506 Ok(OKXWsFrame::Error {
507 arg,
508 code: take_str(obj, "code")?,
509 msg: take_str(obj, "msg")?,
510 })
511}
512
513#[derive(Debug, Serialize, Deserialize)]
514#[serde(rename_all = "camelCase")]
515pub struct OKXWebSocketArg {
516 pub channel: OKXWsChannel,
518 #[serde(default, alias = "sprdId")]
522 pub inst_id: Option<Ustr>,
523 #[serde(default)]
524 pub inst_type: Option<OKXInstrumentType>,
525 #[serde(default)]
526 pub inst_family: Option<Ustr>,
527 #[serde(default)]
528 pub bar: Option<Ustr>,
529}
530
531#[derive(Debug, Serialize, Deserialize)]
533#[serde(rename_all = "camelCase")]
534pub struct OKXTickerMsg {
535 pub inst_type: OKXInstrumentType,
537 pub inst_id: Ustr,
539 #[serde(rename = "last")]
541 pub last_px: String,
542 pub last_sz: String,
544 pub ask_px: String,
546 pub ask_sz: String,
548 pub bid_px: String,
550 pub bid_sz: String,
552 pub open24h: String,
554 pub high24h: String,
556 pub low24h: String,
558 pub vol_ccy_24h: String,
560 pub vol24h: String,
562 pub sod_utc0: String,
564 pub sod_utc8: String,
566 #[serde(deserialize_with = "deserialize_string_to_u64")]
568 pub ts: u64,
569 #[serde(default)]
571 pub source: Option<String>,
572}
573
574#[derive(Debug, Serialize, Deserialize)]
576pub struct OrderBookEntry {
577 pub price: String,
579 pub size: String,
581 #[serde(default)]
587 pub liquidated_orders_count: String,
588 #[serde(default)]
590 pub orders_count: String,
591}
592
593#[derive(Debug, Serialize, Deserialize)]
595#[serde(rename_all = "camelCase")]
596pub struct OKXBookMsg {
597 pub asks: Vec<OrderBookEntry>,
599 pub bids: Vec<OrderBookEntry>,
601 pub checksum: Option<i64>,
603 pub prev_seq_id: Option<i64>,
605 pub seq_id: u64,
607 #[serde(deserialize_with = "deserialize_string_to_u64")]
609 pub ts: u64,
610}
611
612#[derive(Debug, Serialize, Deserialize)]
614#[serde(rename_all = "camelCase", deny_unknown_fields)]
615pub struct OKXRpiBookMsg {
616 pub asks: Vec<OKXRpiBookLevel>,
618 pub bids: Vec<OKXRpiBookLevel>,
620 pub prev_seq_id: i64,
622 pub seq_id: u64,
624 #[serde(deserialize_with = "deserialize_string_to_u64")]
626 pub ts: u64,
627}
628
629#[derive(Debug, Serialize, Deserialize)]
631#[serde(rename_all = "camelCase")]
632pub struct OKXTradeMsg {
633 #[serde(default, alias = "sprdId")]
638 pub inst_id: Ustr,
639 pub trade_id: String,
641 pub px: String,
643 pub sz: String,
645 pub side: OKXSide,
647 #[serde(default)]
649 pub count: String,
650 #[serde(deserialize_with = "deserialize_string_to_u64")]
652 pub ts: u64,
653 #[serde(default)]
655 pub source: Option<String>,
656 #[serde(default)]
658 pub seq_id: Option<u64>,
659}
660
661#[derive(Debug, Serialize, Deserialize)]
663#[serde(rename_all = "camelCase")]
664pub struct OKXFundingRateMsg {
665 #[serde(default)]
667 pub inst_type: Option<OKXInstrumentType>,
668 pub inst_id: Ustr,
670 pub funding_rate: Ustr,
672 pub next_funding_rate: Ustr,
674 #[serde(default)]
676 pub min_funding_rate: Option<String>,
677 #[serde(default)]
679 pub max_funding_rate: Option<String>,
680 #[serde(default)]
682 pub sett_state: OKXSettlementState,
683 #[serde(default)]
685 pub sett_funding_rate: Option<String>,
686 #[serde(default)]
688 pub premium: Option<String>,
689 #[serde(default)]
691 pub method: Option<String>,
692 #[serde(deserialize_with = "deserialize_string_to_u64")]
694 pub funding_time: u64,
695 #[serde(deserialize_with = "deserialize_string_to_u64")]
697 pub next_funding_time: u64,
698 #[serde(deserialize_with = "deserialize_string_to_u64")]
700 pub ts: u64,
701}
702
703#[derive(Debug, Serialize, Deserialize)]
705#[serde(rename_all = "camelCase")]
706pub struct OKXMarkPriceMsg {
707 pub inst_id: Ustr,
709 pub mark_px: String,
711 #[serde(deserialize_with = "deserialize_string_to_u64")]
713 pub ts: u64,
714}
715
716#[derive(Debug, Serialize, Deserialize)]
718#[serde(rename_all = "camelCase")]
719pub struct OKXIndexPriceMsg {
720 pub inst_id: Ustr,
722 pub idx_px: String,
724 pub high24h: String,
726 pub low24h: String,
728 pub open24h: String,
730 pub sod_utc0: String,
732 pub sod_utc8: String,
734 #[serde(deserialize_with = "deserialize_string_to_u64")]
736 pub ts: u64,
737}
738
739#[derive(Debug, Serialize, Deserialize)]
741#[serde(rename_all = "camelCase")]
742pub struct OKXPriceLimitMsg {
743 pub inst_id: Ustr,
745 pub buy_lmt: String,
747 pub sell_lmt: String,
749 #[serde(deserialize_with = "deserialize_string_to_u64")]
751 pub ts: u64,
752}
753
754#[derive(Debug, Serialize, Deserialize)]
756#[serde(rename_all = "camelCase")]
757pub struct OKXCandleMsg {
758 #[serde(deserialize_with = "deserialize_string_to_u64")]
760 pub ts: u64,
761 pub o: String,
763 pub h: String,
765 pub l: String,
767 pub c: String,
769 pub vol: String,
771 pub vol_ccy: String,
773 pub vol_ccy_quote: String,
774 pub confirm: OKXCandleConfirm,
776}
777
778#[derive(Debug, Serialize, Deserialize)]
780#[serde(rename_all = "camelCase")]
781pub struct OKXOpenInterestMsg {
782 pub inst_id: Ustr,
784 pub oi: String,
786 pub oi_ccy: String,
788 #[serde(deserialize_with = "deserialize_string_to_u64")]
790 pub ts: u64,
791}
792
793#[derive(Debug, Serialize, Deserialize)]
795#[serde(rename_all = "camelCase")]
796pub struct OKXOptionSummaryMsg {
797 #[serde(default)]
799 pub inst_type: Option<OKXInstrumentType>,
800 pub inst_id: Ustr,
802 pub uly: String,
804 pub delta: String,
806 pub gamma: String,
808 pub theta: String,
810 pub vega: String,
812 #[serde(alias = "deltaBS")]
814 pub delta_bs: String,
815 #[serde(alias = "gammaBS")]
817 pub gamma_bs: String,
818 #[serde(alias = "thetaBS")]
820 pub theta_bs: String,
821 #[serde(alias = "vegaBS")]
823 pub vega_bs: String,
824 pub real_vol: String,
826 pub bid_vol: String,
828 pub ask_vol: String,
830 pub mark_vol: String,
832 pub lever: String,
834 #[serde(default)]
836 pub fwd_px: Option<String>,
837 #[serde(default)]
839 pub mark_px: Option<String>,
840 #[serde(default)]
842 pub vol_lv: Option<String>,
843 #[serde(deserialize_with = "deserialize_string_to_u64")]
845 pub ts: u64,
846}
847
848#[derive(Debug, Serialize, Deserialize)]
850#[serde(rename_all = "camelCase")]
851pub struct OKXEstimatedPriceMsg {
852 pub inst_id: Ustr,
854 pub settle_px: String,
856 #[serde(deserialize_with = "deserialize_string_to_u64")]
858 pub ts: u64,
859}
860
861#[derive(Debug, Serialize, Deserialize)]
863#[serde(rename_all = "camelCase")]
864pub struct OKXStatusMsg {
865 pub title: Ustr,
867 #[serde(rename = "type")]
869 pub status_type: Ustr,
870 pub state: Ustr,
872 pub end_time: Option<String>,
874 pub begin_time: Option<String>,
876 pub service_type: Option<Ustr>,
878 pub reason: Option<String>,
880 #[serde(deserialize_with = "deserialize_string_to_u64")]
882 pub ts: u64,
883}
884
885pub use crate::common::models::OKXAttachedAlgoOrd;
886
887#[derive(Clone, Debug, Serialize, Deserialize)]
893#[serde(rename_all = "camelCase")]
894pub struct OKXLiquidationWarningMsg {
895 pub inst_type: OKXInstrumentType,
897 #[serde(default)]
899 pub inst_family: Option<Ustr>,
900 pub inst_id: Ustr,
902 pub mgn_mode: OKXMarginMode,
904 #[serde(default)]
906 pub pos_id: Option<Ustr>,
907 pub pos_side: OKXPositionSide,
909 pub pos: String,
911 #[serde(default)]
913 pub pos_ccy: Option<Ustr>,
914 pub lever: String,
916 pub mark_px: String,
918 pub mgn_ratio: String,
920 pub ccy: Ustr,
922 #[serde(deserialize_with = "deserialize_string_to_u64")]
924 pub c_time: u64,
925 #[serde(deserialize_with = "deserialize_string_to_u64")]
927 pub u_time: u64,
928 #[serde(default)]
930 pub p_time: Option<String>,
931}
932
933#[derive(Clone, Debug, Default, Serialize, Deserialize)]
935#[serde(rename_all = "camelCase")]
936pub struct OKXLinkedAlgoOrd {
937 #[serde(default)]
939 pub algo_id: String,
940}
941
942#[derive(Clone, Debug, Serialize, Deserialize)]
944#[serde(rename_all = "camelCase")]
945pub struct OKXOrderMsg {
946 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
948 pub acc_fill_sz: Option<String>,
949 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
951 pub algo_id: Option<String>,
952 pub avg_px: String,
954 #[serde(deserialize_with = "deserialize_string_to_u64")]
956 pub c_time: u64,
957 #[serde(default)]
959 pub cancel_source: Option<String>,
960 #[serde(default)]
962 pub cancel_source_reason: Option<String>,
963 pub category: OKXOrderCategory,
965 pub ccy: Ustr,
967 pub cl_ord_id: String,
969 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
971 pub algo_cl_ord_id: Option<String>,
972 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
974 pub attach_algo_cl_ord_id: Option<String>,
975 #[serde(default)]
977 pub attach_algo_ords: Vec<OKXAttachedAlgoOrd>,
978 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
980 pub outcome: Option<String>,
981 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
983 pub fee: Option<String>,
984 pub fee_ccy: Ustr,
986 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
988 pub fill_fee: Option<String>,
989 #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
991 pub fill_fee_ccy: Option<Ustr>,
992 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
994 pub fill_mark_px: Option<String>,
995 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
997 pub fill_mark_vol: Option<String>,
998 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1000 pub fill_px_vol: Option<String>,
1001 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1003 pub fill_px_usd: Option<String>,
1004 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1006 pub fill_fwd_px: Option<String>,
1007 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1009 pub fill_notional_usd: Option<String>,
1010 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1012 pub fill_pnl: Option<String>,
1013 pub fill_px: String,
1015 pub fill_sz: String,
1017 #[serde(deserialize_with = "deserialize_string_to_u64")]
1019 pub fill_time: u64,
1020 pub inst_id: Ustr,
1022 pub inst_type: OKXInstrumentType,
1024 #[serde(default)]
1026 pub is_tp_limit: Option<String>,
1027 pub lever: String,
1029 #[serde(default)]
1031 pub linked_algo_ord: Option<OKXLinkedAlgoOrd>,
1032 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1034 pub notional_usd: Option<String>,
1035 pub ord_id: Ustr,
1037 pub ord_type: OKXOrderType,
1039 pub pnl: String,
1041 pub pos_side: OKXPositionSide,
1043 #[serde(default)]
1045 pub px: String,
1046 #[serde(default)]
1048 pub px_type: OKXPriceType,
1049 #[serde(default)]
1051 pub px_usd: Option<String>,
1052 #[serde(default)]
1054 pub px_vol: Option<String>,
1055 #[serde(default)]
1057 pub quick_mgn_type: OKXQuickMarginType,
1058 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1060 pub rebate: Option<String>,
1061 #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
1063 pub rebate_ccy: Option<Ustr>,
1064 pub reduce_only: String,
1066 pub side: OKXSide,
1068 #[serde(default)]
1070 pub sl_ord_px: Option<String>,
1071 #[serde(default)]
1073 pub sl_trigger_px: Option<String>,
1074 #[serde(default)]
1076 pub sl_trigger_px_type: Option<OKXTriggerType>,
1077 #[serde(default)]
1079 pub source: Option<String>,
1080 pub state: OKXOrderStatus,
1082 #[serde(default)]
1084 pub stp_id: Option<String>,
1085 #[serde(default)]
1087 pub stp_mode: OKXSelfTradePreventionMode,
1088 pub exec_type: OKXExecType,
1090 pub sz: String,
1092 #[serde(default)]
1094 pub tag: Option<String>,
1095 pub td_mode: OKXTradeMode,
1097 #[serde(default, deserialize_with = "deserialize_target_currency_as_none")]
1099 pub tgt_ccy: Option<OKXTargetCurrency>,
1100 #[serde(default)]
1102 pub tp_ord_px: Option<String>,
1103 #[serde(default)]
1105 pub tp_trigger_px: Option<String>,
1106 #[serde(default)]
1108 pub tp_trigger_px_type: Option<OKXTriggerType>,
1109 pub trade_id: String,
1111 #[serde(deserialize_with = "deserialize_string_to_u64")]
1113 pub u_time: u64,
1114 #[serde(default)]
1116 pub amend_result: Option<String>,
1117 #[serde(default)]
1119 pub req_id: Option<String>,
1120 #[serde(default)]
1122 pub code: Option<String>,
1123 #[serde(default)]
1125 pub msg: Option<String>,
1126}
1127
1128#[derive(Clone, Debug, Deserialize, Serialize)]
1130#[serde(rename_all = "camelCase")]
1131pub struct OKXAlgoOrderMsg {
1132 pub algo_id: String,
1134 #[serde(default)]
1136 pub algo_cl_ord_id: String,
1137 pub cl_ord_id: String,
1139 pub ord_id: String,
1141 #[serde(default)]
1143 pub ord_id_list: Vec<String>,
1144 pub inst_id: Ustr,
1146 pub inst_type: OKXInstrumentType,
1148 pub ord_type: OKXAlgoOrderType,
1150 pub state: OKXAlgoOrderStatus,
1152 pub side: OKXSide,
1154 pub pos_side: OKXPositionSide,
1156 #[serde(default)]
1158 pub sz: String,
1159 #[serde(default)]
1161 pub trigger_px: String,
1162 #[serde(default)]
1164 pub trigger_px_type: OKXTriggerType,
1165 #[serde(default)]
1167 pub sl_trigger_px: String,
1168 #[serde(default)]
1170 pub sl_ord_px: String,
1171 #[serde(default)]
1173 pub sl_trigger_px_type: OKXTriggerType,
1174 #[serde(default)]
1176 pub tp_trigger_px: String,
1177 #[serde(default)]
1179 pub tp_ord_px: String,
1180 #[serde(default)]
1182 pub tp_trigger_px_type: OKXTriggerType,
1183 #[serde(default)]
1185 pub ord_px: String,
1186 pub td_mode: OKXTradeMode,
1188 pub lever: String,
1190 #[serde(default)]
1192 pub reduce_only: String,
1193 #[serde(default)]
1195 pub close_fraction: String,
1196 #[serde(default)]
1198 pub actual_px: String,
1199 #[serde(default)]
1201 pub actual_sz: String,
1202 #[serde(default)]
1204 pub notional_usd: String,
1205 #[serde(deserialize_with = "deserialize_string_to_u64")]
1207 pub c_time: u64,
1208 #[serde(deserialize_with = "deserialize_string_to_u64")]
1210 pub u_time: u64,
1211 #[serde(default)]
1213 pub trigger_time: String,
1214 #[serde(default)]
1216 pub fail_code: String,
1217 #[serde(default)]
1219 pub tag: String,
1220 #[serde(default)]
1222 pub callback_ratio: String,
1223 #[serde(default)]
1225 pub callback_spread: String,
1226 #[serde(default)]
1228 pub active_px: String,
1229 #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
1231 pub ccy: Option<Ustr>,
1232 #[serde(default, deserialize_with = "deserialize_target_currency_as_none")]
1234 pub tgt_ccy: Option<OKXTargetCurrency>,
1235 #[serde(default)]
1237 pub fee: Option<String>,
1238 #[serde(default, deserialize_with = "deserialize_empty_ustr_as_none")]
1240 pub fee_ccy: Option<Ustr>,
1241 #[serde(default, deserialize_with = "deserialize_empty_string_as_none")]
1243 pub advance_ord_type: Option<String>,
1244}
1245
1246#[derive(Clone, Debug, Default, Deserialize, Serialize, Builder)]
1248#[builder(default)]
1249#[builder(setter(into, strip_option))]
1250#[serde(rename_all = "camelCase")]
1251pub struct WsAttachAlgoOrdParams {
1252 #[serde(skip_serializing_if = "Option::is_none")]
1254 pub attach_algo_cl_ord_id: Option<String>,
1255 #[serde(skip_serializing_if = "Option::is_none")]
1257 pub sl_trigger_px: Option<String>,
1258 #[serde(skip_serializing_if = "Option::is_none")]
1260 pub sl_ord_px: Option<String>,
1261 #[serde(skip_serializing_if = "Option::is_none")]
1263 pub sl_trigger_px_type: Option<OKXTriggerType>,
1264 #[serde(skip_serializing_if = "Option::is_none")]
1266 pub tp_trigger_px: Option<String>,
1267 #[serde(skip_serializing_if = "Option::is_none")]
1269 pub tp_ord_px: Option<String>,
1270 #[serde(skip_serializing_if = "Option::is_none")]
1272 pub tp_trigger_px_type: Option<OKXTriggerType>,
1273 #[serde(skip_serializing_if = "Option::is_none")]
1275 pub callback_ratio: Option<String>,
1276 #[serde(skip_serializing_if = "Option::is_none")]
1278 pub callback_spread: Option<String>,
1279 #[serde(skip_serializing_if = "Option::is_none")]
1281 pub active_px: Option<String>,
1282 #[serde(skip_serializing_if = "Option::is_none")]
1284 pub new_callback_ratio: Option<String>,
1285 #[serde(skip_serializing_if = "Option::is_none")]
1287 pub new_callback_spread: Option<String>,
1288 #[serde(skip_serializing_if = "Option::is_none")]
1290 pub new_active_px: Option<String>,
1291}
1292
1293#[derive(Clone, Debug, Deserialize, Serialize, Builder)]
1295#[builder(setter(into, strip_option))]
1296#[serde(rename_all = "camelCase")]
1297pub struct WsPostOrderParams {
1298 #[builder(default)]
1300 #[serde(skip_serializing_if = "Option::is_none")]
1301 pub inst_type: Option<OKXInstrumentType>,
1302 pub inst_id_code: u64,
1304 pub td_mode: OKXTradeMode,
1306 #[builder(default)]
1308 #[serde(skip_serializing_if = "Option::is_none")]
1309 pub ccy: Option<Ustr>,
1310 #[builder(default)]
1312 #[serde(skip_serializing_if = "Option::is_none")]
1313 pub cl_ord_id: Option<String>,
1314 pub side: OKXSide,
1316 #[builder(default)]
1318 #[serde(skip_serializing_if = "Option::is_none")]
1319 pub pos_side: Option<OKXPositionSide>,
1320 pub ord_type: OKXOrderType,
1322 pub sz: String,
1324 #[builder(default)]
1326 #[serde(skip_serializing_if = "Option::is_none")]
1327 pub px: Option<String>,
1328 #[builder(default)]
1330 #[serde(rename = "pxUsd", skip_serializing_if = "Option::is_none")]
1331 pub px_usd: Option<String>,
1332 #[builder(default)]
1335 #[serde(rename = "pxVol", skip_serializing_if = "Option::is_none")]
1336 pub px_vol: Option<String>,
1337 #[builder(default)]
1339 #[serde(skip_serializing_if = "Option::is_none")]
1340 pub reduce_only: Option<bool>,
1341 #[builder(default)]
1343 #[serde(rename = "closePosition", skip_serializing_if = "Option::is_none")]
1344 pub close_position: Option<bool>,
1345 #[builder(default)]
1347 #[serde(skip_serializing_if = "Option::is_none")]
1348 pub tgt_ccy: Option<OKXTargetCurrency>,
1349 #[builder(default)]
1351 #[serde(skip_serializing_if = "Option::is_none")]
1352 pub trade_quote_ccy: Option<Ustr>,
1353 #[builder(default)]
1355 #[serde(skip_serializing_if = "Option::is_none")]
1356 pub tag: Option<String>,
1357 #[builder(default)]
1359 #[serde(skip_serializing_if = "Option::is_none")]
1360 pub attach_algo_ords: Option<Vec<WsAttachAlgoOrdParams>>,
1361 #[builder(default)]
1363 #[serde(skip_serializing_if = "Option::is_none")]
1364 pub outcome: Option<String>,
1365 #[builder(default)]
1370 #[serde(skip_serializing_if = "Option::is_none")]
1371 pub slippage_pct: Option<String>,
1372 #[builder(default)]
1374 #[serde(skip_serializing_if = "Option::is_none")]
1375 pub rpi_taker_access: Option<bool>,
1376 #[builder(default)]
1378 #[serde(skip_serializing_if = "Option::is_none")]
1379 pub rpi_px_round: Option<bool>,
1380}
1381
1382#[derive(Clone, Debug, Default, Deserialize, Serialize, Builder)]
1384#[builder(default)]
1385#[builder(setter(into, strip_option))]
1386#[serde(rename_all = "camelCase")]
1387pub struct WsCancelOrderParams {
1388 pub inst_id_code: u64,
1390 #[serde(skip_serializing_if = "Option::is_none")]
1392 pub ord_id: Option<String>,
1393 #[serde(skip_serializing_if = "Option::is_none")]
1395 pub cl_ord_id: Option<String>,
1396}
1397
1398#[derive(Clone, Debug, Default, Deserialize, Serialize, Builder)]
1400#[builder(default)]
1401#[builder(setter(into, strip_option))]
1402#[serde(rename_all = "camelCase")]
1403pub struct WsMassCancelParams {
1404 pub inst_type: OKXInstrumentType,
1406 pub inst_family: Ustr,
1408}
1409
1410#[derive(Clone, Debug, Default, Deserialize, Serialize, Builder)]
1412#[builder(default)]
1413#[builder(setter(into, strip_option))]
1414#[serde(rename_all = "camelCase")]
1415pub struct WsAmendOrderParams {
1416 pub inst_id_code: u64,
1418 #[serde(skip_serializing_if = "Option::is_none")]
1420 pub ord_id: Option<String>,
1421 #[serde(skip_serializing_if = "Option::is_none")]
1423 pub cl_ord_id: Option<String>,
1424 #[serde(skip_serializing_if = "Option::is_none")]
1426 pub req_id: Option<String>,
1427 #[serde(skip_serializing_if = "Option::is_none")]
1429 pub new_px: Option<String>,
1430 #[serde(rename = "newPxUsd", skip_serializing_if = "Option::is_none")]
1432 pub new_px_usd: Option<String>,
1433 #[serde(rename = "newPxVol", skip_serializing_if = "Option::is_none")]
1436 pub new_px_vol: Option<String>,
1437 #[serde(skip_serializing_if = "Option::is_none")]
1439 pub new_sz: Option<String>,
1440 #[serde(skip_serializing_if = "Option::is_none")]
1442 pub rpi_taker_access: Option<bool>,
1443 #[serde(skip_serializing_if = "Option::is_none")]
1445 pub rpi_px_round: Option<bool>,
1446}
1447
1448#[derive(Clone, Debug, Deserialize, Serialize, Builder)]
1450#[builder(setter(into, strip_option))]
1451#[serde(rename_all = "camelCase")]
1452pub struct WsPostAlgoOrderParams {
1453 pub inst_id_code: u64,
1455 pub td_mode: OKXTradeMode,
1457 pub side: OKXSide,
1459 pub ord_type: OKXAlgoOrderType,
1461 pub sz: String,
1463 #[builder(default)]
1465 #[serde(skip_serializing_if = "Option::is_none")]
1466 pub cl_ord_id: Option<String>,
1467 #[builder(default)]
1469 #[serde(skip_serializing_if = "Option::is_none")]
1470 pub pos_side: Option<OKXPositionSide>,
1471 #[serde(skip_serializing_if = "Option::is_none")]
1473 pub trigger_px: Option<String>,
1474 #[builder(default)]
1476 #[serde(skip_serializing_if = "Option::is_none")]
1477 pub trigger_px_type: Option<OKXTriggerType>,
1478 #[builder(default)]
1480 #[serde(skip_serializing_if = "Option::is_none")]
1481 pub order_px: Option<String>,
1482 #[builder(default)]
1484 #[serde(skip_serializing_if = "Option::is_none")]
1485 pub reduce_only: Option<bool>,
1486 #[builder(default)]
1488 #[serde(skip_serializing_if = "Option::is_none")]
1489 pub tag: Option<String>,
1490 #[builder(default)]
1492 #[serde(skip_serializing_if = "Option::is_none")]
1493 pub callback_ratio: Option<String>,
1494 #[builder(default)]
1496 #[serde(skip_serializing_if = "Option::is_none")]
1497 pub callback_spread: Option<String>,
1498 #[builder(default)]
1500 #[serde(skip_serializing_if = "Option::is_none")]
1501 pub active_px: Option<String>,
1502}
1503
1504#[derive(Clone, Debug, Deserialize, Serialize, Builder)]
1506#[builder(setter(into, strip_option))]
1507#[serde(rename_all = "camelCase")]
1508pub struct WsCancelAlgoOrderParams {
1509 pub inst_id_code: u64,
1511 #[builder(default)]
1513 #[serde(skip_serializing_if = "Option::is_none")]
1514 pub algo_id: Option<String>,
1515 #[builder(default)]
1517 #[serde(skip_serializing_if = "Option::is_none")]
1518 pub algo_cl_ord_id: Option<String>,
1519}
1520
1521#[cfg(test)]
1522mod tests {
1523 use nautilus_core::{string::secret::REDACTED, time::get_atomic_clock_realtime};
1524 use rstest::rstest;
1525 use rust_decimal::Decimal;
1526
1527 use super::*;
1528
1529 #[rstest]
1530 fn authentication_preserves_wire_values_and_redacts_debug() {
1531 let authentication = OKXAuthentication {
1532 op: "login",
1533 args: vec![OKXAuthenticationArg {
1534 api_key: SecretString::from("api-key-value"),
1535 passphrase: SecretString::from("passphrase-value"),
1536 timestamp: "1700000000".to_string(),
1537 sign: SecretString::from("signature-value"),
1538 }],
1539 };
1540
1541 let json = serde_json::to_value(&authentication).unwrap();
1542 let formatted = format!("{authentication:?}");
1543
1544 assert_eq!(json["op"], "login");
1545 assert_eq!(json["args"][0]["apiKey"], "api-key-value");
1546 assert_eq!(json["args"][0]["passphrase"], "passphrase-value");
1547 assert_eq!(json["args"][0]["sign"], "signature-value");
1548 assert!(formatted.contains(REDACTED));
1549 assert!(!formatted.contains("api-key-value"));
1550 assert!(!formatted.contains("passphrase-value"));
1551 assert!(!formatted.contains("signature-value"));
1552 }
1553 use crate::common::testing::load_test_json;
1554
1555 #[rstest]
1556 fn test_deserialize_websocket_arg() {
1557 let json_str = r#"{"channel":"instruments","instType":"SPOT"}"#;
1558
1559 let result: Result<OKXWebSocketArg, _> = serde_json::from_str(json_str);
1560 match result {
1561 Ok(arg) => {
1562 assert_eq!(arg.channel, OKXWsChannel::Instruments);
1563 assert_eq!(arg.inst_type, Some(OKXInstrumentType::Spot));
1564 assert_eq!(arg.inst_id, None);
1565 }
1566 Err(e) => {
1567 panic!("Failed to deserialize WebSocket arg: {e}");
1568 }
1569 }
1570 }
1571
1572 #[rstest]
1573 fn test_deserialize_subscribe_variant_direct() {
1574 #[derive(Debug, Deserialize)]
1575 #[serde(rename_all = "camelCase")]
1576 struct SubscribeMsg {
1577 event: String,
1578 arg: OKXWebSocketArg,
1579 conn_id: String,
1580 }
1581
1582 let json_str = r#"{"event":"subscribe","arg":{"channel":"instruments","instType":"SPOT"},"connId":"380cfa6a"}"#;
1583
1584 let result: Result<SubscribeMsg, _> = serde_json::from_str(json_str);
1585 match result {
1586 Ok(msg) => {
1587 assert_eq!(msg.event, "subscribe");
1588 assert_eq!(msg.arg.channel, OKXWsChannel::Instruments);
1589 assert_eq!(msg.conn_id, "380cfa6a");
1590 }
1591 Err(e) => {
1592 panic!("Failed to deserialize subscribe message directly: {e}");
1593 }
1594 }
1595 }
1596
1597 #[rstest]
1598 fn test_deserialize_subscribe_confirmation() {
1599 let json_str = r#"{"event":"subscribe","arg":{"channel":"instruments","instType":"SPOT"},"connId":"380cfa6a"}"#;
1600
1601 let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1602 match result {
1603 Ok(msg) => {
1604 if let OKXWsFrame::Subscription {
1605 event,
1606 arg,
1607 conn_id,
1608 ..
1609 } = msg
1610 {
1611 assert_eq!(event, OKXSubscriptionEvent::Subscribe);
1612 assert_eq!(arg.channel, OKXWsChannel::Instruments);
1613 assert_eq!(conn_id, "380cfa6a");
1614 } else {
1615 panic!("Expected Subscribe variant, was: {msg:?}");
1616 }
1617 }
1618 Err(e) => {
1619 panic!("Failed to deserialize subscription confirmation: {e}");
1620 }
1621 }
1622 }
1623
1624 #[rstest]
1625 fn test_deserialize_subscribe_with_inst_id() {
1626 let json_str = r#"{"event":"subscribe","arg":{"channel":"candle1m","instId":"ETH-USDT"},"connId":"358602f5"}"#;
1627
1628 let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1629 match result {
1630 Ok(msg) => {
1631 if let OKXWsFrame::Subscription {
1632 event,
1633 arg,
1634 conn_id,
1635 ..
1636 } = msg
1637 {
1638 assert_eq!(event, OKXSubscriptionEvent::Subscribe);
1639 assert_eq!(arg.channel, OKXWsChannel::Candle1Minute);
1640 assert_eq!(conn_id, "358602f5");
1641 } else {
1642 panic!("Expected Subscribe variant, was: {msg:?}");
1643 }
1644 }
1645 Err(e) => {
1646 panic!("Failed to deserialize subscription confirmation: {e}");
1647 }
1648 }
1649 }
1650
1651 #[rstest]
1652 fn test_channel_serialization_for_logging() {
1653 let channel = OKXWsChannel::Candle1Minute;
1654 let serialized = serde_json::to_string(&channel).unwrap();
1655 let cleaned = serialized.trim_matches('"').to_string();
1656 assert_eq!(cleaned, "candle1m");
1657
1658 let channel = OKXWsChannel::BboTbt;
1659 let serialized = serde_json::to_string(&channel).unwrap();
1660 let cleaned = serialized.trim_matches('"').to_string();
1661 assert_eq!(cleaned, "bbo-tbt");
1662
1663 let channel = OKXWsChannel::Trades;
1664 let serialized = serde_json::to_string(&channel).unwrap();
1665 let cleaned = serialized.trim_matches('"').to_string();
1666 assert_eq!(cleaned, "trades");
1667 }
1668
1669 #[rstest]
1670 fn test_order_response_with_enum_operation() {
1671 let json_str = r#"{"id":"req-123","op":"order","code":"0","msg":"","data":[]}"#;
1672 let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1673 match result {
1674 Ok(OKXWsFrame::OrderResponse {
1675 id,
1676 op,
1677 code,
1678 msg,
1679 data,
1680 }) => {
1681 assert_eq!(id, Some("req-123".to_string()));
1682 assert_eq!(op, OKXWsOperation::Order);
1683 assert_eq!(code, "0");
1684 assert_eq!(msg, "");
1685 assert!(data.is_empty());
1686 }
1687 Ok(other) => panic!("Expected OrderResponse, was: {other:?}"),
1688 Err(e) => panic!("Failed to deserialize: {e}"),
1689 }
1690
1691 let json_str = r#"{"id":"cancel-456","op":"cancel-order","code":"50001","msg":"Order not found","data":[]}"#;
1692 let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1693 match result {
1694 Ok(OKXWsFrame::OrderResponse {
1695 id,
1696 op,
1697 code,
1698 msg,
1699 data,
1700 }) => {
1701 assert_eq!(id, Some("cancel-456".to_string()));
1702 assert_eq!(op, OKXWsOperation::CancelOrder);
1703 assert_eq!(code, "50001");
1704 assert_eq!(msg, "Order not found");
1705 assert!(data.is_empty());
1706 }
1707 Ok(other) => panic!("Expected OrderResponse, was: {other:?}"),
1708 Err(e) => panic!("Failed to deserialize: {e}"),
1709 }
1710
1711 let json_str = r#"{"id":"amend-789","op":"amend-order","code":"50002","msg":"Invalid price","data":[]}"#;
1712 let result: Result<OKXWsFrame, _> = serde_json::from_str(json_str);
1713 match result {
1714 Ok(OKXWsFrame::OrderResponse {
1715 id,
1716 op,
1717 code,
1718 msg,
1719 data,
1720 }) => {
1721 assert_eq!(id, Some("amend-789".to_string()));
1722 assert_eq!(op, OKXWsOperation::AmendOrder);
1723 assert_eq!(code, "50002");
1724 assert_eq!(msg, "Invalid price");
1725 assert!(data.is_empty());
1726 }
1727 Ok(other) => panic!("Expected OrderResponse, was: {other:?}"),
1728 Err(e) => panic!("Failed to deserialize: {e}"),
1729 }
1730 }
1731
1732 #[rstest]
1733 fn test_operation_enum_serialization() {
1734 let op = OKXWsOperation::Order;
1735 let serialized = serde_json::to_string(&op).unwrap();
1736 assert_eq!(serialized, "\"order\"");
1737
1738 let op = OKXWsOperation::CancelOrder;
1739 let serialized = serde_json::to_string(&op).unwrap();
1740 assert_eq!(serialized, "\"cancel-order\"");
1741
1742 let op = OKXWsOperation::AmendOrder;
1743 let serialized = serde_json::to_string(&op).unwrap();
1744 assert_eq!(serialized, "\"amend-order\"");
1745
1746 let op = OKXWsOperation::Subscribe;
1747 let serialized = serde_json::to_string(&op).unwrap();
1748 assert_eq!(serialized, "\"subscribe\"");
1749 }
1750
1751 #[rstest]
1752 fn test_order_response_parsing() {
1753 let success_response = r#"{
1754 "id": "req-123",
1755 "op": "order",
1756 "code": "0",
1757 "msg": "",
1758 "data": [{"sMsg": "Order placed successfully"}]
1759 }"#;
1760
1761 let parsed: OKXWsFrame = serde_json::from_str(success_response).unwrap();
1762
1763 match parsed {
1764 OKXWsFrame::OrderResponse {
1765 id,
1766 op,
1767 code,
1768 msg,
1769 data,
1770 } => {
1771 assert_eq!(id, Some("req-123".to_string()));
1772 assert_eq!(op, OKXWsOperation::Order);
1773 assert_eq!(code, "0");
1774 assert_eq!(msg, "");
1775 assert_eq!(data.len(), 1);
1776 }
1777 _ => panic!("Expected OrderResponse variant"),
1778 }
1779
1780 let failure_response = r#"{
1781 "id": "req-456",
1782 "op": "cancel-order",
1783 "code": "50001",
1784 "msg": "Order not found",
1785 "data": [{"sMsg": "Order with client order ID not found"}]
1786 }"#;
1787
1788 let parsed: OKXWsFrame = serde_json::from_str(failure_response).unwrap();
1789
1790 match parsed {
1791 OKXWsFrame::OrderResponse {
1792 id,
1793 op,
1794 code,
1795 msg,
1796 data,
1797 } => {
1798 assert_eq!(id, Some("req-456".to_string()));
1799 assert_eq!(op, OKXWsOperation::CancelOrder);
1800 assert_eq!(code, "50001");
1801 assert_eq!(msg, "Order not found");
1802 assert_eq!(data.len(), 1);
1803 }
1804 _ => panic!("Expected OrderResponse variant"),
1805 }
1806 }
1807
1808 #[rstest]
1809 fn test_subscription_event_parsing() {
1810 let subscription_json = r#"{
1811 "event": "subscribe",
1812 "arg": {
1813 "channel": "tickers",
1814 "instId": "BTC-USDT"
1815 },
1816 "connId": "a4d3ae55"
1817 }"#;
1818
1819 let parsed: OKXWsFrame = serde_json::from_str(subscription_json).unwrap();
1820
1821 match parsed {
1822 OKXWsFrame::Subscription {
1823 event,
1824 arg,
1825 conn_id,
1826 ..
1827 } => {
1828 assert_eq!(
1829 event,
1830 crate::websocket::enums::OKXSubscriptionEvent::Subscribe
1831 );
1832 assert_eq!(arg.channel, OKXWsChannel::Tickers);
1833 assert_eq!(arg.inst_id, Some(Ustr::from("BTC-USDT")));
1834 assert_eq!(conn_id, "a4d3ae55");
1835 }
1836 _ => panic!("Expected Subscription variant"),
1837 }
1838 }
1839
1840 #[rstest]
1841 fn test_login_event_parsing() {
1842 let login_success = r#"{
1843 "event": "login",
1844 "code": "0",
1845 "msg": "Login successful",
1846 "connId": "a4d3ae55"
1847 }"#;
1848
1849 let parsed: OKXWsFrame = serde_json::from_str(login_success).unwrap();
1850
1851 match parsed {
1852 OKXWsFrame::Login {
1853 event,
1854 code,
1855 msg,
1856 conn_id,
1857 } => {
1858 assert_eq!(event, "login");
1859 assert_eq!(code, "0");
1860 assert_eq!(msg, "Login successful");
1861 assert_eq!(conn_id, "a4d3ae55");
1862 }
1863 _ => panic!("Expected Login variant, was: {parsed:?}"),
1864 }
1865 }
1866
1867 #[rstest]
1868 fn test_error_event_parsing() {
1869 let error_json = r#"{
1870 "code": "60012",
1871 "msg": "Invalid request"
1872 }"#;
1873
1874 let parsed: OKXWsFrame = serde_json::from_str(error_json).unwrap();
1875
1876 match parsed {
1877 OKXWsFrame::Error { arg, code, msg } => {
1878 assert!(arg.is_none());
1879 assert_eq!(code, "60012");
1880 assert_eq!(msg, "Invalid request");
1881 }
1882 _ => panic!("Expected Error variant"),
1883 }
1884 }
1885
1886 #[rstest]
1887 fn test_error_event_with_event_field_parsing() {
1888 let error_json = r#"{
1890 "event": "error",
1891 "code": "60018",
1892 "msg": "Invalid sign"
1893 }"#;
1894
1895 let parsed: OKXWsFrame = serde_json::from_str(error_json).unwrap();
1896
1897 match parsed {
1898 OKXWsFrame::Error { arg, code, msg } => {
1899 assert!(arg.is_none());
1900 assert_eq!(code, "60018");
1901 assert_eq!(msg, "Invalid sign");
1902 }
1903 _ => panic!("Expected Error variant, was: {parsed:?}"),
1904 }
1905 }
1906
1907 #[rstest]
1908 fn test_subscription_error_with_arg_field_parsing() {
1909 let error_json = r#"{
1911 "event": "error",
1912 "arg": {"channel": "tickers", "instId": "INVALID-INST"},
1913 "code": "60012",
1914 "msg": "Invalid request: channel not found",
1915 "connId": "a4d3ae55"
1916 }"#;
1917
1918 let parsed: OKXWsFrame = serde_json::from_str(error_json).unwrap();
1919
1920 match parsed {
1921 OKXWsFrame::Error { arg, code, msg } => {
1922 let arg = arg.expect("subscription error arg");
1923 assert_eq!(arg.channel, OKXWsChannel::Tickers);
1924 assert_eq!(arg.inst_id, Some(Ustr::from("INVALID-INST")));
1925 assert_eq!(code, "60012");
1926 assert_eq!(msg, "Invalid request: channel not found");
1927 }
1928 _ => panic!("Expected Error variant, was: {parsed:?}"),
1929 }
1930 }
1931
1932 #[rstest]
1933 fn test_websocket_request_serialization() {
1934 let request = OKXWsRequest {
1935 id: Some("req-123".to_string()),
1936 op: OKXWsOperation::Order,
1937 args: vec![serde_json::json!({
1938 "instId": "BTC-USDT",
1939 "tdMode": "cash",
1940 "side": "buy",
1941 "ordType": "market",
1942 "sz": "0.1"
1943 })],
1944 exp_time: None,
1945 };
1946
1947 let serialized = serde_json::to_string(&request).unwrap();
1948 let parsed: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1949
1950 assert_eq!(parsed["id"], "req-123");
1951 assert_eq!(parsed["op"], "order");
1952 assert!(parsed["args"].is_array());
1953 assert_eq!(parsed["args"].as_array().unwrap().len(), 1);
1954 }
1955
1956 #[rstest]
1957 fn test_subscription_request_serialization() {
1958 let subscription = OKXSubscription {
1959 op: OKXWsOperation::Subscribe,
1960 args: vec![OKXSubscriptionArg {
1961 channel: OKXWsChannel::Tickers,
1962 inst_type: Some(OKXInstrumentType::Spot),
1963 inst_family: None,
1964 inst_id: Some(Ustr::from("BTC-USDT")),
1965 }],
1966 };
1967
1968 let serialized = serde_json::to_string(&subscription).unwrap();
1969 let parsed: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1970
1971 assert_eq!(parsed["op"], "subscribe");
1972 assert!(parsed["args"].is_array());
1973 assert_eq!(parsed["args"][0]["channel"], "tickers");
1974 assert_eq!(parsed["args"][0]["instType"], "SPOT");
1975 assert_eq!(parsed["args"][0]["instId"], "BTC-USDT");
1976 }
1977
1978 #[rstest]
1979 fn test_error_message_extraction() {
1980 let responses = vec![
1981 (
1982 r#"{
1983 "id": "req-123",
1984 "op": "order",
1985 "code": "50001",
1986 "msg": "Order failed",
1987 "data": [{"sMsg": "Insufficient balance"}]
1988 }"#,
1989 "Insufficient balance",
1990 ),
1991 (
1992 r#"{
1993 "id": "req-456",
1994 "op": "cancel-order",
1995 "code": "50002",
1996 "msg": "Cancel failed",
1997 "data": [{}]
1998 }"#,
1999 "Cancel failed",
2000 ),
2001 ];
2002
2003 for (response_json, expected_msg) in responses {
2004 let parsed: OKXWsFrame = serde_json::from_str(response_json).unwrap();
2005
2006 match parsed {
2007 OKXWsFrame::OrderResponse {
2008 id: _,
2009 op: _,
2010 code,
2011 msg,
2012 data,
2013 } => {
2014 assert_ne!(code, "0"); let error_msg = data
2018 .first()
2019 .and_then(|d| d.get("sMsg"))
2020 .and_then(|s| s.as_str())
2021 .filter(|s| !s.is_empty())
2022 .unwrap_or(&msg);
2023
2024 assert_eq!(error_msg, expected_msg);
2025 }
2026 _ => panic!("Expected OrderResponse variant"),
2027 }
2028 }
2029 }
2030
2031 #[rstest]
2032 fn test_book_data_parsing() {
2033 let book_data_json = r#"{
2034 "arg": {
2035 "channel": "books",
2036 "instId": "BTC-USDT"
2037 },
2038 "action": "snapshot",
2039 "data": [{
2040 "asks": [["50000.0", "0.1", "0", "1"]],
2041 "bids": [["49999.0", "0.2", "0", "1"]],
2042 "ts": "1640995200000",
2043 "checksum": 123456789,
2044 "seqId": 1000
2045 }]
2046 }"#;
2047
2048 let parsed: OKXWsFrame = serde_json::from_str(book_data_json).unwrap();
2049
2050 match parsed {
2051 OKXWsFrame::BookData { arg, action, data } => {
2052 assert_eq!(arg.channel, OKXWsChannel::Books);
2053 assert_eq!(arg.inst_id, Some(Ustr::from("BTC-USDT")));
2054 assert_eq!(
2055 action,
2056 super::super::super::common::enums::OKXBookAction::Snapshot
2057 );
2058 assert_eq!(data.len(), 1);
2059 }
2060 _ => panic!("Expected BookData variant"),
2061 }
2062 }
2063
2064 #[rstest]
2065 fn test_rpi_book_fixtures_preserve_depth_types_and_sequence() {
2066 let snapshot: OKXWsFrame =
2067 serde_json::from_str(&load_test_json("ws_books_rpi_snapshot.json")).unwrap();
2068 let update: OKXWsFrame =
2069 serde_json::from_str(&load_test_json("ws_books_rpi_update.json")).unwrap();
2070
2071 let OKXWsFrame::RpiBookData { arg, action, data } = snapshot else {
2072 panic!("Expected RPI book snapshot");
2073 };
2074 let snapshot = &data[0];
2075 assert_eq!(arg.channel, OKXWsChannel::BooksRpi);
2076 assert_eq!(arg.inst_id, Some(Ustr::from("OMI-USD")));
2077 assert_eq!(action, OKXBookAction::Snapshot);
2078 assert_eq!(data.len(), 1);
2079 assert_eq!(snapshot.asks.len(), 4);
2080 assert_eq!(snapshot.bids.len(), 10);
2081 assert_eq!(
2082 snapshot.asks[0],
2083 OKXRpiBookLevel(
2084 Decimal::from_str_exact("0.0001617").unwrap(),
2085 Decimal::from_str_exact("12325166.992").unwrap(),
2086 Decimal::from(1000),
2087 2,
2088 )
2089 );
2090 assert_eq!(snapshot.prev_seq_id, -1);
2091 assert_eq!(snapshot.seq_id, 1_082_831_226);
2092 assert_eq!(snapshot.ts, 1_785_406_442_403);
2093
2094 let OKXWsFrame::RpiBookData { arg, action, data } = update else {
2095 panic!("Expected RPI book update");
2096 };
2097 let update = &data[0];
2098 assert_eq!(arg.channel, OKXWsChannel::BooksRpi);
2099 assert_eq!(arg.inst_id, Some(Ustr::from("OMI-USD")));
2100 assert_eq!(action, OKXBookAction::Update);
2101 assert_eq!(data.len(), 1);
2102 assert_eq!(update.asks.len(), 2);
2103 assert!(update.bids.is_empty());
2104 assert_eq!(
2105 update.asks[1],
2106 OKXRpiBookLevel(
2107 Decimal::from_str_exact("0.0001625").unwrap(),
2108 Decimal::from_str_exact("12324367.786").unwrap(),
2109 Decimal::from(1000),
2110 2,
2111 )
2112 );
2113 assert_eq!(update.prev_seq_id, snapshot.seq_id as i64);
2114 assert_eq!(update.seq_id, 1_082_831_230);
2115 assert_eq!(update.ts, 1_785_406_443_903);
2116 }
2117
2118 #[rstest]
2119 fn test_rpi_book_rejects_checksum_field() {
2120 let mut payload: serde_json::Value =
2121 serde_json::from_str(&load_test_json("ws_books_rpi_update.json")).unwrap();
2122 payload["data"][0]["checksum"] = serde_json::json!(0);
2123
2124 let error = serde_json::from_value::<OKXWsFrame>(payload).unwrap_err();
2125
2126 assert!(error.to_string().contains("checksum"));
2127 }
2128
2129 #[rstest]
2130 fn test_data_event_parsing() {
2131 let data_json = r#"{
2132 "arg": {
2133 "channel": "trades",
2134 "instId": "BTC-USDT"
2135 },
2136 "data": [{
2137 "instId": "BTC-USDT",
2138 "tradeId": "12345",
2139 "px": "50000.0",
2140 "sz": "0.1",
2141 "side": "buy",
2142 "ts": "1640995200000"
2143 }]
2144 }"#;
2145
2146 let parsed: OKXWsFrame = serde_json::from_str(data_json).unwrap();
2147
2148 match parsed {
2149 OKXWsFrame::Data { arg, data } => {
2150 assert_eq!(arg.channel, OKXWsChannel::Trades);
2151 assert_eq!(arg.inst_id, Some(Ustr::from("BTC-USDT")));
2152 assert!(data.is_array());
2153 }
2154 _ => panic!("Expected Data variant"),
2155 }
2156 }
2157
2158 #[rstest]
2159 fn test_nautilus_message_variants() {
2160 let clock = get_atomic_clock_realtime();
2161 let ts_init = clock.get_time_ns();
2162
2163 let error = OKXWebSocketError {
2164 code: "60012".to_string(),
2165 message: "Invalid request".to_string(),
2166 conn_id: None,
2167 timestamp: ts_init.as_u64(),
2168 };
2169 let error_msg = NautilusWsMessage::Error(error);
2170
2171 match error_msg {
2172 NautilusWsMessage::Error(e) => {
2173 assert_eq!(e.code, "60012");
2174 assert_eq!(e.message, "Invalid request");
2175 }
2176 _ => panic!("Expected Error variant"),
2177 }
2178
2179 let raw_scenarios = vec![
2180 ::serde_json::json!({"unknown": "data"}),
2181 ::serde_json::json!({"channel": "unsupported", "data": [1, 2, 3]}),
2182 ::serde_json::json!({"complex": {"nested": {"structure": true}}}),
2183 ];
2184
2185 for raw_data in raw_scenarios {
2186 let raw_msg = NautilusWsMessage::Raw(raw_data.clone());
2187
2188 match raw_msg {
2189 NautilusWsMessage::Raw(data) => {
2190 assert_eq!(data, raw_data);
2191 }
2192 _ => panic!("Expected Raw variant"),
2193 }
2194 }
2195 }
2196
2197 #[rstest]
2198 fn test_order_response_parsing_success() {
2199 let order_response_json = r#"{
2200 "id": "req-123",
2201 "op": "order",
2202 "code": "0",
2203 "msg": "",
2204 "data": [{"sMsg": "Order placed successfully"}]
2205 }"#;
2206
2207 let parsed: OKXWsFrame = serde_json::from_str(order_response_json).unwrap();
2208
2209 match parsed {
2210 OKXWsFrame::OrderResponse {
2211 id,
2212 op,
2213 code,
2214 msg,
2215 data,
2216 } => {
2217 assert_eq!(id, Some("req-123".to_string()));
2218 assert_eq!(op, OKXWsOperation::Order);
2219 assert_eq!(code, "0");
2220 assert_eq!(msg, "");
2221 assert_eq!(data.len(), 1);
2222 }
2223 _ => panic!("Expected OrderResponse variant"),
2224 }
2225 }
2226
2227 #[rstest]
2228 fn test_order_response_parsing_failure() {
2229 let order_response_json = r#"{
2230 "id": "req-456",
2231 "op": "cancel-order",
2232 "code": "50001",
2233 "msg": "Order not found",
2234 "data": [{"sMsg": "Order with client order ID not found"}]
2235 }"#;
2236
2237 let parsed: OKXWsFrame = serde_json::from_str(order_response_json).unwrap();
2238
2239 match parsed {
2240 OKXWsFrame::OrderResponse {
2241 id,
2242 op,
2243 code,
2244 msg,
2245 data,
2246 } => {
2247 assert_eq!(id, Some("req-456".to_string()));
2248 assert_eq!(op, OKXWsOperation::CancelOrder);
2249 assert_eq!(code, "50001");
2250 assert_eq!(msg, "Order not found");
2251 assert_eq!(data.len(), 1);
2252 }
2253 _ => panic!("Expected OrderResponse variant"),
2254 }
2255 }
2256
2257 #[rstest]
2258 fn test_message_request_serialization() {
2259 let request = OKXWsRequest {
2260 id: Some("req-123".to_string()),
2261 op: OKXWsOperation::Order,
2262 args: vec![::serde_json::json!({
2263 "instId": "BTC-USDT",
2264 "tdMode": "cash",
2265 "side": "buy",
2266 "ordType": "market",
2267 "sz": "0.1"
2268 })],
2269 exp_time: None,
2270 };
2271
2272 let serialized = serde_json::to_string(&request).unwrap();
2273 let parsed: serde_json::Value = serde_json::from_str(&serialized).unwrap();
2274
2275 assert_eq!(parsed["id"], "req-123");
2276 assert_eq!(parsed["op"], "order");
2277 assert!(parsed["args"].is_array());
2278 assert_eq!(parsed["args"].as_array().unwrap().len(), 1);
2279 }
2280
2281 #[rstest]
2282 fn test_ws_post_order_params_serializes_inst_id_code() {
2283 use super::WsPostOrderParamsBuilder;
2284 use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2285
2286 let params = WsPostOrderParamsBuilder::default()
2287 .inst_id_code(10459u64)
2288 .td_mode(OKXTradeMode::Cross)
2289 .side(OKXSide::Buy)
2290 .ord_type(OKXOrderType::Limit)
2291 .sz("0.01".to_string())
2292 .px("50000".to_string())
2293 .build()
2294 .unwrap();
2295
2296 let json = serde_json::to_string(¶ms).unwrap();
2297
2298 assert!(json.contains("\"instIdCode\":10459"));
2299 assert!(!json.contains("\"instId\""));
2300 }
2301
2302 #[rstest]
2303 fn test_ws_post_order_params_serializes_slippage_pct() {
2304 use super::WsPostOrderParamsBuilder;
2305 use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2306
2307 let params = WsPostOrderParamsBuilder::default()
2308 .inst_id_code(10459u64)
2309 .td_mode(OKXTradeMode::Cross)
2310 .side(OKXSide::Buy)
2311 .ord_type(OKXOrderType::Market)
2312 .sz("0.01".to_string())
2313 .slippage_pct("0.005".to_string())
2314 .build()
2315 .unwrap();
2316
2317 let json: serde_json::Value = serde_json::to_value(¶ms).unwrap();
2318 assert_eq!(json["slippagePct"], "0.005");
2319 }
2320
2321 #[rstest]
2322 fn test_ws_post_order_params_omits_slippage_pct_when_unset() {
2323 use super::WsPostOrderParamsBuilder;
2324 use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2325
2326 let params = WsPostOrderParamsBuilder::default()
2327 .inst_id_code(10459u64)
2328 .td_mode(OKXTradeMode::Cross)
2329 .side(OKXSide::Buy)
2330 .ord_type(OKXOrderType::Market)
2331 .sz("0.01".to_string())
2332 .build()
2333 .unwrap();
2334
2335 let json = serde_json::to_string(¶ms).unwrap();
2336 assert!(!json.contains("slippagePct"));
2337 assert!(!json.contains("tradeQuoteCcy"));
2338 }
2339
2340 #[rstest]
2341 fn test_ws_post_order_params_serializes_trade_quote_ccy_usd() {
2342 use super::WsPostOrderParamsBuilder;
2343 use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2344
2345 let params = WsPostOrderParamsBuilder::default()
2346 .inst_id_code(20459u64)
2347 .td_mode(OKXTradeMode::Cash)
2348 .side(OKXSide::Buy)
2349 .ord_type(OKXOrderType::Limit)
2350 .sz("0.01".to_string())
2351 .px("100000".to_string())
2352 .trade_quote_ccy("USD")
2353 .build()
2354 .unwrap();
2355
2356 let json: serde_json::Value = serde_json::to_value(¶ms).unwrap();
2357 assert_eq!(json["instIdCode"], 20459);
2358 assert_eq!(json["tradeQuoteCcy"], "USD");
2359 assert!(json.get("instId").is_none());
2360 }
2361
2362 #[rstest]
2363 fn test_ws_post_order_params_serializes_trade_quote_ccy_usdc() {
2364 use super::WsPostOrderParamsBuilder;
2365 use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2366
2367 let params = WsPostOrderParamsBuilder::default()
2368 .inst_id_code(20459u64)
2369 .td_mode(OKXTradeMode::Cash)
2370 .side(OKXSide::Buy)
2371 .ord_type(OKXOrderType::Limit)
2372 .sz("0.01".to_string())
2373 .px("100000".to_string())
2374 .trade_quote_ccy("USDC")
2375 .build()
2376 .unwrap();
2377
2378 let json: serde_json::Value = serde_json::to_value(¶ms).unwrap();
2379 assert_eq!(json["tradeQuoteCcy"], "USDC");
2380 }
2381
2382 #[rstest]
2383 fn test_ws_post_order_params_serializes_attached_tp_sl() {
2384 use super::{WsAttachAlgoOrdParamsBuilder, WsPostOrderParamsBuilder};
2385 use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode, OKXTriggerType};
2386
2387 let params = WsPostOrderParamsBuilder::default()
2388 .inst_id_code(10459u64)
2389 .td_mode(OKXTradeMode::Cross)
2390 .side(OKXSide::Buy)
2391 .ord_type(OKXOrderType::Limit)
2392 .sz("0.01".to_string())
2393 .px("50000".to_string())
2394 .attach_algo_ords(vec![
2395 WsAttachAlgoOrdParamsBuilder::default()
2396 .attach_algo_cl_ord_id("O-bracket-sl")
2397 .sl_trigger_px("39000")
2398 .sl_ord_px("-1")
2399 .sl_trigger_px_type(OKXTriggerType::Last)
2400 .build()
2401 .unwrap(),
2402 WsAttachAlgoOrdParamsBuilder::default()
2403 .attach_algo_cl_ord_id("O-bracket-tp")
2404 .tp_trigger_px("41000")
2405 .tp_ord_px("-1")
2406 .tp_trigger_px_type(OKXTriggerType::Last)
2407 .build()
2408 .unwrap(),
2409 ])
2410 .build()
2411 .unwrap();
2412
2413 let json = serde_json::to_string(¶ms).unwrap();
2414
2415 assert!(json.contains("\"attachAlgoOrds\""));
2416 assert!(json.contains("\"attachAlgoClOrdId\":\"O-bracket-sl\""));
2417 assert!(json.contains("\"slTriggerPx\":\"39000\""));
2418 assert!(json.contains("\"slOrdPx\":\"-1\""));
2419 assert!(json.contains("\"attachAlgoClOrdId\":\"O-bracket-tp\""));
2420 assert!(json.contains("\"tpTriggerPx\":\"41000\""));
2421 assert!(json.contains("\"tpOrdPx\":\"-1\""));
2422 }
2423
2424 #[rstest]
2425 fn test_ws_cancel_order_params_serializes_inst_id_code() {
2426 use super::WsCancelOrderParamsBuilder;
2427
2428 let params = WsCancelOrderParamsBuilder::default()
2429 .inst_id_code(10461u64)
2430 .ord_id("12345678".to_string())
2431 .build()
2432 .unwrap();
2433
2434 let json = serde_json::to_string(¶ms).unwrap();
2435
2436 assert!(json.contains("\"instIdCode\":10461"));
2437 assert!(!json.contains("\"instId\""));
2438 assert!(json.contains("\"ordId\":\"12345678\""));
2439 }
2440
2441 #[rstest]
2442 fn test_ws_amend_order_params_serializes_inst_id_code() {
2443 use super::WsAmendOrderParamsBuilder;
2444
2445 let params = WsAmendOrderParamsBuilder::default()
2446 .inst_id_code(10459u64)
2447 .cl_ord_id("client123".to_string())
2448 .new_px("51000".to_string())
2449 .build()
2450 .unwrap();
2451
2452 let json = serde_json::to_string(¶ms).unwrap();
2453
2454 assert!(json.contains("\"instIdCode\":10459"));
2455 assert!(!json.contains("\"instId\""));
2456 assert!(json.contains("\"newPx\":\"51000\""));
2457 }
2458
2459 #[rstest]
2460 fn test_ws_post_algo_order_params_serializes_inst_id_code() {
2461 use super::WsPostAlgoOrderParamsBuilder;
2462 use crate::common::enums::{OKXAlgoOrderType, OKXSide, OKXTradeMode, OKXTriggerType};
2463
2464 let params = WsPostAlgoOrderParamsBuilder::default()
2465 .inst_id_code(10459u64)
2466 .td_mode(OKXTradeMode::Cross)
2467 .side(OKXSide::Buy)
2468 .ord_type(OKXAlgoOrderType::Trigger)
2469 .sz("0.01".to_string())
2470 .trigger_px("48000".to_string())
2471 .trigger_px_type(OKXTriggerType::Last)
2472 .build()
2473 .unwrap();
2474
2475 let json = serde_json::to_string(¶ms).unwrap();
2476
2477 assert!(json.contains("\"instIdCode\":10459"));
2478 assert!(!json.contains("\"instId\""));
2479 assert!(json.contains("\"triggerPx\":\"48000\""));
2480 }
2481
2482 #[rstest]
2483 fn test_ws_cancel_algo_order_params_serializes_inst_id_code() {
2484 let params = WsCancelAlgoOrderParams {
2485 inst_id_code: 10459,
2486 algo_id: Some("987654321".to_string()),
2487 algo_cl_ord_id: None,
2488 };
2489
2490 let json = serde_json::to_string(¶ms).unwrap();
2491
2492 assert!(json.contains("\"instIdCode\":10459"));
2493 assert!(!json.contains("\"instId\""));
2494 assert!(json.contains("\"algoId\":\"987654321\""));
2495 }
2496
2497 #[rstest]
2498 fn test_ws_cancel_algo_order_params_builder_allows_either_identifier() {
2499 use super::WsCancelAlgoOrderParamsBuilder;
2500
2501 let by_cl_ord_id = WsCancelAlgoOrderParamsBuilder::default()
2502 .inst_id_code(10459u64)
2503 .algo_cl_ord_id("Odstalgocancel0000001".to_string())
2504 .build()
2505 .unwrap();
2506 let json = serde_json::to_value(&by_cl_ord_id).unwrap();
2507 assert_eq!(json["algoClOrdId"], "Odstalgocancel0000001");
2508 assert!(json.get("algoId").is_none());
2509
2510 let by_algo_id = WsCancelAlgoOrderParamsBuilder::default()
2511 .inst_id_code(10459u64)
2512 .algo_id("987654321".to_string())
2513 .build()
2514 .unwrap();
2515 let json = serde_json::to_value(&by_algo_id).unwrap();
2516 assert_eq!(json["algoId"], "987654321");
2517 assert!(json.get("algoClOrdId").is_none());
2518 }
2519
2520 #[rstest]
2521 fn test_ws_post_order_params_serializes_px_usd() {
2522 use super::WsPostOrderParamsBuilder;
2523 use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2524
2525 let params = WsPostOrderParamsBuilder::default()
2526 .inst_id_code(10459u64)
2527 .td_mode(OKXTradeMode::Cross)
2528 .side(OKXSide::Buy)
2529 .ord_type(OKXOrderType::Limit)
2530 .sz("1".to_string())
2531 .px_usd("100.5".to_string())
2532 .build()
2533 .unwrap();
2534
2535 let json = serde_json::to_string(¶ms).unwrap();
2536 assert!(json.contains("\"pxUsd\":\"100.5\""));
2537 assert!(!json.contains("\"pxVol\""));
2538 assert!(!json.contains("\"px\":"));
2539 }
2540
2541 #[rstest]
2542 fn test_ws_post_order_params_serializes_px_vol() {
2543 use super::WsPostOrderParamsBuilder;
2544 use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2545
2546 let params = WsPostOrderParamsBuilder::default()
2547 .inst_id_code(10459u64)
2548 .td_mode(OKXTradeMode::Cross)
2549 .side(OKXSide::Buy)
2550 .ord_type(OKXOrderType::Limit)
2551 .sz("1".to_string())
2552 .px_vol("0.55".to_string())
2553 .build()
2554 .unwrap();
2555
2556 let json = serde_json::to_string(¶ms).unwrap();
2557 assert!(json.contains("\"pxVol\":\"0.55\""));
2558 assert!(!json.contains("\"pxUsd\""));
2559 assert!(!json.contains("\"px\":"));
2560 }
2561
2562 #[rstest]
2563 fn test_ws_amend_order_params_serializes_new_px_usd() {
2564 use super::WsAmendOrderParamsBuilder;
2565
2566 let params = WsAmendOrderParamsBuilder::default()
2567 .inst_id_code(10459u64)
2568 .cl_ord_id("client123".to_string())
2569 .new_px_usd("105.0".to_string())
2570 .build()
2571 .unwrap();
2572
2573 let json = serde_json::to_string(¶ms).unwrap();
2574 assert!(json.contains("\"newPxUsd\":\"105.0\""));
2575 assert!(!json.contains("\"newPx\":"));
2576 assert!(!json.contains("\"newPxVol\""));
2577 }
2578
2579 #[rstest]
2580 fn test_ws_amend_order_params_serializes_new_px_vol() {
2581 use super::WsAmendOrderParamsBuilder;
2582
2583 let params = WsAmendOrderParamsBuilder::default()
2584 .inst_id_code(10459u64)
2585 .cl_ord_id("client123".to_string())
2586 .new_px_vol("0.60".to_string())
2587 .build()
2588 .unwrap();
2589
2590 let json = serde_json::to_string(¶ms).unwrap();
2591 assert!(json.contains("\"newPxVol\":\"0.60\""));
2592 assert!(!json.contains("\"newPx\":"));
2593 assert!(!json.contains("\"newPxUsd\""));
2594 }
2595
2596 #[rstest]
2597 fn test_ws_event_contract_markets_channel_serialization() {
2598 let json = serde_json::to_string(&OKXWsChannel::EventContractMarkets).unwrap();
2599 let channel: OKXWsChannel = serde_json::from_str(&json).unwrap();
2600
2601 assert_eq!(json, "\"event-contract-markets\"");
2602 assert_eq!(channel, OKXWsChannel::EventContractMarkets);
2603 }
2604
2605 #[rstest]
2606 fn test_ws_post_order_params_serializes_event_contract_fields() {
2607 use super::WsPostOrderParamsBuilder;
2608 use crate::common::enums::{OKXOrderType, OKXSide, OKXTradeMode};
2609
2610 let params = WsPostOrderParamsBuilder::default()
2611 .inst_id_code(10459u64)
2612 .td_mode(OKXTradeMode::Cash)
2613 .side(OKXSide::Buy)
2614 .ord_type(OKXOrderType::Limit)
2615 .sz("10".to_string())
2616 .px("0.42".to_string())
2617 .outcome("yes")
2618 .build()
2619 .unwrap();
2620
2621 let json: serde_json::Value = serde_json::to_value(¶ms).unwrap();
2622
2623 assert!(json.get("speedBump").is_none());
2624 assert_eq!(json["outcome"], "yes");
2625 }
2626
2627 #[rstest]
2628 fn test_ws_amend_order_params_omits_speed_bump() {
2629 use super::WsAmendOrderParamsBuilder;
2630
2631 let params = WsAmendOrderParamsBuilder::default()
2632 .inst_id_code(10459u64)
2633 .cl_ord_id("event-1".to_string())
2634 .new_px("0.43".to_string())
2635 .build()
2636 .unwrap();
2637
2638 let json: serde_json::Value = serde_json::to_value(¶ms).unwrap();
2639
2640 assert_eq!(
2641 json,
2642 serde_json::json!({
2643 "instIdCode": 10459,
2644 "clOrdId": "event-1",
2645 "newPx": "0.43",
2646 })
2647 );
2648 }
2649
2650 #[rstest]
2651 fn test_ws_attach_algo_ord_params_serializes_trailing_fields() {
2652 use super::WsAttachAlgoOrdParamsBuilder;
2653
2654 let params = WsAttachAlgoOrdParamsBuilder::default()
2655 .attach_algo_cl_ord_id("trail-1")
2656 .callback_ratio("0.01")
2657 .active_px("64000")
2658 .new_callback_ratio("0.02")
2659 .new_callback_spread("25")
2660 .new_active_px("65000")
2661 .build()
2662 .unwrap();
2663
2664 let json: serde_json::Value = serde_json::to_value(¶ms).unwrap();
2665
2666 assert_eq!(json["callbackRatio"], "0.01");
2667 assert_eq!(json["activePx"], "64000");
2668 assert_eq!(json["newCallbackRatio"], "0.02");
2669 assert_eq!(json["newCallbackSpread"], "25");
2670 assert_eq!(json["newActivePx"], "65000");
2671 assert!(json.get("callbackSpread").is_none());
2672 }
2673
2674 #[rstest]
2675 fn test_subscription_arg_serializes_sprd_id_for_spread_channels() {
2676 let arg = OKXSubscriptionArg {
2677 channel: OKXWsChannel::SprdBooks5,
2678 inst_type: None,
2679 inst_family: None,
2680 inst_id: Some(Ustr::from("ETH-USD-260925_ETH-USD-261225")),
2681 };
2682 let json = serde_json::to_value(&arg).unwrap();
2683 assert_eq!(json["channel"], "sprd-books5");
2684 assert_eq!(json["sprdId"], "ETH-USD-260925_ETH-USD-261225");
2685 assert!(json.get("instId").is_none());
2686 }
2687
2688 #[rstest]
2689 fn test_subscription_arg_serializes_inst_id_for_standard_channels() {
2690 let arg = OKXSubscriptionArg {
2691 channel: OKXWsChannel::BboTbt,
2692 inst_type: None,
2693 inst_family: None,
2694 inst_id: Some(Ustr::from("BTC-USDT")),
2695 };
2696 let json = serde_json::to_value(&arg).unwrap();
2697 assert_eq!(json["instId"], "BTC-USDT");
2698 assert!(json.get("sprdId").is_none());
2699 }
2700
2701 #[rstest]
2702 fn test_websocket_arg_resolves_sprd_id_into_inst_id() {
2703 let arg: OKXWebSocketArg = serde_json::from_value(serde_json::json!({
2704 "channel": "sprd-bbo-tbt",
2705 "sprdId": "ETH-USD-260925_ETH-USD-261225",
2706 }))
2707 .unwrap();
2708 assert_eq!(arg.channel, OKXWsChannel::SprdBboTbt);
2709 assert_eq!(
2710 arg.inst_id,
2711 Some(Ustr::from("ETH-USD-260925_ETH-USD-261225"))
2712 );
2713 }
2714
2715 #[rstest]
2716 fn test_book_msg_parses_three_element_spread_levels() {
2717 let msg: OKXBookMsg = serde_json::from_value(serde_json::json!({
2720 "asks": [["16.7", "100", "1"]],
2721 "bids": [["16.65", "100", "1"]],
2722 "ts": "1780044924909",
2723 "seqId": 1_779_935_772_619_784_u64,
2724 }))
2725 .unwrap();
2726 assert_eq!(msg.asks[0].price, "16.7");
2727 assert_eq!(msg.asks[0].size, "100");
2728 assert_eq!(msg.bids[0].price, "16.65");
2729 }
2730
2731 #[rstest]
2732 fn test_trade_msg_parses_spread_public_trade() {
2733 let msg: OKXTradeMsg = serde_json::from_value(serde_json::json!({
2735 "sprdId": "ETH-USD-260925_ETH-USD-261225",
2736 "tradeId": "3392538740127301632",
2737 "px": "16.9",
2738 "sz": "100",
2739 "side": "sell",
2740 "ts": "1780047866507",
2741 }))
2742 .unwrap();
2743 assert_eq!(msg.inst_id, Ustr::from("ETH-USD-260925_ETH-USD-261225"));
2744 assert_eq!(msg.px, "16.9");
2745 assert_eq!(msg.side, OKXSide::Sell);
2746 assert!(msg.count.is_empty());
2747 }
2748}