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