Skip to main content

nautilus_kraken/websocket/dispatch/
spot_orders.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! WebSocket order-request dispatch state for Kraken Spot v2.
17
18use 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    /// Pending add_order request.
58    Submit,
59    /// Pending amend_order request.
60    Amend,
61    /// Pending cancel_order request.
62    Cancel,
63    /// Pending batch_add request.
64    BatchAdd,
65}
66
67#[derive(Debug, Clone)]
68pub struct PendingRequest {
69    /// The pending operation type.
70    pub operation: PendingOperation,
71    /// Client order IDs associated with this request.
72    pub client_order_ids: Vec<ClientOrderId>,
73    /// Venue order IDs associated with this request, if known.
74    pub venue_order_ids: Vec<Option<VenueOrderId>>,
75    /// UNIX nanosecond timestamp when the request was sent.
76    pub ts_sent_ns: u64,
77    /// New quantity carried for amend operations so the success-path
78    /// `OrderUpdated` event reflects the requested change rather than the
79    /// originally registered quantity. `None` for non-amend operations or
80    /// amends that do not modify quantity.
81    pub new_quantity: Option<Quantity>,
82    /// New limit price carried for amend operations so the success-path
83    /// `OrderUpdated` event surfaces the amended price to strategies.
84    /// `None` for non-amend operations or amends that do not modify price.
85    pub new_price: Option<Price>,
86    /// New trigger price carried for amend operations on conditional orders.
87    /// `None` for non-amend operations or amends that do not modify the
88    /// trigger price.
89    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    /// Shared handle to the WS handler cmd_tx. Read at send time, not at
98    /// construction: `connect()` swaps the inner sender, and the
99    /// post-`new()` placeholder's receiver is already dropped.
100    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    /// WS auth token shared with the client; used to build compensating
106    /// cancels when a submit request times out.
107    auth_token: Arc<tokio::sync::RwLock<Option<SecretString>>>,
108    task_spawner: RwLock<TaskSpawner>,
109    /// Clock used to stamp local timeout diagnostics.
110    clock: &'static AtomicTime,
111}
112
113impl OrderRequestState {
114    /// Creates a new [`OrderRequestState`].
115    #[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    /// Reads a clone of the current cmd_tx. Returns `None` only when the
146    /// handle's `RwLock` is being written to (connect/disconnect window).
147    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    /// Returns the next request ID and advances the counter.
152    pub fn next_req_id(&self) -> u64 {
153        self.req_id_counter.fetch_add(1, Ordering::Relaxed) + 1
154    }
155
156    /// Sends an add_order request over the WebSocket transport.
157    ///
158    /// # Errors
159    ///
160    /// Returns an error if the JSON envelope fails to serialize or if the
161    /// handler command channel is closed.
162    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    /// Sends an amend_order request over the WebSocket transport.
180    ///
181    /// # Errors
182    ///
183    /// Returns an error if serialization fails or the handler command channel is closed.
184    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    /// Sends a cancel_order request over the WebSocket transport.
202    ///
203    /// # Errors
204    ///
205    /// Returns an error if serialization fails or the handler command channel is closed.
206    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    /// Sends a batch_add request over the WebSocket transport.
224    ///
225    /// # Errors
226    ///
227    /// Returns an error if serialization fails or the handler command channel is closed.
228    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    /// Routes an order-method WebSocket response to the correct event emitter.
318    ///
319    /// `ts_event_ns` is the local receipt time. Kraken's `time_in`/`time_out`
320    /// are ignored to keep event ordering monotonic against the local clock.
321    /// Timed-out requests remain correlated until a definitive response arrives.
322    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    /// Sends a best-effort `cancel_order` over the WebSocket after a Submit or
405    /// BatchAdd request has timed out, defending against the case where the
406    /// venue accepted the order but the response was delayed past the
407    /// configured timeout window.
408    ///
409    /// Fire-and-forget: the response is silently dropped because no `pending`
410    /// entry is registered for the cancel `req_id`. If the auth token is not
411    /// available or the command channel is closed the cancel is skipped and
412    /// the original request remains resolvable by a late response or
413    /// reconciliation.
414    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        // Match per-leg results by the cl_ord_id Kraken echoes back rather than
693        // relying purely on positional alignment with the legs we sent. Echoed
694        // matches are robust against the venue reordering or omitting entries
695        // in the response array; positional fallback covers the (currently
696        // dominant) case where the echoed cl_ord_id is missing.
697        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    /// Forwards an order event to the execution-engine channel.
771    ///
772    /// A send failure means the receiver has been dropped, which only happens
773    /// during shutdown after the forwarder task has exited. Events lost in
774    /// that window are recovered on the next start by the reconciliation
775    /// engine (`open_check_interval_secs`), which queries the venue for any
776    /// orders or fills the local cache is missing. The error log preserves
777    /// operator visibility for the (rare) shutdown-race case.
778    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        // Regression: dispatcher built before connect() must observe the
997        // live cmd_tx after the swap, not the dropped placeholder.
998        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        // Vendor returns per-leg results in REVERSE order vs sent. Echo-based
1678        // matching must attribute V-B to cl_b and V-A to cl_a regardless of
1679        // index alignment.
1680        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        // Vendor returns envelope.success=true but only one per-leg entry for three sent.
1757        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    // The compensating cancel is the last command the timeout task emits, so
1822    // awaiting it means timeout handling is complete without racing a fixed
1823    // sleep against the global runtime.
1824    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        // The compensating cancel after a submit timeout is fire-and-forget;
1885        // its response must not surface an event or resolve the original
1886        // request correlation.
1887        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        // Cancellation clears the pending entry from a task on the global runtime,
2137        // so poll the condition rather than racing a fixed sleep.
2138        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}