1use std::{
19 sync::{
20 Arc,
21 atomic::{AtomicU64, Ordering},
22 },
23 time::Duration,
24};
25
26use ahash::AHashMap;
27use dashmap::DashMap;
28use nautilus_core::{UUID4, time::AtomicTime};
29use nautilus_live::task::TaskSpawner;
30use nautilus_model::{
31 events::{
32 OrderAccepted, OrderCancelRejected, OrderEventAny, OrderModifyRejected, OrderRejected,
33 OrderUpdated,
34 },
35 identifiers::{AccountId, ClientOrderId, TraderId, VenueOrderId},
36 types::{Price, Quantity},
37};
38use parking_lot::RwLock;
39use ustr::Ustr;
40
41use super::WsDispatchState;
42use crate::{
43 common::parse::truncate_cl_ord_id,
44 websocket::spot_v2::{
45 enums::KrakenWsMethod,
46 handler::SpotHandlerCommand,
47 messages::{
48 KrakenWsAddOrderParams, KrakenWsAmendOrderParams, KrakenWsBatchAddParams,
49 KrakenWsCancelOrderParams, KrakenWsOrderResponse, KrakenWsOrderResult, KrakenWsParams,
50 KrakenWsRequest,
51 },
52 },
53};
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub enum PendingOperation {
57 Submit,
59 Amend,
61 Cancel,
63 BatchAdd,
65}
66
67#[derive(Debug, Clone)]
68pub struct PendingRequest {
69 pub operation: PendingOperation,
71 pub client_order_ids: Vec<ClientOrderId>,
73 pub venue_order_ids: Vec<Option<VenueOrderId>>,
75 pub ts_sent_ns: u64,
77 pub new_quantity: Option<Quantity>,
82 pub new_price: Option<Price>,
86 pub new_trigger_price: Option<Price>,
90}
91
92#[derive(Debug)]
93pub struct OrderRequestState {
94 req_id_counter: Arc<AtomicU64>,
95 pending: DashMap<u64, PendingRequest>,
96 timeout: Duration,
97 cmd_tx_handle: Arc<tokio::sync::RwLock<tokio::sync::mpsc::UnboundedSender<SpotHandlerCommand>>>,
101 event_tx: tokio::sync::mpsc::UnboundedSender<OrderEventAny>,
102 dispatch_state: Arc<WsDispatchState>,
103 trader_id: TraderId,
104 account_id: AccountId,
105 auth_token: Arc<tokio::sync::RwLock<Option<String>>>,
108 task_spawner: RwLock<TaskSpawner>,
109 clock: &'static AtomicTime,
111}
112
113impl OrderRequestState {
114 #[expect(clippy::too_many_arguments, reason = "all fields are independent")]
116 pub fn new(
117 cmd_tx_handle: Arc<
118 tokio::sync::RwLock<tokio::sync::mpsc::UnboundedSender<SpotHandlerCommand>>,
119 >,
120 event_tx: tokio::sync::mpsc::UnboundedSender<OrderEventAny>,
121 dispatch_state: Arc<WsDispatchState>,
122 req_id_counter: Arc<AtomicU64>,
123 timeout: Duration,
124 trader_id: TraderId,
125 account_id: AccountId,
126 auth_token: Arc<tokio::sync::RwLock<Option<String>>>,
127 task_spawner: TaskSpawner,
128 clock: &'static AtomicTime,
129 ) -> Self {
130 Self {
131 req_id_counter,
132 pending: DashMap::new(),
133 timeout,
134 cmd_tx_handle,
135 event_tx,
136 dispatch_state,
137 trader_id,
138 account_id,
139 auth_token,
140 task_spawner: RwLock::new(task_spawner),
141 clock,
142 }
143 }
144
145 fn cmd_tx(&self) -> Option<tokio::sync::mpsc::UnboundedSender<SpotHandlerCommand>> {
148 self.cmd_tx_handle.try_read().ok().map(|g| g.clone())
149 }
150
151 pub fn next_req_id(&self) -> u64 {
153 self.req_id_counter.fetch_add(1, Ordering::Relaxed) + 1
154 }
155
156 pub fn submit(
163 self: &Arc<Self>,
164 params: KrakenWsAddOrderParams,
165 identity: PendingRequest,
166 ts_now_ns: u64,
167 ) -> anyhow::Result<u64> {
168 self.send(
169 KrakenWsRequest {
170 method: KrakenWsMethod::AddOrder,
171 params: Some(KrakenWsParams::AddOrder(params)),
172 req_id: None,
173 },
174 identity,
175 ts_now_ns,
176 )
177 }
178
179 pub fn amend(
185 self: &Arc<Self>,
186 params: KrakenWsAmendOrderParams,
187 identity: PendingRequest,
188 ts_now_ns: u64,
189 ) -> anyhow::Result<u64> {
190 self.send(
191 KrakenWsRequest {
192 method: KrakenWsMethod::AmendOrder,
193 params: Some(KrakenWsParams::AmendOrder(params)),
194 req_id: None,
195 },
196 identity,
197 ts_now_ns,
198 )
199 }
200
201 pub fn cancel(
207 self: &Arc<Self>,
208 params: KrakenWsCancelOrderParams,
209 identity: PendingRequest,
210 ts_now_ns: u64,
211 ) -> anyhow::Result<u64> {
212 self.send(
213 KrakenWsRequest {
214 method: KrakenWsMethod::CancelOrder,
215 params: Some(KrakenWsParams::CancelOrder(params)),
216 req_id: None,
217 },
218 identity,
219 ts_now_ns,
220 )
221 }
222
223 pub fn batch_add(
229 self: &Arc<Self>,
230 params: KrakenWsBatchAddParams,
231 identity: PendingRequest,
232 ts_now_ns: u64,
233 ) -> anyhow::Result<u64> {
234 self.send(
235 KrakenWsRequest {
236 method: KrakenWsMethod::BatchAdd,
237 params: Some(KrakenWsParams::BatchAdd(params)),
238 req_id: None,
239 },
240 identity,
241 ts_now_ns,
242 )
243 }
244
245 fn send(
246 self: &Arc<Self>,
247 mut envelope: KrakenWsRequest,
248 mut identity: PendingRequest,
249 ts_now_ns: u64,
250 ) -> anyhow::Result<u64> {
251 let req_id = self.next_req_id();
252 envelope.req_id = Some(req_id);
253 identity.ts_sent_ns = ts_now_ns;
254
255 let payload = serde_json::to_string(&envelope)
256 .map_err(|e| anyhow::anyhow!("serialize WS order request: {e}"))?;
257
258 let cmd_tx = self
259 .cmd_tx()
260 .ok_or_else(|| anyhow::anyhow!("WS handler command sender unavailable"))?;
261
262 self.pending.insert(req_id, identity);
263
264 if let Err(e) = cmd_tx.send(SpotHandlerCommand::SendOrderRequest { req_id, payload }) {
265 self.pending.remove(&req_id);
266 anyhow::bail!("handler command channel closed: {e}");
267 }
268
269 let state_for_timeout = Arc::downgrade(self);
270 let task_spawner = self.task_spawner.read().clone();
271 let cancel = task_spawner.cancellation_token();
272 let timeout = self.timeout;
273
274 if let Err(e) = task_spawner.spawn(async move {
275 tokio::select! {
276 biased;
277 () = cancel.cancelled() => {
278 if let Some(state) = state_for_timeout.upgrade() {
279 state.pending.remove(&req_id);
280 }
281 }
282 () = tokio::time::sleep(timeout) => {
283 let Some(state) = state_for_timeout.upgrade() else {
284 return;
285 };
286
287 if let Some(pending) = state.pending.get(&req_id) {
288 if cancel.is_cancelled() {
289 drop(pending);
290 state.pending.remove(&req_id);
291 return;
292 }
293
294 let ts_timeout_ns = state.clock.get_time_ns().as_u64();
295 log::warn!(
296 "Kraken WS response timeout req_id={req_id} op={:?} cl_ord_ids={:?} \
297 ts_timeout_ns={ts_timeout_ns}; awaiting definitive venue evidence",
298 pending.operation,
299 pending.client_order_ids,
300 );
301 state.handle_timeout(&pending);
302 }
303 }
304 }
305 }) {
306 log::warn!(
307 "Kraken order request {req_id} was sent without a local timeout task: {e}; \
308 awaiting definitive venue evidence"
309 );
310 }
311
312 Ok(req_id)
313 }
314
315 pub fn handle_response(&self, response: &KrakenWsOrderResponse, ts_event_ns: u64) {
321 let req_id = match response.req_id {
322 Some(id) => id,
323 None => {
324 log::warn!(
325 "Kraken WS order response without req_id method={:?} success={}",
326 response.method,
327 response.success,
328 );
329 return;
330 }
331 };
332
333 let Some(pending) = self.pending.get(&req_id) else {
334 log::debug!("Kraken WS response without pending request req_id={req_id}");
335 return;
336 };
337
338 let expected_method = pending_op_to_method(pending.operation);
339 if response.method != expected_method {
340 log::error!(
341 "Kraken WS response method {:?} mismatched pending op {:?} req_id={req_id}",
342 response.method,
343 pending.operation,
344 );
345 return;
346 }
347 drop(pending);
348
349 let Some((_, pending)) = self.pending.remove(&req_id) else {
350 log::debug!("Kraken WS duplicate response req_id={req_id}");
351 return;
352 };
353
354 match (pending.operation, response.success) {
355 (PendingOperation::Submit, true) => {
356 self.emit_order_accepted(&pending, response, ts_event_ns);
357 }
358 (PendingOperation::Submit, false) => {
359 self.emit_order_rejected(&pending, response, ts_event_ns);
360 }
361 (PendingOperation::Amend, true) => {
362 self.emit_order_updated(&pending, response, ts_event_ns);
363 }
364 (PendingOperation::Amend, false) => {
365 self.emit_order_modify_rejected(&pending, response, ts_event_ns);
366 }
367 (PendingOperation::Cancel, true) => {
368 log::debug!(
369 "Kraken WS cancel ack req_id={req_id} cl_ord_ids={:?}",
370 pending.client_order_ids,
371 );
372 }
373 (PendingOperation::Cancel, false) => {
374 self.emit_order_cancel_rejected(&pending, response, ts_event_ns);
375 }
376 (PendingOperation::BatchAdd, _) => {
377 self.handle_batch_add_response(&pending, response, ts_event_ns);
378 }
379 }
380 }
381
382 fn handle_timeout(&self, pending: &PendingRequest) {
383 match pending.operation {
384 PendingOperation::Submit => {
385 self.send_compensating_cancel(&pending.client_order_ids);
386 }
387 PendingOperation::Amend | PendingOperation::Cancel => {}
388 PendingOperation::BatchAdd => {
389 self.send_compensating_cancel(&pending.client_order_ids);
390 }
391 }
392 }
393
394 pub(crate) fn clear(&self) {
395 self.pending.clear();
396 }
397
398 pub(crate) fn reset_task_spawner(&self, task_spawner: TaskSpawner) {
399 *self.task_spawner.write() = task_spawner;
400 }
401
402 fn send_compensating_cancel(&self, cl_ord_ids: &[ClientOrderId]) {
413 let Some(token) = self.auth_token.try_read().ok().and_then(|g| g.clone()) else {
414 log::error!(
415 "Submit timeout: no auth token for compensating cancel cl_ord_ids={cl_ord_ids:?}; \
416 relying on reconciliation to recover any orphan order",
417 );
418 return;
419 };
420
421 let req_id = self.next_req_id();
422 let params = KrakenWsCancelOrderParams {
423 token,
424 order_id: None,
425 cl_ord_id: Some(cl_ord_ids.iter().map(truncate_cl_ord_id).collect()),
426 };
427 let envelope = KrakenWsRequest {
428 method: KrakenWsMethod::CancelOrder,
429 params: Some(KrakenWsParams::CancelOrder(params)),
430 req_id: Some(req_id),
431 };
432
433 let payload = match serde_json::to_string(&envelope) {
434 Ok(p) => p,
435 Err(e) => {
436 log::warn!("Submit timeout: compensating cancel serialise failed: {e}");
437 return;
438 }
439 };
440
441 let Some(cmd_tx) = self.cmd_tx() else {
442 log::error!(
443 "Submit timeout: compensating cancel sender unavailable cl_ord_ids={cl_ord_ids:?}; \
444 relying on reconciliation to recover any orphan order",
445 );
446 return;
447 };
448
449 if let Err(e) = cmd_tx.send(SpotHandlerCommand::SendOrderRequest { req_id, payload }) {
450 log::error!(
451 "Submit timeout: compensating cancel channel closed: {e}; \
452 relying on reconciliation to recover any orphan order",
453 );
454 } else {
455 log::debug!(
456 "Submit timeout: compensating cancel sent req_id={req_id} cl_ord_ids={cl_ord_ids:?}",
457 );
458 }
459 }
460
461 fn emit_order_accepted(
462 &self,
463 pending: &PendingRequest,
464 response: &KrakenWsOrderResponse,
465 ts_event_ns: u64,
466 ) {
467 let Some(client_order_id) = pending.client_order_ids.first().copied() else {
468 log::error!("Kraken WS add_order response without client_order_id");
469 return;
470 };
471 let Some(identity) = self.dispatch_state.lookup_identity(&client_order_id) else {
472 log::warn!(
473 "Kraken WS add_order response for untracked order client_order_id={client_order_id}",
474 );
475 return;
476 };
477 let venue_order_id = response
478 .result
479 .as_ref()
480 .and_then(|r| r.order_id.as_deref())
481 .map(VenueOrderId::new);
482 let Some(venue_order_id) = venue_order_id else {
483 log::error!(
484 "Kraken WS add_order success without order_id client_order_id={client_order_id}",
485 );
486 return;
487 };
488
489 if !self.dispatch_state.insert_accepted(client_order_id) {
490 return;
491 }
492
493 let event = OrderAccepted::new(
494 self.trader_id,
495 identity.strategy_id,
496 identity.instrument_id,
497 client_order_id,
498 venue_order_id,
499 self.account_id,
500 UUID4::new(),
501 ts_event_ns.into(),
502 ts_event_ns.into(),
503 false,
504 );
505 self.send_event(OrderEventAny::Accepted(event));
506 }
507
508 fn emit_order_rejected(
509 &self,
510 pending: &PendingRequest,
511 response: &KrakenWsOrderResponse,
512 ts_event_ns: u64,
513 ) {
514 let Some(client_order_id) = pending.client_order_ids.first().copied() else {
515 log::error!("Kraken WS add_order rejection without client_order_id");
516 return;
517 };
518 let Some(identity) = self.dispatch_state.lookup_identity(&client_order_id) else {
519 log::warn!(
520 "Kraken WS add_order rejection for untracked order client_order_id={client_order_id}",
521 );
522 return;
523 };
524 let reason = response
525 .error
526 .as_deref()
527 .filter(|s| !s.is_empty())
528 .unwrap_or("UNKNOWN");
529
530 let event = OrderRejected::new(
531 self.trader_id,
532 identity.strategy_id,
533 identity.instrument_id,
534 client_order_id,
535 self.account_id,
536 Ustr::from(reason),
537 UUID4::new(),
538 ts_event_ns.into(),
539 ts_event_ns.into(),
540 false,
541 false,
542 );
543 self.send_event(OrderEventAny::Rejected(event));
544 }
545
546 fn emit_order_updated(
547 &self,
548 pending: &PendingRequest,
549 response: &KrakenWsOrderResponse,
550 ts_event_ns: u64,
551 ) {
552 let Some(client_order_id) = pending.client_order_ids.first().copied() else {
553 log::error!("Kraken WS amend_order response without client_order_id");
554 return;
555 };
556 let Some(identity) = self.dispatch_state.lookup_identity(&client_order_id) else {
557 log::warn!(
558 "Kraken WS amend_order response for untracked order client_order_id={client_order_id}",
559 );
560 return;
561 };
562 let venue_order_id = response
563 .result
564 .as_ref()
565 .and_then(|r| r.order_id.as_deref())
566 .map(VenueOrderId::new)
567 .or_else(|| pending.venue_order_ids.first().copied().flatten());
568
569 let quantity = pending.new_quantity.unwrap_or(identity.quantity);
570 if pending.new_quantity.is_some() {
571 self.dispatch_state
572 .update_identity_quantity(&client_order_id, quantity);
573 }
574
575 let event = OrderUpdated::new(
576 self.trader_id,
577 identity.strategy_id,
578 identity.instrument_id,
579 client_order_id,
580 quantity,
581 UUID4::new(),
582 ts_event_ns.into(),
583 ts_event_ns.into(),
584 false,
585 venue_order_id,
586 Some(self.account_id),
587 pending.new_price,
588 pending.new_trigger_price,
589 None,
590 false,
591 );
592 self.send_event(OrderEventAny::Updated(event));
593 }
594
595 fn emit_order_modify_rejected(
596 &self,
597 pending: &PendingRequest,
598 response: &KrakenWsOrderResponse,
599 ts_event_ns: u64,
600 ) {
601 let Some(client_order_id) = pending.client_order_ids.first().copied() else {
602 log::error!("Kraken WS amend_order rejection without client_order_id");
603 return;
604 };
605 let Some(identity) = self.dispatch_state.lookup_identity(&client_order_id) else {
606 log::warn!(
607 "Kraken WS amend_order rejection for untracked order client_order_id={client_order_id}",
608 );
609 return;
610 };
611 let venue_order_id = pending.venue_order_ids.first().copied().flatten();
612 let reason = response
613 .error
614 .as_deref()
615 .filter(|s| !s.is_empty())
616 .unwrap_or("UNKNOWN");
617
618 let event = OrderModifyRejected::new(
619 self.trader_id,
620 identity.strategy_id,
621 identity.instrument_id,
622 client_order_id,
623 Ustr::from(reason),
624 UUID4::new(),
625 ts_event_ns.into(),
626 ts_event_ns.into(),
627 false,
628 venue_order_id,
629 Some(self.account_id),
630 );
631 self.send_event(OrderEventAny::ModifyRejected(event));
632 }
633
634 fn emit_order_cancel_rejected(
635 &self,
636 pending: &PendingRequest,
637 response: &KrakenWsOrderResponse,
638 ts_event_ns: u64,
639 ) {
640 let Some(client_order_id) = pending.client_order_ids.first().copied() else {
641 log::error!("Kraken WS cancel_order rejection without client_order_id");
642 return;
643 };
644 let Some(identity) = self.dispatch_state.lookup_identity(&client_order_id) else {
645 log::warn!(
646 "Kraken WS cancel_order rejection for untracked order client_order_id={client_order_id}",
647 );
648 return;
649 };
650 let venue_order_id = pending.venue_order_ids.first().copied().flatten();
651 let reason = response
652 .error
653 .as_deref()
654 .filter(|s| !s.is_empty())
655 .unwrap_or("UNKNOWN");
656
657 let event = OrderCancelRejected::new(
658 self.trader_id,
659 identity.strategy_id,
660 identity.instrument_id,
661 client_order_id,
662 Ustr::from(reason),
663 UUID4::new(),
664 ts_event_ns.into(),
665 ts_event_ns.into(),
666 false,
667 venue_order_id,
668 Some(self.account_id),
669 );
670 self.send_event(OrderEventAny::CancelRejected(event));
671 }
672
673 fn handle_batch_add_response(
674 &self,
675 pending: &PendingRequest,
676 response: &KrakenWsOrderResponse,
677 ts_event_ns: u64,
678 ) {
679 let per_order = response
680 .result
681 .as_ref()
682 .and_then(|r| r.orders.as_ref())
683 .map_or(&[][..], Vec::as_slice);
684
685 let echo_index: AHashMap<&str, usize> = per_order
691 .iter()
692 .enumerate()
693 .filter_map(|(i, r)| r.cl_ord_id.as_deref().map(|cid| (cid, i)))
694 .collect();
695
696 for (idx, client_order_id) in pending.client_order_ids.iter().copied().enumerate() {
697 let leg_venue = pending.venue_order_ids.get(idx).copied().flatten();
698 let truncated = truncate_cl_ord_id(&client_order_id);
699 let leg_result = echo_index
700 .get(truncated.as_str())
701 .and_then(|&i| per_order.get(i))
702 .or_else(|| per_order.get(idx));
703
704 if leg_result.is_none() {
705 log::error!(
706 "Kraken WS batch_add response missing per-leg result for client_order_id={client_order_id} idx={idx} \
707 legs_sent={legs_sent} legs_received={legs_received}; treating as rejection",
708 legs_sent = pending.client_order_ids.len(),
709 legs_received = per_order.len(),
710 );
711 }
712
713 let leg_success = leg_result.is_some_and(|r| r.success);
714 let leg_venue_order_id = leg_result
715 .and_then(|r| r.order_id.as_deref())
716 .map(VenueOrderId::new)
717 .or(leg_venue);
718 let leg_error = leg_result.and_then(|r| r.error.clone()).or_else(|| {
719 if leg_result.is_none() {
720 Some(format!(
721 "batch_add response missing per-leg result (legs_sent={}, legs_received={})",
722 pending.client_order_ids.len(),
723 per_order.len(),
724 ))
725 } else {
726 response.error.clone()
727 }
728 });
729
730 let leg_response = KrakenWsOrderResponse {
731 method: KrakenWsMethod::AddOrder,
732 req_id: response.req_id,
733 success: leg_success,
734 time_in: response.time_in.clone(),
735 time_out: response.time_out.clone(),
736 error: leg_error,
737 result: leg_venue_order_id.map(|v| KrakenWsOrderResult {
738 order_id: Some(v.to_string()),
739 cl_ord_id: leg_result.and_then(|r| r.cl_ord_id.clone()),
740 order_userref: None,
741 warning: None,
742 orders: None,
743 }),
744 };
745 let leg_pending = PendingRequest {
746 operation: PendingOperation::Submit,
747 client_order_ids: vec![client_order_id],
748 venue_order_ids: vec![leg_venue_order_id],
749 ts_sent_ns: pending.ts_sent_ns,
750 new_quantity: None,
751 new_price: None,
752 new_trigger_price: None,
753 };
754
755 if leg_success {
756 self.emit_order_accepted(&leg_pending, &leg_response, ts_event_ns);
757 } else {
758 self.emit_order_rejected(&leg_pending, &leg_response, ts_event_ns);
759 }
760 }
761 }
762
763 fn send_event(&self, event: OrderEventAny) {
772 if let Err(e) = self.event_tx.send(event) {
773 log::error!("Kraken WS order-event channel send failed: {e}");
774 }
775 }
776}
777
778fn pending_op_to_method(op: PendingOperation) -> KrakenWsMethod {
779 match op {
780 PendingOperation::Submit => KrakenWsMethod::AddOrder,
781 PendingOperation::Amend => KrakenWsMethod::AmendOrder,
782 PendingOperation::Cancel => KrakenWsMethod::CancelOrder,
783 PendingOperation::BatchAdd => KrakenWsMethod::BatchAdd,
784 }
785}
786
787#[cfg(test)]
788impl OrderRequestState {
789 pub(crate) fn pending_len(&self) -> usize {
790 self.pending.len()
791 }
792}
793
794#[cfg(test)]
795mod tests {
796 use std::sync::atomic::AtomicU64;
797
798 use nautilus_live::task::TaskGroup;
799 use nautilus_model::{
800 enums::{OrderSide, OrderType},
801 identifiers::{AccountId, ClientOrderId, InstrumentId, StrategyId, TraderId},
802 types::{Price, Quantity},
803 };
804 use rstest::rstest;
805 use rust_decimal_macros::dec;
806
807 use super::*;
808 use crate::{
809 common::enums::{KrakenOrderSide, KrakenOrderType},
810 websocket::{
811 dispatch::OrderIdentity,
812 spot_v2::messages::{
813 KrakenWsAddOrderParams, KrakenWsAmendOrderParams, KrakenWsBatchAddParams,
814 KrakenWsBatchOrderResult, KrakenWsCancelOrderParams, KrakenWsOrderResponse,
815 KrakenWsOrderResult,
816 },
817 },
818 };
819
820 const CLIENT_ORDER_ID: &str = "O-1";
821 const VENUE_ORDER_ID: &str = "O-VENUE";
822 const INSTRUMENT_ID: &str = "BTCUSD.KRAKEN";
823
824 pub(super) struct Harness {
825 pub(super) state: Arc<OrderRequestState>,
826 pub(super) cmd_rx: tokio::sync::mpsc::UnboundedReceiver<SpotHandlerCommand>,
827 pub(super) event_rx: tokio::sync::mpsc::UnboundedReceiver<OrderEventAny>,
828 pub(super) dispatch_state: Arc<WsDispatchState>,
829 pub(super) auth_token: Arc<tokio::sync::RwLock<Option<String>>>,
830 pub(super) pending_tasks: TaskGroup,
831 pub(super) cmd_tx_handle:
832 Arc<tokio::sync::RwLock<tokio::sync::mpsc::UnboundedSender<SpotHandlerCommand>>>,
833 }
834
835 pub(super) fn make_harness(timeout_ms: u64) -> Harness {
836 let (cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
837 let (event_tx, event_rx) = tokio::sync::mpsc::unbounded_channel();
838 let counter = Arc::new(AtomicU64::new(0));
839 let dispatch_state = Arc::new(WsDispatchState::new());
840 let auth_token = Arc::new(tokio::sync::RwLock::new(None));
841 let pending_tasks = TaskGroup::new();
842 let pending_spawner = pending_tasks.spawner().expect("pending task spawner");
843 let cmd_tx_handle = Arc::new(tokio::sync::RwLock::new(cmd_tx));
844 let state = Arc::new(OrderRequestState::new(
845 Arc::clone(&cmd_tx_handle),
846 event_tx,
847 Arc::clone(&dispatch_state),
848 counter,
849 Duration::from_millis(timeout_ms),
850 TraderId::new("TESTER-001"),
851 AccountId::new("KRAKEN-001"),
852 Arc::clone(&auth_token),
853 pending_spawner,
854 nautilus_core::time::get_atomic_clock_realtime(),
855 ));
856 Harness {
857 state,
858 cmd_rx,
859 event_rx,
860 dispatch_state,
861 auth_token,
862 pending_tasks,
863 cmd_tx_handle,
864 }
865 }
866
867 fn register_default_identity(dispatch_state: &WsDispatchState, cl_ord_id: ClientOrderId) {
868 dispatch_state.register_identity(
869 cl_ord_id,
870 OrderIdentity {
871 strategy_id: StrategyId::new("S-1"),
872 instrument_id: InstrumentId::from(INSTRUMENT_ID),
873 order_side: OrderSide::Buy,
874 order_type: OrderType::Limit,
875 quantity: Quantity::from("0.001"),
876 },
877 );
878 }
879
880 fn make_state(
881 timeout_ms: u64,
882 ) -> (
883 Arc<OrderRequestState>,
884 tokio::sync::mpsc::UnboundedReceiver<SpotHandlerCommand>,
885 ) {
886 let harness = make_harness(timeout_ms);
887 (harness.state, harness.cmd_rx)
888 }
889
890 fn make_identity(op: PendingOperation) -> PendingRequest {
891 PendingRequest {
892 operation: op,
893 client_order_ids: vec![ClientOrderId::from(CLIENT_ORDER_ID)],
894 venue_order_ids: vec![None],
895 ts_sent_ns: 0,
896 new_quantity: None,
897 new_price: None,
898 new_trigger_price: None,
899 }
900 }
901
902 fn make_add_order_params(token: &str) -> KrakenWsAddOrderParams {
903 KrakenWsAddOrderParams {
904 order_type: KrakenOrderType::Limit,
905 side: KrakenOrderSide::Buy,
906 order_qty: dec!(0.001),
907 symbol: "BTC/USD".to_string(),
908 token: token.to_string(),
909 limit_price: Some(dec!(50000)),
910 time_in_force: None,
911 expire_time: None,
912 cl_ord_id: Some(CLIENT_ORDER_ID.to_string()),
913 post_only: None,
914 reduce_only: None,
915 leverage: None,
916 trigger: None,
917 conditional: None,
918 }
919 }
920
921 #[rstest]
922 fn test_next_req_id_is_monotonic() {
923 let (state, _rx) = make_state(1_000);
924 let a = state.next_req_id();
925 let b = state.next_req_id();
926 let c = state.next_req_id();
927 assert!(b > a && c > b);
928 }
929
930 #[rstest]
931 fn test_submit_registers_pending_and_sends_command() {
932 let (state, mut rx) = make_state(60_000);
933 let identity = make_identity(PendingOperation::Submit);
934
935 let params = KrakenWsAddOrderParams {
936 order_type: KrakenOrderType::Limit,
937 side: KrakenOrderSide::Buy,
938 order_qty: dec!(0.001),
939 symbol: "BTC/USD".to_string(),
940 token: "test-token".to_string(),
941 limit_price: Some(dec!(50000)),
942 time_in_force: None,
943 expire_time: None,
944 cl_ord_id: Some("O-1".to_string()),
945 post_only: None,
946 reduce_only: None,
947 leverage: None,
948 trigger: None,
949 conditional: None,
950 };
951
952 let req_id = state.submit(params, identity, 1).expect("submit ok");
953 assert_eq!(state.pending_len(), 1);
954
955 let cmd = rx.try_recv().expect("cmd queued");
956 match cmd {
957 SpotHandlerCommand::SendOrderRequest {
958 req_id: rid,
959 payload,
960 } => {
961 assert_eq!(rid, req_id);
962 assert!(payload.contains("\"add_order\""));
963 assert!(payload.contains(&format!("\"req_id\":{req_id}")));
964 }
965 _ => panic!("wrong cmd variant"),
966 }
967 }
968
969 #[tokio::test]
970 async fn test_submit_uses_swapped_command_sender() {
971 let harness = make_harness(60_000);
974
975 drop(harness.cmd_rx);
976 let (new_cmd_tx, mut new_cmd_rx) = tokio::sync::mpsc::unbounded_channel();
977 *harness.cmd_tx_handle.write().await = new_cmd_tx;
978
979 let params = KrakenWsAddOrderParams {
980 order_type: KrakenOrderType::Limit,
981 side: KrakenOrderSide::Buy,
982 order_qty: dec!(0.001),
983 symbol: "BTC/USD".to_string(),
984 token: "TKN".to_string(),
985 limit_price: Some(dec!(50000)),
986 time_in_force: None,
987 expire_time: None,
988 cl_ord_id: Some(CLIENT_ORDER_ID.to_string()),
989 post_only: None,
990 reduce_only: None,
991 leverage: None,
992 trigger: None,
993 conditional: None,
994 };
995 let identity = make_identity(PendingOperation::Submit);
996
997 harness
998 .state
999 .submit(params, identity, 1)
1000 .expect("submit must succeed via the swapped sender");
1001
1002 let cmd = new_cmd_rx
1003 .try_recv()
1004 .expect("swapped receiver must observe the order");
1005
1006 match cmd {
1007 SpotHandlerCommand::SendOrderRequest { payload, .. } => {
1008 assert!(payload.contains("\"add_order\""));
1009 }
1010 other => panic!("expected SendOrderRequest, was {other:?}"),
1011 }
1012 }
1013
1014 #[rstest]
1015 fn test_amend_sends_amend_order_envelope() {
1016 let (state, mut rx) = make_state(60_000);
1017 let params = KrakenWsAmendOrderParams {
1018 order_id: Some("O-VENUE".to_string()),
1019 cl_ord_id: None,
1020 order_qty: Some(dec!(0.005)),
1021 limit_price: None,
1022 trigger_price: None,
1023 token: "TKN".to_string(),
1024 };
1025 let identity = PendingRequest {
1026 operation: PendingOperation::Amend,
1027 client_order_ids: vec![ClientOrderId::from("O-1")],
1028 venue_order_ids: vec![Some(VenueOrderId::from("O-VENUE"))],
1029 ts_sent_ns: 0,
1030 new_quantity: None,
1031 new_price: None,
1032 new_trigger_price: None,
1033 };
1034 let _ = state.amend(params, identity, 1).expect("amend ok");
1035 let cmd = rx.try_recv().unwrap();
1036 if let SpotHandlerCommand::SendOrderRequest { payload, .. } = cmd {
1037 assert!(payload.contains("\"amend_order\""));
1038 } else {
1039 panic!("wrong variant");
1040 }
1041 }
1042
1043 #[rstest]
1044 fn test_cancel_sends_cancel_order_envelope() {
1045 let (state, mut rx) = make_state(60_000);
1046 let params = KrakenWsCancelOrderParams {
1047 order_id: Some(vec!["O-VENUE".to_string()]),
1048 cl_ord_id: None,
1049 token: "TKN".to_string(),
1050 };
1051 let identity = PendingRequest {
1052 operation: PendingOperation::Cancel,
1053 client_order_ids: vec![ClientOrderId::from("O-1")],
1054 venue_order_ids: vec![Some(VenueOrderId::from("O-VENUE"))],
1055 ts_sent_ns: 0,
1056 new_quantity: None,
1057 new_price: None,
1058 new_trigger_price: None,
1059 };
1060 let _ = state.cancel(params, identity, 1).expect("cancel ok");
1061 let cmd = rx.try_recv().unwrap();
1062 if let SpotHandlerCommand::SendOrderRequest { payload, .. } = cmd {
1063 assert!(payload.contains("\"cancel_order\""));
1064 } else {
1065 panic!("wrong variant");
1066 }
1067 }
1068
1069 #[rstest]
1070 fn test_batch_add_sends_batch_add_envelope() {
1071 let (state, mut rx) = make_state(60_000);
1072 let params = KrakenWsBatchAddParams {
1073 symbol: "BTC/USD".to_string(),
1074 orders: vec![],
1075 token: "TKN".to_string(),
1076 };
1077 let identity = PendingRequest {
1078 operation: PendingOperation::BatchAdd,
1079 client_order_ids: vec![ClientOrderId::from("O-A"), ClientOrderId::from("O-B")],
1080 venue_order_ids: vec![None, None],
1081 ts_sent_ns: 0,
1082 new_quantity: None,
1083 new_price: None,
1084 new_trigger_price: None,
1085 };
1086 let _ = state.batch_add(params, identity, 1).expect("batch ok");
1087 let cmd = rx.try_recv().unwrap();
1088 if let SpotHandlerCommand::SendOrderRequest { payload, .. } = cmd {
1089 assert!(payload.contains("\"batch_add\""));
1090 } else {
1091 panic!("wrong variant");
1092 }
1093 }
1094
1095 fn make_response(
1096 method: KrakenWsMethod,
1097 success: bool,
1098 req_id: u64,
1099 order_id: Option<&str>,
1100 error: Option<&str>,
1101 ) -> KrakenWsOrderResponse {
1102 KrakenWsOrderResponse {
1103 method,
1104 req_id: Some(req_id),
1105 success,
1106 time_in: None,
1107 time_out: None,
1108 error: error.map(str::to_string),
1109 result: order_id.map(|id| KrakenWsOrderResult {
1110 order_id: Some(id.to_string()),
1111 cl_ord_id: None,
1112 order_userref: None,
1113 warning: None,
1114 orders: None,
1115 }),
1116 }
1117 }
1118
1119 #[rstest]
1120 fn test_handle_response_submit_success_emits_order_accepted() {
1121 let mut harness = make_harness(60_000);
1122 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1123 register_default_identity(&harness.dispatch_state, cl_ord_id);
1124
1125 let req_id = 42;
1126 harness
1127 .state
1128 .pending
1129 .insert(req_id, make_identity(PendingOperation::Submit));
1130
1131 let response = make_response(
1132 KrakenWsMethod::AddOrder,
1133 true,
1134 req_id,
1135 Some(VENUE_ORDER_ID),
1136 None,
1137 );
1138 harness.state.handle_response(&response, 1_000);
1139
1140 assert_eq!(harness.state.pending_len(), 0);
1141 let event = harness.event_rx.try_recv().expect("event emitted");
1142 match event {
1143 OrderEventAny::Accepted(e) => {
1144 assert_eq!(e.client_order_id, cl_ord_id);
1145 assert_eq!(e.venue_order_id.as_str(), VENUE_ORDER_ID);
1146 assert_eq!(e.account_id.as_str(), "KRAKEN-001");
1147 }
1148 other => panic!("expected Accepted, was {other:?}"),
1149 }
1150 }
1151
1152 #[rstest]
1153 fn test_handle_response_submit_failure_emits_order_rejected() {
1154 let mut harness = make_harness(60_000);
1155 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1156 register_default_identity(&harness.dispatch_state, cl_ord_id);
1157
1158 let req_id = 7;
1159 harness
1160 .state
1161 .pending
1162 .insert(req_id, make_identity(PendingOperation::Submit));
1163
1164 let response = make_response(
1165 KrakenWsMethod::AddOrder,
1166 false,
1167 req_id,
1168 None,
1169 Some("Insufficient funds"),
1170 );
1171 harness.state.handle_response(&response, 2_000);
1172
1173 assert_eq!(harness.state.pending_len(), 0);
1174 let event = harness.event_rx.try_recv().expect("event emitted");
1175 match event {
1176 OrderEventAny::Rejected(e) => {
1177 assert_eq!(e.client_order_id, cl_ord_id);
1178 assert_eq!(e.reason.as_str(), "Insufficient funds");
1179 }
1180 other => panic!("expected Rejected, was {other:?}"),
1181 }
1182 }
1183
1184 #[rstest]
1185 fn test_handle_response_amend_success_emits_order_updated() {
1186 let mut harness = make_harness(60_000);
1187 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1188 register_default_identity(&harness.dispatch_state, cl_ord_id);
1189
1190 let req_id = 11;
1191 let pending = PendingRequest {
1192 operation: PendingOperation::Amend,
1193 client_order_ids: vec![cl_ord_id],
1194 venue_order_ids: vec![Some(VenueOrderId::from(VENUE_ORDER_ID))],
1195 ts_sent_ns: 0,
1196 new_quantity: None,
1197 new_price: None,
1198 new_trigger_price: None,
1199 };
1200 harness.state.pending.insert(req_id, pending);
1201
1202 let response = make_response(
1203 KrakenWsMethod::AmendOrder,
1204 true,
1205 req_id,
1206 Some(VENUE_ORDER_ID),
1207 None,
1208 );
1209 harness.state.handle_response(&response, 3_000);
1210
1211 let event = harness.event_rx.try_recv().expect("event emitted");
1212 match event {
1213 OrderEventAny::Updated(e) => {
1214 assert_eq!(e.client_order_id, cl_ord_id);
1215 assert_eq!(
1216 e.venue_order_id.expect("venue id present").as_str(),
1217 VENUE_ORDER_ID
1218 );
1219 assert_eq!(e.quantity, Quantity::from("0.001"));
1220 }
1221 other => panic!("expected Updated, was {other:?}"),
1222 }
1223 }
1224
1225 #[rstest]
1226 fn test_handle_response_amend_with_new_quantity_emits_new_quantity() {
1227 let mut harness = make_harness(60_000);
1228 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1229 register_default_identity(&harness.dispatch_state, cl_ord_id);
1230
1231 let req_id = 11;
1232 let new_qty = Quantity::from("0.005");
1233 let pending = PendingRequest {
1234 operation: PendingOperation::Amend,
1235 client_order_ids: vec![cl_ord_id],
1236 venue_order_ids: vec![Some(VenueOrderId::from(VENUE_ORDER_ID))],
1237 ts_sent_ns: 0,
1238 new_quantity: Some(new_qty),
1239 new_price: None,
1240 new_trigger_price: None,
1241 };
1242 harness.state.pending.insert(req_id, pending);
1243
1244 let response = make_response(
1245 KrakenWsMethod::AmendOrder,
1246 true,
1247 req_id,
1248 Some(VENUE_ORDER_ID),
1249 None,
1250 );
1251 harness.state.handle_response(&response, 3_500);
1252
1253 let event = harness.event_rx.try_recv().expect("event emitted");
1254 match event {
1255 OrderEventAny::Updated(e) => {
1256 assert_eq!(e.quantity, new_qty, "OrderUpdated must carry new quantity");
1257 }
1258 other => panic!("expected Updated, was {other:?}"),
1259 }
1260
1261 let identity = harness
1262 .dispatch_state
1263 .lookup_identity(&cl_ord_id)
1264 .expect("identity present");
1265 assert_eq!(
1266 identity.quantity, new_qty,
1267 "dispatch identity must be updated for follow-up ops",
1268 );
1269 }
1270
1271 #[rstest]
1272 fn test_handle_response_amend_carries_new_price_and_trigger() {
1273 let mut harness = make_harness(60_000);
1274 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1275 register_default_identity(&harness.dispatch_state, cl_ord_id);
1276
1277 let req_id = 13;
1278 let new_price = Price::from("31000.00");
1279 let new_trigger = Price::from("30500.00");
1280 let pending = PendingRequest {
1281 operation: PendingOperation::Amend,
1282 client_order_ids: vec![cl_ord_id],
1283 venue_order_ids: vec![Some(VenueOrderId::from(VENUE_ORDER_ID))],
1284 ts_sent_ns: 0,
1285 new_quantity: None,
1286 new_price: Some(new_price),
1287 new_trigger_price: Some(new_trigger),
1288 };
1289 harness.state.pending.insert(req_id, pending);
1290
1291 let response = make_response(
1292 KrakenWsMethod::AmendOrder,
1293 true,
1294 req_id,
1295 Some(VENUE_ORDER_ID),
1296 None,
1297 );
1298 harness.state.handle_response(&response, 4_500);
1299
1300 let event = harness.event_rx.try_recv().expect("event emitted");
1301 match event {
1302 OrderEventAny::Updated(e) => {
1303 assert_eq!(
1304 e.price,
1305 Some(new_price),
1306 "OrderUpdated must carry amended price",
1307 );
1308 assert_eq!(
1309 e.trigger_price,
1310 Some(new_trigger),
1311 "OrderUpdated must carry amended trigger price",
1312 );
1313 }
1314 other => panic!("expected Updated, was {other:?}"),
1315 }
1316 }
1317
1318 #[rstest]
1319 fn test_handle_response_amend_failure_emits_order_modify_rejected() {
1320 let mut harness = make_harness(60_000);
1321 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1322 register_default_identity(&harness.dispatch_state, cl_ord_id);
1323
1324 let req_id = 12;
1325 let pending = PendingRequest {
1326 operation: PendingOperation::Amend,
1327 client_order_ids: vec![cl_ord_id],
1328 venue_order_ids: vec![Some(VenueOrderId::from(VENUE_ORDER_ID))],
1329 ts_sent_ns: 0,
1330 new_quantity: None,
1331 new_price: None,
1332 new_trigger_price: None,
1333 };
1334 harness.state.pending.insert(req_id, pending);
1335
1336 let response = make_response(
1337 KrakenWsMethod::AmendOrder,
1338 false,
1339 req_id,
1340 None,
1341 Some("Order not found"),
1342 );
1343 harness.state.handle_response(&response, 4_000);
1344
1345 let event = harness.event_rx.try_recv().expect("event emitted");
1346 match event {
1347 OrderEventAny::ModifyRejected(e) => {
1348 assert_eq!(e.client_order_id, cl_ord_id);
1349 assert_eq!(e.reason.as_str(), "Order not found");
1350 }
1351 other => panic!("expected ModifyRejected, was {other:?}"),
1352 }
1353 }
1354
1355 #[rstest]
1356 fn test_handle_response_cancel_failure_emits_order_cancel_rejected() {
1357 let mut harness = make_harness(60_000);
1358 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1359 register_default_identity(&harness.dispatch_state, cl_ord_id);
1360
1361 let req_id = 21;
1362 let pending = PendingRequest {
1363 operation: PendingOperation::Cancel,
1364 client_order_ids: vec![cl_ord_id],
1365 venue_order_ids: vec![Some(VenueOrderId::from(VENUE_ORDER_ID))],
1366 ts_sent_ns: 0,
1367 new_quantity: None,
1368 new_price: None,
1369 new_trigger_price: None,
1370 };
1371 harness.state.pending.insert(req_id, pending);
1372
1373 let response = make_response(
1374 KrakenWsMethod::CancelOrder,
1375 false,
1376 req_id,
1377 None,
1378 Some("Unknown order"),
1379 );
1380 harness.state.handle_response(&response, 5_000);
1381
1382 let event = harness.event_rx.try_recv().expect("event emitted");
1383 match event {
1384 OrderEventAny::CancelRejected(e) => {
1385 assert_eq!(e.client_order_id, cl_ord_id);
1386 assert_eq!(e.reason.as_str(), "Unknown order");
1387 }
1388 other => panic!("expected CancelRejected, was {other:?}"),
1389 }
1390 }
1391
1392 #[rstest]
1393 fn test_handle_response_cancel_success_is_silent() {
1394 let mut harness = make_harness(60_000);
1395 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1396 register_default_identity(&harness.dispatch_state, cl_ord_id);
1397
1398 let req_id = 22;
1399 let pending = PendingRequest {
1400 operation: PendingOperation::Cancel,
1401 client_order_ids: vec![cl_ord_id],
1402 venue_order_ids: vec![Some(VenueOrderId::from(VENUE_ORDER_ID))],
1403 ts_sent_ns: 0,
1404 new_quantity: None,
1405 new_price: None,
1406 new_trigger_price: None,
1407 };
1408 harness.state.pending.insert(req_id, pending);
1409
1410 let response = make_response(KrakenWsMethod::CancelOrder, true, req_id, None, None);
1411 harness.state.handle_response(&response, 6_000);
1412
1413 assert_eq!(harness.state.pending_len(), 0);
1414 assert!(harness.event_rx.try_recv().is_err());
1415 }
1416
1417 #[rstest]
1418 #[case(PendingOperation::Submit)]
1419 #[case(PendingOperation::Amend)]
1420 #[case(PendingOperation::Cancel)]
1421 #[case(PendingOperation::BatchAdd)]
1422 fn test_timeout_emits_no_order_event(#[case] operation: PendingOperation) {
1423 let mut harness = make_harness(60_000);
1424 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1425 register_default_identity(&harness.dispatch_state, cl_ord_id);
1426
1427 let pending = PendingRequest {
1428 operation,
1429 client_order_ids: vec![cl_ord_id],
1430 venue_order_ids: vec![Some(VenueOrderId::from(VENUE_ORDER_ID))],
1431 ts_sent_ns: 0,
1432 new_quantity: None,
1433 new_price: None,
1434 new_trigger_price: None,
1435 };
1436
1437 harness.state.handle_timeout(&pending);
1438
1439 assert!(harness.event_rx.try_recv().is_err());
1440 }
1441
1442 #[tokio::test]
1443 async fn test_late_submit_response_resolves_timed_out_request() {
1444 let mut harness = make_harness(50);
1445 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1446 register_default_identity(&harness.dispatch_state, cl_ord_id);
1447 *harness.auth_token.write().await = Some("TEST-TOKEN".to_string());
1448
1449 let params = make_add_order_params("TEST-TOKEN");
1450 let req_id = harness
1451 .state
1452 .submit(params, make_identity(PendingOperation::Submit), 1)
1453 .expect("submit ok");
1454
1455 let payloads = recv_send_payloads_until_cancel(&mut harness.cmd_rx).await;
1456 assert!(payloads.iter().any(|p| p.contains("\"cancel_order\"")));
1457 assert_eq!(harness.state.pending_len(), 1);
1458 assert!(harness.event_rx.try_recv().is_err());
1459
1460 let response = make_response(
1461 KrakenWsMethod::AddOrder,
1462 true,
1463 req_id,
1464 Some(VENUE_ORDER_ID),
1465 None,
1466 );
1467 harness.state.handle_response(&response, 7_000);
1468
1469 let event = harness.event_rx.try_recv().expect("late response event");
1470 match event {
1471 OrderEventAny::Accepted(e) => {
1472 assert_eq!(e.client_order_id, cl_ord_id);
1473 assert_eq!(e.venue_order_id.as_str(), VENUE_ORDER_ID);
1474 }
1475 other => panic!("expected Accepted, was {other:?}"),
1476 }
1477 assert_eq!(harness.state.pending_len(), 0);
1478 }
1479
1480 #[tokio::test]
1481 async fn test_late_submit_rejection_resolves_timed_out_request() {
1482 let mut harness = make_harness(50);
1483 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1484 register_default_identity(&harness.dispatch_state, cl_ord_id);
1485 *harness.auth_token.write().await = Some("TEST-TOKEN".to_string());
1486
1487 let params = make_add_order_params("TEST-TOKEN");
1488 let req_id = harness
1489 .state
1490 .submit(params, make_identity(PendingOperation::Submit), 1)
1491 .expect("submit ok");
1492
1493 recv_send_payloads_until_cancel(&mut harness.cmd_rx).await;
1494 assert_eq!(harness.state.pending_len(), 1);
1495 assert!(harness.event_rx.try_recv().is_err());
1496
1497 let response = make_response(
1498 KrakenWsMethod::AddOrder,
1499 false,
1500 req_id,
1501 None,
1502 Some("Insufficient funds"),
1503 );
1504 harness.state.handle_response(&response, 7_000);
1505
1506 let event = harness.event_rx.try_recv().expect("late response event");
1507 match event {
1508 OrderEventAny::Rejected(e) => {
1509 assert_eq!(e.client_order_id, cl_ord_id);
1510 assert_eq!(e.reason.as_str(), "Insufficient funds");
1511 }
1512 other => panic!("expected Rejected, was {other:?}"),
1513 }
1514 assert_eq!(harness.state.pending_len(), 0);
1515 }
1516
1517 #[rstest]
1518 fn test_handle_response_method_op_mismatch_retains_pending() {
1519 let mut harness = make_harness(60_000);
1520 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1521 register_default_identity(&harness.dispatch_state, cl_ord_id);
1522
1523 let req_id = 33;
1524 harness
1525 .state
1526 .pending
1527 .insert(req_id, make_identity(PendingOperation::Submit));
1528
1529 let response = make_response(KrakenWsMethod::CancelOrder, true, req_id, None, None);
1530 harness.state.handle_response(&response, 8_000);
1531
1532 assert_eq!(harness.state.pending_len(), 1);
1533 assert!(harness.event_rx.try_recv().is_err());
1534
1535 let response = make_response(
1536 KrakenWsMethod::AddOrder,
1537 true,
1538 req_id,
1539 Some(VENUE_ORDER_ID),
1540 None,
1541 );
1542 harness.state.handle_response(&response, 8_001);
1543
1544 let event = harness
1545 .event_rx
1546 .try_recv()
1547 .expect("matching response event");
1548 match event {
1549 OrderEventAny::Accepted(e) => {
1550 assert_eq!(e.client_order_id, cl_ord_id);
1551 assert_eq!(e.venue_order_id.as_str(), VENUE_ORDER_ID);
1552 }
1553 other => panic!("expected Accepted, was {other:?}"),
1554 }
1555 assert_eq!(harness.state.pending_len(), 0);
1556 }
1557
1558 #[rstest]
1559 fn test_handle_response_batch_add_emits_per_leg_events() {
1560 let mut harness = make_harness(60_000);
1561 let cl_a = ClientOrderId::from("O-A");
1562 let cl_b = ClientOrderId::from("O-B");
1563 register_default_identity(&harness.dispatch_state, cl_a);
1564 register_default_identity(&harness.dispatch_state, cl_b);
1565
1566 let req_id = 50;
1567 let pending = PendingRequest {
1568 operation: PendingOperation::BatchAdd,
1569 client_order_ids: vec![cl_a, cl_b],
1570 venue_order_ids: vec![None, None],
1571 ts_sent_ns: 0,
1572 new_quantity: None,
1573 new_price: None,
1574 new_trigger_price: None,
1575 };
1576 harness.state.pending.insert(req_id, pending);
1577
1578 let response = KrakenWsOrderResponse {
1579 method: KrakenWsMethod::BatchAdd,
1580 req_id: Some(req_id),
1581 success: true,
1582 time_in: None,
1583 time_out: None,
1584 error: None,
1585 result: Some(KrakenWsOrderResult {
1586 order_id: None,
1587 cl_ord_id: None,
1588 order_userref: None,
1589 warning: None,
1590 orders: Some(vec![
1591 KrakenWsBatchOrderResult {
1592 success: true,
1593 order_id: Some("V-A".to_string()),
1594 cl_ord_id: Some("O-A".to_string()),
1595 error: None,
1596 },
1597 KrakenWsBatchOrderResult {
1598 success: false,
1599 order_id: None,
1600 cl_ord_id: Some("O-B".to_string()),
1601 error: Some("Bad price".to_string()),
1602 },
1603 ]),
1604 }),
1605 };
1606 harness.state.handle_response(&response, 9_000);
1607
1608 let first = harness.event_rx.try_recv().expect("first event");
1609 let second = harness.event_rx.try_recv().expect("second event");
1610
1611 match first {
1612 OrderEventAny::Accepted(e) => {
1613 assert_eq!(e.client_order_id, cl_a);
1614 assert_eq!(e.venue_order_id.as_str(), "V-A");
1615 }
1616 other => panic!("expected Accepted, was {other:?}"),
1617 }
1618
1619 match second {
1620 OrderEventAny::Rejected(e) => {
1621 assert_eq!(e.client_order_id, cl_b);
1622 assert_eq!(e.reason.as_str(), "Bad price");
1623 }
1624 other => panic!("expected Rejected, was {other:?}"),
1625 }
1626 }
1627
1628 #[rstest]
1629 fn test_handle_response_batch_add_matches_legs_by_echoed_cl_ord_id() {
1630 let mut harness = make_harness(60_000);
1631 let cl_a = ClientOrderId::from("O-A");
1632 let cl_b = ClientOrderId::from("O-B");
1633 register_default_identity(&harness.dispatch_state, cl_a);
1634 register_default_identity(&harness.dispatch_state, cl_b);
1635
1636 let req_id = 60;
1637 let pending = PendingRequest {
1638 operation: PendingOperation::BatchAdd,
1639 client_order_ids: vec![cl_a, cl_b],
1640 venue_order_ids: vec![None, None],
1641 ts_sent_ns: 0,
1642 new_quantity: None,
1643 new_price: None,
1644 new_trigger_price: None,
1645 };
1646 harness.state.pending.insert(req_id, pending);
1647
1648 let response = KrakenWsOrderResponse {
1652 method: KrakenWsMethod::BatchAdd,
1653 req_id: Some(req_id),
1654 success: true,
1655 time_in: None,
1656 time_out: None,
1657 error: None,
1658 result: Some(KrakenWsOrderResult {
1659 order_id: None,
1660 cl_ord_id: None,
1661 order_userref: None,
1662 warning: None,
1663 orders: Some(vec![
1664 KrakenWsBatchOrderResult {
1665 success: true,
1666 order_id: Some("V-B".to_string()),
1667 cl_ord_id: Some("O-B".to_string()),
1668 error: None,
1669 },
1670 KrakenWsBatchOrderResult {
1671 success: true,
1672 order_id: Some("V-A".to_string()),
1673 cl_ord_id: Some("O-A".to_string()),
1674 error: None,
1675 },
1676 ]),
1677 }),
1678 };
1679 harness.state.handle_response(&response, 14_000);
1680
1681 let first = harness.event_rx.try_recv().expect("first event");
1682 let second = harness.event_rx.try_recv().expect("second event");
1683 let mut by_cl_ord = std::collections::HashMap::new();
1684
1685 for event in [first, second] {
1686 match event {
1687 OrderEventAny::Accepted(e) => {
1688 by_cl_ord.insert(e.client_order_id, e.venue_order_id);
1689 }
1690 other => panic!("expected Accepted, was {other:?}"),
1691 }
1692 }
1693 assert_eq!(
1694 by_cl_ord.get(&cl_a).map(|v| v.as_str()),
1695 Some("V-A"),
1696 "cl_a must be paired with V-A despite reversed response order",
1697 );
1698 assert_eq!(
1699 by_cl_ord.get(&cl_b).map(|v| v.as_str()),
1700 Some("V-B"),
1701 "cl_b must be paired with V-B despite reversed response order",
1702 );
1703 }
1704
1705 #[rstest]
1706 fn test_handle_response_batch_add_truncated_per_leg_results_rejects_trailing_legs() {
1707 let mut harness = make_harness(60_000);
1708 let cl_a = ClientOrderId::from("O-A");
1709 let cl_b = ClientOrderId::from("O-B");
1710 let cl_c = ClientOrderId::from("O-C");
1711 register_default_identity(&harness.dispatch_state, cl_a);
1712 register_default_identity(&harness.dispatch_state, cl_b);
1713 register_default_identity(&harness.dispatch_state, cl_c);
1714
1715 let req_id = 51;
1716 let pending = PendingRequest {
1717 operation: PendingOperation::BatchAdd,
1718 client_order_ids: vec![cl_a, cl_b, cl_c],
1719 venue_order_ids: vec![None, None, None],
1720 ts_sent_ns: 0,
1721 new_quantity: None,
1722 new_price: None,
1723 new_trigger_price: None,
1724 };
1725 harness.state.pending.insert(req_id, pending);
1726
1727 let response = KrakenWsOrderResponse {
1729 method: KrakenWsMethod::BatchAdd,
1730 req_id: Some(req_id),
1731 success: true,
1732 time_in: None,
1733 time_out: None,
1734 error: None,
1735 result: Some(KrakenWsOrderResult {
1736 order_id: None,
1737 cl_ord_id: None,
1738 order_userref: None,
1739 warning: None,
1740 orders: Some(vec![KrakenWsBatchOrderResult {
1741 success: true,
1742 order_id: Some("V-A".to_string()),
1743 cl_ord_id: Some("O-A".to_string()),
1744 error: None,
1745 }]),
1746 }),
1747 };
1748 harness.state.handle_response(&response, 12_000);
1749
1750 let first = harness.event_rx.try_recv().expect("first event");
1751 let second = harness.event_rx.try_recv().expect("second event");
1752 let third = harness.event_rx.try_recv().expect("third event");
1753
1754 match first {
1755 OrderEventAny::Accepted(e) => assert_eq!(e.client_order_id, cl_a),
1756 other => panic!("expected Accepted for present leg, was {other:?}"),
1757 }
1758
1759 for (event, cl_id) in [(second, cl_b), (third, cl_c)] {
1760 match event {
1761 OrderEventAny::Rejected(e) => {
1762 assert_eq!(
1763 e.client_order_id, cl_id,
1764 "missing-leg rejection cl_ord_id mismatch",
1765 );
1766 assert!(
1767 e.reason.as_str().contains("missing per-leg result"),
1768 "expected truncation reason, was {}",
1769 e.reason,
1770 );
1771 }
1772 other => panic!(
1773 "missing per-leg result must reject (not inherit envelope.success), was {other:?}",
1774 ),
1775 }
1776 }
1777 }
1778
1779 fn drain_send_payloads(
1780 rx: &mut tokio::sync::mpsc::UnboundedReceiver<SpotHandlerCommand>,
1781 ) -> Vec<String> {
1782 let mut out = Vec::new();
1783
1784 while let Ok(cmd) = rx.try_recv() {
1785 if let SpotHandlerCommand::SendOrderRequest { payload, .. } = cmd {
1786 out.push(payload);
1787 }
1788 }
1789 out
1790 }
1791
1792 async fn recv_send_payloads_until_cancel(
1796 rx: &mut tokio::sync::mpsc::UnboundedReceiver<SpotHandlerCommand>,
1797 ) -> Vec<String> {
1798 let mut out = Vec::new();
1799
1800 loop {
1801 let cmd = tokio::time::timeout(Duration::from_secs(5), rx.recv())
1802 .await
1803 .expect("timed out awaiting compensating cancel")
1804 .expect("command channel closed");
1805
1806 let SpotHandlerCommand::SendOrderRequest { payload, .. } = cmd else {
1807 continue;
1808 };
1809 let is_cancel = payload.contains("\"cancel_order\"");
1810 out.push(payload);
1811
1812 if is_cancel {
1813 return out;
1814 }
1815 }
1816 }
1817
1818 #[tokio::test]
1819 async fn test_submit_timeout_sends_compensating_cancel() {
1820 let mut harness = make_harness(50);
1821 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1822 register_default_identity(&harness.dispatch_state, cl_ord_id);
1823 *harness.auth_token.write().await = Some("TEST-TOKEN".to_string());
1824
1825 let params = make_add_order_params("TEST-TOKEN");
1826 let identity = make_identity(PendingOperation::Submit);
1827 harness
1828 .state
1829 .submit(params, identity, 1)
1830 .expect("submit ok");
1831
1832 let payloads = recv_send_payloads_until_cancel(&mut harness.cmd_rx).await;
1833 assert!(
1834 payloads.iter().any(|p| p.contains("\"add_order\"")),
1835 "original add_order missing: {payloads:?}",
1836 );
1837 let cancel = payloads
1838 .iter()
1839 .find(|p| p.contains("\"cancel_order\""))
1840 .expect("compensating cancel_order missing");
1841 assert!(
1842 cancel.contains(CLIENT_ORDER_ID),
1843 "compensating cancel must reference cl_ord_id, was {cancel}",
1844 );
1845
1846 assert_eq!(harness.state.pending_len(), 1);
1847 assert!(harness.event_rx.try_recv().is_err());
1848 }
1849
1850 #[tokio::test]
1851 async fn test_compensating_cancel_response_is_silently_dropped() {
1852 let mut harness = make_harness(50);
1856 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1857 register_default_identity(&harness.dispatch_state, cl_ord_id);
1858 *harness.auth_token.write().await = Some("TEST-TOKEN".to_string());
1859
1860 let params = make_add_order_params("TEST-TOKEN");
1861 let identity = make_identity(PendingOperation::Submit);
1862 harness
1863 .state
1864 .submit(params, identity, 1)
1865 .expect("submit ok");
1866
1867 let payloads = recv_send_payloads_until_cancel(&mut harness.cmd_rx).await;
1868
1869 let cancel = payloads
1870 .iter()
1871 .find(|p| p.contains("\"cancel_order\""))
1872 .expect("compensating cancel missing");
1873 let cancel_value: serde_json::Value = serde_json::from_str(cancel).expect("valid json");
1874 let cancel_req_id = cancel_value["req_id"].as_u64().expect("req_id present");
1875
1876 for success in [true, false] {
1877 let response = KrakenWsOrderResponse {
1878 method: KrakenWsMethod::CancelOrder,
1879 req_id: Some(cancel_req_id),
1880 success,
1881 time_in: None,
1882 time_out: None,
1883 error: (!success).then(|| "Unknown order".to_string()),
1884 result: None,
1885 };
1886 harness.state.handle_response(&response, 9_000);
1887 }
1888
1889 assert!(
1890 harness.event_rx.try_recv().is_err(),
1891 "compensating-cancel responses must not surface events to strategies",
1892 );
1893 assert_eq!(harness.state.pending_len(), 1);
1894 }
1895
1896 #[rstest]
1897 fn test_submit_timeout_without_token_skips_compensating_cancel() {
1898 let mut harness = make_harness(60_000);
1899 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
1900 register_default_identity(&harness.dispatch_state, cl_ord_id);
1901 let identity = make_identity(PendingOperation::Submit);
1902
1903 harness.state.handle_timeout(&identity);
1904
1905 let payloads = drain_send_payloads(&mut harness.cmd_rx);
1906 assert!(
1907 payloads.iter().all(|p| !p.contains("\"cancel_order\"")),
1908 "no compensating cancel expected without token, was {payloads:?}",
1909 );
1910 }
1911
1912 #[tokio::test]
1913 async fn test_batch_add_timeout_sends_compensating_cancel_for_all_legs() {
1914 let mut harness = make_harness(50);
1915 let cl_a = ClientOrderId::from("O-A");
1916 let cl_b = ClientOrderId::from("O-B");
1917 register_default_identity(&harness.dispatch_state, cl_a);
1918 register_default_identity(&harness.dispatch_state, cl_b);
1919 *harness.auth_token.write().await = Some("TEST-TOKEN".to_string());
1920
1921 let params = KrakenWsBatchAddParams {
1922 symbol: "BTC/USD".to_string(),
1923 orders: vec![],
1924 token: "TEST-TOKEN".to_string(),
1925 };
1926 let identity = PendingRequest {
1927 operation: PendingOperation::BatchAdd,
1928 client_order_ids: vec![cl_a, cl_b],
1929 venue_order_ids: vec![None, None],
1930 ts_sent_ns: 0,
1931 new_quantity: None,
1932 new_price: None,
1933 new_trigger_price: None,
1934 };
1935 harness
1936 .state
1937 .batch_add(params, identity, 1)
1938 .expect("batch ok");
1939
1940 let payloads = recv_send_payloads_until_cancel(&mut harness.cmd_rx).await;
1941 let cancel = payloads
1942 .iter()
1943 .find(|p| p.contains("\"cancel_order\""))
1944 .expect("compensating cancel missing");
1945 assert!(cancel.contains("O-A") && cancel.contains("O-B"));
1946 assert_eq!(harness.state.pending_len(), 1);
1947 assert!(harness.event_rx.try_recv().is_err());
1948
1949 let batch_req_id = payloads
1950 .iter()
1951 .find(|p| p.contains("\"batch_add\""))
1952 .and_then(|p| serde_json::from_str::<serde_json::Value>(p).ok())
1953 .and_then(|v| v["req_id"].as_u64())
1954 .expect("batch req_id missing");
1955 let response = KrakenWsOrderResponse {
1956 method: KrakenWsMethod::BatchAdd,
1957 req_id: Some(batch_req_id),
1958 success: true,
1959 time_in: None,
1960 time_out: None,
1961 error: None,
1962 result: Some(KrakenWsOrderResult {
1963 order_id: None,
1964 cl_ord_id: None,
1965 order_userref: None,
1966 warning: None,
1967 orders: Some(vec![
1968 KrakenWsBatchOrderResult {
1969 success: false,
1970 order_id: None,
1971 cl_ord_id: Some("O-B".to_string()),
1972 error: Some("Bad price".to_string()),
1973 },
1974 KrakenWsBatchOrderResult {
1975 success: true,
1976 order_id: Some("V-A".to_string()),
1977 cl_ord_id: Some("O-A".to_string()),
1978 error: None,
1979 },
1980 ]),
1981 }),
1982 };
1983 harness.state.handle_response(&response, 10_000);
1984
1985 let first = harness.event_rx.try_recv().expect("first late event");
1986 let second = harness.event_rx.try_recv().expect("second late event");
1987 let mut accepted = None;
1988 let mut rejected = None;
1989
1990 for event in [first, second] {
1991 match event {
1992 OrderEventAny::Accepted(e) => {
1993 accepted = Some((e.client_order_id, e.venue_order_id));
1994 }
1995 OrderEventAny::Rejected(e) => {
1996 rejected = Some((e.client_order_id, e.reason));
1997 }
1998 other => panic!("expected Accepted or Rejected, was {other:?}"),
1999 }
2000 }
2001 assert_eq!(
2002 accepted.map(|(cl, venue)| (cl, venue.to_string())),
2003 Some((cl_a, "V-A".to_string()))
2004 );
2005 assert_eq!(
2006 rejected.map(|(cl, reason)| (cl, reason.to_string())),
2007 Some((cl_b, "Bad price".to_string()))
2008 );
2009 assert_eq!(harness.state.pending_len(), 0);
2010 }
2011
2012 #[tokio::test]
2013 async fn test_clear_removes_timed_out_request() {
2014 let mut harness = make_harness(50);
2015 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
2016 register_default_identity(&harness.dispatch_state, cl_ord_id);
2017 *harness.auth_token.write().await = Some("TEST-TOKEN".to_string());
2018
2019 let params = make_add_order_params("TEST-TOKEN");
2020 harness
2021 .state
2022 .submit(params, make_identity(PendingOperation::Submit), 1)
2023 .expect("submit ok");
2024
2025 recv_send_payloads_until_cancel(&mut harness.cmd_rx).await;
2026 assert_eq!(harness.state.pending_len(), 1);
2027
2028 harness.state.clear();
2029
2030 assert_eq!(harness.state.pending_len(), 0);
2031 }
2032
2033 #[tokio::test]
2034 async fn test_reset_cancellation_token_keeps_new_timeout_pending() {
2035 let mut harness = make_harness(50);
2036 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
2037 register_default_identity(&harness.dispatch_state, cl_ord_id);
2038 *harness.auth_token.write().await = Some("TEST-TOKEN".to_string());
2039 harness.pending_tasks.begin_shutdown();
2040 harness
2041 .pending_tasks
2042 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(1))
2043 .await
2044 .expect("finish old timeout generation");
2045 harness
2046 .pending_tasks
2047 .start_generation()
2048 .expect("start timeout generation");
2049 let pending_spawner = harness
2050 .pending_tasks
2051 .spawner()
2052 .expect("pending task spawner");
2053 harness.state.reset_task_spawner(pending_spawner);
2054
2055 harness
2056 .state
2057 .submit(
2058 make_add_order_params("TEST-TOKEN"),
2059 make_identity(PendingOperation::Submit),
2060 1,
2061 )
2062 .expect("submit ok");
2063
2064 recv_send_payloads_until_cancel(&mut harness.cmd_rx).await;
2065
2066 assert_eq!(harness.state.pending_len(), 1);
2067 assert!(harness.event_rx.try_recv().is_err());
2068 }
2069
2070 #[tokio::test]
2071 async fn test_task_group_shutdown_aborts_pending_timeout() {
2072 let harness = make_harness(60_000);
2073 let cl_ord_id = ClientOrderId::from(CLIENT_ORDER_ID);
2074 register_default_identity(&harness.dispatch_state, cl_ord_id);
2075
2076 let params = KrakenWsAddOrderParams {
2077 order_type: KrakenOrderType::Limit,
2078 side: KrakenOrderSide::Buy,
2079 order_qty: dec!(0.001),
2080 symbol: "BTC/USD".to_string(),
2081 token: String::new(),
2082 limit_price: Some(dec!(50000)),
2083 time_in_force: None,
2084 expire_time: None,
2085 cl_ord_id: Some(CLIENT_ORDER_ID.to_string()),
2086 post_only: None,
2087 reduce_only: None,
2088 leverage: None,
2089 trigger: None,
2090 conditional: None,
2091 };
2092 let identity = make_identity(PendingOperation::Submit);
2093 harness
2094 .state
2095 .submit(params, identity, 1)
2096 .expect("submit ok");
2097 assert_eq!(harness.state.pending_len(), 1);
2098
2099 harness.pending_tasks.begin_shutdown();
2100
2101 nautilus_common::testing::wait_until_async(
2104 || {
2105 let state = Arc::clone(&harness.state);
2106 async move { state.pending_len() == 0 }
2107 },
2108 Duration::from_secs(5),
2109 )
2110 .await;
2111
2112 assert_eq!(
2113 harness.state.pending_len(),
2114 0,
2115 "pending entry must be cleared on cancellation",
2116 );
2117 }
2118}
2119
2120#[cfg(test)]
2121mod property_tests {
2122 use nautilus_model::identifiers::ClientOrderId;
2123 use proptest::prelude::*;
2124 use rstest::rstest;
2125
2126 use super::{tests::make_harness, *};
2127 use crate::websocket::spot_v2::messages::{KrakenWsBatchOrderResult, KrakenWsOrderResult};
2128
2129 proptest! {
2130 #[rstest]
2131 fn no_pending_leak_on_random_response_interleaving(
2132 ops in proptest::collection::vec(
2133 prop_oneof![
2134 Just((PendingOperation::Submit, true, 1usize)),
2135 Just((PendingOperation::Submit, false, 1usize)),
2136 Just((PendingOperation::Amend, true, 1usize)),
2137 Just((PendingOperation::Amend, false, 1usize)),
2138 Just((PendingOperation::Cancel, true, 1usize)),
2139 Just((PendingOperation::Cancel, false, 1usize)),
2140 (2usize..=4usize).prop_map(|n| (PendingOperation::BatchAdd, true, n)),
2141 (2usize..=4usize).prop_map(|n| (PendingOperation::BatchAdd, false, n)),
2142 ],
2143 1..50usize,
2144 )
2145 ) {
2146 let h = make_harness(60_000);
2147 let state = &h.state;
2148 let mut req_ids = Vec::new();
2149
2150 for (op, success, leg_count) in &ops {
2151 let req_id = state.next_req_id();
2152 let client_order_ids: Vec<ClientOrderId> = (0..*leg_count)
2153 .map(|i| ClientOrderId::from(format!("O-{req_id}-{i}").as_str()))
2154 .collect();
2155 let venue_order_ids = vec![None; *leg_count];
2156 state.pending.insert(req_id, PendingRequest {
2157 operation: *op,
2158 client_order_ids: client_order_ids.clone(),
2159 venue_order_ids,
2160 ts_sent_ns: 0,
2161 new_quantity: None,
2162 new_price: None,
2163 new_trigger_price: None,
2164 });
2165 req_ids.push((req_id, *op, *success, client_order_ids));
2166 }
2167 let mut shuffled = req_ids.clone();
2168 shuffled.reverse();
2169 for (req_id, op, success, client_order_ids) in shuffled {
2170 let method = match op {
2171 PendingOperation::Submit => KrakenWsMethod::AddOrder,
2172 PendingOperation::Amend => KrakenWsMethod::AmendOrder,
2173 PendingOperation::Cancel => KrakenWsMethod::CancelOrder,
2174 PendingOperation::BatchAdd => KrakenWsMethod::BatchAdd,
2175 };
2176 let result = if op == PendingOperation::BatchAdd {
2177 Some(KrakenWsOrderResult {
2178 order_id: None,
2179 cl_ord_id: None,
2180 order_userref: None,
2181 warning: None,
2182 orders: Some(client_order_ids
2183 .iter()
2184 .enumerate()
2185 .map(|(i, cid)| KrakenWsBatchOrderResult {
2186 success,
2187 order_id: success.then(|| format!("V-{req_id}-{i}")),
2188 cl_ord_id: Some(cid.as_str().to_string()),
2189 error: (!success).then(|| "test-error".to_string()),
2190 })
2191 .collect()),
2192 })
2193 } else {
2194 None
2195 };
2196 let response = KrakenWsOrderResponse {
2197 method,
2198 req_id: Some(req_id),
2199 success,
2200 time_in: None,
2201 time_out: None,
2202 error: if success {
2203 None
2204 } else {
2205 Some("test-error".to_string())
2206 },
2207 result,
2208 };
2209 state.handle_response(&response, 1);
2210 }
2211 prop_assert_eq!(state.pending_len(), 0);
2212 }
2213 }
2214}