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