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