Skip to main content

nautilus_kraken/execution/
spot.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//! Kraken Spot execution client implementation.
17
18use std::{
19    collections::HashSet,
20    future::Future,
21    sync::Arc,
22    time::{Duration, Instant},
23};
24
25use anyhow::Context;
26use async_trait::async_trait;
27use futures_util::StreamExt;
28use jiff::Timestamp;
29use nautilus_common::{
30    cache::InstrumentLookupError,
31    clients::ExecutionClient,
32    live::runner::get_exec_event_sender,
33    messages::execution::{
34        BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
35        GenerateOrderStatusReport, GenerateOrderStatusReports, GeneratePositionStatusReports,
36        ModifyOrder, QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList,
37    },
38};
39use nautilus_core::{
40    AtomicMap, Params, UnixNanos,
41    time::{AtomicTime, get_atomic_clock_realtime},
42};
43use nautilus_live::{
44    ExecutionClientCore, ExecutionEventEmitter, SocketControl, execution::failure::CommandFailure,
45    task::TaskGroup,
46};
47use nautilus_model::{
48    accounts::AccountAny,
49    enums::{
50        AccountType, OmsType, OrderSide, OrderType, PositionSide, TimeInForce, TrailingOffsetType,
51        TriggerType,
52    },
53    events::OrderEventAny,
54    identifiers::{
55        AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, Venue, VenueOrderId,
56    },
57    instruments::{Instrument, InstrumentAny},
58    orders::{Order, OrderAny},
59    reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
60    types::{AccountBalance, MarginBalance, Price, Quantity},
61};
62use rust_decimal::Decimal;
63use tokio_util::sync::CancellationToken;
64use ustr::Ustr;
65
66use super::{
67    command_failure_from_cancel_error, command_failure_from_modify_error,
68    command_failure_from_spot_batch_error, command_failure_from_spot_batch_item,
69    command_failure_from_spot_cancel_error, command_failure_from_submit_error,
70};
71use crate::{
72    common::{
73        consts::{KRAKEN_SPOT_POST_ONLY_ERROR, KRAKEN_VENUE},
74        enums::{
75            KrakenOrderSide, KrakenOrderType, KrakenProductType, KrakenSpotTrigger,
76            KrakenTimeInForce, product_type_from_symbol,
77        },
78        order_params::{
79            build_add_order_params, build_amend_order_params, build_cancel_order_params,
80            compute_ws_time_in_force, format_expire_time,
81        },
82        parse::truncate_cl_ord_id,
83    },
84    config::KrakenExecutionClientConfig,
85    http::{
86        KrakenSpotCancelOrderBatchParams, KrakenSpotCancelOrderParamsBuilder, KrakenSpotHttpClient,
87        spot::client::KRAKEN_SPOT_DEFAULT_RATE_LIMIT_PER_SECOND,
88    },
89    websocket::{
90        dispatch::{
91            self, OrderIdentity, WsDispatchState,
92            spot_orders::{OrderRequestState, PendingOperation, PendingRequest},
93        },
94        spot_v2::{
95            client::KrakenSpotWebSocketClient,
96            messages::{
97                KrakenSpotWsMessage, KrakenWsBatchAddOrder, KrakenWsBatchAddParams,
98                KrakenWsTriggerParams,
99            },
100        },
101    },
102};
103
104/// Kraken Spot execution client.
105///
106/// Provides order management and account operations for Kraken Spot markets.
107#[allow(dead_code)]
108#[derive(Debug)]
109pub struct KrakenSpotExecutionClient {
110    core: ExecutionClientCore,
111    clock: &'static AtomicTime,
112    config: KrakenExecutionClientConfig,
113    emitter: ExecutionEventEmitter,
114    http: KrakenSpotHttpClient,
115    ws: KrakenSpotWebSocketClient,
116    cancellation_token: CancellationToken,
117    session_tasks: TaskGroup,
118    pending_tasks: TaskGroup,
119    instruments: Arc<AtomicMap<InstrumentId, InstrumentAny>>,
120    order_qty_cache: Arc<AtomicMap<String, Decimal>>,
121    truncated_id_map: Arc<AtomicMap<String, ClientOrderId>>,
122    ws_dispatch_state: Arc<WsDispatchState>,
123    order_request_state: Arc<OrderRequestState>,
124    order_event_rx: Option<tokio::sync::mpsc::UnboundedReceiver<OrderEventAny>>,
125}
126
127impl KrakenSpotExecutionClient {
128    /// Creates a new [`KrakenSpotExecutionClient`].
129    pub fn new(
130        core: ExecutionClientCore,
131        config: KrakenExecutionClientConfig,
132    ) -> anyhow::Result<Self> {
133        let clock = get_atomic_clock_realtime();
134        let emitter = ExecutionEventEmitter::new(
135            clock,
136            core.trader_id,
137            core.account_id,
138            config.spot_account_type,
139            None,
140        );
141
142        let session_tasks = TaskGroup::new();
143        let cancellation_token = session_tasks.cancellation_token();
144        let pending_tasks = TaskGroup::new();
145        let pending_spawner = pending_tasks
146            .spawner()
147            .context("Kraken Spot execution task admission is closed")?;
148
149        let http = KrakenSpotHttpClient::with_credentials(
150            config.api_key.clone(),
151            config.api_secret.clone(),
152            config.environment,
153            config.base_url.clone(),
154            config.timeout_secs,
155            None,
156            None,
157            None,
158            config.proxy_url.clone(),
159            config
160                .max_requests_per_second
161                .unwrap_or(KRAKEN_SPOT_DEFAULT_RATE_LIMIT_PER_SECOND),
162        )?;
163
164        let data_config = crate::config::KrakenDataClientConfig {
165            api_key: Some(config.api_key.clone()),
166            api_secret: Some(config.api_secret.clone()),
167            product_type: config.product_type,
168            environment: config.environment,
169            base_url: config.base_url.clone(),
170            ws_public_url: None,
171            ws_private_url: Some(config.ws_url()),
172            ws_l3_url: None,
173            validate_l3_checksum: true,
174            proxy_url: config.proxy_url.clone(),
175            timeout_secs: config.timeout_secs,
176            heartbeat_interval_secs: config.heartbeat_interval_secs,
177            // The authenticated order-routing WS keeps its prior behavior (no
178            // idle timeout); issue #4255 concerns the public spot data WS.
179            ws_idle_timeout_ms: 0,
180            max_requests_per_second: config.max_requests_per_second,
181            transport_backend: config.transport_backend,
182        };
183        let ws = KrakenSpotWebSocketClient::new(
184            data_config,
185            cancellation_token.clone(),
186            config.proxy_url.clone(),
187        )
188        .with_socket_control(SocketControl::new(
189            core.client_id,
190            Some(*KRAKEN_VENUE),
191            "kraken-spot-user-streams",
192        ));
193
194        let ws_dispatch_state = Arc::new(WsDispatchState::new());
195        // Connect() swaps in a live cmd_tx; capture the shared handle so the
196        // dispatcher reads the current sender, not the dropped placeholder.
197        let cmd_tx_handle = ws.handler_command_handle();
198        let (order_event_tx, order_event_rx) = tokio::sync::mpsc::unbounded_channel();
199        let order_request_state = Arc::new(OrderRequestState::new(
200            cmd_tx_handle,
201            order_event_tx,
202            Arc::clone(&ws_dispatch_state),
203            ws.req_id_counter(),
204            Duration::from_secs(config.ws_request_timeout_secs),
205            core.trader_id,
206            core.account_id,
207            ws.auth_token_handle(),
208            pending_spawner,
209            clock,
210        ));
211
212        Ok(Self {
213            core,
214            clock,
215            config,
216            emitter,
217            http,
218            ws,
219            cancellation_token,
220            session_tasks,
221            pending_tasks,
222            instruments: Arc::new(AtomicMap::new()),
223            order_qty_cache: Arc::new(AtomicMap::new()),
224            truncated_id_map: Arc::new(AtomicMap::new()),
225            ws_dispatch_state,
226            order_request_state,
227            order_event_rx: Some(order_event_rx),
228        })
229    }
230
231    fn register_order_identity(&self, order: &OrderAny) {
232        // Quote-quantity orders submit a quote amount (e.g. 100 USD), but the
233        // venue reports fills in base units (e.g. 0.001 BTC). Registering the
234        // raw `order.quantity()` would make the cumulative-fill comparison in
235        // the fill-side dispatch mismatch base against quote, leaving the
236        // order "open" forever. These orders instead flow through the
237        // untracked path and the engine reconciles them from status reports.
238        if order.is_quote_quantity() {
239            return;
240        }
241        self.ws_dispatch_state.register_identity(
242            order.client_order_id(),
243            OrderIdentity {
244                strategy_id: order.strategy_id(),
245                instrument_id: order.instrument_id(),
246                order_side: order.order_side(),
247                order_type: order.order_type(),
248                quantity: order.quantity(),
249            },
250        );
251    }
252
253    /// Returns a reference to the clock.
254    #[must_use]
255    pub fn clock(&self) -> &'static AtomicTime {
256        self.clock
257    }
258
259    /// Returns a reference to the event emitter.
260    #[must_use]
261    pub fn emitter(&self) -> &ExecutionEventEmitter {
262        &self.emitter
263    }
264
265    fn spawn_task<F>(&self, description: &'static str, fut: F)
266    where
267        F: Future<Output = anyhow::Result<()>> + Send + 'static,
268    {
269        let future = async move {
270            if let Err(e) = fut.await {
271                log::warn!("{description} failed: {e:?}");
272            }
273        };
274
275        if let Err(e) = self.pending_tasks.spawn(future) {
276            log::warn!("Skipping Kraken Spot {description} after shutdown began: {e}");
277        }
278    }
279
280    async fn finish_tasks(&self) -> anyhow::Result<()> {
281        let (session_result, pending_result) = tokio::join!(
282            self.session_tasks
283                .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
284            self.pending_tasks
285                .finish_shutdown(Duration::from_secs(2), Duration::from_secs(2)),
286        );
287        session_result.context("failed to finish Kraken Spot execution session tasks")?;
288        pending_result.context("failed to finish Kraken Spot execution command tasks")?;
289        Ok(())
290    }
291
292    async fn prepare_task_groups(&mut self) -> anyhow::Result<()> {
293        if !self.session_tasks.is_open() || !self.pending_tasks.is_open() {
294            self.session_tasks.begin_shutdown();
295            self.pending_tasks.begin_shutdown();
296            self.finish_tasks().await?;
297            self.session_tasks
298                .start_generation()
299                .context("failed to start Kraken Spot execution session task generation")?;
300            self.pending_tasks
301                .start_generation()
302                .context("failed to start Kraken Spot execution command task generation")?;
303            self.cancellation_token = self.session_tasks.cancellation_token();
304            let pending_spawner = self
305                .pending_tasks
306                .spawner()
307                .context("Kraken Spot execution task admission is closed")?;
308            self.order_request_state.reset_task_spawner(pending_spawner);
309        }
310        Ok(())
311    }
312
313    async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
314        self.http.cancel_all_requests();
315        self.pending_tasks.begin_shutdown();
316        self.order_request_state.clear();
317        let ws_result = self.ws.close().await;
318        self.session_tasks.begin_shutdown();
319        let tasks_result = self.finish_tasks().await;
320        self.core.set_disconnected();
321        tasks_result?;
322        Ok(ws_result?)
323    }
324
325    fn submit_single_order(
326        &self,
327        command: &SubmitOrder,
328        order: &OrderAny,
329        task_name: &'static str,
330        leverage: Option<u16>,
331    ) {
332        if order.is_closed() {
333            log::warn!(
334                "Cannot submit closed order: client_order_id={}",
335                order.client_order_id()
336            );
337            return;
338        }
339
340        let order_type = order.order_type();
341        let time_in_force = order.time_in_force();
342
343        if time_in_force == TimeInForce::Fok && order_type != OrderType::Limit {
344            self.emitter.emit_order_denied(
345                order,
346                "FOK time in force only supported for LIMIT orders on Kraken Spot",
347            );
348            return;
349        }
350
351        if matches!(
352            order_type,
353            OrderType::TrailingStopMarket | OrderType::TrailingStopLimit
354        ) && let Some(offset_type) = order.trailing_offset_type()
355            && offset_type != TrailingOffsetType::Price
356        {
357            self.emitter.emit_order_denied(
358                order,
359                &format!(
360                    "Kraken Spot only supports Price trailing offset type: received {offset_type:?}"
361                ),
362            );
363            return;
364        }
365
366        if order.is_reduce_only() && self.config.spot_account_type == AccountType::Cash {
367            self.emitter
368                .emit_order_denied(order, "reduce_only requires spot_account_type=Margin");
369            return;
370        }
371
372        let client_order_id = order.client_order_id();
373
374        log::debug!("OrderSubmitted: client_order_id={client_order_id}");
375        self.register_order_identity(order);
376        self.emitter.emit_order_submitted(order);
377
378        let kraken_cl_ord_id = truncate_cl_ord_id(&client_order_id);
379
380        if !order.is_quote_quantity() {
381            self.order_qty_cache
382                .insert(kraken_cl_ord_id.clone(), order.quantity().as_decimal());
383        }
384
385        if kraken_cl_ord_id != client_order_id.as_str() {
386            self.truncated_id_map
387                .insert(kraken_cl_ord_id, client_order_id);
388        }
389
390        // Quote-quantity orders submit a quote-currency amount but the venue echoes
391        // fills in base units. The WS dispatch identity is intentionally not
392        // registered for these (see register_order_identity), so a WS round-trip
393        // would never emit OrderAccepted to the strategy. Plus order.quantity()
394        // for quote-qty is a quote amount that would be wrongly sent as `order_qty`
395        // (base units) on the WS path. Force REST.
396        let use_ws_trade = resolve_use_ws_trade(command.params.as_ref(), self.config.use_ws_trade);
397        if use_ws_trade && self.ws.is_active() && !order.is_quote_quantity() {
398            match self.submit_via_ws(command, order, leverage) {
399                Ok(()) => return,
400                Err(e) => log::warn!("Kraken WS submit_order fallback to REST: {e}"),
401            }
402        }
403
404        self.submit_via_rest(order, task_name, leverage);
405    }
406
407    fn submit_via_rest(&self, order: &OrderAny, task_name: &'static str, leverage: Option<u16>) {
408        let account_id = self.core.account_id;
409        let client_order_id = order.client_order_id();
410        let strategy_id = order.strategy_id();
411        let instrument_id = order.instrument_id();
412        let order_side = order.order_side();
413        let order_type = order.order_type();
414        let quantity = order.quantity();
415        let time_in_force = order.time_in_force();
416        let expire_time = order.expire_time();
417        let price = order.price();
418        let trigger_price = order.trigger_price();
419        let trigger_type = order.trigger_type();
420        let trailing_offset = order.trailing_offset();
421        let limit_offset = order.limit_offset();
422        let is_reduce_only = order.is_reduce_only();
423        let is_post_only = order.is_post_only();
424        let is_quote_quantity = order.is_quote_quantity();
425        let display_qty = order.display_qty();
426
427        let http = self.http.clone();
428        let emitter = self.emitter.clone();
429        let clock = self.clock;
430        let dispatch_state = self.ws_dispatch_state.clone();
431        let spot_account_type = self.config.spot_account_type;
432
433        self.spawn_task(task_name, async move {
434            let result = http
435                .submit_order(
436                    account_id,
437                    instrument_id,
438                    client_order_id,
439                    order_side,
440                    order_type,
441                    quantity,
442                    time_in_force,
443                    expire_time,
444                    price,
445                    trigger_price,
446                    trigger_type,
447                    trailing_offset,
448                    limit_offset,
449                    is_reduce_only,
450                    is_post_only,
451                    is_quote_quantity,
452                    display_qty,
453                    leverage,
454                    spot_account_type,
455                )
456                .await;
457
458            match result {
459                Ok(_) => {}
460                Err(e) => match command_failure_from_submit_error(&e) {
461                    CommandFailure::Ambiguous(reason) => {
462                        log::warn!(
463                            "{task_name} outcome is ambiguous for client_order_id={client_order_id}: {reason}"
464                        );
465                    }
466                    CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
467                        let ts_event = clock.get_time_ns();
468                        let error_msg = format!("{task_name} error: {reason}");
469                        let due_post_only = error_msg.contains("POST_ONLY_REJECTED")
470                            || error_msg.contains(KRAKEN_SPOT_POST_ONLY_ERROR);
471                        dispatch_state.cleanup_terminal(&client_order_id);
472                        emitter.emit_order_rejected_event(
473                            strategy_id,
474                            instrument_id,
475                            client_order_id,
476                            &error_msg,
477                            ts_event,
478                            due_post_only,
479                        );
480                    }
481                },
482            }
483
484            Ok(())
485        });
486    }
487    fn submit_via_ws(
488        &self,
489        command: &SubmitOrder,
490        order: &OrderAny,
491        leverage: Option<u16>,
492    ) -> anyhow::Result<()> {
493        let token = self
494            .ws
495            .auth_token_blocking()
496            .ok_or_else(|| anyhow::anyhow!("missing WS auth token"))?;
497
498        let params = build_add_order_params(command, order, token, leverage)?;
499        let identity = PendingRequest {
500            operation: PendingOperation::Submit,
501            client_order_ids: vec![command.client_order_id],
502            venue_order_ids: vec![None],
503            ts_sent_ns: 0,
504            new_quantity: None,
505            new_price: None,
506            new_trigger_price: None,
507        };
508        self.order_request_state
509            .submit(params, identity, self.clock.get_time_ns().as_u64())?;
510        Ok(())
511    }
512
513    fn cancel_single_order(&self, cmd: &CancelOrder) {
514        let use_ws_trade = resolve_use_ws_trade(cmd.params.as_ref(), self.config.use_ws_trade);
515        if use_ws_trade && self.ws.is_active() {
516            match self.cancel_via_ws(cmd) {
517                Ok(()) => return,
518                Err(e) => log::warn!("Kraken WS cancel_order fallback to REST: {e}"),
519            }
520        }
521
522        self.cancel_via_rest(cmd);
523    }
524
525    fn cancel_via_rest(&self, cmd: &CancelOrder) {
526        let account_id = self.core.account_id;
527        let client_order_id = cmd.client_order_id;
528        let venue_order_id = cmd.venue_order_id;
529        let strategy_id = cmd.strategy_id;
530        let instrument_id = cmd.instrument_id;
531
532        log::debug!(
533            "Canceling order: venue_order_id={venue_order_id:?}, client_order_id={client_order_id}"
534        );
535
536        let http = self.http.clone();
537        let emitter = self.emitter.clone();
538        let clock = self.clock;
539
540        self.spawn_task("cancel_order", async move {
541            if let Err(failure) = cancel_order_for_spot(
542                &http,
543                account_id,
544                instrument_id,
545                Some(client_order_id),
546                venue_order_id,
547            )
548            .await
549            {
550                handle_cancel_failure(
551                    &emitter,
552                    clock,
553                    strategy_id,
554                    instrument_id,
555                    client_order_id,
556                    venue_order_id,
557                    failure,
558                );
559            }
560            Ok(())
561        });
562    }
563
564    fn cancel_via_ws(&self, cmd: &CancelOrder) -> anyhow::Result<()> {
565        let token = self
566            .ws
567            .auth_token_blocking()
568            .ok_or_else(|| anyhow::anyhow!("missing WS auth token"))?;
569
570        let params = build_cancel_order_params(cmd, token);
571        let identity = PendingRequest {
572            operation: PendingOperation::Cancel,
573            client_order_ids: vec![cmd.client_order_id],
574            venue_order_ids: vec![cmd.venue_order_id],
575            ts_sent_ns: 0,
576            new_quantity: None,
577            new_price: None,
578            new_trigger_price: None,
579        };
580        self.order_request_state
581            .cancel(params, identity, self.clock.get_time_ns().as_u64())?;
582        Ok(())
583    }
584
585    fn spawn_message_handler(&mut self) -> anyhow::Result<()> {
586        let stream = self.ws.stream().map_err(|e| anyhow::anyhow!("{e}"))?;
587        let emitter = self.emitter.clone();
588        let instruments = self.instruments.clone();
589        let order_qty_cache = self.order_qty_cache.clone();
590        let truncated_id_map = self.truncated_id_map.clone();
591        let dispatch_state = self.ws_dispatch_state.clone();
592        let order_request_state = self.order_request_state.clone();
593        let account_id = self.core.account_id;
594        let clock = self.clock;
595        let cancellation_token = self.cancellation_token.clone();
596
597        let future = async move {
598            tokio::pin!(stream);
599
600            loop {
601                tokio::select! {
602                    () = cancellation_token.cancelled() => {
603                        log::debug!("Spot execution message handler cancelled");
604                        break;
605                    }
606                    msg = stream.next() => {
607                        match msg {
608                            Some(ws_msg) => {
609                                Self::handle_ws_message(
610                                    ws_msg,
611                                    &emitter,
612                                    &dispatch_state,
613                                    &order_request_state,
614                                    &instruments,
615                                    &order_qty_cache,
616                                    &truncated_id_map,
617                                    account_id,
618                                    clock,
619                                );
620                            }
621                            None => {
622                                log::debug!("Spot execution WebSocket stream ended");
623                                break;
624                            }
625                        }
626                    }
627                }
628            }
629        };
630
631        self.session_tasks
632            .spawn(future)
633            .context("failed to register Kraken Spot execution stream task")?;
634
635        let event_rx = self.order_event_rx.take();
636
637        if let Some(mut event_rx) = event_rx {
638            let emitter = self.emitter.clone();
639            let cancellation_token = self.cancellation_token.clone();
640
641            let future = async move {
642                loop {
643                    tokio::select! {
644                        () = cancellation_token.cancelled() => {
645                            log::debug!("Spot execution order-event forwarder cancelled");
646                            break;
647                        }
648                        event = event_rx.recv() => {
649                            match event {
650                                Some(event) => emitter.send_order_event(event),
651                                None => {
652                                    log::debug!("Spot execution order-event channel closed");
653                                    break;
654                                }
655                            }
656                        }
657                    }
658                }
659            };
660            self.session_tasks
661                .spawn(future)
662                .context("failed to register Kraken Spot order event task")?;
663        }
664
665        Ok(())
666    }
667
668    fn modify_single_order(&self, cmd: &ModifyOrder) {
669        let use_ws_trade = resolve_use_ws_trade(cmd.params.as_ref(), self.config.use_ws_trade);
670        if use_ws_trade && self.ws.is_active() {
671            match self.amend_via_ws(cmd) {
672                Ok(()) => return,
673                Err(e) => log::warn!("Kraken WS amend_order fallback to REST: {e}"),
674            }
675        }
676
677        self.amend_via_rest(cmd);
678    }
679
680    fn amend_via_rest(&self, cmd: &ModifyOrder) {
681        let client_order_id = cmd.client_order_id;
682        let venue_order_id = cmd.venue_order_id;
683        let strategy_id = cmd.strategy_id;
684        let instrument_id = cmd.instrument_id;
685        let quantity = cmd.quantity;
686        let price = cmd.price;
687
688        log::debug!(
689            "Modifying order: venue_order_id={venue_order_id:?}, client_order_id={client_order_id}"
690        );
691
692        let http = self.http.clone();
693        let emitter = self.emitter.clone();
694        let clock = self.clock;
695
696        self.spawn_task("modify_order", async move {
697            match http
698                .modify_order(
699                    instrument_id,
700                    Some(client_order_id),
701                    venue_order_id,
702                    quantity,
703                    price,
704                    None,
705                )
706                .await
707            {
708                Ok(_) => {}
709                Err(e) => match command_failure_from_modify_error(&e) {
710                    CommandFailure::Ambiguous(reason) => {
711                        log::warn!(
712                            "modify_order outcome is ambiguous for client_order_id={client_order_id}: {reason}"
713                        );
714                    }
715                    CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
716                        let ts_event = clock.get_time_ns();
717                        emitter.emit_order_modify_rejected_event(
718                            strategy_id,
719                            instrument_id,
720                            client_order_id,
721                            venue_order_id,
722                            &format!("modify-order error: {reason}"),
723                            ts_event,
724                        );
725                    }
726                },
727            }
728            Ok(())
729        });
730    }
731
732    fn amend_via_ws(&self, cmd: &ModifyOrder) -> anyhow::Result<()> {
733        let token = self
734            .ws
735            .auth_token_blocking()
736            .ok_or_else(|| anyhow::anyhow!("missing WS auth token"))?;
737
738        let params = build_amend_order_params(cmd, token);
739        let identity = PendingRequest {
740            operation: PendingOperation::Amend,
741            client_order_ids: vec![cmd.client_order_id],
742            venue_order_ids: vec![cmd.venue_order_id],
743            ts_sent_ns: 0,
744            new_quantity: cmd.quantity,
745            new_price: cmd.price,
746            new_trigger_price: cmd.trigger_price,
747        };
748        self.order_request_state
749            .amend(params, identity, self.clock.get_time_ns().as_u64())?;
750        Ok(())
751    }
752
753    /// Polls the cache until the account is registered or timeout is reached.
754    async fn await_account_registered(&self, timeout_secs: f64) -> anyhow::Result<()> {
755        let account_id = self.core.account_id;
756
757        if self.core.cache().account(&account_id).is_some() {
758            log::info!("Account {account_id} registered");
759            return Ok(());
760        }
761
762        let start = Instant::now();
763        let timeout = Duration::from_secs_f64(timeout_secs);
764        let interval = Duration::from_millis(10);
765
766        loop {
767            tokio::time::sleep(interval).await;
768
769            if self.core.cache().account(&account_id).is_some() {
770                log::info!("Account {account_id} registered");
771                return Ok(());
772            }
773
774            if start.elapsed() >= timeout {
775                anyhow::bail!(
776                    "Timeout waiting for account {account_id} to be registered after {timeout_secs}s"
777                );
778            }
779        }
780    }
781
782    #[expect(clippy::too_many_arguments)]
783    fn handle_ws_message(
784        msg: KrakenSpotWsMessage,
785        emitter: &ExecutionEventEmitter,
786        dispatch_state: &Arc<WsDispatchState>,
787        order_request_state: &Arc<OrderRequestState>,
788        instruments: &Arc<AtomicMap<InstrumentId, InstrumentAny>>,
789        order_qty_cache: &Arc<AtomicMap<String, Decimal>>,
790        truncated_id_map: &Arc<AtomicMap<String, ClientOrderId>>,
791        account_id: AccountId,
792        clock: &'static AtomicTime,
793    ) {
794        match msg {
795            KrakenSpotWsMessage::Execution(executions) => {
796                let ts_init = clock.get_time_ns();
797
798                for exec in &executions {
799                    dispatch::spot::execution(
800                        exec,
801                        dispatch_state,
802                        emitter,
803                        instruments,
804                        truncated_id_map,
805                        order_qty_cache,
806                        account_id,
807                        ts_init,
808                    );
809                }
810            }
811            KrakenSpotWsMessage::OrderResponse(response) => {
812                let ts_event = clock.get_time_ns().as_u64();
813                order_request_state.handle_response(&response, ts_event);
814            }
815            KrakenSpotWsMessage::Reconnected => {
816                log::info!("Spot execution WebSocket reconnected");
817            }
818            KrakenSpotWsMessage::Ticker(_)
819            | KrakenSpotWsMessage::Trade(_)
820            | KrakenSpotWsMessage::Book { .. }
821            | KrakenSpotWsMessage::Ohlc(_)
822            | KrakenSpotWsMessage::L3Snapshot(_)
823            | KrakenSpotWsMessage::L3Update(_) => {}
824        }
825    }
826
827    fn sweep_stale_margin_positions(
828        &self,
829        account_id: AccountId,
830        reports: &mut Vec<PositionStatusReport>,
831    ) {
832        let reported: HashSet<InstrumentId> = reports
833            .iter()
834            .filter(|r| r.position_side != PositionSide::Flat)
835            .map(|r| r.instrument_id)
836            .collect();
837
838        let ts_now = self.clock.get_time_ns();
839        let cache = self.core.cache();
840        let open_positions =
841            cache.positions_open(Some(&*KRAKEN_VENUE), None, None, Some(&account_id), None);
842
843        for pos in open_positions {
844            let inst_id = pos.instrument_id;
845
846            if product_type_from_symbol(inst_id.symbol.inner().as_str()) != KrakenProductType::Spot
847            {
848                continue;
849            }
850
851            if reported.contains(&inst_id) {
852                continue;
853            }
854
855            let precision = cache.instrument(&inst_id).map_or(0, |i| i.size_precision());
856            log::debug!("Emitting synthetic FLAT for closed margin position {inst_id}");
857            reports.push(PositionStatusReport::new(
858                account_id,
859                inst_id,
860                PositionSide::Flat,
861                Quantity::zero(precision),
862                ts_now,
863                ts_now,
864                None,
865                None,
866                None,
867            ));
868        }
869    }
870
871    fn batch_add_via_rest(
872        &self,
873        order_tuples: Vec<BatchOrderTuple>,
874        order_meta: Vec<(StrategyId, InstrumentId, ClientOrderId)>,
875    ) {
876        let http = self.http.clone();
877        let emitter = self.emitter.clone();
878        let clock = self.clock;
879        let dispatch_state = self.ws_dispatch_state.clone();
880        let spot_account_type = self.config.spot_account_type;
881
882        self.spawn_task("submit_order_list", async move {
883            let results = http
884                .send_order_batches(order_tuples, spot_account_type)
885                .await;
886
887            for (result, (strategy_id, instrument_id, client_order_id)) in
888                results.into_iter().zip(&order_meta)
889            {
890                let outcome = match result {
891                    Ok(item) => command_failure_from_spot_batch_item(item),
892                    Err(e) => Err(command_failure_from_spot_batch_error(&e)),
893                };
894
895                match outcome {
896                    Ok(()) => {}
897                    Err(CommandFailure::Ambiguous(reason)) => {
898                        log::warn!(
899                            "submit_order_list outcome is ambiguous for client_order_id={client_order_id}: {reason}"
900                        );
901                    }
902                    Err(
903                        CommandFailure::NotSent(reason)
904                        | CommandFailure::VenueRejected(reason),
905                    ) => {
906                        let ts_event = clock.get_time_ns();
907                        let due_post_only = reason.contains("POST_ONLY_REJECTED")
908                            || reason.contains(KRAKEN_SPOT_POST_ONLY_ERROR);
909                        dispatch_state.cleanup_terminal(client_order_id);
910                        emitter.emit_order_rejected_event(
911                            *strategy_id,
912                            *instrument_id,
913                            *client_order_id,
914                            &format!("submit_order_list batch item rejected: {reason}"),
915                            ts_event,
916                            due_post_only,
917                        );
918                    }
919                }
920            }
921            Ok(())
922        });
923    }
924
925    fn batch_add_via_ws(&self, orders: &[OrderAny], leverage: Option<u16>) -> anyhow::Result<()> {
926        let token = self
927            .ws
928            .auth_token_blocking()
929            .ok_or_else(|| anyhow::anyhow!("missing WS auth token"))?;
930
931        let first = orders
932            .first()
933            .ok_or_else(|| anyhow::anyhow!("batch_add requires at least one order"))?;
934        let symbol = first.instrument_id().symbol.inner().to_string();
935
936        let mut batch_orders = Vec::with_capacity(orders.len());
937        let mut client_order_ids = Vec::with_capacity(orders.len());
938        for order in orders {
939            batch_orders.push(build_batch_order(order, leverage)?);
940            client_order_ids.push(order.client_order_id());
941        }
942        let venue_order_ids = vec![None; orders.len()];
943
944        let params = KrakenWsBatchAddParams {
945            symbol,
946            orders: batch_orders,
947            token,
948        };
949        let identity = PendingRequest {
950            operation: PendingOperation::BatchAdd,
951            client_order_ids,
952            venue_order_ids,
953            ts_sent_ns: 0,
954            new_quantity: None,
955            new_price: None,
956            new_trigger_price: None,
957        };
958        self.order_request_state
959            .batch_add(params, identity, self.clock.get_time_ns().as_u64())?;
960        Ok(())
961    }
962}
963
964type BatchOrderTuple = (
965    InstrumentId,
966    ClientOrderId,
967    OrderSide,
968    OrderType,
969    Quantity,
970    TimeInForce,
971    Option<UnixNanos>,
972    Option<Price>,
973    Option<Price>,
974    Option<TriggerType>,
975    Option<Decimal>,
976    Option<Decimal>,
977    bool,
978    bool,
979    bool,
980    Option<Quantity>,
981    Option<u16>,
982);
983
984fn build_batch_order(
985    order: &OrderAny,
986    leverage: Option<u16>,
987) -> anyhow::Result<KrakenWsBatchAddOrder> {
988    let order_type = order.order_type();
989    let side = match order.order_side() {
990        OrderSide::Buy => KrakenOrderSide::Buy,
991        OrderSide::Sell => KrakenOrderSide::Sell,
992    };
993
994    if matches!(
995        order_type,
996        OrderType::TrailingStopMarket | OrderType::TrailingStopLimit
997    ) {
998        anyhow::bail!(
999            "Trailing stop orders are not yet supported on the Kraken WS batch path; use REST",
1000        );
1001    }
1002
1003    if order.display_qty().is_some() {
1004        anyhow::bail!(
1005            "Iceberg (display_qty) orders are not supported on the Kraken WS batch path; use REST",
1006        );
1007    }
1008
1009    let kraken_order_type = match order_type {
1010        OrderType::Market => KrakenOrderType::Market,
1011        OrderType::Limit => KrakenOrderType::Limit,
1012        OrderType::StopMarket => KrakenOrderType::StopLoss,
1013        OrderType::StopLimit => KrakenOrderType::StopLossLimit,
1014        OrderType::MarketIfTouched => KrakenOrderType::TakeProfit,
1015        OrderType::LimitIfTouched => KrakenOrderType::TakeProfitLimit,
1016        ty => anyhow::bail!("Unsupported order type for Kraken WS batch: {ty:?}"),
1017    };
1018
1019    let is_limit_order = matches!(
1020        order_type,
1021        OrderType::Limit | OrderType::StopLimit | OrderType::LimitIfTouched
1022    );
1023
1024    if is_limit_order && order.price().is_none() {
1025        anyhow::bail!("limit_price is required for batch order type {order_type:?}");
1026    }
1027
1028    let ws_tif =
1029        compute_ws_time_in_force(is_limit_order, order.time_in_force(), order.expire_time())?;
1030    let expire_time = match (ws_tif, order.expire_time()) {
1031        (Some(KrakenTimeInForce::GoodTilDate), Some(ts)) => Some(format_expire_time(ts)),
1032        _ => None,
1033    };
1034
1035    let is_conditional = matches!(
1036        order_type,
1037        OrderType::StopMarket
1038            | OrderType::StopLimit
1039            | OrderType::MarketIfTouched
1040            | OrderType::LimitIfTouched
1041    );
1042
1043    let trigger = if is_conditional {
1044        let trigger_ref = match order.trigger_type() {
1045            Some(TriggerType::IndexPrice) => KrakenSpotTrigger::Index,
1046            Some(TriggerType::LastPrice | TriggerType::Default) | None => KrakenSpotTrigger::Last,
1047            Some(other) => anyhow::bail!(
1048                "Unsupported trigger type for Kraken Spot WS batch: {other:?} (only LastPrice and IndexPrice supported)",
1049            ),
1050        };
1051        order.trigger_price().map(|tp| KrakenWsTriggerParams {
1052            reference: trigger_ref,
1053            price: tp.as_decimal(),
1054            price_type: None,
1055        })
1056    } else {
1057        None
1058    };
1059
1060    if is_conditional && trigger.is_none() {
1061        anyhow::bail!(
1062            "Conditional order type {order_type:?} requires trigger_price for Kraken WS batch",
1063        );
1064    }
1065
1066    Ok(KrakenWsBatchAddOrder {
1067        order_type: kraken_order_type,
1068        side,
1069        order_qty: order.quantity().as_decimal(),
1070        limit_price: order.price().map(|p| p.as_decimal()),
1071        cl_ord_id: Some(truncate_cl_ord_id(&order.client_order_id())),
1072        time_in_force: ws_tif,
1073        expire_time,
1074        post_only: order.is_post_only().then_some(true),
1075        reduce_only: order.is_reduce_only().then_some(true),
1076        leverage,
1077        trigger,
1078    })
1079}
1080
1081#[async_trait(?Send)]
1082impl ExecutionClient for KrakenSpotExecutionClient {
1083    fn is_connected(&self) -> bool {
1084        self.core.is_connected()
1085    }
1086
1087    fn client_id(&self) -> ClientId {
1088        self.core.client_id
1089    }
1090
1091    fn account_id(&self) -> AccountId {
1092        self.core.account_id
1093    }
1094
1095    fn venue(&self) -> Venue {
1096        *KRAKEN_VENUE
1097    }
1098
1099    fn oms_type(&self) -> OmsType {
1100        self.core.oms_type
1101    }
1102
1103    fn get_account(&self) -> Option<AccountAny> {
1104        self.core.cache().account_owned(&self.core.account_id)
1105    }
1106
1107    fn generate_account_state(
1108        &self,
1109        balances: Vec<AccountBalance>,
1110        margins: Vec<MarginBalance>,
1111        reported: bool,
1112        ts_event: UnixNanos,
1113        info: Option<Params>,
1114    ) -> anyhow::Result<()> {
1115        self.emitter
1116            .emit_account_state(balances, margins, reported, ts_event, info);
1117        Ok(())
1118    }
1119
1120    fn start(&mut self) -> anyhow::Result<()> {
1121        if self.core.is_started() {
1122            return Ok(());
1123        }
1124
1125        self.emitter.set_sender(get_exec_event_sender());
1126        self.core.set_started();
1127
1128        log::info!(
1129            "Started: client_id={}, account_id={}, product_type=Spot, environment={:?}",
1130            self.core.client_id,
1131            self.core.account_id,
1132            self.config.environment
1133        );
1134        Ok(())
1135    }
1136
1137    fn stop(&mut self) -> anyhow::Result<()> {
1138        if self.core.is_stopped() {
1139            return Ok(());
1140        }
1141
1142        self.http.cancel_all_requests();
1143        self.session_tasks.begin_shutdown();
1144        self.pending_tasks.begin_shutdown();
1145        self.ws.begin_shutdown();
1146        self.core.set_stopped();
1147        self.core.set_disconnected();
1148        log::info!("Stopped: client_id={}", self.core.client_id);
1149        Ok(())
1150    }
1151
1152    async fn connect(&mut self) -> anyhow::Result<()> {
1153        if self.core.is_connected() && self.session_tasks.is_open() && self.pending_tasks.is_open()
1154        {
1155            return Ok(());
1156        }
1157
1158        self.http.reset_cancellation_token();
1159        self.prepare_task_groups().await?;
1160
1161        if !self.core.instruments_initialized() {
1162            let instruments = self
1163                .http
1164                .request_instruments(None)
1165                .await
1166                .context("Failed to load Kraken spot instruments")?;
1167            log::debug!("Loaded {} Spot instruments", instruments.len());
1168            self.http.cache_instruments(&instruments);
1169            self.core.set_instruments_initialized();
1170        }
1171
1172        let session_result = async {
1173            self.ws
1174                .connect()
1175                .await
1176                .context("Failed to connect spot WebSocket")?;
1177            self.ws
1178                .wait_until_active(10.0)
1179                .await
1180                .context("Spot WebSocket failed to become active")?;
1181
1182            self.ws
1183                .authenticate()
1184                .await
1185                .context("Failed to authenticate spot WebSocket")?;
1186
1187            // Request initial account state and await registration before spawning
1188            // the message handler. Report events from execution snapshots conflict
1189            // with ExecEngine borrows during startup, so account registration must
1190            // complete first.
1191            let account_state = self
1192                .http
1193                .request_account_state(
1194                    self.core.account_id,
1195                    self.config.spot_account_type,
1196                    self.config.margin_balance_asset.as_deref(),
1197                )
1198                .await
1199                .context("Failed to request Kraken account state")?;
1200
1201            if !account_state.balances.is_empty() {
1202                log::debug!(
1203                    "Received account state with {} balance(s)",
1204                    account_state.balances.len()
1205                );
1206            }
1207
1208            self.emitter.send_account_state(account_state);
1209            self.await_account_registered(30.0).await?;
1210
1211            self.spawn_message_handler()?;
1212
1213            self.instruments.rcu(|m| {
1214                for instrument in self.http.instruments_cache.load().values() {
1215                    m.insert(instrument.id(), instrument.clone());
1216                }
1217            });
1218
1219            self.ws
1220                .subscribe_executions(false, false)
1221                .await
1222                .context("Failed to subscribe to executions")?;
1223
1224            log::debug!("Spot WebSocket authenticated and subscribed to executions");
1225
1226            Ok::<(), anyhow::Error>(())
1227        }
1228        .await;
1229
1230        if let Err(e) = session_result {
1231            if let Err(teardown_error) = self.teardown_partial_connect().await {
1232                return Err(e.context(format!(
1233                    "Kraken Spot execution startup teardown failed: {teardown_error}"
1234                )));
1235            }
1236            return Err(e);
1237        }
1238
1239        self.core.set_connected();
1240        log::info!("Connected: client_id={}", self.core.client_id);
1241        Ok(())
1242    }
1243
1244    async fn disconnect(&mut self) -> anyhow::Result<()> {
1245        self.teardown_partial_connect().await?;
1246        log::info!("Disconnected: client_id={}", self.core.client_id);
1247        Ok(())
1248    }
1249
1250    async fn generate_order_status_report(
1251        &self,
1252        cmd: &GenerateOrderStatusReport,
1253    ) -> anyhow::Result<Option<OrderStatusReport>> {
1254        log::debug!(
1255            "Generating order status report: venue_order_id={:?}, client_order_id={:?}",
1256            cmd.venue_order_id,
1257            cmd.client_order_id
1258        );
1259
1260        let account_id = self.core.account_id;
1261        let reports = self
1262            .http
1263            .request_order_status_reports(account_id, None, None, None, false)
1264            .await?;
1265
1266        // Match by venue_order_id or client_order_id (comparing truncated form
1267        // since Kraken stores the truncated cl_ord_id for long IDs)
1268        Ok(reports.into_iter().find(|r| {
1269            cmd.venue_order_id
1270                .is_some_and(|id| r.venue_order_id.as_str() == id.as_str())
1271                || cmd.client_order_id.is_some_and(|id| {
1272                    r.client_order_id
1273                        .as_ref()
1274                        .is_some_and(|r_id| r_id.as_str() == truncate_cl_ord_id(&id))
1275                })
1276        }))
1277    }
1278
1279    async fn generate_order_status_reports(
1280        &self,
1281        cmd: &GenerateOrderStatusReports,
1282    ) -> anyhow::Result<Vec<OrderStatusReport>> {
1283        log::debug!(
1284            "Generating order status reports: instrument_id={:?}, open_only={}",
1285            cmd.instrument_id,
1286            cmd.open_only
1287        );
1288
1289        let account_id = self.core.account_id;
1290        let start = cmd.start.map(Timestamp::from);
1291        let end = cmd.end.map(Timestamp::from);
1292        self.http
1293            .request_order_status_reports(account_id, cmd.instrument_id, start, end, cmd.open_only)
1294            .await
1295    }
1296
1297    async fn generate_fill_reports(
1298        &self,
1299        cmd: GenerateFillReports,
1300    ) -> anyhow::Result<Vec<FillReport>> {
1301        log::debug!(
1302            "Generating fill reports: instrument_id={:?}",
1303            cmd.instrument_id
1304        );
1305
1306        let account_id = self.core.account_id;
1307        let start = cmd.start.map(Timestamp::from);
1308        let end = cmd.end.map(Timestamp::from);
1309        self.http
1310            .request_fill_reports(account_id, cmd.instrument_id, start, end)
1311            .await
1312    }
1313
1314    async fn generate_position_status_reports(
1315        &self,
1316        cmd: &GeneratePositionStatusReports,
1317    ) -> anyhow::Result<Vec<PositionStatusReport>> {
1318        log::debug!(
1319            "Generating position status reports: instrument_id={:?}",
1320            cmd.instrument_id
1321        );
1322
1323        let account_id = self.core.account_id;
1324        let mut reports = self
1325            .http
1326            .request_position_status_reports(
1327                account_id,
1328                cmd.instrument_id,
1329                self.config.spot_account_type,
1330                self.config.use_spot_position_reports,
1331                Ustr::from(self.config.spot_positions_quote_currency.as_str()),
1332            )
1333            .await?;
1334
1335        if cmd.instrument_id.is_none() && self.config.spot_account_type == AccountType::Margin {
1336            self.sweep_stale_margin_positions(account_id, &mut reports);
1337        }
1338
1339        Ok(reports)
1340    }
1341
1342    async fn generate_mass_status(
1343        &self,
1344        lookback_mins: Option<u64>,
1345    ) -> anyhow::Result<Option<ExecutionMassStatus>> {
1346        log::debug!("Generating mass status: lookback_mins={lookback_mins:?}");
1347
1348        let start = lookback_mins.map(|mins| Timestamp::now() - Duration::from_secs(mins * 60));
1349
1350        let account_id = self.core.account_id;
1351        let order_reports = self
1352            .http
1353            .request_order_status_reports(account_id, None, start, None, true)
1354            .await?;
1355        let fill_reports = self
1356            .http
1357            .request_fill_reports(account_id, None, start, None)
1358            .await?;
1359        let mut position_reports = self
1360            .http
1361            .request_position_status_reports(
1362                account_id,
1363                None,
1364                self.config.spot_account_type,
1365                self.config.use_spot_position_reports,
1366                Ustr::from(self.config.spot_positions_quote_currency.as_str()),
1367            )
1368            .await?;
1369
1370        if self.config.spot_account_type == AccountType::Margin {
1371            self.sweep_stale_margin_positions(account_id, &mut position_reports);
1372        }
1373
1374        let mut mass_status = ExecutionMassStatus::new(
1375            self.core.client_id,
1376            self.core.account_id,
1377            *KRAKEN_VENUE,
1378            self.clock.get_time_ns(),
1379            None,
1380        );
1381        mass_status.add_order_reports(order_reports);
1382        mass_status.add_fill_reports(fill_reports);
1383        mass_status.add_position_reports(position_reports);
1384
1385        Ok(Some(mass_status))
1386    }
1387
1388    fn query_account(&self, cmd: QueryAccount) -> anyhow::Result<()> {
1389        log::debug!("Querying account: {cmd:?}");
1390
1391        let account_id = self.core.account_id;
1392        let http = self.http.clone();
1393        let emitter = self.emitter.clone();
1394
1395        let spot_account_type = self.config.spot_account_type;
1396        let margin_balance_asset = self.config.margin_balance_asset.clone();
1397        self.spawn_task("query_account", async move {
1398            let account_state = http
1399                .request_account_state(
1400                    account_id,
1401                    spot_account_type,
1402                    margin_balance_asset.as_deref(),
1403                )
1404                .await?;
1405            emitter.emit_account_state(
1406                account_state.balances.clone(),
1407                account_state.margins.clone(),
1408                account_state.is_reported,
1409                account_state.ts_event,
1410                account_state.info,
1411            );
1412            Ok(())
1413        });
1414
1415        Ok(())
1416    }
1417
1418    fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
1419        log::debug!("Querying order: {cmd:?}");
1420
1421        let venue_order_id = cmd
1422            .venue_order_id
1423            .context("venue_order_id required for query_order")?;
1424        let account_id = self.core.account_id;
1425        let http = self.http.clone();
1426        let emitter = self.emitter.clone();
1427
1428        self.spawn_task("query_order", async move {
1429            let reports = http
1430                .request_order_status_reports(account_id, None, None, None, true)
1431                .await
1432                .context("Failed to query order")?;
1433
1434            if let Some(report) = reports
1435                .into_iter()
1436                .find(|r| r.venue_order_id == venue_order_id)
1437            {
1438                emitter.send_order_status_report(report);
1439            }
1440            Ok(())
1441        });
1442
1443        Ok(())
1444    }
1445
1446    fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
1447        let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
1448        let leverage = match resolve_leverage(cmd.params.as_ref(), self.config.default_leverage) {
1449            Ok(lev) => lev,
1450            Err(reason) => {
1451                self.emitter.emit_order_denied(&order, &reason);
1452                return Ok(());
1453            }
1454        };
1455        self.submit_single_order(&cmd, &order, "submit_order", leverage);
1456        Ok(())
1457    }
1458
1459    fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
1460        let orders = self.core.get_orders_for_list(&cmd.order_list)?;
1461
1462        log::debug!(
1463            "Submitting order list: order_list_id={}, count={}",
1464            cmd.order_list.id,
1465            orders.len()
1466        );
1467
1468        let leverage = match resolve_leverage(cmd.params.as_ref(), self.config.default_leverage) {
1469            Ok(lev) => lev,
1470            Err(reason) => {
1471                for order in &orders {
1472                    self.emitter.emit_order_denied(order, &reason);
1473                }
1474                return Ok(());
1475            }
1476        };
1477
1478        let mut order_tuples = Vec::with_capacity(orders.len());
1479        let mut order_meta = Vec::with_capacity(orders.len());
1480        let mut prepared_orders = Vec::with_capacity(orders.len());
1481
1482        for order in &orders {
1483            if order.is_closed() {
1484                log::warn!(
1485                    "Cannot submit closed order: client_order_id={}",
1486                    order.client_order_id()
1487                );
1488                continue;
1489            }
1490
1491            if order.time_in_force() == TimeInForce::Fok && order.order_type() != OrderType::Limit {
1492                self.emitter.emit_order_denied(
1493                    order,
1494                    "FOK time in force only supported for LIMIT orders on Kraken Spot",
1495                );
1496                continue;
1497            }
1498
1499            if matches!(
1500                order.order_type(),
1501                OrderType::TrailingStopMarket | OrderType::TrailingStopLimit
1502            ) && let Some(offset_type) = order.trailing_offset_type()
1503                && offset_type != TrailingOffsetType::Price
1504            {
1505                self.emitter.emit_order_denied(
1506                    order,
1507                    &format!(
1508                        "Kraken Spot only supports Price trailing offset type: received {offset_type:?}"
1509                    ),
1510                );
1511                continue;
1512            }
1513
1514            if order.is_reduce_only() && self.config.spot_account_type == AccountType::Cash {
1515                self.emitter
1516                    .emit_order_denied(order, "reduce_only requires spot_account_type=Margin");
1517                continue;
1518            }
1519
1520            let client_order_id = order.client_order_id();
1521            let kraken_cl_ord_id = truncate_cl_ord_id(&client_order_id);
1522
1523            self.register_order_identity(order);
1524            self.emitter.emit_order_submitted(order);
1525
1526            if !order.is_quote_quantity() {
1527                self.order_qty_cache
1528                    .insert(kraken_cl_ord_id.clone(), order.quantity().as_decimal());
1529            }
1530
1531            if kraken_cl_ord_id != client_order_id.as_str() {
1532                self.truncated_id_map
1533                    .insert(kraken_cl_ord_id, client_order_id);
1534            }
1535            order_tuples.push((
1536                order.instrument_id(),
1537                client_order_id,
1538                order.order_side(),
1539                order.order_type(),
1540                order.quantity(),
1541                order.time_in_force(),
1542                order.expire_time(),
1543                order.price(),
1544                order.trigger_price(),
1545                order.trigger_type(),
1546                order.trailing_offset(),
1547                order.limit_offset(),
1548                order.is_reduce_only(),
1549                order.is_post_only(),
1550                order.is_quote_quantity(),
1551                order.display_qty(),
1552                leverage,
1553            ));
1554
1555            order_meta.push((order.strategy_id(), order.instrument_id(), client_order_id));
1556            prepared_orders.push(order.clone());
1557        }
1558
1559        if order_tuples.is_empty() {
1560            return Ok(());
1561        }
1562
1563        let use_ws_trade = resolve_use_ws_trade(cmd.params.as_ref(), self.config.use_ws_trade);
1564        if use_ws_trade && self.ws.is_active() {
1565            let any_quote_qty = prepared_orders.iter().any(|o| o.is_quote_quantity());
1566            let symbols_match = prepared_orders
1567                .windows(2)
1568                .all(|w| w[0].instrument_id() == w[1].instrument_id());
1569
1570            if any_quote_qty {
1571                log::warn!(
1572                    "Kraken WS batch_add does not support quote-quantity orders, falling back to REST for order_list_id={}",
1573                    cmd.order_list.id,
1574                );
1575            } else if symbols_match {
1576                match self.batch_add_via_ws(&prepared_orders, leverage) {
1577                    Ok(()) => return Ok(()),
1578                    Err(e) => log::warn!("Kraken WS batch_add fallback to REST: {e}"),
1579                }
1580            } else {
1581                log::warn!(
1582                    "Kraken WS batch_add requires single shared symbol, falling back to REST for order_list_id={}",
1583                    cmd.order_list.id,
1584                );
1585            }
1586        }
1587
1588        self.batch_add_via_rest(order_tuples, order_meta);
1589
1590        Ok(())
1591    }
1592
1593    fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
1594        self.modify_single_order(&cmd);
1595        Ok(())
1596    }
1597
1598    fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
1599        self.cancel_single_order(&cmd);
1600        Ok(())
1601    }
1602
1603    fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
1604        let instrument_id = cmd.instrument_id;
1605
1606        if cmd.order_side.is_none() {
1607            log::debug!("Canceling all orders: instrument_id={instrument_id} (bulk)");
1608
1609            let http = self.http.clone();
1610
1611            self.spawn_task("cancel_all_orders", async move {
1612                if let Err(e) = http.inner.cancel_all_orders().await {
1613                    match command_failure_from_cancel_error(e) {
1614                        CommandFailure::NotSent(reason) => {
1615                            log::warn!("Cancel-all failed local validation: {reason}");
1616                        }
1617                        CommandFailure::Ambiguous(reason)
1618                        | CommandFailure::VenueRejected(reason) => {
1619                            log::warn!(
1620                                "Cancel-all ambiguous failure, awaiting reconciliation: {reason}"
1621                            );
1622                        }
1623                    }
1624                }
1625                Ok(())
1626            });
1627
1628            return Ok(());
1629        }
1630
1631        log::debug!(
1632            "Canceling all orders: instrument_id={instrument_id}, side={:?}",
1633            cmd.order_side
1634        );
1635
1636        let orders_to_cancel: Vec<_> = {
1637            let cache = self.core.cache();
1638            let open_orders = cache.orders_open(None, Some(&instrument_id), None, None, None);
1639
1640            open_orders
1641                .into_iter()
1642                .filter(|order| Some(order.order_side()) == cmd.order_side)
1643                .filter_map(|order| {
1644                    Some((
1645                        order.venue_order_id()?,
1646                        order.client_order_id(),
1647                        order.instrument_id(),
1648                        order.strategy_id(),
1649                    ))
1650                })
1651                .collect()
1652        };
1653
1654        let account_id = self.core.account_id;
1655
1656        for (venue_order_id, client_order_id, order_instrument_id, strategy_id) in orders_to_cancel
1657        {
1658            let http = self.http.clone();
1659            let emitter = self.emitter.clone();
1660            let clock = self.clock;
1661
1662            self.spawn_task("cancel_order_by_side", async move {
1663                if let Err(failure) = cancel_order_for_spot(
1664                    &http,
1665                    account_id,
1666                    order_instrument_id,
1667                    Some(client_order_id),
1668                    Some(venue_order_id),
1669                )
1670                .await
1671                {
1672                    handle_cancel_failure(
1673                        &emitter,
1674                        clock,
1675                        strategy_id,
1676                        order_instrument_id,
1677                        client_order_id,
1678                        Some(venue_order_id),
1679                        failure,
1680                    );
1681                }
1682                Ok(())
1683            });
1684        }
1685
1686        Ok(())
1687    }
1688
1689    fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
1690        log::debug!(
1691            "Batch canceling orders: instrument_id={}, count={}",
1692            cmd.instrument_id,
1693            cmd.cancels.len()
1694        );
1695
1696        let use_ws_trade = resolve_use_ws_trade(cmd.params.as_ref(), self.config.use_ws_trade);
1697        if use_ws_trade && self.ws.is_active() {
1698            for cancel in &cmd.cancels {
1699                self.cancel_single_order(cancel);
1700            }
1701
1702            return Ok(());
1703        }
1704
1705        let http = self.http.clone();
1706        let cancels = cmd.cancels;
1707
1708        self.spawn_task("batch_cancel_orders", async move {
1709            batch_cancel_orders_for_spot(&http, &cancels).await;
1710            Ok(())
1711        });
1712
1713        Ok(())
1714    }
1715}
1716
1717async fn cancel_order_for_spot(
1718    http: &KrakenSpotHttpClient,
1719    _account_id: AccountId,
1720    instrument_id: InstrumentId,
1721    client_order_id: Option<ClientOrderId>,
1722    venue_order_id: Option<VenueOrderId>,
1723) -> Result<(), CommandFailure> {
1724    http.get_cached_instrument(&instrument_id.symbol.inner())
1725        .ok_or_else(|| {
1726            CommandFailure::not_sent(InstrumentLookupError::not_found(instrument_id).to_string())
1727        })?;
1728
1729    let txid = venue_order_id.as_ref().map(ToString::to_string);
1730    let cl_ord_id = client_order_id.as_ref().map(truncate_cl_ord_id);
1731
1732    if txid.is_none() && cl_ord_id.is_none() {
1733        return Err(CommandFailure::not_sent(
1734            "Either client_order_id or venue_order_id must be provided",
1735        ));
1736    }
1737
1738    let mut builder = KrakenSpotCancelOrderParamsBuilder::default();
1739
1740    if let Some(ref id) = txid {
1741        builder.txid(id.clone());
1742    } else if let Some(ref id) = cl_ord_id {
1743        builder.cl_ord_id(id.clone());
1744    }
1745
1746    let params = builder
1747        .build()
1748        .map_err(|e| CommandFailure::not_sent(format!("Failed to build cancel params: {e}")))?;
1749
1750    http.inner
1751        .cancel_order(&params)
1752        .await
1753        .map_err(command_failure_from_spot_cancel_error)?;
1754
1755    Ok(())
1756}
1757
1758async fn batch_cancel_orders_for_spot(http: &KrakenSpotHttpClient, cancels: &[CancelOrder]) {
1759    let mut orders = Vec::new();
1760
1761    for cancel in cancels {
1762        match batch_cancel_item_for_spot(http, cancel) {
1763            Ok(order) => orders.push(order),
1764            Err(CommandFailure::NotSent(reason)) => {
1765                log::warn!(
1766                    "Batch cancel command failed local validation for {}: {reason}",
1767                    cancel.client_order_id
1768                );
1769            }
1770            Err(CommandFailure::Ambiguous(reason) | CommandFailure::VenueRejected(reason)) => {
1771                log::warn!(
1772                    "Batch cancel command ambiguous failure for {}, awaiting reconciliation: {reason}",
1773                    cancel.client_order_id
1774                );
1775            }
1776        }
1777    }
1778
1779    for chunk in orders.chunks(50) {
1780        let params = KrakenSpotCancelOrderBatchParams {
1781            orders: chunk.to_vec(),
1782        };
1783
1784        match http.inner.cancel_order_batch(&params).await {
1785            Ok(response) => {
1786                if response.count < chunk.len() as i32 {
1787                    log::warn!(
1788                        "Batch cancel accepted {} of {} request(s) without per-order results; awaiting reconciliation",
1789                        response.count,
1790                        chunk.len()
1791                    );
1792                }
1793            }
1794            Err(e) => match command_failure_from_cancel_error(e) {
1795                CommandFailure::NotSent(reason) => {
1796                    log::warn!("Batch cancel failed local validation: {reason}");
1797                }
1798                CommandFailure::Ambiguous(reason) | CommandFailure::VenueRejected(reason) => {
1799                    log::warn!(
1800                        "Batch cancel failed without per-order results, awaiting reconciliation: {reason}"
1801                    );
1802                }
1803            },
1804        }
1805    }
1806}
1807
1808fn batch_cancel_item_for_spot(
1809    http: &KrakenSpotHttpClient,
1810    cancel: &CancelOrder,
1811) -> Result<String, CommandFailure> {
1812    http.get_cached_instrument(&cancel.instrument_id.symbol.inner())
1813        .ok_or_else(|| {
1814            CommandFailure::not_sent(
1815                InstrumentLookupError::not_found(cancel.instrument_id).to_string(),
1816            )
1817        })?;
1818
1819    if let Some(venue_order_id) = cancel.venue_order_id {
1820        Ok(venue_order_id.to_string())
1821    } else {
1822        Ok(truncate_cl_ord_id(&cancel.client_order_id))
1823    }
1824}
1825
1826fn handle_cancel_failure(
1827    emitter: &ExecutionEventEmitter,
1828    clock: &'static AtomicTime,
1829    strategy_id: StrategyId,
1830    instrument_id: InstrumentId,
1831    client_order_id: ClientOrderId,
1832    venue_order_id: Option<VenueOrderId>,
1833    failure: CommandFailure,
1834) {
1835    match failure {
1836        CommandFailure::VenueRejected(reason) => {
1837            emitter.emit_order_cancel_rejected_event(
1838                strategy_id,
1839                instrument_id,
1840                client_order_id,
1841                venue_order_id,
1842                &reason,
1843                clock.get_time_ns(),
1844            );
1845        }
1846        CommandFailure::NotSent(reason) => {
1847            log::warn!("Cancel command failed local validation for {client_order_id}: {reason}");
1848        }
1849        CommandFailure::Ambiguous(reason) => {
1850            log::warn!(
1851                "Ambiguous cancel failure for {client_order_id}, awaiting reconciliation: {reason}"
1852            );
1853        }
1854    }
1855}
1856
1857fn resolve_leverage(params: Option<&Params>, default: Option<u16>) -> Result<Option<u16>, String> {
1858    let Some(p) = params else {
1859        return Ok(default);
1860    };
1861    let Some(raw) = p.get("leverage") else {
1862        return Ok(default);
1863    };
1864    let n = raw.as_u64().ok_or_else(|| {
1865        format!("Invalid leverage param: expected unsigned integer, received {raw}")
1866    })?;
1867    let lev =
1868        u16::try_from(n).map_err(|_| format!("leverage {n} exceeds maximum ({})", u16::MAX))?;
1869    Ok(Some(lev))
1870}
1871
1872/// Resolves the per-call `params["use_ws_trade"]` override against the
1873/// configured default. Non-boolean values warn and fall back to the default.
1874fn resolve_use_ws_trade(params: Option<&Params>, default: bool) -> bool {
1875    let Some(p) = params else {
1876        return default;
1877    };
1878    let Some(raw) = p.get("use_ws_trade") else {
1879        return default;
1880    };
1881
1882    match raw.as_bool() {
1883        Some(b) => b,
1884        None => {
1885            log::warn!(
1886                "Invalid use_ws_trade param: expected boolean, received {raw}; using default {default}",
1887            );
1888            default
1889        }
1890    }
1891}
1892
1893#[cfg(test)]
1894mod tests {
1895    use std::{cell::RefCell, rc::Rc};
1896
1897    use nautilus_common::{
1898        cache::{Cache, InstrumentLookupError},
1899        clock::TestClock,
1900        factories::ExecutionClientFactory,
1901        messages::execution::CancelOrder,
1902    };
1903    use nautilus_core::{Params, UUID4, UnixNanos};
1904    use nautilus_model::identifiers::{
1905        AccountId, ClientOrderId, InstrumentId, StrategyId, TraderId, VenueOrderId,
1906    };
1907    use rstest::rstest;
1908    use serde_json::json;
1909
1910    use super::{
1911        CommandFailure, batch_cancel_item_for_spot, cancel_order_for_spot, resolve_leverage,
1912        resolve_use_ws_trade,
1913    };
1914    use crate::{
1915        common::enums::KrakenProductType, config::KrakenExecutionClientConfig,
1916        factories::KrakenExecutionClientFactory, http::KrakenSpotHttpClient,
1917    };
1918
1919    const TEST_INSTRUMENT_ID: &str = "BTC/USDT.KRAKEN";
1920
1921    fn params_with(key: &str, val: serde_json::Value) -> Params {
1922        let mut map = indexmap::IndexMap::new();
1923        map.insert(key.to_owned(), val);
1924        Params::from_index_map(map)
1925    }
1926
1927    #[tokio::test]
1928    async fn test_cancel_order_for_spot_missing_cached_instrument_returns_canonical_error() {
1929        let http = KrakenSpotHttpClient::default();
1930        let instrument_id = InstrumentId::from(TEST_INSTRUMENT_ID);
1931        let client_order_id = ClientOrderId::from("C-001");
1932        let venue_order_id = VenueOrderId::from("V-001");
1933
1934        let result = cancel_order_for_spot(
1935            &http,
1936            AccountId::from("KRAKEN-001"),
1937            instrument_id,
1938            Some(client_order_id),
1939            Some(venue_order_id),
1940        )
1941        .await;
1942
1943        match result {
1944            Err(CommandFailure::NotSent(reason)) => {
1945                assert_eq!(
1946                    reason,
1947                    InstrumentLookupError::not_found(instrument_id).to_string()
1948                );
1949            }
1950            _ => panic!("Expected local validation failure"),
1951        }
1952    }
1953
1954    #[rstest]
1955    fn test_batch_cancel_item_for_spot_missing_cached_instrument_returns_canonical_error() {
1956        let http = KrakenSpotHttpClient::default();
1957        let instrument_id = InstrumentId::from(TEST_INSTRUMENT_ID);
1958        let cancel = CancelOrder::new(
1959            TraderId::from("TESTER-001"),
1960            None,
1961            StrategyId::from("S-001"),
1962            instrument_id,
1963            ClientOrderId::from("C-001"),
1964            Some(VenueOrderId::from("V-001")),
1965            UUID4::new(),
1966            UnixNanos::default(),
1967            None,
1968            None,
1969        );
1970
1971        let result = batch_cancel_item_for_spot(&http, &cancel);
1972
1973        match result {
1974            Err(CommandFailure::NotSent(reason)) => {
1975                assert_eq!(
1976                    reason,
1977                    InstrumentLookupError::not_found(instrument_id).to_string()
1978                );
1979            }
1980            _ => panic!("Expected local validation failure"),
1981        }
1982    }
1983
1984    #[rstest]
1985    fn test_resolve_leverage_absent_uses_default() {
1986        let p = params_with("other", json!(1));
1987        assert_eq!(resolve_leverage(Some(&p), Some(3)).unwrap(), Some(3));
1988        assert_eq!(resolve_leverage(None, Some(5)).unwrap(), Some(5));
1989        assert_eq!(resolve_leverage(None, None).unwrap(), None);
1990    }
1991
1992    #[rstest]
1993    fn test_resolve_leverage_valid_integer() {
1994        let p = params_with("leverage", json!(5u64));
1995        assert_eq!(resolve_leverage(Some(&p), Some(3)).unwrap(), Some(5));
1996    }
1997
1998    #[rstest]
1999    fn test_resolve_leverage_string_value_errors() {
2000        let p = params_with("leverage", json!("5"));
2001        let err = resolve_leverage(Some(&p), Some(3)).unwrap_err();
2002        assert!(err.contains("Invalid leverage param"), "unexpected: {err}");
2003    }
2004
2005    #[rstest]
2006    fn test_resolve_leverage_overflow_errors() {
2007        let p = params_with("leverage", json!(65539u64));
2008        let err = resolve_leverage(Some(&p), None).unwrap_err();
2009        assert!(err.contains("exceeds maximum"), "unexpected: {err}");
2010    }
2011
2012    #[rstest]
2013    fn test_resolve_use_ws_trade_absent_uses_default() {
2014        let p = params_with("other", json!(1));
2015        assert!(resolve_use_ws_trade(Some(&p), true));
2016        assert!(!resolve_use_ws_trade(Some(&p), false));
2017        assert!(resolve_use_ws_trade(None, true));
2018        assert!(!resolve_use_ws_trade(None, false));
2019    }
2020
2021    #[rstest]
2022    fn test_resolve_use_ws_trade_overrides_default() {
2023        let p_false = params_with("use_ws_trade", json!(false));
2024        let p_true = params_with("use_ws_trade", json!(true));
2025        assert!(!resolve_use_ws_trade(Some(&p_false), true));
2026        assert!(resolve_use_ws_trade(Some(&p_true), false));
2027    }
2028
2029    #[rstest]
2030    fn test_resolve_use_ws_trade_non_boolean_falls_back_to_default() {
2031        let p = params_with("use_ws_trade", json!("true"));
2032        assert!(resolve_use_ws_trade(Some(&p), true));
2033        assert!(!resolve_use_ws_trade(Some(&p), false));
2034    }
2035
2036    #[rstest]
2037    fn test_execution_client_constructs_with_ws_trade_enabled() {
2038        let factory = KrakenExecutionClientFactory::new();
2039        let config = KrakenExecutionClientConfig {
2040            product_type: KrakenProductType::Spot,
2041            use_ws_trade: true,
2042            ws_request_timeout_secs: 7,
2043            ..Default::default()
2044        };
2045        let cache = Rc::new(RefCell::new(Cache::default()));
2046        let _clock = Rc::new(RefCell::new(TestClock::new()));
2047
2048        let result = factory.create(
2049            TraderId::from("TRADER-001"),
2050            "KRAKEN-WS",
2051            &config,
2052            cache.into(),
2053        );
2054        assert!(result.is_ok(), "construction failed: {:?}", result.err());
2055    }
2056}