1use 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#[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 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 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 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 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 #[must_use]
257 pub fn clock(&self) -> &'static AtomicTime {
258 self.clock
259 }
260
261 #[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 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 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 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 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(¶ms)
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(¶ms).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
1875fn 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(¶ms).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}