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